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>
This commit is contained in:
SaarMesh-Bot
2026-09-02 11:07:08 +02:00
committed by GitHub
co-authored by SaarMesh-Bot Claude
parent 1441734991
commit 1720060284
5 changed files with 149 additions and 19 deletions
+15 -10
View File
@@ -10,10 +10,15 @@
// (see #1481 P0-2). Helpers do NOT touch store.mu — the handler owns lock
// scoping.
//
// Concurrency note (#1839 MINOR): the RLock snapshot only covers the
// *StoreObs pointer slice; reads of store.byTxID inside buildPacketTypes /
// buildNodesTimeline are unsynchronized concurrent-map reads (pre-existing
// behavior — enrichObs did the same). Not introduced by this refactor.
// 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
@@ -93,10 +98,10 @@ func buildTimeline(filtered []*StoreObs, days int) []TimeBucket {
// 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(store *PacketStore, filtered []*StoreObs) map[string]int {
func buildPacketTypes(filtered []*StoreObs, txByID map[int]*StoreTx) map[string]int {
out := map[string]int{}
for _, obs := range filtered {
tx := store.byTxID[obs.TransmissionID]
tx := txByID[obs.TransmissionID]
if tx == nil || tx.PayloadType == nil {
continue
}
@@ -107,7 +112,7 @@ func buildPacketTypes(store *PacketStore, filtered []*StoreObs) map[string]int {
// buildNodesTimeline builds the distinct-node-per-bucket timeline aggregate.
// Nodes = path-json hops decoded_json pubKey/srcHash/destHash.
func buildNodesTimeline(store *PacketStore, filtered []*StoreObs, days int) []TimeBucket {
func buildNodesTimeline(filtered []*StoreObs, days int, txByID map[int]*StoreTx) []TimeBucket {
bucketDur := observerAnalyticsBucketDur(days)
nodeBucketSets := map[int64]map[string]struct{}{}
for _, obs := range filtered {
@@ -121,7 +126,7 @@ func buildNodesTimeline(store *PacketStore, filtered []*StoreObs, days int) []Ti
}
// Legacy handler read decoded_json via enrichObs (which pulls it off
// tx.DecodedJSON). Read tx directly for parity + savings.
if tx := store.byTxID[obs.TransmissionID]; tx != nil && tx.DecodedJSON != "" {
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"} {
@@ -189,7 +194,7 @@ func buildSnrDistribution(filtered []*StoreObs) []SnrDistributionEntry {
// 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) []map[string]interface{} {
func buildRecentPackets(store *PacketStore, filtered []*StoreObs, limit int, txByID map[int]*StoreTx) []map[string]interface{} {
if limit <= 0 {
return []map[string]interface{}{}
}
@@ -201,7 +206,7 @@ func buildRecentPackets(store *PacketStore, filtered []*StoreObs, limit int) []m
if _, ok := obs.ParsedTime(); !ok {
continue
}
out = append(out, store.enrichObs(obs))
out = append(out, store.enrichObsWithTx(obs, txByID[obs.TransmissionID]))
}
return out
}
+4 -4
View File
@@ -78,7 +78,7 @@ func TestBuildPacketTypesDirectRead(t *testing.T) {
buildObsForTest(3, now, nil, "[]"),
buildObsForTest(4, now, nil, "[]"),
}
got := buildPacketTypes(store, filtered)
got := buildPacketTypes(filtered, store.byTxID)
if got["3"] != 2 {
t.Errorf("packetTypes[3] = %d, want 2", got["3"])
}
@@ -146,7 +146,7 @@ func TestBuildRecentPacketsLimit(t *testing.T) {
for i := 0; i < 25; i++ {
filtered = append(filtered, buildObsForTest(1, base.Add(-time.Duration(i)*time.Minute), nil, "[]"))
}
got := buildRecentPackets(store, filtered, 20)
got := buildRecentPackets(store, filtered, 20, store.byTxID)
if len(got) != 20 {
t.Errorf("recentPackets len = %d, want 20", len(got))
}
@@ -177,7 +177,7 @@ func TestBuildRecentPacketsSkipsUnparseableTimestamp(t *testing.T) {
filtered = append(filtered, buildObsForTest(1, base.Add(-time.Duration(i)*time.Minute), nil, "[]"))
}
got := buildRecentPackets(store, filtered, 20)
got := buildRecentPackets(store, filtered, 20, store.byTxID)
// Legacy loop: index 0-2 skipped (bad ts), index 3-19 appended under the
// i<20 gate (17 entries), index 20-22 dropped (i>=20). Result len = 17.
@@ -208,7 +208,7 @@ func TestBuildNodesTimelineDistinct(t *testing.T) {
buildObsForTest(1, base, nil, `["bb","cc"]`), // bucket A: nodes {aaaa, bb, cc}
buildObsForTest(1, base.Add(10*time.Minute), nil, `["bb"]`), // same bucket: dedup
}
got := buildNodesTimeline(store, filtered, 1)
got := buildNodesTimeline(filtered, 1, store.byTxID)
if len(got) != 1 {
t.Fatalf("nodes timeline entries = %d, want 1 (got %+v)", len(got), got)
}
+15 -3
View File
@@ -2850,6 +2850,18 @@ func (s *Server) handleObserverAnalytics(w http.ResponseWriter, r *http.Request)
obsList := s.store.byObserver[id]
obsSnapshot := make([]*StoreObs, len(obsList))
copy(obsSnapshot, obsList)
// #1830: also resolve each referenced transmission's *StoreTx under
// this same RLock. s.store.byTxID is guarded by s.store.mu (writes
// from ingest/eviction); reading it after RUnlock — as the loop below
// used to via s.store.byTxID[...] and enrichObs() — races with those
// writers. Keyed by TransmissionID (not by observation index) since
// multiple observations can share one transmission.
txByID := make(map[int]*StoreTx, len(obsSnapshot))
for _, obs := range obsSnapshot {
if _, ok := txByID[obs.TransmissionID]; !ok {
txByID[obs.TransmissionID] = s.store.byTxID[obs.TransmissionID]
}
}
s.store.mu.RUnlock()
filtered := make([]*StoreObs, 0, len(obsSnapshot))
for _, obs := range obsSnapshot {
@@ -2865,10 +2877,10 @@ func (s *Server) handleObserverAnalytics(w http.ResponseWriter, r *http.Request)
writeJSON(w, ObserverAnalyticsResponse{
Timeline: buildTimeline(filtered, days),
PacketTypes: buildPacketTypes(s.store, filtered),
NodesTimeline: buildNodesTimeline(s.store, filtered, days),
PacketTypes: buildPacketTypes(filtered, txByID),
NodesTimeline: buildNodesTimeline(filtered, days, txByID),
SnrDistribution: buildSnrDistribution(filtered),
RecentPackets: buildRecentPackets(s.store, filtered, 20),
RecentPackets: buildRecentPackets(s.store, filtered, 20, txByID),
})
}
+98 -1
View File
@@ -9,6 +9,7 @@ import (
"path/filepath"
"strconv"
"strings"
"sync"
"testing"
"time"
@@ -1031,9 +1032,44 @@ func TestObserverAnalytics(t *testing.T) {
if body["recentPackets"] == nil {
t.Error("expected recentPackets")
}
if recent, ok := body["recentPackets"].([]interface{}); !ok || len(recent) == 0 {
recent, ok := body["recentPackets"].([]interface{})
if !ok || len(recent) == 0 {
t.Errorf("expected non-empty recentPackets, got %v", body["recentPackets"])
}
// #1827: packetTypes must still reflect the real per-transmission
// payload_type after switching the hot loop from enrichObs() to a
// direct s.store.byTxID read. seedTestData gives obs1 three
// observations: tx1 (payload_type=4), tx2 (payload_type=5), tx3
// (payload_type=4) — i.e. counts {"4":2, "5":1}.
pt, ok := body["packetTypes"].(map[string]interface{})
if !ok {
t.Fatalf("expected packetTypes to be an object, got %T", body["packetTypes"])
}
if pt["4"] != float64(2) {
t.Errorf("expected packetTypes[4]=2, got %v", pt["4"])
}
if pt["5"] != float64(1) {
t.Errorf("expected packetTypes[5]=1, got %v", pt["5"])
}
// recentPackets is still built via the unchanged enrichObs() path
// and should retain resolved_path for observations that have one
// (tx1's first observation, resolved via obs1).
foundResolved := false
for _, rp := range recent {
m, ok := rp.(map[string]interface{})
if !ok {
continue
}
if m["resolved_path"] != nil {
foundResolved = true
break
}
}
if !foundResolved {
t.Error("expected at least one recentPackets entry to carry resolved_path")
}
})
t.Run("custom days", func(t *testing.T) {
@@ -1057,6 +1093,67 @@ func TestObserverAnalytics(t *testing.T) {
})
}
// TestObserverAnalytics_ConcurrentByTxIDMutation_1830 is the -race
// regression requested in #1830: handleObserverAnalytics reads
// s.store.byTxID (directly, and via enrichObsWithTx) for every observation
// in the window. Before the #1830 fix those reads happened after the
// snapshot RLock was released, racing with concurrent ingest/eviction
// writers to that same map. This test hammers the endpoint concurrently
// with goroutines that mutate byTxID exactly like ingest/eviction would,
// and must pass under `go test -race`.
func TestObserverAnalytics_ConcurrentByTxIDMutation_1830(t *testing.T) {
srv, router := setupTestServer(t)
stop := make(chan struct{})
var writers sync.WaitGroup
var readers sync.WaitGroup
// Writer goroutines: mutate byTxID under s.store.mu, like ingest
// (adds) and eviction (deletes) do in production. They run until
// stop is closed, which happens once all readers are done below.
for w := 0; w < 4; w++ {
writers.Add(1)
go func(base int) {
defer writers.Done()
n := 0
for {
select {
case <-stop:
return
default:
}
srv.store.mu.Lock()
id := base*100000 + n
srv.store.byTxID[id] = &StoreTx{ID: id, PayloadType: intPtr(4)}
delete(srv.store.byTxID, id-1)
srv.store.mu.Unlock()
n++
}
}(w)
}
// Reader goroutines: repeatedly hit the endpoint under test.
for r := 0; r < 8; r++ {
readers.Add(1)
go func() {
defer readers.Done()
for i := 0; i < 50; i++ {
req := httptest.NewRequest("GET", "/api/observers/obs1/analytics", nil)
w := httptest.NewRecorder()
router.ServeHTTP(w, req)
if w.Code != 200 {
t.Errorf("expected 200, got %d", w.Code)
return
}
}
}()
}
readers.Wait()
close(stop)
writers.Wait()
}
func TestChannelMessages(t *testing.T) {
_, router := setupTestServer(t)
req := httptest.NewRequest("GET", "/api/channels/%23test/messages", nil)
+17 -1
View File
@@ -3747,9 +3747,25 @@ func (s *PacketStore) computeNodeHomeRegions() map[string]string {
}
// enrichObs returns a map with observation fields + transmission fields.
// Looks up the transmission in s.byTxID itself — safe only when the caller
// already holds s.mu (directly, or via a defer'd RLock spanning the call).
// Callers that snapshot observations under RLock and then release it before
// iterating (e.g. handleObserverAnalytics, #1830) must use enrichObsWithTx
// with a tx pointer resolved during that same snapshot instead.
func (s *PacketStore) enrichObs(obs *StoreObs) map[string]interface{} {
tx := s.byTxID[obs.TransmissionID]
return s.enrichObsWithTx(obs, s.byTxID[obs.TransmissionID])
}
// enrichObsWithTx is enrichObs with the transmission pointer already
// resolved by the caller, instead of looking it up in s.byTxID here. #1830:
// s.byTxID is guarded by s.mu (writes from ingest/eviction); reading it
// without holding at least RLock races with those writers — Go maps can
// panic with "concurrent map read and map write" during a rehash, not just
// fail under -race. Callers that need to read byTxID after releasing their
// RLock (to keep JSON decode / enrichment off the hot lock, per #1481)
// should resolve the *StoreTx for each observation during their RLock-held
// snapshot and pass it in here.
func (s *PacketStore) enrichObsWithTx(obs *StoreObs, tx *StoreTx) map[string]interface{} {
m := map[string]interface{}{
"id": obs.ID,
"timestamp": strOrNil(obs.Timestamp),