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 }