From b4126f0df39f5be37d966d22e46099b38e1ab45a Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Mon, 20 Jul 2026 11:01:10 -0500 Subject: [PATCH] measure backfill durability --- docs/grafana-dashboard.md | 29 +++- docs/semantic-parity.md | 2 +- docs/upstream-bootstrap-spec.md | 13 ++ src/internal/bootstrap/engine.zig | 227 +++++++++++++++++++++++---- src/internal/bootstrap/lifecycle.zig | 1 + src/internal/bootstrap/resync.zig | 20 ++- src/internal/bootstrap/retry.zig | 42 ++++- src/internal/metrics.zig | 168 +++++++++++++++++++- tests/oracle.py | 50 ++++++ 9 files changed, 504 insertions(+), 48 deletions(-) diff --git a/docs/grafana-dashboard.md b/docs/grafana-dashboard.md index edf35b2..1ac69b4 100644 --- a/docs/grafana-dashboard.md +++ b/docs/grafana-dashboard.md @@ -117,11 +117,26 @@ The following dashboard families already have non-placeholder producers: - `jetstream_orchestrator_merge_repo_revs_updated_total` - `jetstream_orchestrator_merge_dids_discovered_post_bootstrap_total` - `jetstream_verifier_failures_total` +- `jetstream_backfill_discovered_total` - `jetstream_backfill_completed_total` - `jetstream_backfill_failed_total` +- `jetstream_backfill_active_flips_total` +- `jetstream_backfill_on_fail_store_errors_total` +- `jetstream_backfill_handle_repo_duration_seconds` +- `jetstream_backfill_progress_completed` +- `jetstream_backfill_completion_queued_total` +- `jetstream_backfill_completion_queue_depth` +- `jetstream_backfill_completion_durable_batches_total` +- `jetstream_backfill_completion_durable_repos_total` +- `jetstream_backfill_completion_stage_errors_total` +- `jetstream_backfill_completion_queue_wait_seconds` +- `jetstream_backfill_forced_checkpoint_flushes_total` +- `jetstream_backfill_failed_repo_retry_passes_total` +- `jetstream_backfill_failed_repo_retry_candidates_total` - `jetstream_backfill_failed_repo_retry_attempts_total` - `jetstream_backfill_failed_repo_retry_succeeded_total` - `jetstream_backfill_failed_repo_retry_failed_total` +- `jetstream_backfill_failed_repo_retry_skipped_host_parked_total` - `jetstream_subscribe_subscribers` - `jetstream_subscribe_events_sent_total` - `jetstream_subscribe_bytes_sent_total` @@ -151,6 +166,14 @@ The following dashboard families already have non-placeholder producers: - `jetstream_import_segments_patched_total` - `jetstream_import_bytes_rewritten_total` -The checksum-pinned deployment contract remains red until every other queried -family is backed by its real lifecycle boundary. The current binary is not a -dashboard-complete release. +The backfill families are measured at Stream's actual page checkpoint: worker +results enter the completion queue after download/handling, archive bytes are +flushed first, and only the succeeding RocksDB transition batch removes them +from queue depth and reports completion. This is the semantic counterpart of +upstream's writer-hook completion batcher without manufacturing a second queue. +Offline loopback receipts cover a terminal RepoNotFound, a generated complete +CAR, redirected rate limiting, and a restart-safe parked-host skip. + +The checksum-pinned deployment contract remains red until the remaining +compaction families are backed by their real lifecycle boundaries. The current +binary is not a dashboard-complete release. diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 90c15e6..24a85b5 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -20,7 +20,7 @@ this document or `bootstrap-semantic-parity.md` is open. | Sync 1.1 resync ordering | Upstream serializes per DID, drops stale async repairs, and emits a sync tombstone plus authoritative replacements | **closed at the implementation boundary.** The live coordinator has 32 real fetch workers, a 64-job queue, a five-minute fetch budget plus the 1,000 B/s-for-30-seconds slow-transfer guard, a bounded 16,384-entry 5/minute-per-DID limiter, 2,048 pending commits with drop-oldest overflow, fetch outside the DID lane, authenticated apply with older/equal-contradictory head rejection, ordered pending replay, and trigger-ticket-gated outbox delivery. Commit divergence schedules async repair; `#sync` divergence attempts inline repair and queues one retry only for transient failure, preserving the original relay cursor only on inline success. The writer emits `sync` then bounded 1,024-row `create_resync` batches and stages the fetched/pending chain checkpoints under the archive lock before any matching fsync. Loopback-only tests exercise a real HTTP getRepo mmap, real DID resolution/signature, buffering while the response is held open, slow-transfer cancellation, durable RocksDB promotion/cursor coupling, v2 tail publication, and sealed JSS row order. | | Timestamp import | Canonical dashboard and API include timestamp import lifecycle/rewrite metrics and behavior | **Closed at the implementation boundary.** The format preserves distinct `witnessed_at` and sentinel-zero `indexed_at` columns. Subscriber v1/v2 encoding applies upstream's `indexed_at != 0 ? indexed_at : witnessed_at` display rule while ranges and timestamp cursors remain anchored to immutable witness envelopes. The strict seekable RFC 4180 parser matches upstream's header, row rejection, byte-offset, CID, RFC3339, 64 KiB, and bounded-sampling contracts in a forward-only allocation-free pass. A separate Bloom-filtered RocksDB rule store bulk-loads sorted external SST chunks with CSV-order last-write-wins, reconstructs its resident collection gate on open, and implements specific-CID precedence with all-version fallback. Its central archive-lock hook stamps live, repair, bootstrap, and merge materializations before sequence assignment, and feeds the mutated rows to immediate delivery so hot and cold encodings agree. The sealed-JSS patch primitive preserves block topology/envelopes, clean frames, and the opaque bloom/collection footer while atomically changing only `indexed_at`; it is exercised against an upstream-produced segment. A concrete sealed-segment bloom catalog, bounded DID/FD bucketer, fsynced packed offsets, positioned revalidation, and collision-safe per-segment patch plan implement Phase B/C with specific-CID precedence and idempotent disk-to-disk coverage. A dedicated archive rewrite mutex serializes the real delete compactor with timestamp patching while leaving live append available. The unified metadata DB synchronously persists the current-job pointer, complete job record, and per-segment done set. The real runner activates rules before force-seal, durably crosses the bucketed handoff, checkpoints only after patch fsync/rename/dir-sync, retains partial cancellation progress, and resumes after full archive/rule/metadata reopen. The background manager canonicalizes symlinks, confines regular files, enforces one durable nonterminal job, and adopts it on startup. Bearer-gated import/status XRPC matches the pinned lexicons, fixed-width digest comparison, error names/statuses, steady-state admission, and current/by-id status behavior; a real loopback HTTP run exercised it against the local simulator. Every canonical import metric is produced at its real parse, route, patch, or durable-terminal boundary; cancellation remains a resumable pause and is not counted as terminal failure. | | Status/diagnostics | Repo attempts/error class/final PDS and per-host aggregates survive restart and are operator-queryable | **closed at the implementation boundary.** Durable v5 repo rows distinguish initial Backfill.Rev from latest Rev/UpdatedAt, retain resolved Handle/PDS and reserved record/byte fields, and atomically maintain `handle/`. Merge source cursor commits refresh latest revisions without mutating backfill watermarks. Same-batch `host/` totals, active/status counts, cumulative error classes, and five bounded recent samples survive restart; host moves and active flips are atomic with the source repo row. `/status?tab=hosts` reads those rows without a whole-network scan and matches upstream ordering and error presentation. Account queries preserve resolver-first handle semantics, the legacy `did`/`handle` aliases, and missing-identity hydration, then reconstruct records last-writer-wins from checksum-verified, DID-bloom-pruned sealed JSS plus rotation-safe active/pending and bootstrap-live snapshots. The resulting canonical MST is compared with the PDS commit root through real Sync 1.1 `getLatestCommit`/`getBlocks`; the page renders the same match fields and error state as upstream. `HEAD` performs lookup without verification; successful status responses are `no-store` and carry a generation timestamp. Expensive verification uses upstream's exact 4/source-IP/minute fixed-window limiter, 4,096-entry bound, stale pruning, oldest eviction, and explicit disable option. The offline receipt uses a real HTTP Stream route, real RocksDB row, real JSS archive, signed commit, and loopback PLC/PDS; four separate connections verify successfully, a fifth changing only its ephemeral port is blocked without another upstream request, HEAD remains side-effect-free, and disabling the limiter performs a new authoritative verification. | -| Prometheus/Grafana | Every canonical dashboard query maps to a real, same-semantics producer; runtime-specific panels use honest Zig/process metrics | exact checksum-pinned upstream dashboard is vendored with an offline fail-closed scrape contract; canonical build/live-firehose/archive-append/drop/verifier/backfill/timestamp-import producers are present. The complete canonical subscriber family set is exported from real scan, delivery, rejection, cursor, and disconnect boundaries; `none`, `deflate`, and `zstd` labels now arise only from the actual negotiated delivery path. Live frame inspection distinguishes gaps, bounded-label server errors, forward-compatible unknowns, and malformed CBOR at the raw transport boundary; a Zat reconnect callback measures attempts rather than connection failures. The only dashboard divergence is now a deterministic, checksum-pinned process-runtime row: same-semantics start/CPU/RSS/FD measurements, real VM/thread/page-fault/context-switch signals, no `go_*` aliases, and no host-network counters mislabeled as process I/O. A production-binary loopback receipt proves values and cumulative monotonicity on independent scrapes. The hot readable log now matches upstream's durability pin: rows above the archive watermark may exceed the byte budget but cannot be evicted, publication releases their pinned-byte accounting, and only then may retention advance the floor; all three canonical readable-log gauges come from that exact state. Data-directory free bytes are collected fresh from the archive directory's filesystem at scrape time, and the canonical seal histogram measures only successful flush-through-footer/header-fsync work with upstream's exact 15 slow-latency buckets. HTTP duration uses upstream's public-mux boundary and stable labels, including inline WebSocket lifetime and the single `xrpc/` subtree; debug and unmatched routes remain unobserved. getBlock outcomes, duration, and served bytes come from actual response completion, with bytes counted only for complete 200 responses. A ReleaseSafe offline receipt proves 200/206/304/400/404/416 behavior and WebSocket close accounting. Sealed headers and block indexes are now genuinely resident: startup and verified refresh latency, loaded count, and lookup hits/misses use upstream's exact boundaries. Refcounted snapshots remain memory-safe across refresh, while fresh-file offsets prevent cross-generation splicing; the pinned official client proves cold replay produces hits. RocksDB point reads, durable point writes/deletes, and atomic batch commits populate upstream's exact `{op,status}` store histogram and fast-latency buckets; real-database and production-startup receipts prove the non-error paths. Durable lifecycle phase/transition metrics, all six cutover-state timings, and merge counters are wired to the actual bootstrap/merge state machine. Merge now performs one cached Backfill.Rev lookup per DID, including misses; real JSS/RocksDB and eight-crash production receipts prove counter and restart behavior. **Remaining resident bloom/collection planning, backfill, and compaction families open.** | +| Prometheus/Grafana | Every canonical dashboard query maps to a real, same-semantics producer; runtime-specific panels use honest Zig/process metrics | exact checksum-pinned upstream dashboard is vendored with an offline fail-closed scrape contract; canonical build/live-firehose/archive-append/drop/verifier/backfill/timestamp-import producers are present. The complete canonical subscriber family set is exported from real scan, delivery, rejection, cursor, and disconnect boundaries; `none`, `deflate`, and `zstd` labels now arise only from the actual negotiated delivery path. Live frame inspection distinguishes gaps, bounded-label server errors, forward-compatible unknowns, and malformed CBOR at the raw transport boundary; a Zat reconnect callback measures attempts rather than connection failures. The only dashboard divergence is now a deterministic, checksum-pinned process-runtime row: same-semantics start/CPU/RSS/FD measurements, real VM/thread/page-fault/context-switch signals, no `go_*` aliases, and no host-network counters mislabeled as process I/O. A production-binary loopback receipt proves values and cumulative monotonicity on independent scrapes. The hot readable log now matches upstream's durability pin: rows above the archive watermark may exceed the byte budget but cannot be evicted, publication releases their pinned-byte accounting, and only then may retention advance the floor; all three canonical readable-log gauges come from that exact state. Data-directory free bytes are collected fresh from the archive directory's filesystem at scrape time, and the canonical seal histogram measures only successful flush-through-footer/header-fsync work with upstream's exact 15 slow-latency buckets. HTTP duration uses upstream's public-mux boundary and stable labels, including inline WebSocket lifetime and the single `xrpc/` subtree; debug and unmatched routes remain unobserved. getBlock outcomes, duration, and served bytes come from actual response completion, with bytes counted only for complete 200 responses. A ReleaseSafe offline receipt proves 200/206/304/400/404/416 behavior and WebSocket close accounting. Sealed headers and block indexes are now genuinely resident: startup and verified refresh latency, loaded count, and lookup hits/misses use upstream's exact boundaries. Refcounted snapshots remain memory-safe across refresh, while fresh-file offsets prevent cross-generation splicing; the pinned official client proves cold replay produces hits. RocksDB point reads, durable point writes/deletes, and atomic batch commits populate upstream's exact `{op,status}` store histogram and fast-latency buckets; real-database and production-startup receipts prove the non-error paths. Durable lifecycle phase/transition metrics, all six cutover-state timings, and merge counters are wired to the actual bootstrap/merge state machine. Merge now performs one cached Backfill.Rev lookup per DID, including misses; real JSS/RocksDB and eight-crash production receipts prove counter and restart behavior. Backfill discovery, active flips, handler timing, completion queue/durability, current-run progress, forced checkpoints, failure persistence, and steady retry scans now use Stream's real worker/page/RocksDB boundaries. Generated-CAR, RepoNotFound, redirected-429, durable-host-park, and production crash receipts prevent zero-series or API-shape substitutions. **Only the compaction dashboard families remain open.** | | Oracle | Event-log equivalence, final-state convergence, crash/power-loss, verifier repair, hostile input, and anti-vacuity receipts | **open — current crash matrix is necessary but substantially smaller than upstream's oracle** | ## Experiment admission diff --git a/docs/upstream-bootstrap-spec.md b/docs/upstream-bootstrap-spec.md index 97c2c15..f1aaf43 100644 --- a/docs/upstream-bootstrap-spec.md +++ b/docs/upstream-bootstrap-spec.md @@ -102,4 +102,17 @@ until steady_state. failed-repo retry loop: 4h interval, 16 workers, 4/host, are updated only at these durable lifecycle boundaries. The offline crash matrix checks their restart semantics, while a deterministic real-JSS and real-RocksDB merge receipt proves kept/dropped/source/lookup/revision counts. +- Stream's page result list is the real completion queue: a successful worker + queues exactly once, the page archive checkpoint makes every appended row + durable, and the following atomic RocksDB batch makes the queued repo states + durable. Queue depth, wait time, durable batch/repo totals, completed totals, + and current-run progress advance only after that batch succeeds. A failed + completion batch remains queued and increments the stage-error signal. + RepoNotFound follows the same completion path with zero appended rows. +- Discovery and active-flip counters advance only after their RocksDB writes; + successful HandleRepo timing covers validated CAR traversal and archive + append but excludes download and serialization wait. Failed-repo retry pass, + due-candidate, attempt, outcome, and parked-host counters use the same scan + and post-gate boundaries as upstream. Failure-persistence errors are counted + per affected repo rather than being disguised as remote getRepo failures. - reproduce crashpoint seams at every commit boundary for future oracle tests. diff --git a/src/internal/bootstrap/engine.zig b/src/internal/bootstrap/engine.zig index 4dd9009..49aea95 100644 --- a/src/internal/bootstrap/engine.zig +++ b/src/internal/bootstrap/engine.zig @@ -76,6 +76,10 @@ const Result = struct { error_class: repo_store.ErrorClass = .unknown, attempts: u32 = 0, started_us: i64 = 0, + /// Monotonic time when a successful completion entered the real worker + /// result queue. It remains queued until the page's archive checkpoint + /// and metadata batch are both durable. + completion_queued_awake_us: i64 = 0, rows: u64 = 0, drops: repos.DropCounts = .{}, }; @@ -335,6 +339,7 @@ const Work = struct { job.did, self.archive, &self.convert_mu, + self.config.metrics, ) catch |err| { if (self.config.metrics) |m| _ = m.backfill_workers_active.fetchSub(1, .monotonic); if (isFatalDownloadError(err)) { @@ -361,14 +366,24 @@ const Work = struct { fn finish(self: *Work, maybe_result: ?Result) void { self.result_mu.lockUncancelable(self.io); defer self.result_mu.unlock(self.io); - if (maybe_result) |result| { + defer { + self.finished += 1; + self.result_ready.broadcast(self.io); + } + if (maybe_result) |raw_result| { + var result = raw_result; + if (result.status == .complete) + result.completion_queued_awake_us = Io.Timestamp.now(self.io, .awake).toMicroseconds(); self.results.append(self.allocator, result) catch { freeResult(self.allocator, result); if (self.fatal_error == null) self.fatal_error = error.OutOfMemory; + return; + }; + if (result.status == .complete) if (self.config.metrics) |m| { + _ = m.backfill_completion_queued_total.fetchAdd(1, .monotonic); + _ = m.backfill_completion_queue_depth.fetchAdd(1, .monotonic); }; } - self.finished += 1; - self.result_ready.broadcast(self.io); } fn publishGauges(self: *Work) void { @@ -418,6 +433,7 @@ pub fn run( defer work.shutdown(); var stats: Stats = .{}; + if (config.metrics) |m| m.backfill_progress_completed.store(0, .monotonic); var selected: usize = 0; // An explicit selection is a self-contained debug run. It must not leave // an old full-network discovery watermark for merge to consume later. @@ -446,6 +462,7 @@ pub fn run( var jobs: std.ArrayList(Job) = .empty; var dispatched: std.ArrayList(repo_store.Transition) = .empty; var discovered_inactive: std.ArrayList(repo_store.Transition) = .empty; + var discovered_active: usize = 0; for (page.repos) |ref| { if (config.max_repos > 0 and selected >= config.max_repos) break; @@ -453,6 +470,7 @@ pub fn run( if (selected_page) { if (existing == null) { try store.put(ref.did, .{ .status = .not_started, .active = true, .started_us = Io.Timestamp.now(io, .real).toMicroseconds() }); + if (config.metrics) |m| _ = m.backfill_discovered_total.fetchAdd(1, .monotonic); existing = try store.get(alloc, ref.did); } // Upstream reconciles an inactive selected row before @@ -462,32 +480,36 @@ pub fn run( var next = state; next.active = ref.active; try store.put(ref.did, next); + if (config.metrics) |m| _ = m.backfill_active_flips_total.fetchAdd(1, .monotonic); existing = try store.get(alloc, ref.did); }; try resolveAndRecordSelectedIdentity(alloc, io, config.plc_url.?, store, ref.did, existing.?); existing = try store.get(alloc, ref.did); } if (existing) |state| { - if (state.active != ref.active) try store.put(ref.did, .{ - .status = state.status, - .active = ref.active, - .rev = state.rev, - .latest_rev = state.latest_rev, - .host = state.host, - .handle = state.handle, - .pds = state.pds, - .last_error = state.last_error, - .error_class = state.error_class, - .attempts = state.attempts, - .last_attempt_us = state.last_attempt_us, - .started_us = state.started_us, - .completed_us = state.completed_us, - .retry_count = state.retry_count, - .next_attempt_us = state.next_attempt_us, - .updated_us = state.updated_us, - .record_count = state.record_count, - .total_bytes = state.total_bytes, - }); + if (state.active != ref.active) { + try store.put(ref.did, .{ + .status = state.status, + .active = ref.active, + .rev = state.rev, + .latest_rev = state.latest_rev, + .host = state.host, + .handle = state.handle, + .pds = state.pds, + .last_error = state.last_error, + .error_class = state.error_class, + .attempts = state.attempts, + .last_attempt_us = state.last_attempt_us, + .started_us = state.started_us, + .completed_us = state.completed_us, + .retry_count = state.retry_count, + .next_attempt_us = state.next_attempt_us, + .updated_us = state.updated_us, + .record_count = state.record_count, + .total_bytes = state.total_bytes, + }); + if (config.metrics) |m| _ = m.backfill_active_flips_total.fetchAdd(1, .monotonic); + } } if (!ref.active) { if (existing == null) @@ -498,6 +520,7 @@ pub fn run( if (state.status != .failed and state.status != .not_started) continue; } try jobs.append(alloc, .{ .did = ref.did }); + if (existing == null) discovered_active += 1; try dispatched.append(alloc, .{ .did = ref.did, .state = if (existing) |state| stateWithStatus(state, .not_started) else .{ .status = .not_started } }); selected += 1; } @@ -505,7 +528,9 @@ pub fn run( // One synced dispatch batch bounds crash fallout to this page and // makes every worker-owned DID visible before any archive append. try store.putMany(discovered_inactive.items); + if (config.metrics) |m| _ = m.backfill_discovered_total.fetchAdd(discovered_inactive.items.len, .monotonic); try store.putMany(dispatched.items); + if (config.metrics) |m| _ = m.backfill_discovered_total.fetchAdd(discovered_active, .monotonic); if (jobs.items.len > 0 and crashpoint.wouldHit(.mid_backfill_download)) crashpoint.hit(.mid_backfill_download); @@ -525,25 +550,58 @@ pub fn run( } // Archive first, metadata second. All worker appends have drained. + if (config.metrics) |m| _ = m.backfill_forced_checkpoint_flushes_total.fetchAdd(1, .monotonic); try archive.flushBlock(); const transitions = try alloc.alloc(repo_store.Transition, results.items.len); for (results.items, transitions) |result, *transition| { const existing = try store.get(alloc, result.did); const terminal_us = Io.Timestamp.now(io, .real).toMicroseconds(); - const downloaded = result.status == .complete and result.rev.len > 0; transition.* = .{ .did = result.did, .state = terminalState(result, existing, terminal_us) }; - if (downloaded) { + } + store.putMany(transitions) catch |err| { + var has_completion = false; + var failed_persists: u64 = 0; + for (results.items) |result| switch (result.status) { + .complete => has_completion = true, + .failed => failed_persists += 1, + else => {}, + }; + if (config.metrics) |m| { + if (has_completion) _ = m.backfill_completion_stage_errors_total.fetchAdd(1, .monotonic); + if (failed_persists > 0) _ = m.backfill_on_fail_store_errors_total.fetchAdd(failed_persists, .monotonic); + } + return err; + }; + + var durable_completions: usize = 0; + const durable_awake_us = Io.Timestamp.now(io, .awake).toMicroseconds(); + for (results.items) |result| switch (result.status) { + .complete => { stats.completed += 1; stats.rows += result.rows; stats.drops.invalid_path += result.drops.invalid_path; stats.drops.field_too_long += result.drops.field_too_long; - if (config.metrics) |m| _ = m.backfill_completed_total.fetchAdd(1, .monotonic); - } else if (result.status == .failed) { + durable_completions += 1; + if (config.metrics) |m| { + _ = m.backfill_completed_total.fetchAdd(1, .monotonic); + if (result.completion_queued_awake_us > 0 and durable_awake_us >= result.completion_queued_awake_us) + m.observeBackfillCompletionWait(@intCast(durable_awake_us - result.completion_queued_awake_us)); + } + }, + .failed => { stats.failed += 1; if (config.metrics) |m| _ = m.backfill_failed_total.fetchAdd(1, .monotonic); + }, + else => {}, + }; + if (config.metrics) |m| { + m.backfill_progress_completed.store(@intCast(stats.completed), .monotonic); + if (durable_completions > 0) { + _ = m.backfill_completion_queue_depth.fetchSub(durable_completions, .monotonic); + _ = m.backfill_completion_durable_batches_total.fetchAdd(1, .monotonic); + _ = m.backfill_completion_durable_repos_total.fetchAdd(durable_completions, .monotonic); } } - try store.putMany(transitions); if (config.max_repos > 0 and selected >= config.max_repos) break :paging; const next = page.cursor orelse break; @@ -657,6 +715,7 @@ test "explicit selected backfill resolves identity, bypasses listRepos, and pers defer store.deinit(); try meta.putDurable(bootstrap_cursor_key, "stale-full-network-cursor"); const selected = [_][]const u8{"did:plc:abcdefghijklmnopqrstuvwx"}; + var instrumentation: metrics_mod.Stats = .{}; const stats = try run(testing.allocator, io, .{ .relay_http = url, .plc_url = url, @@ -664,10 +723,21 @@ test "explicit selected backfill resolves identity, bypasses listRepos, and pers .selected_repos = &selected, .workers = 1, .max_inflight_bytes = 1024 * 1024, + .metrics = &instrumentation, }, &archive, &store); - // RepoNotFound is terminal-complete state, but upstream does not count it - // as a downloaded repository completion. - try testing.expectEqual(@as(u64, 0), stats.completed); + // A RepoNotFound response is a successful durable completion upstream: + // there is no CAR to handle, but the DID must not be retried forever. + try testing.expectEqual(@as(u64, 1), stats.completed); + try testing.expectEqual(@as(u64, 1), instrumentation.backfill_discovered_total.load(.monotonic)); + try testing.expectEqual(@as(u64, 1), instrumentation.backfill_completion_queued_total.load(.monotonic)); + try testing.expectEqual(@as(usize, 0), instrumentation.backfill_completion_queue_depth.load(.monotonic)); + try testing.expectEqual(@as(u64, 1), instrumentation.backfill_completion_durable_batches_total.load(.monotonic)); + try testing.expectEqual(@as(u64, 1), instrumentation.backfill_completion_durable_repos_total.load(.monotonic)); + try testing.expectEqual(@as(u64, 1), instrumentation.backfill_completion_queue_wait_count.load(.monotonic)); + try testing.expectEqual(@as(u64, 1), instrumentation.backfill_completed_total.load(.monotonic)); + try testing.expectEqual(@as(usize, 1), instrumentation.backfill_progress_completed.load(.monotonic)); + try testing.expectEqual(@as(u64, 1), instrumentation.backfill_forced_checkpoint_flushes_total.load(.monotonic)); + try testing.expectEqual(@as(u64, 0), instrumentation.backfill_handle_repo_duration_count.load(.monotonic)); const state = (try store.get(testing.allocator, selected[0])).?; defer store.freeState(testing.allocator, state); try testing.expectEqual(repo_store.Status.complete, state.status); @@ -724,6 +794,7 @@ test "explicit selected backfill repairs identity metadata without redownloading const did = "did:plc:abcdefghijklmnopqrstuvwx"; try store.put(did, .{ .status = .complete, .active = false, .rev = "backfill", .latest_rev = "latest" }); const selected = [_][]const u8{did}; + var instrumentation: metrics_mod.Stats = .{}; const stats = try run(testing.allocator, io, .{ .relay_http = url, .plc_url = url, @@ -731,8 +802,11 @@ test "explicit selected backfill repairs identity metadata without redownloading .selected_repos = &selected, .workers = 1, .max_inflight_bytes = 1024 * 1024, + .metrics = &instrumentation, }, &archive, &store); try testing.expectEqual(@as(u64, 0), stats.completed); + try testing.expectEqual(@as(u64, 1), instrumentation.backfill_active_flips_total.load(.monotonic)); + try testing.expectEqual(@as(u64, 0), instrumentation.backfill_discovered_total.load(.monotonic)); const state = (try store.get(testing.allocator, did)).?; defer store.freeState(testing.allocator, state); try testing.expectEqual(repo_store.Status.complete, state.status); @@ -880,6 +954,7 @@ fn downloadOne( did: []const u8, archive: *archive_mod.Archive, convert_mu: *Io.Mutex, + instrumentation: ?*metrics_mod.Stats, ) !Result { var scratch_buf: [Io.Dir.max_path_bytes]u8 = undefined; const scratch_dir = try std.fmt.bufPrint(&scratch_buf, "{s}/backfill/repo-scratch", .{data_dir}); @@ -923,6 +998,7 @@ fn downloadOne( defer prepared.deinit(io); convert_mu.lockUncancelable(io); defer convert_mu.unlock(io); + const handle_started = Io.Timestamp.now(io, .awake); var local: Stats = .{}; const now_us: i64 = Io.Timestamp.now(io, .real).toMicroseconds(); const Sink = struct { @@ -946,6 +1022,10 @@ fn downloadOne( }, else => return failureResult(result_allocator, did, &final_host, attempts, started_us, .car, @errorName(err)), }; + if (instrumentation) |m| m.observeBackfillHandleRepo(@intCast(@max( + handle_started.durationTo(Io.Timestamp.now(io, .awake)).toMicroseconds(), + 0, + ))); const owned_rev = try result_allocator.dupe(u8, head_rev); errdefer result_allocator.free(owned_rev); const owned_latest_rev = try result_allocator.dupe(u8, head_rev); @@ -1116,6 +1196,7 @@ test "bootstrap redirects preserve independent transient and rate-limit retry bu var origin_url_buf: [96]u8 = undefined; const origin_url = try std.fmt.bufPrint(&origin_url_buf, "http://127.0.0.1:{d}", .{origin.socket.address.getPort()}); + var instrumentation: metrics_mod.Stats = .{}; const result = try downloadOne( testing.allocator, testing.allocator, @@ -1125,10 +1206,12 @@ test "bootstrap redirects preserve independent transient and rate-limit retry bu "did:plc:retry-fixture", &archive, &convert_mu, + &instrumentation, ); defer freeResult(testing.allocator, result); try testing.expectEqual(repo_store.Status.complete, result.status); + try testing.expectEqual(@as(u64, 0), instrumentation.backfill_handle_repo_duration_count.load(.monotonic)); try testing.expectEqual(@as(u32, 4), result.attempts); var expected_host_buf: [64]u8 = undefined; const expected_host = try std.fmt.bufPrint(&expected_host_buf, "127.0.0.1:{d}", .{target.socket.address.getPort()}); @@ -1136,3 +1219,85 @@ test "bootstrap redirects preserve independent transient and rate-limit retry bu try origin_future.await(io); try target_future.await(io); } + +test "successful real CAR handling records only the handler boundary" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + + var fixture_arena = std.heap.ArenaAllocator.init(testing.allocator); + defer fixture_arena.deinit(); + const fixture_alloc = fixture_arena.allocator(); + const record = "offline backfill handler timing fixture"; + const record_cid = try zat.cbor.Cid.forDagCbor(fixture_alloc, record); + var tree = zat.mst.Mst.init(fixture_alloc); + try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); + const data_cid = try tree.rootCid(); + const keypair = try zat.Keypair.fromSecretKey(.p256, .{17} ** 32); + const embedded_did = try keypair.did(fixture_alloc); + const signed = try zat.signCommit(fixture_alloc, .{ + .did = embedded_did, + .rev = "3k2abcdefghij", + .data = data_cid, + }, &keypair); + var blocks: std.ArrayList(zat.car.Block) = .empty; + try blocks.append(fixture_alloc, .{ .cid_raw = signed.cid.raw, .data = signed.bytes }); + try tree.collectBlocks(&blocks); + try blocks.append(fixture_alloc, .{ .cid_raw = record_cid.raw, .data = record }); + const car_bytes = try zat.car.writeAlloc(fixture_alloc, .{ + .roots = &.{signed.cid}, + .blocks = blocks.items, + }); + + var listener = (Io.net.IpAddress{ .ip4 = .loopback(0) }).listen(io, .{ .reuse_address = true }) catch unreachable; + defer listener.deinit(io); + const Fixture = struct { + fn serve(server: *Io.net.Server, task_io: Io, body: []const u8) !void { + var stream = try server.accept(task_io); + defer stream.close(task_io); + var read_buf: [2048]u8 = undefined; + var write_buf: [4096]u8 = undefined; + var reader = stream.reader(task_io, &read_buf); + var writer = stream.writer(task_io, &write_buf); + var http_server = std.http.Server.init(&reader.interface, &writer.interface); + var request = try http_server.receiveHead(); + try request.respond(body, .{ + .status = .ok, + .keep_alive = false, + .extra_headers = &.{.{ .name = "content-type", .value = "application/vnd.ipld.car" }}, + }); + } + }; + var server_future = try io.concurrent(Fixture.serve, .{ &listener, io, car_bytes }); + defer _ = server_future.cancel(io) catch {}; + + var tmp = testing.tmpDir(.{}); + defer tmp.cleanup(); + var data_dir_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const data_dir_len = try tmp.dir.realPath(io, &data_dir_buf); + const data_dir = data_dir_buf[0..data_dir_len]; + var archive = try archive_mod.Archive.init(testing.allocator, io, data_dir); + defer archive.deinit(); + var convert_mu: Io.Mutex = .init; + var url_buf: [96]u8 = undefined; + const url = try std.fmt.bufPrint(&url_buf, "http://127.0.0.1:{d}", .{listener.socket.address.getPort()}); + var instrumentation: metrics_mod.Stats = .{}; + const result = try downloadOne( + testing.allocator, + testing.allocator, + io, + data_dir, + url, + "did:plc:authoritative-listrepos-did", + &archive, + &convert_mu, + &instrumentation, + ); + defer freeResult(testing.allocator, result); + + try testing.expectEqual(repo_store.Status.complete, result.status); + try testing.expectEqual(@as(u64, 1), result.rows); + try testing.expectEqual(@as(u64, 1), instrumentation.backfill_handle_repo_duration_count.load(.monotonic)); + try testing.expectEqual(@as(u64, 0), instrumentation.backfill_completion_queued_total.load(.monotonic)); + try server_future.await(io); +} diff --git a/src/internal/bootstrap/lifecycle.zig b/src/internal/bootstrap/lifecycle.zig index ce5db61..a67ad4d 100644 --- a/src/internal/bootstrap/lifecycle.zig +++ b/src/internal/bootstrap/lifecycle.zig @@ -225,6 +225,7 @@ fn runPendingPass( for (dids) |did| { var details: resync.AttemptDetails = .{}; _ = resync.resyncRepo(allocator, io, scratch_dir, relay_http, did, dst, repo_st, &drops, &details) catch |err| { + if (details.local_failure) return err; switch (err) { error.OutOfMemory => return err, else => {}, diff --git a/src/internal/bootstrap/resync.zig b/src/internal/bootstrap/resync.zig index e199578..a40f4bb 100644 --- a/src/internal/bootstrap/resync.zig +++ b/src/internal/bootstrap/resync.zig @@ -34,6 +34,10 @@ pub const AttemptDetails = struct { final_host: fetched_car.FinalHost = .{}, retry_after_s: ?u64 = null, error_class: repo_store.ErrorClass = .unknown, + /// True once the network/CAR input has been accepted and the failure is + /// attributable to Stream's archive or metadata path. Callers must abort + /// lifecycle work instead of recording this as a remote repo failure. + local_failure: bool = false, }; /// download `did` via the relay and append its replacement stream @@ -56,7 +60,10 @@ pub fn resyncRepo( defer arena.deinit(); const alloc = arena.allocator(); - const previous = try store.get(allocator, did); + const previous = store.get(allocator, did) catch |err| { + details.local_failure = true; + return err; + }; defer if (previous) |state| store.freeState(allocator, state); const active = if (previous) |state| state.active else true; const outcome = prepared_fetch.fetch(allocator, alloc, io, scratch_dir, relay_http, did, &details.final_host) catch |err| { @@ -71,11 +78,17 @@ pub fn resyncRepo( return error.RateLimited; }, .not_found => { - try store.put(did, terminalState(previous, .complete, active, details.final_host.slice(), Io.Timestamp.now(io, .real).toMicroseconds())); + store.put(did, terminalState(previous, .complete, active, details.final_host.slice(), Io.Timestamp.now(io, .real).toMicroseconds())) catch |err| { + details.local_failure = true; + return err; + }; return .{}; }, .unavailable => { - try store.put(did, terminalState(previous, .unavailable, active, details.final_host.slice(), Io.Timestamp.now(io, .real).toMicroseconds())); + store.put(did, terminalState(previous, .unavailable, active, details.final_host.slice(), Io.Timestamp.now(io, .real).toMicroseconds())) catch |err| { + details.local_failure = true; + return err; + }; return .{}; }, .transient => { @@ -88,6 +101,7 @@ pub fn resyncRepo( }, }; defer prepared_fetch_result.deinit(io); + details.local_failure = true; const now_us: i64 = Io.Timestamp.now(io, .real).toMicroseconds(); diff --git a/src/internal/bootstrap/retry.zig b/src/internal/bootstrap/retry.zig index 138a8fd..6043ef7 100644 --- a/src/internal/bootstrap/retry.zig +++ b/src/internal/bootstrap/retry.zig @@ -110,6 +110,7 @@ pub const Retry = struct { /// one scan over failed repos; also the test seam pub fn runPass(self: *Retry) !Stats { + if (self.stats) |s| _ = s.retry_passes_total.fetchAdd(1, .monotonic); var store = try repo_store.Store.init(self.allocator, self.io, self.meta); defer store.deinit(); @@ -147,7 +148,7 @@ pub const Retry = struct { group.await(self.io) catch |err| return err; if (work.fatal_error) |err| return err; return .{ - .candidates = dids.len, + .candidates = work.candidates.load(.monotonic), .succeeded = work.succeeded.load(.monotonic), .failed = work.failed.load(.monotonic), .deferred = work.deferred.load(.monotonic), @@ -187,6 +188,7 @@ const RetryWork = struct { dids: []const []const u8, gates: *const std.StringHashMap(*HostGate), next: std.atomic.Value(usize) = .init(0), + candidates: std.atomic.Value(u64) = .init(0), succeeded: std.atomic.Value(u64) = .init(0), failed: std.atomic.Value(u64) = .init(0), deferred: std.atomic.Value(u64) = .init(0), @@ -205,11 +207,18 @@ const RetryWork = struct { const now_us: i64 = Io.Timestamp.now(self.retry.io, .real).toMicroseconds(); const state = (try self.store.get(self.retry.allocator, did)) orelse return; defer self.store.freeState(self.retry.allocator, state); - if (state.next_attempt_us > now_us) { + if (!state.active or state.next_attempt_us > now_us) { _ = self.deferred.fetchAdd(1, .monotonic); return; } + _ = self.candidates.fetchAdd(1, .monotonic); + if (self.retry.stats) |s| _ = s.retry_candidates_total.fetchAdd(1, .monotonic); const candidate_host = if (state.host.len > 0) state.host else "unknown"; + if (self.retry.loadHostPark(candidate_host) > Io.Timestamp.now(self.retry.io, .real).toMicroseconds()) { + _ = self.deferred.fetchAdd(1, .monotonic); + if (self.retry.stats) |s| _ = s.retry_skipped_host_parked_total.fetchAdd(1, .monotonic); + return; + } const gate = self.gates.get(candidate_host) orelse return error.MissingHostGate; gate.acquire(self.retry.io, @max(self.retry.host_workers, 1)); defer gate.release(self.retry.io); @@ -218,6 +227,7 @@ const RetryWork = struct { // 429 and parked the bucket while this job was queued at the gate. if (self.retry.loadHostPark(candidate_host) > Io.Timestamp.now(self.retry.io, .real).toMicroseconds()) { _ = self.deferred.fetchAdd(1, .monotonic); + if (self.retry.stats) |s| _ = s.retry_skipped_host_parked_total.fetchAdd(1, .monotonic); return; } @@ -227,6 +237,7 @@ const RetryWork = struct { const scratch_dir = try std.fmt.bufPrint(&scratch_buf, "{s}/repair-scratch/failed-repo", .{self.retry.data_dir}); if (self.retry.stats) |s| _ = s.retry_attempts_total.fetchAdd(1, .monotonic); _ = resync.resyncRepo(self.retry.allocator, self.retry.io, scratch_dir, self.retry.relay_http, did, self.retry.archive, self.store, &drops, &details) catch |err| { + if (details.local_failure) return err; switch (err) { error.OutOfMemory, error.ArchiveAppendFailed => return err, else => {}, @@ -238,7 +249,7 @@ const RetryWork = struct { self.retry.backoffDelayUs(count); const next_attempt = now_us + @as(i64, @intCast(retry_delay_us)); const final_host = if (details.final_host.len > 0) details.final_host.slice() else candidate_host; - try self.store.put(did, .{ + self.store.put(did, .{ .status = .failed, .active = state.active, .rev = state.rev, @@ -257,7 +268,10 @@ const RetryWork = struct { .updated_us = state.updated_us, .record_count = state.record_count, .total_bytes = state.total_bytes, - }); + }) catch |store_err| { + if (self.retry.stats) |s| _ = s.backfill_on_fail_store_errors_total.fetchAdd(1, .monotonic); + return store_err; + }; if (err == error.RateLimited) try self.retry.saveHostPark(final_host, next_attempt); if (self.retry.emit_failure_logs) { log.warn("failed-repo retry {s} via {s}: {s} (attempt {d}, next in {d}s)", .{ @@ -321,8 +335,12 @@ test "runPass with no failed repos is a no-op" { defer meta.deinit(); var r = Retry.init(testing.allocator, io, &a, &meta, data_dir, "http://127.0.0.1:9"); defer r.deinit(); + var instrumentation: metrics.Stats = .{}; + r.stats = &instrumentation; const stats = try r.runPass(); try testing.expectEqual(@as(u64, 0), stats.candidates); + try testing.expectEqual(@as(u64, 1), instrumentation.retry_passes_total.load(.monotonic)); + try testing.expectEqual(@as(u64, 0), instrumentation.retry_candidates_total.load(.monotonic)); } test "host parking is isolated and durable by final PDS host" { @@ -426,11 +444,17 @@ test "redirected 429 persists and parks only the final PDS host" { var retry = Retry.init(testing.allocator, io, &archive, &meta, data_dir, origin_url); defer retry.deinit(); retry.emit_failure_logs = false; + var instrumentation: metrics.Stats = .{}; + retry.stats = &instrumentation; const pass = try retry.runPass(); try testing.expectEqual(@as(u64, 1), pass.candidates); try testing.expectEqual(@as(u64, 1), pass.failed); try testing.expectEqual(@as(u64, 0), pass.succeeded); + try testing.expectEqual(@as(u64, 1), instrumentation.retry_passes_total.load(.monotonic)); + try testing.expectEqual(@as(u64, 1), instrumentation.retry_candidates_total.load(.monotonic)); + try testing.expectEqual(@as(u64, 1), instrumentation.retry_attempts_total.load(.monotonic)); + try testing.expectEqual(@as(u64, 1), instrumentation.retry_failed_total.load(.monotonic)); const state = (try store.get(testing.allocator, "did:plc:park-fixture")).?; defer store.freeState(testing.allocator, state); @@ -442,6 +466,16 @@ test "redirected 429 persists and parks only the final PDS host" { try testing.expectEqual(@as(u32, 1), state.retry_count); try testing.expect(retry.loadHostPark(final_host) > 0); try testing.expectEqual(@as(i64, 0), retry.loadHostPark("predicted-pds.example")); + + // A later due row for that exact final host is rejected from the real + // retry path without issuing another request. The park survives in + // RocksDB, so this also proves the restart-safe half of the contract. + try store.put("did:plc:parked-sibling", .{ .status = .failed, .host = final_host }); + const parked_pass = try retry.runPass(); + try testing.expectEqual(@as(u64, 1), parked_pass.candidates); + try testing.expectEqual(@as(u64, 2), parked_pass.deferred); + try testing.expectEqual(@as(u64, 1), instrumentation.retry_skipped_host_parked_total.load(.monotonic)); + try testing.expectEqual(@as(u64, 1), instrumentation.retry_attempts_total.load(.monotonic)); try origin_future.await(io); try target_future.await(io); } diff --git a/src/internal/metrics.zig b/src/internal/metrics.zig index 3bfd6e8..f8c335b 100644 --- a/src/internal/metrics.zig +++ b/src/internal/metrics.zig @@ -120,9 +120,28 @@ pub const Stats = struct { compaction_watermark: Counter = .init(0), retry_resynced_total: Counter = .init(0), retry_failed_total: Counter = .init(0), + backfill_discovered_total: Counter = .init(0), backfill_completed_total: Counter = .init(0), backfill_failed_total: Counter = .init(0), + backfill_active_flips_total: Counter = .init(0), + backfill_on_fail_store_errors_total: Counter = .init(0), + backfill_handle_repo_duration_count: Counter = .init(0), + backfill_handle_repo_duration_sum_us: Counter = .init(0), + backfill_handle_repo_duration_buckets: [15]Counter = @splat(.init(0)), + backfill_progress_completed: std.atomic.Value(usize) = .init(0), + backfill_completion_queued_total: Counter = .init(0), + backfill_completion_queue_depth: std.atomic.Value(usize) = .init(0), + backfill_completion_durable_batches_total: Counter = .init(0), + backfill_completion_durable_repos_total: Counter = .init(0), + backfill_completion_stage_errors_total: Counter = .init(0), + backfill_completion_queue_wait_count: Counter = .init(0), + backfill_completion_queue_wait_sum_us: Counter = .init(0), + backfill_completion_queue_wait_buckets: [15]Counter = @splat(.init(0)), + backfill_forced_checkpoint_flushes_total: Counter = .init(0), + retry_passes_total: Counter = .init(0), + retry_candidates_total: Counter = .init(0), retry_attempts_total: Counter = .init(0), + retry_skipped_host_parked_total: Counter = .init(0), backfill_workers_active: std.atomic.Value(usize) = .init(0), backfill_queue_depth: std.atomic.Value(usize) = .init(0), backfill_inflight_bytes: std.atomic.Value(usize) = .init(0), @@ -258,8 +277,35 @@ pub const Stats = struct { _ = self.orchestrator_state_duration_buckets[index][bucket_index].fetchAdd(1, .monotonic); } } + + pub fn observeBackfillHandleRepo(self: *Stats, duration_us: u64) void { + observeSlowHistogram( + &self.backfill_handle_repo_duration_count, + &self.backfill_handle_repo_duration_sum_us, + &self.backfill_handle_repo_duration_buckets, + duration_us, + ); + } + + pub fn observeBackfillCompletionWait(self: *Stats, duration_us: u64) void { + observeSlowHistogram( + &self.backfill_completion_queue_wait_count, + &self.backfill_completion_queue_wait_sum_us, + &self.backfill_completion_queue_wait_buckets, + duration_us, + ); + } }; +fn observeSlowHistogram(count: *Counter, sum_us: *Counter, buckets: *[15]Counter, duration_us: u64) void { + _ = count.fetchAdd(1, .monotonic); + _ = sum_us.fetchAdd(duration_us, .monotonic); + inline for (0..15) |index| { + const upper_us = @as(u64, 10_000) << @intCast(index); + if (duration_us <= upper_us) _ = buckets[index].fetchAdd(1, .monotonic); + } +} + fn storeIndex(op: StoreOp, status: StoreStatus) usize { return @intFromEnum(op) * @typeInfo(StoreStatus).@"enum".fields.len + @intFromEnum(status); } @@ -339,6 +385,7 @@ pub fn format( process: process_metrics.Snapshot, data_dir_free_bytes: ?u64, ) []const u8 { + @setEvalBranchQuota(2_000); var w: Io.Writer = .fixed(buf); // canary: proves what binary is running (lessons #8) @@ -719,10 +766,6 @@ pub fn format( \\jetstream_ingest_append_errors_total {d} \\# TYPE jetstream_ingest_active_segment_bytes gauge \\jetstream_ingest_active_segment_bytes {d} - \\# TYPE jetstream_backfill_completed_total counter - \\jetstream_backfill_completed_total {d} - \\# TYPE jetstream_backfill_failed_total counter - \\jetstream_backfill_failed_total {d} \\ , .{ build_options.git_sha, @@ -745,8 +788,6 @@ pub fn format( stats.ingest_segments_rotated_total.load(.monotonic), stats.ingest_append_errors_total.load(.monotonic), stats.ingest_active_segment_bytes.load(.monotonic), - stats.backfill_completed_total.load(.monotonic), - stats.backfill_failed_total.load(.monotonic), }) catch {}; w.print("# TYPE jetstream_subscribe_subscribers gauge\n", .{}) catch {}; @@ -862,12 +903,42 @@ pub fn format( \\# TYPE stream_retry_repos_total counter \\stream_retry_repos_total{{result="resynced"}} {d} \\stream_retry_repos_total{{result="failed"}} {d} + \\# TYPE jetstream_backfill_discovered_total counter + \\jetstream_backfill_discovered_total {d} + \\# TYPE jetstream_backfill_completed_total counter + \\jetstream_backfill_completed_total {d} + \\# TYPE jetstream_backfill_failed_total counter + \\jetstream_backfill_failed_total {d} + \\# TYPE jetstream_backfill_active_flips_total counter + \\jetstream_backfill_active_flips_total {d} + \\# TYPE jetstream_backfill_on_fail_store_errors_total counter + \\jetstream_backfill_on_fail_store_errors_total {d} + \\# TYPE jetstream_backfill_progress_completed gauge + \\jetstream_backfill_progress_completed {d} + \\# TYPE jetstream_backfill_completion_queued_total counter + \\jetstream_backfill_completion_queued_total {d} + \\# TYPE jetstream_backfill_completion_queue_depth gauge + \\jetstream_backfill_completion_queue_depth {d} + \\# TYPE jetstream_backfill_completion_durable_batches_total counter + \\jetstream_backfill_completion_durable_batches_total {d} + \\# TYPE jetstream_backfill_completion_durable_repos_total counter + \\jetstream_backfill_completion_durable_repos_total {d} + \\# TYPE jetstream_backfill_completion_stage_errors_total counter + \\jetstream_backfill_completion_stage_errors_total {d} + \\# TYPE jetstream_backfill_forced_checkpoint_flushes_total counter + \\jetstream_backfill_forced_checkpoint_flushes_total {d} + \\# TYPE jetstream_backfill_failed_repo_retry_passes_total counter + \\jetstream_backfill_failed_repo_retry_passes_total {d} + \\# TYPE jetstream_backfill_failed_repo_retry_candidates_total counter + \\jetstream_backfill_failed_repo_retry_candidates_total {d} \\# TYPE jetstream_backfill_failed_repo_retry_attempts_total counter \\jetstream_backfill_failed_repo_retry_attempts_total {d} \\# TYPE jetstream_backfill_failed_repo_retry_succeeded_total counter \\jetstream_backfill_failed_repo_retry_succeeded_total {d} \\# TYPE jetstream_backfill_failed_repo_retry_failed_total counter \\jetstream_backfill_failed_repo_retry_failed_total {d} + \\# TYPE jetstream_backfill_failed_repo_retry_skipped_host_parked_total counter + \\jetstream_backfill_failed_repo_retry_skipped_host_parked_total {d} \\# TYPE stream_backfill_repos_total counter \\stream_backfill_repos_total{{result="complete"}} {d} \\stream_backfill_repos_total{{result="failed"}} {d} @@ -884,9 +955,24 @@ pub fn format( stats.compaction_watermark.load(.monotonic), stats.retry_resynced_total.load(.monotonic), stats.retry_failed_total.load(.monotonic), + stats.backfill_discovered_total.load(.monotonic), + stats.backfill_completed_total.load(.monotonic), + stats.backfill_failed_total.load(.monotonic), + stats.backfill_active_flips_total.load(.monotonic), + stats.backfill_on_fail_store_errors_total.load(.monotonic), + stats.backfill_progress_completed.load(.monotonic), + stats.backfill_completion_queued_total.load(.monotonic), + stats.backfill_completion_queue_depth.load(.monotonic), + stats.backfill_completion_durable_batches_total.load(.monotonic), + stats.backfill_completion_durable_repos_total.load(.monotonic), + stats.backfill_completion_stage_errors_total.load(.monotonic), + stats.backfill_forced_checkpoint_flushes_total.load(.monotonic), + stats.retry_passes_total.load(.monotonic), + stats.retry_candidates_total.load(.monotonic), stats.retry_attempts_total.load(.monotonic), stats.retry_resynced_total.load(.monotonic), stats.retry_failed_total.load(.monotonic), + stats.retry_skipped_host_parked_total.load(.monotonic), stats.backfill_completed_total.load(.monotonic), stats.backfill_failed_total.load(.monotonic), stats.backfill_workers_active.load(.monotonic), @@ -895,6 +981,44 @@ pub fn format( stats.backfill_inflight_peak_bytes.load(.monotonic), }) catch {}; + const backfill_slow_bounds = [_][]const u8{ + "0.01", "0.02", "0.04", "0.08", "0.16", "0.32", "0.64", "1.28", + "2.56", "5.12", "10.24", "20.48", "40.96", "81.92", "163.84", + }; + const handle_count = stats.backfill_handle_repo_duration_count.load(.monotonic); + const handle_sum_us = stats.backfill_handle_repo_duration_sum_us.load(.monotonic); + w.print( + "# TYPE jetstream_backfill_handle_repo_duration_seconds histogram\n", + .{}, + ) catch {}; + inline for (backfill_slow_bounds, 0..) |upper, index| w.print( + "jetstream_backfill_handle_repo_duration_seconds_bucket{{le=\"{s}\"}} {d}\n", + .{ upper, stats.backfill_handle_repo_duration_buckets[index].load(.monotonic) }, + ) catch {}; + w.print( + \\jetstream_backfill_handle_repo_duration_seconds_bucket{{le="+Inf"}} {d} + \\jetstream_backfill_handle_repo_duration_seconds_sum {d}.{d:0>6} + \\jetstream_backfill_handle_repo_duration_seconds_count {d} + \\ + , .{ handle_count, handle_sum_us / std.time.us_per_s, handle_sum_us % std.time.us_per_s, handle_count }) catch {}; + + const completion_wait_count = stats.backfill_completion_queue_wait_count.load(.monotonic); + const completion_wait_sum_us = stats.backfill_completion_queue_wait_sum_us.load(.monotonic); + w.print( + "# TYPE jetstream_backfill_completion_queue_wait_seconds histogram\n", + .{}, + ) catch {}; + inline for (backfill_slow_bounds, 0..) |upper, index| w.print( + "jetstream_backfill_completion_queue_wait_seconds_bucket{{le=\"{s}\"}} {d}\n", + .{ upper, stats.backfill_completion_queue_wait_buckets[index].load(.monotonic) }, + ) catch {}; + w.print( + \\jetstream_backfill_completion_queue_wait_seconds_bucket{{le="+Inf"}} {d} + \\jetstream_backfill_completion_queue_wait_seconds_sum {d}.{d:0>6} + \\jetstream_backfill_completion_queue_wait_seconds_count {d} + \\ + , .{ completion_wait_count, completion_wait_sum_us / std.time.us_per_s, completion_wait_sum_us % std.time.us_per_s, completion_wait_count }) catch {}; + w.print("# TYPE stream_dropped_events_total counter\n", .{}) catch {}; w.print("# TYPE jetstream_ingest_dropped_events_total counter\n", .{}) catch {}; inline for (comptime std.enums.values(convert.DropReason)) |reason| { @@ -1104,6 +1228,21 @@ test "format renders counters" { _ = stats.replayed_account_events_dropped_total.fetchAdd(3, .monotonic); _ = stats.replayed_identity_events_dropped_total.fetchAdd(4, .monotonic); _ = stats.retry_attempts_total.fetchAdd(5, .monotonic); + _ = stats.backfill_discovered_total.fetchAdd(2, .monotonic); + _ = stats.backfill_active_flips_total.fetchAdd(3, .monotonic); + _ = stats.backfill_on_fail_store_errors_total.fetchAdd(4, .monotonic); + stats.observeBackfillHandleRepo(20_000); + stats.backfill_progress_completed.store(11, .monotonic); + _ = stats.backfill_completion_queued_total.fetchAdd(12, .monotonic); + stats.backfill_completion_queue_depth.store(1, .monotonic); + _ = stats.backfill_completion_durable_batches_total.fetchAdd(6, .monotonic); + _ = stats.backfill_completion_durable_repos_total.fetchAdd(11, .monotonic); + _ = stats.backfill_completion_stage_errors_total.fetchAdd(1, .monotonic); + stats.observeBackfillCompletionWait(40_000); + _ = stats.backfill_forced_checkpoint_flushes_total.fetchAdd(7, .monotonic); + _ = stats.retry_passes_total.fetchAdd(8, .monotonic); + _ = stats.retry_candidates_total.fetchAdd(9, .monotonic); + _ = stats.retry_skipped_host_parked_total.fetchAdd(10, .monotonic); stats.backfill_workers_active.store(4, .monotonic); stats.backfill_queue_depth.store(23, .monotonic); stats.backfill_inflight_bytes.store(4096, .monotonic); @@ -1224,6 +1363,23 @@ test "format renders counters" { try std.testing.expect(std.mem.indexOf(u8, out, "stream_backfill_inflight_peak_bytes 8192") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_build_info{version=") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_completed_total 11") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_discovered_total 2") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_active_flips_total 3") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_on_fail_store_errors_total 4") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_handle_repo_duration_seconds_bucket{le=\"0.02\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_handle_repo_duration_seconds_count 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_progress_completed 11") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_completion_queued_total 12") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_completion_queue_depth 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_completion_durable_batches_total 6") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_completion_durable_repos_total 11") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_completion_stage_errors_total 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_completion_queue_wait_seconds_bucket{le=\"0.04\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_completion_queue_wait_seconds_count 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_forced_checkpoint_flushes_total 7") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_failed_repo_retry_passes_total 8") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_failed_repo_retry_candidates_total 9") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_backfill_failed_repo_retry_skipped_host_parked_total 10") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_ingest_blocks_flushed_total 6") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_ingest_segments_rotated_total 7") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_ingest_append_errors_total 8") != null); diff --git a/tests/oracle.py b/tests/oracle.py index 7989226..a5c2b6e 100644 --- a/tests/oracle.py +++ b/tests/oracle.py @@ -240,6 +240,52 @@ def run_case(point): log.close() +def run_backfill_metrics_receipt(): + """A clean production run must exercise, not merely register, backfill metrics.""" + shutil.rmtree(DATA, ignore_errors=True) + log = open("/tmp/oracle-backfill-metrics.log", "w") + p = subprocess.Popen(ARGS, stdout=log, stderr=log) + try: + wait_serving(p) + samples = metrics_snapshot() + completed = metric(samples, "jetstream_backfill_completed_total") + queued = metric(samples, "jetstream_backfill_completion_queued_total") + durable_repos = metric(samples, "jetstream_backfill_completion_durable_repos_total") + assert metric(samples, "jetstream_backfill_discovered_total") > 0 + assert completed > 0 + assert queued == completed + assert durable_repos == completed + assert metric(samples, "jetstream_backfill_progress_completed") == completed + assert metric(samples, "jetstream_backfill_completion_queue_depth") == 0 + assert metric(samples, "jetstream_backfill_completion_durable_batches_total") > 0 + assert metric(samples, "jetstream_backfill_completion_queue_wait_seconds_count") == completed + assert metric(samples, "jetstream_backfill_handle_repo_duration_seconds_count") > 0 + assert metric(samples, "jetstream_backfill_forced_checkpoint_flushes_total") > 0 + assert metric(samples, "jetstream_backfill_completion_stage_errors_total") == 0 + assert metric(samples, "jetstream_backfill_on_fail_store_errors_total") == 0 + + # Disabled retry still has registered canonical families; it must not + # fabricate work to make the dashboard look active. + for name in ( + "jetstream_backfill_failed_repo_retry_passes_total", + "jetstream_backfill_failed_repo_retry_candidates_total", + "jetstream_backfill_failed_repo_retry_attempts_total", + "jetstream_backfill_failed_repo_retry_succeeded_total", + "jetstream_backfill_failed_repo_retry_failed_total", + "jetstream_backfill_failed_repo_retry_skipped_host_parked_total", + ): + assert (name, ()) in samples, f"missing registered retry family {name}" + assert metric(samples, name) == 0 + finally: + p.send_signal(signal.SIGTERM) + try: + p.wait(timeout=15) + except subprocess.TimeoutExpired: + p.kill() + log.close() + shutil.rmtree(DATA, ignore_errors=True) + + def main(): try: with urllib.request.urlopen(SIM + "/xrpc/com.atproto.sync.listRepos?limit=1", timeout=5) as r: @@ -247,6 +293,10 @@ def main(): except OSError: sys.exit("simulator not reachable on :7777 — run `just simulator` first") + t0 = time.time() + run_backfill_metrics_receipt() + print(f"ok backfill metrics ({time.time() - t0:.1f}s)") + for point in CRASHPOINTS: t0 = time.time() run_case(point) -- 2.51.2