diff --git a/cmd/ingestor/db.go b/cmd/ingestor/db.go index b110317e..4d03b952 100644 --- a/cmd/ingestor/db.go +++ b/cmd/ingestor/db.go @@ -8,6 +8,7 @@ import ( "os" "path/filepath" "strings" + "sync" "sync/atomic" "time" @@ -44,6 +45,7 @@ type Store struct { stmtUpsertMetrics *sql.Stmt sampleIntervalSec int + backfillWg sync.WaitGroup } // OpenStore opens or creates a SQLite DB at the given path, applying the @@ -445,54 +447,8 @@ func applySchema(db *sql.DB) error { } // 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')`) - } + // NOTE: This runs ASYNC via BackfillPathJSONAsync() to avoid blocking MQTT startup. + // See staging outage where ~502K rows blocked ingest for 15+ hours. // One-time cleanup: delete legacy packets with empty hash or empty first_seen (#994) row = db.QueryRow("SELECT 1 FROM _migrations WHERE name = 'cleanup_legacy_null_hash_ts'") @@ -800,6 +756,7 @@ func (s *Store) UpsertObserver(id, name, iata string, meta *ObserverMeta) error // Close checkpoints the WAL and closes the database. func (s *Store) Close() error { + s.backfillWg.Wait() s.Checkpoint() return s.db.Close() } @@ -939,6 +896,92 @@ func (s *Store) Checkpoint() { } } +// BackfillPathJSONAsync launches the path_json backfill in a background goroutine. +// It processes observations with NULL/empty path_json that have raw_hex available, +// decoding hop paths and updating the column. Safe to run concurrently with ingest +// because new observations get path_json at write time; this only touches NULL rows. +// Idempotent: skips if migration already recorded. +func (s *Store) BackfillPathJSONAsync() { + s.backfillWg.Add(1) + go func() { + defer s.backfillWg.Done() + defer func() { + if r := recover(); r != nil { + log.Printf("[backfill] path_json async panic recovered: %v", r) + } + }() + + var migDone int + row := s.db.QueryRow("SELECT 1 FROM _migrations WHERE name = 'backfill_path_json_from_raw_hex_v1'") + if row.Scan(&migDone) == nil { + return // already done + } + + log.Println("[backfill] Starting async path_json backfill from raw_hex...") + updated := 0 + errored := false + const batchSize = 1000 + batchNum := 0 + for { + rows, err := s.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("[backfill] path_json query error: %v", err) + errored = true + 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 { + if _, execErr := s.db.Exec(`UPDATE observations SET path_json = '[]' WHERE id = ?`, r.id); execErr != nil { + log.Printf("[backfill] write error (id=%d): %v", r.id, execErr) + } + continue + } + b, _ := json.Marshal(hops) + if _, execErr := s.db.Exec(`UPDATE observations SET path_json = ? WHERE id = ?`, string(b), r.id); execErr != nil { + log.Printf("[backfill] write error (id=%d): %v", r.id, execErr) + } else { + updated++ + } + } + batchNum++ + if batchNum%50 == 0 { + log.Printf("[backfill] progress: %d observations updated so far (%d batches)", updated, batchNum) + } + // Throttle: yield to ingest writers between batches + time.Sleep(50 * time.Millisecond) + } + log.Printf("[backfill] Async path_json backfill complete: %d observations updated", updated) + if !errored { + s.db.Exec(`INSERT INTO _migrations (name) VALUES ('backfill_path_json_from_raw_hex_v1')`) + } else { + log.Printf("[backfill] NOT recording migration due to errors — will retry on next restart") + } + }() +} + // LogStats logs current operational metrics. func (s *Store) LogStats() { log.Printf("[stats] tx_inserted=%d tx_dupes=%d obs_inserted=%d node_upserts=%d observer_upserts=%d write_errors=%d sig_drops=%d", diff --git a/cmd/ingestor/db_test.go b/cmd/ingestor/db_test.go index c3f95933..82a5ce15 100644 --- a/cmd/ingestor/db_test.go +++ b/cmd/ingestor/db_test.go @@ -2210,16 +2210,24 @@ func TestBackfillPathJsonFromRawHex(t *testing.T) { s.db.Exec(`DELETE FROM _migrations WHERE name = 'backfill_path_json_from_raw_hex_v1'`) s.Close() - // Reopen — migration should run + // Reopen — backfill is now async, must trigger explicitly s2, err := OpenStore(dbPath) if err != nil { t.Fatal(err) } defer s2.Close() - // Check migration ran + // Trigger async backfill and wait for completion + s2.BackfillPathJSONAsync() + deadline := time.Now().Add(10 * time.Second) var migCount int - s2.db.QueryRow("SELECT COUNT(*) FROM _migrations WHERE name = 'backfill_path_json_from_raw_hex_v1'").Scan(&migCount) + for time.Now().Before(deadline) { + s2.db.QueryRow("SELECT COUNT(*) FROM _migrations WHERE name = 'backfill_path_json_from_raw_hex_v1'").Scan(&migCount) + if migCount == 1 { + break + } + time.Sleep(50 * time.Millisecond) + } if migCount != 1 { t.Fatalf("migration not recorded") } @@ -2376,3 +2384,186 @@ func TestBuildPacketDataRegionFallsBackToTopic(t *testing.T) { t.Fatalf("expected region SJC from topic, got %q", pkt.Region) } } + + +// TestBackfillPathJSONAsync verifies that the path_json backfill does NOT block +// OpenStore from returning. MQTT connect happens immediately after OpenStore; +// if the backfill is synchronous, MQTT would be delayed indefinitely on large DBs. +// This test creates pending backfill rows, opens the store, and asserts that +// OpenStore returns before the migration is recorded — proving async execution. +func TestBackfillPathJSONAsync(t *testing.T) { + dir := t.TempDir() + dbPath := filepath.Join(dir, "async_test.db") + + // Bootstrap schema manually so we can insert test data BEFORE OpenStore + db, err := sql.Open("sqlite", dbPath+"?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)") + if err != nil { + t.Fatal(err) + } + // Create tables manually (minimal schema for this test) + _, err = db.Exec(` + CREATE TABLE _migrations (name TEXT PRIMARY KEY); + CREATE TABLE transmissions ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + raw_hex TEXT NOT NULL, + hash TEXT NOT NULL UNIQUE, + first_seen TEXT NOT NULL, + route_type INTEGER, + payload_type INTEGER, + payload_version INTEGER, + decoded_json TEXT, + created_at TEXT DEFAULT (datetime('now')), + channel_hash TEXT + ); + CREATE TABLE observers ( + id TEXT PRIMARY KEY, + name TEXT, + iata TEXT, + last_seen TEXT, + first_seen TEXT, + packet_count INTEGER DEFAULT 0, + model TEXT, + firmware TEXT, + client_version TEXT, + radio TEXT, + battery_mv INTEGER, + uptime_secs INTEGER, + noise_floor REAL, + inactive INTEGER DEFAULT 0, + last_packet_at TEXT + ); + CREATE TABLE nodes ( + public_key TEXT PRIMARY KEY, + name TEXT, role TEXT, lat REAL, lon REAL, + last_seen TEXT, first_seen TEXT, advert_count INTEGER DEFAULT 0, + battery_mv INTEGER, temperature_c REAL + ); + CREATE TABLE inactive_nodes ( + public_key TEXT PRIMARY KEY, + name TEXT, role TEXT, lat REAL, lon REAL, + last_seen TEXT, first_seen TEXT, advert_count INTEGER DEFAULT 0, + battery_mv INTEGER, temperature_c REAL + ); + CREATE TABLE observations ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + transmission_id INTEGER NOT NULL REFERENCES transmissions(id), + observer_idx INTEGER, + direction TEXT, + snr REAL, rssi REAL, score INTEGER, + path_json TEXT, + timestamp INTEGER NOT NULL, + raw_hex TEXT + ); + CREATE UNIQUE INDEX idx_observations_dedup ON observations(transmission_id, observer_idx, COALESCE(path_json, '')); + CREATE INDEX idx_observations_transmission_id ON observations(transmission_id); + CREATE INDEX idx_observations_observer_idx ON observations(observer_idx); + CREATE INDEX idx_observations_timestamp ON observations(timestamp); + CREATE TABLE observer_metrics ( + observer_id TEXT NOT NULL, + timestamp TEXT NOT NULL, + noise_floor REAL, tx_air_secs INTEGER, rx_air_secs INTEGER, + recv_errors INTEGER, battery_mv INTEGER, + packets_sent INTEGER, packets_recv INTEGER, + PRIMARY KEY (observer_id, timestamp) + ); + CREATE TABLE dropped_packets ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + hash TEXT, raw_hex TEXT, reason TEXT NOT NULL, + observer_id TEXT, observer_name TEXT, + node_pubkey TEXT, node_name TEXT, + dropped_at DATETIME DEFAULT CURRENT_TIMESTAMP + ); + `) + if err != nil { + t.Fatal("bootstrap schema:", err) + } + + // Mark all migrations as done EXCEPT the path_json backfill + for _, m := range []string{ + "advert_count_unique_v1", "noise_floor_real_v1", "node_telemetry_v1", + "obs_timestamp_index_v1", "observer_metrics_v1", "observer_metrics_ts_idx", + "observers_inactive_v1", "observer_metrics_packets_v1", "channel_hash_v1", + "dropped_packets_v1", "observations_raw_hex_v1", "observers_last_packet_at_v1", + "cleanup_legacy_null_hash_ts", + } { + db.Exec(`INSERT INTO _migrations (name) VALUES (?)`, m) + } + + // Insert a transmission + observations with NULL path_json and valid raw_hex + // raw_hex "0102AABBCCDD0000" has 2-hop path decodable by packetpath + rawHex := "41020304AABBCCDD05060708" + _, err = db.Exec(`INSERT INTO transmissions (raw_hex, hash, first_seen, payload_type) VALUES (?, 'hash1', '2025-01-01T00:00:00Z', 4)`, rawHex) + if err != nil { + t.Fatal("insert tx:", err) + } + // Insert 100 observations needing backfill + for i := 0; i < 100; i++ { + _, err = db.Exec(`INSERT INTO observations (transmission_id, observer_idx, timestamp, raw_hex, path_json) VALUES (1, ?, ?, ?, NULL)`, + i+1, 1700000000+i, rawHex) + if err != nil { + // dedup index might fire — use unique observer_idx + t.Fatalf("insert obs %d: %v", i, err) + } + } + db.Close() + + // Now open store via OpenStore — this must return QUICKLY (non-blocking) + start := time.Now() + store, err := OpenStoreWithInterval(dbPath, 300) + elapsed := time.Since(start) + if err != nil { + t.Fatal("OpenStore:", err) + } + defer store.Close() + + // OpenStore must return in under 2 seconds (backfill is no longer in applySchema) + if elapsed > 2*time.Second { + t.Fatalf("OpenStore blocked for %v — backfill must not run in applySchema", elapsed) + } + + // Backfill must NOT be recorded yet — it hasn't been triggered + var done int + err = store.db.QueryRow("SELECT 1 FROM _migrations WHERE name = 'backfill_path_json_from_raw_hex_v1'").Scan(&done) + if err == nil { + t.Fatal("migration recorded during OpenStore — backfill must be async via BackfillPathJSONAsync()") + } + + // Now trigger the async backfill (simulates what main.go does after OpenStore) + store.BackfillPathJSONAsync() + + // Wait for backfill to complete (should be very fast with 100 rows) + deadline := time.Now().Add(10 * time.Second) + for time.Now().Before(deadline) { + err = store.db.QueryRow("SELECT 1 FROM _migrations WHERE name = 'backfill_path_json_from_raw_hex_v1'").Scan(&done) + if err == nil { + break + } + time.Sleep(100 * time.Millisecond) + } + if err != nil { + t.Fatal("backfill never completed within 10s") + } + + // Verify backfill actually worked — observations should have non-NULL path_json + var nullCount int + store.db.QueryRow("SELECT COUNT(*) FROM observations WHERE path_json IS NULL").Scan(&nullCount) + if nullCount > 0 { + t.Errorf("backfill left %d observations with NULL path_json", nullCount) + } +} + +// TestBackfillPathJSONAsyncMethodExists verifies the async backfill API surface +// exists — BackfillPathJSONAsync must be callable independently from OpenStore. +func TestBackfillPathJSONAsyncMethodExists(t *testing.T) { + dir := t.TempDir() + dbPath := filepath.Join(dir, "method_test.db") + store, err := OpenStoreWithInterval(dbPath, 300) + if err != nil { + t.Fatal(err) + } + defer store.Close() + + // BackfillPathJSONAsync must exist as a method on *Store + // This is a compile-time check — if the method doesn't exist, the test won't compile. + store.BackfillPathJSONAsync() +} diff --git a/cmd/ingestor/main.go b/cmd/ingestor/main.go index 6dda8b99..e23fb2b3 100644 --- a/cmd/ingestor/main.go +++ b/cmd/ingestor/main.go @@ -57,6 +57,9 @@ func main() { defer store.Close() log.Printf("SQLite opened: %s", cfg.DBPath) + // Async backfill: path_json from raw_hex (#888) — must not block MQTT startup + store.BackfillPathJSONAsync() + // Check auto_vacuum mode and optionally migrate (#919) store.CheckAutoVacuum(cfg)