diff --git a/cmd/server/pathlen_fast_test.go b/cmd/server/pathlen_fast_test.go new file mode 100644 index 000000000..fb182d56c --- /dev/null +++ b/cmd/server/pathlen_fast_test.go @@ -0,0 +1,83 @@ +package main + +import ( + "math/rand" + "strings" + "testing" +) + +// pathLen is called once per observation of every candidate transmission in +// /api/nodes/{pk}/paths and /hop_analytics. The fast path must agree with the +// original json.Unmarshal-based implementation (pathLenSlow) for every input. + +func TestPathLenFast_MatchesReference(t *testing.T) { + cases := []string{ + "", " ", "invalid", "null", "true", "0", `""`, `"aa"`, `{}`, `{"a":1}`, + `[]`, ` [] `, "\t[\n]\r\n", `[ ]`, `[,]`, `[""]`, `["",""]`, + `["aa"]`, `["aa","bb"]`, `["aa","bb","cc"]`, ` [ "aa" , "bb" ] `, + `["aa",]`, `[,"aa"]`, `["aa" "bb"]`, `["aa","bb"`, `["aa","bb"]]`, + `["aa","bb"] x`, `x ["aa"]`, `["aa"]["bb"]`, `["aa"`, `["aa`, `[`, `]`, + `[1,2,3]`, `[null]`, `[true,false]`, `["aa",1]`, `[["aa"],["bb"]]`, + `[{"h":"aa"}]`, `["a\"b"]`, `["a\\b","c"]`, `["a,b","c]d","[e"]`, + `["é"]`, `["日本"]`, "[\"a\x01b\"]", "[\"a\nb\"]", "[\"a\x7fb\"]", + `["aa","bb","cc","dd","ee","ff","00","11"]`, + `["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"]`, + } + for _, in := range cases { + if got, want := pathLen(in), pathLenSlow(in); got != want { + t.Errorf("pathLen(%q) = %d, want %d (reference)", in, got, want) + } + } +} + +// TestPathLenFast_RandomisedAgainstReference throws structurally "almost +// valid" JSON at both implementations. The token alphabet is chosen so most +// strings are close to the fast path's grammar, which is where a divergence +// would hide. +func TestPathLenFast_RandomisedAgainstReference(t *testing.T) { + tokens := []string{ + "[", "]", ",", `"`, `"aa"`, `"bb"`, `""`, " ", "\n", "\t", `\`, `\"`, + "a", "1", "null", "é", "\x00", "{", "}", ":", "-", "e", + } + rng := rand.New(rand.NewSource(1352)) + for i := 0; i < 300000; i++ { + var b strings.Builder + for j, n := 0, rng.Intn(14); j < n; j++ { + b.WriteString(tokens[rng.Intn(len(tokens))]) + } + in := b.String() + if got, want := pathLen(in), pathLenSlow(in); got != want { + t.Fatalf("pathLen(%q) = %d, want %d (reference)", in, got, want) + } + } +} + +// TestPathLenFast_WellFormedPathsStayOnFastPath pins that the shapes real +// observations have (hex hop prefixes) are actually served by the +// allocation-free scanner, not the json.Unmarshal fallback. +func TestPathLenFast_WellFormedPathsStayOnFastPath(t *testing.T) { + for _, in := range []string{`[]`, `["aa"]`, `["aa","bb","cc"]`, `["a1b2","c3d4"]`} { + if _, ok := pathLenFast(in); !ok { + t.Errorf("pathLenFast(%q) fell off the fast path", in) + } + } + if allocs := testing.AllocsPerRun(100, func() { _ = pathLen(`["aa","bb","cc"]`) }); allocs != 0 { + t.Errorf("pathLen allocates %.0f times per call on a well-formed path, want 0", allocs) + } +} + +func BenchmarkPathLen(b *testing.B) { + in := `["aa","bb","cc","dd","ee"]` + b.Run("fast", func(b *testing.B) { + b.ReportAllocs() + for i := 0; i < b.N; i++ { + _ = pathLen(in) + } + }) + b.Run("reference_json_unmarshal", func(b *testing.B) { + b.ReportAllocs() + for i := 0; i < b.N; i++ { + _ = pathLenSlow(in) + } + }) +} diff --git a/cmd/server/paths_confirm_deferred_test.go b/cmd/server/paths_confirm_deferred_test.go new file mode 100644 index 000000000..7428f5e98 --- /dev/null +++ b/cmd/server/paths_confirm_deferred_test.go @@ -0,0 +1,156 @@ +package main + +import ( + "encoding/json" + "fmt" + "net/http/httptest" + "strings" + "testing" + "time" +) + +// /api/nodes/{pk}/paths and /hop_analytics used to confirm every hash-index +// candidate with confirmResolvedPathContains: one SQL query per candidate +// transmission, each scanning all observation rows of that tx. For a busy +// node that was thousands of sequential queries and the dominant CPU cost of +// the endpoint. Membership is decided from the canonical resolved_path +// anyway, so the query is now only issued for candidates that have no +// canonical path. These tests pin both halves: the queries are gone in the +// normal case, and results are unchanged when the index is stale/colliding. + +const confirmTestTarget = "aabbccdd11223344" // seeded TestRepeater + +// seedConfirmTx inserts a transmission with one observation and returns its +// id. An empty rpJSON stores a NULL resolved_path. +func seedConfirmTx(t *testing.T, srv *Server, hash, pathJSON, rpJSON string) int { + t.Helper() + now := time.Now().UTC().Add(-20 * time.Minute) + if _, err := srv.db.conn.Exec(`INSERT INTO transmissions (raw_hex, hash, first_seen, route_type, payload_type, decoded_json) + VALUES ('FF01', ?, ?, 1, 4, '{}')`, hash, now.Format(time.RFC3339)); err != nil { + t.Fatalf("insert tx %s: %v", hash, err) + } + var txID int + if err := srv.db.conn.QueryRow(`SELECT id FROM transmissions WHERE hash = ?`, hash).Scan(&txID); err != nil { + t.Fatalf("lookup tx %s: %v", hash, err) + } + var rp interface{} + if rpJSON != "" { + rp = rpJSON + } + if _, err := srv.db.conn.Exec(`INSERT INTO observations (transmission_id, observer_idx, snr, rssi, path_json, timestamp, resolved_path) + VALUES (?, 1, 10.0, -90, ?, ?, ?)`, txID, pathJSON, now.Unix(), rp); err != nil { + t.Fatalf("insert obs %s: %v", hash, err) + } + return txID +} + +// reloadConfirmStore swaps in a freshly loaded store so the seeded rows and +// their resolved-path index entries are visible to the handlers. +func reloadConfirmStore(t *testing.T, srv *Server) *PacketStore { + t.Helper() + store := NewPacketStore(srv.db, nil) + if err := store.Load(); err != nil { + t.Fatalf("store.Load: %v", err) + } + if !store.WaitIndexesReady(5 * time.Second) { + t.Fatal("indexes never became ready") + } + srv.store = store + return store +} + +type confirmPathsResp struct { + Paths []json.RawMessage `json:"paths"` + TotalTransmissions int `json:"totalTransmissions"` +} + +func TestNodePaths_NoPerCandidateSQLConfirm(t *testing.T) { + srv, router := setupTestServer(t) + const n = 25 + for i := 0; i < n; i++ { + seedConfirmTx(t, srv, fmt.Sprintf("confirm_def_%02d", i), + `["aa","bb"]`, `["`+confirmTestTarget+`","eeff00112233aabb"]`) + } + store := reloadConfirmStore(t, srv) + + w := httptest.NewRecorder() + router.ServeHTTP(w, httptest.NewRequest("GET", "/api/nodes/"+confirmTestTarget+"/paths", nil)) + if w.Code != 200 { + t.Fatalf("expected 200, got %d: %s", w.Code, w.Body.String()) + } + var resp confirmPathsResp + if err := json.Unmarshal(w.Body.Bytes(), &resp); err != nil { + t.Fatalf("bad JSON: %v", err) + } + if resp.TotalTransmissions < n { + t.Errorf("totalTransmissions = %d, want >= %d (all seeded txs resolve through the target)", resp.TotalTransmissions, n) + } + if q := store.confirmResolvedPathQueries.Load(); q != 0 { + t.Errorf("confirmResolvedPathContains ran %d times for %d candidates, want 0 (every candidate has a canonical resolved_path)", q, n) + } +} + +// A hash-index entry that points at a tx whose stored resolved_path does NOT +// contain the queried pubkey (hash collision, or an index entry made stale by +// a later resolved_path overwrite) must still be excluded — now by the +// canonical-path check instead of the SQL pre-filter. +func TestNodePaths_StaleIndexEntryStillExcluded(t *testing.T) { + srv, router := setupTestServer(t) + staleID := seedConfirmTx(t, srv, "confirm_stale_hash", + `["aa","bb"]`, `["aacafe0000000000","eeff00112233aabb"]`) + store := reloadConfirmStore(t, srv) + + store.mu.Lock() + h := resolvedPubkeyHash(confirmTestTarget) + store.resolvedPubkeyIndex[h] = append(store.resolvedPubkeyIndex[h], staleID) + store.mu.Unlock() + + w := httptest.NewRecorder() + router.ServeHTTP(w, httptest.NewRequest("GET", "/api/nodes/"+confirmTestTarget+"/paths", nil)) + if w.Code != 200 { + t.Fatalf("expected 200, got %d: %s", w.Code, w.Body.String()) + } + if strings.Contains(w.Body.String(), "confirm_stale_hash") { + t.Error("tx whose resolved_path does not contain the target leaked into /paths via a stale index entry") + } + + w = httptest.NewRecorder() + router.ServeHTTP(w, httptest.NewRequest("GET", "/api/nodes/"+confirmTestTarget+"/hop_analytics", nil)) + if w.Code != 200 { + t.Fatalf("hop_analytics: expected 200, got %d: %s", w.Code, w.Body.String()) + } + if strings.Contains(w.Body.String(), "confirm_stale_hash") { + t.Error("tx whose resolved_path does not contain the target leaked into /hop_analytics via a stale index entry") + } + if q := store.confirmResolvedPathQueries.Load(); q != 0 { + t.Errorf("confirmResolvedPathContains ran %d times, want 0", q) + } +} + +// Candidates with no canonical resolved_path are decided by the legacy +// fallback arm, which still relies on the SQL confirmation. That behaviour is +// preserved: the query runs (once) and a NULL resolved_path does not confirm. +func TestNodePaths_NoCanonicalPathStillConfirmedBySQL(t *testing.T) { + srv, router := setupTestServer(t) + nullID := seedConfirmTx(t, srv, "confirm_nullrp_hash", `["aa","bb"]`, "") + store := reloadConfirmStore(t, srv) + + // Force the "in the hash index but no canonical path" state. + store.mu.Lock() + h := resolvedPubkeyHash(confirmTestTarget) + store.resolvedPubkeyIndex[h] = append(store.resolvedPubkeyIndex[h], nullID) + store.resolvedPubkeyReverse[nullID] = []uint64{h ^ 1} // has *some* indexed pubkey + store.mu.Unlock() + + w := httptest.NewRecorder() + router.ServeHTTP(w, httptest.NewRequest("GET", "/api/nodes/"+confirmTestTarget+"/paths", nil)) + if w.Code != 200 { + t.Fatalf("expected 200, got %d: %s", w.Code, w.Body.String()) + } + if strings.Contains(w.Body.String(), "confirm_nullrp_hash") { + t.Error("tx with NULL resolved_path was confirmed for the target") + } + if q := store.confirmResolvedPathQueries.Load(); q != 1 { + t.Errorf("confirmResolvedPathContains ran %d times, want exactly 1 (the tx without a canonical path)", q) + } +} diff --git a/cmd/server/resolved_index.go b/cmd/server/resolved_index.go index ccc901903..d2bcafe31 100644 --- a/cmd/server/resolved_index.go +++ b/cmd/server/resolved_index.go @@ -153,6 +153,7 @@ func (s *PacketStore) confirmResolvedPathContains(txID int, pubkey string) bool if s.db == nil || s.db.conn == nil { return true } + s.confirmResolvedPathQueries.Add(1) // Use INSTR with surrounding quotes for exact match — avoids LIKE escape issues. // resolved_path format: ["pubkey1","pubkey2",...] needle := `"` + strings.ToLower(pubkey) + `"` diff --git a/cmd/server/routes.go b/cmd/server/routes.go index ab8b2a50d..2df908bc6 100644 --- a/cmd/server/routes.go +++ b/cmd/server/routes.go @@ -2221,37 +2221,48 @@ func (s *Server) handleNodePaths(w http.ResponseWriter, r *http.Request) { inIndex bool } checks := make([]candidateCheck, len(candidates)) + // Membership set for the queried pubkey, built once. The previous + // per-candidate scan of the index list was O(candidates × list length). + var indexedForTarget map[int]struct{} + if s.store.useResolvedPathIndex { + ids := s.store.resolvedPubkeyIndex[resolvedPubkeyHash(lowerPK)] + indexedForTarget = make(map[int]struct{}, len(ids)) + for _, id := range ids { + indexedForTarget[id] = struct{}{} + } + } for i, tx := range candidates { cc := candidateCheck{tx: tx} if !s.store.useResolvedPathIndex { cc.inIndex = true // flag off — keep all } else if _, hasRev := s.store.resolvedPubkeyReverse[tx.ID]; !hasRev { cc.inIndex = true // no indexed pubkeys — keep (conservative) - } else { - h := resolvedPubkeyHash(lowerPK) - for _, id := range s.store.resolvedPubkeyIndex[h] { - if id == tx.ID { - cc.hasReverse = true // needs SQL confirmation - break - } - } - // If not in index at all, it's a definite no + } else if _, ok := indexedForTarget[tx.ID]; ok { + cc.hasReverse = true // hash-index hit; exact pubkey confirmed below } + // If not in index at all, it's a definite no checks[i] = cc } s.store.mu.RUnlock() - // Now run SQL checks outside the lock for candidates that need confirmation. - confirmedBySQL := make(map[int]bool) + // Candidates admitted by the hash index (hasReverse) used to be confirmed + // one by one with confirmResolvedPathContains — a SQL query per candidate + // that scans every observation row of the tx. For a busy node that is + // thousands of sequential queries and dominated /paths CPU (~43% of a + // 60 s profile). The confirmation only guards against hash collisions + // and a stale index; for every candidate that has a canonical persisted + // resolved_path, membership is decided again below from that exact path + // (resolvedPK == lowerPK), which gives the same answer. So defer the SQL + // check to the few candidates with no canonical path at all, where the + // legacy fallback still needs confirmedBySQL. + needsConfirm := make(map[int]bool) filtered := candidates[:0] for _, cc := range checks { if cc.inIndex { filtered = append(filtered, cc.tx) } else if cc.hasReverse { - if s.store.confirmResolvedPathContains(cc.tx.ID, lowerPK) { - filtered = append(filtered, cc.tx) - confirmedBySQL[cc.tx.ID] = true - } + filtered = append(filtered, cc.tx) + needsConfirm[cc.tx.ID] = true } // else: not in index → exclude } @@ -2278,6 +2289,27 @@ func (s *Server) handleNodePaths(w http.ResponseWriter, r *http.Request) { } } + // Deferred exact-pubkey confirmation (see needsConfirm above): only for + // hash-index candidates that have no canonical resolved_path, because + // those are the ones decided by the legacy fallback arm, which consumes + // confirmedBySQL. + confirmedBySQL := make(map[int]bool) + if len(needsConfirm) > 0 { + kept := candidates[:0] + for _, tx := range candidates { + if needsConfirm[tx.ID] { + if _, hasCanonical := canonicalRP[tx.ID]; !hasCanonical { + if !s.store.confirmResolvedPathContains(tx.ID, lowerPK) { + continue + } + confirmedBySQL[tx.ID] = true + } + } + kept = append(kept, tx) + } + candidates = kept + } + // Re-acquire read lock for the aggregation phase that reads store data. s.store.mu.RLock() diff --git a/cmd/server/store.go b/cmd/server/store.go index 2e4a945b8..ba0c37787 100644 --- a/cmd/server/store.go +++ b/cmd/server/store.go @@ -457,6 +457,12 @@ type PacketStore struct { lruOrder []int // FIFO order for LRU eviction lruMu sync.RWMutex // guards apiResolvedPathLRU + lruOrder + // confirmResolvedPathQueries counts SQL round-trips made by + // confirmResolvedPathContains. Each one scans the tx's observation rows, + // so tests use it to pin that /paths and /hop_analytics no longer issue + // one per candidate transmission. + confirmResolvedPathQueries atomic.Uint64 + // Persisted neighbor graph for hop resolution at ingest time. // Accessed via atomic.Pointer because async rebuilds (path_inspect.go // ensureNeighborGraph) and ingest-time readers race on the pointer @@ -1901,7 +1907,27 @@ func (s *PacketStore) txChargedBytes(tx *StoreTx) int64 { return int64(tx.chargedBytes) + resolvedRelayBytes(len(s.pathHopResolved[tx])) + fallbackRelayBytes(len(s.fallbackByNode[tx])) } +// pathLen returns the number of hops in a JSON-encoded path array (0 for an +// empty or invalid input). +// +// It runs once per observation of every candidate transmission in +// /api/nodes/{pk}/paths and /hop_analytics (via fetchResolvedPathForTxBest), +// i.e. hundreds of thousands of times per request, so it must not allocate. +// The common shape — an array of plain ASCII strings such as ["aa","bb"] — is +// counted with a single byte scan. Anything outside that restricted grammar +// (escapes, non-ASCII or control bytes, nested values, numbers, null, +// trailing data, malformed input) falls back to pathLenSlow, which keeps the +// original json.Unmarshal semantics exactly. func pathLen(pathJSON string) int { + if n, ok := pathLenFast(pathJSON); ok { + return n + } + return pathLenSlow(pathJSON) +} + +// pathLenSlow is the reference implementation: parse the whole array and +// count it. Kept as the fallback for inputs pathLenFast does not recognise. +func pathLenSlow(pathJSON string) int { if pathJSON == "" { return 0 } @@ -1912,6 +1938,79 @@ func pathLen(pathJSON string) int { return len(hops) } +func isJSONSpace(c byte) bool { + return c == ' ' || c == '\t' || c == '\n' || c == '\r' +} + +// pathLenFast counts the elements of a JSON array of simple strings without +// allocating. ok is false when s is anything other than +// +// ws* "[" ws* ( "]" | str ( ws* "," ws* str )* ws* "]" ) ws* +// +// where str is a double-quoted run of printable ASCII (0x20-0x7e) with no +// backslash. For every input it accepts, the count equals what +// json.Unmarshal into []interface{} would give. +func pathLenFast(s string) (n int, ok bool) { + l := len(s) + i := 0 + for i < l && isJSONSpace(s[i]) { + i++ + } + if i == l || s[i] != '[' { + return 0, false + } + i++ + for i < l && isJSONSpace(s[i]) { + i++ + } + if i < l && s[i] == ']' { + i++ + for i < l && isJSONSpace(s[i]) { + i++ + } + return 0, i == l + } + count := 0 + for { + if i >= l || s[i] != '"' { + return 0, false + } + i++ + for i < l && s[i] != '"' { + if c := s[i]; c < 0x20 || c >= 0x7f || c == '\\' { + return 0, false + } + i++ + } + if i >= l { + return 0, false + } + i++ // closing quote + count++ + for i < l && isJSONSpace(s[i]) { + i++ + } + if i >= l { + return 0, false + } + switch s[i] { + case ',': + i++ + for i < l && isJSONSpace(s[i]) { + i++ + } + case ']': + i++ + for i < l && isJSONSpace(s[i]) { + i++ + } + return count, i == l + default: + return 0, false + } + } +} + // pathFirstHop returns path[0] (the entry-point repeater's hex prefix), or // "" when pathJSON is empty/invalid/has no hops. func pathFirstHop(pathJSON string) string { @@ -10530,30 +10629,36 @@ func (s *PacketStore) GetNodeHopAnalytics(pubkey string, days int) (*NodeHopAnal inIndex bool } checks := make([]candidateCheck, len(candidates)) + // Built once instead of scanning the index list for every candidate. + var indexedForTarget map[int]struct{} + if s.useResolvedPathIndex { + ids := s.resolvedPubkeyIndex[resolvedPubkeyHash(lowerPK)] + indexedForTarget = make(map[int]struct{}, len(ids)) + for _, id := range ids { + indexedForTarget[id] = struct{}{} + } + } for i, tx := range candidates { cc := candidateCheck{tx: tx} if !s.useResolvedPathIndex { cc.inIndex = true } else if _, hasRev := s.resolvedPubkeyReverse[tx.ID]; !hasRev { cc.inIndex = true - } else { - h := resolvedPubkeyHash(lowerPK) - for _, id := range s.resolvedPubkeyIndex[h] { - if id == tx.ID { - cc.hasReverse = true - break - } - } + } else if _, ok := indexedForTarget[tx.ID]; ok { + cc.hasReverse = true } checks[i] = cc } s.mu.RUnlock() + // No per-candidate confirmResolvedPathContains SQL query here: the loop + // below already decides membership from the canonical resolved_path + // (rp == nil / idx < 0 → skip), which also rejects hash collisions and + // stale index entries. The query only scanned every observation row of + // each candidate tx to reach the same conclusion. confirmed := candidates[:0] for _, cc := range checks { - if cc.inIndex { - confirmed = append(confirmed, cc.tx) - } else if cc.hasReverse && s.confirmResolvedPathContains(cc.tx.ID, lowerPK) { + if cc.inIndex || cc.hasReverse { confirmed = append(confirmed, cc.tx) } }