Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
83 changes: 83 additions & 0 deletions cmd/server/pathlen_fast_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
})
}
156 changes: 156 additions & 0 deletions cmd/server/paths_confirm_deferred_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
1 change: 1 addition & 0 deletions cmd/server/resolved_index.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) + `"`
Expand Down
62 changes: 47 additions & 15 deletions cmd/server/routes.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand All @@ -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()

Expand Down
Loading
Loading