perf: optimize QueryGroupedPackets — cache observer count, defer map construction (#594)

## 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 <you@example.com>
This commit is contained in:
Kpa-clawbot
2026-04-04 10:39:04 -07:00
committed by GitHub
co-authored by you
parent b0862f7a41
commit 6e2f79c0ad
+81 -62
View File
@@ -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 {