diff --git a/docs/grafana-dashboard.md b/docs/grafana-dashboard.md index de0b972..87182c5 100644 --- a/docs/grafana-dashboard.md +++ b/docs/grafana-dashboard.md @@ -79,6 +79,8 @@ The following dashboard families already have non-placeholder producers: - `jetstream_ingest_readable_log_pinned_bytes` - `jetstream_ingest_readable_log_pinned_overrun_bytes` - `jetstream_ingest_dropped_events_total` +- `jetstream_segment_seal_duration_seconds` +- `jetstream_data_dir_free_bytes` - `jetstream_verifier_failures_total` - `jetstream_backfill_completed_total` - `jetstream_backfill_failed_total` diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 81ecd77..8494917 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. **Remaining disk/API/manifest/store/orchestrator/backfill/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. **Remaining API/manifest/store/orchestrator/backfill/compaction families 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-harness.md b/docs/upstream-harness.md index 9868074..025bf3b 100644 --- a/docs/upstream-harness.md +++ b/docs/upstream-harness.md @@ -30,6 +30,7 @@ driver boots real server against simulator, walks lifecycle gated on durable-app - field limits (upstream columnar format): did ≤65535, collection/rkey/rev ≤255, payload ≤u32 - invariants: fsync data before committing cursor; seq assigned at append under writer mutex, starts at 1 (0 = sentinel); per-DID order preserved; never crash on upstream data / crash loud on own corruption - readable log (hot tail): writer-owned deque, deep-copied entries at seq allocation, 256MiB budget, evicts only below durable watermark (pinned above), notify channel per append, encode-once wire memo per entry shared across fan-out; cursor below floor → cold reader over sealed segments. Stream's tail enforces the same pin: an undurable suffix is retained even when it overruns the byte budget, `publishDurable` releases only the newly durable byte suffix before eviction, and scrape-time readable/pinned/overrun gauges are O(1) accumulators rather than a scan under the hot lock. +- archive health: free space is collected from the filesystem containing the open archive directory on every scrape; segment-seal latency starts before the pending-block flush and is observed only after footer write/fsync plus finalized-header pwrite/fsync succeeds, using upstream's `0.01 × 2^n` bucket boundaries. ## zig oracle v1 (2026-07-12) diff --git a/src/internal/archive.zig b/src/internal/archive.zig index 2827598..56ae076 100644 --- a/src/internal/archive.zig +++ b/src/internal/archive.zig @@ -305,6 +305,7 @@ pub const Archive = struct { } fn rotateLocked(self: *Archive) !void { + const seal_start = Io.Timestamp.now(self.io, .awake); try self.flushBlockLocked(); if (self.writer.block_count == 0) return; // nothing to seal try self.writer.seal(); @@ -315,7 +316,11 @@ pub const Archive = struct { try self.file.sync(self.io); self.file.close(self.io); self.writer.deinit(); - if (self.stats) |stats| _ = stats.ingest_segments_rotated_total.fetchAdd(1, .monotonic); + if (self.stats) |stats| { + const elapsed_us = seal_start.durationTo(Io.Timestamp.now(self.io, .awake)).toMicroseconds(); + stats.observeSegmentSeal(@intCast(@max(0, elapsed_us))); + _ = stats.ingest_segments_rotated_total.fetchAdd(1, .monotonic); + } self.seg_index += 1; try self.openNextSegment(); } @@ -330,6 +335,7 @@ pub const Archive = struct { // Cross the ordinary fsync + metadata hook before sealing. Bypassing // syncToDisk here would make the final bootstrap-live block durable // without its relay cursor and promoted verifier state. + const seal_start = Io.Timestamp.now(self.io, .awake); try self.flushBlockLocked(); if (self.writer.block_count == 0) { var name_buf: [64]u8 = undefined; @@ -344,6 +350,10 @@ pub const Archive = struct { try self.file.writePositionalAll(self.io, all[0..segment.header_size], 0); try self.file.sync(self.io); self.file.close(self.io); + if (self.stats) |stats| { + const elapsed_us = seal_start.durationTo(Io.Timestamp.now(self.io, .awake)).toMicroseconds(); + stats.observeSegmentSeal(@intCast(@max(0, elapsed_us))); + } } self.writer.deinit(); self.dir.close(self.io); @@ -420,10 +430,12 @@ test "archive: append, rotate, recover across restart" { var data_dir_buf: [std.Io.Dir.max_path_bytes]u8 = undefined; const data_dir = try std.fmt.bufPrint(&data_dir_buf, "{s}/arch", .{tmp_path}); + var seal_stats: metrics.Stats = .{}; { var a = try Archive.init(testing.allocator, io, data_dir); defer a.deinit(); + a.stats = &seal_stats; a.max_events_per_block = 10; a.writer.max_events_per_block = 10; a.max_segment_bytes = 2000; // force rotation @@ -435,6 +447,10 @@ test "archive: append, rotate, recover across restart" { } try testing.expect(a.seg_index > 0); // rotated at least once try testing.expect(a.committed_seq.load(.acquire) == 100); + try testing.expectEqual( + seal_stats.ingest_segments_rotated_total.load(.monotonic), + seal_stats.segment_seal_duration_count.load(.monotonic), + ); a.close(); // simulate crash: no seal of the current active segment } diff --git a/src/internal/disk_space.zig b/src/internal/disk_space.zig new file mode 100644 index 0000000..ce06b07 --- /dev/null +++ b/src/internal/disk_space.zig @@ -0,0 +1,23 @@ +//! Scrape-time filesystem capacity for the archive data directory. + +const std = @import("std"); + +const c = @cImport({ + @cInclude("sys/statvfs.h"); +}); + +/// Bytes available to this unprivileged process on the filesystem containing +/// dir. Failure is represented as unavailable so /metrics never lies with a +/// zero that looks like a full disk. +pub fn freeBytes(dir: std.Io.Dir) ?u64 { + var stat: c.struct_statvfs = undefined; + if (c.fstatvfs(dir.handle, &stat) != 0) return null; + return std.math.mul(u64, @intCast(stat.f_bavail), @intCast(stat.f_frsize)) catch null; +} + +test "free bytes comes from the real temporary filesystem" { + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + const free = freeBytes(tmp.dir); + try std.testing.expect(free != null and free.? > 0); +} diff --git a/src/internal/metrics.zig b/src/internal/metrics.zig index 529d280..e5641b3 100644 --- a/src/internal/metrics.zig +++ b/src/internal/metrics.zig @@ -30,6 +30,9 @@ pub const Stats = struct { ingest_segments_rotated_total: Counter = .init(0), ingest_append_errors_total: Counter = .init(0), ingest_active_segment_bytes: std.atomic.Value(usize) = .init(0), + segment_seal_duration_count: Counter = .init(0), + segment_seal_duration_sum_us: Counter = .init(0), + segment_seal_duration_buckets: [15]Counter = @splat(.init(0)), upstream_seq: std.atomic.Value(i64) = .init(0), drops: std.enums.EnumArray(convert.DropReason, Counter) = .initFill(.init(0)), verify_valid: Counter = .init(0), @@ -120,6 +123,16 @@ pub const Stats = struct { if (duration_us <= upper) _ = self.subscribe_cursor_resolve_buckets[i].fetchAdd(1, .monotonic); } } + + pub fn observeSegmentSeal(self: *Stats, duration_us: u64) void { + _ = self.segment_seal_duration_count.fetchAdd(1, .monotonic); + _ = self.segment_seal_duration_sum_us.fetchAdd(duration_us, .monotonic); + inline for (0..15) |i| { + const upper_us = @as(u64, 10_000) << @intCast(i); + if (duration_us <= upper_us) + _ = self.segment_seal_duration_buckets[i].fetchAdd(1, .monotonic); + } + } }; pub const TailGauges = struct { @@ -138,6 +151,7 @@ pub fn format( tail: TailGauges, now_s: i64, process: process_metrics.Snapshot, + data_dir_free_bytes: ?u64, ) []const u8 { var w: Io.Writer = .fixed(buf); @@ -242,6 +256,11 @@ pub fn format( \\process_involuntary_context_switches_total {d} \\ , .{value}) catch {}; + if (data_dir_free_bytes) |value| w.print( + \\# TYPE jetstream_data_dir_free_bytes gauge + \\jetstream_data_dir_free_bytes {d} + \\ + , .{value}) catch {}; w.print( \\# TYPE stream_resync_total counter @@ -273,6 +292,26 @@ pub fn format( stats.verify_queue_dropped_total.load(.monotonic), }) catch {}; + const seal_sum_us = stats.segment_seal_duration_sum_us.load(.monotonic); + const seal_count = stats.segment_seal_duration_count.load(.monotonic); + w.print( + \\# TYPE jetstream_segment_seal_duration_seconds histogram + \\jetstream_segment_seal_duration_seconds_sum {d}.{d:0>6} + \\jetstream_segment_seal_duration_seconds_count {d} + \\ + , .{ seal_sum_us / std.time.us_per_s, seal_sum_us % std.time.us_per_s, seal_count }) catch {}; + const seal_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", + }; + inline for (seal_bounds, 0..) |upper, i| { + w.print("jetstream_segment_seal_duration_seconds_bucket{{le=\"{s}\"}} {d}\n", .{ + upper, + stats.segment_seal_duration_buckets[i].load(.monotonic), + }) catch {}; + } + w.print("jetstream_segment_seal_duration_seconds_bucket{{le=\"+Inf\"}} {d}\n", .{seal_count}) catch {}; + // Canonical upstream names. These are aliases only where Stream owns the // same lifecycle boundary; absent families stay absent so the dashboard // contract cannot be satisfied by placeholder zeroes. @@ -622,6 +661,7 @@ test "format renders counters" { _ = stats.subscribe_adversarial_drops_total.fetchAdd(14, .monotonic); _ = stats.subscribe_clean_disconnects_total.fetchAdd(15, .monotonic); _ = stats.subscribe_options_updates_total.fetchAdd(16, .monotonic); + stats.observeSegmentSeal(20_000); var buf: [16 * 1024]u8 = undefined; const out = format(&buf, &stats, 2, .{ @@ -642,7 +682,7 @@ test "format renders counters" { .major_page_faults = 2, .voluntary_context_switches = 41, .involuntary_context_switches = 3, - }); + }, 123_456); try std.testing.expect(std.mem.indexOf(u8, out, "stream_events_total 7") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_verify_total{result=\"valid\"} 3") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_dropped_events_total{reason=\"invalid_signature\"} 1") != null); @@ -662,6 +702,10 @@ test "format renders counters" { try std.testing.expect(std.mem.indexOf(u8, out, "process_major_page_faults_total 2") != null); try std.testing.expect(std.mem.indexOf(u8, out, "process_voluntary_context_switches_total 41") != null); try std.testing.expect(std.mem.indexOf(u8, out, "process_involuntary_context_switches_total 3") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_data_dir_free_bytes 123456") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_segment_seal_duration_seconds_sum 0.020000") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_segment_seal_duration_seconds_bucket{le=\"0.02\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_segment_seal_duration_seconds_count 1") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_backfill_repos_total{result=\"complete\"} 11") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_backfill_workers_active 4") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_backfill_queue_depth 23") != null); diff --git a/src/internal/server.zig b/src/internal/server.zig index 8c9f36e..de74e3b 100644 --- a/src/internal/server.zig +++ b/src/internal/server.zig @@ -13,6 +13,7 @@ const websocket = @import("websocket"); const zat = @import("zat"); const archive_mod = @import("archive.zig"); const cold = @import("cold.zig"); +const disk_space = @import("disk_space.zig"); const filter_mod = @import("filter.zig"); const homepage = @import("homepage.zig"); const metrics = @import("metrics.zig"); @@ -773,7 +774,7 @@ pub const Handler = struct { .pinned_overrun_bytes = g.pinned_overrun_bytes, .base = g.base, .tip = g.tip, - }, now_s, process_metrics.collect(hub.io)); + }, now_s, process_metrics.collect(hub.io), if (hub.archive) |archive| disk_space.freeBytes(archive.dir) else null); var extra_buf: [256]u8 = undefined; var extra: []const u8 = ""; if (hub.pipeline_gauges) |pg| {