From 58484ad924b03bd41653aaf81bbbcfff238cdb9e Mon Sep 17 00:00:00 2001 From: Kpa-clawbot Date: Sat, 2 May 2026 19:52:43 -0700 Subject: [PATCH] feat(ingestor): backfill observations.path_json from raw_hex (closes #888) (#983) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Summary Adds an idempotent startup migration to the ingestor that backfills `observations.path_json` from per-observation `raw_hex` (added in #882). **Approach: Server-side migration (Option B)** — runs automatically at startup, chunked in batches of 1000, tracked via `_migrations` table. Chosen over a standalone script because: 1. Follows existing migration pattern (channel_hash, last_packet_at, etc.) 2. Zero operator action required — just deploy 3. Idempotent — safe to restart mid-migration (uncommitted rows get picked up next run) ## What it does - Selects observations where `raw_hex` is populated but `path_json` is NULL/empty/`[]` - Excludes TRACE packets (`payload_type = 9`) at the SQL level — their header bytes are SNR values, not hops - Decodes hops via `packetpath.DecodePathFromRawHex` (reuses existing helper) - Updates `path_json` with the decoded JSON array - Marks rows with undecoded/empty hops as `'[]'` to prevent infinite re-scanning - Records `backfill_path_json_from_raw_hex_v1` in `_migrations` when complete ## Safety - **Never overwrites** existing non-empty `path_json` — only fills where missing - **Batched** (1000 rows per iteration) — won't OOM on large DBs - **TRACE-safe** — excluded at query level per `packetpath.PathBytesAreHops` semantics ## Test `TestBackfillPathJsonFromRawHex` — creates synthetic observations with: - Empty path_json + valid raw_hex → verifies backfill populates correctly - NULL path_json → verifies backfill populates - Existing path_json → verifies NO overwrite - TRACE packet → verifies skip Anti-tautology: test asserts specific decoded values (`["AABB","CCDD"]`) from known raw_hex input, not just "something changed." Closes #888 Co-authored-by: you --- cmd/ingestor/db.go | 50 ++++++++++++++++++++++++++++ cmd/ingestor/db_test.go | 74 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 124 insertions(+) diff --git a/cmd/ingestor/db.go b/cmd/ingestor/db.go index 6fc0f9a1..f9ba558d 100644 --- a/cmd/ingestor/db.go +++ b/cmd/ingestor/db.go @@ -444,6 +444,56 @@ func applySchema(db *sql.DB) error { log.Println("[migration] observers.last_packet_at column added") } + // Migration: backfill observations.path_json from raw_hex (#888) + row = db.QueryRow("SELECT 1 FROM _migrations WHERE name = 'backfill_path_json_from_raw_hex_v1'") + if row.Scan(&migDone) != nil { + log.Println("[migration] Backfilling observations.path_json from raw_hex...") + updated := 0 + const batchSize = 1000 + for { + rows, err := db.Query(` + SELECT o.id, o.raw_hex + FROM observations o + JOIN transmissions t ON o.transmission_id = t.id + WHERE o.raw_hex IS NOT NULL AND o.raw_hex != '' + AND (o.path_json IS NULL OR o.path_json = '' OR o.path_json = '[]') + AND t.payload_type != 9 + LIMIT ?`, batchSize) + if err != nil { + log.Printf("[migration] backfill_path_json query error: %v", err) + break + } + type pendingRow struct { + id int64 + rawHex string + } + var batch []pendingRow + for rows.Next() { + var r pendingRow + if err := rows.Scan(&r.id, &r.rawHex); err == nil { + batch = append(batch, r) + } + } + rows.Close() + if len(batch) == 0 { + break + } + for _, r := range batch { + hops, err := packetpath.DecodePathFromRawHex(r.rawHex) + if err != nil || len(hops) == 0 { + // Mark as processed with empty path to avoid re-scanning + db.Exec(`UPDATE observations SET path_json = '[]' WHERE id = ?`, r.id) + continue + } + b, _ := json.Marshal(hops) + db.Exec(`UPDATE observations SET path_json = ? WHERE id = ?`, string(b), r.id) + updated++ + } + } + log.Printf("[migration] Backfilled path_json for %d observations from raw_hex", updated) + db.Exec(`INSERT INTO _migrations (name) VALUES ('backfill_path_json_from_raw_hex_v1')`) + } + return nil } diff --git a/cmd/ingestor/db_test.go b/cmd/ingestor/db_test.go index 957b7bdb..86d2a52a 100644 --- a/cmd/ingestor/db_test.go +++ b/cmd/ingestor/db_test.go @@ -2178,3 +2178,77 @@ func TestBuildPacketData_NonTracePathJSON(t *testing.T) { t.Errorf("path_json = %s, want %s", pd.PathJSON, expectedPathJSON) } } + +// --- Issue #888: Backfill path_json from raw_hex --- + +func TestBackfillPathJsonFromRawHex(t *testing.T) { + dbPath := tempDBPath(t) + s, err := OpenStore(dbPath) + if err != nil { + t.Fatal(err) + } + + // Insert a transmission with payload_type != TRACE (e.g. 0x01) + // raw_hex: header 0x05 (route FLOOD, payload 0x01), path byte 0x42 (hash_size=2, count=2), + // hops: AABB, CCDD, then some payload bytes + rawHex := "0542AABBCCDD0000000000000000000000000000" + s.db.Exec(`INSERT INTO transmissions (raw_hex, hash, first_seen, payload_type) VALUES (?, 'h1', '2025-01-01T00:00:00Z', 1)`, rawHex) + + // Insert observation with raw_hex but empty path_json + s.db.Exec(`INSERT INTO observations (transmission_id, timestamp, raw_hex, path_json) VALUES (1, 1000, ?, '[]')`, rawHex) + // Insert observation with raw_hex and NULL path_json + s.db.Exec(`INSERT INTO observations (transmission_id, timestamp, raw_hex, path_json) VALUES (1, 1001, ?, NULL)`, rawHex) + // Insert observation with existing path_json (should NOT be overwritten) + s.db.Exec(`INSERT INTO observations (transmission_id, timestamp, raw_hex, path_json) VALUES (1, 1002, ?, '["XX","YY"]')`, rawHex) + + // Insert a TRACE transmission (payload_type = 0x09) — should be skipped + traceRaw := "2604302D0D2359FEE7B100000000006733D63367" + s.db.Exec(`INSERT INTO transmissions (raw_hex, hash, first_seen, payload_type) VALUES (?, 'h2', '2025-01-01T00:00:00Z', 9)`, traceRaw) + s.db.Exec(`INSERT INTO observations (transmission_id, timestamp, raw_hex, path_json) VALUES (2, 1003, ?, '[]')`, traceRaw) + + // Remove the migration marker so it runs again on reopen + s.db.Exec(`DELETE FROM _migrations WHERE name = 'backfill_path_json_from_raw_hex_v1'`) + s.Close() + + // Reopen — migration should run + s2, err := OpenStore(dbPath) + if err != nil { + t.Fatal(err) + } + defer s2.Close() + + // Check migration ran + var migCount int + s2.db.QueryRow("SELECT COUNT(*) FROM _migrations WHERE name = 'backfill_path_json_from_raw_hex_v1'").Scan(&migCount) + if migCount != 1 { + t.Fatalf("migration not recorded") + } + + // Row 1 (was '[]') should now have decoded hops + var pj1 string + s2.db.QueryRow("SELECT path_json FROM observations WHERE id = 1").Scan(&pj1) + if pj1 != `["AABB","CCDD"]` { + t.Errorf("row 1 path_json = %q, want %q", pj1, `["AABB","CCDD"]`) + } + + // Row 2 (was NULL) should now have decoded hops + var pj2 string + s2.db.QueryRow("SELECT path_json FROM observations WHERE id = 2").Scan(&pj2) + if pj2 != `["AABB","CCDD"]` { + t.Errorf("row 2 path_json = %q, want %q", pj2, `["AABB","CCDD"]`) + } + + // Row 3 (had existing data) should NOT be overwritten + var pj3 string + s2.db.QueryRow("SELECT path_json FROM observations WHERE id = 3").Scan(&pj3) + if pj3 != `["XX","YY"]` { + t.Errorf("row 3 path_json = %q, want %q (should not be overwritten)", pj3, `["XX","YY"]`) + } + + // Row 4 (TRACE) should NOT be updated + var pj4 string + s2.db.QueryRow("SELECT path_json FROM observations WHERE id = 4").Scan(&pj4) + if pj4 != "[]" { + t.Errorf("row 4 (TRACE) path_json = %q, want %q (should be skipped)", pj4, "[]") + } +}