diff --git a/cmd/ingestor/advert_evidence_migration_race_test.go b/cmd/ingestor/advert_evidence_migration_race_test.go new file mode 100644 index 000000000..b7965e111 --- /dev/null +++ b/cmd/ingestor/advert_evidence_migration_race_test.go @@ -0,0 +1,217 @@ +package main + +import ( + "context" + "fmt" + "path/filepath" + "strings" + "testing" +) + +func seedUnbackfilledAdvert(t *testing.T, canonical, observed string) (*Store, string) { + t.Helper() + path := filepath.Join(t.TempDir(), "legacy-advert.db") + s, err := OpenStore(path) + if err != nil { + t.Fatal(err) + } + if err := s.UpsertObserver("fixture-observer", "Fixture observer", "", nil); err != nil { + t.Fatal(err) + } + if _, err := s.db.Exec(`INSERT INTO transmissions(id,hash,raw_hex,first_seen,payload_type,route_type) VALUES(1,'legacy-advert',?,'2026-01-01T00:00:00Z',4,?); + INSERT INTO observations(id,transmission_id,observer_idx,raw_hex,path_json,timestamp) VALUES(1,1,(SELECT rowid FROM observers WHERE id='fixture-observer'),?,'[]',1)`, canonical, int(canonical[1]-'0')&3, observed); err != nil { + s.Close() + t.Fatal(err) + } + return s, path +} + +func TestAdvertRouteEvidenceLegacyProtectionErrorsDoNotDropIncomingRaw(t *testing.T) { + for _, failure := range []string{"lookup", "write"} { + t.Run(failure, func(t *testing.T) { + s, _ := seedUnbackfilledAdvert(t, "1100aa", "1200aa") + defer s.Close() + stmt := `DROP TABLE advert_evidence_backfill` + if failure == "write" { + stmt = `CREATE TRIGGER fail_old_evidence BEFORE INSERT ON advert_route_evidence WHEN NEW.bit=2 BEGIN SELECT RAISE(ABORT,'fixture old evidence failure'); END` + } + if _, err := s.db.Exec(stmt); err != nil { + t.Fatal(err) + } + _, err := s.InsertTransmission(&PacketData{Hash: "legacy-advert", ObserverID: "fixture-observer", PayloadType: 4, RouteType: 1, RawHex: "1100aa", PathJSON: "[]"}) + if err != nil { + t.Errorf("analytics preservation failure dropped incoming frame: %v", err) + } + var old string + if err := s.db.QueryRow(`SELECT raw_hex FROM observations WHERE id=1`).Scan(&old); err != nil { + t.Fatal(err) + } + if old != "1100aa" { + t.Fatalf("failed analytics preservation prevented incoming raw: %q", old) + } + }) + } +} + +func TestAdvertRouteEvidenceLegacyCoalescedPathAndIndex(t *testing.T) { + s, _ := seedUnbackfilledAdvert(t, "1100aa", "1200aa") + defer s.Close() + if _, err := s.db.Exec(`UPDATE observations SET path_json=NULL WHERE id=1`); err != nil { + t.Fatal(err) + } + if _, err := s.InsertTransmission(&PacketData{Hash: "legacy-advert", ObserverID: "fixture-observer", PayloadType: 4, RouteType: 1, RawHex: "1100aa", PathJSON: ""}); err != nil { + t.Fatal(err) + } + var count, mask int + if err := s.db.QueryRow(`SELECT COUNT(*) FROM observations`).Scan(&count); err != nil { + t.Fatal(err) + } + if count != 1 { + t.Fatalf("NULL and empty path must conflict, got %d observations", count) + } + if err := s.db.QueryRow(`SELECT SUM(bit) FROM advert_route_evidence`).Scan(&mask); err != nil { + t.Fatal(err) + } + if mask != 3 { + t.Fatalf("coalesced path lookup lost old evidence: mask=%d", mask) + } + rows, err := s.db.Query(`EXPLAIN QUERY PLAN `+legacyAdvertObservationSQL, 1, 1, "") + if err != nil { + t.Fatal(err) + } + defer rows.Close() + indexed := false + for rows.Next() { + var id, parent, unused int + var detail string + if err := rows.Scan(&id, &parent, &unused, &detail); err != nil { + t.Fatal(err) + } + if strings.Contains(detail, "SEARCH observations USING INDEX idx_observations_dedup") { + indexed = true + } + } + if err := rows.Err(); err != nil { + t.Fatal(err) + } + if !indexed { + t.Fatal("legacy preservation must use the unique observation conflict index") + } +} + +func TestAdvertRouteEvidenceCompletionSkipsLookupAfterRestart(t *testing.T) { + s, path := seedUnbackfilledAdvert(t, "1100aa", "1200aa") + defer func() { s.Close() }() + if err := s.RunAsyncMigration(context.Background(), "advert_route_evidence_v1", s.backfillAdvertEvidence); err != nil { + t.Fatal(err) + } + s.WaitForAsyncMigrations() + if !s.advertEvidenceComplete.Load() { + t.Fatal("backfill did not close preservation window") + } + s.Close() + var err error + s, err = OpenStore(path) + if err != nil { + t.Fatal(err) + } + if !s.advertEvidenceComplete.Load() { + t.Fatal("restart did not restore persisted migration completion") + } + // Any accidental legacy lookup now fails. Completed stores must incur + // no extra SQL on this steady-state overwrite path. + if _, err := s.db.Exec(`DROP TABLE advert_evidence_backfill`); err != nil { + t.Fatal(err) + } + if _, err := s.InsertTransmission(&PacketData{Hash: "legacy-advert", ObserverID: "fixture-observer", PayloadType: 4, RouteType: 1, RawHex: "1100aa", PathJSON: "[]"}); err != nil { + t.Fatal(err) + } +} + +func BenchmarkAdvertEvidenceLegacyPreservation(b *testing.B) { + s, err := OpenStore(filepath.Join(b.TempDir(), "legacy-preservation.db")) + if err != nil { + b.Fatal(err) + } + defer s.Close() + s.WaitForAsyncMigrations() + if err := s.UpsertObserver("fixture-observer", "Fixture observer", "", nil); err != nil { + b.Fatal(err) + } + if _, err := s.InsertTransmission(&PacketData{Hash: "bench-legacy", ObserverID: "fixture-observer", PayloadType: 4, RouteType: 1, RawHex: "1100aa", PathJSON: "[]"}); err != nil { + b.Fatal(err) + } + var observerIdx int64 + if err := s.db.QueryRow(`SELECT rowid FROM observers WHERE id='fixture-observer'`).Scan(&observerIdx); err != nil { + b.Fatal(err) + } + for _, state := range []string{"pending", "checkpointed", "complete"} { + b.Run(state, func(b *testing.B) { + cursor := 0 + if state != "pending" { + cursor = 1 + } + if _, err := s.db.Exec(`INSERT INTO advert_evidence_backfill(id,obs_cursor) VALUES(1,?) ON CONFLICT(id) DO UPDATE SET obs_cursor=excluded.obs_cursor`, cursor); err != nil { + b.Fatal(err) + } + s.advertEvidenceComplete.Store(state == "complete") + writerMu.Lock() + defer writerMu.Unlock() + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + if err := s.preserveLegacyAdvertObservation(1, observerIdx, "[]"); err != nil { + b.Fatal(err) + } + } + }) + } +} + +func TestAdvertRouteEvidencePreservesLegacyConflictBeforeBackfill(t *testing.T) { + for _, tc := range []struct{ name, canonical, oldRaw, incoming string }{ + {"flood_zero_then_flood", "1100aa", "1200aa", "1100aa"}, + {"zero_flood_then_zero", "1200aa", "1100aa", "1200aa"}, + {"flood_zero_then_malformed", "1100aa", "1200aa", "11zzaa"}, + {"zero_flood_then_malformed", "1200aa", "1100aa", "12zzaa"}, + } { + for _, restart := range []bool{false, true} { + t.Run(fmt.Sprintf("%s/restart=%v", tc.name, restart), func(t *testing.T) { + s, path := seedUnbackfilledAdvert(t, tc.canonical, tc.oldRaw) + defer func() { s.Close() }() + // Do not replay oldRaw: it only exists in the legacy row and + // this one incoming frame will replace it before backfill. + data := &PacketData{Hash: "legacy-advert", ObserverID: "fixture-observer", PayloadType: 4, RouteType: int(tc.canonical[1]-'0') & 3, RawHex: tc.incoming, PathJSON: "[]"} + if _, err := s.InsertTransmission(data); err != nil { + t.Fatal(err) + } + var count, id int + var surviving string + if err := s.db.QueryRow(`SELECT COUNT(*),MIN(id),MIN(raw_hex) FROM observations`).Scan(&count, &id, &surviving); err != nil { + t.Fatal(err) + } + if count != 1 || id != 1 || surviving != tc.incoming { + t.Fatalf("fixture did not replace same observation: count=%d id=%d raw=%s", count, id, surviving) + } + if restart { + s.Close() + var err error + s, err = OpenStore(path) + if err != nil { + t.Fatal(err) + } + } + if err := s.backfillAdvertEvidence(context.Background(), s.db); err != nil { + t.Fatal(err) + } + var mask int + if err := s.db.QueryRow(`SELECT COUNT(*),COALESCE(SUM(bit),0) FROM advert_route_evidence WHERE tx_id=1`).Scan(&count, &mask); err != nil { + t.Fatal(err) + } + if count != 2 || mask != 3 { + t.Fatalf("upgrade overwrote available legacy route evidence: rows=%d mask=%d, want 2/3", count, mask) + } + }) + } + } +} diff --git a/cmd/ingestor/advert_evidence_scale_test.go b/cmd/ingestor/advert_evidence_scale_test.go new file mode 100644 index 000000000..bd716a373 --- /dev/null +++ b/cmd/ingestor/advert_evidence_scale_test.go @@ -0,0 +1,120 @@ +//go:build advert_evidence_scale + +package main + +import ( + "context" + "os" + "path/filepath" + "sort" + "strings" + "testing" + "time" +) + +// Opt-in disk benchmark: a new database with 1M transmissions and 11M +// observations (10% adverts, 11 observations per transmission, 120-byte frames). +// ADVERT_SCALE_DB retains the fixture for the server catch-up benchmark. +func TestAdvertBackfillScale(t *testing.T) { + path := os.Getenv("ADVERT_SCALE_DB") + if path == "" { + path = filepath.Join(t.TempDir(), "advert-scale.db") + } + if _, err := os.Stat(path); err == nil { + t.Fatal("ADVERT_SCALE_DB must be a new temporary database") + } else if !os.IsNotExist(err) { + t.Fatal(err) + } + s, err := OpenStore(path) + if err != nil { + t.Fatal(err) + } + defer s.Close() + s.WaitForAsyncMigrations() + var existing int + if err := s.db.QueryRow(`SELECT COUNT(*) FROM observations`).Scan(&existing); err != nil { + t.Fatal(err) + } + if existing != 0 { + t.Fatal("scale benchmark requires a new empty database") + } + started := time.Now() + { + raw := "1100" + strings.Repeat("ab", 118) + for start := 1; start <= 1000000; start += 10000 { + _, err = s.db.Exec(`WITH RECURSIVE ids(n) AS (VALUES(?) UNION ALL SELECT n+1 FROM ids WHERE n COALESCE((SELECT obs_cursor FROM advert_evidence_backfill WHERE id=1),0)` + +// Caller holds writerMu, including through the subsequent observation UPSERT. +// The backfill commits evidence and obs_cursor together under that same lock. +func (s *Store) preserveLegacyAdvertObservation(txID, observerIdx int64, path string) error { + if s.advertEvidenceComplete.Load() { + return nil + } + var raw sql.NullString + err := s.stmtGetLegacyAdvertObservation.QueryRow(txID, observerIdx, path).Scan(&raw) + if err == sql.ErrNoRows { + return nil + } + if err != nil { + return err + } + if bit := packetpath.AdvertRouteEvidence(raw.String); bit != 0 { + _, err = s.stmtInsertAdvertEvidence.Exec(txID, bit, txID, bit) + } + return err +} + +// backfillAdvertEvidence recovers only evidence still present in canonical and +// observation frames. An overwritten middle frame is unrecoverable. Persisted +// cursors and 500-row transactions make this cancellable/resumable; new traffic +// records evidence synchronously, so rows beyond either scan need no replay. +func (s *Store) backfillAdvertEvidence(ctx context.Context, db *sql.DB) error { + if _, err := db.ExecContext(ctx, `INSERT OR IGNORE INTO advert_evidence_backfill(id) VALUES(1)`); err != nil { + return err + } + for _, source := range []struct{ cursor, table, query string }{ + {"tx_cursor", "transmissions", `SELECT id,id,COALESCE(raw_hex,''),payload_type FROM transmissions WHERE id>? AND id<=? ORDER BY id LIMIT 500`}, + {"obs_cursor", "observations", `SELECT o.id,o.transmission_id,COALESCE(o.raw_hex,''),t.payload_type FROM observations o JOIN transmissions t ON t.id=o.transmission_id WHERE o.id>? AND o.id<=? ORDER BY o.id LIMIT 500`}, + } { + // A finite horizon prevents live traffic from extending this scan. + var upper int64 + if err := db.QueryRowContext(ctx, `SELECT COALESCE(MAX(id),0) FROM `+source.table).Scan(&upper); err != nil { + return err + } + for { + if err := ctx.Err(); err != nil { + return err + } + var cursor int64 + if err := db.QueryRowContext(ctx, `SELECT `+source.cursor+` FROM advert_evidence_backfill WHERE id=1`).Scan(&cursor); err != nil { + return err + } + rows, err := db.QueryContext(ctx, source.query, cursor, upper) + if err != nil { + return err + } + type evidence struct { + txID int64 + bit uint8 + } + batch := make([]evidence, 0, 500) + lastID := cursor + for rows.Next() { + var id, txID int64 + var raw string + var payloadType sql.NullInt64 + if err := rows.Scan(&id, &txID, &raw, &payloadType); err != nil { + rows.Close() + return err + } + lastID = id + if bit := packetpath.AdvertRouteEvidence(raw); bit != 0 && payloadType.Valid && payloadType.Int64 == 4 { + batch = append(batch, evidence{txID, bit}) + } + } + err = rows.Err() + rows.Close() + if err != nil { + return err + } + if lastID == cursor { + break + } + // The read cursor is closed before the single writer is acquired. + // Retention may have removed a transmission meanwhile; do not + // resurrect evidence for it or violate its foreign key. + err = func() error { + writerMu.Lock() + defer writerMu.Unlock() + tx, err := db.BeginTx(ctx, nil) + if err != nil { + return err + } + defer tx.Rollback() + for _, item := range batch { + if _, err := tx.ExecContext(ctx, insertAdvertEvidenceSQL+` AND EXISTS(SELECT 1 FROM transmissions WHERE id=?)`, item.txID, item.bit, item.txID, item.bit, item.txID); err != nil { + return err + } + } + if _, err := tx.ExecContext(ctx, `UPDATE advert_evidence_backfill SET `+source.cursor+`=? WHERE id=1`, lastID); err != nil { + return err + } + return tx.Commit() + }() + if err != nil { + return fmt.Errorf("advert evidence backfill: %w", err) + } + // Yield between bounded batches so live ingestion gets the writer. + timer := time.NewTimer(10 * time.Millisecond) + select { + case <-ctx.Done(): + timer.Stop() + return ctx.Err() + case <-timer.C: + } + } + } + // RunAsyncMigration persists "done" after this returns. A crash before + // that status write simply re-enables protection and resumes the cursors. + s.advertEvidenceComplete.Store(true) + return nil +} diff --git a/cmd/ingestor/advert_route_evidence_test.go b/cmd/ingestor/advert_route_evidence_test.go new file mode 100644 index 000000000..bccbf8a11 --- /dev/null +++ b/cmd/ingestor/advert_route_evidence_test.go @@ -0,0 +1,362 @@ +package main + +import ( + "context" + "errors" + "fmt" + "path/filepath" + "testing" +) + +// These frames have the same advert payload. Firmware hashes payload/type, +// excluding the route header and path, so they belong to one transmission. +func TestAdvertRouteEvidenceSurvivesObservationUpsert(t *testing.T) { + for _, sequence := range []struct { + name string + raws []string + }{ + {"flood_zero", []string{"1100aa", "1200aa"}}, + {"zero_flood", []string{"1200aa", "1100aa"}}, + {"flood_zero_flood", []string{"1100aa", "1200aa", "1100aa"}}, + {"zero_flood_zero", []string{"1200aa", "1100aa", "1200aa"}}, + {"flood_empty_2byte_flood", []string{"1100aa", "1240aa", "1100aa"}}, + {"empty_3byte_flood_empty", []string{"1280aa", "1100aa", "1280aa"}}, + {"transport_empty_2byte_flood", []string{"130102030440aa", "100102030400aa"}}, + {"transport_flood_empty_3byte", []string{"100102030400aa", "130102030480aa"}}, + {"transport_flood_zero_flood", []string{"100102030400aa", "130102030400aa", "100102030400aa"}}, + {"transport_zero_flood_zero", []string{"130102030400aa", "100102030400aa", "130102030400aa"}}, + } { + t.Run(sequence.name, func(t *testing.T) { + path := filepath.Join(t.TempDir(), "evidence.db") + s, err := OpenStore(path) + if err != nil { + t.Fatal(err) + } + defer func() { s.Close() }() + if err := s.UpsertObserver("test-observer", "Fixture observer", "", nil); err != nil { + t.Fatal(err) + } + data := &PacketData{Hash: "advert-evidence", PayloadType: 4, ObserverID: "test-observer", PathJSON: "[]", DecodedJSON: `{"type":"ADVERT"}`} + for i, raw := range sequence.raws { + data.RawHex = raw + data.RouteType = int(raw[1]-'0') & 3 + // Same second, then an older receive-time: timestamps cannot be + // used as a reliable cursor for evidence changes. + data.Timestamp = "2026-01-02T00:00:00Z" + if i == 2 { + data.Timestamp = "2026-01-01T00:00:00Z" + } + if _, err := s.InsertTransmission(data); err != nil { + t.Fatal(err) + } + } + var count, obsID int + var canonical, surviving string + if err := s.db.QueryRow(`SELECT COUNT(*), MIN(id), MIN(raw_hex) FROM observations`).Scan(&count, &obsID, &surviving); err != nil { + t.Fatal(err) + } + if count != 1 || obsID != 1 || surviving != sequence.raws[len(sequence.raws)-1] { + t.Fatalf("fixture must overwrite one observation in place: count=%d id=%d raw=%s", count, obsID, surviving) + } + if err := s.db.QueryRow(`SELECT raw_hex FROM transmissions`).Scan(&canonical); err != nil { + t.Fatal(err) + } + if canonical != sequence.raws[0] { + t.Fatalf("canonical raw changed: %s", canonical) + } + if err := s.db.QueryRow(`SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='advert_route_evidence'`).Scan(&count); err != nil { + t.Fatal(err) + } + if count != 1 { + t.Fatal("missing durable advert route evidence: expected both route families to survive same-row observation replacement") + } + assertMixed := func() int64 { + t.Helper() + var rows, mask int + var maxID int64 + if err := s.db.QueryRow(`SELECT COUNT(*), COALESCE(SUM(bit),0), COALESCE(MAX(id),0) FROM advert_route_evidence`).Scan(&rows, &mask, &maxID); err != nil { + t.Fatal(err) + } + if rows != 2 || mask != 3 { + t.Fatalf("durable evidence rows=%d mask=%d, want exactly two rows and mixed mask=3", rows, mask) + } + return maxID + } + maxID := assertMixed() + for i := 0; i < 100; i++ { + data.RawHex = sequence.raws[i%len(sequence.raws)] + data.RouteType = int(data.RawHex[1]-'0') & 3 + if _, err := s.InsertTransmission(data); err != nil { + t.Fatal(err) + } + } + if assertMixed() != maxID { + t.Fatal("duplicate traffic appended route evidence") + } + var sequenceID int64 + if err := s.db.QueryRow(`SELECT seq FROM sqlite_sequence WHERE name='advert_route_evidence'`).Scan(&sequenceID); err != nil { + t.Fatal(err) + } + if sequenceID != maxID { + t.Fatalf("duplicate traffic advanced evidence sequence: %d, want %d", sequenceID, maxID) + } + s.Close() + s, err = OpenStore(path) + if err != nil { + t.Fatal(err) + } + assertMixed() + // The evidence lifetime is bounded by the retained transmission. + if _, err := s.db.Exec(`DELETE FROM observations; DELETE FROM transmissions`); err != nil { + t.Fatal(err) + } + if err := s.db.QueryRow(`SELECT COUNT(*) FROM advert_route_evidence`).Scan(&count); err != nil { + t.Fatal(err) + } + if count != 0 { + t.Fatalf("retention left %d orphan evidence rows", count) + } + }) + } +} + +func TestAdvertRouteEvidenceBackfillResumeAndLiveUnion(t *testing.T) { + s, err := OpenStore(filepath.Join(t.TempDir(), "backfill.db")) + if err != nil { + t.Fatal(err) + } + defer s.Close() + s.WaitForAsyncMigrations() + if err := s.UpsertObserver("fixture-observer", "Fixture observer", "", nil); err != nil { + t.Fatal(err) + } + // Simulate pre-upgrade history without going through the new writer. + tx, err := s.db.Begin() + if err != nil { + t.Fatal(err) + } + for id := 1; id <= 1200; id++ { + if _, err := tx.Exec(`INSERT INTO transmissions(id,hash,raw_hex,first_seen,payload_type,route_type) VALUES(?,?, '1100aa','2026-01-01T00:00:00Z',4,1)`, id, fmt.Sprintf("history-%d", id)); err != nil { + t.Fatal(err) + } + if _, err := tx.Exec(`INSERT INTO observations(transmission_id,observer_idx,raw_hex,path_json,timestamp) VALUES(?,(SELECT rowid FROM observers WHERE id='fixture-observer'),'1200aa','[]',1)`, id); err != nil { + t.Fatal(err) + } + } + if err := tx.Commit(); err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := s.backfillAdvertEvidence(ctx, s.db); !errors.Is(err, context.Canceled) { + t.Fatalf("cancelled backfill returned %v", err) + } + // Abort after one committed batch; the failed batch must not move its + // persisted cursor, so a retry can recover every remaining frame. + if _, err := s.db.Exec(`CREATE TRIGGER fail_evidence_batch BEFORE INSERT ON advert_route_evidence WHEN NEW.tx_id=501 BEGIN SELECT RAISE(ABORT,'fixture failure'); END`); err != nil { + t.Fatal(err) + } + if err := s.backfillAdvertEvidence(context.Background(), s.db); err == nil { + t.Fatal("expected injected batch failure") + } + var cursor, count int + if err := s.db.QueryRow(`SELECT tx_cursor FROM advert_evidence_backfill WHERE id=1`).Scan(&cursor); err != nil { + t.Fatal(err) + } + if cursor != 500 { + t.Fatalf("cursor=%d after failure, want committed batch boundary 500", cursor) + } + if _, err := s.db.Exec(`DROP TRIGGER fail_evidence_batch`); err != nil { + t.Fatal(err) + } + done := make(chan error, 1) + go func() { done <- s.backfillAdvertEvidence(context.Background(), s.db) }() + // Live processing can overwrite the only surviving zero-hop raw before + // the backfill reaches it; its synchronous evidence must preserve it. + data := &PacketData{Hash: "history-1200", ObserverID: "fixture-observer", PayloadType: 4, RouteType: 1, RawHex: "1100aa", Timestamp: "2026-01-01T00:00:00Z", PathJSON: "[]"} + if _, err := s.InsertTransmission(data); err != nil { + t.Fatal(err) + } + if err := <-done; err != nil { + t.Fatal(err) + } + if err := s.db.QueryRow(`SELECT COUNT(*) FROM advert_route_evidence`).Scan(&count); err != nil { + t.Fatal(err) + } + if count != 2400 { + t.Fatalf("got %d evidence rows, want two per history transmission", count) + } + var before, after int64 + if err := s.db.QueryRow(`SELECT seq FROM sqlite_sequence WHERE name='advert_route_evidence'`).Scan(&before); err != nil { + t.Fatal(err) + } + if err := s.backfillAdvertEvidence(context.Background(), s.db); err != nil { + t.Fatal(err) + } + if err := s.db.QueryRow(`SELECT seq FROM sqlite_sequence WHERE name='advert_route_evidence'`).Scan(&after); err != nil { + t.Fatal(err) + } + if before != after { + t.Fatalf("idempotent backfill changed evidence sequence %d -> %d", before, after) + } +} + +// Analytics is best effort: failure must not prevent the core observation, +// resolved relay path, relay liveness, or transmission timestamp from landing. +func TestAdvertRouteEvidenceFailureKeepsCoreIngestion(t *testing.T) { + for _, failure := range []string{"incoming-write", "legacy-read", "legacy-write"} { + t.Run(failure, func(t *testing.T) { + s, err := OpenStore(filepath.Join(t.TempDir(), "write-failure.db")) + if err != nil { + t.Fatal(err) + } + defer s.Close() + const relay = "bbbbbbbbbb" + seedRelayNode(t, s, relay, "Fixture relay", "2026-01-01T00:00:00Z") + if err := s.RefreshPrefixIndex(); err != nil { + t.Fatal(err) + } + if err := s.UpsertObserver("fixture-observer", "Fixture observer", "", nil); err != nil { + t.Fatal(err) + } + data := &PacketData{Hash: "write-failure", PayloadType: 4, RouteType: 1, RawHex: "1101bbaa", PathJSON: `["bb"]`, ObserverID: "fixture-observer", Timestamp: "2026-01-01T00:00:00Z"} + if _, err := s.InsertTransmission(data); err != nil { + t.Fatal(err) + } + if _, err := s.db.Exec(`DELETE FROM advert_route_evidence; UPDATE observations SET raw_hex='1200aa',resolved_path=NULL`); err != nil { + t.Fatal(err) + } + stmt := `CREATE TRIGGER fail_evidence BEFORE INSERT ON advert_route_evidence WHEN NEW.bit=1 BEGIN SELECT RAISE(ABORT,'fixture evidence failure'); END` + if failure == "legacy-read" { + stmt = `DROP TABLE advert_evidence_backfill` + } + if failure == "legacy-write" { + stmt = `CREATE TRIGGER fail_evidence BEFORE INSERT ON advert_route_evidence WHEN NEW.bit=2 BEGIN SELECT RAISE(ABORT,'fixture evidence failure'); END` + } + if _, err := s.db.Exec(stmt); err != nil { + t.Fatal(err) + } + before := s.Stats.WriteErrors.Load() + data.Timestamp = "2026-01-02T00:00:00Z" + if _, err := s.InsertTransmission(data); err != nil { + t.Errorf("analytics failure aborted core ingestion: %v", err) + } + // The observation UPSERT preserves its original timestamp; tx last_seen advances. + var raw, resolved string + var ts, lastSeen int64 + if err := s.db.QueryRow(`SELECT raw_hex,COALESCE(resolved_path,''),timestamp FROM observations`).Scan(&raw, &resolved, &ts); err != nil { + t.Fatal(err) + } + if raw != data.RawHex || resolved != `["bbbbbbbbbb"]` || ts != 1767225600 { + t.Errorf("core observation not updated: raw=%s resolved=%s timestamp=%d", raw, resolved, ts) + } + if err := s.db.QueryRow(`SELECT last_seen FROM transmissions`).Scan(&lastSeen); err != nil { + t.Fatal(err) + } + if lastSeen != 1767312000 || nodeLastSeen(t, s, relay) != data.Timestamp { + t.Errorf("core liveness not updated: tx=%d relay=%s", lastSeen, nodeLastSeen(t, s, relay)) + } + if s.Stats.WriteErrors.Load() != before+1 { + t.Errorf("analytics failure not counted once: before=%d after=%d", before, s.Stats.WriteErrors.Load()) + } + var mask int + if err := s.db.QueryRow(`SELECT COALESCE(SUM(bit),0) FROM advert_route_evidence`).Scan(&mask); err != nil { + t.Fatal(err) + } + want := 1 + if failure == "incoming-write" { + want = 2 + } + if mask != want { + t.Errorf("independent evidence not preserved: mask=%d want %d", mask, want) + } + }) + } +} + +func TestAdvertRouteEvidenceConstraintsAndFeedIndex(t *testing.T) { + s, err := OpenStore(filepath.Join(t.TempDir(), "constraints.db")) + if err != nil { + t.Fatal(err) + } + defer s.Close() + if _, err := s.InsertTransmission(&PacketData{Hash: "constraints", PayloadType: 4, RouteType: 1, RawHex: "1100aa", PathJSON: "[]"}); err != nil { + t.Fatal(err) + } + for _, stmt := range []string{ + `INSERT INTO advert_route_evidence(tx_id,bit) VALUES(1,1)`, + `INSERT INTO advert_route_evidence(tx_id,bit) VALUES(1,3)`, + `INSERT INTO advert_route_evidence(tx_id,bit) VALUES(999,2)`, + } { + if _, err := s.db.Exec(stmt); err == nil { + t.Errorf("constraint accepted: %s", stmt) + } + } + var id, parent, unused int + var detail string + if err := s.db.QueryRow(`EXPLAIN QUERY PLAN SELECT id,tx_id,bit FROM advert_route_evidence WHERE id>1 ORDER BY id LIMIT 500`).Scan(&id, &parent, &unused, &detail); err != nil { + t.Fatal(err) + } + if detail != "SEARCH advert_route_evidence USING INTEGER PRIMARY KEY (rowid>?)" { + t.Fatalf("feed must seek by primary key: %s", detail) + } + if _, err := s.db.Exec(`DELETE FROM observations; DELETE FROM transmissions`); err != nil { + t.Fatal(err) + } + if _, err := s.InsertTransmission(&PacketData{Hash: "after-retention", PayloadType: 4, RouteType: 1, RawHex: "1100aa", PathJSON: "[]"}); err != nil { + t.Fatal(err) + } + var nextID int + if err := s.db.QueryRow(`SELECT MIN(id) FROM advert_route_evidence`).Scan(&nextID); err != nil { + t.Fatal(err) + } + if nextID <= 1 { + t.Fatalf("feed ID reused after retention: %d", nextID) + } +} + +func BenchmarkAdvertEvidenceRepeatedWrite(b *testing.B) { + s, err := OpenStore(filepath.Join(b.TempDir(), "repeat.db")) + if err != nil { + b.Fatal(err) + } + defer s.Close() + s.WaitForAsyncMigrations() + if _, err := s.InsertTransmission(&PacketData{Hash: "benchmark-repeat", PayloadType: 4, RouteType: 1, RawHex: "1100aa", PathJSON: "[]"}); err != nil { + b.Fatal(err) + } + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + if _, err := s.stmtInsertAdvertEvidence.Exec(1, 1, 1, 1); err != nil { + b.Fatal(err) + } + } +} + +func TestAdvertRouteEvidenceRejectsUnprovenFrames(t *testing.T) { + s, err := OpenStore(filepath.Join(t.TempDir(), "evidence.db")) + if err != nil { + t.Fatal(err) + } + defer s.Close() + for _, raw := range []string{"", "12", "1200", "1200a", "12zzaa", "12c0aa", "1201ffaa", "130102030401ffaa", "13zz00000000aa", "1600aa", "1100zz", "1101"} { + data := &PacketData{Hash: "unproven-" + raw, PayloadType: 4, RouteType: 2, RawHex: raw, Timestamp: "2026-01-01T00:00:00Z", PathJSON: "[]"} + if _, err := s.InsertTransmission(data); err != nil { + t.Fatal(err) + } + } + var count int + if err := s.db.QueryRow(`SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='advert_route_evidence'`).Scan(&count); err != nil { + t.Fatal(err) + } + if count != 1 { + t.Fatal("missing durable advert route evidence table") + } + if err := s.db.QueryRow(`SELECT COUNT(*) FROM advert_route_evidence`).Scan(&count); err != nil { + t.Fatal(err) + } + if count != 0 { + t.Fatalf("malformed/nonzero-path/unrelated frames contributed %d evidence rows", count) + } +} diff --git a/cmd/ingestor/db.go b/cmd/ingestor/db.go index ff80b2827..cf5024bc7 100644 --- a/cmd/ingestor/db.go +++ b/cmd/ingestor/db.go @@ -76,6 +76,7 @@ type Store struct { stmtUpdateTxFirstSeen *sql.Stmt stmtBumpTxLastSeen *sql.Stmt stmtInsertObservation *sql.Stmt + stmtInsertAdvertEvidence *sql.Stmt stmtUpsertNode *sql.Stmt stmtIncrementAdvertCount *sql.Stmt stmtUpsertObserver *sql.Stmt @@ -85,9 +86,13 @@ type Store struct { stmtTouchNodeLastSeen *sql.Stmt stmtUpsertMetrics *sql.Stmt + stmtGetLegacyAdvertObservation *sql.Stmt + sampleIntervalSec int backfillWg sync.WaitGroup + advertEvidenceComplete atomic.Bool // restored from durable migration status + // prefixIdx holds the prefix → pubkey index used by the // resolved_path writer (#1547). Rebuilt on startup and once per // neighbor-edges builder tick (60s). @@ -213,6 +218,11 @@ func OpenStoreWithInterval(dbPath string, sampleIntervalSec int) (*Store, error) log.Printf("[migration/async] scheduling tx_last_seen_backfill_v1 failed: %v", err) } + // A missing/failed completion lookup leaves preservation enabled. This + // lifecycle state is restored on restart, independently of main's startup. + var evidenceStatus string + _ = db.QueryRow(`SELECT status FROM _async_migrations WHERE name='advert_route_evidence_v1'`).Scan(&evidenceStatus) + s.advertEvidenceComplete.Store(evidenceStatus == "done") return s, nil } @@ -906,6 +916,14 @@ func applySchema(db *sql.DB) error { func (s *Store) prepareStatements() error { var err error + s.stmtInsertAdvertEvidence, err = s.db.Prepare(insertAdvertEvidenceSQL) + if err != nil { + return err + } + s.stmtGetLegacyAdvertObservation, err = s.db.Prepare(legacyAdvertObservationSQL) + if err != nil { + return err + } s.stmtGetTxByHash, err = s.db.Prepare("SELECT id, first_seen FROM transmissions WHERE hash = ?") if err != nil { @@ -1106,6 +1124,18 @@ func (s *Store) InsertTransmission(data *PacketData) (bool, error) { if !isNew { s.Stats.DuplicateTransmissions.Add(1) } + // Capture route evidence BEFORE the observation conflict update can erase + // a different route. Duplicate evidence is a read-only indexed probe. + // Analytics failures must not drop core observations or liveness updates; + // known evidence is a lower bound when an evidence read/write fails. + if data.PayloadType == 4 { + if bit := packetpath.AdvertRouteEvidence(data.RawHex); bit != 0 { + if _, err := s.stmtInsertAdvertEvidence.Exec(txID, bit, txID, bit); err != nil { + s.Stats.WriteErrors.Add(1) + log.Printf("[db] record advert route evidence (non-fatal): %v", err) + } + } + } // Resolve observer_idx and update last_seen var observerIdx *int64 @@ -1122,6 +1152,17 @@ func (s *Store) InsertTransmission(data *PacketData) (bool, error) { } } + // Until backfill commits this observation's evidence, preserve its old + // frame before UPSERT can destroy it. writerMu also guards checkpoints. + // Run even for malformed incoming raw: the surviving old frame is valid + // evidence independently of whether the new frame contributes a bit. + if !isNew && data.PayloadType == 4 && observerIdx != nil && data.RawHex != "" { + if err := s.preserveLegacyAdvertObservation(txID, *observerIdx, data.PathJSON); err != nil { + s.Stats.WriteErrors.Add(1) + log.Printf("[db] preserve legacy advert evidence (non-fatal): %v", err) + } + } + // Insert observation epochTs := time.Now().Unix() if t, err := time.Parse(time.RFC3339, rxTime); err == nil { diff --git a/cmd/ingestor/main.go b/cmd/ingestor/main.go index ce1ad6a08..a107e2956 100644 --- a/cmd/ingestor/main.go +++ b/cmd/ingestor/main.go @@ -1,6 +1,7 @@ package main import ( + "context" "crypto/hmac" "crypto/rand" "crypto/sha256" @@ -342,6 +343,13 @@ func main() { // explicit. Now drain everything the subscription buffered during startup. store.WaitForAsyncMigrations() ingestBuffer.Ready() + // History recovery must not join the startup readiness wait above. + // Its known evidence is a lower bound: old UPSERTs erased some frames. + evidenceCtx, stopEvidenceBackfill := context.WithCancel(context.Background()) + defer stopEvidenceBackfill() + if err := store.RunAsyncMigration(evidenceCtx, "advert_route_evidence_v1", store.backfillAdvertEvidence); err != nil { + log.Printf("[migration] scheduling advert evidence backfill: %v", err) + } if d := ingestBuffer.Dropped(); d > 0 { log.Printf("[ingest-buffer] write path ready; draining backlog (dropped %d during startup — consider raising ingestBufferSize)", d) } else { @@ -616,6 +624,7 @@ func main() { <-sig log.Println("Shutting down...") + stopEvidenceBackfill() retentionTicker.Stop() metricsRetentionTicker.Stop() if packetRetentionTicker != nil { diff --git a/cmd/server/advert_evidence_scale_test.go b/cmd/server/advert_evidence_scale_test.go new file mode 100644 index 000000000..b569d30b6 --- /dev/null +++ b/cmd/server/advert_evidence_scale_test.go @@ -0,0 +1,105 @@ +//go:build advert_evidence_scale + +package main + +import ( + "fmt" + "os" + "testing" + "time" +) + +// Run after the ingestor's TestAdvertBackfillScale with the same ADVERT_SCALE_DB. +// Models a server starting just before backfill: its cursor is zero and retained +// packets initially have no evidence. Uses the default 1-second poll cadence, +// with one real cached analytics request per tick and a 300K-packet retention. +func TestAdvertEvidenceCatchupScale(t *testing.T) { + path := os.Getenv("ADVERT_SCALE_DB") + if path == "" { + t.Fatal("ADVERT_SCALE_DB must name the completed ingestor scale fixture") + } + db, err := OpenDB(path) + if err != nil { + t.Fatal(err) + } + defer db.conn.Close() + var horizon int64 + var evidenceCount int + if err := db.conn.QueryRow(`SELECT MAX(id),COUNT(*) FROM advert_route_evidence`).Scan(&horizon, &evidenceCount); err != nil { + t.Fatal(err) + } + if evidenceCount < 200000 { + t.Fatal("fixture must contain the completed 11M-observation evidence set") + } + const size = 300000 + packets := make([]*StoreTx, size) + for i := range packets { + pt := PayloadGRP_TXT + if (i+1)%10 == 0 { + pt = PayloadADVERT + } + tx := makeRelayAirtimeTx(i+1, pt, 120, 0, fmt.Sprintf("scale-%d", i+1)) + tx.FirstSeen = time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC).Add(time.Duration(i) * (7 * 24 * time.Hour / size)).Format(time.RFC3339) + packets[i] = tx + } + s := newRelayAirtimeShareTestStore(packets) + s.db = db + s.rfCacheTTL = time.Hour + window := TimeWindow{Since: "2026-01-07T23:00:00Z", Until: "2026-01-08T00:00:00Z", Label: "1h"} + initial := s.GetRelayAirtimeShareWithWindow(window) + if initial["total_count"] != 1785 { + t.Fatalf("window fixture selected %v, want1785", initial["total_count"]) + } + var initialAdverts int + for _, row := range initial["rows"].([]map[string]interface{}) { + if row["advert_kind"] == "other" { + initialAdverts = row["count"].(int) + } + } + if initialAdverts != 179 { + t.Fatalf("initial unknown adverts=%d want179", initialAdverts) + } + started := time.Now() + ticker := time.NewTicker(time.Second) + defer ticker.Stop() + ticks := 0 + var computeTotal, timeMax, retainedReady time.Duration + for s.advertEvidenceCursor < horizon { + <-ticker.C + ticks++ + if err := s.pollAdvertEvidence(500); err != nil { + t.Fatal(err) + } + if retainedReady == 0 && s.byTxID[size].AdvertRouteEvidence == 3 { + retainedReady = time.Since(started) + } + at := time.Now() + s.GetRelayAirtimeShareWithWindow(window) + d := time.Since(at) + computeTotal += d + if d > timeMax { + timeMax = d + } + } + mixed := 0 + for _, tx := range packets { + if tx.AdvertRouteEvidence == 3 { + mixed++ + } + } + final := s.GetRelayAirtimeShareWithWindow(window) + var finalAdverts int + for _, row := range final["rows"].([]map[string]interface{}) { + if row["advert_kind"] == "mixed" { + finalAdverts = row["count"].(int) + } + } + if finalAdverts != 179 || final["total_count"] != 1785 { + t.Fatalf("final classification/window changed: %v", final) + } + t.Logf("eligible=1785 eligible_adverts=179 final_cursor=%d horizon=%d", s.advertEvidenceCursor, horizon) + t.Logf("CATCHUP evidence_rows=%d retained_packets=%d adverts=30000 poll_rows=500 cadence=1s elapsed=%s ticks=%d cache_misses=%d revisions=%d request_total=%s request_max=%s retained_ready=%s mixed_retained=%d", evidenceCount, size, time.Since(started), ticks, s.cacheMisses, s.advertEvidenceRevision, computeTotal, timeMax, retainedReady, mixed) + if mixed != 30000 { + t.Fatalf("lost retained evidence: mixed=%d want30000", mixed) + } +} diff --git a/cmd/server/advert_route_evidence.go b/cmd/server/advert_route_evidence.go new file mode 100644 index 000000000..fc8336cbc --- /dev/null +++ b/cmd/server/advert_route_evidence.go @@ -0,0 +1,170 @@ +package main + +import ( + "fmt" + "log" + "strings" +) + +// The table is optional for legacy read-only databases. Re-probe while absent +// so a concurrently starting ingestor can create it without a server restart. +func (db *DB) advertEvidencePresent() bool { + if db == nil || db.conn == nil { + return false + } + if db.advertEvidenceTable.Load() { + return true + } + var exists int + if err := db.conn.QueryRow(`SELECT 1 FROM sqlite_master WHERE type='table' AND name='advert_route_evidence'`).Scan(&exists); err != nil { + return false + } + db.advertEvidenceTable.Store(true) + return true +} + +// At most two indexed evidence rows per selected transmission, fetched in +// batches rather than per-packet round trips. No raw frames enter the store. +func (db *DB) advertEvidenceForIDs(ids []int) (map[int]uint8, error) { + result := make(map[int]uint8, len(ids)) + if len(ids) == 0 || !db.advertEvidencePresent() { + return result, nil + } + for start := 0; start < len(ids); start += 500 { + end := start + 500 + if end > len(ids) { + end = len(ids) + } + args := make([]interface{}, end-start) + for i, id := range ids[start:end] { + args[i] = id + } + query := `SELECT tx_id,SUM(bit) FROM advert_route_evidence WHERE tx_id IN (` + strings.TrimSuffix(strings.Repeat("?,", len(args)), ",") + `) GROUP BY tx_id` + if db.advertEvidenceReadHook != nil { + db.advertEvidenceReadHook() + } + rows, err := db.conn.Query(query, args...) + if err != nil { + return nil, err + } + for rows.Next() { + var id int + var mask uint8 + if err := rows.Scan(&id, &mask); err != nil { + rows.Close() + return nil, err + } + result[id] = mask + } + err = rows.Err() + rows.Close() + if err != nil { + return nil, err + } + } + return result, nil +} + +// Caller holds mu. Unioning evidence is independent of cache invalidation so +// loaders and polls can invalidate once per batch, not once per transmission. +func unionAdvertEvidence(tx *StoreTx, mask uint8) bool { + if tx == nil || tx.AdvertRouteEvidence|mask == tx.AdvertRouteEvidence { + return false + } + tx.AdvertRouteEvidence |= mask + return true +} + +// Caller holds mu. The revision prevents an in-flight old computation from +// repopulating an entry after the evidence that produced it has changed. +func (s *PacketStore) invalidateAdvertEvidence() { + s.cacheMu.Lock() + s.advertEvidenceRevision++ + for key := range s.rfCache { + if strings.HasPrefix(key, "relay-airtime-share|") { + delete(s.rfCache, key) + } + } + s.cacheMu.Unlock() +} + +func advertTxIDs(txs []*StoreTx) []int { + ids := make([]int, 0, len(txs)) + for _, tx := range txs { + if tx != nil && tx.PayloadType != nil && *tx.PayloadType == PayloadADVERT { + ids = append(ids, tx.ID) + } + } + return ids +} + +// Caller holds mu. A transmission evicted during the mask read is skipped; +// its persisted union will be fetched if it is loaded again. +func (s *PacketStore) mergeAdvertEvidence(masks map[int]uint8) { + changed := false + for id, mask := range masks { + if unionAdvertEvidence(s.byTxID[id], mask) { + changed = true + } + } + if changed { + s.invalidateAdvertEvidence() + } +} + +func (s *PacketStore) pollAdvertEvidence(limit int) error { + s.advertEvidenceMu.Lock() + defer s.advertEvidenceMu.Unlock() + if !s.db.advertEvidencePresent() { + return nil + } + if limit <= 0 || limit > 500 { + limit = 500 + } + rows, err := s.db.conn.Query(`SELECT id,tx_id,bit FROM advert_route_evidence WHERE id>? ORDER BY id LIMIT ?`, s.advertEvidenceCursor, limit) + if err != nil { + return err + } + type event struct { + id int64 + txID int + bit uint8 + } + var batch []event + for rows.Next() { + var e event + if err := rows.Scan(&e.id, &e.txID, &e.bit); err != nil { + rows.Close() + return err + } + if e.bit != 1 && e.bit != 2 { + rows.Close() + return fmt.Errorf("invalid advert evidence bit %d", e.bit) + } + batch = append(batch, e) + } + err = rows.Err() + rows.Close() + if err != nil { + return err + } + s.mu.Lock() + defer s.mu.Unlock() + changed := false + for _, e := range batch { + if unionAdvertEvidence(s.byTxID[e.txID], e.bit) { + changed = true + } + s.advertEvidenceCursor = e.id + } + if changed { + s.invalidateAdvertEvidence() + } + return nil +} + +func (s *PacketStore) refreshAdvertEvidence() { + if err := s.pollAdvertEvidence(500); err != nil { + log.Printf("[store] advert evidence poll: %v", err) + } +} diff --git a/cmd/server/advert_route_evidence_test.go b/cmd/server/advert_route_evidence_test.go new file mode 100644 index 000000000..3fc57e06b --- /dev/null +++ b/cmd/server/advert_route_evidence_test.go @@ -0,0 +1,484 @@ +package main + +import ( + "database/sql" + "fmt" + "testing" + "time" + + "github.com/meshcore-analyzer/lora" +) + +// Writer-side ingestion is exercised in cmd/ingestor. This fixture represents +// its durable output after F->Z->F (or Z->F->Z), where surviving raw frames +// alone cannot reconstruct the complete set of known route families. +func advertEvidenceFixture(t *testing.T, raw string, bits ...int) (*DB, *sql.DB) { + t.Helper() + path := createTestDBWithResolvedPath(t, 1, []string{"fixture-relay"}) + w, err := sql.Open("sqlite3", path+"?_journal_mode=WAL&_foreign_keys=on") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { w.Close() }) + for _, stmt := range []string{ + `ALTER TABLE transmissions ADD COLUMN from_pubkey TEXT`, + `CREATE TABLE advert_route_evidence (id INTEGER PRIMARY KEY AUTOINCREMENT, tx_id INTEGER NOT NULL REFERENCES transmissions(id) ON DELETE CASCADE, bit INTEGER NOT NULL CHECK (bit IN (1,2)), UNIQUE(tx_id,bit))`, + } { + if _, err := w.Exec(stmt); err != nil { + t.Fatal(err) + } + } + now := time.Now().UTC().Add(-time.Minute).Format(time.RFC3339) + if _, err := w.Exec(`UPDATE transmissions SET raw_hex=?, route_type=?, first_seen=?, from_pubkey='fixture-origin'; UPDATE observations SET raw_hex=?, timestamp=?`, raw, int(raw[1]-'0')&3, now, raw, now); err != nil { + t.Fatal(err) + } + for _, bit := range bits { + if _, err := w.Exec(`INSERT INTO advert_route_evidence(tx_id,bit) VALUES(1,?)`, bit); err != nil { + t.Fatal(err) + } + } + db, err := OpenDB(path) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { db.conn.Close() }) + return db, w +} + +func TestAdvertRouteEvidenceNodeAPIHonorsRequestedLimit(t *testing.T) { + db, w := advertEvidenceFixture(t, "1100aa", 1) + for id := 2; id <= 24; id++ { + if _, err := w.Exec(`INSERT INTO transmissions(id,hash,raw_hex,first_seen,payload_type,route_type,from_pubkey) SELECT ?,?,'1100aa',first_seen,4,1,'fixture-origin' FROM transmissions WHERE id=1`, id, fmt.Sprintf("node-advert-%d", id)); err != nil { + t.Fatal(err) + } + mask := id % 4 + for _, bit := range []int{1, 2} { + if mask&bit != 0 { + if _, err := w.Exec(`INSERT INTO advert_route_evidence(tx_id,bit) VALUES(?,?)`, id, bit); err != nil { + t.Fatal(err) + } + } + } + } + rows, err := db.GetRecentTransmissionsForNode("fixture-origin", 1000) + if err != nil { + t.Fatal(err) + } + if len(rows) != 24 { + t.Fatalf("node recent adverts returned %d, want requested 24 available rows", len(rows)) + } + wants := []string{"other", "flood", "zero_hop", "mixed"} + for _, row := range rows { + id := row["id"].(int) + if row["advert_kind"] != wants[id%4] { + t.Errorf("id=%d kind=%v want %s", id, row["advert_kind"], wants[id%4]) + } + } +} + +// Measures the added bulk-load work and an idle feed at a 30K-packet scale. +// Both paths touch indexed compact evidence, never observation raw frames. +func BenchmarkAdvertEvidence30K(b *testing.B) { + conn, err := sql.Open("sqlite3", ":memory:") + if err != nil { + b.Fatal(err) + } + conn.SetMaxOpenConns(1) + defer conn.Close() + if _, err := conn.Exec(`CREATE TABLE advert_route_evidence(id INTEGER PRIMARY KEY AUTOINCREMENT,tx_id INTEGER,bit INTEGER,UNIQUE(tx_id,bit)); + WITH RECURSIVE ids(n) AS (VALUES(1) UNION ALL SELECT n+1 FROM ids WHERE n<30000) + INSERT INTO advert_route_evidence(tx_id,bit) SELECT n,1 FROM ids UNION ALL SELECT n,2 FROM ids`); err != nil { + b.Fatal(err) + } + db := &DB{conn: conn} + db.advertEvidenceTable.Store(true) + ids := make([]int, 30000) + for i := range ids { + ids[i] = i + 1 + } + b.Run("bulk_load", func(b *testing.B) { + b.ReportAllocs() + for i := 0; i < b.N; i++ { + masks, err := db.advertEvidenceForIDs(ids) + if err != nil || len(masks) != 30000 { + b.Fatalf("mask load: %d %v", len(masks), err) + } + } + }) + b.Run("idle_poll", func(b *testing.B) { + s := NewPacketStore(db, &PacketStoreConfig{}) + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + if err := s.pollAdvertEvidence(500); err != nil { + b.Fatal(err) + } + } + }) +} + +func assertAdvertEvidenceViews(t *testing.T, s *PacketStore, want string) { + t.Helper() + result := s.GetRelayAirtimeShareWithWindow(TimeWindow{}) + rows := result["rows"].([]map[string]interface{}) + if len(rows) != 1 || rows[0]["advert_kind"] != want || rows[0]["count"] != 1 || result["total_count"] != 1 { + t.Errorf("relay airtime must count one hash as %s: %v", want, result) + } + if len(s.packets) != 1 { + t.Fatalf("loaded %d transmissions, want 1", len(s.packets)) + } + tx := s.packets[0] + // Classification must not change the established airtime formula, even + // when evidence includes zero-hop and the packet has a resolved relay. + wantScore := int64(lora.TimeOnAir(len(tx.RawHex)/2, defaultLoRaPreset())) * int64(s.distinctRelayCount(tx)) + if result["total_score"] != wantScore { + t.Errorf("airtime changed: %v, want %d", result["total_score"], wantScore) + } + for _, obs := range tx.Observations { + if obs.RawHex != "" { + t.Error("route evidence must not retain raw frames on StoreObs") + } + } + adverts, err := s.db.GetRecentTransmissionsForNode("fixture-origin", 20) + if err != nil { + t.Fatal(err) + } + if len(adverts) != 1 || adverts[0]["advert_kind"] != want { + t.Errorf("node API must agree on %s: %v", want, adverts) + } +} + +func TestAdvertRouteEvidenceLoadPaths(t *testing.T) { + for _, raw := range []string{"1100aa", "1200aa", "100102030400aa", "130102030400aa"} { + for _, mode := range []string{"cold", "chunk", "chunked_startup", "new_transmission"} { + t.Run(raw+"/"+mode, func(t *testing.T) { + db, _ := advertEvidenceFixture(t, raw, 1, 2) + s := NewPacketStore(db, &PacketStoreConfig{}) + s.useResolvedPathIndex = true + s.initResolvedPathIndex() + maskReads := 0 + db.advertEvidenceReadHook = func() { + maskReads++ + if !s.mu.TryRLock() { + t.Error("bulk evidence SQL must run outside the global store lock") + } else { + s.mu.RUnlock() + } + } + switch mode { + case "chunked_startup": + if err := s.LoadChunked(1); err != nil { + t.Fatal(err) + } + case "cold": + if err := s.Load(); err != nil { + t.Fatal(err) + } + case "chunk": + if err := s.loadChunk(time.Now().Add(-time.Hour), time.Now()); err != nil { + t.Fatal(err) + } + case "new_transmission": + s.IngestNewFromDB(0, 100) + } + if maskReads == 0 { + t.Fatal("load path never fetched durable evidence") + } + db.advertEvidenceReadHook = nil + assertAdvertEvidenceViews(t, s, "mixed") + }) + } + } +} + +func TestAdvertRouteEvidencePollRetryAndUnknownTransmission(t *testing.T) { + db, w := advertEvidenceFixture(t, "1100aa", 1) + s := NewPacketStore(db, &PacketStoreConfig{}) + if err := s.Load(); err != nil { + t.Fatal(err) + } + before := s.advertEvidenceCursor + if _, err := w.Exec(`ALTER TABLE advert_route_evidence RENAME TO evidence_unavailable`); err != nil { + t.Fatal(err) + } + if err := s.pollAdvertEvidence(1); err == nil { + t.Fatal("expected failed feed read") + } + if s.advertEvidenceCursor != before { + t.Fatal("failed feed read advanced cursor") + } + if _, err := w.Exec(`ALTER TABLE evidence_unavailable RENAME TO advert_route_evidence; + INSERT INTO transmissions(id,hash,raw_hex,first_seen,payload_type,route_type,from_pubkey) SELECT 2,'late-advert','1100aa',first_seen,4,1,'fixture-late' FROM transmissions WHERE id=1; + INSERT INTO advert_route_evidence(tx_id,bit) VALUES(2,1),(2,2),(1,2)`); err != nil { + t.Fatal(err) + } + if err := s.pollAdvertEvidence(1); err != nil { + t.Fatal(err) + } + if s.advertEvidenceCursor != before+1 { + t.Fatal("feed ignored batch limit") + } + if s.byTxID[2] != nil { + t.Fatal("feed should not retain unknown transmissions") + } + if err := s.pollAdvertEvidence(500); err != nil { + t.Fatal(err) + } + if s.byTxID[1].AdvertRouteEvidence != 3 { + t.Fatal("retry missed known transmission's late evidence") + } + s.IngestNewFromDB(1, 100) + if tx := s.byTxID[2]; tx == nil || tx.AdvertRouteEvidence != 3 { + t.Fatalf("load after feed consumed unknown event lost mask: %+v", tx) + } +} + +func TestAdvertRouteEvidenceStartupHandoffAndChunkMerge(t *testing.T) { + db, w := advertEvidenceFixture(t, "1100aa", 1) + s := NewPacketStore(db, &PacketStoreConfig{}) + // Arrives after the startup watermark, before the first packet load. + if _, err := w.Exec(`INSERT INTO advert_route_evidence(tx_id,bit) VALUES(1,2)`); err != nil { + t.Fatal(err) + } + if err := s.pollAdvertEvidence(500); err != nil { + t.Fatal(err) + } + if err := s.LoadChunked(1); err != nil { + t.Fatal(err) + } + if s.byTxID[1].AdvertRouteEvidence != 3 { + t.Fatal("startup handoff lost pre-load event") + } + if err := s.loadChunk(time.Now().Add(-time.Hour), time.Now()); err != nil { + t.Fatal(err) + } + for _, tx := range s.packets { + if tx.AdvertRouteEvidence != 3 { + t.Fatal("chunk merge discarded durable bits") + } + } + rows := s.computeRelayAirtimeShare(TimeWindow{})["rows"].([]map[string]interface{}) + if len(rows) != 1 || rows[0]["advert_kind"] != "mixed" || rows[0]["count"] != 1 { + t.Fatalf("chunk merge double-counted mixed hash: %v", rows) + } +} + +func TestAdvertRouteEvidenceEvictionReload(t *testing.T) { + db, _ := advertEvidenceFixture(t, "1200aa", 1, 2) + s := NewPacketStore(db, &PacketStoreConfig{RetentionHours: 1}) + if err := s.Load(); err != nil { + t.Fatal(err) + } + s.mu.Lock() + s.packets[0].FirstSeen = time.Now().Add(-48 * time.Hour).UTC().Format(time.RFC3339) + s.packets[0].LatestSeen = s.packets[0].FirstSeen + n := s.EvictStale() + s.mu.Unlock() + if n != 1 || s.byTxID[1] != nil { + t.Fatalf("eviction removed %d, want 1", n) + } + s.IngestNewFromDB(0, 100) + if tx := s.byTxID[1]; tx == nil || tx.AdvertRouteEvidence != 3 { + t.Fatalf("reload lost persisted mask: %+v", tx) + } +} + +func TestAdvertRouteEvidencePollSameObservationID(t *testing.T) { + for _, tc := range []struct { + name, first, second, kind string + firstBit, secondBit int + }{ + {"flood_zero_flood", "1100aa", "1200aa", "flood", 1, 2}, + {"zero_flood_zero", "1200aa", "1100aa", "zero_hop", 2, 1}, + } { + t.Run(tc.name, func(t *testing.T) { + db, w := advertEvidenceFixture(t, tc.first, tc.firstBit) + s := NewPacketStore(db, &PacketStoreConfig{}) + s.useResolvedPathIndex = true + s.initResolvedPathIndex() + s.rfCacheTTL = time.Hour + if err := s.Load(); err != nil { + t.Fatal(err) + } + initial := s.GetRelayAirtimeShareWithWindow(TimeWindow{}) + if rows := initial["rows"].([]map[string]interface{}); len(rows) != 1 || rows[0]["advert_kind"] != tc.kind { + t.Fatalf("initial fixture classification: %v", initial) + } + unrelated := &cachedResult{expiresAt: time.Now().Add(time.Hour)} + s.rfCache["unrelated-rf-result"] = unrelated + maxObs, maxTx := db.GetMaxObservationID(), db.GetMaxTransmissionID() + // This matches the ingestor's conflict update: neither observation + // ID nor timestamp advances, and the final raw matches the first. + if _, err := w.Exec(`UPDATE observations SET raw_hex=? WHERE id=1; INSERT INTO advert_route_evidence(tx_id,bit) VALUES(1,?); UPDATE observations SET raw_hex=? WHERE id=1`, tc.second, tc.secondBit, tc.first); err != nil { + t.Fatal(err) + } + s.IngestNewFromDB(maxTx, 100) + s.IngestNewObservations(maxObs, 100) + assertAdvertEvidenceViews(t, s, "mixed") + if s.rfCache["unrelated-rf-result"] != unrelated { + t.Error("route evidence invalidated unrelated RF cache") + } + if db.GetMaxObservationID() != maxObs || db.GetMaxTransmissionID() != maxTx { + t.Fatal("fixture unexpectedly created a new observation/transmission") + } + // A fresh store must recover the same classification from disk. + restarted := NewPacketStore(db, &PacketStoreConfig{}) + restarted.useResolvedPathIndex = true + restarted.initResolvedPathIndex() + if err := restarted.Load(); err != nil { + t.Fatal(err) + } + assertAdvertEvidenceViews(t, restarted, "mixed") + }) + } +} + +func TestAdvertRouteEvidenceMissingTableIsUnknown(t *testing.T) { + db, w := advertEvidenceFixture(t, "1100aa") + if _, err := w.Exec(`DROP TABLE advert_route_evidence`); err != nil { + t.Fatal(err) + } + // Reopen after removing the table to exercise legacy read-only detection. + legacy, err := OpenDB(db.path) + if err != nil { + t.Fatal(err) + } + defer legacy.conn.Close() + s := NewPacketStore(legacy, &PacketStoreConfig{}) + s.useResolvedPathIndex = true + s.initResolvedPathIndex() + if err := s.Load(); err != nil { + t.Fatal(err) + } + assertAdvertEvidenceViews(t, s, "other") + var count int + if err := w.QueryRow(`SELECT COUNT(*) FROM sqlite_master WHERE name='advert_route_evidence'`).Scan(&count); err != nil { + t.Fatal(err) + } + if count != 0 { + t.Fatal("read-only server created legacy evidence table") + } +} + +func TestAdvertRouteEvidenceTableAppearsAfterStartup(t *testing.T) { + db, w := advertEvidenceFixture(t, "1100aa") + var schema string + if err := w.QueryRow(`SELECT sql FROM sqlite_master WHERE name='advert_route_evidence'`).Scan(&schema); err != nil { + t.Fatal(err) + } + if _, err := w.Exec(`DROP TABLE advert_route_evidence`); err != nil { + t.Fatal(err) + } + s := NewPacketStore(db, &PacketStoreConfig{}) + if err := s.Load(); err != nil { + t.Fatal(err) + } + if s.byTxID[1].AdvertRouteEvidence != 0 { + t.Fatal("legacy load invented evidence") + } + s.GetRelayAirtimeShareWithWindow(TimeWindow{}) + if _, err := w.Exec(schema + `; INSERT INTO advert_route_evidence(tx_id,bit) VALUES(1,1),(1,2)`); err != nil { + t.Fatal(err) + } + if err := s.pollAdvertEvidence(500); err != nil { + t.Fatal(err) + } + rows := s.GetRelayAirtimeShareWithWindow(TimeWindow{})["rows"].([]map[string]interface{}) + if len(rows) != 1 || rows[0]["advert_kind"] != "mixed" { + t.Fatalf("late migration did not refresh legacy store: %v", rows) + } +} + +func TestAdvertRouteEvidencePollRejectsPartialBatch(t *testing.T) { + db, w := advertEvidenceFixture(t, "1100aa", 1) + w.SetMaxOpenConns(1) + s := NewPacketStore(db, &PacketStoreConfig{}) + if err := s.Load(); err != nil { + t.Fatal(err) + } + before := s.advertEvidenceCursor + if _, err := w.Exec(`PRAGMA ignore_check_constraints=ON; INSERT INTO advert_route_evidence(tx_id,bit) VALUES(1,2),(1,3)`); err != nil { + t.Fatal(err) + } + if err := s.pollAdvertEvidence(500); err == nil { + t.Fatal("corrupt feed row accepted") + } + if s.advertEvidenceCursor != before || s.byTxID[1].AdvertRouteEvidence != 1 { + t.Fatal("partial failed batch changed evidence or cursor") + } + if _, err := w.Exec(`DELETE FROM advert_route_evidence WHERE bit=3; PRAGMA ignore_check_constraints=OFF`); err != nil { + t.Fatal(err) + } + if err := s.pollAdvertEvidence(500); err != nil { + t.Fatal(err) + } + if s.byTxID[1].AdvertRouteEvidence != 3 { + t.Fatal("retry lost valid event preceding corrupt row") + } +} + +func TestAdvertRouteEvidenceInvalidatesOncePerBatch(t *testing.T) { + s := newRelayAirtimeShareTestStore([]*StoreTx{ + makeRelayAirtimeTx(1, PayloadADVERT, 3, 0, "batch-1"), + makeRelayAirtimeTx(2, PayloadADVERT, 3, 0, "batch-2"), + }) + other := &cachedResult{} + s.rfCache["unrelated-rf-result"] = other + s.rfCache["relay-airtime-share|"] = &cachedResult{} + s.rfCache["relay-airtime-share|7d"] = &cachedResult{} + masks := map[int]uint8{1: 3, 2: 3} + s.mu.Lock() + s.mergeAdvertEvidence(masks) + s.mergeAdvertEvidence(masks) + s.mu.Unlock() + if s.advertEvidenceRevision != 1 { + t.Fatalf("batch invalidated %d times, want one; unchanged replay must not invalidate", s.advertEvidenceRevision) + } + if len(s.rfCache) != 1 || s.rfCache["unrelated-rf-result"] != other { + t.Fatal("batch must invalidate only relay airtime entries") + } +} + +func TestRelayAirtimeShareAdvertDuplicateHashEvidenceUnion(t *testing.T) { + for _, reverse := range []bool{false, true} { + flood, direct := RouteFlood, RouteDirect + first := advertAirtimeTx(1, &flood, "1100aa") + second := advertAirtimeTx(2, &direct, "1200aa") + second.Hash = first.Hash + packets := []*StoreTx{first, second} + if reverse { + packets[0], packets[1] = packets[1], packets[0] + } + s := newRelayAirtimeShareTestStore(packets) + // Keep relay evidence identical so only classification changes. + s.addToResolvedPubkeyIndex(1, []string{"fixture-relay"}) + s.addToResolvedPubkeyIndex(2, []string{"fixture-relay"}) + result := s.computeRelayAirtimeShare(TimeWindow{}) + rows := result["rows"].([]map[string]interface{}) + if len(rows) != 1 || rows[0]["advert_kind"] != "mixed" || rows[0]["count"] != 1 || result["total_count"] != 1 { + t.Errorf("reverse=%v: duplicate hash must have one mixed bucket: %v", reverse, result) + } + if result["total_score"] != int64(lora.TimeOnAir(3, defaultLoRaPreset())) { + t.Errorf("reverse=%v: duplicate hash inflated score: %v", reverse, result) + } + } +} + +func TestRelayAirtimeShareAdvertEvidenceIndependentOfReportWindow(t *testing.T) { + flood, direct := RouteFlood, RouteDirect + inside := advertAirtimeTx(1, &direct, "1200aa") + outside := advertAirtimeTx(2, &flood, "1100aa") + outside.Hash = inside.Hash + outside.FirstSeen = "2025-12-01T00:00:00Z" + s := newRelayAirtimeShareTestStore([]*StoreTx{inside, outside}) + s.addToResolvedPubkeyIndex(inside.ID, []string{"fixture-relay"}) + s.addToResolvedPubkeyIndex(outside.ID, []string{"fixture-relay", "excluded-relay"}) + result := s.computeRelayAirtimeShare(TimeWindow{Since: "2026-01-01T00:00:00Z", Until: "2026-01-02T00:00:00Z"}) + rows := result["rows"].([]map[string]interface{}) + if len(rows) != 1 || rows[0]["advert_kind"] != "mixed" || result["total_count"] != 1 { + t.Fatalf("report window changed known hash evidence: %v", result) + } + if result["total_score"] != int64(lora.TimeOnAir(3, defaultLoRaPreset())) { + t.Fatalf("excluded record contributed airtime score: %v", result) + } +} diff --git a/cmd/server/advert_stats.go b/cmd/server/advert_stats.go index a67d790f6..fcde4b54a 100644 --- a/cmd/server/advert_stats.go +++ b/cmd/server/advert_stats.go @@ -1,30 +1,30 @@ package main -import "time" +import ( + "errors" + "time" -// advertRouteTypeFlood is ROUTE_TYPE_FLOOD from the MeshCore packet header. -// Named distinctly from the equivalent constant in the (still open) unscoped- -// relay PR so the two changes merge independently. -const advertRouteTypeFlood = 1 + "github.com/meshcore-analyzer/packetpath" +) // floodAdvertEntry is one advert transmission originated by a node, reduced to -// what the windowed flood-advert count needs: first-seen timestamp, route type +// what the windowed flood-advert count needs: first-seen timestamp, known route evidence // and packet hash (for dedup across re-ingests / multi-observer rows). type floodAdvertEntry struct { ts string - rt int + mask uint8 hash string } -// countFloodAdverts counts distinct flood adverts (route_type == -// advertRouteTypeFlood) whose first-seen lies within the past windowHours. Entries +// countFloodAdverts counts distinct adverts with known flood evidence whose +// first-seen lies within the past windowHours. Entries // with unparseable timestamps are skipped, matching relay-liveness behaviour; // entries without a hash fall back to their timestamp as the dedup key. func countFloodAdverts(entries []floodAdvertEntry, now time.Time, windowHours float64) int { cutoff := now.Add(-time.Duration(windowHours * float64(time.Hour))) seen := map[string]struct{}{} for _, e := range entries { - if e.rt != advertRouteTypeFlood { + if e.mask&packetpath.AdvertFlood == 0 { continue } t, ok := parseRelayTS(e.ts) @@ -40,16 +40,13 @@ func countFloodAdverts(entries []floodAdvertEntry, now time.Time, windowHours fl return len(seen) } -// CountFloodAdvertsForNode returns how many distinct FLOOD adverts pubkey -// originated in the last windowHours - the mesh-wide-airtime kind. Zero-hop -// adverts (route_type DIRECT) are excluded, so a nearby observer hearing a -// node's cheap local adverts does not inflate the number. +// CountFloodAdvertsForNode counts distinct adverts with recorded flood evidence, +// including mixed and transport-flood frames, independently of the first route. +// It is a lower bound while backfill is pending or evidence writes have failed. +// With no evidence table the count is unavailable, not a proven zero. // -// route_type is filtered in SQL so an advert-spamming node cannot truncate -// the flood count (review feedback on the earlier LIMIT approach). The time -// floor is a DATE-ONLY string with one day of slack: a date prefix compares -// lexically the same across every first_seen format parseRelayTS accepts -// ('T' and ' ' separators alike); the exact window check stays in Go. +// SQL filters known flood evidence before the row cap. The date-only floor has +// one day of slack for mixed timestamp formats; the exact window stays in Go. // // The row cap is a pure safety valve on per-request allocation: it applies to // flood adverts inside the floor window only, and 50000 in ~8 days is ~4 per @@ -61,10 +58,16 @@ func countFloodAdverts(entries []floodAdvertEntry, now time.Time, windowHours fl const floodAdvertRowCap = 50000 func (db *DB) CountFloodAdvertsForNode(pubkey string, windowHours float64, rowCap int) (int, error) { + if !db.advertEvidencePresent() { + return 0, errors.New("advert route evidence unavailable") + } floor := time.Now().UTC().Add(-time.Duration(windowHours*float64(time.Hour))).AddDate(0, 0, -1).Format("2006-01-02") rows, err := db.conn.Query( - "SELECT COALESCE(first_seen, ''), COALESCE(route_type, -1), COALESCE(hash, '') FROM transmissions WHERE from_pubkey = ? AND payload_type = ? AND route_type = ? AND first_seen >= ? ORDER BY id DESC LIMIT ?", - pubkey, payloadTypeAdvert, advertRouteTypeFlood, floor, rowCap) + `SELECT COALESCE(first_seen, ''), ?, COALESCE(hash, '') FROM transmissions t + WHERE from_pubkey=? AND payload_type=? AND first_seen>=? + AND EXISTS(SELECT 1 FROM advert_route_evidence e WHERE e.tx_id=t.id AND e.bit=?) + ORDER BY id DESC LIMIT ?`, + packetpath.AdvertFlood, pubkey, payloadTypeAdvert, floor, packetpath.AdvertFlood, rowCap) if err != nil { return 0, err } @@ -72,10 +75,10 @@ func (db *DB) CountFloodAdvertsForNode(pubkey string, windowHours float64, rowCa var entries []floodAdvertEntry for rows.Next() { var e floodAdvertEntry - if err := rows.Scan(&e.ts, &e.rt, &e.hash); err != nil { + if err := rows.Scan(&e.ts, &e.mask, &e.hash); err != nil { return 0, err } entries = append(entries, e) } - return countFloodAdverts(entries, time.Now(), windowHours), nil + return countFloodAdverts(entries, time.Now(), windowHours), rows.Err() } diff --git a/cmd/server/advert_stats_test.go b/cmd/server/advert_stats_test.go index 099fac2cd..942a094ae 100644 --- a/cmd/server/advert_stats_test.go +++ b/cmd/server/advert_stats_test.go @@ -2,9 +2,12 @@ package main import ( "encoding/json" + "fmt" "net/http/httptest" "testing" "time" + + "github.com/meshcore-analyzer/packetpath" ) // insertAdvertTx seeds one advert transmission row - the single place that @@ -16,6 +19,15 @@ func insertAdvertTx(t *testing.T, db *DB, pubkey, hash string, rt int, ts time.T hash, ts.Format("2006-01-02T15:04:05.000Z"), rt, payloadTypeAdvert, pubkey); err != nil { t.Fatalf("insert transmission: %v", err) } + // Fixtures represent ingestor output: canonical routing alone is not evidence. + if _, err := db.conn.Exec(`CREATE TABLE IF NOT EXISTS advert_route_evidence(id INTEGER PRIMARY KEY AUTOINCREMENT,tx_id INTEGER,bit INTEGER,UNIQUE(tx_id,bit))`); err != nil { + t.Fatal(err) + } + if rt == RouteFlood { + if _, err := db.conn.Exec(`INSERT INTO advert_route_evidence(tx_id,bit) SELECT id,1 FROM transmissions WHERE hash=?`, hash); err != nil { + t.Fatal(err) + } + } } func advertTS(hoursAgo float64) string { @@ -27,21 +39,21 @@ func advertTS(hoursAgo float64) string { func TestCountFloodAdverts(t *testing.T) { now := time.Now() entries := []floodAdvertEntry{ - {ts: advertTS(1), rt: advertRouteTypeFlood, hash: "a1"}, - {ts: advertTS(2), rt: advertRouteTypeFlood, hash: "a1"}, // dup hash: one advert, two rows - {ts: advertTS(3), rt: advertRouteTypeFlood, hash: "a2"}, - {ts: advertTS(4), rt: 0, hash: "a3"}, // zero-hop (DIRECT): excluded - {ts: advertTS(9 * 24), rt: advertRouteTypeFlood, hash: "a4"}, // outside 7d window - {ts: "not-a-time", rt: advertRouteTypeFlood, hash: "a5"}, // unparseable: skipped - {ts: advertTS(5), rt: -1, hash: "a6"}, // route type absent: excluded + {ts: advertTS(1), mask: packetpath.AdvertFlood, hash: "a1"}, + {ts: advertTS(2), mask: packetpath.AdvertFlood, hash: "a1"}, // dup hash: one advert, two rows + {ts: advertTS(3), mask: packetpath.AdvertFlood, hash: "a2"}, + {ts: advertTS(4), mask: packetpath.AdvertDirectEmptyPath, hash: "a3"}, // direct-only evidence: excluded + {ts: advertTS(9 * 24), mask: packetpath.AdvertFlood, hash: "a4"}, // outside 7d window + {ts: "not-a-time", mask: packetpath.AdvertFlood, hash: "a5"}, // unparseable: skipped + {ts: advertTS(5), mask: 0, hash: "a6"}, // no known evidence: excluded } if got := countFloodAdverts(entries, now, 7*24); got != 2 { t.Fatalf("want 2 flood adverts in window, got %d", got) } // Hash-less entries dedup by timestamp instead of collapsing into one. hashless := []floodAdvertEntry{ - {ts: advertTS(1), rt: advertRouteTypeFlood}, - {ts: advertTS(2), rt: advertRouteTypeFlood}, + {ts: advertTS(1), mask: packetpath.AdvertFlood}, + {ts: advertTS(2), mask: packetpath.AdvertFlood}, } if got := countFloodAdverts(hashless, now, 7*24); got != 2 { t.Fatalf("want 2 hash-less flood adverts, got %d", got) @@ -52,6 +64,9 @@ func TestCountFloodAdverts(t *testing.T) { // counting recent flood adverts only (zero-hop and out-of-window excluded). func TestNodeDetailIncludesFloodAdvertCount(t *testing.T) { srv, router := setupTestServer(t) + if _, err := srv.db.conn.Exec(`CREATE TABLE advert_route_evidence(id INTEGER PRIMARY KEY AUTOINCREMENT,tx_id INTEGER,bit INTEGER,UNIQUE(tx_id,bit))`); err != nil { + t.Fatal(err) + } now := time.Now().UTC() ins := func(hash string, rt int, ts time.Time) { insertAdvertTx(t, srv.db, "aabbccdd11223344", hash, rt, ts) @@ -82,13 +97,13 @@ func TestNodeDetailIncludesFloodAdvertCount(t *testing.T) { } before := fetch() - ins("fa1", advertRouteTypeFlood, now.Add(-2*time.Hour)) - ins("fa2", advertRouteTypeFlood, now.Add(-30*time.Hour)) - ins("za1", 0, now.Add(-1*time.Hour)) // zero-hop: excluded - ins("fa3", advertRouteTypeFlood, now.Add(-9*24*time.Hour)) // outside the 7d window + ins("fa1", RouteFlood, now.Add(-2*time.Hour)) + ins("fa2", RouteFlood, now.Add(-30*time.Hour)) + ins("za1", 0, now.Add(-1*time.Hour)) // zero-hop: excluded + ins("fa3", RouteFlood, now.Add(-9*24*time.Hour)) // outside the 7d window // Inside the SQL date floor (window + 1d slack) but outside the exact 7d // window - only the Go-side check rejects this one. - ins("fa4", advertRouteTypeFlood, now.Add(-time.Duration(7.5*24)*time.Hour)) + ins("fa4", RouteFlood, now.Add(-time.Duration(7.5*24)*time.Hour)) if got := fetch(); got != before+2 { t.Fatalf("want flood_advert_count_7d = %v+2, got %v", before, got) @@ -101,7 +116,7 @@ func TestCountFloodAdvertsForNode_RowCapSaturates(t *testing.T) { db := setupTestDB(t) now := time.Now().UTC() for i, h := range []string{"cap1", "cap2", "cap3"} { - insertAdvertTx(t, db, "capnode11223344", h, advertRouteTypeFlood, now.Add(-time.Duration(i+1)*time.Hour)) + insertAdvertTx(t, db, "capnode11223344", h, RouteFlood, now.Add(-time.Duration(i+1)*time.Hour)) } n, err := db.CountFloodAdvertsForNode("capnode11223344", 7*24, 2) if err != nil { @@ -111,3 +126,57 @@ func TestCountFloodAdvertsForNode_RowCapSaturates(t *testing.T) { t.Fatalf("want saturated count 2, got %d", n) } } + +func TestNodeDetailFloodCountAgreesWithRouteEvidence(t *testing.T) { + for _, firstRoute := range []int{RouteDirect, RouteFlood, RouteTransportFlood} { + t.Run(fmt.Sprint(firstRoute), func(t *testing.T) { + srv, router := setupTestServer(t) + if _, err := srv.db.conn.Exec(`CREATE TABLE IF NOT EXISTS advert_route_evidence(id INTEGER PRIMARY KEY AUTOINCREMENT, tx_id INTEGER, bit INTEGER, UNIQUE(tx_id,bit))`); err != nil { + t.Fatal(err) + } + const pubkey = "aabbccdd11223344" + if _, err := srv.db.conn.Exec(`DELETE FROM transmissions WHERE from_pubkey=?`, pubkey); err != nil { + t.Fatal(err) + } + insertAdvertTx(t, srv.db, pubkey, "mixed-advert", firstRoute, time.Now().Add(-time.Hour)) + if _, err := srv.db.conn.Exec(`INSERT INTO advert_route_evidence(tx_id,bit) SELECT id,1 FROM transmissions WHERE hash='mixed-advert' AND NOT EXISTS(SELECT 1 FROM advert_route_evidence WHERE tx_id=transmissions.id AND bit=1) UNION ALL SELECT id,2 FROM transmissions WHERE hash='mixed-advert'`); err != nil { + t.Fatal(err) + } + w := httptest.NewRecorder() + router.ServeHTTP(w, httptest.NewRequest("GET", "/api/nodes/"+pubkey, nil)) + var response NodeDetailResponse + if err := json.Unmarshal(w.Body.Bytes(), &response); err != nil { + t.Fatal(err) + } + if w.Code != 200 || len(response.RecentAdverts) != 1 || response.RecentAdverts[0]["advert_kind"] != "mixed" { + t.Fatalf("unexpected node response: %d %+v", w.Code, response) + } + if response.Node["flood_advert_count_7d"] != float64(1) { + t.Errorf("mixed advert must contribute one flood count regardless of first route %d: %v", firstRoute, response.Node["flood_advert_count_7d"]) + } + }) + } +} + +func TestNodeDetailLegacyFloodCountIsUnavailable(t *testing.T) { + srv, router := setupTestServer(t) + if _, err := srv.db.conn.Exec(`DROP TABLE IF EXISTS advert_route_evidence`); err != nil { + t.Fatal(err) + } + insertAdvertTx(t, srv.db, "aabbccdd11223344", "legacy-flood", RouteFlood, time.Now().Add(-time.Hour)) + if _, err := srv.db.conn.Exec(`DROP TABLE advert_route_evidence`); err != nil { + t.Fatal(err) + } + w := httptest.NewRecorder() + router.ServeHTTP(w, httptest.NewRequest("GET", "/api/nodes/aabbccdd11223344", nil)) + var response NodeDetailResponse + if err := json.Unmarshal(w.Body.Bytes(), &response); err != nil { + t.Fatal(err) + } + if w.Code != 200 { + t.Fatal(w.Code) + } + if _, ok := response.Node["flood_advert_count_7d"]; ok { + t.Errorf("missing evidence table must not report a proven count: %v", response.Node["flood_advert_count_7d"]) + } +} diff --git a/cmd/server/chunked_load.go b/cmd/server/chunked_load.go index 1b9c043b9..0a82b33f3 100644 --- a/cmd/server/chunked_load.go +++ b/cmd/server/chunked_load.go @@ -433,13 +433,18 @@ func (s *PacketStore) LoadChunked(chunkSize int) error { ORDER BY t.id ASC, o.timestamp DESC` } + // Acquire before opening the cursor: a feed waiting for the same + // connection must never hold this gate while we hold its cursor. + s.advertEvidenceMu.Lock() rows, err := s.db.conn.Query(chunkSQL) if err != nil { + s.advertEvidenceMu.Unlock() return fmt.Errorf("chunk %d: query: %w", chunkIdx, err) } chunkTxCount, lastID, err := s.scanAndMergeChunk(rows, relayPM, &coldLoadAmbiguousHopsSkipped) rows.Close() + s.advertEvidenceMu.Unlock() if err != nil { return fmt.Errorf("chunk %d: scan: %w", chunkIdx, err) } @@ -617,10 +622,9 @@ func (s *PacketStore) scanAndMergeChunk(rows *sql.Rows, relayPM *prefixMap, cold RSSI: nullFloatPtr(rssi), Score: nullIntPtr(score), PathJSON: obsPJ, - // obs.RawHex deliberately NOT stored: it duplicates the parent - // tx.RawHex (same content hash ⇒ same frame) and enrichObs falls - // back to tx.RawHex when obs.RawHex == "". obsRawHex is still - // scanned to keep scanArgs aligned with the o.raw_hex column. + // Raw frames stay in SQLite; hash equality does not imply route + // equality. Only compact advert evidence is retained per tx. + // Packet-detail queries can read the original observation raw. Timestamp: normalizeTimestamp(nullStrVal(obsTimestamp)), } @@ -670,6 +674,19 @@ func (s *PacketStore) scanAndMergeChunk(rows *sql.Rows, relayPM *prefixMap, cold if err := rows.Err(); err != nil { return len(seenTxIDs), maxID, err } + rows.Close() + txs := make([]*StoreTx, 0, len(seenTxIDs)) + for id := range seenTxIDs { + txs = append(txs, s.byTxID[id]) + } + ids := advertTxIDs(txs) + s.mu.Unlock() + masks, err := s.db.advertEvidenceForIDs(ids) + s.mu.Lock() + if err != nil { + return len(seenTxIDs), maxID, err + } + s.mergeAdvertEvidence(masks) return len(seenTxIDs), maxID, nil } diff --git a/cmd/server/db.go b/cmd/server/db.go index e6f1a3dd3..115d2d5a2 100644 --- a/cmd/server/db.go +++ b/cmd/server/db.go @@ -17,6 +17,7 @@ import ( _ "github.com/mattn/go-sqlite3" "github.com/meshcore-analyzer/dbschema" "github.com/meshcore-analyzer/geofilter" + "github.com/meshcore-analyzer/packetpath" "golang.org/x/sync/singleflight" ) @@ -37,6 +38,8 @@ const routeTypeNonTransportSQL = "route_type IN (1, 2)" // DB wraps a read-only connection to the MeshCore SQLite database. type DB struct { + advertEvidenceTable atomic.Bool + advertEvidenceReadHook func() // test-only: immediately before a bulk mask query conn *sql.DB path string // filesystem path to the database file isV3 bool // v3 schema: observer_idx in observations (vs observer_id in v2) @@ -1323,6 +1326,20 @@ func (db *DB) GetRecentTransmissionsForNode(pubkey string, limit int) ([]map[str } } + if err := rows.Err(); err != nil { + return nil, err + } + rows.Close() + masks, err := db.advertEvidenceForIDs(txIDs) + if err != nil { + return nil, err + } + for _, p := range packets { + if p["payload_type"] == 4 { + p["advert_kind"] = packetpath.AdvertKind(masks[p["id"].(int)]) + } + } + // Fetch observations for all transmissions if len(txIDs) > 0 { obsMap := db.getObservationsForTransmissions(txIDs) diff --git a/cmd/server/relay_airtime_adverts_test.go b/cmd/server/relay_airtime_adverts_test.go new file mode 100644 index 000000000..b20d99aa8 --- /dev/null +++ b/cmd/server/relay_airtime_adverts_test.go @@ -0,0 +1,246 @@ +package main + +import ( + "fmt" + "math" + "reflect" + "strings" + "testing" + "time" + + "github.com/meshcore-analyzer/lora" + "github.com/meshcore-analyzer/packetpath" +) + +func advertAirtimeTx(id int, route *int, raw string) *StoreTx { + tx := makeRelayAirtimeTx(id, PayloadADVERT, 120, 0, fmt.Sprintf("advert-%d", id)) + tx.RouteType = route + tx.RawHex = raw + tx.AdvertRouteEvidence = packetpath.AdvertRouteEvidence(raw) + return tx +} + +func TestRelayAirtimeShare_AdvertRouting(t *testing.T) { + route := func(n int) *int { return &n } + for _, tc := range []struct { + name string + route *int + raw, path, want string + }{ + {"flood without hops", route(1), "1100aa", "[]", "flood"}, + {"transport flood without hops", route(0), "100000000000aa", "[]", "flood"}, + {"flood with hops", route(1), "1101ffaa", "[]", "flood"}, + {"direct zero hop", route(2), "1200aa", `["ff"]`, "zero_hop"}, + {"transport direct zero hop", route(3), "130102030400aa", `["ff"]`, "zero_hop"}, + {"two byte zero count", route(2), "1240aa", "[]", "zero_hop"}, + {"three byte zero count", route(2), "1280aa", "[]", "zero_hop"}, + {"direct nonempty original path", route(2), "1201ffaa", "[]", "other"}, + {"transport direct nonempty original path", route(3), "130000000001ffaa", "[]", "other"}, + {"raw evidence without canonical route", nil, "1200aa", "[]", "zero_hop"}, + {"raw evidence with invalid canonical route", route(9), "1200aa", "[]", "zero_hop"}, + {"missing raw", route(2), "", "[]", "other"}, + {"missing path byte", route(2), "12", "[]", "other"}, + {"truncated transport", route(3), "130000", "[]", "other"}, + {"missing transport path byte", route(3), "1300000000", "[]", "other"}, + {"invalid header", route(2), "zz00aa", "[]", "other"}, + {"raw evidence overrides canonical route", route(2), "1100aa", "[]", "flood"}, + {"conflicting raw payload", route(2), "1600aa", "[]", "other"}, + {"invalid path hex", route(2), "12zzaa", "[]", "other"}, + {"invalid transport hex", route(3), "13zz00000000aa", "[]", "other"}, + {"reserved path encoding", route(2), "12c0aa", "[]", "other"}, + {"odd raw hex", route(2), "1200a", "[]", "other"}, + {"no payload bytes", route(2), "1200", "[]", "other"}, + } { + t.Run(tc.name, func(t *testing.T) { + tx := advertAirtimeTx(1, tc.route, tc.raw) + tx.PathJSON = tc.path // Longest observation's display path is not the original path. + result := newRelayAirtimeShareTestStore([]*StoreTx{tx}).computeRelayAirtimeShare(TimeWindow{}) + rows := result["rows"].([]map[string]interface{}) + if len(rows) != 1 || rows[0]["advert_kind"] != tc.want { + t.Fatalf("advert_kind = %v, want %q", rows, tc.want) + } + if rows[0]["payload_type"] != "ADVERT" || rows[0]["type"] != PayloadADVERT || rows[0]["count"] != 1 { + t.Fatalf("legacy advert fields changed: %v", rows[0]) + } + }) + } +} + +func TestRelayAirtimeShare_AdvertSplitConservesMetric(t *testing.T) { + flood, direct := RouteFlood, RouteDirect + packets := []*StoreTx{ + advertAirtimeTx(1, &flood, "1100"+strings.Repeat("ab", 118)), + advertAirtimeTx(2, &direct, "1200"+strings.Repeat("ab", 118)), + advertAirtimeTx(3, &direct, "1200"+strings.Repeat("ab", 118)), + advertAirtimeTx(4, nil, strings.Repeat("ab", 120)), + makeRelayAirtimeTx(5, PayloadACK, 10, 0, "ack"), + } + store := newRelayAirtimeShareTestStore(packets) + store.addToResolvedPubkeyIndex(1, []string{"relay-a", "relay-b", "relay-a"}) + store.addToResolvedPubkeyIndex(3, []string{"relay-c"}) // Do not force direct evidence to zero. + store.addToResolvedPubkeyIndex(4, []string{"relay-d"}) + store.addToResolvedPubkeyIndex(5, []string{"relay-e"}) + result := store.computeRelayAirtimeShare(TimeWindow{}) + rows := result["rows"].([]map[string]interface{}) + if len(rows) != 4 { + t.Fatalf("got %d buckets, want flood, zero-hop, other and ACK: %v", len(rows), rows) + } + toa := int64(lora.TimeOnAir(120, defaultLoRaPreset())) + ackToA := int64(lora.TimeOnAir(10, defaultLoRaPreset())) + wantScore := 4*toa + ackToA + if result["total_count"] != 5 || result["total_score"] != wantScore { + t.Fatalf("totals changed: %v", result) + } + countSum, scoreSum, countPctSum, airtimePctSum := 0, int64(0), 0.0, 0.0 + for _, row := range rows { + count := row["count"].(int) + score := row["score"].(int64) + countSum += count + scoreSum += score + countPctSum += row["count_pct"].(float64) + airtimePctSum += row["airtime_pct"].(float64) + if math.Abs(row["count_pct"].(float64)-float64(count)/5*100) > 1e-9 || math.Abs(row["airtime_pct"].(float64)-float64(score)/float64(wantScore)*100) > 1e-9 { + t.Fatalf("incorrect percentage: %v", row) + } + switch row["advert_kind"] { + case "flood": + if count != 1 || score != 2*toa { + t.Fatalf("flood: %v", row) + } + case "zero_hop": + if count != 2 || score != toa { + t.Fatalf("zero-hop: %v", row) + } + case "other": + if count != 1 || score != toa { + t.Fatalf("other: %v", row) + } + default: + if row["payload_type"] != "ACK" || row["type"] != PayloadACK || count != 1 || score != ackToA { + t.Fatalf("non-advert changed: %v", row) + } + } + } + if countSum != 5 || scoreSum != wantScore || math.Abs(countPctSum-100) > 1e-9 || math.Abs(airtimePctSum-100) > 1e-9 { + t.Fatalf("shares not conserved: %d %d %f %f", countSum, scoreSum, countPctSum, airtimePctSum) + } +} + +func TestRelayAirtimeShare_AdvertWindowDedupCacheAndZeroRelays(t *testing.T) { + flood, direct := RouteFlood, RouteDirect + first := advertAirtimeTx(1, &direct, "1200aa") + duplicate := advertAirtimeTx(2, &flood, "1100aa") + duplicate.Hash = first.Hash + outside := advertAirtimeTx(3, &flood, "1100aa") + outside.FirstSeen = "2025-12-01T00:00:00Z" + ack := makeRelayAirtimeTx(4, PayloadACK, 10, 1, "ack") + store := newRelayAirtimeShareTestStore([]*StoreTx{first, duplicate, outside, ack}) + store.rfCacheTTL = time.Minute + store.addToResolvedPubkeyIndex(2, []string{"relay-a"}) + store.addToResolvedPubkeyIndex(3, []string{"relay-b"}) + store.addToResolvedPubkeyIndex(4, []string{"relay-c"}) + window := TimeWindow{Since: "2026-01-01T00:00:00Z", Until: "2026-01-02T00:00:00Z", Label: "test"} + result := store.GetRelayAirtimeShareWithWindow(window) + cached := store.GetRelayAirtimeShareWithWindow(window) + if result["cached"] != false || cached["cached"] != true || result["window"] != "test" || !reflect.DeepEqual(result["rows"], cached["rows"]) { + t.Fatalf("cache contract changed: %v / %v", result, cached) + } + rows := result["rows"].([]map[string]interface{}) + if len(rows) != 2 || result["total_count"] != 2 || rows[1]["advert_kind"] != "mixed" || rows[1]["count"] != 1 || rows[1]["score"] != int64(0) || rows[1]["airtime_pct"] != float64(0) || rows[1]["count_pct"] != float64(50) { + t.Fatalf("zero relay row/window/dedup changed: %v", result) + } + if store.GetRelayAirtimeShareWithWindow(TimeWindow{})["total_count"] != 3 { + t.Fatal("time-window caches collided") + } +} + +func TestRelayAirtimeShare_AdvertStableTies(t *testing.T) { + flood, direct := RouteFlood, RouteDirect + store := newRelayAirtimeShareTestStore([]*StoreTx{advertAirtimeTx(1, &direct, "1200aa"), advertAirtimeTx(2, nil, "000000"), advertAirtimeTx(3, &flood, "1100aa")}) + for i := 0; i < 25; i++ { + rows := store.computeRelayAirtimeShare(TimeWindow{})["rows"].([]map[string]interface{}) + if len(rows) != 3 { + t.Fatalf("got %d rows, want 3", len(rows)) + } + for j, want := range []string{"flood", "other", "zero_hop"} { + if rows[j]["advert_kind"] != want { + t.Fatalf("unstable tie order: %v", rows) + } + } + } +} + +// Same workload is measured before/after #2041; construction is outside timing. +func BenchmarkRelayAirtimeShare30K(b *testing.B) { + packets := make([]*StoreTx, 30000) + for i := range packets { + pt := PayloadACK + if i%3 == 0 { + pt = PayloadADVERT + } + tx := makeRelayAirtimeTx(i+1, pt, 120, 0, fmt.Sprintf("packet-%d", i)) + route := i % 4 + tx.RouteType = &route + header := fmt.Sprintf("%02x", pt<<2|route) + if route == 0 || route == 3 { + header += "01020304" + } + tx.RawHex = header + "00" + strings.Repeat("ab", 120-len(header)/2-1) + tx.AdvertRouteEvidence = packetpath.AdvertRouteEvidence(tx.RawHex) + packets[i] = tx + } + store := newRelayAirtimeShareTestStore(packets) + for _, tx := range packets { + if tx.ID%4 != 0 { + store.addToResolvedPubkeyIndex(tx.ID, []string{"relay-a", "relay-b", "relay-c"}) + } + } + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + store.computeRelayAirtimeShare(TimeWindow{}) + } +} + +// Short windows are the allocation worst case: retained history is much larger +// than the selected 1-hour slice. All adverts exercise the maximum evidence map. +func BenchmarkRelayAirtimeShareRetentionWindow(b *testing.B) { + for _, size := range []int{30000, 300000} { + packets := make([]*StoreTx, size) + for i := range packets { + tx := makeRelayAirtimeTx(i+1, PayloadADVERT, 120, 0, fmt.Sprintf("retained-%d", i)) + tx.AdvertRouteEvidence = uint8(i%3 + 1) + tx.FirstSeen = time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC).Add(time.Duration(i) * (7 * 24 * time.Hour / time.Duration(size))).Format(time.RFC3339) + packets[i] = tx + } + store := newRelayAirtimeShareTestStore(packets) + for _, hours := range []int{1, 24, 168} { + window := TimeWindow{Since: time.Date(2026, 1, 8, 0, 0, 0, 0, time.UTC).Add(-time.Duration(hours) * time.Hour).Format(time.RFC3339), Until: "2026-01-08T00:00:00Z"} + if got := store.computeRelayAirtimeShare(window)["total_count"]; got != size*hours/168 { + b.Fatalf("fixture window %dh selected %v, want%d", hours, got, size*hours/168) + } + b.Run(fmt.Sprintf("packets_%d/window_%dh", size, hours), func(b *testing.B) { + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + store.computeRelayAirtimeShare(window) + } + }) + } + } +} + +func TestRelayAirtimeShareWindowDoesNotMutateStoreSlice(t *testing.T) { + first := makeRelayAirtimeTx(1, PayloadACK, 10, 0, "first") + excluded := makeRelayAirtimeTx(2, PayloadACK, 10, 0, "excluded") + excluded.FirstSeen = "2025-12-01T00:00:00Z" + last := makeRelayAirtimeTx(3, PayloadACK, 10, 0, "last") + store := newRelayAirtimeShareTestStore([]*StoreTx{first, excluded, last}) + store.packets = []*StoreTx{first, nil, excluded, last} + if got := store.computeRelayAirtimeShare(TimeWindow{Since: "2026-01-01T00:00:00Z"})["total_count"]; got != 2 { + t.Fatalf("selected count=%v want2", got) + } + if len(store.packets) != 4 || store.packets[0] != first || store.packets[1] != nil || store.packets[2] != excluded || store.packets[3] != last { + t.Fatal("window filtering mutated retained store") + } +} diff --git a/cmd/server/relay_airtime_share.go b/cmd/server/relay_airtime_share.go index 9eb85d9bd..9e13a12d7 100644 --- a/cmd/server/relay_airtime_share.go +++ b/cmd/server/relay_airtime_share.go @@ -6,6 +6,7 @@ import ( "time" "github.com/meshcore-analyzer/lora" + "github.com/meshcore-analyzer/packetpath" ) // relay_airtime_share.go — issues #1359 + #1768 @@ -25,7 +26,7 @@ import ( // preset 869.6 MHz / BW 62.5 kHz / SF 8 / CR 4/5, with the // SF-dependent preamble pulled from internal/lora.PreambleForSF. // -// Aggregated by payload_type. Originator TX is deliberately excluded — a +// Aggregated by payload_type and ADVERT routing kind. Originator TX is excluded — a // never-relayed direct message scores 0, which is the correct framing for a // "relay amplification" metric. In-memory only; no SQL, no new index. @@ -121,7 +122,8 @@ func (s *PacketStore) distinctRelayCount(tx *StoreTx) int { return len(s.resolvedPubkeyReverse[tx.ID]) } -// computeRelayAirtimeShare aggregates relay-airtime-share per payload_type. +// computeRelayAirtimeShare aggregates by payload_type, splitting only ADVERT +// by recorded routing. Other payload-mix analytics retain their usual grouping. // // Returns: // @@ -143,32 +145,69 @@ func (s *PacketStore) computeRelayAirtimeShare(window TimeWindow) map[string]int count int score int64 // sum of ToA(payload) × relays, in nanoseconds } - buckets := make(map[int]*bucket) - seenHash := make(map[string]bool, len(s.packets)) - totalCount := 0 - var totalScore int64 - - for _, tx := range s.packets { - if tx == nil || tx.PayloadType == nil { + // The additional key has four possible values, all for ADVERT. + type bucketKey struct { + payloadType int + advertKind string + } + buckets := make(map[bucketKey]*bucket) + // Filter before allocating the hash map. Reuse the packet slice when every + // record is eligible; otherwise copy only the selected pointers. This keeps + // short-window scratch space small without a full-window allocation penalty. + selected := s.packets + filtered := false + for i, tx := range s.packets { + if tx == nil || tx.PayloadType == nil || !window.Includes(tx.FirstSeen) { + if !filtered { + selected = append([]*StoreTx(nil), s.packets[:i]...) + filtered = true + } continue } - if !window.Includes(tx.FirstSeen) { - continue + if filtered { + selected = append(selected, tx) + } + } + // Low bits union known evidence; the high bit deduplicates scores below. + seenHash := make(map[string]uint8, len(selected)) + for _, tx := range selected { + if *tx.PayloadType == PayloadADVERT && tx.Hash != "" { + seenHash[tx.Hash] |= tx.AdvertRouteEvidence + } + } + if filtered { + // An eligible hash keeps its known older evidence, but neither hashes + // found only outside the window nor their scores enter the result. + for _, tx := range s.packets { + if tx != nil && tx.PayloadType != nil && *tx.PayloadType == PayloadADVERT && tx.Hash != "" { + if mask, eligible := seenHash[tx.Hash]; eligible { + seenHash[tx.Hash] = mask | tx.AdvertRouteEvidence + } + } } - // Dedup per-hash: each distinct packet counted once. ACKs in the - // test fixture have unique hashes so this only collapses true - // re-observations of the same packet. + } + totalCount := 0 + var totalScore int64 + for _, tx := range selected { if tx.Hash != "" { - if seenHash[tx.Hash] { + if seenHash[tx.Hash]&128 != 0 { continue } - seenHash[tx.Hash] = true + seenHash[tx.Hash] |= 128 } pt := *tx.PayloadType - b := buckets[pt] + key := bucketKey{payloadType: pt} + if pt == PayloadADVERT { + mask := tx.AdvertRouteEvidence + if tx.Hash != "" { + mask = seenHash[tx.Hash] & 3 + } + key.advertKind = packetpath.AdvertKind(mask) + } + b := buckets[key] if b == nil { b = &bucket{} - buckets[pt] = b + buckets[key] = b } b.count++ totalCount++ @@ -185,7 +224,8 @@ func (s *PacketStore) computeRelayAirtimeShare(window TimeWindow) map[string]int } rows := make([]map[string]interface{}, 0, len(buckets)) - for pt, b := range buckets { + for key, b := range buckets { + pt := key.payloadType name := ptNames[pt] if name == "" { name = "UNK" @@ -197,17 +237,21 @@ func (s *PacketStore) computeRelayAirtimeShare(window TimeWindow) map[string]int if totalScore > 0 { airtimePct = float64(b.score) / float64(totalScore) * 100.0 } - rows = append(rows, map[string]interface{}{ + row := map[string]interface{}{ "payload_type": name, "type": pt, "count": b.count, "count_pct": countPct, "score": b.score, "airtime_pct": airtimePct, - }) + } + if key.advertKind != "" { + row["advert_kind"] = key.advertKind + } + rows = append(rows, row) } - // Sort descending by airtime_pct; tiebreak count desc, then name asc + // Sort descending by airtime_pct; tiebreak count desc, then name/kind asc // for deterministic ordering. sort.SliceStable(rows, func(i, j int) bool { ai, _ := rows[i]["airtime_pct"].(float64) @@ -222,7 +266,12 @@ func (s *PacketStore) computeRelayAirtimeShare(window TimeWindow) map[string]int } ni, _ := rows[i]["payload_type"].(string) nj, _ := rows[j]["payload_type"].(string) - return ni < nj + if ni != nj { + return ni < nj + } + ki, _ := rows[i]["advert_kind"].(string) + kj, _ := rows[j]["advert_kind"].(string) + return ki < kj }) label := "" @@ -262,12 +311,15 @@ func (s *PacketStore) GetRelayAirtimeShareWithWindow(window TimeWindow) map[stri return out } s.cacheMisses++ + revision := s.advertEvidenceRevision s.cacheMu.Unlock() result := s.computeRelayAirtimeShare(window) s.cacheMu.Lock() - s.rfCache[cacheKey] = &cachedResult{data: result, expiresAt: time.Now().Add(s.rfCacheTTL)} + if revision == s.advertEvidenceRevision { + s.rfCache[cacheKey] = &cachedResult{data: result, expiresAt: time.Now().Add(s.rfCacheTTL)} + } s.cacheMu.Unlock() return result diff --git a/cmd/server/routes.go b/cmd/server/routes.go index b9dbcdc00..76e6fd24b 100644 --- a/cmd/server/routes.go +++ b/cmd/server/routes.go @@ -1656,10 +1656,9 @@ func (s *Server) handleNodeDetail(w http.ResponseWriter, r *http.Request) { // attribution is strict exact-match on the indexed from_pubkey column. recentAdverts, _ := s.db.GetRecentTransmissionsForNode(pubkey, 20) - // Windowed flood-advert count (7d): only the mesh-wide-airtime advert kind, - // separated from zero-hop adverts so a nearby observer hearing a node's - // cheap local adverts does not inflate the number. Consumed by the ArcScope - // repeater advisor to rate advert hygiene. + // Windowed known flood-advert count (7d), including mixed evidence. + // Like recentAdverts, this is a lower bound on available route history, + // not a classification from the canonical first-ingested route. if n, err := s.db.CountFloodAdvertsForNode(pubkey, 7*24, floodAdvertRowCap); err == nil { node["flood_advert_count_7d"] = n } else { diff --git a/cmd/server/store.go b/cmd/server/store.go index df94d8c61..cc0fde835 100644 --- a/cmd/server/store.go +++ b/cmd/server/store.go @@ -58,10 +58,11 @@ type StoreTx struct { LatestSeen string // max observation timestamp (or FirstSeen if no observations) UniqueObserverCount int // cached count of distinct observer IDs // Cached parsed fields (set once, read many) - parsedPath []string // cached parsePathJSON result - pathParsed bool // whether parsedPath has been set - decodedOnce sync.Once // guards parsedDecoded - parsedDecoded map[string]interface{} // cached json.Unmarshal of DecodedJSON + parsedPath []string // cached parsePathJSON result + pathParsed bool // whether parsedPath has been set + AdvertRouteEvidence uint8 // union of known flood/direct-empty-path evidence + decodedOnce sync.Once // guards parsedDecoded + parsedDecoded map[string]interface{} // cached json.Unmarshal of DecodedJSON // Dedup map: "observerID|pathJSON" → true for O(1) duplicate checks obsKeys map[string]bool observerSet map[string]bool // unique observer IDs (for UniqueObserverCount) @@ -172,6 +173,12 @@ func (tx *StoreTx) ParsedDecoded() map[string]interface{} { // All other locks are acquired independently (no nesting). // When adding new lock acquisitions, respect this ordering. type PacketStore struct { + // Lock order: advertEvidenceMu -> mu -> cacheMu. The feed and mask + // loading share this gate, preventing a late event/chunk-merge race. + advertEvidenceMu sync.Mutex + advertEvidenceCursor int64 + advertEvidenceRevision uint64 // cacheMu + mu sync.RWMutex db *DB packets []*StoreTx // sorted by first_seen ASC (oldest first; newest at tail) @@ -731,6 +738,11 @@ func NewPacketStore(db *DB, cfg *PacketStoreConfig, cacheTTLs ...map[string]inte ps.invCooldown = v } } + // Capture BEFORE loading any packets. Evidence arriving during startup + // remains above this watermark; each loaded tx also reads its full mask. + if db.advertEvidencePresent() { + _ = db.conn.QueryRow(`SELECT COALESCE(MAX(id),0) FROM advert_route_evidence`).Scan(&ps.advertEvidenceCursor) + } return ps } @@ -738,6 +750,8 @@ func NewPacketStore(db *DB, cfg *PacketStoreConfig, cacheTTLs ...map[string]inte // When maxMemoryMB > 0, loads only the newest N transmissions that fit // within the memory budget, avoiding OOM on large databases. func (s *PacketStore) Load() error { + s.advertEvidenceMu.Lock() + defer s.advertEvidenceMu.Unlock() s.mu.Lock() defer s.mu.Unlock() @@ -1007,6 +1021,21 @@ func (s *PacketStore) Load() error { } } + if err := rows.Err(); err != nil { + return err + } + rows.Close() + ids := advertTxIDs(s.packets) + // Keep the evidence gate across read/merge, but do not add SQL work + // under the global store lock. Eviction during the read is harmless. + s.mu.Unlock() + masks, err := s.db.advertEvidenceForIDs(ids) + s.mu.Lock() + if err != nil { + return err + } + s.mergeAdvertEvidence(masks) + // Post-load: pick best observation (longest path) for each transmission, // then re-index so relay hops from resolved_path land in byNode. // indexByNode was called earlier (on StoreTx creation) before observations @@ -1340,6 +1369,17 @@ func (s *PacketStore) loadChunk(from, to time.Time) error { if len(localPackets) == 0 { return nil } + rows.Close() + s.advertEvidenceMu.Lock() + defer s.advertEvidenceMu.Unlock() + ids := advertTxIDs(localPackets) + masks, err := s.db.advertEvidenceForIDs(ids) + if err != nil { + return err + } + for _, tx := range localPackets { + tx.AdvertRouteEvidence = masks[tx.ID] + } // PR #1187 r3 MUST-FIX 1: index↔slice consistency. // @@ -1398,11 +1438,22 @@ func (s *PacketStore) loadChunk(from, to time.Time) error { newObsIDs[k] = true } } + evidenceChanged := false for k, v := range batchHashes { if s.byHash[k] == nil { s.byHash[k] = v + } else { + // A live-loaded copy may already carry newer bits. Union in + // both directions before mergeChunkIntoPackets chooses it. + if unionAdvertEvidence(s.byHash[k], v.AdvertRouteEvidence) { + evidenceChanged = true + } + v.AdvertRouteEvidence |= s.byHash[k].AdvertRouteEvidence } } + if evidenceChanged { + s.invalidateAdvertEvidence() + } for k, v := range batchTxIDs { if s.byTxID[k] == nil { s.byTxID[k] = v @@ -1452,6 +1503,7 @@ func (s *PacketStore) loadChunk(from, to time.Time) error { // before it readers see the old slice (which is still fully indexed). s.mu.Lock() s.packets = mergeChunkIntoPackets(localPackets, s.packets) + s.invalidateAdvertEvidence() s.totalObs += localTotalObs s.trackedBytes += localTrackedBytes if localMaxTxID > s.maxTxID { @@ -2778,12 +2830,30 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac if len(tempRows) == 0 { return nil, sinceID } + if err := rows.Err(); err != nil { + return nil, sinceID + } + rows.Close() + s.advertEvidenceMu.Lock() + defer s.advertEvidenceMu.Unlock() + ids := make([]int, 0, txCount) + for _, row := range tempRows { + if row.payloadType != nil && *row.payloadType == PayloadADVERT && (len(ids) == 0 || ids[len(ids)-1] != row.txID) { + ids = append(ids, row.txID) + } + } + masks, err := s.db.advertEvidenceForIDs(ids) + if err != nil { + log.Printf("[store] load advert evidence: %v", err) + return nil, sinceID + } // Now lock and merge into store s.mu.Lock() defer s.mu.Unlock() newMaxID := sinceID + evidenceChanged := false broadcastTxs := make(map[int]*StoreTx) // track new transmissions for broadcast hasNewNodes := false // track genuinely new node pubkeys var broadcastOrder []int @@ -2844,6 +2914,9 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac } } + if unionAdvertEvidence(tx, masks[r.txID]) { + evidenceChanged = true + } if r.obsID != nil { oid := *r.obsID // Dedup (O(1) map lookup) @@ -3052,6 +3125,9 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac // sees every observation) owns that write too. _ = broadcastRP // resolved path is still computed in-memory (above) for live broadcast; no SQL write. + if evidenceChanged { + s.invalidateAdvertEvidence() + } return result, newMaxID } @@ -3059,6 +3135,8 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac // store. This catches observations that arrive after IngestNewFromDB has already // advanced past the transmission's ID (fixes #174). func (s *PacketStore) IngestNewObservations(sinceObsID, limit int) []map[string]interface{} { + // Must run even when observation IDs/timestamps have not changed. + s.refreshAdvertEvidence() if limit <= 0 { limit = 500 } diff --git a/internal/dbschema/advert_evidence.go b/internal/dbschema/advert_evidence.go new file mode 100644 index 000000000..dc6e2bacc --- /dev/null +++ b/internal/dbschema/advert_evidence.go @@ -0,0 +1,24 @@ +package dbschema + +import "database/sql" + +// Empty on creation; history is backfilled after the ingestor becomes ready. +// At most two rows per retained transmission. AUTOINCREMENT prevents reuse of +// feed IDs after retention removes the newest rows. +// PREFLIGHT: async=true reason="creates empty bounded evidence and single-row progress tables; history backfill is batched after live ingest readiness" +func ensureAdvertEvidence(rw *sql.DB) error { + _, err := rw.Exec(` + CREATE TABLE IF NOT EXISTS advert_route_evidence ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + tx_id INTEGER NOT NULL REFERENCES transmissions(id) ON DELETE CASCADE, + bit INTEGER NOT NULL CHECK(bit IN (1,2)), + UNIQUE(tx_id,bit) + ); + CREATE TABLE IF NOT EXISTS advert_evidence_backfill ( + id INTEGER PRIMARY KEY CHECK(id=1), + tx_cursor INTEGER NOT NULL DEFAULT 0, + obs_cursor INTEGER NOT NULL DEFAULT 0 + ); + `) + return err +} diff --git a/internal/dbschema/dbschema.go b/internal/dbschema/dbschema.go index 433ee8de0..ba0533ae7 100644 --- a/internal/dbschema/dbschema.go +++ b/internal/dbschema/dbschema.go @@ -108,6 +108,12 @@ func Apply(rw *sql.DB, logf Logger) error { if err := ensureTransmissionsLastSeenColumn(rw, logf); err != nil { return fmt.Errorf("ensure transmissions.last_seen: %w", err) } + if err := ensureAdvertEvidence(rw); err != nil { + return fmt.Errorf("ensure advert route evidence: %w", err) + } + // Advert evidence is intentionally not required by AssertReady: a legacy + // read-only server reports unknown until this table/backfill is available, + // and re-probes absence instead of caching the startup race permanently. return nil } diff --git a/internal/packetpath/advert.go b/internal/packetpath/advert.go new file mode 100644 index 000000000..3e6b581bc --- /dev/null +++ b/internal/packetpath/advert.go @@ -0,0 +1,65 @@ +package packetpath + +import "encoding/hex" + +const ( + AdvertFlood uint8 = 1 + AdvertDirectEmptyPath uint8 = 2 +) + +// AdvertRouteEvidence classifies one received wire frame, independently of a +// canonical transmission's metadata. Mesh::sendZeroHop writes path_len = 0; +// Mesh::onRecvPacket -> removeSelfFromPath -> retransmit can also leave an +// empty remaining path, preserving the hash-size flags. Neither encoding +// proves an original zero-hop send or RF distance. Transport +// routes carry four bytes before path_len (firmware Packet.cpp / Mesh.cpp). +// The fixed buffer bounds both work and allocation to a radio-sized frame. +func AdvertRouteEvidence(raw string) uint8 { + var frame [256]byte + if len(raw)%2 != 0 || len(raw) > len(frame)*2 { + return 0 + } + n, err := hex.Decode(frame[:], []byte(raw)) + if err != nil || n < 3 || (frame[0]>>2)&15 != 4 { + return 0 + } + route := int(frame[0] & 3) + pathOffset := 1 + if IsTransportRoute(route) { + pathOffset += 4 + } + if n <= pathOffset+1 { + return 0 + } + path := frame[pathOffset] + hashSize := int(path>>6) + 1 + pathBytes := int(path&63) * hashSize + payloadBytes := n - pathOffset - 1 - pathBytes + if hashSize == 4 || pathBytes > 64 || payloadBytes < 1 || payloadBytes > 184 { + return 0 + } + if route == RouteFlood || route == RouteTransportFlood { + return AdvertFlood + } + if path&63 == 0 { + return AdvertDirectEmptyPath + } + return 0 +} + +// AdvertKind describes the union of known evidence, not exclusive historical +// use: legacy observations may already have been overwritten before upgrade, +// or evidence writes may have failed. The API spelling "zero_hop" is retained +// for compatibility and denotes observed direct empty-path frames only. +func AdvertKind(evidence uint8) string { + switch evidence & (AdvertFlood | AdvertDirectEmptyPath) { + case AdvertFlood: + return "flood" + case AdvertDirectEmptyPath: + return "zero_hop" + case AdvertFlood | AdvertDirectEmptyPath: + return "mixed" + default: + return "other" + } +} diff --git a/internal/packetpath/advert_test.go b/internal/packetpath/advert_test.go new file mode 100644 index 000000000..8fca6f2f0 --- /dev/null +++ b/internal/packetpath/advert_test.go @@ -0,0 +1,40 @@ +package packetpath + +import ( + "strings" + "testing" +) + +func TestAdvertRouteEvidence(t *testing.T) { + for _, tc := range []struct { + raw string + want uint8 + }{ + {"1100aa", 1}, {"1101ffaa", 1}, {"1140aa", 1}, {"100102030400aa", 1}, + {"1200aa", 2}, {"130102030400aa", 2}, + {"1240aa", 2}, {"1280aa", 2}, {"130102030440aa", 2}, {"130102030480aa", 2}, {"12c0aa", 0}, {"1201ffaa", 0}, + {"130102030401ffaa", 0}, {"13zz00000000aa", 0}, {"16ff", 0}, + {"1100zz", 0}, {"1101", 0}, {"1101ff", 0}, {"11c0aa", 0}, + {"1100a", 0}, {"1100", 0}, {"1001020304", 0}, {"100102030400", 0}, + {"", 0}, {"1200" + strings.Repeat("aa", 255), 0}, + {"1100" + strings.Repeat("aa", 185), 0}, + {"1100" + strings.Repeat("aa", 184), 1}, + } { + if got := AdvertRouteEvidence(tc.raw); got != tc.want { + t.Errorf("%s: got %d want %d", tc.raw, got, tc.want) + } + } + for mask, want := range []string{"other", "flood", "zero_hop", "mixed"} { + if got := AdvertKind(uint8(mask)); got != want { + t.Errorf("mask %d: %s want %s", mask, got, want) + } + } +} + +func BenchmarkAdvertRouteEvidence(b *testing.B) { + raw := "1100" + strings.Repeat("ab", 118) + b.ReportAllocs() + for i := 0; i < b.N; i++ { + AdvertRouteEvidence(raw) + } +} diff --git a/public/analytics.js b/public/analytics.js index 62c135653..c39d18189 100644 --- a/public/analytics.js +++ b/public/analytics.js @@ -509,7 +509,7 @@ } var totalScore = (data && typeof data.total_score === 'number') ? data.total_score : 0; if (totalScore <= 0) { - return '
No relay activity observed in this window (all packets direct).
'; + return '
No relay activity observed in this window.
'; } // Issue #1768 — surface the LoRa preset baked into the ToA score. Share // numbers are only meaningful relative to one PHY preset; operators must @@ -538,8 +538,17 @@ var palette = ['#ef4444','#f59e0b','#22c55e','#3b82f6','#8b5cf6','#ec4899','#14b8a6','#64748b','#f97316','#06b6d4','#84cc16']; var html = '
'; if (presetCaption) html += presetCaption; + if (rows.some(function (r) { return r.payload_type === 'ADVERT'; })) { + html += '
Known route evidence: Mixed adverts were observed as both flood and direct with an empty remaining path. An empty path does not prove an original zero-hop send. Other adverts have no classified evidence yet. Background backfill improves classification automatically. Available history is a lower bound; older overwritten observations cannot be recovered, and failed evidence writes may leave gaps.
'; + } rows.forEach(function (r, i) { var name = r.payload_type || 'UNK'; + if (name === 'ADVERT') { + if (r.advert_kind === 'flood') name = 'Flood adverts'; + else if (r.advert_kind === 'zero_hop') name = 'Direct adverts (empty path)'; + else if (r.advert_kind === 'mixed') name = 'Mixed adverts'; + else if (r.advert_kind === 'other') name = 'Other adverts'; + } var cnt = Number(r.count || 0); var cpct = Number(r.count_pct || 0); var apct = Number(r.airtime_pct || 0); @@ -560,7 +569,7 @@ 'Count: ' + cnt.toLocaleString() + ' (' + cpct.toFixed(2) + '%)\n' + 'Airtime: ' + apct.toFixed(2) + '% (score ' + scoreStr + ' · airtime × repeaters)\n' + 'Score = LoRa Time-on-Air × distinct repeaters. Within-mesh only.'; - html += '
' + + html += '
' + '
' + esc(name) + '
' + '
' + '
' + @@ -575,7 +584,7 @@ '
'; }); // Axis legend. - html += '
' + + html += '
' + '
' + '
0%50%100%
' + '
' + diff --git a/public/style.css b/public/style.css index 7b95eda7b..1d24a9705 100644 --- a/public/style.css +++ b/public/style.css @@ -3438,6 +3438,11 @@ tr[data-hops]:hover { background: rgba(59,130,246,0.1); } } .analytics-chart-card { background: var(--card-bg); border: 1px solid var(--border); border-radius: 6px; padding: var(--space-sm); min-width: 0; } .analytics-chart-card.full { grid-column: 1 / -1; } +/* Keep relay labels, percentage values and a visible track inside narrow cards. */ +.dumbbell-chart { container-type: inline-size; } +@container (max-width: 480px) { + .dumbbell-row, .dumbbell-axis { --relay-airtime-columns: 70px minmax(20px, 1fr) 90px; } +} /* Constrain chart media inside the card (svg/canvas at any depth). The `.analytics-chart-card svg, .analytics-chart-card canvas` descendant selector is robust to wrapper elements (legends, tooltips, axis diff --git a/tests/e2e/test-e2e-playwright.js b/tests/e2e/test-e2e-playwright.js index 1bdf3e1d6..d38f7a0b2 100644 --- a/tests/e2e/test-e2e-playwright.js +++ b/tests/e2e/test-e2e-playwright.js @@ -809,6 +809,68 @@ async function run() { }); // Test 8b (#842): time-window picker triggers requests with ?window=… param. + // #2041: exercise the overview's actual API-to-renderer path at both sizes. + await test('Relay airtime chart splits adverts without hiding zero-relay rows', async () => { + const chartPage = await context.newPage(); + try { + await chartPage.route('**/api/analytics/relay-airtime-share*', route => route.fulfill({ + json: { total_count: 5, total_score: 300, rows: [ + { payload_type: 'ADVERT', type: 4, advert_kind: 'flood', count: 1, count_pct: 20, score: 100, airtime_pct: 33.333 }, + { payload_type: 'ADVERT', type: 4, advert_kind: 'other', count: 1, count_pct: 20, score: 100, airtime_pct: 33.333 }, + { payload_type: 'ADVERT', type: 4, advert_kind: 'zero_hop', count: 1, count_pct: 20, score: 0, airtime_pct: 0 }, + { payload_type: 'ADVERT', type: 4, advert_kind: 'mixed', count: 1, count_pct: 20, score: 100, airtime_pct: 33.333 }, + { payload_type: 'ACK', type: 3, count: 1, count_pct: 20, score: 0, airtime_pct: 0 }, + ] }, + })); + for (const width of [1280, 320]) { + await chartPage.setViewportSize({ width, height: 900 }); + await chartPage.goto(BASE + '/#/analytics'); + await chartPage.reload({ waitUntil: 'domcontentloaded' }); + await chartPage.waitForSelector('.dumbbell-row'); + const labels = await chartPage.locator('.dumbbell-label').allTextContents(); + assert(JSON.stringify(labels) === JSON.stringify(['Flood adverts', 'Other adverts', 'Direct adverts (empty path)', 'Mixed adverts', 'ACK']), 'advert chart labels: ' + labels.join(', ')); + assert((await chartPage.locator('.dumbbell-evidence-note').textContent()).includes('older overwritten observations cannot be recovered'), 'known-evidence caveat is visible'); + const zero = chartPage.locator('.dumbbell-row').filter({ hasText: 'Direct adverts (empty path)' }); + assert((await zero.textContent()).includes('air 0.0%'), 'zero-relay advert row must remain visible'); + assert((await zero.getAttribute('title')).includes('Count: 1 (20.00%)'), 'tooltip retains count'); + const layout = await chartPage.locator('.dumbbell-chart').evaluate(chart => { + const box = chart.getBoundingClientRect(); + const axisLabels = [...chart.querySelectorAll('.dumbbell-axis span')]; + const expectedAxis = ['0%', '50%', '100%']; + return { + axisFits: axisLabels.length === expectedAxis.length && axisLabels.every((label, i) => { + const bounds = label.getBoundingClientRect(); + return label.textContent.trim() === expectedAxis[i] && bounds.width > 0 && bounds.height > 0 && + (i === 0 || axisLabels[i - 1].getBoundingClientRect().right <= bounds.left); + }), + overflow: chart.scrollWidth > chart.clientWidth + 1, + outside: box.left < -1 || box.right > window.innerWidth + 1, + rowsFit: [...chart.querySelectorAll('.dumbbell-row')].every(row => { + const label = row.querySelector('.dumbbell-label').getBoundingClientRect(); + const track = row.querySelector('.dumbbell-track').getBoundingClientRect(); + const values = row.querySelector('.dumbbell-values').getBoundingClientRect(); + return label.right <= track.left && track.width >= 20 && track.right <= values.left && values.right <= box.right + 1; + }), + }; + }); + assert(!layout.overflow && !layout.outside && layout.rowsFit && layout.axisFits, `relay chart layout at ${width}px: ${JSON.stringify(layout)}`); + } + } finally { await chartPage.close(); } + }); + + await test('Relay airtime zero-activity state does not infer direct routing', async () => { + const chartPage = await context.newPage(); + try { + await chartPage.route('**/api/analytics/relay-airtime-share*', route => route.fulfill({ + json: { total_count: 1, total_score: 0, rows: [ + { payload_type: 'ADVERT', type: 4, advert_kind: 'flood', count: 1, count_pct: 100, score: 0, airtime_pct: 0 }, + ] }, + })); + await chartPage.goto(BASE + '/#/analytics'); + await chartPage.waitForFunction(() => document.body.textContent.includes('No relay activity observed')); + assert(!(await chartPage.locator('body').textContent()).includes('all packets direct'), 'no resolved relays does not imply direct packets'); + } finally { await chartPage.close(); } + }); await test('Analytics time-window picker refetches with window param', async () => { // Picker must be rendered. await page.waitForSelector('#analyticsTimeWindow', { timeout: 5000 }); diff --git a/tests/unit/test-frontend-helpers.js b/tests/unit/test-frontend-helpers.js index ad54cb4f2..ed747ed82 100644 --- a/tests/unit/test-frontend-helpers.js +++ b/tests/unit/test-frontend-helpers.js @@ -2071,6 +2071,51 @@ console.log('\n=== app.js: isTransportRoute + transportBadge ==='); test('transportBadge(1) returns empty string', () => assert.strictEqual(transportBadge(1), '')); } +// #2041: execute the actual IIFE renderer with a test-only export. +console.log('\n=== Relay airtime advert labels ==='); +{ + const ctx = makeSandbox(); + ctx.registerPage = () => {}; + ctx.getComputedStyle = () => ({ getPropertyValue: () => '' }); + ctx.esc = s => String(s).replace(/&/g, '&').replace(/ { + const html = render({ rows, total_score: 300 }); + for (const label of ['Flood adverts', 'Direct adverts (empty path)', 'Other adverts', 'Mixed adverts', 'ACK']) { + assert.ok(html.includes('>' + label + '
'), 'missing chart label: ' + label); + assert.ok(html.includes('title="' + label + '\n'), 'missing tooltip label: ' + label); + } + assert.strictEqual((html.match(/class="dumbbell-row"/g) || []).length, 5); + assert.ok(html.includes('Known route evidence'), 'classification must not imply complete history'); + assert.ok(html.includes('Background backfill improves classification automatically'), 'upgrade progress is explained'); + assert.ok(html.includes('does not prove an original zero-hop send'), 'direct empty paths do not prove origin'); + assert.ok(html.includes('failed evidence writes may leave gaps'), 'best-effort evidence limitation is explained'); + assert.ok(html.includes('older overwritten observations cannot be recovered'), 'legacy caveat must remain visible'); + assert.ok(html.includes('air 0.0%'), 'zero-relay advert share remains visible'); + }); + test('relay chart supports legacy payload labels and preserves row input', () => { + const legacyRow = { payload_type: 'ADVERT', type: 4, score: 100 }; + const before = JSON.stringify(legacyRow); + const html = render({ rows: [legacyRow], total_score: 100 }); + assert.ok(html.includes('>ADVERT
')); + assert.strictEqual(JSON.stringify(legacyRow), before); + }); + test('no relay evidence does not imply all packets were direct', () => { + const html = render({ rows: [rows[0]], total_score: 0 }); + assert.ok(html.includes('No relay activity observed')); + assert.ok(!html.includes('all packets direct'), 'flood adverts can have no resolved relays'); + assert.ok(render({ rows: [] }).includes('No relay-airtime data')); + }); +} // ===== ANALYTICS.JS: Channel Sort ===== console.log('\n=== analytics.js: sortChannels ==='); {