diff --git a/cmd/server/neighbor_graph.go b/cmd/server/neighbor_graph.go index e05bb83c..1bc93bad 100644 --- a/cmd/server/neighbor_graph.go +++ b/cmd/server/neighbor_graph.go @@ -18,7 +18,7 @@ const ( // Time-decay half-life: 7 days. affinityHalfLifeHours = 168.0 // Cache TTL for the built graph. - neighborGraphTTL = 60 * time.Second + neighborGraphTTL = 5 * time.Minute // Auto-resolve confidence: best must be >= this factor × second-best. affinityConfidenceRatio = 3.0 // Minimum observation count to auto-resolve. @@ -130,6 +130,17 @@ func BuildFromStore(store *PacketStore) *NeighborGraph { return BuildFromStoreWithLog(store, false) } +// cachedToLower returns strings.ToLower(s), caching results to avoid +// repeated allocations for the same pubkey string. +func cachedToLower(cache map[string]string, s string) string { + if v, ok := cache[s]; ok { + return v + } + v := strings.ToLower(s) + cache[s] = v + return v +} + // BuildFromStoreWithLog constructs the neighbor graph, optionally logging disambiguation decisions. func BuildFromStoreWithLog(store *PacketStore, enableLog bool) *NeighborGraph { g := NewNeighborGraph() @@ -149,30 +160,27 @@ func BuildFromStoreWithLog(store *PacketStore, enableLog bool) *NeighborGraph { // Use cached nodes+PM (avoids DB call if cache is fresh). _, pm := store.getCachedNodesAndPM() + // Local cache for strings.ToLower — pubkeys are immutable and repeat + // across hundreds of thousands of observations. + lowerCache := make(map[string]string, 256) + // Phase 1: Extract edges from every transmission + observation. for _, tx := range packets { isAdvert := tx.PayloadType != nil && *tx.PayloadType == 4 - fromNode := "" // originator pubkey (from byNode index key) - // Find the originator pubkey — it's the key in store.byNode. - // StoreTx doesn't store from_node directly; we find it via decoded JSON - // or the byNode index. However, iterating byNode is expensive. - // The originator pubkey is in the decoded JSON "from_node" field, - // but parsing JSON per tx is expensive too. - // Actually, let's look at how byNode is keyed. - // Looking at store.go, byNode maps pubkey → transmissions where that - // pubkey is the "from" node. We need the reverse: tx → from_node. - // The from_node is embedded in DecodedJSON. - // For efficiency, let's extract it once. - fromNode = extractFromNode(tx) + fromNode := extractFromNode(tx) + // Pre-compute lowered originator once per tx (not per observation). + fromLower := "" + if fromNode != "" { + fromLower = cachedToLower(lowerCache, fromNode) + } for _, obs := range tx.Observations { path := parsePathJSON(obs.PathJSON) - observerPK := strings.ToLower(obs.ObserverID) + observerPK := cachedToLower(lowerCache, obs.ObserverID) if len(path) == 0 { // Zero-hop - if isAdvert && fromNode != "" { - fromLower := strings.ToLower(fromNode) + if isAdvert && fromLower != "" { if fromLower != observerPK { // self-edge guard g.upsertEdge(fromLower, observerPK, "", observerPK, obs.SNR, parseTimestamp(obs.Timestamp)) } @@ -181,20 +189,19 @@ func BuildFromStoreWithLog(store *PacketStore, enableLog bool) *NeighborGraph { } // Edge 1: originator ↔ path[0] — ADVERTs only - if isAdvert && fromNode != "" { - firstHop := strings.ToLower(path[0]) - fromLower := strings.ToLower(fromNode) + if isAdvert && fromLower != "" { + firstHop := cachedToLower(lowerCache, path[0]) if fromLower != firstHop { // self-edge guard (shouldn't happen but spec says check) candidates := pm.m[firstHop] - g.upsertEdgeWithCandidates(fromLower, firstHop, candidates, observerPK, obs.SNR, parseTimestamp(obs.Timestamp)) + g.upsertEdgeWithCandidates(fromLower, firstHop, candidates, observerPK, obs.SNR, parseTimestamp(obs.Timestamp), lowerCache) } } // Edge 2: observer ↔ path[last] — ALL packet types - lastHop := strings.ToLower(path[len(path)-1]) + lastHop := cachedToLower(lowerCache, path[len(path)-1]) if observerPK != lastHop { // self-edge guard candidates := pm.m[lastHop] - g.upsertEdgeWithCandidates(observerPK, lastHop, candidates, observerPK, obs.SNR, parseTimestamp(obs.Timestamp)) + g.upsertEdgeWithCandidates(observerPK, lastHop, candidates, observerPK, obs.SNR, parseTimestamp(obs.Timestamp), lowerCache) } } } @@ -211,12 +218,10 @@ func BuildFromStoreWithLog(store *PacketStore, enableLog bool) *NeighborGraph { // extractFromNode pulls the originator pubkey from a StoreTx's DecodedJSON. // ADVERTs use "pubKey", other packets may use "from_node" or "from". +// Uses the cached ParsedDecoded() accessor to avoid repeated json.Unmarshal. func extractFromNode(tx *StoreTx) string { - if tx.DecodedJSON == "" { - return "" - } - var decoded map[string]interface{} - if err := jsonUnmarshalFast(tx.DecodedJSON, &decoded); err != nil { + decoded := tx.ParsedDecoded() + if decoded == nil { return "" } // ADVERTs store the originator pubkey as "pubKey"; other packets may use @@ -275,9 +280,9 @@ func (g *NeighborGraph) upsertEdge(pubkeyA, pubkeyB, prefix, observer string, sn } // upsertEdgeWithCandidates handles prefix-based edges that may be ambiguous. -func (g *NeighborGraph) upsertEdgeWithCandidates(knownPK, prefix string, candidates []nodeInfo, observer string, snr *float64, ts time.Time) { +func (g *NeighborGraph) upsertEdgeWithCandidates(knownPK, prefix string, candidates []nodeInfo, observer string, snr *float64, ts time.Time, lc map[string]string) { if len(candidates) == 1 { - resolved := strings.ToLower(candidates[0].PublicKey) + resolved := cachedToLower(lc, candidates[0].PublicKey) if resolved == knownPK { return // self-edge guard } @@ -288,7 +293,7 @@ func (g *NeighborGraph) upsertEdgeWithCandidates(knownPK, prefix string, candida // Filter out self from candidates filtered := make([]string, 0, len(candidates)) for _, c := range candidates { - pk := strings.ToLower(c.PublicKey) + pk := cachedToLower(lc, c.PublicKey) if pk != knownPK { filtered = append(filtered, pk) } diff --git a/cmd/server/neighbor_graph_test.go b/cmd/server/neighbor_graph_test.go index 809ba959..9500a134 100644 --- a/cmd/server/neighbor_graph_test.go +++ b/cmd/server/neighbor_graph_test.go @@ -717,3 +717,120 @@ func TestNeighborGraph_CacheTTL(t *testing.T) { t.Error("old graph should be stale") } } + +func TestNeighborGraph_TTLIsReasonable(t *testing.T) { + // TTL must be long enough to avoid rebuild storms on busy meshes, + // but short enough to reflect topology changes within minutes. + if neighborGraphTTL < 1*time.Minute { + t.Errorf("neighborGraphTTL too short (%v), will cause rebuild storms", neighborGraphTTL) + } + if neighborGraphTTL > 10*time.Minute { + t.Errorf("neighborGraphTTL too long (%v), topology changes will be stale", neighborGraphTTL) + } +} + +func TestCachedToLower(t *testing.T) { + cache := make(map[string]string) + // Basic lowercasing + if got := cachedToLower(cache, "AABB"); got != "aabb" { + t.Errorf("expected 'aabb', got %q", got) + } + // Verify it was cached + if _, ok := cache["AABB"]; !ok { + t.Error("expected 'AABB' to be in cache") + } + // Same input returns cached result + if got := cachedToLower(cache, "AABB"); got != "aabb" { + t.Errorf("expected cached 'aabb', got %q", got) + } + // Already lowercase stays the same + if got := cachedToLower(cache, "aabb"); got != "aabb" { + t.Errorf("expected 'aabb', got %q", got) + } + // Empty string + if got := cachedToLower(cache, ""); got != "" { + t.Errorf("expected empty, got %q", got) + } +} + +func TestParsedDecoded_Caching(t *testing.T) { + tx := &StoreTx{DecodedJSON: `{"pubKey":"abc123","name":"test"}`} + // First call parses + d1 := tx.ParsedDecoded() + if d1 == nil { + t.Fatal("expected non-nil parsed result") + } + if d1["pubKey"] != "abc123" { + t.Errorf("expected pubKey=abc123, got %v", d1["pubKey"]) + } + // Second call must return the exact same map (pointer equality proves caching) + d2 := tx.ParsedDecoded() + if &d1 == nil || &d2 == nil { + t.Fatal("unexpected nil") + } + // Mutate d1 and verify d2 sees the mutation — proves same underlying map + d1["_sentinel"] = true + if d2["_sentinel"] != true { + t.Error("expected same map instance from second call (caching broken)") + } + delete(d1, "_sentinel") // clean up +} + +func TestParsedDecoded_EmptyJSON(t *testing.T) { + tx := &StoreTx{DecodedJSON: ""} + d := tx.ParsedDecoded() + if d != nil { + t.Errorf("expected nil for empty DecodedJSON, got %v", d) + } +} + +func TestParsedDecoded_InvalidJSON(t *testing.T) { + tx := &StoreTx{DecodedJSON: "not json"} + d := tx.ParsedDecoded() + if d != nil { + t.Errorf("expected nil for invalid JSON, got %v", d) + } +} + +func TestExtractFromNode_UsesCachedParse(t *testing.T) { + tx := &StoreTx{DecodedJSON: `{"pubKey":"aabb1122"}`} + // First call to extractFromNode should use ParsedDecoded + from := extractFromNode(tx) + if from != "aabb1122" { + t.Errorf("expected aabb1122, got %q", from) + } + // ParsedDecoded should now be cached + d := tx.ParsedDecoded() + if d == nil || d["pubKey"] != "aabb1122" { + t.Error("expected ParsedDecoded to return cached result") + } +} + +func BenchmarkBuildFromStore(b *testing.B) { + // Simulate a dataset with many packets and repeated pubkeys + nodes := []nodeInfo{ + {PublicKey: "aaaa1111", Name: "NodeA"}, + {PublicKey: "bbbb2222", Name: "NodeB"}, + {PublicKey: "cccc3333", Name: "NodeC"}, + {PublicKey: "dddd4444", Name: "NodeD"}, + } + const numPackets = 1000 + packets := make([]*StoreTx, 0, numPackets) + for i := 0; i < numPackets; i++ { + pt := 4 // ADVERT + packets = append(packets, &StoreTx{ + ID: i, + PayloadType: &pt, + DecodedJSON: `{"pubKey":"aaaa1111"}`, + Observations: []*StoreObs{ + {ObserverID: "bbbb2222", PathJSON: `["cccc"]`, Timestamp: nowStr, SNR: ngFloatPtr(-5.0)}, + }, + }) + } + store := ngTestStore(nodes, packets) + + b.ResetTimer() + for i := 0; i < b.N; i++ { + BuildFromStore(store) + } +} diff --git a/cmd/server/store.go b/cmd/server/store.go index 111f5bca..0c5da5eb 100644 --- a/cmd/server/store.go +++ b/cmd/server/store.go @@ -42,8 +42,10 @@ type StoreTx struct { ResolvedPath []*string // resolved path from best observation LatestSeen string // max observation timestamp (or FirstSeen if no observations) // Cached parsed fields (set once, read many) - parsedPath []string // cached parsePathJSON result - pathParsed bool // whether parsedPath has been set + 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 } @@ -63,6 +65,17 @@ type StoreObs struct { Timestamp string } +// ParsedDecoded returns the parsed DecodedJSON map, caching the result. +// Thread-safe via sync.Once — the first call parses, subsequent calls return cached. +func (tx *StoreTx) ParsedDecoded() map[string]interface{} { + tx.decodedOnce.Do(func() { + if tx.DecodedJSON != "" { + json.Unmarshal([]byte(tx.DecodedJSON), &tx.parsedDecoded) + } + }) + return tx.parsedDecoded +} + // distRebuildInterval is the minimum time between distance index rebuilds // to avoid hot-looping on busy meshes. const distRebuildInterval = 30 * time.Second @@ -406,8 +419,8 @@ func (s *PacketStore) indexByNode(tx *StoreTx) { if !strings.Contains(tx.DecodedJSON, "ubKey") { return } - var decoded map[string]interface{} - if json.Unmarshal([]byte(tx.DecodedJSON), &decoded) != nil { + decoded := tx.ParsedDecoded() + if decoded == nil { return } for _, field := range []string{"pubKey", "destPubKey", "srcPubKey"} { @@ -430,8 +443,8 @@ func (s *PacketStore) trackAdvertPubkey(tx *StoreTx) { if tx.PayloadType == nil || *tx.PayloadType != 4 || tx.DecodedJSON == "" { return } - var d map[string]interface{} - if json.Unmarshal([]byte(tx.DecodedJSON), &d) != nil { + d := tx.ParsedDecoded() + if d == nil { return } pk := ""