diff --git a/cmd/ingestor/resolved_path_backfill.go b/cmd/ingestor/resolved_path_backfill.go index 7463a9715..a29a4f1a5 100644 --- a/cmd/ingestor/resolved_path_backfill.go +++ b/cmd/ingestor/resolved_path_backfill.go @@ -229,6 +229,11 @@ func (s *Store) RunResolvedPathBackfill(ctx context.Context, batchSize int, paus return res, nil } +// resolvedPathBackfillResolve is the resolver a batch calls. It is a +// variable so the #267 test can check that no batch resolves while it holds +// the write transaction. +var resolvedPathBackfillResolve = resolveObservationPath + type resolvedPathBackfillRow struct { id int64 pathJSON string @@ -299,7 +304,7 @@ func (s *Store) resolvedPathBackfillBatch(ctx context.Context, after, ceiling in if r.payloadType == int(payloadADVERT) { from = strings.ToLower(r.fromPubkey) } - if rp := marshalResolvedPath(resolveObservationPath(hops, from, r.observerID, r.routeType, graph, idx)); rp != "" { + if rp := marshalResolvedPath(resolvedPathBackfillResolve(hops, from, r.observerID, r.routeType, graph, idx)); rp != "" { updates = append(updates, update{r.id, rp}) } } diff --git a/cmd/ingestor/resolved_path_backfill_188_test.go b/cmd/ingestor/resolved_path_backfill_188_test.go index cd8417c72..49bfd03cf 100644 --- a/cmd/ingestor/resolved_path_backfill_188_test.go +++ b/cmd/ingestor/resolved_path_backfill_188_test.go @@ -3,13 +3,17 @@ package main import ( "context" "database/sql" + "database/sql/driver" "encoding/json" "fmt" "os" "path/filepath" "strings" + "sync" "testing" "time" + + "modernc.org/sqlite" ) // Issue #188 point 4: backfill of NULL observations.resolved_path. @@ -17,7 +21,7 @@ import ( // backfillFixture188 opens a store at path with two relays sharing the // 1-byte prefix "c3" and an observer whose neighbour edge picks c3a, so a // non-advert flood observation with path ["c3"] resolves to c3a. -func backfillFixture188(t *testing.T, path string, seed bool) *Store { +func backfillFixture188(t testing.TB, path string, seed bool) *Store { t.Helper() store, err := OpenStore(path) if err != nil { @@ -43,7 +47,7 @@ func backfillFixture188(t *testing.T, path string, seed bool) *Store { // primeIndexAndGraph188 does what StartNeighborEdgesBuilder's warm-up does // before the backfill may run: prime the index, build neighbor_edges, and // publish the post-build graph. -func primeIndexAndGraph188(t *testing.T, store *Store) { +func primeIndexAndGraph188(t testing.TB, store *Store) { t.Helper() if err := store.RefreshPrefixIndex(); err != nil { t.Fatal(err) @@ -304,27 +308,238 @@ func TestResolvedPathBackfill_StopKeepsWatermark_188(t *testing.T) { } } -// resolvedPathBackfillHoldBudget is the stated per-batch budget for the write -// transaction at the default batch size (500 rows). It bounds how long a live -// InsertTransmission can wait behind one batch. -const resolvedPathBackfillHoldBudget = 250 * time.Millisecond - +// One batch bounds how long the backfill holds the single write connection, +// and so how long a live InsertTransmission can wait behind it. The bound is +// structural, so the test asserts the structure: each write transaction +// updates at most batchSize rows, there is one write transaction per batch, +// and no row is resolved while the write lock is held. +// +// It used to assert a wall-clock hold under 250 ms. That failed under CI load +// (#267): about three quarters of a batch's hold is the WAL fsync in COMMIT, +// whose latency depends on the runner's disk, not on this code. Resolution +// costs about 1 ms per batch, so moving it into the transaction did not trip +// the timing check either. BenchmarkResolvedPathBackfillBatch reports the +// hold. The hold logged here also includes the probe's triggers. func TestResolvedPathBackfill_WriteHoldUnderBudget_188(t *testing.T) { + const rows, batchSize = 5000, defaultResolvedPathBackfillBatchSize + probe := watchBackfillWrites267(t) store := backfillFixture188(t, filepath.Join(t.TempDir(), "ingest.db"), true) defer store.Close() - seedNullRows188(t, store, 5000) + seedNullRows188(t, store, rows) primeIndexAndGraph188(t, store) - res, err := store.RunResolvedPathBackfill(context.Background(), defaultResolvedPathBackfillBatchSize, 0) + probe.install(t, store) + txBefore := writerStatsAgg.get(resolvedPathBackfillComponent).snapshot().Count + + res, err := store.RunResolvedPathBackfill(context.Background(), batchSize, 0) if err != nil { t.Fatal(err) } - if res.Resolved != 5000 || res.Batches != 10 { - t.Fatalf("pass = %+v, want 5000 resolved in 10 batches", res) + if res.Resolved != rows || res.Batches != rows/batchSize { + t.Fatalf("pass = %+v, want %d resolved in %d batches", res, rows, rows/batchSize) + } + t.Logf("max write hold per %d-row batch: %s (with the probe's triggers)", batchSize, res.MaxHold) + + if n := writerStatsAgg.get(resolvedPathBackfillComponent).snapshot().Count - txBefore; n != int64(res.Batches) { + t.Fatalf("%d write transactions for %d batches, want one per batch", n, res.Batches) + } + perTx, unmarked := probe.rowsPerWatermark() + if unmarked != 0 { + t.Fatalf("%d row updates were not followed by a watermark write in their transaction", unmarked) + } + if len(perTx) != res.Batches { + t.Fatalf("%d watermark writes for %d batches, want one per batch", len(perTx), res.Batches) + } + total := 0 + for i, n := range perTx { + if n > batchSize { + t.Fatalf("write transaction %d updated %d rows, want at most %d", i, n, batchSize) + } + total += n + } + if total != rows { + t.Fatalf("write transactions updated %d rows in all, want %d", total, rows) + } + resolved, underLock, writesOutsideLock := probe.counts() + if writesOutsideLock != 0 { + t.Fatalf("%d backfill writes ran without writerMu", writesOutsideLock) + } + if resolved != rows { + t.Fatalf("the probe saw %d resolutions, want %d", resolved, rows) + } + if underLock != 0 { + t.Fatalf("%d of %d rows were resolved while the write transaction held writerMu", underLock, resolved) + } +} + +// backfillWriteProbe267 records, in order, every row the backfill resolves and +// every row and watermark it writes, and whether writerMu was held at the +// time. SQLite triggers report the writes from inside the transaction that +// makes them. +type backfillWriteProbe267 struct { + mu sync.Mutex + events []string // "r" row update, "w" watermark write + resolved int + underLock int + writesOutsideLock int +} + +var ( + backfillProbeOnce267 sync.Once + backfillProbeMu267 sync.Mutex + backfillProbe267 *backfillWriteProbe267 // the probe the SQL function reports to +) + +// writerMuHeld reports whether some goroutine holds writerMu. Only the +// backfill writes during the test, so a held lock is its write transaction. +func writerMuHeld() bool { + if writerMu.TryLock() { + writerMu.Unlock() + return false + } + return true +} + +// watchBackfillWrites267 starts a probe. It registers the SQL function the +// triggers call (before the store opens its connection, as the driver +// requires) and swaps in a resolver that reports to the probe. +func watchBackfillWrites267(t *testing.T) *backfillWriteProbe267 { + t.Helper() + backfillProbeOnce267.Do(func() { + err := sqlite.RegisterScalarFunction("backfill_probe_267", 1, func(_ *sqlite.FunctionContext, args []driver.Value) (driver.Value, error) { + backfillProbeMu267.Lock() + p := backfillProbe267 + backfillProbeMu267.Unlock() + if p == nil { + return nil, nil + } + switch kind := args[0].(type) { + case string: + p.write(kind) + case []byte: + p.write(string(kind)) + default: + return nil, fmt.Errorf("backfill_probe_267: kind %T", args[0]) + } + return nil, nil + }) + if err != nil { + t.Fatal(err) + } + }) + p := &backfillWriteProbe267{} + backfillProbeMu267.Lock() + backfillProbe267 = p + backfillProbeMu267.Unlock() + resolve := resolvedPathBackfillResolve + resolvedPathBackfillResolve = func(hops []string, fromPubkey, observerID string, routeType int, graph *NeighborGraph, idx prefixIndex) []*string { + p.resolve() + return resolve(hops, fromPubkey, observerID, routeType, graph, idx) + } + t.Cleanup(func() { + resolvedPathBackfillResolve = resolve + backfillProbeMu267.Lock() + backfillProbe267 = nil + backfillProbeMu267.Unlock() + }) + return p +} + +// install adds the triggers that report the backfill's writes: an observation +// whose resolved_path goes from NULL to a value, and a watermark write. +func (p *backfillWriteProbe267) install(t *testing.T, store *Store) { + t.Helper() + if err := ensureResolvedPathBackfillState(store.db); err != nil { + t.Fatal(err) + } + for _, q := range []string{ + `CREATE TRIGGER probe267_row AFTER UPDATE OF resolved_path ON observations + WHEN OLD.resolved_path IS NULL AND NEW.resolved_path IS NOT NULL + BEGIN SELECT backfill_probe_267('r'); END`, + `CREATE TRIGGER probe267_wm_insert AFTER INSERT ON resolved_path_backfill_state + BEGIN SELECT backfill_probe_267('w'); END`, + `CREATE TRIGGER probe267_wm_update AFTER UPDATE ON resolved_path_backfill_state + BEGIN SELECT backfill_probe_267('w'); END`, + } { + if _, err := store.db.Exec(q); err != nil { + t.Fatal(err) + } + } +} + +func (p *backfillWriteProbe267) resolve() { + held := writerMuHeld() + p.mu.Lock() + defer p.mu.Unlock() + p.resolved++ + if held { + p.underLock++ + } +} + +func (p *backfillWriteProbe267) write(kind string) { + held := writerMuHeld() + p.mu.Lock() + defer p.mu.Unlock() + p.events = append(p.events, kind) + if !held { + p.writesOutsideLock++ } - t.Logf("max write hold per %d-row batch: %s (budget %s)", defaultResolvedPathBackfillBatchSize, res.MaxHold, resolvedPathBackfillHoldBudget) - if res.MaxHold > resolvedPathBackfillHoldBudget { - t.Fatalf("max write hold %s exceeds the budget %s", res.MaxHold, resolvedPathBackfillHoldBudget) +} + +// rowsPerWatermark splits the writes at each watermark write. The batch writes +// its watermark last in its transaction, so each part is the rows of one +// transaction. unmarked counts rows after the last watermark write. +func (p *backfillWriteProbe267) rowsPerWatermark() (perTx []int, unmarked int) { + p.mu.Lock() + defer p.mu.Unlock() + for _, e := range p.events { + if e == "w" { + perTx = append(perTx, unmarked) + unmarked = 0 + } else { + unmarked++ + } + } + return perTx, unmarked +} + +func (p *backfillWriteProbe267) counts() (resolved, underLock, writesOutsideLock int) { + p.mu.Lock() + defer p.mu.Unlock() + return p.resolved, p.underLock, p.writesOutsideLock +} + +// BenchmarkResolvedPathBackfillBatch reports the write hold of one batch at +// the default batch size (hold-ms/batch). It is the measurement the wall-clock +// assertion in TestResolvedPathBackfill_WriteHoldUnderBudget_188 used to make +// (#267). +func BenchmarkResolvedPathBackfillBatch(b *testing.B) { + const batchSize = defaultResolvedPathBackfillBatchSize + store := backfillFixture188(b, filepath.Join(b.TempDir(), "ingest.db"), true) + defer store.Close() + ids := seedNullRows188(b, store, batchSize) + primeIndexAndGraph188(b, store) + if err := ensureResolvedPathBackfillState(store.db); err != nil { + b.Fatal(err) + } + var hold time.Duration + b.ResetTimer() + for i := 0; i < b.N; i++ { + b.StopTimer() + if _, err := store.db.Exec(`UPDATE observations SET resolved_path = NULL`); err != nil { + b.Fatal(err) + } + b.StartTimer() + r, err := store.resolvedPathBackfillBatch(context.Background(), 0, ids[len(ids)-1], batchSize) + if err != nil { + b.Fatal(err) + } + if r.Resolved != batchSize { + b.Fatalf("batch = %+v, want %d resolved", r, batchSize) + } + hold += r.Hold } + b.ReportMetric(float64(hold.Microseconds())/1000/float64(b.N), "hold-ms/batch") } func TestResolvedPathBackfillSettings_Defaults_188(t *testing.T) {