Files
meshcore-analyzer/cmd/server/observer_analytics.go
T
1720060284 fix(#1827): avoid per-observation SQL fetch in handleObserverAnalytics hot loop (#1829)
## Summary

Fixes the CPU/DoS issue in #1827: observer detail pages were saturating
CPU on busy observers — 6-7 concurrently loaded tabs pegged 12 cores for
seconds, and auto-refresh made it self-sustaining.

## Root cause

`handleObserverAnalytics` iterated every observation in the requested
window and called `enrichObs()` per observation just to read
`payload_type` and `decoded_json` for the `packetTypes`/`nodesTimeline`
aggregates. `enrichObs()` also runs an on-demand SQL `SELECT
resolved_path FROM observations WHERE id=?` (`fetchResolvedPathForObs`)
and builds a full response map — both of which are unused by this
aggregation loop. `resolved_path` is only actually consumed by the
`<=20` kept `recentPackets` entries.

Per the triage in #1827 (@carmack): *"Replacing `enrichObs(obs)` with a
direct `s.store.byTxID[obs.TransmissionID].PayloadType` read (as
sketched in the body) drops a map alloc + interface boxes per obs on the
loop that saturated the operator's 12 cores. Byte-identical output.
That's ~90% of the value."*

This PR implements exactly that fast-path.

## Change

- Aggregate loop (`packetTypes`, `nodesTimeline`): read
`payload_type`/`decoded_json` directly off the transmission via
`s.store.byTxID[obs.TransmissionID]` — no SQL, no per-obs map
allocation.
- `recentPackets` (`<=20` entries): unchanged, still calls `enrichObs()`
since it needs `resolved_path`/`raw_hex`/etc. for display.
- Output is unchanged: `packetTypes`/`nodesTimeline` are computed from
the exact same underlying fields (`tx.PayloadType`, `tx.DecodedJSON`),
just without the O(N) SQL round-trips.

## Scope

This is the concrete hot-path fix from #1827's triage — not the broader
`/api/observers/{id}/analytics` endpoint-split proposal in #1828, which
(per that issue's discussion) is a separate P3 follow-up. #1828's own
triage converged on this same `byTxID` fast-path as "the ground-work
minimum" before any endpoint splitting.

## Testing

- Existing `TestObserverAnalytics` passes unchanged.
- Extended `TestObserverAnalytics/default` to assert `packetTypes`
counts come out correct (`{"4":2,"5":1}` for the seeded fixture) via the
new `byTxID` path, and that `recentPackets` still carries
`resolved_path` where present (confirming the `enrichObs()` path for
those 20 entries is untouched).
- `go build ./...` and `go vet ./...` clean in `cmd/server`.
- Full `go test ./...` in `cmd/server`: passes except 4 pre-existing
test-order-dependent failures in `TestHandleNodePaths_*` (unrelated to
this change — reproduced identically on a fresh, unpatched clone of
`upstream/master`).

---------

Co-authored-by: SaarMesh-Bot <300107934+SaarMesh-Bot@users.noreply.github.com>
Co-authored-by: Claude <noreply@anthropic.com>
2026-09-02 11:07:08 +02:00

213 lines
7.3 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Package main — observer analytics helpers.
//
// #1828 Phase A: extracted from cmd/server/routes.go handleObserverAnalytics
// (routes.go:2819-2953). Split the 5 aggregate builders into standalone
// functions for readability + isolated tests. No behavior change relative to
// the pre-#1828 handler; JSON output is byte-identical.
//
// The `filtered` argument is the time-window-filtered, timestamp-desc sorted
// slice of *StoreObs produced by the handler after snapshotting under RLock
// (see #1481 P0-2). Helpers do NOT touch store.mu — the handler owns lock
// scoping.
//
// Concurrency note (#1830): buildPacketTypes / buildNodesTimeline /
// buildRecentPackets take a txByID map instead of *PacketStore for
// transmission lookups. store.byTxID is guarded by store.mu (writes from
// ingest/eviction); reading it here — after the handler's RLock snapshot has
// already been released, to keep JSON decode/enrichment off the hot lock per
// #1481 P0-2 — would race with those writers (Go maps can panic with
// "concurrent map read and map write" during a rehash, not just fail under
// -race). The handler resolves txByID once, under the same RLock that
// snapshots the observation slice, and passes it down instead.
//
// Perf note (#1839 MINOR): dropping enrichObs on the histogram / nodes-
// timeline paths also eliminates N fetchResolvedPathForObs SQL calls per
// /analytics request (store.go:fetchResolvedPathForObs), not just the alloc/
// boxing win — /analytics SQL-load drops materially.
package main
import (
"encoding/json"
"fmt"
"sort"
"strconv"
"time"
)
// observerAnalyticsBucketDur returns the timeline bucket size for a given day
// range. Mirrors the constant table in the legacy handler.
func observerAnalyticsBucketDur(days int) time.Duration {
bucketDur := 24 * time.Hour
if days <= 1 {
bucketDur = time.Hour
} else if days <= 7 {
bucketDur = 4 * time.Hour
}
return bucketDur
}
// observerAnalyticsFormatLabel returns a closure that formats a bucket time
// as a human-readable label appropriate for the day range.
func observerAnalyticsFormatLabel(days int) func(time.Time) string {
return func(t time.Time) string {
if days <= 1 {
return t.UTC().Format("15:04")
}
if days <= 7 {
return t.UTC().Format("Mon 15:04")
}
return t.UTC().Format("Jan 02")
}
}
// buildTimelineBuckets is the shared kernel used by both buildTimeline and
// buildNodesTimeline. Sorts the count map by bucket-key and returns
// TimeBucket entries with labels formatted per days.
func buildTimelineBuckets(counts map[int64]int, days int) []TimeBucket {
formatLabel := observerAnalyticsFormatLabel(days)
keys := make([]int64, 0, len(counts))
for k := range counts {
keys = append(keys, k)
}
sort.Slice(keys, func(i, j int) bool { return keys[i] < keys[j] })
out := make([]TimeBucket, 0, len(keys))
for _, k := range keys {
lbl := formatLabel(time.Unix(k, 0))
out = append(out, TimeBucket{Label: &lbl, Count: counts[k]})
}
return out
}
// buildTimeline builds the packet-timeline aggregate for /analytics.
func buildTimeline(filtered []*StoreObs, days int) []TimeBucket {
bucketDur := observerAnalyticsBucketDur(days)
counts := map[int64]int{}
for _, obs := range filtered {
ts, ok := obs.ParsedTime()
if !ok {
continue
}
bucketStart := ts.UTC().Truncate(bucketDur).Unix()
counts[bucketStart]++
}
return buildTimelineBuckets(counts, days)
}
// buildPacketTypes builds the payload-type histogram. Uses a direct
// store.byTxID read (see #1828 triage) rather than enrichObs — the loop only
// needs payload_type, which is a single indirection off the transmission.
// This avoids one map allocation + several interface-boxing conversions per
// observation (routes.go:2886 hot path pre-#1828).
func buildPacketTypes(filtered []*StoreObs, txByID map[int]*StoreTx) map[string]int {
out := map[string]int{}
for _, obs := range filtered {
tx := txByID[obs.TransmissionID]
if tx == nil || tx.PayloadType == nil {
continue
}
out[strconv.Itoa(*tx.PayloadType)]++
}
return out
}
// buildNodesTimeline builds the distinct-node-per-bucket timeline aggregate.
// Nodes = path-json hops decoded_json pubKey/srcHash/destHash.
func buildNodesTimeline(filtered []*StoreObs, days int, txByID map[int]*StoreTx) []TimeBucket {
bucketDur := observerAnalyticsBucketDur(days)
nodeBucketSets := map[int64]map[string]struct{}{}
for _, obs := range filtered {
ts, ok := obs.ParsedTime()
if !ok {
continue
}
bucketStart := ts.UTC().Truncate(bucketDur).Unix()
if nodeBucketSets[bucketStart] == nil {
nodeBucketSets[bucketStart] = map[string]struct{}{}
}
// Legacy handler read decoded_json via enrichObs (which pulls it off
// tx.DecodedJSON). Read tx directly for parity + savings.
if tx := txByID[obs.TransmissionID]; tx != nil && tx.DecodedJSON != "" {
var decoded map[string]interface{}
if json.Unmarshal([]byte(tx.DecodedJSON), &decoded) == nil {
for _, k := range []string{"pubKey", "srcHash", "destHash"} {
if v, ok := decoded[k].(string); ok && v != "" {
nodeBucketSets[bucketStart][v] = struct{}{}
}
}
}
}
for _, hop := range parsePathJSON(obs.PathJSON) {
if hop != "" {
nodeBucketSets[bucketStart][hop] = struct{}{}
}
}
}
nodeCounts := make(map[int64]int, len(nodeBucketSets))
for k, nodes := range nodeBucketSets {
nodeCounts[k] = len(nodes)
}
return buildTimelineBuckets(nodeCounts, days)
}
// buildSnrDistribution builds the SNR histogram (2-unit buckets, floor).
func buildSnrDistribution(filtered []*StoreObs) []SnrDistributionEntry {
snrBuckets := map[int]*SnrDistributionEntry{}
for _, obs := range filtered {
if obs.SNR == nil {
continue
}
bucket := int(*obs.SNR) / 2 * 2
if *obs.SNR < 0 && int(*obs.SNR) != bucket {
bucket -= 2
}
if snrBuckets[bucket] == nil {
snrBuckets[bucket] = &SnrDistributionEntry{Range: fmt.Sprintf("%d to %d", bucket, bucket+2)}
}
snrBuckets[bucket].Count++
}
keys := make([]int, 0, len(snrBuckets))
for k := range snrBuckets {
keys = append(keys, k)
}
sort.Ints(keys)
out := make([]SnrDistributionEntry, 0, len(keys))
for _, k := range keys {
out = append(out, *snrBuckets[k])
}
return out
}
// buildRecentPackets builds the "first N enriched observations" list. This is
// the only aggregate that needs the full enrichObs map — recentPackets is a
// UI-facing payload and the extra fields matter here.
//
// Legacy parity (#1839): the pre-#1828 routes.go loop was
//
// for i, obs := range filtered {
// if _, ok := obs.ParsedTime(); !ok { continue }
// ...
// if i < limit { recentPackets = append(..., enriched) }
// }
//
// The `i` in the gate is the RAW slice index — a bad-ts obs at position k<limit
// consumed its slot (via `continue`) and could not be replaced by a later
// good-ts obs. Result can be <limit when bad-ts obs sit in the head of
// `filtered`. We reproduce that exact semantic here to keep output byte-
// identical.
func buildRecentPackets(store *PacketStore, filtered []*StoreObs, limit int, txByID map[int]*StoreTx) []map[string]interface{} {
if limit <= 0 {
return []map[string]interface{}{}
}
out := make([]map[string]interface{}, 0, limit)
for i, obs := range filtered {
if i >= limit {
break
}
if _, ok := obs.ParsedTime(); !ok {
continue
}
out = append(out, store.enrichObsWithTx(obs, txByID[obs.TransmissionID]))
}
return out
}