From 6e2f79c0adf82bea0c0d53bf2b9f130cc0f6f3b0 Mon Sep 17 00:00:00 2001 From: Kpa-clawbot Date: Sat, 4 Apr 2026 10:39:04 -0700 Subject: [PATCH] =?UTF-8?q?perf:=20optimize=20QueryGroupedPackets=20?= =?UTF-8?q?=E2=80=94=20cache=20observer=20count,=20defer=20map=20construct?= =?UTF-8?q?ion=20(#594)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Summary Optimizes `QueryGroupedPackets()` in `store.go` to eliminate two major inefficiencies on every grouped packet list request: ### Changes 1. **Cache `UniqueObserverCount` on `StoreTx`** — Instead of iterating all observations to count unique observers on every query (O(total_observations) per request), we now track unique observers at ingest time via an `observerSet` map and pre-computed `UniqueObserverCount` field. This is updated incrementally as observations arrive. 2. **Defer map construction until after pagination** — Previously, `map[string]interface{}` was built for ALL 30K+ filtered results before sorting and paginating. Now the grouped cache stores sorted `[]*StoreTx` pointers (lightweight), and `groupedTxsToPage()` builds maps only for the requested page (typically 50 items). This eliminates ~30K map allocations per cache miss. 3. **Lighter cache footprint** — The grouped cache now stores `[]*StoreTx` instead of `*PacketResult` with pre-built maps, reducing memory pressure and GC work. ### Complexity - Observer counting: O(1) per query (was O(total_observations)) - Map construction: O(page_size) per query (was O(n) where n = all filtered results) - Sort remains O(n log n) on cache miss, but the cache (3s TTL) absorbs repeated requests ### Testing - `cd cmd/server && go test ./...` — all tests pass - `cd cmd/ingestor && go build ./...` — builds clean Fixes #370 --------- Co-authored-by: you --- cmd/server/store.go | 143 +++++++++++++++++++++++++------------------- 1 file changed, 81 insertions(+), 62 deletions(-) diff --git a/cmd/server/store.go b/cmd/server/store.go index 84143b5e..b3d1bcc7 100644 --- a/cmd/server/store.go +++ b/cmd/server/store.go @@ -41,14 +41,16 @@ type StoreTx struct { PathJSON string Direction string ResolvedPath []*string // resolved path from best observation - LatestSeen string // max observation timestamp (or FirstSeen if no observations) + LatestSeen string // max observation timestamp (or FirstSeen if no observations) + UniqueObserverCount int // cached count of distinct observer IDs // Cached parsed fields (set once, read many) parsedPath []string // cached parsePathJSON result pathParsed bool // whether parsedPath has been set decodedOnce sync.Once // guards parsedDecoded parsedDecoded map[string]interface{} // cached json.Unmarshal of DecodedJSON // Dedup map: "observerID|pathJSON" → true for O(1) duplicate checks - obsKeys map[string]bool + obsKeys map[string]bool + observerSet map[string]bool // unique observer IDs (for UniqueObserverCount) } // StoreObs is a lean in-memory observation (no duplication of transmission fields). @@ -114,10 +116,11 @@ type PacketStore struct { pendingInv *cacheInvalidation // accumulated dirty flags during cooldown invCooldown time.Duration // minimum time between invalidations // Short-lived cache for QueryGroupedPackets (avoids repeated full sort) - groupedCacheMu sync.Mutex - groupedCacheKey string - groupedCacheExp time.Time - groupedCacheRes *PacketResult + groupedCacheMu sync.Mutex + groupedCacheKey string + groupedCacheExp time.Time + groupedCacheTxs []*StoreTx // sorted by LatestSeen DESC + groupedCacheTotal int // Short-lived cache for GetChannels (avoids repeated full scan + JSON unmarshal) channelsCacheMu sync.Mutex channelsCacheKey string @@ -305,6 +308,7 @@ func (s *PacketStore) Load() error { PayloadType: nullIntPtr(payloadType), DecodedJSON: nullStrVal(decodedJSON), obsKeys: make(map[string]bool), + observerSet: make(map[string]bool), } s.byHash[hashStr] = tx s.packets = append(s.packets, tx) @@ -347,6 +351,10 @@ func (s *PacketStore) Load() error { tx.Observations = append(tx.Observations, obs) tx.obsKeys[dk] = true + if obs.ObserverID != "" && !tx.observerSet[obs.ObserverID] { + tx.observerSet[obs.ObserverID] = true + tx.UniqueObserverCount++ + } tx.ObservationCount++ if obs.Timestamp > tx.LatestSeen { tx.LatestSeen = obs.Timestamp @@ -561,77 +569,37 @@ func (s *PacketStore) QueryGroupedPackets(q PacketQuery) *PacketResult { // Return cached sorted list if still fresh (3s TTL) s.groupedCacheMu.Lock() - if s.groupedCacheRes != nil && s.groupedCacheKey == cacheKey && time.Now().Before(s.groupedCacheExp) { - cached := s.groupedCacheRes + if s.groupedCacheTxs != nil && s.groupedCacheKey == cacheKey && time.Now().Before(s.groupedCacheExp) { + cachedTxs := s.groupedCacheTxs + cachedTotal := s.groupedCacheTotal s.groupedCacheMu.Unlock() - return pagePacketResult(cached, q.Offset, q.Limit) + return groupedTxsToPage(cachedTxs, cachedTotal, q.Offset, q.Limit) } s.groupedCacheMu.Unlock() - // Build entries under read lock (observer scan needs lock), sort outside it. - type groupEntry struct { - latest map[string]interface{} - ts string - } - var entries []groupEntry - + // Collect StoreTx pointers under read lock; sort outside it. s.mu.RLock() results := s.filterPackets(q) - entries = make([]groupEntry, 0, len(results)) - for _, tx := range results { - observerCount := 0 - seen := make(map[string]bool) - for _, obs := range tx.Observations { - if obs.ObserverID != "" && !seen[obs.ObserverID] { - seen[obs.ObserverID] = true - observerCount++ - } - } - entries = append(entries, groupEntry{ - ts: tx.LatestSeen, - latest: map[string]interface{}{ - "hash": strOrNil(tx.Hash), - "first_seen": strOrNil(tx.FirstSeen), - "count": tx.ObservationCount, - "observer_count": observerCount, - "observation_count": tx.ObservationCount, - "latest": strOrNil(tx.LatestSeen), - "observer_id": strOrNil(tx.ObserverID), - "observer_name": strOrNil(tx.ObserverName), - "path_json": strOrNil(tx.PathJSON), - "payload_type": intPtrOrNil(tx.PayloadType), - "route_type": intPtrOrNil(tx.RouteType), - "raw_hex": strOrNil(tx.RawHex), - "decoded_json": strOrNil(tx.DecodedJSON), - "snr": floatPtrOrNil(tx.SNR), - "rssi": floatPtrOrNil(tx.RSSI), - }, - }) - if tx.ResolvedPath != nil { - entries[len(entries)-1].latest["resolved_path"] = tx.ResolvedPath - } - } + txs := make([]*StoreTx, len(results)) + copy(txs, results) s.mu.RUnlock() - // Sort outside the lock — only touches our local slice. - sort.Slice(entries, func(i, j int) bool { - return entries[i].ts > entries[j].ts + total := len(txs) + + // Full sort by LatestSeen DESC so the cached slice supports all page offsets. + sort.Slice(txs, func(i, j int) bool { + return txs[i].LatestSeen > txs[j].LatestSeen }) - packets := make([]map[string]interface{}, len(entries)) - for i, e := range entries { - packets[i] = e.latest - } - - full := &PacketResult{Packets: packets, Total: len(packets)} - + // Cache the sorted StoreTx slice (not maps) — lightweight and reusable for any page. s.groupedCacheMu.Lock() - s.groupedCacheRes = full + s.groupedCacheTxs = txs + s.groupedCacheTotal = total s.groupedCacheKey = cacheKey s.groupedCacheExp = time.Now().Add(3 * time.Second) s.groupedCacheMu.Unlock() - return pagePacketResult(full, q.Offset, q.Limit) + return groupedTxsToPage(txs, total, q.Offset, q.Limit) } // pagePacketResult returns a window of a PacketResult without re-allocating the slice. @@ -647,6 +615,46 @@ func pagePacketResult(r *PacketResult, offset, limit int) *PacketResult { return &PacketResult{Packets: r.Packets[offset:end], Total: total} } +// groupedTxsToPage builds map representations only for the requested page of sorted StoreTx pointers. +// This avoids allocating maps for all 30K+ transmissions when only 50 are needed. +func groupedTxsToPage(txs []*StoreTx, total, offset, limit int) *PacketResult { + if offset >= len(txs) { + return &PacketResult{Packets: []map[string]interface{}{}, Total: total} + } + end := offset + limit + if end > len(txs) { + end = len(txs) + } + page := txs[offset:end] + + packets := make([]map[string]interface{}, len(page)) + for i, tx := range page { + m := map[string]interface{}{ + "hash": strOrNil(tx.Hash), + "first_seen": strOrNil(tx.FirstSeen), + "count": tx.ObservationCount, + "observer_count": tx.UniqueObserverCount, + "observation_count": tx.ObservationCount, + "latest": strOrNil(tx.LatestSeen), + "observer_id": strOrNil(tx.ObserverID), + "observer_name": strOrNil(tx.ObserverName), + "path_json": strOrNil(tx.PathJSON), + "payload_type": intPtrOrNil(tx.PayloadType), + "route_type": intPtrOrNil(tx.RouteType), + "raw_hex": strOrNil(tx.RawHex), + "decoded_json": strOrNil(tx.DecodedJSON), + "snr": floatPtrOrNil(tx.SNR), + "rssi": floatPtrOrNil(tx.RSSI), + } + if tx.ResolvedPath != nil { + m["resolved_path"] = tx.ResolvedPath + } + packets[i] = m + } + + return &PacketResult{Packets: packets, Total: total} +} + // GetStoreStats returns aggregate counts (packet data from memory, node/observer from DB). func (s *PacketStore) GetStoreStats() (*Stats, error) { s.mu.RLock() @@ -1190,6 +1198,7 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac PayloadType: r.payloadType, DecodedJSON: r.decodedJSON, obsKeys: make(map[string]bool), + observerSet: make(map[string]bool), } s.byHash[r.hash] = tx s.packets = append(s.packets, tx) // oldest-first; new items go to tail @@ -1218,6 +1227,7 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac dk := r.observerID + "|" + r.pathJSON if tx.obsKeys == nil { tx.obsKeys = make(map[string]bool) + tx.observerSet = make(map[string]bool) } if tx.obsKeys[dk] { continue @@ -1244,6 +1254,10 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac tx.Observations = append(tx.Observations, obs) tx.obsKeys[dk] = true + if obs.ObserverID != "" && !tx.observerSet[obs.ObserverID] { + tx.observerSet[obs.ObserverID] = true + tx.UniqueObserverCount++ + } tx.ObservationCount++ if obs.Timestamp > tx.LatestSeen { tx.LatestSeen = obs.Timestamp @@ -1512,6 +1526,7 @@ func (s *PacketStore) IngestNewObservations(sinceObsID, limit int) []map[string] dk := r.observerID + "|" + r.pathJSON if tx.obsKeys == nil { tx.obsKeys = make(map[string]bool) + tx.observerSet = make(map[string]bool) } if tx.obsKeys[dk] { continue @@ -1539,6 +1554,10 @@ func (s *PacketStore) IngestNewObservations(sinceObsID, limit int) []map[string] tx.Observations = append(tx.Observations, obs) tx.obsKeys[dk] = true + if obs.ObserverID != "" && !tx.observerSet[obs.ObserverID] { + tx.observerSet[obs.ObserverID] = true + tx.UniqueObserverCount++ + } tx.ObservationCount++ newObs = append(newObs, obs) if obs.Timestamp > tx.LatestSeen {