diff --git a/README.md b/README.md index eb26458..1493830 100644 --- a/README.md +++ b/README.md @@ -14,7 +14,9 @@ the live, phone-sized field report is at dictionary-negotiated v2 zstd - a disk-backed historical archive exposed through `listSegments`, `getSegment`, `getBlock`, and `planBackfill`; Stream's sealed `jss` v1 segments are - byte-compatible with upstream Jetstream + byte-compatible with upstream Jetstream. Sealed headers and block indexes + are resident and refresh atomically after seal, compaction, and timestamp + rewrite; cold replay reads only selected compressed blocks - full-network bootstrap with concurrent live capture, deterministic merge, and crash-safe cutover into steady-state ingestion - Sync 1.1 commit verification, including PLC key rotation, MST inversion, @@ -43,7 +45,7 @@ zig build test && zig build # unit tests + exe just simulator # local fake atproto network on :7777 just run-sim # stream against it just e2e # python wire checks -just archive-contract # four archive XRPCs + pinned official Go client +just archive-contract # archive XRPCs + resident manifest + pinned Go client just status-contract # durable host rows + public HTTP view, offline just oracle # crash-matrix: abort at lifecycle seams, assert recovery ``` diff --git a/docs/grafana-dashboard.md b/docs/grafana-dashboard.md index 6042a32..4a547ec 100644 --- a/docs/grafana-dashboard.md +++ b/docs/grafana-dashboard.md @@ -49,6 +49,12 @@ boundary and `xrpc/` label as upstream, all four getBlock outcomes, exact histogram counts, and full-200-only served-byte accounting. Upstream's debug `/metrics` and `/healthz` endpoints and unmatched public requests remain outside the HTTP histogram rather than manufacturing extra handler labels. +The pinned official-client archive receipt also proves the resident manifest: +startup loads exactly the sealed-segment count, archive-to-live replay produces +real block-index hits, unknown-index unit coverage produces misses, and +refcounted readers survive a verified metadata refresh. Compaction and +timestamp rewrite tests prove the resident checksum/index is refreshed before +the rewritten generation is published to later serving work. ## Implemented canonical producers @@ -92,6 +98,10 @@ The following dashboard families already have non-placeholder producers: - `jetstream_getblock_requests_total` - `jetstream_getblock_served_bytes_total` - `jetstream_getblock_duration_seconds` +- `jetstream_manifest_segments_loaded` +- `jetstream_manifest_block_index_cache_hits_total` +- `jetstream_manifest_block_index_cache_misses_total` +- `jetstream_manifest_block_index_load_seconds` - `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 7969b8c..956a158 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. **Remaining 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. 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. **Remaining resident bloom/collection planning, store, orchestrator, backfill, and 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-compaction-spec.md b/docs/upstream-compaction-spec.md index 372e058..6c0b6ca 100644 --- a/docs/upstream-compaction-spec.md +++ b/docs/upstream-compaction-spec.md @@ -205,9 +205,11 @@ observes lag and continues. ## port notes maps naturally onto stream: -- `src/internal/archive.zig` owns the sealed-segment list — it is the - manifest analog; `OnSegmentCompacted` becomes an archive method that - reopens + re-validates one sealed jss file and swaps its resident entry. +- `src/internal/manifest.zig` now owns sorted sealed headers and refcounted + block indexes. Seal publishes a new entry; compaction and timestamp import + verify header/footer checksum and atomically replace it before their durable + lifecycle boundary advances. Blooms and collection indexes remain on-disk + planner work and are not mislabeled as resident yet. - jss sealed segments already carry per-block indexes and blooms, so the window fold (block seq-bound pruning) and candidate-DID bloom prefilter translate directly; keep the 100k-DID narrowing cutoff and the diff --git a/docs/upstream-harness.md b/docs/upstream-harness.md index 1a9fda3..010da96 100644 --- a/docs/upstream-harness.md +++ b/docs/upstream-harness.md @@ -73,6 +73,13 @@ toolchain when necessary, and forces `GOTOOLCHAIN=local`, `GOPROXY=off`, and Zig's checksum-pinned dependencies must already be in its global cache. A missing pin, Go toolchain/module, Python environment, row, sentinel, or cursor transition fails the command rather than weakening the assertion. +The same official-client run now checks the resident manifest after its +archive-to-live cutover: the sealed gauge equals `listSegments`, the load +histogram has one observation per startup segment, replay has produced a real +resident block-index hit, and no unknown-index lookup occurred. Unit adversity +holds an old refcounted index across a verified rewrite and proves cold cursor +translation combines its historical envelope with offsets read from the newly +opened file generation, matching upstream's topology-preserving race contract. ## HTTP and getBlock metrics contract diff --git a/justfile b/justfile index 3d94d65..b44c72e 100644 --- a/justfile +++ b/justfile @@ -71,7 +71,7 @@ http-metrics-contract: if ! kill -0 "$pid" 2>/dev/null; then cat "$LOG" >&2; exit 1; fi sleep 0.1 done - STREAM_BASE_URL="http://127.0.0.1:$PORT" python3 tests/http_metrics_contract.py + STREAM_BASE_URL="http://127.0.0.1:$PORT" STREAM_DATA_DIR="$DATA" python3 tests/http_metrics_contract.py # Seed a real sealed archive, serve it, then exercise all four archive XRPCs # plus whole-segment, block-plan, bounded, and live-cutover paths through the diff --git a/src/internal/archive.zig b/src/internal/archive.zig index 56ae076..2c54e2d 100644 --- a/src/internal/archive.zig +++ b/src/internal/archive.zig @@ -9,13 +9,15 @@ //! tails, never trust an unvalidated record). //! //! durability watermark: committed_seq is the last local seq whose block -//! is fsynced. anything above it exists only in memory. the manifest is -//! deliberately just a directory scan — files are self-describing. +//! is fsynced. anything above it exists only in memory. sealed headers and +//! block indexes are held by the same refreshable resident manifest as +//! upstream; files remain the source of truth for frame bytes. const std = @import("std"); const segment = @import("segment.zig"); const writer_mod = @import("segment_writer.zig"); const metrics = @import("metrics.zig"); +const manifest_mod = @import("manifest.zig"); const Io = std.Io; const Allocator = std.mem.Allocator; @@ -27,6 +29,7 @@ pub const Archive = struct { allocator: Allocator, io: Io, dir: Io.Dir, + manifest: manifest_mod.Manifest, writer: writer_mod.ActiveWriter, file: Io.File, /// bytes of writer.buf already written to disk @@ -76,6 +79,10 @@ pub const Archive = struct { poisoned: bool = false, pub fn init(allocator: Allocator, io: Io, data_dir: []const u8) !Archive { + return initWithStats(allocator, io, data_dir, null); + } + + pub fn initWithStats(allocator: Allocator, io: Io, data_dir: []const u8, stats: ?*metrics.Stats) !Archive { var buf: [256]u8 = undefined; const seg_path = try std.fmt.bufPrint(&buf, "{s}/segments", .{data_dir}); var dir = try Io.Dir.cwd().createDirPathOpen(io, seg_path, .{ .open_options = .{ .iterate = true } }); @@ -85,8 +92,10 @@ pub const Archive = struct { .allocator = allocator, .io = io, .dir = dir, + .manifest = undefined, .writer = undefined, .file = undefined, + .stats = stats, }; try a.recover(); // Recovery only advances next_seq from bytes that are already sealed @@ -94,6 +103,8 @@ pub const Archive = struct { // watermark (and /status) falsely reports zero until the first new // block happens to flush after a restart. a.committed_seq.store(a.next_seq - 1, .release); + a.manifest = try manifest_mod.Manifest.init(allocator, io, dir, stats); + errdefer a.manifest.deinit(); try a.openNextSegment(); return a; } @@ -101,6 +112,7 @@ pub const Archive = struct { pub fn deinit(self: *Archive) void { self.writer.deinit(); self.file.close(self.io); + self.manifest.deinit(); self.dir.close(self.io); } @@ -314,6 +326,7 @@ pub const Archive = struct { try self.file.sync(self.io); try self.file.writePositionalAll(self.io, all[0..segment.header_size], 0); try self.file.sync(self.io); + try self.manifest.refresh(self.dir, self.seg_index, false); self.file.close(self.io); self.writer.deinit(); if (self.stats) |stats| { @@ -349,6 +362,7 @@ pub const Archive = struct { try self.file.sync(self.io); try self.file.writePositionalAll(self.io, all[0..segment.header_size], 0); try self.file.sync(self.io); + try self.manifest.refresh(self.dir, self.seg_index, false); self.file.close(self.io); if (self.stats) |stats| { const elapsed_us = seal_start.durationTo(Io.Timestamp.now(self.io, .awake)).toMicroseconds(); @@ -356,6 +370,7 @@ pub const Archive = struct { } } self.writer.deinit(); + self.manifest.deinit(); self.dir.close(self.io); } @@ -371,24 +386,8 @@ pub const Archive = struct { } }; -pub fn formatSegmentName(buf: []u8, index: u64) []const u8 { - // seg_<10-char zero-padded base36>.jss (lowercase), matching upstream - var digits: [10]u8 = @splat('0'); - var v = index; - var i: usize = 10; - while (v > 0) : (v /= 36) { - i -= 1; - digits[i] = std.fmt.digitToChar(@intCast(v % 36), .lower); - } - return std.fmt.bufPrint(buf, "seg_{s}.jss", .{digits}) catch unreachable; -} - -pub fn parseSegmentIndex(name: []const u8) ?u64 { - if (!std.mem.startsWith(u8, name, "seg_") or !std.mem.endsWith(u8, name, ".jss")) return null; - const digits = name[4 .. name.len - 4]; - if (digits.len != 10) return null; - return std.fmt.parseInt(u64, digits, 36) catch null; -} +pub const formatSegmentName = manifest_mod.formatSegmentName; +pub const parseSegmentIndex = manifest_mod.parseSegmentIndex; // === tests === diff --git a/src/internal/cold.zig b/src/internal/cold.zig index 48586f2..e88f708 100644 --- a/src/internal/cold.zig +++ b/src/internal/cold.zig @@ -8,15 +8,15 @@ //! stops at `until_us` — the hot tail's floor — where the caller hands off //! to the live loop. overlap at the boundary is at-least-once, per contract. //! -//! v0 scope: sealed segments only. the flushed-but-unsealed active segment -//! is not read; the hot tail's 256 MiB budget covers that window in -//! practice. revisit alongside compaction. +//! sealed segments only: resident metadata covers finalized files; the hot +//! tail owns the flushed-but-unsealed active window. const std = @import("std"); const zat = @import("zat"); const filter_mod = @import("filter.zig"); const segment = @import("segment.zig"); const archive_mod = @import("archive.zig"); +const rewrite_mod = @import("segment_rewrite.zig"); const wire = @import("wire.zig"); const tail = @import("tail.zig"); @@ -32,77 +32,54 @@ pub const ResolvedCursor = struct { clamped: bool = false, }; -const SegmentMeta = struct { - name: []u8, - header: segment.Header, -}; - /// Translate a witnessed-at timestamp to the first archived sequence at or -/// after it. This is the pre-upgrade cursor path: bounded positional reads, -/// no whole-segment allocation, and no silent corruption skipping. +/// after it. Segment bounds and block indexes are resident; only the selected +/// compressed block is read from disk. pub fn resolveTimeToSeq( allocator: Allocator, io: Io, archive: *archive_mod.Archive, time_us: i64, ) !ResolvedCursor { - var metas: std.ArrayList(SegmentMeta) = .empty; - defer { - for (metas.items) |meta| allocator.free(meta.name); - metas.deinit(allocator); - } + const metas = try archive.manifest.all(allocator); + defer allocator.free(metas); - var it = archive.dir.iterate(); - while (try it.next(io)) |dir_entry| { - if (archive_mod.parseSegmentIndex(dir_entry.name) == null) continue; - var file = try archive.dir.openFile(io, dir_entry.name, .{}); - defer file.close(io); - var header_buf: [segment.header_size]u8 = undefined; - _ = try file.readPositionalAll(io, &header_buf, 0); - const header = segment.Header.decode(&header_buf) catch |err| switch (err) { - error.ActiveSegment => continue, - else => return err, - }; - try metas.append(allocator, .{ - .name = try allocator.dupe(u8, dir_entry.name), - .header = header, - }); - } - std.mem.sort(SegmentMeta, metas.items, {}, struct { - fn lessThan(_: void, a: SegmentMeta, b: SegmentMeta) bool { - return std.mem.lessThan(u8, a.name, b.name); - } - }.lessThan); - - if (metas.items.len == 0) return .{ .seq = 1, .clamped = true }; - const first = metas.items[0]; + if (metas.len == 0) return .{ .seq = 1, .clamped = true }; + const first = metas[0]; if (time_us <= first.header.min_witnessed_at) return .{ .seq = @max(first.header.min_seq, 1), .clamped = true }; - var candidate: ?SegmentMeta = null; - for (metas.items) |meta| { + var candidate: ?@TypeOf(first) = null; + for (metas) |meta| { if (meta.header.max_witnessed_at >= time_us) { candidate = meta; break; } } if (candidate == null) return .{ - .seq = metas.items[metas.items.len - 1].header.max_seq +| 1, + .seq = metas[metas.len - 1].header.max_seq +| 1, }; const meta = candidate.?; - var file = try archive.dir.openFile(io, meta.name, .{}); + var name_buffer: [64]u8 = undefined; + const name = archive_mod.formatSegmentName(&name_buffer, meta.idx); + var file = try archive.dir.openFile(io, name, .{}); defer file.close(io); + var blocks = try archive.manifest.blockIndex(meta.idx); + defer blocks.deinit(); + const fresh_header = try segment.readHeaderFile(io, &file); + if (@as(usize, fresh_header.block_count) != blocks.entries.len) return error.BlockTopologyChanged; var lo: u32 = 0; - var hi = meta.header.block_count; + var hi: u32 = @intCast(blocks.entries.len); while (lo < hi) { const mid = lo + (hi - lo) / 2; - const entry = try readBlockIndexEntry(io, &file, meta.header, mid); + const entry = blocks.entries[@intCast(mid)]; if (entry.max_witnessed_at >= time_us) hi = mid else lo = mid + 1; } - if (lo == meta.header.block_count) return .{ .seq = @max(meta.header.min_seq, 1) }; - const entry = try readBlockIndexEntry(io, &file, meta.header, lo); - var block = try segment.readBlockFile(allocator, io, &file, entry); + if (lo == blocks.entries.len) return .{ .seq = @max(meta.header.min_seq, 1) }; + const entry = blocks.entries[@intCast(lo)]; + const fresh_entry = try segment.readBlockIndexFile(io, &file, fresh_header, @intCast(lo)); + var block = try segment.readBlockFile(allocator, io, &file, fresh_entry); defer block.deinit(allocator); for (block.events) |event| { if (event.witnessed_at >= time_us) return .{ .seq = @max(event.seq, 1) }; @@ -110,18 +87,6 @@ pub fn resolveTimeToSeq( return .{ .seq = @max(entry.max_seq, 1) }; } -fn readBlockIndexEntry( - io: Io, - file: *Io.File, - header: segment.Header, - index: u32, -) !segment.BlockIndexEntry { - var entry_buf: [segment.BlockIndexEntry.encoded_size]u8 = undefined; - const offset = header.block_index_offset + @as(u64, index) * segment.BlockIndexEntry.encoded_size; - _ = try file.readPositionalAll(io, &entry_buf, offset); - return segment.BlockIndexEntry.decodePublic(&entry_buf); -} - /// stream all events with witnessed_at in [from_us, until_us) matching /// `filter` to `sink` (fn (ctx, json: []const u8) !void). returns the /// number of frames written. @@ -137,37 +102,30 @@ pub fn replay( from_us: i64, until_us: i64, ) !u64 { - var names: std.ArrayList([]u8) = .empty; - defer { - for (names.items) |n| allocator.free(n); - names.deinit(allocator); - } - var it = archive.dir.iterate(); - while (try it.next(io)) |entry| { - if (archive_mod.parseSegmentIndex(entry.name) == null) continue; - try names.append(allocator, try allocator.dupe(u8, entry.name)); - } - std.mem.sort([]u8, names.items, {}, lessThan); + const summaries = try archive.manifest.all(allocator); + defer allocator.free(summaries); var frames: u64 = 0; - for (names.items) |name| { - const bytes = archive.dir.readFileAlloc(io, name, allocator, .limited(1 << 31)) catch continue; - defer allocator.free(bytes); - var sealed = segment.Sealed.parse(allocator, bytes) catch |err| switch (err) { - error.ActiveSegment => continue, // v0: hot tail covers this window - else => { - log.warn("cold: skipping unreadable segment {s}: {s}", .{ name, @errorName(err) }); - continue; - }, + for (summaries) |summary| { + if (summary.header.max_witnessed_at < from_us) continue; + if (summary.header.min_witnessed_at >= until_us) break; + var blocks = try archive.manifest.blockIndex(summary.idx); + defer blocks.deinit(); + var name_buffer: [64]u8 = undefined; + const name = archive_mod.formatSegmentName(&name_buffer, summary.idx); + var file = archive.dir.openFile(io, name, .{}) catch |err| { + log.warn("cold: open resident segment {s}: {s}", .{ name, @errorName(err) }); + continue; }; - defer sealed.deinit(allocator); - if (sealed.header.max_witnessed_at < from_us) continue; - if (sealed.header.min_witnessed_at >= until_us) break; + defer file.close(io); + const fresh_header = try segment.readHeaderFile(io, &file); + if (@as(usize, fresh_header.block_count) != blocks.entries.len) return error.BlockTopologyChanged; - for (sealed.block_index, 0..) |info, i| { + for (blocks.entries, 0..) |info, block_index| { if (info.max_witnessed_at < from_us) continue; if (info.min_witnessed_at >= until_us) break; - var block = try sealed.readBlock(allocator, i); + const fresh_entry = try segment.readBlockIndexFile(io, &file, fresh_header, block_index); + var block = try segment.readBlockFile(allocator, io, &file, fresh_entry); defer block.deinit(allocator); for (block.events) |ev| { if (ev.witnessed_at < from_us or ev.witnessed_at >= until_us) continue; @@ -229,29 +187,25 @@ pub fn replaySeq( v2: bool, from_seq: u64, ) !u64 { - var names: std.ArrayList([]u8) = .empty; - defer { - for (names.items) |n| allocator.free(n); - names.deinit(allocator); - } - var it = archive.dir.iterate(); - while (try it.next(io)) |entry| { - if (archive_mod.parseSegmentIndex(entry.name) == null) continue; - try names.append(allocator, try allocator.dupe(u8, entry.name)); - } - std.mem.sort([]u8, names.items, {}, lessThan); + const summaries = try archive.manifest.all(allocator); + defer allocator.free(summaries); var frames: u64 = 0; - for (names.items) |name| { - const bytes = archive.dir.readFileAlloc(io, name, allocator, .limited(1 << 31)) catch continue; - defer allocator.free(bytes); - var sealed = segment.Sealed.parse(allocator, bytes) catch continue; - defer sealed.deinit(allocator); - if (sealed.header.max_seq < from_seq) continue; - - for (sealed.block_index, 0..) |info, i| { + for (summaries) |summary| { + if (summary.header.max_seq < from_seq) continue; + var blocks = try archive.manifest.blockIndex(summary.idx); + defer blocks.deinit(); + var name_buffer: [64]u8 = undefined; + const name = archive_mod.formatSegmentName(&name_buffer, summary.idx); + var file = archive.dir.openFile(io, name, .{}) catch continue; + defer file.close(io); + const fresh_header = try segment.readHeaderFile(io, &file); + if (@as(usize, fresh_header.block_count) != blocks.entries.len) return error.BlockTopologyChanged; + + for (blocks.entries, 0..) |info, block_index| { if (info.max_seq < from_seq) continue; - var block = try sealed.readBlock(allocator, i); + const fresh_entry = try segment.readBlockIndexFile(io, &file, fresh_header, block_index); + var block = try segment.readBlockFile(allocator, io, &file, fresh_entry); defer block.deinit(allocator); for (block.events) |ev| { if (ev.seq < from_seq) continue; @@ -363,10 +317,6 @@ pub fn encodeEventV2(alloc: Allocator, ev: segment.Event) !?[]u8 { return try out.toOwnedSlice(); } -fn lessThan(_: void, a: []u8, b: []u8) bool { - return std.mem.lessThan(u8, a, b); -} - /// one archived row → one v1 wire frame (null for rows v1 never emits) pub fn encodeEvent(alloc: Allocator, ev: segment.Event) !?[]u8 { var out: Io.Writer.Allocating = .init(alloc); @@ -590,3 +540,62 @@ test "timestamp cursor surfaces corrupt selected block" { resolveTimeToSeq(testing.allocator, io, &archive, 1_700_000_000_002_000), ); } + +test "cursor combines resident envelopes with fresh generation offsets" { + const testing = std.testing; + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tmp = testing.tmpDir(.{}); + defer tmp.cleanup(); + var path_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const path_len = try tmp.dir.realPath(io, &path_buf); + var archive = try archive_mod.Archive.init(testing.allocator, io, path_buf[0..path_len]); + defer archive.deinit(); + archive.max_events_per_block = 1; + archive.writer.max_events_per_block = 1; + + for (0..3) |index| { + _ = try archive.append(.{ + .seq = 0, + .witnessed_at = 1_700_000_000_001_000 + @as(i64, @intCast(index)) * 1_000, + .indexed_at = 0, + .kind = .create, + .did = "did:plc:generationfixture", + .collection = "app.bsky.feed.post", + .rkey = "3k2abcdefghij", + .rev = "3k2abcdefghij", + .payload = "\xa1\x64text\x62hi", + }, 1_700_000_000_001_000); + } + try archive.rotate(); + var resident = try archive.manifest.blockIndex(0); + defer resident.deinit(); + const old_offset = resident.entries[2].offset; + + const Keep = struct { + fn row(_: void, event: segment.Event) bool { + return event.seq >= 2; + } + }; + const rewritten = try rewrite_mod.rewriteFile( + testing.allocator, + io, + archive.dir, + "seg_0000000000.jss", + {}, + Keep.row, + ); + try testing.expect(rewritten.rewritten); + + var file = try archive.dir.openFile(io, "seg_0000000000.jss", .{}); + defer file.close(io); + const fresh_header = try segment.readHeaderFile(io, &file); + const fresh_third = try segment.readBlockIndexFile(io, &file, fresh_header, 2); + try testing.expect(old_offset != fresh_third.offset); + + // Deliberately do not refresh the manifest: this models an in-flight + // reader that retained the old immutable snapshot across the rename. + const resolved = try resolveTimeToSeq(testing.allocator, io, &archive, 1_700_000_000_003_000); + try testing.expectEqual(@as(u64, 3), resolved.seq); +} diff --git a/src/internal/compact/pass.zig b/src/internal/compact/pass.zig index e168a02..ce437ea 100644 --- a/src/internal/compact/pass.zig +++ b/src/internal/compact/pass.zig @@ -39,6 +39,8 @@ pub const Options = struct { /// snapshot entry cap per chunk; reaching it closes the chunk at the /// current segment's max_seq tombstone_cap: usize = default_tombstone_cap, + on_rewritten_ctx: ?*anyopaque = null, + on_rewritten: ?*const fn (?*anyopaque, u64) anyerror!void = null, }; pub const Stats = struct { @@ -105,6 +107,7 @@ pub fn run( const name = archive_mod.formatSegmentName(&name_buf, m.idx); const r = try rewrite_mod.rewriteFile(allocator, io, seg_dir, name, ctx, decide); if (r.rewritten) { + if (opts.on_rewritten) |refresh| try refresh(opts.on_rewritten_ctx, m.idx); stats.segments_rewritten += 1; stats.rows_dropped += r.dropped; log.info("compaction: rewrote {s} ({d} rows dropped)", .{ name, r.dropped }); diff --git a/src/internal/compact/steady.zig b/src/internal/compact/steady.zig index 098a349..9105175 100644 --- a/src/internal/compact/steady.zig +++ b/src/internal/compact/steady.zig @@ -12,10 +12,10 @@ //! windows from disk. a failed pass logs and retries next tick — it must //! never tear down the daemon. //! -//! no manifest refresh / cache invalidation: stream's serving paths read -//! sealed files from disk per request, and rename is atomic — an -//! in-flight read sees either generation, both valid. revisit if a -//! resident manifest or decoded-block cache is added. +//! every successful rewrite refreshes the resident header/block-index entry +//! before its watermark may commit. In-flight readers retain the old +//! immutable index, but obtain physical offsets from their freshly opened +//! file generation; topology-preserving rewrites make that combination safe. const std = @import("std"); const archive_mod = @import("../archive.zig"); @@ -152,6 +152,8 @@ pub const Compactor = struct { // never race; entries <= W only cover physically-compacted rows) const stats = try pass.run(self.allocator, self.io, seg_dir, &wm, null, .{ .tombstone_cap = self.tombstone_cap, + .on_rewritten_ctx = self.archive, + .on_rewritten = refreshManifest, }); if (wm.load(self.io)) |w| { self.mu.lockUncancelable(self.io); @@ -167,6 +169,11 @@ pub const Compactor = struct { const seg_path = try std.fmt.bufPrint(&buf, "{s}/segments", .{self.data_dir}); return Io.Dir.cwd().openDir(self.io, seg_path, .{ .iterate = true }); } + + fn refreshManifest(ctx: ?*anyopaque, idx: u64) !void { + const archive: *archive_mod.Archive = @ptrCast(@alignCast(ctx.?)); + try archive.manifest.refresh(archive.dir, idx, true); + } }; // === tests === @@ -219,6 +226,11 @@ test "steady: hook feeds live set, runOnce force-rotates, compacts, evicts" { try testing.expectEqual(@as(u64, 3), stats.target_watermark); try testing.expectEqual(@as(u64, 1), stats.segments_rewritten); try testing.expectEqual(@as(u64, 1), stats.rows_dropped); + const resident = a.manifest.segmentByIndex(0).?; + try testing.expectEqual(@as(u32, 2), resident.header.event_count); + var resident_blocks = try a.manifest.blockIndex(0); + defer resident_blocks.deinit(); + try testing.expectEqual(@as(u32, 2), resident_blocks.entries[0].event_count); // eviction cleared the covered entry (delete seq 3 <= watermark 3) try testing.expectEqual(@as(usize, 0), c.live.len()); diff --git a/src/internal/manifest.zig b/src/internal/manifest.zig new file mode 100644 index 0000000..fc70f14 --- /dev/null +++ b/src/internal/manifest.zig @@ -0,0 +1,371 @@ +//! Resident sealed-segment manifest. +//! +//! Upstream Jetstream V2 loads every sealed segment header and block index at +//! startup, then refreshes that immutable metadata after seal or rewrite. +//! Serving code never treats repeated disk reads as a "cache". Stream follows +//! that boundary here: metadata is resident, lookups are synchronized, and a +//! refcounted block-index handle remains valid across concurrent refreshes. + +const std = @import("std"); +const metrics = @import("metrics.zig"); +const segment = @import("segment.zig"); + +const Allocator = std.mem.Allocator; +const Io = std.Io; + +const ResidentBlocks = struct { + allocator: Allocator, + refs: std.atomic.Value(usize) = .init(1), + entries: []segment.BlockIndexEntry, + + fn create(allocator: Allocator, entries: []segment.BlockIndexEntry) !*ResidentBlocks { + const resident = try allocator.create(ResidentBlocks); + resident.* = .{ .allocator = allocator, .entries = entries }; + return resident; + } + + fn retain(self: *ResidentBlocks) void { + _ = self.refs.fetchAdd(1, .monotonic); + } + + fn release(self: *ResidentBlocks) void { + if (self.refs.fetchSub(1, .acq_rel) != 1) return; + self.allocator.free(self.entries); + self.allocator.destroy(self); + } +}; + +pub const BlockIndex = struct { + resident: *ResidentBlocks, + entries: []const segment.BlockIndexEntry, + + pub fn deinit(self: *BlockIndex) void { + self.resident.release(); + self.* = undefined; + } +}; + +pub const Summary = struct { + idx: u64, + size: u64, + header: segment.Header, +}; + +const Metadata = struct { + summary: Summary, + blocks: *ResidentBlocks, + + fn deinit(self: *Metadata) void { + self.blocks.release(); + self.* = undefined; + } +}; + +pub const Manifest = struct { + allocator: Allocator, + io: Io, + mu: Io.Mutex = .init, + segments: std.ArrayList(Metadata) = .empty, + stats: ?*metrics.Stats, + generation: u64 = 0, + + pub fn init(allocator: Allocator, io: Io, dir: Io.Dir, stats: ?*metrics.Stats) !Manifest { + var manifest: Manifest = .{ .allocator = allocator, .io = io, .stats = stats }; + errdefer manifest.deinit(); + + var iterator = dir.iterate(); + while (try iterator.next(io)) |entry| { + const idx = parseSegmentIndex(entry.name) orelse continue; + const started = Io.Timestamp.now(io, .awake); + var metadata = (try loadMetadata(allocator, io, dir, idx, false)) orelse continue; + errdefer metadata.deinit(); + try manifest.segments.append(allocator, metadata); + if (stats) |s| { + const elapsed = started.durationTo(Io.Timestamp.now(io, .awake)).toMicroseconds(); + s.observeManifestLoad(@intCast(@max(0, elapsed))); + } + } + std.mem.sort(Metadata, manifest.segments.items, {}, lessThan); + try validateMonotonic(manifest.segments.items); + manifest.generation = 1; + if (stats) |s| s.manifest_segments_loaded.store(manifest.segments.items.len, .release); + return manifest; + } + + pub fn deinit(self: *Manifest) void { + for (self.segments.items) |*metadata| metadata.deinit(); + self.segments.deinit(self.allocator); + } + + pub fn all(self: *Manifest, allocator: Allocator) ![]Summary { + self.mu.lockUncancelable(self.io); + defer self.mu.unlock(self.io); + const result = try allocator.alloc(Summary, self.segments.items.len); + for (self.segments.items, result) |metadata, *summary| summary.* = metadata.summary; + return result; + } + + pub fn segmentByIndex(self: *Manifest, idx: u64) ?Summary { + self.mu.lockUncancelable(self.io); + defer self.mu.unlock(self.io); + const at = self.find(idx); + if (at >= self.segments.items.len or self.segments.items[at].summary.idx != idx) return null; + return self.segments.items[at].summary; + } + + /// Return a stable view of resident metadata. The handle owns one + /// reference so a concurrent compaction/timestamp refresh may swap the + /// manifest entry without invalidating an in-flight replay. + pub fn blockIndex(self: *Manifest, idx: u64) !BlockIndex { + self.mu.lockUncancelable(self.io); + defer self.mu.unlock(self.io); + const at = self.find(idx); + if (at >= self.segments.items.len or self.segments.items[at].summary.idx != idx) { + if (self.stats) |stats| _ = stats.manifest_block_index_cache_misses_total.fetchAdd(1, .monotonic); + return error.UnknownSegment; + } + if (self.stats) |stats| _ = stats.manifest_block_index_cache_hits_total.fetchAdd(1, .monotonic); + const resident = self.segments.items[at].blocks; + resident.retain(); + return .{ .resident = resident, .entries = resident.entries }; + } + + pub fn refresh(self: *Manifest, dir: Io.Dir, idx: u64, verify_checksum: bool) !void { + const started = Io.Timestamp.now(self.io, .awake); + var next = (try loadMetadata(self.allocator, self.io, dir, idx, verify_checksum)) orelse + return error.ActiveSegment; + errdefer next.deinit(); + + self.mu.lockUncancelable(self.io); + defer self.mu.unlock(self.io); + const at = self.find(idx); + try validateReplacement(self.segments.items, at, next.summary); + if (at < self.segments.items.len and self.segments.items[at].summary.idx == idx) { + var old = self.segments.items[at]; + self.segments.items[at] = next; + old.deinit(); + } else { + try self.segments.insert(self.allocator, at, next); + } + self.generation +%= 1; + if (self.stats) |stats| { + stats.manifest_segments_loaded.store(self.segments.items.len, .release); + const elapsed = started.durationTo(Io.Timestamp.now(self.io, .awake)).toMicroseconds(); + stats.observeManifestLoad(@intCast(@max(0, elapsed))); + } + } + + pub fn currentGeneration(self: *Manifest) u64 { + self.mu.lockUncancelable(self.io); + defer self.mu.unlock(self.io); + return self.generation; + } + + fn find(self: *const Manifest, idx: u64) usize { + return std.sort.lowerBound(Metadata, self.segments.items, idx, struct { + fn order(key: u64, metadata: Metadata) std.math.Order { + return std.math.order(key, metadata.summary.idx); + } + }.order); + } +}; + +fn loadMetadata(allocator: Allocator, io: Io, dir: Io.Dir, idx: u64, verify_checksum: bool) !?Metadata { + var name_buffer: [64]u8 = undefined; + const name = formatSegmentName(&name_buffer, idx); + var file = try dir.openFile(io, name, .{}); + defer file.close(io); + const stat = try file.stat(io); + + var header_bytes: [segment.header_size]u8 = undefined; + const header_read = try file.readPositionalAll(io, &header_bytes, 0); + if (header_read != header_bytes.len) return error.Truncated; + const header = segment.Header.decode(&header_bytes) catch |err| switch (err) { + error.ActiveSegment => return null, + else => return err, + }; + + const index_len_u64 = @as(u64, header.block_count) * segment.BlockIndexEntry.encoded_size; + if (header.block_index_offset + index_len_u64 > stat.size) return error.Truncated; + const index_len: usize = @intCast(index_len_u64); + const encoded = try allocator.alloc(u8, index_len); + defer allocator.free(encoded); + if (try file.readPositionalAll(io, encoded, header.block_index_offset) != index_len) + return error.Truncated; + + const entries = try allocator.alloc(segment.BlockIndexEntry, header.block_count); + errdefer allocator.free(entries); + for (entries, 0..) |*entry, block| { + const start = block * segment.BlockIndexEntry.encoded_size; + entry.* = segment.BlockIndexEntry.decodePublic(encoded[start..][0..segment.BlockIndexEntry.encoded_size]); + if (entry.offset < segment.header_size or + entry.offset + 8 + entry.compressed_size > header.footer_offset) + return error.BadBlockFrame; + } + + if (verify_checksum) { + if (header.footer_offset > stat.size) return error.Truncated; + const footer_len: usize = @intCast(stat.size - header.footer_offset); + const footer = try allocator.alloc(u8, footer_len); + defer allocator.free(footer); + if (try file.readPositionalAll(io, footer, header.footer_offset) != footer.len) + return error.Truncated; + var hasher = std.hash.XxHash3.init(0); + hasher.update(header_bytes[12..]); + hasher.update(footer); + if (hasher.final() != header.checksum) return error.BadChecksum; + } + + const resident = try ResidentBlocks.create(allocator, entries); + return .{ + .summary = .{ .idx = idx, .size = stat.size, .header = header }, + .blocks = resident, + }; +} + +fn lessThan(_: void, left: Metadata, right: Metadata) bool { + return left.summary.idx < right.summary.idx; +} + +fn validateMonotonic(items: []const Metadata) !void { + var prior_max: u64 = 0; + var have_prior = false; + for (items) |metadata| { + const summary = metadata.summary; + if (summary.header.event_count == 0) continue; + if (have_prior and summary.header.min_seq <= prior_max) return error.SegmentSeqOverlap; + prior_max = summary.header.max_seq; + have_prior = true; + } +} + +fn validateReplacement(items: []const Metadata, at: usize, candidate: Summary) !void { + if (candidate.header.event_count == 0) return; + var before = at; + while (before > 0) { + before -= 1; + if (items[before].summary.idx == candidate.idx) continue; + if (items[before].summary.header.event_count == 0) continue; + if (items[before].summary.header.max_seq >= candidate.header.min_seq) + return error.SegmentSeqOverlap; + break; + } + var after = at; + if (after < items.len and items[after].summary.idx == candidate.idx) after += 1; + while (after < items.len) : (after += 1) { + if (items[after].summary.header.event_count == 0) continue; + if (candidate.header.max_seq >= items[after].summary.header.min_seq) + return error.SegmentSeqOverlap; + break; + } +} + +pub fn formatSegmentName(buf: []u8, index: u64) []const u8 { + var digits: [10]u8 = @splat('0'); + var value = index; + var at: usize = digits.len; + while (value > 0) : (value /= 36) { + at -= 1; + digits[at] = std.fmt.digitToChar(@intCast(value % 36), .lower); + } + return std.fmt.bufPrint(buf, "seg_{s}.jss", .{digits}) catch unreachable; +} + +pub fn parseSegmentIndex(name: []const u8) ?u64 { + if (!std.mem.startsWith(u8, name, "seg_") or !std.mem.endsWith(u8, name, ".jss")) return null; + const digits = name[4 .. name.len - 4]; + if (digits.len != 10) return null; + return std.fmt.parseInt(u64, digits, 36) catch null; +} + +test "segment name round trip" { + var buffer: [64]u8 = undefined; + try std.testing.expectEqualStrings("seg_0000000000.jss", formatSegmentName(&buffer, 0)); + try std.testing.expectEqualStrings("seg_000000000z.jss", formatSegmentName(&buffer, 35)); + try std.testing.expectEqual(@as(?u64, 35), parseSegmentIndex(formatSegmentName(&buffer, 35))); + try std.testing.expectEqual(@as(?u64, null), parseSegmentIndex("seg_x.jss")); +} + +const writer_mod = @import("segment_writer.zig"); + +fn testEvent(seq: u64) segment.Event { + return .{ + .seq = seq, + .witnessed_at = @intCast(seq * 10), + .indexed_at = 0, + .kind = .create, + .did = "did:plc:manifesttestaccount", + .collection = "app.bsky.feed.post", + .rkey = "3lmanifesttest", + .rev = "3lmanifestrev", + .payload = "\xa1\x64text\x62hi", + }; +} + +fn writeSegment(allocator: Allocator, io: Io, dir: Io.Dir, idx: u64, seqs: []const u64, seal: bool) !void { + var writer = try writer_mod.ActiveWriter.init(allocator); + defer writer.deinit(); + writer.max_events_per_block = 1; + for (seqs) |seq| try writer.append(testEvent(seq)); + if (seal) try writer.seal(); + var name_buffer: [64]u8 = undefined; + try dir.writeFile(io, .{ .sub_path = formatSegmentName(&name_buffer, idx), .data = writer.bytes() }); +} + +test "resident manifest loads sealed metadata and measures real lookups" { + var threaded: Io.Threaded = .init(std.testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + var dir = try tmp.dir.createDirPathOpen(io, "segments", .{ .open_options = .{ .iterate = true } }); + defer dir.close(io); + try writeSegment(std.testing.allocator, io, dir, 0, &.{ 1, 2 }, true); + try writeSegment(std.testing.allocator, io, dir, 1, &.{3}, false); + + var stats: metrics.Stats = .{}; + var manifest = try Manifest.init(std.testing.allocator, io, dir, &stats); + defer manifest.deinit(); + const summaries = try manifest.all(std.testing.allocator); + defer std.testing.allocator.free(summaries); + try std.testing.expectEqual(@as(usize, 1), summaries.len); + try std.testing.expectEqual(@as(u64, 0), summaries[0].idx); + try std.testing.expectEqual(@as(u32, 2), summaries[0].header.block_count); + try std.testing.expectEqual(@as(usize, 1), stats.manifest_segments_loaded.load(.acquire)); + try std.testing.expectEqual(@as(u64, 1), stats.manifest_block_index_load_count.load(.monotonic)); + + var blocks = try manifest.blockIndex(0); + defer blocks.deinit(); + try std.testing.expectEqual(@as(usize, 2), blocks.entries.len); + try std.testing.expectError(error.UnknownSegment, manifest.blockIndex(42)); + try std.testing.expectEqual(@as(u64, 1), stats.manifest_block_index_cache_hits_total.load(.monotonic)); + try std.testing.expectEqual(@as(u64, 1), stats.manifest_block_index_cache_misses_total.load(.monotonic)); +} + +test "block index handle remains valid across verified refresh" { + var threaded: Io.Threaded = .init(std.testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + var dir = try tmp.dir.createDirPathOpen(io, "segments", .{ .open_options = .{ .iterate = true } }); + defer dir.close(io); + try writeSegment(std.testing.allocator, io, dir, 0, &.{1}, true); + + var stats: metrics.Stats = .{}; + var manifest = try Manifest.init(std.testing.allocator, io, dir, &stats); + defer manifest.deinit(); + var old = try manifest.blockIndex(0); + defer old.deinit(); + try std.testing.expectEqual(@as(usize, 1), old.entries.len); + + try writeSegment(std.testing.allocator, io, dir, 0, &.{ 1, 2 }, true); + try manifest.refresh(dir, 0, true); + var current = try manifest.blockIndex(0); + defer current.deinit(); + try std.testing.expectEqual(@as(usize, 2), current.entries.len); + try std.testing.expectEqual(@as(usize, 1), old.entries.len); + try std.testing.expectEqual(@as(u64, 1), old.entries[0].max_seq); + try std.testing.expectEqual(@as(usize, 1), stats.manifest_segments_loaded.load(.acquire)); + try std.testing.expectEqual(@as(u64, 2), stats.manifest_block_index_load_count.load(.monotonic)); +} diff --git a/src/internal/metrics.zig b/src/internal/metrics.zig index 0516e56..9e6d535 100644 --- a/src/internal/metrics.zig +++ b/src/internal/metrics.zig @@ -54,6 +54,12 @@ pub const Stats = struct { getblock_duration_count: Counter = .init(0), getblock_duration_sum_us: Counter = .init(0), getblock_duration_buckets: [12]Counter = @splat(.init(0)), + manifest_segments_loaded: std.atomic.Value(usize) = .init(0), + manifest_block_index_cache_hits_total: Counter = .init(0), + manifest_block_index_cache_misses_total: Counter = .init(0), + manifest_block_index_load_count: Counter = .init(0), + manifest_block_index_load_sum_us: Counter = .init(0), + manifest_block_index_load_buckets: [8]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), @@ -184,6 +190,16 @@ pub const Stats = struct { _ = self.getblock_duration_buckets[i].fetchAdd(1, .monotonic); } } + + pub fn observeManifestLoad(self: *Stats, duration_us: u64) void { + _ = self.manifest_block_index_load_count.fetchAdd(1, .monotonic); + _ = self.manifest_block_index_load_sum_us.fetchAdd(duration_us, .monotonic); + inline for (0..8) |index| { + const upper_us = @as(u64, 100) << @intCast(index * 2); + if (duration_us <= upper_us) + _ = self.manifest_block_index_load_buckets[index].fetchAdd(1, .monotonic); + } + } }; fn httpIndex(handler: HttpHandler, method: HttpMethod, code: HttpCode) usize { @@ -473,6 +489,34 @@ pub fn format( stats.getblock_duration_count.load(.monotonic), }) catch {}; + const manifest_load_sum_us = stats.manifest_block_index_load_sum_us.load(.monotonic); + const manifest_load_count = stats.manifest_block_index_load_count.load(.monotonic); + w.print( + \\# TYPE jetstream_manifest_segments_loaded gauge + \\jetstream_manifest_segments_loaded {d} + \\# TYPE jetstream_manifest_block_index_cache_hits_total counter + \\jetstream_manifest_block_index_cache_hits_total {d} + \\# TYPE jetstream_manifest_block_index_cache_misses_total counter + \\jetstream_manifest_block_index_cache_misses_total {d} + \\# TYPE jetstream_manifest_block_index_load_seconds histogram + \\jetstream_manifest_block_index_load_seconds_sum {d}.{d:0>6} + \\jetstream_manifest_block_index_load_seconds_count {d} + \\ + , .{ + stats.manifest_segments_loaded.load(.monotonic), + stats.manifest_block_index_cache_hits_total.load(.monotonic), + stats.manifest_block_index_cache_misses_total.load(.monotonic), + manifest_load_sum_us / std.time.us_per_s, + manifest_load_sum_us % std.time.us_per_s, + manifest_load_count, + }) catch {}; + const manifest_load_bounds = [_][]const u8{ "0.0001", "0.0004", "0.0016", "0.0064", "0.0256", "0.1024", "0.4096", "1.6384" }; + inline for (manifest_load_bounds, 0..) |upper, index| w.print( + "jetstream_manifest_block_index_load_seconds_bucket{{le=\"{s}\"}} {d}\n", + .{ upper, stats.manifest_block_index_load_buckets[index].load(.monotonic) }, + ) catch {}; + w.print("jetstream_manifest_block_index_load_seconds_bucket{{le=\"+Inf\"}} {d}\n", .{manifest_load_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. @@ -908,6 +952,10 @@ test "format renders counters" { stats.observeHttp(.xrpc, .get, 200, 12_000); stats.observeGetBlock(200, 321, 800); stats.observeGetBlock(404, 0, 200); + stats.manifest_segments_loaded.store(3, .monotonic); + _ = stats.manifest_block_index_cache_hits_total.fetchAdd(4, .monotonic); + _ = stats.manifest_block_index_cache_misses_total.fetchAdd(5, .monotonic); + stats.observeManifestLoad(400); var buf: [64 * 1024]u8 = undefined; const out = format(&buf, &stats, 2, .{ @@ -958,6 +1006,11 @@ test "format renders counters" { try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_getblock_requests_total{result=\"not_found\"} 1") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_getblock_served_bytes_total 321") != null); try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_getblock_duration_seconds_count 2") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_manifest_segments_loaded 3") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_manifest_block_index_cache_hits_total 4") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_manifest_block_index_cache_misses_total 5") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_manifest_block_index_load_seconds_bucket{le=\"0.0004\"} 1") != null); + try std.testing.expect(std.mem.indexOf(u8, out, "jetstream_manifest_block_index_load_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/segment.zig b/src/internal/segment.zig index b7a1e3d..e925e14 100644 --- a/src/internal/segment.zig +++ b/src/internal/segment.zig @@ -324,6 +324,23 @@ pub fn readBlockFile( return decodeBlock(allocator, buffer); } +pub fn readHeaderFile(io: std.Io, file: *std.Io.File) !Header { + var bytes: [header_size]u8 = undefined; + if (try file.readPositionalAll(io, &bytes, 0) != bytes.len) return error.Truncated; + return Header.decode(&bytes); +} + +pub fn readBlockIndexFile(io: std.Io, file: *std.Io.File, header: Header, index: usize) !BlockIndexEntry { + if (index >= header.block_count) return error.BadBlockFrame; + var bytes: [BlockIndexEntry.encoded_size]u8 = undefined; + const offset = header.block_index_offset + @as(u64, @intCast(index)) * BlockIndexEntry.encoded_size; + if (try file.readPositionalAll(io, &bytes, offset) != bytes.len) return error.Truncated; + const entry = BlockIndexEntry.decodePublic(&bytes); + if (entry.offset < header_size or entry.offset + 8 + entry.compressed_size > header.footer_offset) + return error.BadBlockFrame; + return entry; +} + /// decode a decompressed columnar block. takes ownership of `buffer` (freed /// via Block.deinit; events borrow from it). pub fn decodeBlock(allocator: Allocator, buffer: []u8) Error!Block { diff --git a/src/internal/timestamp/runner.zig b/src/internal/timestamp/runner.zig index 29a8d98..a8705c5 100644 --- a/src/internal/timestamp/runner.zig +++ b/src/internal/timestamp/runner.zig @@ -198,6 +198,7 @@ fn runInner(config: Config, record: *jobs_mod.Record) !void { record.specific_cids_unmatched += applied.specific_cids_unmatched; var rewritten_bytes: u64 = 0; if (patch.patched) { + try config.archive.manifest.refresh(config.archive.dir, item.index, true); record.segments_patched += 1; record.rows_mutated += patch.rows_mutated; var file = try config.archive.dir.openFile(config.io, item.segment_name, .{}); @@ -440,6 +441,12 @@ test "real import pauses after durable patch and resumes after full reopen" { defer testing.allocator.free(bytes); var sealed = segment.Sealed.parse(testing.allocator, bytes) catch continue; defer sealed.deinit(testing.allocator); + const segment_index = archive_mod.parseSegmentIndex(entry.name).?; + const resident = archive.manifest.segmentByIndex(segment_index).?; + try testing.expectEqual(sealed.header.checksum, resident.header.checksum); + var resident_blocks = try archive.manifest.blockIndex(segment_index); + defer resident_blocks.deinit(); + try testing.expectEqual(sealed.block_index.len, resident_blocks.entries.len); for (0..sealed.block_index.len) |block_index| { var block_data = try sealed.readBlock(testing.allocator, block_index); defer block_data.deinit(testing.allocator); diff --git a/src/internal/xrpcapi.zig b/src/internal/xrpcapi.zig index 6ffbb25..84e0475 100644 --- a/src/internal/xrpcapi.zig +++ b/src/internal/xrpcapi.zig @@ -168,14 +168,8 @@ pub const Api = struct { defer arena.deinit(); const alloc = arena.allocator(); - var indices: std.ArrayList(u64) = .empty; - var it = self.archive.dir.iterate(); - while (it.next(self.io) catch null) |entry| { - const idx = archive_mod.parseSegmentIndex(entry.name) orelse continue; - if (after != null and idx <= after.?) continue; - indices.append(alloc, idx) catch return respond.err(500, "InternalError", "oom"); - } - std.mem.sort(u64, indices.items, {}, std.sort.asc(u64)); + const summaries = self.archive.manifest.all(alloc) catch + return respond.err(500, "InternalError", "oom"); var out: Io.Writer.Allocating = .init(alloc); var s: std.json.Stringify = .{ .writer = &out.writer }; @@ -184,11 +178,16 @@ pub const Api = struct { s.beginArray() catch return; var emitted: usize = 0; var last_idx: u64 = 0; - for (indices.items) |idx| { - if (emitted >= limit) break; + var remaining = false; + for (summaries) |meta| { + const idx = meta.idx; + if (after != null and idx <= after.?) continue; + if (emitted >= limit) { + remaining = true; + break; + } var name_buf: [64]u8 = undefined; const name = archive_mod.formatSegmentName(&name_buf, idx); - const meta = self.readSealedMeta(name) orelse continue; // active or unreadable s.beginObject() catch return; s.objectField("name") catch return; s.write(name) catch return; @@ -216,7 +215,7 @@ pub const Api = struct { last_idx = idx; } s.endArray() catch return; - if (emitted >= limit and indices.items.len > limit) { + if (remaining) { s.objectField("cursor") catch return; var cur_buf: [24]u8 = undefined; s.write(std.fmt.bufPrint(&cur_buf, "{d}", .{last_idx}) catch return) catch return; @@ -228,12 +227,22 @@ pub const Api = struct { fn getSegment(self: *Api, respond: anytype, query: []const u8, range_header: ?[]const u8, if_none_match: ?[]const u8, if_range: ?[]const u8) void { const name = queryParam(query, "name") orelse return respond.err(400, "InvalidRequest", "missing name"); - if (archive_mod.parseSegmentIndex(name) == null) + const idx = archive_mod.parseSegmentIndex(name) orelse return respond.err(400, "InvalidRequest", "bad segment name"); - const meta = self.readSealedMeta(name) orelse + _ = self.archive.manifest.segmentByIndex(idx) orelse return respond.err(404, "SegmentNotFound", name); + // Validators, size, and bytes all come from this one freshly-opened + // descriptor. A compaction rename between manifest lookup and open + // can otherwise pair an old ETag with a new generation's bytes. + var file = self.archive.dir.openFile(self.io, name, .{}) catch + return respond.err(500, "InternalError", "failed to open segment"); + defer file.close(self.io); + const stat = file.stat(self.io) catch + return respond.err(500, "InternalError", "failed to stat segment"); + const header = segment.readHeaderFile(self.io, &file) catch + return respond.err(500, "InternalError", "failed to read segment header"); var hex_buf: [16]u8 = undefined; - const etag = checksumHex(&hex_buf, meta.header.checksum); + const etag = checksumHex(&hex_buf, header.checksum); if (if_none_match) |candidate| { if (ifNoneMatchMatches(candidate, etag)) return respond.notModified(etag); } @@ -244,19 +253,13 @@ pub const Api = struct { if (if_range) |candidate| { if (!etagMatches(candidate, etag)) break :range; } - const r = parseRange(raw, meta.size) orelse - return respond.rangeNotSatisfiable(meta.size); - var file = self.archive.dir.openFile(self.io, name, .{}) catch - return respond.err(404, "SegmentNotFound", name); - defer file.close(self.io); - respond.file(self.io, file, etag, r.start, r.end, meta.size, true); + const r = parseRange(raw, stat.size) orelse + return respond.rangeNotSatisfiable(stat.size); + respond.file(self.io, file, etag, r.start, r.end, stat.size, true); return; } - var file = self.archive.dir.openFile(self.io, name, .{}) catch - return respond.err(404, "SegmentNotFound", name); - defer file.close(self.io); - respond.file(self.io, file, etag, 0, meta.size - 1, meta.size, false); + respond.file(self.io, file, etag, 0, stat.size - 1, stat.size, false); } fn getBlock(self: *Api, respond: anytype, query: []const u8, range_header: ?[]const u8, if_none_match: ?[]const u8, if_range: ?[]const u8) void { @@ -268,32 +271,35 @@ pub const Api = struct { break :blk std.fmt.parseInt(u32, raw, 10) catch return respond.err(400, "InvalidRequest", "bad blockIndex"); }; - if (archive_mod.parseSegmentIndex(name) == null) + const segment_idx = archive_mod.parseSegmentIndex(name) orelse return respond.err(400, "InvalidRequest", "bad segment name"); + if (self.archive.manifest.segmentByIndex(segment_idx) == null) + return respond.err(404, "SegmentNotFound", name); var file = self.archive.dir.openFile(self.io, name, .{}) catch - return respond.err(404, "SegmentNotFound", name); + return respond.err(500, "InternalError", "failed to open segment"); defer file.close(self.io); - var header_buf: [segment.header_size]u8 = undefined; - _ = file.readPositionalAll(self.io, &header_buf, 0) catch - return respond.err(404, "SegmentNotFound", name); - const header = segment.Header.decode(&header_buf) catch - return respond.err(404, "SegmentNotFound", name); // active or corrupt + const header = segment.readHeaderFile(self.io, &file) catch + return respond.err(500, "InternalError", "failed to read segment header"); if (block_idx >= header.block_count) return respond.err(404, "BlockNotFound", "blockIndex out of range"); - var entry_buf: [segment.BlockIndexEntry.encoded_size]u8 = undefined; - const entry_off = header.block_index_offset + @as(u64, block_idx) * segment.BlockIndexEntry.encoded_size; - _ = file.readPositionalAll(self.io, &entry_buf, entry_off) catch - return respond.err(404, "BlockNotFound", name); - const entry = segment.BlockIndexEntry.decodePublic(&entry_buf); + const entry = segment.readBlockIndexFile(self.io, &file, header, block_idx) catch + return respond.err(500, "InternalError", "failed to read block"); const frame = self.allocator.alloc(u8, entry.compressed_size) catch return respond.err(500, "InternalError", "oom"); defer self.allocator.free(frame); - _ = file.readPositionalAll(self.io, frame, entry.offset + 8) catch - return respond.err(404, "BlockNotFound", name); + var frame_len_bytes: [8]u8 = undefined; + const frame_len_read = file.readPositionalAll(self.io, &frame_len_bytes, entry.offset) catch + return respond.err(500, "InternalError", "failed to read block"); + if (frame_len_read != frame_len_bytes.len or std.mem.readInt(u64, &frame_len_bytes, .little) != entry.compressed_size) + return respond.err(500, "InternalError", "failed to read block"); + const frame_read = file.readPositionalAll(self.io, frame, entry.offset + 8) catch + return respond.err(500, "InternalError", "failed to read block"); + if (frame_read != frame.len) + return respond.err(500, "InternalError", "failed to read block"); var checksum_buf: [16]u8 = undefined; var etag_buf: [32]u8 = undefined; const etag = std.fmt.bufPrint(&etag_buf, "{s}:{d}", .{ checksumHex(&checksum_buf, header.checksum), block_idx }) catch unreachable; @@ -311,18 +317,6 @@ pub const Api = struct { } respond.raw(200, frame, etag); } - - const SealedMeta = struct { header: segment.Header, size: u64 }; - - fn readSealedMeta(self: *Api, name: []const u8) ?SealedMeta { - var file = self.archive.dir.openFile(self.io, name, .{}) catch return null; - defer file.close(self.io); - var header_buf: [segment.header_size]u8 = undefined; - _ = file.readPositionalAll(self.io, &header_buf, 0) catch return null; - const header = segment.Header.decode(&header_buf) catch return null; // skips active - const st = file.stat(self.io) catch return null; - return .{ .header = header, .size = st.size }; - } }; pub fn checksumHex(buf: *[16]u8, checksum: u64) []const u8 { diff --git a/src/main.zig b/src/main.zig index 98b78ac..dcfaec2 100644 --- a/src/main.zig +++ b/src/main.zig @@ -169,13 +169,12 @@ pub fn main(init: std.process.Init.Minimal) !void { }; defer hub.deinit(); - var archive = try archive_mod.Archive.init(allocator, io, data_dir); + var archive = try archive_mod.Archive.initWithStats(allocator, io, data_dir, &stats); defer { archive.close(); archive.deinit(); } archive.max_segment_bytes = max_segment_bytes; - archive.stats = &stats; var rule_store = try timestamp_rules.Store.open(allocator, io, data_dir); defer rule_store.deinit(); archive.stamper_ctx = &rule_store; diff --git a/src/tests.zig b/src/tests.zig index 2339cae..de3903d 100644 --- a/src/tests.zig +++ b/src/tests.zig @@ -26,6 +26,7 @@ comptime { _ = @import("internal/bootstrap/byte_budget.zig"); _ = @import("internal/bootstrap/lifecycle.zig"); _ = @import("internal/archive.zig"); + _ = @import("internal/manifest.zig"); _ = @import("internal/tombstone.zig"); _ = @import("internal/segment_writer.zig"); _ = @import("internal/segment_rewrite.zig"); diff --git a/tests/archive_client_e2e.py b/tests/archive_client_e2e.py index 984c9a6..743dbfb 100644 --- a/tests/archive_client_e2e.py +++ b/tests/archive_client_e2e.py @@ -16,6 +16,7 @@ import pathlib import platform import re import subprocess +import urllib.request UPSTREAM_PIN = "f6d220f3551efd32bfc1a8757abc93f51a713c12" @@ -146,4 +147,30 @@ assert re.search(r"events=5,000\b", final), final assert re.search(r"last_cursor=5000\b", final), final print("backfill-live-cutover: PASS (5,000 unique rows, cursor 5000)") +# The official client's archive-to-live cutover must traverse the same +# all-resident manifest as upstream, not rescan block indexes from disk. +segments = json.load( + urllib.request.urlopen( + HOST + "/xrpc/network.bsky.jetstream.listSegments", timeout=5 + ) +)["segments"] +metrics = urllib.request.urlopen(HOST + "/metrics", timeout=5).read().decode() + + +def scalar(name): + prefix = name + " " + line = next((line for line in metrics.splitlines() if line.startswith(prefix)), None) + assert line is not None, f"missing metric {name}" + return float(line[len(prefix) :]) + + +assert scalar("jetstream_manifest_segments_loaded") == len(segments) +assert scalar("jetstream_manifest_block_index_load_seconds_count") == len(segments) +assert scalar("jetstream_manifest_block_index_cache_hits_total") > 0 +assert scalar("jetstream_manifest_block_index_cache_misses_total") == 0 +print( + "resident-manifest: PASS " + f"({len(segments)} segments, {int(scalar('jetstream_manifest_block_index_cache_hits_total'))} block-index hits)" +) + print(f"OFFICIAL CLIENT ARCHIVE CONTRACT PASS ({UPSTREAM_PIN[:7]}, {GO})") diff --git a/tests/http_metrics_contract.py b/tests/http_metrics_contract.py index 3ea4b56..64ab3e0 100644 --- a/tests/http_metrics_contract.py +++ b/tests/http_metrics_contract.py @@ -7,6 +7,7 @@ import base64 import hashlib import json import os +import pathlib import re import socket import struct @@ -17,6 +18,7 @@ import urllib.request BASE = os.environ.get("STREAM_BASE_URL", "http://127.0.0.1:6021").rstrip("/") +DATA_DIR = pathlib.Path(os.environ["STREAM_DATA_DIR"]) XRPC = BASE + "/xrpc/network.bsky.jetstream." SAMPLE_RE = re.compile(r"^([a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{([^}]*)\})?\s+([^\s]+)$") LABEL_RE = re.compile(r'(\w+)="([^"]*)"') @@ -183,9 +185,20 @@ def main() -> None: assert delta(before, after, "jetstream_getblock_served_bytes_total") == len(frame) assert delta(before, after, "jetstream_getblock_duration_seconds_count") == 6 + # Match upstream's manifest/file race semantics through the real server: + # a segment which remains known to the resident manifest but disappears + # from disk is an internal storage failure, not SegmentNotFound. + (DATA_DIR / "segments" / segment["name"]).unlink() + segment_path = f"/xrpc/network.bsky.jetstream.getSegment?name={name}" + assert request(segment_path)[0] == 500 + assert request(block_path)[0] == 500 + failed = scrape() + assert delta(after, failed, "jetstream_http_request_duration_seconds_count", handler="xrpc/", method="GET", code="500") == 2 + assert delta(after, failed, "jetstream_getblock_requests_total", result="error") == 1 + print( "HTTP/getBlock metrics PASS: " - f"http={int(total_http_delta)} getBlock=6 full_bytes={len(frame)} websocket=1" + f"http={int(total_http_delta)} getBlock=7 full_bytes={len(frame)} websocket=1 storage_race=500" )