From 7abd05ff7f3d20b0cb55d97c7a8244ae10e26d7d Mon Sep 17 00:00:00 2001 From: you Date: Sun, 19 Apr 2026 00:03:13 +0000 Subject: [PATCH] fix(#791): eliminate spTxIndex to save ~280MB at 3.4M subpaths Remove the spTxIndex map (map[string][]*StoreTx) that stored per-subpath transmission pointer slices. This was the largest single heap consumer during startup on databases with many unique subpaths. Replace the only read site (GetSubpathDetail) with an O(packets) scan that matches subpaths on the fly. This is acceptable because the endpoint is only called on drill-down, not listing. Remove all write sites (addTxToSubpathIndexFull, removeTxFromSubpathIndexFull) and collapse them back into the simpler addTxToSubpathIndex/removeTxFromSubpathIndex functions that only maintain counts. --- cmd/server/store.go | 66 +++++++++++++++++++-------------------------- 1 file changed, 28 insertions(+), 38 deletions(-) diff --git a/cmd/server/store.go b/cmd/server/store.go index fc905f1d..52167978 100644 --- a/cmd/server/store.go +++ b/cmd/server/store.go @@ -176,8 +176,7 @@ type PacketStore struct { // Precomputed subpath index: raw comma-joined hops → occurrence count. // Built during Load(), incrementally updated on ingest. Avoids full // packet iteration at query time (O(unique_subpaths) vs O(total_packets)). - spIndex map[string]int // "hop1,hop2" → count - spTxIndex map[string][]*StoreTx // "hop1,hop2" → transmissions containing this subpath + spIndex map[string]int // "hop1,hop2" → count spTotalPaths int // transmissions with paths >= 2 hops // Precomputed distance analytics: hop distances and path totals // computed during Load() and incrementally updated on ingest. @@ -311,7 +310,7 @@ func NewPacketStore(db *DB, cfg *PacketStoreConfig, cacheTTLs ...map[string]inte collisionCacheTTL: 3600 * time.Second, invCooldown: 300 * time.Second, spIndex: make(map[string]int, 4096), - spTxIndex: make(map[string][]*StoreTx, 4096), + advertPubkeys: make(map[string]int), lastSeenTouched: make(map[string]time.Time), clockSkew: NewClockSkewEngine(), @@ -1537,7 +1536,7 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac // Incrementally update precomputed subpath index with new transmissions for _, tx := range broadcastTxs { - if addTxToSubpathIndexFull(s.spIndex, s.spTxIndex, tx) { + if addTxToSubpathIndex(s.spIndex, tx) { s.spTotalPaths++ } addTxToPathHopIndex(s.byPathHop, tx) @@ -1906,7 +1905,7 @@ func (s *PacketStore) IngestNewObservations(sinceObsID, limit int) []map[string] // Temporarily set parsedPath to old hops for removal. saved, savedFlag := tx.parsedPath, tx.pathParsed tx.parsedPath, tx.pathParsed = oldHops, true - if removeTxFromSubpathIndexFull(s.spIndex, s.spTxIndex, tx) { + if removeTxFromSubpathIndex(s.spIndex, tx) { s.spTotalPaths-- } tx.parsedPath, tx.pathParsed = saved, savedFlag @@ -1923,7 +1922,7 @@ func (s *PacketStore) IngestNewObservations(sinceObsID, limit int) []map[string] } // pickBestObservation already set pathParsed=false so // addTxToSubpathIndex will re-parse the new path. - if addTxToSubpathIndexFull(s.spIndex, s.spTxIndex, tx) { + if addTxToSubpathIndex(s.spIndex, tx) { s.spTotalPaths++ } addTxToPathHopIndex(s.byPathHop, tx) @@ -2398,12 +2397,6 @@ func txGetParsedPath(tx *StoreTx) []string { // increments their counts in the index. Returns true if the tx contributed // (path had ≥ 2 hops). func addTxToSubpathIndex(idx map[string]int, tx *StoreTx) bool { - return addTxToSubpathIndexFull(idx, nil, tx) -} - -// addTxToSubpathIndexFull is like addTxToSubpathIndex but also appends -// tx to txIdx for each subpath key (if txIdx is non-nil). -func addTxToSubpathIndexFull(idx map[string]int, txIdx map[string][]*StoreTx, tx *StoreTx) bool { hops := txGetParsedPath(tx) if len(hops) < 2 { return false @@ -2413,9 +2406,6 @@ func addTxToSubpathIndexFull(idx map[string]int, txIdx map[string][]*StoreTx, tx for start := 0; start <= len(hops)-l; start++ { key := strings.ToLower(strings.Join(hops[start:start+l], ",")) idx[key]++ - if txIdx != nil { - txIdx[key] = append(txIdx[key], tx) - } } } return true @@ -2425,12 +2415,6 @@ func addTxToSubpathIndexFull(idx map[string]int, txIdx map[string][]*StoreTx, tx // decrements counts for all raw subpaths of tx. Returns true if the tx // had a path. func removeTxFromSubpathIndex(idx map[string]int, tx *StoreTx) bool { - return removeTxFromSubpathIndexFull(idx, nil, tx) -} - -// removeTxFromSubpathIndexFull is like removeTxFromSubpathIndex but also -// removes tx from txIdx for each subpath key (if txIdx is non-nil). -func removeTxFromSubpathIndexFull(idx map[string]int, txIdx map[string][]*StoreTx, tx *StoreTx) bool { hops := txGetParsedPath(tx) if len(hops) < 2 { return false @@ -2443,18 +2427,6 @@ func removeTxFromSubpathIndexFull(idx map[string]int, txIdx map[string][]*StoreT if idx[key] <= 0 { delete(idx, key) } - if txIdx != nil { - txs := txIdx[key] - for i, t := range txs { - if t == tx { - txIdx[key] = append(txs[:i], txs[i+1:]...) - break - } - } - if len(txIdx[key]) == 0 { - delete(txIdx, key) - } - } } } return true @@ -2464,10 +2436,9 @@ func removeTxFromSubpathIndexFull(idx map[string]int, txIdx map[string][]*StoreT // Must be called with s.mu held. func (s *PacketStore) buildSubpathIndex() { s.spIndex = make(map[string]int, 4096) - s.spTxIndex = make(map[string][]*StoreTx, 4096) s.spTotalPaths = 0 for _, tx := range s.packets { - if addTxToSubpathIndexFull(s.spIndex, s.spTxIndex, tx) { + if addTxToSubpathIndex(s.spIndex, tx) { s.spTotalPaths++ } } @@ -2906,7 +2877,7 @@ func (s *PacketStore) EvictStale() int { } // Remove from subpath index - removeTxFromSubpathIndexFull(s.spIndex, s.spTxIndex, tx) + removeTxFromSubpathIndex(s.spIndex, tx) // Remove from path-hop index removeTxFromPathHopIndex(s.byPathHop, tx) } @@ -7078,8 +7049,27 @@ func (s *PacketStore) GetSubpathDetail(rawHops []string) map[string]interface{} // Build the subpath key the same way the index does (lowercase, comma-joined) spKey := strings.ToLower(strings.Join(rawHops, ",")) - // Direct lookup instead of scanning all packets - matchedTxs := s.spTxIndex[spKey] + // Scan all transmissions for matching subpaths (replaces former spTxIndex + // lookup — the per-tx pointer index was eliminated in #791 to save ~280MB + // at scale). This is O(packets) but only called on drill-down, not listing. + var matchedTxs []*StoreTx + for _, tx := range s.packets { + hops := txGetParsedPath(tx) + if len(hops) < 2 { + continue + } + maxL := min(8, len(hops)) + for l := 2; l <= maxL; l++ { + for start := 0; start <= len(hops)-l; start++ { + key := strings.ToLower(strings.Join(hops[start:start+l], ",")) + if key == spKey { + matchedTxs = append(matchedTxs, tx) + goto nextTx + } + } + } + nextTx: + } hourBuckets := make([]int, 24) var snrSum, rssiSum float64