514 lines
12 KiB
Go
514 lines
12 KiB
Go
package web
|
|||
|
|
|
||
|
|
import (
|
||
|
|
"reflect"
|
||
|
|
"sort"
|
||
|
|
"strconv"
|
||
|
|
"sync"
|
||
|
|
"sync/atomic"
|
||
|
|
"time"
|
||
|
|
)
|
||
|
|
|
||
|
|
// PropMap is a set of properties produced by a single analyzer.
|
||
|
|
type PropMap = map[string]interface{}
|
||
|
|
|
||
|
|
// Props is the combined property map of all analyzers of a stream.
|
||
|
|
type Props = map[string]PropMap
|
||
|
|
|
||
|
|
// Event kinds.
|
||
|
|
const (
|
||
|
|
EventKindAction = "action"
|
||
|
|
EventKindLog = "log"
|
||
|
|
EventKindError = "error"
|
||
|
|
)
|
||
|
|
|
||
|
|
// Event is a single entry of the live event feed shown by the web UI.
|
||
|
|
type Event struct {
|
||
|
|
Seq uint64 `json:"seq"`
|
||
|
|
Time int64 `json:"time"` // Unix milliseconds
|
||
|
|
Kind string `json:"kind"`
|
||
|
|
ID int64 `json:"id"`
|
||
|
|
Proto string `json:"proto"`
|
||
|
|
SrcIP string `json:"srcIP"`
|
||
|
|
SrcPort uint16 `json:"srcPort"`
|
||
|
|
DstIP string `json:"dstIP"`
|
||
|
|
DstPort uint16 `json:"dstPort"`
|
||
|
|
Action string `json:"action,omitempty"`
|
||
|
|
Rule string `json:"rule,omitempty"`
|
||
|
|
Host string `json:"host,omitempty"`
|
||
|
|
Message string `json:"message,omitempty"`
|
||
|
|
Props Props `json:"props,omitempty"`
|
||
|
|
}
|
||
|
|
|
||
|
|
// Bucket is a single time slot of the traffic time series.
|
||
|
|
type Bucket struct {
|
||
|
|
Time int64 `json:"time"` // Unix milliseconds, start of the bucket
|
||
|
|
TCP uint64 `json:"tcp"`
|
||
|
|
UDP uint64 `json:"udp"`
|
||
|
|
Blocked uint64 `json:"blocked"`
|
||
|
|
Allowed uint64 `json:"allowed"`
|
||
|
|
}
|
||
|
|
|
||
|
|
// NameCount is a name/count pair used by the various "top N" lists.
|
||
|
|
type NameCount struct {
|
||
|
|
Name string `json:"name"`
|
||
|
|
Count uint64 `json:"count"`
|
||
|
|
}
|
||
|
|
|
||
|
|
const (
|
||
|
|
defaultEventBufferSize = 512
|
||
|
|
subscriberQueueSize = 256
|
||
|
|
maxCounterEntries = 4096
|
||
|
|
topNLimit = 10
|
||
|
|
)
|
||
|
|
|
||
|
|
// Hub collects statistics from the engine and fans live events out to the
|
||
|
|
// connected web UI clients. All methods are safe for concurrent use, and the
|
||
|
|
// ones called from the engine hot path are kept as cheap as possible.
|
||
|
|
type Hub struct {
|
||
|
|
startTime time.Time
|
||
|
|
|
||
|
|
// Hot path counters. Updated with atomics, no locking involved.
|
||
|
|
tcpStreams atomic.Uint64
|
||
|
|
udpStreams atomic.Uint64
|
||
|
|
actAllow atomic.Uint64
|
||
|
|
actBlock atomic.Uint64
|
||
|
|
actDrop atomic.Uint64
|
||
|
|
actModify atomic.Uint64
|
||
|
|
actMaybe atomic.Uint64
|
||
|
|
errors atomic.Uint64
|
||
|
|
workersAlive atomic.Int64
|
||
|
|
|
||
|
|
mu sync.Mutex
|
||
|
|
seq uint64
|
||
|
|
events []*Event // ring buffer, oldest first once wrapped
|
||
|
|
evHead int
|
||
|
|
evFilled bool
|
||
|
|
|
||
|
|
short *series // 10 seconds per bucket, 10 minutes total
|
||
|
|
long *series // 1 minute per bucket, 1 hour total
|
||
|
|
|
||
|
|
hosts *counterMap
|
||
|
|
dstIPs *counterMap
|
||
|
|
rules *counterMap
|
||
|
|
analyzers *counterMap
|
||
|
|
|
||
|
|
subs map[int]chan *Event
|
||
|
|
next int
|
||
|
|
}
|
||
|
|
|
||
|
|
// NewHub creates a new statistics hub.
|
||
|
|
func NewHub() *Hub {
|
||
|
|
return &Hub{
|
||
|
|
startTime: time.Now(),
|
||
|
|
events: make([]*Event, defaultEventBufferSize),
|
||
|
|
short: newSeries(10*time.Second, 60),
|
||
|
|
long: newSeries(time.Minute, 60),
|
||
|
|
hosts: newCounterMap(),
|
||
|
|
dstIPs: newCounterMap(),
|
||
|
|
rules: newCounterMap(),
|
||
|
|
analyzers: newCounterMap(),
|
||
|
|
subs: make(map[int]chan *Event),
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// WorkerStarted / WorkerStopped track the number of running workers.
|
||
|
|
func (h *Hub) WorkerStarted() { h.workersAlive.Add(1) }
|
||
|
|
func (h *Hub) WorkerStopped() { h.workersAlive.Add(-1) }
|
||
|
|
|
||
|
|
// StreamNew records a newly seen stream.
|
||
|
|
func (h *Hub) StreamNew(proto string) {
|
||
|
|
if proto == "udp" {
|
||
|
|
h.udpStreams.Add(1)
|
||
|
|
} else {
|
||
|
|
h.tcpStreams.Add(1)
|
||
|
|
}
|
||
|
|
h.mu.Lock()
|
||
|
|
now := time.Now()
|
||
|
|
h.short.add(now, proto, "")
|
||
|
|
h.long.add(now, proto, "")
|
||
|
|
h.mu.Unlock()
|
||
|
|
}
|
||
|
|
|
||
|
|
// PropUpdate records the analyzer properties of a stream. It is used to build
|
||
|
|
// the "top hosts" list and the per-analyzer hit counters.
|
||
|
|
func (h *Hub) PropUpdate(props Props) {
|
||
|
|
if len(props) == 0 {
|
||
|
|
return
|
||
|
|
}
|
||
|
|
host := HostOf(props)
|
||
|
|
h.mu.Lock()
|
||
|
|
for name := range props {
|
||
|
|
h.analyzers.inc(name)
|
||
|
|
}
|
||
|
|
if host != "" {
|
||
|
|
h.hosts.inc(host)
|
||
|
|
}
|
||
|
|
h.mu.Unlock()
|
||
|
|
}
|
||
|
|
|
||
|
|
// StreamAction records the verdict of a stream and, unless the verdict is the
|
||
|
|
// neutral "maybe", publishes an event to the live feed.
|
||
|
|
func (h *Hub) StreamAction(info StreamInfo, action string) {
|
||
|
|
switch action {
|
||
|
|
case "allow":
|
||
|
|
h.actAllow.Add(1)
|
||
|
|
case "block":
|
||
|
|
h.actBlock.Add(1)
|
||
|
|
case "drop":
|
||
|
|
h.actDrop.Add(1)
|
||
|
|
case "modify":
|
||
|
|
h.actModify.Add(1)
|
||
|
|
default:
|
||
|
|
h.actMaybe.Add(1)
|
||
|
|
return
|
||
|
|
}
|
||
|
|
now := time.Now()
|
||
|
|
h.mu.Lock()
|
||
|
|
h.short.add(now, "", action)
|
||
|
|
h.long.add(now, "", action)
|
||
|
|
if action == "block" || action == "drop" {
|
||
|
|
if info.DstIP != "" {
|
||
|
|
h.dstIPs.inc(info.DstIP)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
h.mu.Unlock()
|
||
|
|
h.publish(now, EventKindAction, info, action, "", "")
|
||
|
|
}
|
||
|
|
|
||
|
|
// RuleLog records a `log: true` hit of a ruleset rule.
|
||
|
|
func (h *Hub) RuleLog(info StreamInfo, rule string) {
|
||
|
|
h.mu.Lock()
|
||
|
|
h.rules.inc(rule)
|
||
|
|
h.mu.Unlock()
|
||
|
|
h.publish(time.Now(), EventKindLog, info, "", rule, "")
|
||
|
|
}
|
||
|
|
|
||
|
|
// Error records an engine or ruleset error.
|
||
|
|
func (h *Hub) Error(info StreamInfo, rule, message string) {
|
||
|
|
h.errors.Add(1)
|
||
|
|
h.publish(time.Now(), EventKindError, info, "", rule, message)
|
||
|
|
}
|
||
|
|
|
||
|
|
// StreamInfo is the subset of ruleset.StreamInfo the hub cares about.
|
||
|
|
type StreamInfo struct {
|
||
|
|
ID int64
|
||
|
|
Proto string
|
||
|
|
SrcIP string
|
||
|
|
SrcPort uint16
|
||
|
|
DstIP string
|
||
|
|
DstPort uint16
|
||
|
|
Props Props
|
||
|
|
}
|
||
|
|
|
||
|
|
func (h *Hub) publish(now time.Time, kind string, info StreamInfo, action, rule, message string) {
|
||
|
|
h.mu.Lock()
|
||
|
|
h.seq++
|
||
|
|
ev := &Event{
|
||
|
|
Seq: h.seq,
|
||
|
|
Time: now.UnixMilli(),
|
||
|
|
Kind: kind,
|
||
|
|
ID: info.ID,
|
||
|
|
Proto: info.Proto,
|
||
|
|
SrcIP: info.SrcIP,
|
||
|
|
SrcPort: info.SrcPort,
|
||
|
|
DstIP: info.DstIP,
|
||
|
|
DstPort: info.DstPort,
|
||
|
|
Action: action,
|
||
|
|
Rule: rule,
|
||
|
|
Message: message,
|
||
|
|
Host: HostOf(info.Props),
|
||
|
|
Props: info.Props,
|
||
|
|
}
|
||
|
|
h.events[h.evHead] = ev
|
||
|
|
h.evHead = (h.evHead + 1) % len(h.events)
|
||
|
|
if h.evHead == 0 {
|
||
|
|
h.evFilled = true
|
||
|
|
}
|
||
|
|
for _, ch := range h.subs {
|
||
|
|
select {
|
||
|
|
case ch <- ev:
|
||
|
|
default: // Slow client, drop the event rather than blocking the engine
|
||
|
|
}
|
||
|
|
}
|
||
|
|
h.mu.Unlock()
|
||
|
|
}
|
||
|
|
|
||
|
|
// Subscribe returns a channel of live events and a function to unsubscribe.
|
||
|
|
func (h *Hub) Subscribe() (<-chan *Event, func()) {
|
||
|
|
ch := make(chan *Event, subscriberQueueSize)
|
||
|
|
h.mu.Lock()
|
||
|
|
id := h.next
|
||
|
|
h.next++
|
||
|
|
h.subs[id] = ch
|
||
|
|
h.mu.Unlock()
|
||
|
|
return ch, func() {
|
||
|
|
h.mu.Lock()
|
||
|
|
if c, ok := h.subs[id]; ok {
|
||
|
|
delete(h.subs, id)
|
||
|
|
close(c)
|
||
|
|
}
|
||
|
|
h.mu.Unlock()
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// Events returns up to limit most recent events, newest first.
|
||
|
|
func (h *Hub) Events(limit int) []*Event {
|
||
|
|
h.mu.Lock()
|
||
|
|
defer h.mu.Unlock()
|
||
|
|
size := len(h.events)
|
||
|
|
n := h.evHead
|
||
|
|
if h.evFilled {
|
||
|
|
n = size
|
||
|
|
}
|
||
|
|
if limit <= 0 || limit > n {
|
||
|
|
limit = n
|
||
|
|
}
|
||
|
|
out := make([]*Event, 0, limit)
|
||
|
|
for i := 0; i < limit; i++ {
|
||
|
|
idx := (h.evHead - 1 - i + size*2) % size
|
||
|
|
if h.events[idx] == nil {
|
||
|
|
break
|
||
|
|
}
|
||
|
|
out = append(out, h.events[idx])
|
||
|
|
}
|
||
|
|
return out
|
||
|
|
}
|
||
|
|
|
||
|
|
// Metrics is the snapshot of everything the dashboard needs.
|
||
|
|
type Metrics struct {
|
||
|
|
UptimeSeconds int64 `json:"uptimeSeconds"`
|
||
|
|
Workers int64 `json:"workers"`
|
||
|
|
Streams StreamStats `json:"streams"`
|
||
|
|
Actions ActionStats `json:"actions"`
|
||
|
|
Errors uint64 `json:"errors"`
|
||
|
|
Short []Bucket `json:"short"`
|
||
|
|
Long []Bucket `json:"long"`
|
||
|
|
TopHosts []NameCount `json:"topHosts"`
|
||
|
|
TopBlockedIPs []NameCount `json:"topBlockedIPs"`
|
||
|
|
TopRules []NameCount `json:"topRules"`
|
||
|
|
Analyzers []NameCount `json:"analyzers"`
|
||
|
|
}
|
||
|
|
|
||
|
|
type StreamStats struct {
|
||
|
|
TCP uint64 `json:"tcp"`
|
||
|
|
UDP uint64 `json:"udp"`
|
||
|
|
Total uint64 `json:"total"`
|
||
|
|
}
|
||
|
|
|
||
|
|
type ActionStats struct {
|
||
|
|
Allow uint64 `json:"allow"`
|
||
|
|
Block uint64 `json:"block"`
|
||
|
|
Drop uint64 `json:"drop"`
|
||
|
|
Modify uint64 `json:"modify"`
|
||
|
|
Maybe uint64 `json:"maybe"`
|
||
|
|
}
|
||
|
|
|
||
|
|
// Metrics returns a consistent snapshot of the collected statistics.
|
||
|
|
func (h *Hub) Metrics() Metrics {
|
||
|
|
tcp, udp := h.tcpStreams.Load(), h.udpStreams.Load()
|
||
|
|
m := Metrics{
|
||
|
|
UptimeSeconds: int64(time.Since(h.startTime).Seconds()),
|
||
|
|
Workers: h.workersAlive.Load(),
|
||
|
|
Streams: StreamStats{TCP: tcp, UDP: udp, Total: tcp + udp},
|
||
|
|
Actions: ActionStats{
|
||
|
|
Allow: h.actAllow.Load(),
|
||
|
|
Block: h.actBlock.Load(),
|
||
|
|
Drop: h.actDrop.Load(),
|
||
|
|
Modify: h.actModify.Load(),
|
||
|
|
Maybe: h.actMaybe.Load(),
|
||
|
|
},
|
||
|
|
Errors: h.errors.Load(),
|
||
|
|
}
|
||
|
|
now := time.Now()
|
||
|
|
h.mu.Lock()
|
||
|
|
m.Short = h.short.snapshot(now)
|
||
|
|
m.Long = h.long.snapshot(now)
|
||
|
|
m.TopHosts = h.hosts.top(topNLimit)
|
||
|
|
m.TopBlockedIPs = h.dstIPs.top(topNLimit)
|
||
|
|
m.TopRules = h.rules.top(topNLimit)
|
||
|
|
m.Analyzers = h.analyzers.top(maxCounterEntries)
|
||
|
|
h.mu.Unlock()
|
||
|
|
return m
|
||
|
|
}
|
||
|
|
|
||
|
|
// HostOf extracts the most meaningful hostname out of a property map, so that
|
||
|
|
// the UI can show something more useful than a bare IP address.
|
||
|
|
func HostOf(props Props) string {
|
||
|
|
if props == nil {
|
||
|
|
return ""
|
||
|
|
}
|
||
|
|
for _, path := range [][]string{
|
||
|
|
{"tls", "req", "sni"},
|
||
|
|
{"quic", "req", "sni"},
|
||
|
|
{"http", "req", "headers", "host"},
|
||
|
|
{"dns", "questions", "0", "name"},
|
||
|
|
} {
|
||
|
|
if s := nestedString(props, path...); s != "" {
|
||
|
|
return s
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return ""
|
||
|
|
}
|
||
|
|
|
||
|
|
// nestedString walks a path of keys (or numeric slice indices) and returns the
|
||
|
|
// value at the end of it if it is a string. Analyzers use named map types, so
|
||
|
|
// the lookups go through reflection instead of plain type assertions.
|
||
|
|
func nestedString(root interface{}, path ...string) string {
|
||
|
|
cur := root
|
||
|
|
for _, key := range path {
|
||
|
|
next, ok := index(cur, key)
|
||
|
|
if !ok {
|
||
|
|
return ""
|
||
|
|
}
|
||
|
|
cur = next
|
||
|
|
}
|
||
|
|
s, _ := cur.(string)
|
||
|
|
return s
|
||
|
|
}
|
||
|
|
|
||
|
|
func index(container interface{}, key string) (interface{}, bool) {
|
||
|
|
if container == nil {
|
||
|
|
return nil, false
|
||
|
|
}
|
||
|
|
if m, ok := container.(map[string]interface{}); ok {
|
||
|
|
v, ok := m[key]
|
||
|
|
return v, ok
|
||
|
|
}
|
||
|
|
v := reflect.ValueOf(container)
|
||
|
|
switch v.Kind() {
|
||
|
|
case reflect.Map:
|
||
|
|
if v.Type().Key().Kind() != reflect.String {
|
||
|
|
return nil, false
|
||
|
|
}
|
||
|
|
mv := v.MapIndex(reflect.ValueOf(key).Convert(v.Type().Key()))
|
||
|
|
if !mv.IsValid() {
|
||
|
|
return nil, false
|
||
|
|
}
|
||
|
|
return mv.Interface(), true
|
||
|
|
case reflect.Slice, reflect.Array:
|
||
|
|
i, err := strconv.Atoi(key)
|
||
|
|
if err != nil || i < 0 || i >= v.Len() {
|
||
|
|
return nil, false
|
||
|
|
}
|
||
|
|
return v.Index(i).Interface(), true
|
||
|
|
default:
|
||
|
|
return nil, false
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// series is a fixed size ring of time buckets.
|
||
|
|
type series struct {
|
||
|
|
dur time.Duration
|
||
|
|
buckets []Bucket
|
||
|
|
head int // index of the newest bucket
|
||
|
|
}
|
||
|
|
|
||
|
|
func newSeries(dur time.Duration, size int) *series {
|
||
|
|
s := &series{dur: dur, buckets: make([]Bucket, size)}
|
||
|
|
now := time.Now()
|
||
|
|
for i := range s.buckets {
|
||
|
|
s.buckets[i].Time = s.align(now.Add(-time.Duration(len(s.buckets)-1-i) * dur))
|
||
|
|
}
|
||
|
|
s.head = len(s.buckets) - 1
|
||
|
|
return s
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *series) align(t time.Time) int64 {
|
||
|
|
return t.Truncate(s.dur).UnixMilli()
|
||
|
|
}
|
||
|
|
|
||
|
|
// rotate advances the ring so that the newest bucket matches now.
|
||
|
|
func (s *series) rotate(now time.Time) *Bucket {
|
||
|
|
ts := s.align(now)
|
||
|
|
head := &s.buckets[s.head]
|
||
|
|
if head.Time == ts {
|
||
|
|
return head
|
||
|
|
}
|
||
|
|
if ts < head.Time {
|
||
|
|
// Clock went backwards, keep using the current bucket.
|
||
|
|
return head
|
||
|
|
}
|
||
|
|
steps := int((ts - head.Time) / s.dur.Milliseconds())
|
||
|
|
if steps > len(s.buckets) {
|
||
|
|
steps = len(s.buckets)
|
||
|
|
}
|
||
|
|
for i := 0; i < steps; i++ {
|
||
|
|
s.head = (s.head + 1) % len(s.buckets)
|
||
|
|
s.buckets[s.head] = Bucket{Time: head.Time + int64(i+1)*s.dur.Milliseconds()}
|
||
|
|
}
|
||
|
|
s.buckets[s.head].Time = ts
|
||
|
|
return &s.buckets[s.head]
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *series) add(now time.Time, proto, action string) {
|
||
|
|
b := s.rotate(now)
|
||
|
|
switch proto {
|
||
|
|
case "tcp":
|
||
|
|
b.TCP++
|
||
|
|
case "udp":
|
||
|
|
b.UDP++
|
||
|
|
}
|
||
|
|
switch action {
|
||
|
|
case "block", "drop":
|
||
|
|
b.Blocked++
|
||
|
|
case "allow":
|
||
|
|
b.Allowed++
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// snapshot returns the buckets in chronological order.
|
||
|
|
func (s *series) snapshot(now time.Time) []Bucket {
|
||
|
|
s.rotate(now)
|
||
|
|
out := make([]Bucket, 0, len(s.buckets))
|
||
|
|
for i := 1; i <= len(s.buckets); i++ {
|
||
|
|
out = append(out, s.buckets[(s.head+i)%len(s.buckets)])
|
||
|
|
}
|
||
|
|
return out
|
||
|
|
}
|
||
|
|
|
||
|
|
// counterMap is a bounded string counter used for the "top N" lists.
|
||
|
|
type counterMap struct {
|
||
|
|
m map[string]uint64
|
||
|
|
}
|
||
|
|
|
||
|
|
func newCounterMap() *counterMap {
|
||
|
|
return &counterMap{m: make(map[string]uint64)}
|
||
|
|
}
|
||
|
|
|
||
|
|
func (c *counterMap) inc(key string) {
|
||
|
|
if len(c.m) >= maxCounterEntries {
|
||
|
|
if _, ok := c.m[key]; !ok {
|
||
|
|
c.prune()
|
||
|
|
}
|
||
|
|
}
|
||
|
|
c.m[key]++
|
||
|
|
}
|
||
|
|
|
||
|
|
// prune keeps the most frequent half of the entries.
|
||
|
|
func (c *counterMap) prune() {
|
||
|
|
keep := c.top(maxCounterEntries / 2)
|
||
|
|
m := make(map[string]uint64, len(keep))
|
||
|
|
for _, kv := range keep {
|
||
|
|
m[kv.Name] = kv.Count
|
||
|
|
}
|
||
|
|
c.m = m
|
||
|
|
}
|
||
|
|
|
||
|
|
func (c *counterMap) top(n int) []NameCount {
|
||
|
|
out := make([]NameCount, 0, len(c.m))
|
||
|
|
for k, v := range c.m {
|
||
|
|
out = append(out, NameCount{Name: k, Count: v})
|
||
|
|
}
|
||
|
|
sort.Slice(out, func(i, j int) bool {
|
||
|
|
if out[i].Count != out[j].Count {
|
||
|
|
return out[i].Count > out[j].Count
|
||
|
|
}
|
||
|
|
return out[i].Name < out[j].Name
|
||
|
|
})
|
||
|
|
if len(out) > n {
|
||
|
|
out = out[:n]
|
||
|
|
}
|
||
|
|
return out
|
||
|
|
}
|