diff --git a/docs/grafana-dashboard.md b/docs/grafana-dashboard.md index e03f54a..de0b972 100644 --- a/docs/grafana-dashboard.md +++ b/docs/grafana-dashboard.md @@ -75,6 +75,9 @@ The following dashboard families already have non-placeholder producers: - `jetstream_ingest_segments_rotated_total` - `jetstream_ingest_append_errors_total` - `jetstream_ingest_active_segment_bytes` +- `jetstream_ingest_readable_log_bytes` +- `jetstream_ingest_readable_log_pinned_bytes` +- `jetstream_ingest_readable_log_pinned_overrun_bytes` - `jetstream_ingest_dropped_events_total` - `jetstream_verifier_failures_total` - `jetstream_backfill_completed_total` diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 40c5960..81ecd77 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. **Remaining storage/API/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. **Remaining disk/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 bbd5bac..9868074 100644 --- a/docs/upstream-harness.md +++ b/docs/upstream-harness.md @@ -29,7 +29,7 @@ driver boots real server against simulator, walks lifecycle gated on durable-app - validation gate (drop, count, never crash): non-TID rev → whole event (no log line — hostile input must not drive log volume); invalid NSID/rkey → that op; missing CAR block for create/update → that op (partial CARs are spec-legal; survivors still archived). metrics: dropped_events_total{reason=invalid_rev|invalid_collection|invalid_rkey|field_too_long|missing_block} - 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 +- 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. ## zig oracle v1 (2026-07-12) diff --git a/src/internal/metrics.zig b/src/internal/metrics.zig index 90b07a4..529d280 100644 --- a/src/internal/metrics.zig +++ b/src/internal/metrics.zig @@ -125,6 +125,8 @@ pub const Stats = struct { pub const TailGauges = struct { entries: usize, bytes: usize, + pinned_bytes: usize, + pinned_overrun_bytes: usize, base: u64, tip: u64, }; @@ -157,6 +159,12 @@ pub fn format( \\stream_tail_entries {d} \\# TYPE stream_tail_bytes gauge \\stream_tail_bytes {d} + \\# TYPE jetstream_ingest_readable_log_bytes gauge + \\jetstream_ingest_readable_log_bytes {d} + \\# TYPE jetstream_ingest_readable_log_pinned_bytes gauge + \\jetstream_ingest_readable_log_pinned_bytes {d} + \\# TYPE jetstream_ingest_readable_log_pinned_overrun_bytes gauge + \\jetstream_ingest_readable_log_pinned_overrun_bytes {d} \\# TYPE stream_tail_base gauge \\stream_tail_base {d} \\# TYPE stream_tail_tip gauge @@ -172,6 +180,9 @@ pub fn format( subscribers_active, tail.entries, tail.bytes, + tail.bytes, + tail.pinned_bytes, + tail.pinned_overrun_bytes, tail.base, tail.tip, }) catch {}; @@ -613,7 +624,14 @@ test "format renders counters" { _ = stats.subscribe_options_updates_total.fetchAdd(16, .monotonic); var buf: [16 * 1024]u8 = undefined; - const out = format(&buf, &stats, 2, .{ .entries = 5, .bytes = 100, .base = 1, .tip = 6 }, 60, .{ + const out = format(&buf, &stats, 2, .{ + .entries = 5, + .bytes = 100, + .pinned_bytes = 70, + .pinned_overrun_bytes = 6, + .base = 1, + .tip = 6, + }, 60, .{ .cpu_us = 1_250_000, .resident_memory_bytes = 4096, .virtual_memory_bytes = 8192, @@ -630,6 +648,9 @@ test "format renders counters" { try std.testing.expect(std.mem.indexOf(u8, out, "stream_dropped_events_total{reason=\"invalid_signature\"} 1") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_build_info{git_sha=") != null); try std.testing.expect(std.mem.indexOf(u8, out, "stream_subscribers_active 2") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_ingest_readable_log_bytes 100") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_ingest_readable_log_pinned_bytes 70") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_ingest_readable_log_pinned_overrun_bytes 6") != null); try std.testing.expect(std.mem.indexOf(u8, out, "process_start_time_seconds 0") != null); try std.testing.expect(std.mem.indexOf(u8, out, "process_cpu_seconds_total 1.250000") != null); try std.testing.expect(std.mem.indexOf(u8, out, "process_resident_memory_bytes 4096") != null); diff --git a/src/internal/server.zig b/src/internal/server.zig index 10f151d..8c9f36e 100644 --- a/src/internal/server.zig +++ b/src/internal/server.zig @@ -769,6 +769,8 @@ pub const Handler = struct { const out = metrics.format(&buf, hub.stats, hub.active.load(.monotonic), .{ .entries = g.entries, .bytes = g.bytes, + .pinned_bytes = g.pinned_bytes, + .pinned_overrun_bytes = g.pinned_overrun_bytes, .base = g.base, .tip = g.tip, }, now_s, process_metrics.collect(hub.io)); diff --git a/src/internal/tail.zig b/src/internal/tail.zig index f0e87a3..0955f4f 100644 --- a/src/internal/tail.zig +++ b/src/internal/tail.zig @@ -46,6 +46,9 @@ pub const Tail = struct { /// absolute index of entries.items[0] base: u64 = 0, cur_bytes: usize = 0, + /// Bytes belonging to seq-bearing rows newer than published_seq. Kept as + /// an accumulator so a Prometheus scrape stays O(1) under the tail lock. + pinned_bytes: usize = 0, max_bytes: usize, closed: bool = false, /// When enabled, archived entries remain private until their segment @@ -62,13 +65,33 @@ pub const Tail = struct { defer self.mutex.unlock(self.io); self.durability_gated = true; self.published_seq = durable_seq; + self.pinned_bytes = 0; + for (self.entries.items) |entry| { + if (entry.seq != 0 and entry.seq > durable_seq) + self.pinned_bytes += entryBytes(entry); + } + self.evictLocked(); } /// Publish frames only after Archive's fsync + metadata transaction. pub fn publishDurable(self: *Tail, durable_seq: u64) void { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); - if (durable_seq > self.published_seq) self.published_seq = durable_seq; + const old_durable = self.published_seq; + if (durable_seq > old_durable) { + // Seq-bearing rows are ordered. Walk only the undurable suffix, + // not the entire retained log, then stop at the prior watermark. + var i = self.entries.items.len; + while (i > 0) { + i -= 1; + const entry = self.entries.items[i]; + if (entry.seq == 0 or entry.seq > durable_seq) continue; + if (entry.seq <= old_durable) break; + self.pinned_bytes -= entryBytes(entry); + } + self.published_seq = durable_seq; + } + self.evictLocked(); self.cond.broadcast(self.io); } @@ -92,22 +115,41 @@ pub const Tail = struct { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); try self.entries.append(self.allocator, entry); - self.cur_bytes += entryBytes(entry); - while (self.cur_bytes > self.max_bytes and self.entries.items.len > 1) { + const entry_bytes = entryBytes(entry); + self.cur_bytes += entry_bytes; + if (self.durability_gated and entry.seq != 0 and entry.seq > self.published_seq) + self.pinned_bytes += entry_bytes; + self.evictLocked(); + if (!self.durability_gated or entry.seq == 0 or entry.seq <= self.published_seq) + self.cond.broadcast(self.io); + } + + fn evictLocked(self: *Tail) void { + while (self.cur_bytes > self.max_bytes and self.entries.items.len > 0) { + const oldest = self.entries.items[0]; + // Rows newer than the archive watermark are the only copy a hot + // subscriber can read. Match upstream's ReadableLog: exceed the + // retention budget rather than evicting an undurable row. + if (self.durability_gated and oldest.seq != 0 and oldest.seq > self.published_seq) break; const evicted = self.entries.orderedRemove(0); self.cur_bytes -= entryBytes(evicted); evicted.deinit(self.allocator); self.base += 1; } - if (!self.durability_gated or entry.seq == 0 or entry.seq <= self.published_seq) - self.cond.broadcast(self.io); } fn entryBytes(e: Entry) usize { return e.json.len + e.json_v2.len + e.did.len + e.collection.len + 64; } - pub const Gauges = struct { entries: usize, bytes: usize, base: u64, tip: u64 }; + pub const Gauges = struct { + entries: usize, + bytes: usize, + pinned_bytes: usize, + pinned_overrun_bytes: usize, + base: u64, + tip: u64, + }; pub fn gauges(self: *Tail) Gauges { self.mutex.lockUncancelable(self.io); @@ -115,6 +157,8 @@ pub const Tail = struct { return .{ .entries = self.entries.items.len, .bytes = self.cur_bytes, + .pinned_bytes = self.pinned_bytes, + .pinned_overrun_bytes = self.pinned_bytes -| self.max_bytes, .base = self.base, .tip = self.base + self.entries.items.len, }; @@ -308,6 +352,37 @@ test "durability gate hides appended frames until publication" { try testing.expectEqualStrings("{\"n\":1}", json); } +test "undurable rows pin readable log beyond budget until publication" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + var t = Tail.init(testing.allocator, threaded.io(), 1); + defer t.deinit(); + t.enableDurabilityGate(0); + + try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 100, 1, "{\"n\":1}", "{\"v2\":1}", .none); + try t.append(.commit, "did:plc:a", "app.bsky.feed.post", 200, 2, "{\"n\":2}", "{\"v2\":2}", .none); + const pinned = t.gauges(); + try testing.expectEqual(@as(usize, 2), pinned.entries); + try testing.expectEqual(pinned.bytes, pinned.pinned_bytes); + try testing.expectEqual(pinned.pinned_bytes - 1, pinned.pinned_overrun_bytes); + try testing.expectEqual(@as(u64, 0), pinned.base); + + t.publishDurable(1); + const one_pinned = t.gauges(); + try testing.expectEqual(@as(usize, 1), one_pinned.entries); + try testing.expectEqual(@as(u64, 1), one_pinned.base); + try testing.expectEqual(one_pinned.bytes, one_pinned.pinned_bytes); + + t.publishDurable(2); + const drained = t.gauges(); + try testing.expectEqual(@as(usize, 0), drained.entries); + try testing.expectEqual(@as(usize, 0), drained.bytes); + try testing.expectEqual(@as(usize, 0), drained.pinned_bytes); + try testing.expectEqual(@as(usize, 0), drained.pinned_overrun_bytes); + try testing.expectEqual(@as(u64, 2), drained.base); + try testing.expectEqual(@as(u64, 2), drained.tip); +} + test "v1 observes skipped rows while v2 receives them" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit();