diff --git a/cmd/server/advert_pubkey_test.go b/cmd/server/advert_pubkey_test.go new file mode 100644 index 00000000..d15b1d71 --- /dev/null +++ b/cmd/server/advert_pubkey_test.go @@ -0,0 +1,181 @@ +package main + +import ( + "encoding/json" + "fmt" + "testing" +) + +// TestAdvertPubkeyTracking verifies that advertPubkeys is maintained +// incrementally during ingest and eviction, and that GetPerfStoreStats +// returns the correct count without per-request JSON parsing. +func TestAdvertPubkeyTracking(t *testing.T) { + ps := NewPacketStore(nil, nil) + ps.mu.Lock() + + // Helper to create an ADVERT StoreTx with a given pubkey. + pt4 := 4 + mkAdvert := func(id int, pubkey string) *StoreTx { + d := map[string]interface{}{"pubKey": pubkey} + j, _ := json.Marshal(d) + return &StoreTx{ + ID: id, + Hash: fmt.Sprintf("hash%d", id), + PayloadType: &pt4, + DecodedJSON: string(j), + } + } + + // Add 3 adverts: 2 distinct pubkeys + tx1 := mkAdvert(1, "pk_alpha") + tx2 := mkAdvert(2, "pk_beta") + tx3 := mkAdvert(3, "pk_alpha") // duplicate pubkey + + for _, tx := range []*StoreTx{tx1, tx2, tx3} { + ps.packets = append(ps.packets, tx) + ps.byHash[tx.Hash] = tx + ps.byTxID[tx.ID] = tx + ps.byPayloadType[4] = append(ps.byPayloadType[4], tx) + ps.trackAdvertPubkey(tx) + } + ps.mu.Unlock() + + // GetPerfStoreStats should report 2 distinct pubkeys + stats := ps.GetPerfStoreStats() + indexes := stats["indexes"].(map[string]interface{}) + got := indexes["advertByObserver"].(int) + if got != 2 { + t.Errorf("advertByObserver = %d, want 2", got) + } + + // GetPerfStoreStatsTyped should agree + typed := ps.GetPerfStoreStatsTyped() + if typed.Indexes.AdvertByObserver != 2 { + t.Errorf("typed AdvertByObserver = %d, want 2", typed.Indexes.AdvertByObserver) + } + + // Evict tx3 (pk_alpha duplicate) — count should stay 2 + ps.mu.Lock() + ps.untrackAdvertPubkey(tx3) + ps.mu.Unlock() + + stats2 := ps.GetPerfStoreStats() + idx2 := stats2["indexes"].(map[string]interface{}) + if idx2["advertByObserver"].(int) != 2 { + t.Errorf("after evicting duplicate: advertByObserver = %d, want 2", idx2["advertByObserver"].(int)) + } + + // Evict tx1 (last pk_alpha) — count should drop to 1 + ps.mu.Lock() + ps.untrackAdvertPubkey(tx1) + ps.mu.Unlock() + + stats3 := ps.GetPerfStoreStats() + idx3 := stats3["indexes"].(map[string]interface{}) + if idx3["advertByObserver"].(int) != 1 { + t.Errorf("after evicting last pk_alpha: advertByObserver = %d, want 1", idx3["advertByObserver"].(int)) + } + + // Evict tx2 (last remaining) — count should be 0 + ps.mu.Lock() + ps.untrackAdvertPubkey(tx2) + ps.mu.Unlock() + + stats4 := ps.GetPerfStoreStats() + idx4 := stats4["indexes"].(map[string]interface{}) + if idx4["advertByObserver"].(int) != 0 { + t.Errorf("after evicting all: advertByObserver = %d, want 0", idx4["advertByObserver"].(int)) + } +} + +// TestAdvertPubkeyPublicKeyField tests the "public_key" JSON field variant. +func TestAdvertPubkeyPublicKeyField(t *testing.T) { + ps := NewPacketStore(nil, nil) + ps.mu.Lock() + pt4 := 4 + d, _ := json.Marshal(map[string]interface{}{"public_key": "pk_legacy"}) + tx := &StoreTx{ID: 1, Hash: "h1", PayloadType: &pt4, DecodedJSON: string(d)} + ps.trackAdvertPubkey(tx) + ps.mu.Unlock() + + stats := ps.GetPerfStoreStats() + idx := stats["indexes"].(map[string]interface{}) + if idx["advertByObserver"].(int) != 1 { + t.Errorf("public_key field: advertByObserver = %d, want 1", idx["advertByObserver"].(int)) + } +} + +// TestAdvertPubkeyNonAdvert ensures non-ADVERT packets don't affect the count. +func TestAdvertPubkeyNonAdvert(t *testing.T) { + ps := NewPacketStore(nil, nil) + ps.mu.Lock() + pt2 := 2 + d, _ := json.Marshal(map[string]interface{}{"pubKey": "pk_text"}) + tx := &StoreTx{ID: 1, Hash: "h1", PayloadType: &pt2, DecodedJSON: string(d)} + ps.trackAdvertPubkey(tx) + ps.mu.Unlock() + + stats := ps.GetPerfStoreStats() + idx := stats["indexes"].(map[string]interface{}) + if idx["advertByObserver"].(int) != 0 { + t.Errorf("non-ADVERT should not be tracked: advertByObserver = %d, want 0", idx["advertByObserver"].(int)) + } +} + +// BenchmarkGetPerfStoreStats benchmarks the perf stats endpoint with many adverts. +// Before the fix, this did O(N) JSON unmarshals per call. +// After the fix, it's O(1) — just len(map). +func BenchmarkGetPerfStoreStats(b *testing.B) { + ps := NewPacketStore(nil, nil) + ps.mu.Lock() + pt4 := 4 + for i := 0; i < 5000; i++ { + pk := fmt.Sprintf("pk_%04d", i%200) // 200 distinct pubkeys + d, _ := json.Marshal(map[string]interface{}{"pubKey": pk}) + tx := &StoreTx{ + ID: i + 1, + Hash: fmt.Sprintf("hash%d", i+1), + PayloadType: &pt4, + DecodedJSON: string(d), + } + ps.packets = append(ps.packets, tx) + ps.byHash[tx.Hash] = tx + ps.byTxID[tx.ID] = tx + ps.byPayloadType[4] = append(ps.byPayloadType[4], tx) + ps.trackAdvertPubkey(tx) + } + ps.mu.Unlock() + + b.ResetTimer() + for i := 0; i < b.N; i++ { + ps.GetPerfStoreStats() + } +} + +// BenchmarkGetPerfStoreStatsTyped benchmarks the typed variant. +func BenchmarkGetPerfStoreStatsTyped(b *testing.B) { + ps := NewPacketStore(nil, nil) + ps.mu.Lock() + pt4 := 4 + for i := 0; i < 5000; i++ { + pk := fmt.Sprintf("pk_%04d", i%200) + d, _ := json.Marshal(map[string]interface{}{"pubKey": pk}) + tx := &StoreTx{ + ID: i + 1, + Hash: fmt.Sprintf("hash%d", i+1), + PayloadType: &pt4, + DecodedJSON: string(d), + } + ps.packets = append(ps.packets, tx) + ps.byHash[tx.Hash] = tx + ps.byTxID[tx.ID] = tx + ps.byPayloadType[4] = append(ps.byPayloadType[4], tx) + ps.trackAdvertPubkey(tx) + } + ps.mu.Unlock() + + b.ResetTimer() + for i := 0; i < b.N; i++ { + ps.GetPerfStoreStatsTyped() + } +} diff --git a/cmd/server/store.go b/cmd/server/store.go index 5723fd3d..61a38a78 100644 --- a/cmd/server/store.go +++ b/cmd/server/store.go @@ -119,6 +119,10 @@ type PacketStore struct { hashSizeInfoCache map[string]*hashSizeNodeInfo hashSizeInfoAt time.Time + // Precomputed distinct advert pubkey count (refcounted for eviction correctness). + // Updated incrementally during Load/Ingest/Evict — avoids JSON parsing in GetPerfStoreStats. + advertPubkeys map[string]int // pubkey → number of advert packets referencing it + // Eviction config and stats retentionHours float64 // 0 = unlimited maxMemoryMB int // 0 = unlimited @@ -185,6 +189,7 @@ func NewPacketStore(db *DB, cfg *PacketStoreConfig) *PacketStore { rfCacheTTL: 15 * time.Second, collisionCacheTTL: 60 * time.Second, spIndex: make(map[string]int, 4096), + advertPubkeys: make(map[string]int), } if cfg != nil { ps.retentionHours = cfg.RetentionHours @@ -265,6 +270,7 @@ func (s *PacketStore) Load() error { pt := *tx.PayloadType s.byPayloadType[pt] = append(s.byPayloadType[pt], tx) } + s.trackAdvertPubkey(tx) } if obsID.Valid { @@ -390,6 +396,52 @@ func (s *PacketStore) indexByNode(tx *StoreTx) { } } +// trackAdvertPubkey increments the advertPubkeys refcount for ADVERT packets. +// Must be called under s.mu write lock. +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 { + return + } + pk := "" + if v, ok := d["pubKey"].(string); ok { + pk = v + } else if v, ok := d["public_key"].(string); ok { + pk = v + } + if pk != "" { + s.advertPubkeys[pk]++ + } +} + +// untrackAdvertPubkey decrements the advertPubkeys refcount for ADVERT packets. +// Must be called under s.mu write lock. +func (s *PacketStore) untrackAdvertPubkey(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 { + return + } + pk := "" + if v, ok := d["pubKey"].(string); ok { + pk = v + } else if v, ok := d["public_key"].(string); ok { + pk = v + } + if pk != "" { + if s.advertPubkeys[pk] <= 1 { + delete(s.advertPubkeys, pk) + } else { + s.advertPubkeys[pk]-- + } + } +} + // QueryPackets returns filtered, paginated packets from memory. func (s *PacketStore) QueryPackets(q PacketQuery) *PacketResult { atomic.AddInt64(&s.queryCount, 1) @@ -577,30 +629,8 @@ func (s *PacketStore) GetPerfStoreStats() map[string]interface{} { nodeIdx := len(s.byNode) ptIdx := len(s.byPayloadType) - // Count distinct pubkeys with ADVERT observations (matches Node.js _advertByObserver.size) - advertByObsCount := 0 - if adverts, ok := s.byPayloadType[4]; ok { - seen := make(map[string]bool) - for _, tx := range adverts { - if tx.DecodedJSON == "" { - continue - } - var d map[string]interface{} - if json.Unmarshal([]byte(tx.DecodedJSON), &d) != nil { - continue - } - pk := "" - if v, ok := d["pubKey"].(string); ok { - pk = v - } else if v, ok := d["public_key"].(string); ok { - pk = v - } - if pk != "" && !seen[pk] { - seen[pk] = true - advertByObsCount++ - } - } - } + // Distinct advert pubkey count — precomputed incrementally (see trackAdvertPubkey). + advertByObsCount := len(s.advertPubkeys) s.mu.RUnlock() // Realistic estimate: ~5KB per packet + ~500 bytes per observation @@ -740,29 +770,7 @@ func (s *PacketStore) GetPerfStoreStatsTyped() PerfPacketStoreStats { observerIdx := len(s.byObserver) nodeIdx := len(s.byNode) - advertByObsCount := 0 - if adverts, ok := s.byPayloadType[4]; ok { - seen := make(map[string]bool) - for _, tx := range adverts { - if tx.DecodedJSON == "" { - continue - } - var d map[string]interface{} - if json.Unmarshal([]byte(tx.DecodedJSON), &d) != nil { - continue - } - pk := "" - if v, ok := d["pubKey"].(string); ok { - pk = v - } else if v, ok := d["public_key"].(string); ok { - pk = v - } - if pk != "" && !seen[pk] { - seen[pk] = true - advertByObsCount++ - } - } - } + advertByObsCount := len(s.advertPubkeys) s.mu.RUnlock() estimatedMB := math.Round(float64(totalLoaded*5120+totalObs*500)/1048576*10) / 10 @@ -1071,6 +1079,7 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac // so GetChannelMessages reverse iteration stays correct s.byPayloadType[pt] = append(s.byPayloadType[pt], tx) } + s.trackAdvertPubkey(tx) if _, exists := broadcastTxs[r.txID]; !exists { broadcastTxs[r.txID] = tx @@ -1938,6 +1947,7 @@ func (s *PacketStore) EvictStale() int { } // Remove from byPayloadType + s.untrackAdvertPubkey(tx) if tx.PayloadType != nil { pt := *tx.PayloadType ptList := s.byPayloadType[pt]