From 514a68db5c7e0c29ad5f9290d182898e01320825 Mon Sep 17 00:00:00 2001 From: dborup Date: Fri, 2 Oct 2026 19:07:06 +0200 Subject: [PATCH 1/2] perf(nodes): cache region membership for node listings --- cmd/server/db.go | 37 +-- cmd/server/nodes_region_cache.go | 312 +++++++++++++++++++++ cmd/server/nodes_region_cache_port_test.go | 119 ++++++++ 3 files changed, 445 insertions(+), 23 deletions(-) create mode 100644 cmd/server/nodes_region_cache.go create mode 100644 cmd/server/nodes_region_cache_port_test.go diff --git a/cmd/server/db.go b/cmd/server/db.go index 5485fae9c..bee3e0c99 100644 --- a/cmd/server/db.go +++ b/cmd/server/db.go @@ -18,6 +18,7 @@ import ( "github.com/meshcore-analyzer/dbschema" "github.com/meshcore-analyzer/geofilter" regionutil "github.com/meshcore-analyzer/regions" + "golang.org/x/sync/singleflight" _ "modernc.org/sqlite" ) @@ -112,6 +113,13 @@ type DB struct { // channelsRowsHook wraps the result rows of each real query, so a test // can fail the iteration part-way. channelsRowsHook func(kind string, rows channelRows) channelRows + + // Region-membership cache for GetNodes (see nodes_region_cache.go). + nodeRegionCacheMu sync.Mutex + nodeRegionCache map[string]*nodeRegionEntry + nodeRegionSF singleflight.Group + nodeRegionFullMu sync.Mutex + nodeRegionQueryHook func() } // channelRows is the part of *sql.Rows the channel list scans use. @@ -1126,30 +1134,13 @@ func (db *DB) GetNodes(limit, offset int, role, search, before, lastHeard, sortB } } - if region != "" { - codes := normalizeRegionCodes(region) - if len(codes) > 0 { - placeholders := make([]string, len(codes)) - regionArgs := make([]interface{}, len(codes)) - for i, c := range codes { - placeholders[i] = "?" - regionArgs[i] = c - } - joinCond := "obs.rowid = o.observer_idx" - if !db.isV3() { - joinCond = "obs.id = o.observer_id" - } - subq := fmt.Sprintf(`public_key IN ( - SELECT DISTINCT JSON_EXTRACT(t.decoded_json, '$.pubKey') - FROM transmissions t - JOIN observations o ON o.transmission_id = t.id - JOIN observers obs ON %s - WHERE t.payload_type = 4 - AND UPPER(TRIM(obs.iata)) IN (%s) - )`, joinCond, strings.Join(placeholders, ",")) - where = append(where, subq) - args = append(args, regionArgs...) + if codes := normalizeRegionCodes(region); len(codes) > 0 { + keysJSON, err := db.nodeRegionKeysJSON(codes) + if err != nil { + return nil, 0, nil, err } + where = append(where, "public_key IN (SELECT value FROM json_each(?))") + args = append(args, keysJSON) } w := "" diff --git a/cmd/server/nodes_region_cache.go b/cmd/server/nodes_region_cache.go new file mode 100644 index 000000000..dc549a5a9 --- /dev/null +++ b/cmd/server/nodes_region_cache.go @@ -0,0 +1,312 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "log" + "sort" + "strings" + "time" +) + +// Region-scoped node membership cache for GetNodes (#2101). +// +// `/api/nodes?region=X` restricts the node list to nodes whose ADVERTs were +// heard by an observer in region X. That used to be an inline +// `public_key IN (SELECT DISTINCT from_pubkey FROM transmissions ⋈ +// observations ⋈ observers ...)` subquery, uncached, evaluated twice per +// request (COUNT(*) and the page). On a 1.4 GB database (1.47M observations) +// each evaluation took ~10s, and fetchAllNodes() pages at 500, so a single +// open Nodes or Map tab kept a 2-vCPU host busy. +// +// Membership only grows as observations arrive, so it is cached per region +// set with an observations.id watermark: +// - fresh (< nodeRegionFreshTTL): served as-is; +// - stale: served as-is while ONE background refresh scans only +// observations with id > watermark (a rowid range, milliseconds); +// - every nodeRegionRebuildInterval the refresh is a full rebuild instead, +// so retention pruning, observer IATA changes and late from_pubkey +// backfills are picked up. +// +// Only the first request for a region set waits on a full scan, and +// concurrent first requests share it (the #1910 singleflight pattern, as +// channelsSF). + +const ( + nodeRegionFreshTTL = 30 * time.Second + nodeRegionRebuildInterval = 30 * time.Minute + // nodeRegionMaxEntries bounds the cache. Entries hold up to a few thousand + // pubkeys each, so this is far below the shared maxCacheEntries; in + // practice an instance sees a handful of distinct region sets. + nodeRegionMaxEntries = 32 +) + +// errNodeRegionFlightResult means nodeRegionSF returned something other than +// a *nodeRegionEntry, which only a programming error can cause. +var errNodeRegionFlightResult = errors.New("unexpected node region flight result") + +// nodeRegionEntry is immutable once stored: refreshes build a new entry, so +// readers can use one without holding nodeRegionCacheMu. +type nodeRegionEntry struct { + keys []string // sorted from_pubkeys heard in the region set + keysJSON string // keys as a JSON array, bound to json_each(?) + watermark int64 // max observations.id covered by keys + refreshed time.Time + built time.Time // last full rebuild +} + +func (e *nodeRegionEntry) fresh() bool { + return time.Since(e.refreshed) < nodeRegionFreshTTL +} + +func (e *nodeRegionEntry) rebuildDue() bool { + return time.Since(e.built) >= nodeRegionRebuildInterval +} + +// newNodeRegionEntry snapshots set into an entry. +func newNodeRegionEntry(set map[string]struct{}, watermark int64, built time.Time) (*nodeRegionEntry, error) { + keys := make([]string, 0, len(set)) + for k := range set { + keys = append(keys, k) + } + sort.Strings(keys) + buf, err := json.Marshal(keys) + if err != nil { + return nil, fmt.Errorf("encode region keys: %w", err) + } + return &nodeRegionEntry{ + keys: keys, + keysJSON: string(buf), + watermark: watermark, + refreshed: time.Now(), + built: built, + }, nil +} + +// nodeRegionKey canonicalises normalised region codes so "EDI,GLA", +// "gla, edi" and "EDI,EDI,GLA" share one cache entry. +func nodeRegionKey(codes []string) string { + seen := make(map[string]struct{}, len(codes)) + uniq := make([]string, 0, len(codes)) + for _, c := range codes { + if _, ok := seen[c]; !ok { + seen[c] = struct{}{} + uniq = append(uniq, c) + } + } + sort.Strings(uniq) + return strings.Join(uniq, ",") +} + +func (db *DB) getNodeRegionEntry(key string) *nodeRegionEntry { + db.nodeRegionCacheMu.Lock() + defer db.nodeRegionCacheMu.Unlock() + return db.nodeRegionCache[key] +} + +func (db *DB) setNodeRegionEntry(key string, e *nodeRegionEntry) { + db.nodeRegionCacheMu.Lock() + defer db.nodeRegionCacheMu.Unlock() + _, replacing := db.nodeRegionCache[key] + if db.nodeRegionCache == nil || (!replacing && len(db.nodeRegionCache) >= nodeRegionMaxEntries) { + db.nodeRegionCache = make(map[string]*nodeRegionEntry) + } + db.nodeRegionCache[key] = e +} + +// nodeRegionKeysJSON returns the JSON array of node public keys heard in the +// given (normalised, non-empty) region codes, for binding to +// `public_key IN (SELECT value FROM json_each(?))`. +func (db *DB) nodeRegionKeysJSON(codes []string) (string, error) { + key := nodeRegionKey(codes) + + if e := db.getNodeRegionEntry(key); e != nil { + if !e.fresh() { + // Stale-while-revalidate: kick one refresh, don't wait on it. + // DoChan dedups against an in-flight refresh for the same key, + // and its buffered result channel may be left unread. + db.nodeRegionSF.DoChan(key, func() (any, error) { + return db.refreshNodeRegion(key, codes) + }) + } + return e.keysJSON, nil + } + + v, err, _ := db.nodeRegionSF.Do(key, func() (any, error) { + return db.refreshNodeRegion(key, codes) + }) + if err != nil { + return "", err + } + e, ok := v.(*nodeRegionEntry) + if !ok { + return "", fmt.Errorf("nodes-region %s: %w: %T", key, errNodeRegionFlightResult, v) + } + return e.keysJSON, nil +} + +// refreshNodeRegion brings the cache entry for key up to date and stores it: +// a delta scan past the watermark when an entry exists and is not due a +// rebuild, a full scan otherwise. It runs inside nodeRegionSF, so at most one +// refresh per key is in flight. Failures are logged here because a background +// refresh has no caller to report to; the previous entry stays in place. +func (db *DB) refreshNodeRegion(key string, codes []string) (*nodeRegionEntry, error) { + prev := db.getNodeRegionEntry(key) + if prev != nil && prev.fresh() { + return prev, nil // a flight that finished just before this one got here + } + if db.nodeRegionQueryHook != nil { + db.nodeRegionQueryHook() + } + + // A scan still running when the next rebuild would be due is abandoned + // rather than left holding one of the pooled connections. + ctx, cancel := context.WithTimeout(context.Background(), nodeRegionRebuildInterval) + defer cancel() + + var e *nodeRegionEntry + var err error + if prev == nil || prev.rebuildDue() { + e, err = db.buildNodeRegion(ctx, key, codes) + } else { + e, err = db.extendNodeRegion(ctx, key, codes, prev) + } + if err != nil { + err = fmt.Errorf("nodes-region %s: %w", key, err) + log.Printf("[nodes-region] refresh failed: %v", err) + return nil, err + } + db.setNodeRegionEntry(key, e) + return e, nil +} + +// buildNodeRegion scans the full observation history for the region set. +// Full builds are serialised across region sets: each holds a pooled +// connection for seconds on a large database, and entries built together at +// startup fall due together. +func (db *DB) buildNodeRegion(ctx context.Context, key string, codes []string) (*nodeRegionEntry, error) { + db.nodeRegionFullMu.Lock() + defer db.nodeRegionFullMu.Unlock() + + wm, err := db.maxObservationID(ctx) + if err != nil { + return nil, err + } + start := time.Now() + set := make(map[string]struct{}) + if _, err := db.scanNodeRegionKeys(ctx, codes, 0, wm, false, set); err != nil { + return nil, err + } + log.Printf("[nodes-region] %s: full build, %d nodes in %v (obs id <= %d)", + key, len(set), time.Since(start).Round(time.Millisecond), wm) + return newNodeRegionEntry(set, wm, start) +} + +// extendNodeRegion adds nodes heard since prev's watermark. +func (db *DB) extendNodeRegion( + ctx context.Context, key string, codes []string, prev *nodeRegionEntry, +) (*nodeRegionEntry, error) { + wm, err := db.maxObservationID(ctx) + if err != nil { + return nil, err + } + set := make(map[string]struct{}, len(prev.keys)) + for _, k := range prev.keys { + set[k] = struct{}{} + } + start := time.Now() + added := 0 + if wm > prev.watermark { + added, err = db.scanNodeRegionKeys(ctx, codes, prev.watermark, wm, true, set) + if err != nil { + return nil, err + } + } + if added == 0 { + // Same membership: keep the encoded keys, advance the watermark. + e := *prev + e.watermark = max(wm, prev.watermark) + e.refreshed = time.Now() + return &e, nil + } + log.Printf("[nodes-region] %s: +%d nodes from obs %d..%d in %v", + key, added, prev.watermark+1, wm, time.Since(start).Round(time.Millisecond)) + return newNodeRegionEntry(set, wm, prev.built) +} + +// maxObservationID pins the upper bound of a scan. The ingestor is the single +// writer and observations.id is AUTOINCREMENT, so every id <= the returned +// value is committed and visible to a later read; newer rows are left for the +// next delta. +func (db *DB) maxObservationID(ctx context.Context) (int64, error) { + const q = "SELECT COALESCE(MAX(id), 0) FROM observations" + // Issue exactly one query: an unscanned *sql.Row keeps its connection + // checked out, and the pool is 4 (OpenDB). + row := db.conn.QueryRowContext(ctx, q) + var wm int64 + if err := row.Scan(&wm); err != nil { + return 0, fmt.Errorf("max observation id: %w", err) + } + return wm, nil +} + +// scanNodeRegionKeys adds to set every ADVERT from_pubkey observed by an +// observer in codes with lo < observations.id <= hi, returning how many were +// new. For a delta scan the join order is forced (CROSS JOIN) to drive from +// the observations rowid range; left to the planner, sqlite_stat1 (#2058) +// makes it start from the region's observers and walk their whole history. +func (db *DB) scanNodeRegionKeys( + ctx context.Context, codes []string, lo, hi int64, delta bool, set map[string]struct{}, +) (int, error) { + placeholders := make([]string, len(codes)) + args := []any{payloadTypeAdvert} + for i, c := range codes { + placeholders[i] = "?" + args = append(args, c) + } + args = append(args, lo, hi) + + joinCond := "obs.rowid = o.observer_idx" + if !db.isV3() { + joinCond = "obs.id = o.observer_id" + } + join := "JOIN" + if delta { + join = "CROSS JOIN" + } + // Only fixed fragments and "?" placeholders are formatted in; every value + // is bound. Use the indexed from_pubkey for current ADVERTs, but legacy + // rows can predate its asynchronous backfill. JSON_EXTRACT preserves the + // fork's existing region membership for those rows. + q := fmt.Sprintf(`SELECT DISTINCT COALESCE(t.from_pubkey, JSON_EXTRACT(t.decoded_json, '$.pubKey')) + FROM observations o + %[1]s transmissions t ON t.id = o.transmission_id + %[1]s observers obs ON %[2]s + WHERE t.payload_type = ? + AND COALESCE(t.from_pubkey, JSON_EXTRACT(t.decoded_json, '$.pubKey')) IS NOT NULL + AND UPPER(TRIM(obs.iata)) IN (%[3]s) + AND o.id > ? AND o.id <= ?`, join, joinCond, strings.Join(placeholders, ",")) + + rows, err := db.conn.QueryContext(ctx, q, args...) + if err != nil { + return 0, fmt.Errorf("scan region keys: %w", err) + } + defer rows.Close() + added := 0 + for rows.Next() { + var pk string + if err := rows.Scan(&pk); err != nil { + return added, fmt.Errorf("scan region keys: %w", err) + } + if _, ok := set[pk]; !ok { + set[pk] = struct{}{} + added++ + } + } + if err := rows.Err(); err != nil { + return added, fmt.Errorf("scan region keys: %w", err) + } + return added, nil +} diff --git a/cmd/server/nodes_region_cache_port_test.go b/cmd/server/nodes_region_cache_port_test.go new file mode 100644 index 000000000..eb86d0a9a --- /dev/null +++ b/cmd/server/nodes_region_cache_port_test.go @@ -0,0 +1,119 @@ +package main + +import ( + "fmt" + "strings" + "sync/atomic" + "testing" + "time" +) + +func TestNodeRegionCacheReplacementAtCapacityKeepsOtherKeys(t *testing.T) { + db := &DB{} + for i := range nodeRegionMaxEntries { + key := fmt.Sprintf("R%02d", i) + db.setNodeRegionEntry(key, &nodeRegionEntry{keysJSON: `["` + key + `"]`}) + } + db.setNodeRegionEntry("R00", &nodeRegionEntry{keysJSON: `["updated"]`}) + if got := len(db.nodeRegionCache); got != nodeRegionMaxEntries { + t.Fatalf("replacing one of %d entries left %d", nodeRegionMaxEntries, got) + } + if got := db.getNodeRegionEntry("R01"); got == nil || got.keysJSON != `["R01"]` { + t.Fatalf("unrelated cache entry lost on replacement: %#v", got) + } + if got := db.getNodeRegionEntry("R00").keysJSON; got != `["updated"]` { + t.Fatalf("replacement missing: %q", got) + } +} + +func TestNodeRegionCachePortMembershipAndReuse(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + seedTestData(t, db) + var scans atomic.Int32 + db.nodeRegionQueryHook = func() { scans.Add(1) } + for _, region := range []string{"SJC", "sjc", "SJC,SJC"} { + nodes, total, _, err := db.GetNodes(10, 0, "", "", "", "", "", region) + if err != nil { + t.Fatal(err) + } + if total != 1 || len(nodes) != 1 || nodes[0]["public_key"] != "aabbccdd11223344" { + t.Fatalf("region %q: total=%d nodes=%v", region, total, nodes) + } + } + if got := scans.Load(); got != 1 { + t.Fatalf("same region should scan once, got %d", got) + } +} + +func TestNodeRegionCachePortLegacyBackfillAndRefresh(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + seedTestData(t, db) + _, _, _, err := db.GetNodes(10, 0, "", "", "", "", "", "SJC") + if err != nil { + t.Fatal(err) + } + now := time.Now().UTC().Format(time.RFC3339) + _, err = db.conn.Exec(`INSERT INTO nodes(public_key,name,role,last_seen,first_seen) VALUES ('legacykey','Legacy','repeater',?,?)`, now, now) + if err != nil { + t.Fatal(err) + } + _, err = db.conn.Exec(`INSERT INTO transmissions(raw_hex,hash,first_seen,route_type,payload_type,decoded_json) VALUES ('AA','legacy-hash',?,1,4,'{"pubKey":"legacykey"}')`, now) + if err != nil { + t.Fatal(err) + } + // Simulate an old row not yet processed by the asynchronous backfill. + _, err = db.conn.Exec(`UPDATE transmissions SET from_pubkey=NULL WHERE hash='legacy-hash'`) + if err != nil { + t.Fatal(err) + } + _, err = db.conn.Exec(`INSERT INTO observations(transmission_id,observer_idx,timestamp) VALUES ((SELECT id FROM transmissions WHERE hash='legacy-hash'),1,?)`, time.Now().Unix()) + if err != nil { + t.Fatal(err) + } + prev := db.getNodeRegionEntry("SJC") + aged := *prev + aged.refreshed = time.Now().Add(-2 * nodeRegionFreshTTL) + db.setNodeRegionEntry("SJC", &aged) + if _, err := db.refreshNodeRegion("SJC", []string{"SJC"}); err != nil { + t.Fatal(err) + } + nodes, total, _, err := db.GetNodes(10, 0, "", "", "", "", "", "SJC") + if err != nil { + t.Fatal(err) + } + if total != 2 || len(nodes) != 2 { + t.Fatalf("legacy row missing after delta refresh: total=%d nodes=%v", total, nodes) + } + if !strings.Contains(db.getNodeRegionEntry("SJC").keysJSON, "legacykey") { + t.Fatal("legacy key absent from cached set") + } +} + +func TestNodeRegionCachePortFullRebuildDropsPrunedRegion(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + seedTestData(t, db) + if _, _, _, err := db.GetNodes(10, 0, "", "", "", "", "", "SJC"); err != nil { + t.Fatal(err) + } + if _, err := db.conn.Exec(`DELETE FROM observations WHERE observer_idx=1`); err != nil { + t.Fatal(err) + } + prev := db.getNodeRegionEntry("SJC") + aged := *prev + aged.refreshed = time.Now().Add(-2 * nodeRegionFreshTTL) + aged.built = time.Now().Add(-2 * nodeRegionRebuildInterval) + db.setNodeRegionEntry("SJC", &aged) + if _, err := db.refreshNodeRegion("SJC", []string{"SJC"}); err != nil { + t.Fatal(err) + } + nodes, total, _, err := db.GetNodes(10, 0, "", "", "", "", "", "SJC") + if err != nil { + t.Fatal(err) + } + if total != 0 || len(nodes) != 0 { + t.Fatalf("full rebuild retained pruned membership: total=%d nodes=%v", total, nodes) + } +} From a8541c0970d5c3cd07c52bf9dc2a92f9c985d245 Mon Sep 17 00:00:00 2001 From: dborup Date: Wed, 7 Oct 2026 13:40:46 +0200 Subject: [PATCH 2/2] fix(nodes): bound, refresh and prove the region membership cache MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review follow-ups on the #2102 port. 1. Keep #38's from_pubkey contract inside the cache. The port had reintroduced COALESCE(t.from_pubkey, JSON_EXTRACT(decoded_json, '$.pubKey')), which #38 measured and rejected: it rescues no rows, and one corrupt ADVERT fails the whole query. Behind a cache that is worse than before, because the failed refresh leaves the entry stale for every later request too, not just the one that triggered it. #38's comment block moves to the query it documents. 2. Bound the cache by evicting one entry, not by clearing the map. The old setNodeRegionEntry replaced the whole map when a 33rd region set arrived, so a client cycling through region sets — ?region= is a query parameter — emptied the cache on every request and every following request paid a full observation scan, serialised behind nodeRegionFullMu. Entries now evict least-recently-used, and a region set too long to key is stored under its SHA-256 digest (channelListMaxKeyBytes's rule). 3. Make the documented freshness the delivered freshness. An entry older than nodeRegionMaxStale is no longer served blind: the caller waits for its refresh, which at that age is a full rebuild. Without the wait a refresh that keeps failing serves membership of unbounded age with only a log line to show it. A failing refresh still falls back to the stale entry rather than 500-ing /api/nodes?region=. Tests: LRU eviction and the bound under 4x churn, the digest key, an addition visible within nodeRegionFreshTTL, removal by retention and by an observer IATA change gone within nodeRegionRebuildInterval and not by the delta scan, the over-stale wait and its fallback, and the cached result set against the uncached subquery for six region sets before and after new ADVERTs. The port's legacy-backfill test is replaced by one that locks #38's contract on the delta path. nodeRegionQueryHook now returns an error so a test can fail a scan. Perf (rule 0), nodes_region_cache_perf_176_test.go: 498MB / 1.5M observations / 749k transmissions, ANALYZE run; one Nodes page load with a region selected is 4 GetNodes pages of 500, 7 interleaved rounds of the same test file compiled against master and against this branch. Warm page load median 3729ms -> 31ms (121x; per GetNodes call 932ms -> 7.7ms, master's 932ms matching the 1-1.4s a region call costs on prod). Cold page load median 4292ms -> 705ms (6.1x): the one full membership scan replaces eight evaluations of the subquery. --- cmd/server/db.go | 18 +- cmd/server/nodes_region_cache.go | 162 ++++++- cmd/server/nodes_region_cache_176_test.go | 447 ++++++++++++++++++ .../nodes_region_cache_perf_176_test.go | 254 ++++++++++ cmd/server/nodes_region_cache_port_test.go | 73 +-- 5 files changed, 901 insertions(+), 53 deletions(-) create mode 100644 cmd/server/nodes_region_cache_176_test.go create mode 100644 cmd/server/nodes_region_cache_perf_176_test.go diff --git a/cmd/server/db.go b/cmd/server/db.go index 2e5e179db..7527de093 100644 --- a/cmd/server/db.go +++ b/cmd/server/db.go @@ -120,12 +120,18 @@ type DB struct { schemaProbeHook func() error // Region-membership cache for GetNodes (see nodes_region_cache.go). - nodeRegionCacheMu sync.Mutex - nodeRegionCache map[string]*nodeRegionEntry - nodeRegionLRU int64 - nodeRegionSF singleflight.Group - nodeRegionFullMu sync.Mutex - nodeRegionQueryHook func() + // nodeRegionUsed holds the LRU tick of each cached key and is written + // together with nodeRegionCache under nodeRegionCacheMu. + nodeRegionCacheMu sync.Mutex + nodeRegionCache map[string]*nodeRegionEntry + nodeRegionUsed map[string]int64 + nodeRegionTick int64 + nodeRegionSF singleflight.Group + nodeRegionFullMu sync.Mutex + // nodeRegionQueryHook (test seam, nil in production) runs once per real + // membership scan, before the query; a non-nil error fails that refresh + // the way a failing scan does. + nodeRegionQueryHook func() error } // channelRows is the part of *sql.Rows the channel list scans use. diff --git a/cmd/server/nodes_region_cache.go b/cmd/server/nodes_region_cache.go index dc549a5a9..d7f5bd719 100644 --- a/cmd/server/nodes_region_cache.go +++ b/cmd/server/nodes_region_cache.go @@ -2,6 +2,8 @@ package main import ( "context" + "crypto/sha256" + "encoding/hex" "encoding/json" "errors" "fmt" @@ -33,14 +35,40 @@ import ( // Only the first request for a region set waits on a full scan, and // concurrent first requests share it (the #1910 singleflight pattern, as // channelsSF). +// +// Two bounds make that safe to expose to a query parameter (AGENTS.md rule 0, +// "no unbounded data structures"): +// +// - at most nodeRegionMaxEntries entries, kept by evicting the least +// recently used one (setNodeRegionEntry), and a stored key longer than +// nodeRegionMaxKeyBytes is replaced by its digest (nodeRegionKey); +// - an entry older than nodeRegionMaxStale is not served blind: the caller +// waits for its refresh, so the ages promised above are the ages +// actually served even when refreshes keep failing. const ( nodeRegionFreshTTL = 30 * time.Second nodeRegionRebuildInterval = 30 * time.Minute // nodeRegionMaxEntries bounds the cache. Entries hold up to a few thousand - // pubkeys each, so this is far below the shared maxCacheEntries; in - // practice an instance sees a handful of distinct region sets. + // pubkeys each (one per node in the region set, so the nodes table caps + // an entry), so this is far below the shared maxCacheEntries; in practice + // an instance sees a handful of distinct region sets. nodeRegionMaxEntries = 32 + // nodeRegionMaxKeyBytes caps a stored key, as channelListMaxKeyBytes does + // for the channel list cache: ?region= is caller-controlled and + // normalizeRegionCodes does not limit how many codes it accepts, so a + // long set is stored under "sha256:" + 64 hex digits instead. Normalized + // codes are upper-case, so no short key can collide with a digest. + nodeRegionMaxKeyBytes = 256 + // nodeRegionMaxStale bounds the age of a served entry. Below it a stale + // entry is served while one refresh runs in the background — the fast + // path every page uses. At or above it the caller waits for the refresh, + // which is a full rebuild at this age, so the documented bounds hold: + // without the wait a refresh that keeps failing serves membership of + // unbounded age with only a log line to show it. The cost is that the + // first region request after nodeRegionMaxStale with no region traffic + // pays one full scan, exactly as a cold start does. + nodeRegionMaxStale = nodeRegionRebuildInterval ) // errNodeRegionFlightResult means nodeRegionSF returned something other than @@ -86,7 +114,10 @@ func newNodeRegionEntry(set map[string]struct{}, watermark int64, built time.Tim } // nodeRegionKey canonicalises normalised region codes so "EDI,GLA", -// "gla, edi" and "EDI,EDI,GLA" share one cache entry. +// "gla, edi" and "EDI,EDI,GLA" share one cache entry. A set long enough to +// bloat the map (or the log line) is keyed by its digest instead; the codes +// themselves are passed to the scan separately, so the key is only an +// identity. func nodeRegionKey(codes []string) string { seen := make(map[string]struct{}, len(codes)) uniq := make([]string, 0, len(codes)) @@ -97,23 +128,70 @@ func nodeRegionKey(codes []string) string { } } sort.Strings(uniq) - return strings.Join(uniq, ",") + key := strings.Join(uniq, ",") + if len(key) <= nodeRegionMaxKeyBytes { + return key + } + sum := sha256.Sum256([]byte(key)) + return "sha256:" + hex.EncodeToString(sum[:]) } +// getNodeRegionEntry returns the entry for key and marks it as the most +// recently used one, so eviction drops the sets nobody is paging through. func (db *DB) getNodeRegionEntry(key string) *nodeRegionEntry { db.nodeRegionCacheMu.Lock() defer db.nodeRegionCacheMu.Unlock() - return db.nodeRegionCache[key] + e := db.nodeRegionCache[key] + if e != nil { + db.touchNodeRegionLocked(key) + } + return e } +// setNodeRegionEntry stores e, evicting the least recently used entry when +// the cache is full. +// +// It must evict exactly one entry, never clear the map: ?region= is +// caller-controlled, so a client cycling through more than +// nodeRegionMaxEntries distinct region sets would otherwise drop the entries +// the real pages depend on, and every following request would pay a full +// observation scan — serialised behind nodeRegionFullMu, seconds each on a +// large database. That is the one thing this cache exists to prevent. func (db *DB) setNodeRegionEntry(key string, e *nodeRegionEntry) { db.nodeRegionCacheMu.Lock() defer db.nodeRegionCacheMu.Unlock() - _, replacing := db.nodeRegionCache[key] - if db.nodeRegionCache == nil || (!replacing && len(db.nodeRegionCache) >= nodeRegionMaxEntries) { - db.nodeRegionCache = make(map[string]*nodeRegionEntry) + if db.nodeRegionCache == nil { + db.nodeRegionCache = make(map[string]*nodeRegionEntry, nodeRegionMaxEntries) + db.nodeRegionUsed = make(map[string]int64, nodeRegionMaxEntries) + } + if _, replacing := db.nodeRegionCache[key]; !replacing { + for len(db.nodeRegionCache) >= nodeRegionMaxEntries { + victim, oldest, found := "", int64(0), false + for k, used := range db.nodeRegionUsed { + if !found || used < oldest { + victim, oldest, found = k, used, true + } + } + if !found { + break // cannot happen: the two maps are written together + } + delete(db.nodeRegionCache, victim) + delete(db.nodeRegionUsed, victim) + } } db.nodeRegionCache[key] = e + db.touchNodeRegionLocked(key) +} + +// touchNodeRegionLocked records key as the most recent use. Caller holds +// nodeRegionCacheMu. The tick is a counter, not a clock, so entries touched +// inside one coarse clock tick still order. +func (db *DB) touchNodeRegionLocked(key string) { + db.nodeRegionTick++ + if db.nodeRegionUsed == nil { + db.nodeRegionUsed = make(map[string]int64, nodeRegionMaxEntries) + } + db.nodeRegionUsed[key] = db.nodeRegionTick } // nodeRegionKeysJSON returns the JSON array of node public keys heard in the @@ -123,15 +201,35 @@ func (db *DB) nodeRegionKeysJSON(codes []string) (string, error) { key := nodeRegionKey(codes) if e := db.getNodeRegionEntry(key); e != nil { - if !e.fresh() { + age := time.Since(e.refreshed) + switch { + case age < nodeRegionFreshTTL: + return e.keysJSON, nil + case age < nodeRegionMaxStale: // Stale-while-revalidate: kick one refresh, don't wait on it. // DoChan dedups against an in-flight refresh for the same key, // and its buffered result channel may be left unread. db.nodeRegionSF.DoChan(key, func() (any, error) { return db.refreshNodeRegion(key, codes) }) + return e.keysJSON, nil + } + // Past nodeRegionMaxStale the entry is too old to serve blind, so + // wait for the refresh (a full rebuild at this age). A refresh that + // fails still falls back to the stale entry: an outage of the + // observation scan must not turn /api/nodes?region= into a 500. + v, err, _ := db.nodeRegionSF.Do(key, func() (any, error) { + return db.refreshNodeRegion(key, codes) + }) + if err != nil { + log.Printf("[nodes-region] %s: serving entry %v old, refresh failed: %v", + key, age.Round(time.Second), err) + return e.keysJSON, nil } - return e.keysJSON, nil + if fresh, ok := v.(*nodeRegionEntry); ok { + return fresh.keysJSON, nil + } + return "", fmt.Errorf("nodes-region %s: %w: %T", key, errNodeRegionFlightResult, v) } v, err, _ := db.nodeRegionSF.Do(key, func() (any, error) { @@ -158,7 +256,9 @@ func (db *DB) refreshNodeRegion(key string, codes []string) (*nodeRegionEntry, e return prev, nil // a flight that finished just before this one got here } if db.nodeRegionQueryHook != nil { - db.nodeRegionQueryHook() + if err := db.nodeRegionQueryHook(); err != nil { + return nil, fmt.Errorf("nodes-region %s: %w", key, err) + } } // A scan still running when the next rebuild would be due is abandoned @@ -277,15 +377,45 @@ func (db *DB) scanNodeRegionKeys( join = "CROSS JOIN" } // Only fixed fragments and "?" placeholders are formatted in; every value - // is bound. Use the indexed from_pubkey for current ADVERTs, but legacy - // rows can predate its asynchronous backfill. JSON_EXTRACT preserves the - // fork's existing region membership for those rows. - q := fmt.Sprintf(`SELECT DISTINCT COALESCE(t.from_pubkey, JSON_EXTRACT(t.decoded_json, '$.pubKey')) + // is bound. + // + // #1143 / PR #38: from_pubkey is a dedicated column the ingestor fills at + // write time for ADVERT rows, so this lookup does not have to + // JSON_EXTRACT and parse decoded_json for every candidate row. Measured + // against a live 4.9GB database, 9 interleaved rounds so concurrent + // ingest cannot bias one variant: 934ms -> 835ms median (1.12x) for a + // single-region filter, identical result set both ways. The win is purely + // the avoided JSON parse — EXPLAIN QUERY PLAN is byte-identical before + // and after, both driving off idx_transmissions_payload_type. Despite the + // column being indexed, idx_transmissions_from_pubkey does not + // participate in this plan at all. + // + // THE TRAP, for whoever touches this next: from_pubkey is only written + // for ADVERTs that actually carry a pubkey — see the guard in + // cmd/ingestor/db.go. A NULL here drops out of IN (...) silently, so a + // write path that fills decoded_json but forgets from_pubkey would cost + // nodes from this filter with no error to notice. That is safe today only + // because it is measured to be: on live staging 25 of 62104 ADVERTs have + // NULL from_pubkey, and all 25 have no pubKey in decoded_json either — + // nothing is lost, because there was never a pubkey to record. A COALESCE + // fallback to JSON_EXTRACT was measured and rejected: it rescued 0 rows, + // and it parses decoded_json again for every NULL row, so one corrupt + // ADVERT that the backfill has not reached yet fails the whole query with + // "malformed JSON" (a 500 on /api/nodes?region=). Behind this cache that + // is worse, not better: the failed refresh leaves the cached entry stale + // for every later request too. Its speed is not the reason: on a 4.3GB + // synthetic DB it times within noise of from_pubkey. If you add a new + // write path for transmissions, populate from_pubkey there rather than + // reintroducing a per-row JSON parse here. nodes_region_from_pubkey_test.go + // locks the result set against the JSON_EXTRACT expression; + // TestBackfillFromPubkey_* in cmd/ingestor locks the backfill's half of + // the contract. + q := fmt.Sprintf(`SELECT DISTINCT t.from_pubkey FROM observations o %[1]s transmissions t ON t.id = o.transmission_id %[1]s observers obs ON %[2]s WHERE t.payload_type = ? - AND COALESCE(t.from_pubkey, JSON_EXTRACT(t.decoded_json, '$.pubKey')) IS NOT NULL + AND t.from_pubkey IS NOT NULL AND UPPER(TRIM(obs.iata)) IN (%[3]s) AND o.id > ? AND o.id <= ?`, join, joinCond, strings.Join(placeholders, ",")) diff --git a/cmd/server/nodes_region_cache_176_test.go b/cmd/server/nodes_region_cache_176_test.go new file mode 100644 index 000000000..494ecf381 --- /dev/null +++ b/cmd/server/nodes_region_cache_176_test.go @@ -0,0 +1,447 @@ +package main + +// PR #176 review follow-ups for the region membership cache: +// +// - the cache is bounded by evicting one entry, not by clearing the map, +// and a caller-supplied region set cannot grow a stored key without +// limit (AGENTS.md rule 0: no unbounded data structures); +// - the freshness contract the cache documents is the freshness it +// delivers: additions within nodeRegionFreshTTL, removals and observer +// IATA changes within nodeRegionRebuildInterval, and nothing older than +// nodeRegionMaxStale served at all; +// - a cached region set returns exactly the nodes the uncached subquery +// returns, for several region sets. + +import ( + "fmt" + "slices" + "sort" + "strings" + "testing" + "time" +) + +// ageNodeRegionEntry rewrites the cached entry for key with its refresh and +// build times pushed back, so a test can stand where a 30-second-old or +// 30-minute-old entry stands without waiting. The entry is immutable once +// stored, so this stores a copy. +func ageNodeRegionEntry(t *testing.T, db *DB, key string, refreshedAgo, builtAgo time.Duration) *nodeRegionEntry { + t.Helper() + prev := db.getNodeRegionEntry(key) + if prev == nil { + t.Fatalf("no cache entry for %q to age", key) + } + aged := *prev + aged.refreshed = time.Now().Add(-refreshedAgo) + if builtAgo > 0 { + aged.built = time.Now().Add(-builtAgo) + } + db.setNodeRegionEntry(key, &aged) + return &aged +} + +// waitForNodeRegionRefresh waits for the background refresh kicked by a +// stale read to replace the entry, and fails if it never does. +func waitForNodeRegionRefresh(t *testing.T, db *DB, key string, was time.Time) { + t.Helper() + deadline := time.Now().Add(5 * time.Second) + for time.Now().Before(deadline) { + if e := db.getNodeRegionEntry(key); e != nil && e.refreshed.After(was) { + return + } + time.Sleep(2 * time.Millisecond) + } + t.Fatalf("no background refresh of %q within 5s of a stale read", key) +} + +func nodeRegionPublicKeys(t *testing.T, db *DB, region string) []string { + t.Helper() + nodes, total, _, err := db.GetNodes(10000, 0, "", "", "", "", "", region) + if err != nil { + t.Fatalf("GetNodes(region=%q): %v", region, err) + } + if total != len(nodes) { + t.Fatalf("GetNodes(region=%q): total=%d but %d rows", region, total, len(nodes)) + } + out := make([]string, 0, len(nodes)) + for _, n := range nodes { + pk, _ := n["public_key"].(string) + out = append(out, pk) + } + sort.Strings(out) + return out +} + +// ---------------------------------------------------------------- bounded + +// Filling the cache past nodeRegionMaxEntries must evict exactly one entry — +// the least recently used one — and leave every other entry in place. +// +// The bound used to be enforced by replacing the whole map, so the 33rd +// distinct region set dropped the 32 before it. ?region= is a query +// parameter: a client walking 33 region sets would have emptied the cache on +// every request, and every request after it would have paid a full +// observation scan (~10s on a 1.4GB database), serialised behind +// nodeRegionFullMu. That is the exact collapse this cache exists to prevent, +// reachable from the outside. +func TestNodeRegionCacheEvictsOneLRUEntryAtCapacity_176(t *testing.T) { + db := &DB{} + for i := range nodeRegionMaxEntries { + key := fmt.Sprintf("R%02d", i) + db.setNodeRegionEntry(key, &nodeRegionEntry{keysJSON: `["` + key + `"]`}) + } + // R00 is the least recently *written*; read it so R01 becomes the victim. + // That also locks the LRU order against a plain insertion-order policy. + if got := db.getNodeRegionEntry("R00"); got == nil { + t.Fatal("R00 missing before eviction") + } + + db.setNodeRegionEntry("NEW", &nodeRegionEntry{keysJSON: `["NEW"]`}) + + if got := len(db.nodeRegionCache); got != nodeRegionMaxEntries { + t.Fatalf("cache holds %d entries, want the bound %d", got, nodeRegionMaxEntries) + } + if got := db.getNodeRegionEntry("R01"); got != nil { + t.Errorf("least recently used entry R01 survived: %#v", got) + } + var lost []string + for i := range nodeRegionMaxEntries { + key := fmt.Sprintf("R%02d", i) + if key == "R01" { + continue + } + if e := db.nodeRegionCache[key]; e == nil || e.keysJSON != `["`+key+`"]` { + lost = append(lost, key) + } + } + if len(lost) > 0 { + t.Errorf("eviction of one entry also dropped %d others: %v", len(lost), lost) + } + if e := db.getNodeRegionEntry("NEW"); e == nil || e.keysJSON != `["NEW"]` { + t.Errorf("new entry not stored: %#v", e) + } + if got, want := len(db.nodeRegionUsed), len(db.nodeRegionCache); got != want { + t.Errorf("LRU bookkeeping holds %d keys for %d entries", got, want) + } +} + +// Repeated eviction keeps the bound: 4x the capacity in distinct region sets +// never grows the cache, and the most recent ones are the ones kept. +func TestNodeRegionCacheStaysBoundedUnderChurn_176(t *testing.T) { + db := &DB{} + const churn = 4 * nodeRegionMaxEntries + for i := range churn { + key := fmt.Sprintf("Q%03d", i) + db.setNodeRegionEntry(key, &nodeRegionEntry{keysJSON: `["` + key + `"]`}) + if got := len(db.nodeRegionCache); got > nodeRegionMaxEntries { + t.Fatalf("after %d inserts the cache holds %d entries, bound is %d", + i+1, got, nodeRegionMaxEntries) + } + } + if got := len(db.nodeRegionCache); got != nodeRegionMaxEntries { + t.Fatalf("cache holds %d entries after churn, want %d", got, nodeRegionMaxEntries) + } + for i := churn - nodeRegionMaxEntries; i < churn; i++ { + key := fmt.Sprintf("Q%03d", i) + if db.nodeRegionCache[key] == nil { + t.Errorf("most recent entry %s was evicted", key) + } + } +} + +// normalizeRegionCodes does not limit how many codes ?region= may carry, so +// the stored key must be bounded on its own. A long set is keyed by its +// digest; distinct long sets still get distinct keys, and the same long set +// still shares one entry. +func TestNodeRegionCacheLongRegionSetKeyIsBounded_176(t *testing.T) { + long := make([]string, 0, 2000) + for i := range 2000 { + long = append(long, fmt.Sprintf("X%04d", i)) + } + key := nodeRegionKey(long) + if len(key) > nodeRegionMaxKeyBytes { + t.Fatalf("key for %d codes is %d bytes, bound is %d", len(long), len(key), nodeRegionMaxKeyBytes) + } + if !strings.HasPrefix(key, "sha256:") { + t.Fatalf("long key %q is not a digest", key) + } + // Order and duplicates still collapse to one identity. + shuffled := slices.Clone(long) + slices.Reverse(shuffled) + shuffled = append(shuffled, long[0], long[1]) + if got := nodeRegionKey(shuffled); got != key { + t.Errorf("same region set keyed twice: %q vs %q", got, key) + } + other := nodeRegionKey(append(slices.Clone(long), "ZZZZ")) + if other == key { + t.Error("two different long region sets share one cache key") + } + // A short set is still stored verbatim, so the log line and the tests + // above keep reading region codes. + if got := nodeRegionKey([]string{"SJC", "SFO"}); got != "SFO,SJC" { + t.Errorf("short key = %q, want SFO,SJC", got) + } +} + +// ---------------------------------------------------------------- freshness + +// seedRegionFreshnessDB returns a v3 test DB with one observer in SJC, one +// in SFO, and one node heard in SJC. +func seedRegionFreshnessDB(t *testing.T) *DB { + t.Helper() + db := setupTestDB(t) + t.Cleanup(func() { db.Close() }) + mustExec := func(q string, args ...any) { + t.Helper() + if _, err := db.conn.Exec(q, args...); err != nil { + t.Fatalf("%s: %v", q, err) + } + } + mustExec(`INSERT INTO observers (rowid, id, name, iata) VALUES (1,'obs-sjc','SJC obs','SJC'), (2,'obs-sfo','SFO obs','SFO')`) + regionSeedNode(t, db, "aa00000000000001", 1) + return db +} + +// regionSeedNode inserts a node, a stamped ADVERT for it and one observation +// by the given observer rowid, and returns the node's public key. +func regionSeedNode(t *testing.T, db *DB, pubkey string, observerIdx int) string { + t.Helper() + now := time.Now().UTC().Format(time.RFC3339) + mustExec := func(q string, args ...any) { + t.Helper() + if _, err := db.conn.Exec(q, args...); err != nil { + t.Fatalf("%s: %v", q, err) + } + } + mustExec(`INSERT INTO nodes (public_key, name, role, last_seen, first_seen) + VALUES (?, ?, 'repeater', ?, ?)`, pubkey, "n-"+pubkey[:4], now, now) + mustExec(`INSERT INTO transmissions (raw_hex, hash, first_seen, route_type, payload_type, decoded_json, from_pubkey) + VALUES ('00', ?, ?, 1, 4, ?, ?)`, "adv-"+pubkey, now, + `{"type":"ADVERT","pubKey":"`+pubkey+`"}`, pubkey) + mustExec(`INSERT INTO observations (transmission_id, observer_idx, timestamp) + VALUES ((SELECT id FROM transmissions WHERE hash=?), ?, ?)`, "adv-"+pubkey, observerIdx, time.Now().Unix()) + return pubkey +} + +// A node first heard after the entry was built appears within +// nodeRegionFreshTTL: the read at the TTL boundary is still answered from the +// cache (that is the contract), it kicks exactly one refresh, and the next +// read has the new node. Nothing waits longer than the promised 30 seconds +// plus one delta scan. +func TestNodeRegionCacheAdditionVisibleWithinFreshTTL_176(t *testing.T) { + db := seedRegionFreshnessDB(t) + if got := nodeRegionPublicKeys(t, db, "SJC"); !slices.Equal(got, []string{"aa00000000000001"}) { + t.Fatalf("cold SJC = %v", got) + } + + newKey := regionSeedNode(t, db, "bb00000000000002", 1) + + // Inside the TTL the cached set is served unchanged — the cache is real. + if got := nodeRegionPublicKeys(t, db, "SJC"); slices.Contains(got, newKey) { + t.Fatalf("fresh entry already reflects a node added after it was built: %v", got) + } + + aged := ageNodeRegionEntry(t, db, "SJC", nodeRegionFreshTTL, 0) + if got := nodeRegionPublicKeys(t, db, "SJC"); slices.Contains(got, newKey) { + t.Fatalf("read at the TTL boundary blocked on a refresh instead of serving stale: %v", got) + } + waitForNodeRegionRefresh(t, db, "SJC", aged.refreshed) + + got := nodeRegionPublicKeys(t, db, "SJC") + if !slices.Equal(got, []string{"aa00000000000001", newKey}) { + t.Fatalf("node added %v after the entry was built is still missing: %v", + nodeRegionFreshTTL, got) + } + if e := db.getNodeRegionEntry("SJC"); !e.fresh() { + t.Error("entry not fresh after its refresh landed") + } +} + +// A node removed from a region — its observations pruned by retention, or its +// observer's IATA changed — is gone within nodeRegionRebuildInterval and not +// before: the delta scan cannot see a removal, so the documented 30 minutes +// is what governs. Both removal shapes are checked. +func TestNodeRegionCacheRemovalGoneWithinRebuildInterval_176(t *testing.T) { + for _, tc := range []struct { + name string + remove func(t *testing.T, db *DB) + }{ + {"retention pruned the observations", func(t *testing.T, db *DB) { + if _, err := db.conn.Exec(`DELETE FROM observations WHERE observer_idx = 1`); err != nil { + t.Fatal(err) + } + }}, + {"observer moved to another region", func(t *testing.T, db *DB) { + if _, err := db.conn.Exec(`UPDATE observers SET iata='LAX' WHERE rowid=1`); err != nil { + t.Fatal(err) + } + }}, + } { + t.Run(tc.name, func(t *testing.T) { + db := seedRegionFreshnessDB(t) + if got := nodeRegionPublicKeys(t, db, "SJC"); len(got) != 1 { + t.Fatalf("cold SJC = %v", got) + } + tc.remove(t, db) + + // Stale, but not yet due a rebuild: the delta refresh runs and + // must not drop the node. This is the documented behaviour, and + // it is what makes the rebuild interval the real bound. + aged := ageNodeRegionEntry(t, db, "SJC", nodeRegionFreshTTL, nodeRegionRebuildInterval/2) + nodeRegionPublicKeys(t, db, "SJC") + waitForNodeRegionRefresh(t, db, "SJC", aged.refreshed) + if got := nodeRegionPublicKeys(t, db, "SJC"); len(got) != 1 { + t.Fatalf("delta refresh changed membership: %v", got) + } + + // At the rebuild interval the refresh is a full rebuild. + aged = ageNodeRegionEntry(t, db, "SJC", nodeRegionFreshTTL, nodeRegionRebuildInterval) + nodeRegionPublicKeys(t, db, "SJC") + waitForNodeRegionRefresh(t, db, "SJC", aged.refreshed) + if got := nodeRegionPublicKeys(t, db, "SJC"); len(got) != 0 { + t.Fatalf("node still in the region %v after it was removed: %v", + nodeRegionRebuildInterval, got) + } + }) + } +} + +// Past nodeRegionMaxStale an entry is not served blind: the caller waits for +// the refresh, so one call — not two — returns membership inside the +// documented age. Without the wait, membership of unbounded age is served +// whenever refreshes have been failing, with only a log line to show it. +func TestNodeRegionCacheDoesNotServeBeyondMaxStale_176(t *testing.T) { + db := seedRegionFreshnessDB(t) + if got := nodeRegionPublicKeys(t, db, "SJC"); len(got) != 1 { + t.Fatalf("cold SJC = %v", got) + } + newKey := regionSeedNode(t, db, "bb00000000000002", 1) + ageNodeRegionEntry(t, db, "SJC", nodeRegionMaxStale, nodeRegionMaxStale) + + got := nodeRegionPublicKeys(t, db, "SJC") + if !slices.Contains(got, newKey) { + t.Fatalf("entry %v old was served without waiting for its refresh: %v", + nodeRegionMaxStale, got) + } +} + +// ...and a refresh that fails does not turn an over-stale entry into a 500: +// the stale set is served, the failure is logged, and the scan was attempted +// before the call returned. +func TestNodeRegionCacheOverStaleFallsBackWhenRefreshFails_176(t *testing.T) { + db := seedRegionFreshnessDB(t) + want := nodeRegionPublicKeys(t, db, "SJC") + if len(want) != 1 { + t.Fatalf("cold SJC = %v", want) + } + regionSeedNode(t, db, "bb00000000000002", 1) + ageNodeRegionEntry(t, db, "SJC", nodeRegionMaxStale, nodeRegionMaxStale) + + scans := 0 + db.nodeRegionQueryHook = func() error { scans++; return fmt.Errorf("observation scan unavailable") } + got := nodeRegionPublicKeys(t, db, "SJC") + if scans != 1 { + t.Errorf("over-stale read attempted %d scans before returning, want 1", scans) + } + if !slices.Equal(got, want) { + t.Errorf("failed refresh of an over-stale entry returned %v, want the stale set %v", got, want) + } +} + +// ---------------------------------------------------------------- oracle + +// regionNodesUncached is the oracle: the node set the region filter returned +// before this cache existed, built from the uncached subquery (#38's form) so +// the cache is compared to something it cannot be derived from. +func regionNodesUncached(t *testing.T, db *DB, region string) []string { + t.Helper() + codes := normalizeRegionCodes(region) + if len(codes) == 0 { + t.Fatalf("region %q normalizes to no codes", region) + } + joinCond := "obs.rowid = o.observer_idx" + if !db.isV3() { + joinCond = "obs.id = o.observer_id" + } + args := make([]any, len(codes)) + for i, c := range codes { + args[i] = c + } + q := `SELECT public_key FROM nodes WHERE public_key IN ( + SELECT DISTINCT t.from_pubkey + FROM transmissions t + JOIN observations o ON o.transmission_id = t.id + JOIN observers obs ON ` + joinCond + ` + WHERE t.payload_type = 4 + AND UPPER(TRIM(obs.iata)) IN (?` + strings.Repeat(",?", len(codes)-1) + `) + )` + rows, err := db.conn.Query(q, args...) + if err != nil { + t.Fatalf("oracle query for %q: %v", region, err) + } + defer rows.Close() + var out []string + for rows.Next() { + var pk string + if err := rows.Scan(&pk); err != nil { + t.Fatal(err) + } + out = append(out, pk) + } + if err := rows.Err(); err != nil { + t.Fatal(err) + } + sort.Strings(out) + return out +} + +// The cached region filter returns exactly the uncached node set, for single +// regions, multi-region sets, a set written in mixed case with padding and +// duplicates, and a region no observer carries — before and after new +// ADVERTs arrive. Every set is queried twice, so the second answer comes +// from the cache rather than a fresh scan. +func TestNodeRegionCacheMatchesUncachedForManyRegionSets_176(t *testing.T) { + db := seedRegionFreshnessDB(t) + regionSeedNode(t, db, "bb00000000000002", 2) + regionSeedNode(t, db, "cc00000000000003", 1) + // cc is heard in both regions. + if _, err := db.conn.Exec(`INSERT INTO observations (transmission_id, observer_idx, timestamp) + VALUES ((SELECT id FROM transmissions WHERE hash='adv-cc00000000000003'), 2, ?)`, time.Now().Unix()); err != nil { + t.Fatal(err) + } + + sets := []string{"SJC", "SFO", "SJC,SFO", " sfo , SJC , sjc ", "LAX", "LAX,SJC"} + check := func(stage string) { + t.Helper() + nonEmpty := 0 + for _, region := range sets { + want := regionNodesUncached(t, db, region) + for pass := range 2 { + got := nodeRegionPublicKeys(t, db, region) + if !slices.Equal(got, want) { + t.Errorf("%s: region %q pass %d: cached %v, uncached %v", + stage, region, pass, got, want) + } + } + if len(want) > 0 { + nonEmpty++ + } + } + if nonEmpty < 3 { + t.Fatalf("%s: only %d region sets matched any node; the comparison is vacuous", + stage, nonEmpty) + } + } + check("initial") + + // New ADVERTs land, every entry is taken past its TTL, and each set's + // refresh is awaited before it is compared again. + regionSeedNode(t, db, "dd00000000000004", 2) + for _, region := range sets { + key := nodeRegionKey(normalizeRegionCodes(region)) + aged := ageNodeRegionEntry(t, db, key, 2*nodeRegionFreshTTL, 0) + nodeRegionPublicKeys(t, db, region) + waitForNodeRegionRefresh(t, db, key, aged.refreshed) + } + check("after new adverts") +} diff --git a/cmd/server/nodes_region_cache_perf_176_test.go b/cmd/server/nodes_region_cache_perf_176_test.go new file mode 100644 index 000000000..dec49164d --- /dev/null +++ b/cmd/server/nodes_region_cache_perf_176_test.go @@ -0,0 +1,254 @@ +package main + +// Perf proof for PR #176 (AGENTS.md rule 0: a perf claim needs data). +// +// Not a pass/fail test. It uses only identifiers that exist on master +// (OpenDB, DB.GetNodes, DB.Close), so the same file compiles against master +// and against this branch and the two binaries can be run alternately — the +// before/after difference is then the region cache itself, not a rewritten +// measurement: +// +// # one large database, reused by both binaries +// CORESCOPE_PERF_176_GEN=/path/big.db CORESCOPE_PERF_176_OBS=1500000 \ +// go test -run TestPerf_GenerateRegionDB_176 -timeout 2h . +// +// go test -c -o after.test . # on this branch +// go test -c -o before.test . # on master, same file copied in +// for i in 1 2 3 4 5 6 7; do +// CORESCOPE_PERF_176=1 CORESCOPE_PERF_176_DB=/path/big.db ./before.test -test.run TestPerf_NodesRegionPageLoad_176 +// CORESCOPE_PERF_176=1 CORESCOPE_PERF_176_DB=/path/big.db ./after.test -test.run TestPerf_NodesRegionPageLoad_176 +// done +// +// The measured unit is one Nodes page load with a region selected: what +// public/nodes.js does in fetchAllNodes(), four pages of 500. GetNodes +// evaluates the region predicate twice per call (COUNT(*) and the page), so +// the uncached shape pays for it eight times per page load. + +import ( + "database/sql" + "fmt" + "math/rand" + "os" + "sort" + "strconv" + "testing" + "time" +) + +const ( + perf176PageSize = 500 + perf176Pages = 4 + perf176Region = "SJC" +) + +func perf176Median(d []time.Duration) time.Duration { + sort.Slice(d, func(i, j int) bool { return d[i] < d[j] }) + return d[len(d)/2] +} + +// perf176Regions are the observer regions the generated database uses. SJC +// (the measured one) holds a quarter of the observers. +var perf176Regions = []string{"SJC", "SJC", "SJC", "SJC", "SFO", "OAK", "MRY", "LAX", "PDX", "SEA", "BUR", "SAN", "SMF", "RNO", "LAS", "PHX"} + +// TestPerf_GenerateRegionDB_176 writes a database shaped like a busy +// instance: many observations per ADVERT, most transmissions not ADVERTs, +// and ANALYZE run at the end so the planner sees sqlite_stat1 the way a live +// database does (#2058 — that is what makes the uncached subquery drive off +// the region's observers). +func TestPerf_GenerateRegionDB_176(t *testing.T) { + path := os.Getenv("CORESCOPE_PERF_176_GEN") + if path == "" { + t.Skip("set CORESCOPE_PERF_176_GEN= to generate the perf database") + } + obsTarget, _ := strconv.Atoi(os.Getenv("CORESCOPE_PERF_176_OBS")) + if obsTarget <= 0 { + obsTarget = 1500000 + } + if err := os.RemoveAll(path); err != nil { + t.Fatal(err) + } + conn, err := sql.Open("sqlite", "file:"+path+"?_journal_mode=WAL&_synchronous=OFF") + if err != nil { + t.Fatal(err) + } + defer conn.Close() + conn.SetMaxOpenConns(1) + + mustExec := func(q string, args ...any) { + t.Helper() + if _, err := conn.Exec(q, args...); err != nil { + t.Fatalf("%s: %v", q, err) + } + } + // The v3 schema of the columns this query touches, with the indexes + // dbschema.go creates for them. + for _, stmt := range []string{ + `CREATE TABLE nodes (public_key TEXT PRIMARY KEY, name TEXT, role TEXT, lat REAL, lon REAL, + last_seen TEXT, first_seen TEXT, advert_count INTEGER DEFAULT 0, battery_mv INTEGER, + temperature_c REAL, foreign_advert INTEGER DEFAULT 0)`, + `CREATE TABLE observers (id TEXT PRIMARY KEY, name TEXT, iata TEXT, last_seen TEXT, + first_seen TEXT, packet_count INTEGER DEFAULT 0, inactive INTEGER DEFAULT 0, + last_packet_at TEXT DEFAULT NULL)`, + `CREATE TABLE transmissions (id INTEGER PRIMARY KEY AUTOINCREMENT, raw_hex TEXT NOT NULL, + hash TEXT NOT NULL UNIQUE, first_seen TEXT NOT NULL, route_type INTEGER, + payload_type INTEGER, payload_version INTEGER, decoded_json TEXT, + channel_hash TEXT DEFAULT NULL, from_pubkey TEXT DEFAULT NULL, + created_at TEXT DEFAULT (datetime('now')))`, + `CREATE TABLE observations (id INTEGER PRIMARY KEY AUTOINCREMENT, + transmission_id INTEGER NOT NULL REFERENCES transmissions(id), observer_idx INTEGER, + direction TEXT, snr REAL, rssi REAL, score INTEGER, path_json TEXT, + timestamp INTEGER NOT NULL, resolved_path TEXT, raw_hex TEXT)`, + `CREATE INDEX idx_transmissions_payload_type ON transmissions(payload_type)`, + `CREATE INDEX idx_transmissions_from_pubkey ON transmissions(from_pubkey)`, + `CREATE INDEX idx_transmissions_hash ON transmissions(hash)`, + `CREATE INDEX idx_observations_timestamp ON observations(timestamp)`, + `CREATE INDEX idx_observations_transmission_id ON observations(transmission_id)`, + `CREATE INDEX idx_observations_observer_idx ON observations(observer_idx)`, + `CREATE INDEX idx_observations_tx_ts ON observations(transmission_id, timestamp)`, + `CREATE INDEX idx_nodes_last_seen ON nodes(last_seen)`, + } { + mustExec(stmt) + } + + const observers = 64 + const nodes = 4000 + rng := rand.New(rand.NewSource(1764)) + now := time.Now().UTC() + begin := func() *sql.Tx { + tx, err := conn.Begin() + if err != nil { + t.Fatal(err) + } + return tx + } + + tx := begin() + for i := range observers { + iata := perf176Regions[i%len(perf176Regions)] + if _, err := tx.Exec(`INSERT INTO observers (rowid, id, name, iata, last_seen, first_seen) + VALUES (?, ?, ?, ?, ?, ?)`, i+1, fmt.Sprintf("obs-%03d", i), fmt.Sprintf("Observer %03d", i), + iata, now.Format(time.RFC3339), now.Add(-90*24*time.Hour).Format(time.RFC3339)); err != nil { + t.Fatal(err) + } + } + pubkeys := make([]string, nodes) + for i := range nodes { + pubkeys[i] = fmt.Sprintf("%016x", 0x1000000000000000+i) + if _, err := tx.Exec(`INSERT INTO nodes (public_key, name, role, last_seen, first_seen, advert_count) + VALUES (?, ?, 'repeater', ?, ?, 10)`, pubkeys[i], fmt.Sprintf("node-%04d", i), + now.Add(-time.Duration(i)*time.Minute).Format(time.RFC3339), + now.Add(-90*24*time.Hour).Format(time.RFC3339)); err != nil { + t.Fatal(err) + } + } + if err := tx.Commit(); err != nil { + t.Fatal(err) + } + + // Four fifths of the transmissions are not ADVERTs, which is what makes + // the payload_type index worth anything, and each ADVERT is heard by a + // handful of observers. + rawHex := make([]byte, 240) + for i := range rawHex { + rawHex[i] = "0123456789abcdef"[i%16] + } + start := time.Now() + obs, txCount := 0, 0 + tx = begin() + for obs < obsTarget { + isAdvert := txCount%5 == 0 + payload := 2 + var from any + var decoded string + if isAdvert { + payload = 4 + pk := pubkeys[rng.Intn(nodes)] + from = pk + decoded = `{"type":"ADVERT","pubKey":"` + pk + `","name":"n"}` + } else { + decoded = `{"type":"TXT_MSG","text":"x"}` + } + res, err := tx.Exec(`INSERT INTO transmissions (raw_hex, hash, first_seen, route_type, payload_type, decoded_json, from_pubkey) + VALUES (?, ?, ?, 1, ?, ?, ?)`, string(rawHex), fmt.Sprintf("h%012d", txCount), + now.Add(-time.Duration(txCount)*time.Second).Format(time.RFC3339), payload, decoded, from) + if err != nil { + t.Fatal(err) + } + txID, _ := res.LastInsertId() + for range 1 + rng.Intn(3) { + if _, err := tx.Exec(`INSERT INTO observations (transmission_id, observer_idx, snr, rssi, path_json, timestamp) + VALUES (?, ?, 8.5, -95, '[]', ?)`, txID, 1+rng.Intn(observers), + now.Add(-time.Duration(txCount)*time.Second).Unix()); err != nil { + t.Fatal(err) + } + obs++ + } + txCount++ + if txCount%20000 == 0 { + if err := tx.Commit(); err != nil { + t.Fatal(err) + } + t.Logf("… %d transmissions, %d observations (%v)", txCount, obs, time.Since(start).Round(time.Second)) + tx = begin() + } + } + if err := tx.Commit(); err != nil { + t.Fatal(err) + } + mustExec(`ANALYZE`) + mustExec(`PRAGMA wal_checkpoint(TRUNCATE)`) + st, err := os.Stat(path) + if err != nil { + t.Fatal(err) + } + fmt.Printf("RESULT generated db=%s transmissions=%d observations=%d nodes=%d observers=%d bytes=%d in %v\n", + path, txCount, obs, nodes, observers, st.Size(), time.Since(start).Round(time.Second)) +} + +// TestPerf_NodesRegionPageLoad_176 times one Nodes page load with a region +// selected: four GetNodes pages of 500. It prints the cold load (the first +// one on a freshly opened database, which on this branch includes the one +// full membership scan) and then the median of the following loads, which is +// what a reload or a second tab pays. +func TestPerf_NodesRegionPageLoad_176(t *testing.T) { + runs, _ := strconv.Atoi(os.Getenv("CORESCOPE_PERF_176")) + if runs <= 0 { + t.Skip("set CORESCOPE_PERF_176= and CORESCOPE_PERF_176_DB= to measure") + } + path := os.Getenv("CORESCOPE_PERF_176_DB") + if path == "" { + t.Skip("set CORESCOPE_PERF_176_DB= to measure") + } + db, err := OpenDB(path) + if err != nil { + t.Fatal(err) + } + defer db.Close() + + pageLoad := func() (time.Duration, int) { + t0 := time.Now() + total := 0 + for page := range perf176Pages { + _, n, _, err := db.GetNodes(perf176PageSize, page*perf176PageSize, + "", "", "", "", "", perf176Region) + if err != nil { + t.Fatal(err) + } + total = n + } + return time.Since(t0), total + } + + cold, total := pageLoad() + var warm []time.Duration + for range runs { + d, got := pageLoad() + if got != total { + t.Fatalf("region total changed between page loads: %d then %d", total, got) + } + warm = append(warm, d) + } + fmt.Printf("RESULT nodes-region-pageload region=%s pages=%dx%d matched_nodes=%d cold_ns=%d warm_median_ns=%d runs=%d\n", + perf176Region, perf176Pages, perf176PageSize, total, cold.Nanoseconds(), + perf176Median(warm).Nanoseconds(), runs) +} diff --git a/cmd/server/nodes_region_cache_port_test.go b/cmd/server/nodes_region_cache_port_test.go index eb86d0a9a..13fdb2f4d 100644 --- a/cmd/server/nodes_region_cache_port_test.go +++ b/cmd/server/nodes_region_cache_port_test.go @@ -31,7 +31,7 @@ func TestNodeRegionCachePortMembershipAndReuse(t *testing.T) { defer db.Close() seedTestData(t, db) var scans atomic.Int32 - db.nodeRegionQueryHook = func() { scans.Add(1) } + db.nodeRegionQueryHook = func() error { scans.Add(1); return nil } for _, region := range []string{"SJC", "sjc", "SJC,SJC"} { nodes, total, _, err := db.GetNodes(10, 0, "", "", "", "", "", region) if err != nil { @@ -46,48 +46,59 @@ func TestNodeRegionCachePortMembershipAndReuse(t *testing.T) { } } -func TestNodeRegionCachePortLegacyBackfillAndRefresh(t *testing.T) { +// The delta refresh obeys the same from_pubkey contract as the full scan +// (PR #38): an ADVERT the ingestor left with a NULL from_pubkey does not +// place its node in a region, and a row whose decoded_json is not valid JSON +// does not fail the refresh — which, behind this cache, would leave the entry +// stale for every later request too, not just the one that triggered it. +func TestNodeRegionCacheDeltaFollowsFromPubkeyContract_176(t *testing.T) { db := setupTestDB(t) defer db.Close() seedTestData(t, db) - _, _, _, err := db.GetNodes(10, 0, "", "", "", "", "", "SJC") - if err != nil { - t.Fatal(err) - } - now := time.Now().UTC().Format(time.RFC3339) - _, err = db.conn.Exec(`INSERT INTO nodes(public_key,name,role,last_seen,first_seen) VALUES ('legacykey','Legacy','repeater',?,?)`, now, now) - if err != nil { + // The helper installs a trigger that copies decoded_json.pubKey into + // from_pubkey; drop it so the test writes from_pubkey itself, the way the + // real write paths do. + if _, err := db.conn.Exec(`DROP TRIGGER test_from_pubkey_advert`); err != nil { t.Fatal(err) } - _, err = db.conn.Exec(`INSERT INTO transmissions(raw_hex,hash,first_seen,route_type,payload_type,decoded_json) VALUES ('AA','legacy-hash',?,1,4,'{"pubKey":"legacykey"}')`, now) - if err != nil { + if _, _, _, err := db.GetNodes(10, 0, "", "", "", "", "", "SJC"); err != nil { t.Fatal(err) } - // Simulate an old row not yet processed by the asynchronous backfill. - _, err = db.conn.Exec(`UPDATE transmissions SET from_pubkey=NULL WHERE hash='legacy-hash'`) - if err != nil { - t.Fatal(err) + + now := time.Now().UTC().Format(time.RFC3339) + mustExec := func(q string, args ...any) { + t.Helper() + if _, err := db.conn.Exec(q, args...); err != nil { + t.Fatalf("%s: %v", q, err) + } } - _, err = db.conn.Exec(`INSERT INTO observations(transmission_id,observer_idx,timestamp) VALUES ((SELECT id FROM transmissions WHERE hash='legacy-hash'),1,?)`, time.Now().Unix()) - if err != nil { - t.Fatal(err) + for _, pk := range []string{"nullkey000000001", "stampedkey000002"} { + mustExec(`INSERT INTO nodes(public_key,name,role,last_seen,first_seen) VALUES (?,?,'repeater',?,?)`, + pk, "n-"+pk[:4], now, now) + } + // An ADVERT with a pubkey in decoded_json the backfill has not stamped, + // an ADVERT that is stamped, and an ADVERT whose decoded_json is corrupt. + mustExec(`INSERT INTO transmissions(raw_hex,hash,first_seen,route_type,payload_type,decoded_json,from_pubkey) + VALUES ('AA','d-null',?,1,4,'{"pubKey":"nullkey000000001"}',NULL)`, now) + mustExec(`INSERT INTO transmissions(raw_hex,hash,first_seen,route_type,payload_type,decoded_json,from_pubkey) + VALUES ('BB','d-stamped',?,1,4,'{"pubKey":"stampedkey000002"}','stampedkey000002')`, now) + mustExec(`INSERT INTO transmissions(raw_hex,hash,first_seen,route_type,payload_type,decoded_json,from_pubkey) + VALUES ('CC','d-corrupt',?,1,4,'NOT JSON',NULL)`, now) + for _, hash := range []string{"d-null", "d-stamped", "d-corrupt"} { + mustExec(`INSERT INTO observations(transmission_id,observer_idx,timestamp) + VALUES ((SELECT id FROM transmissions WHERE hash=?),1,?)`, hash, time.Now().Unix()) } - prev := db.getNodeRegionEntry("SJC") - aged := *prev - aged.refreshed = time.Now().Add(-2 * nodeRegionFreshTTL) - db.setNodeRegionEntry("SJC", &aged) + + ageNodeRegionEntry(t, db, "SJC", 2*nodeRegionFreshTTL, 0) if _, err := db.refreshNodeRegion("SJC", []string{"SJC"}); err != nil { - t.Fatal(err) - } - nodes, total, _, err := db.GetNodes(10, 0, "", "", "", "", "", "SJC") - if err != nil { - t.Fatal(err) + t.Fatalf("delta refresh over a corrupt ADVERT: %v", err) } - if total != 2 || len(nodes) != 2 { - t.Fatalf("legacy row missing after delta refresh: total=%d nodes=%v", total, nodes) + keys := db.getNodeRegionEntry("SJC").keysJSON + if !strings.Contains(keys, "stampedkey000002") { + t.Errorf("stamped ADVERT missing from cached set: %s", keys) } - if !strings.Contains(db.getNodeRegionEntry("SJC").keysJSON, "legacykey") { - t.Fatal("legacy key absent from cached set") + if strings.Contains(keys, "nullkey000000001") { + t.Errorf("NULL from_pubkey ADVERT placed its node in the region: %s", keys) } }