diff --git a/README.md b/README.md index b500823..5ec30b9 100644 --- a/README.md +++ b/README.md @@ -24,8 +24,10 @@ experiment. - a whole-network bootstrap implementation with concurrent live capture and deterministic merge. Repository completion is coupled to ordinary archive writer durability independently of the 100,000-entry dispatch/checkpoint - batch; several recovery, pending-repair, and live-ingest paths remain blocked by - the deployment-gate audit + batch; cold replay, pending repair, and live-ingest paths remain blocked by + the deployment-gate audit. Startup resumes the existing active segment, + truncating and fsyncing only a framing-torn suffix while failing closed on + complete-frame corruption and sealed high-water checksum failure - Sync 1.1 commit verification, including PLC key rotation, MST inversion, op-CID checks, replay protection, durable repository-chain state, and transparent authenticated whole-repo repair with ordered pending replay; diff --git a/docs/bootstrap-semantic-parity.md b/docs/bootstrap-semantic-parity.md index 82d03f9..6b8c79e 100644 --- a/docs/bootstrap-semantic-parity.md +++ b/docs/bootstrap-semantic-parity.md @@ -50,7 +50,7 @@ live-encoding, and reconnect defects still prevent another deployment. | Per-repository completion durability | **verified** | Stream now ports upstream's watermark queue: data-bearing and genuinely empty successful CARs queue independently, duplicate DIDs replace in place, RepoNotFound follows upstream's immediate terminal OnFail path, and captured completion time survives deferred durability. | Re-run the process-level oracle before artifact admission; the production mechanism and focused boundary are established here. | | Slow or hung straggler isolation | **verified** | A production-boundary test leaves one repository without a worker result, crosses a writer durability boundary, proves its completed sibling durable while the listRepos cursor is absent, fully reopens RocksDB, and proves only the interrupted repository moves to pending. | The admission oracle must scale this invariant beyond a two-repository fixture and inject an actual process kill. | | Completion metrics | **verified** | Queue count/depth change at queue/replacement; durable batch/repo counters change only after the synced repo-state batch; queue wait uses the captured completion time; RepoNotFound bypasses those queue metrics as upstream does. | Prometheus counters are process-local by definition. Restart-stable progress comes from durable repository states, not counter continuity. | -| Resource exhaustion classification | **partial** | Preparation allocation failure is treated as lifecycle-fatal rather than a durable “repo too large” result. | Fault coverage does not establish cleanup and restart behavior at each preparation/emission ownership transfer. Archive active-tail recovery also mishandles iterator allocation failure. | +| Resource exhaustion classification | **partial** | Preparation allocation failure is treated as lifecycle-fatal rather than a durable “repo too large” result. Archive startup now propagates iterator allocation failure rather than treating it as a torn tail. | Fault coverage does not establish cleanup and restart behavior at each preparation/emission ownership transfer. | ## Cursors and lifecycle transitions @@ -75,7 +75,7 @@ live-encoding, and reconnect defects still prevent another deployment. | Merge row filtering | **verified** | Rows are dropped only when the repository is complete and the source revision is at/below its backfill watermark; missing lookups are cached. Kept/dropped fixtures exercise this. | This predicate proof does not establish source-file recovery behavior. | | Latest revision refresh | **partial** | Kept rev-bearing rows update latest revision in the same batch as the source cursor while preserving the backfill revision. | Missing repo rows are skipped; this matches the intended defensive behavior, but corrupt cursor state can replay updates and rows. | | Pending-repository pass | **blocked** | Pending repositories are repaired after captured live rows are merged, preserving intended row order. | Stream materializes pending DIDs and processes them through a separate path rather than upstream's bounded retry runner with the same global/per-host gates and retry semantics. | -| Merge-tail compaction and manifest reconcile | **partial** | Destination sealing precedes delete/update compaction; manifest reconciliation happens before serving is enabled. Physical rewrite tests cover named write failures; compaction watermark and merge source admission now fail closed on metadata/filesystem errors. | Archive startup/recovery and exact cutover failure coverage remain independent blockers. | +| Merge-tail compaction and manifest reconcile | **partial** | Destination sealing precedes delete/update compaction; manifest reconciliation happens before serving is enabled. Physical rewrite tests cover named write failures; compaction watermark, merge source admission, and active-tail startup recovery fail closed. | Exact cutover and post-startup rewrite failure coverage remain independent blockers. | | Post-bootstrap discovery | **verified** | It resumes from the last non-empty bootstrap cursor, writes every previously unknown active or inactive DID as failed while preserving the relay's active flag, follows dynamically owned cursors, rejects cursor loops, and propagates relay/store errors before cleanup. A real-process relay fixture returns an inactive DID, a 4 KiB cursor, and a second DID; both durable account rows are inspected after steady-state admission. | This row proves discovery completeness and failure behavior; the downstream retry runner remains independently blocked. | | Cleanup durability | **verified** | On the successful path, the backfill tree is removed, the data directory is synced, and merge/discovery cursors are deleted in a synced RocksDB batch afterward. Restart-after-cleanup repeats the directory sync, and the source-existence guard fails closed on non-not-found errors. | None known for this cleanup ordering invariant. | diff --git a/docs/configuration-parity.md b/docs/configuration-parity.md index 6512a3f..fa5face 100644 --- a/docs/configuration-parity.md +++ b/docs/configuration-parity.md @@ -20,7 +20,7 @@ semantically correct. Status vocabulary is defined in | Complete upstream flag inventory | **unverified** | The environment contract contains a hand-maintained pinned source map and rejects unknown `JETSTREAM_*` names. | There is no generated comparison proving every current upstream flag is represented here. Before another admission, generate the inventory from the pinned Go command/options definitions and fail on an unclassified addition/removal. | | Relay URL | **partial** | `--relay-url` reaches the WebSocket consumer and relay-fronted HTTP listRepos/getRepo/repair clients; CLI and environment precedence are tested. | Correct wiring still feeds a blocked Zat reconnect loop. | | PLC URL | **verified** | `--plc-url` configures the real DID resolver used by live verification and selected-repo metadata; loopback resolution tests observe requests. | This row does not prove every verification/repair outcome. | -| Data directory | **partial** | `--data-dir` owns JSS, scratch, metadata, bootstrap-live, import, and temporary trees in production-process tests. | Startup recovery and several metadata reads are blocked; path ownership alone is verified. | +| Data directory | **partial** | `--data-dir` owns JSS, scratch, metadata, bootstrap-live, import, and temporary trees in production-process tests. Active-tail startup recovery now resumes in place and fails closed on corruption. | Cold replay, timestamp-import visibility, and retry metadata still have independent blocked paths; path ownership alone does not close them. | | `--max-backfill-repos` | **partial** | Canonical spelling is parsed, limits dispatch, avoids committing a whole-network resume cursor, and automatically skips merge discovery. | The lifecycle and completion paths it invokes are blocked. It is a debug mode, not evidence for whole-network parity. | | `--backfill-repos` | **partial** | Canonical spelling and compatibility alias parse; it is mutually exclusive with max-repos, preserves selected order, bypasses listRepos, and enables selected identity resolution. | The exact current artifact still needs the complete selected-repository admission receipt; configuration wiring itself is established. | | Backfill workers | **verified** | Explicit positive values set physical getRepo concurrency; zero/omitted default to 100. Held real requests observe 7/100/100. | The 200-worker production experiment showed that accepting a value does not mean it is efficient. | diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 65cd748..8e10c14 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -7,8 +7,8 @@ shape, or because a test written around Stream's implementation passes. ## Audit basis -- Stream: `a345f4054df568e364aff08a9b1726275fc48e29` plus the discovery - changes and audit updates in this commit +- Stream: `acb10fd46948236c9a3740303a847f95d67cb2f6` plus the archive + recovery changes and audit updates in this commit - upstream Jetstream: `f29815c391fc2644f8a3dd36b899fb3697dd1ea6` - Atmos: `v0.2.14` - Zat: `649dae356c576a205ec97f2b23671af365756b55` @@ -39,11 +39,11 @@ test, and this audit must be reviewed together. **Stream is not at semantic parity and is not admitted for another whole-network experiment.** -Archive recovery, cold replay, pending-repository/retry orchestration, -ordinary live encoding, retry-supervisor health, and the firehose reconnect -loop remain deployment blockers. The completed bootstrap durability, -correctness-metadata, and discovery work does not admit an experiment while -those independent blockers remain. The detailed bootstrap audit is in +Cold replay, pending-repository/retry orchestration, ordinary live encoding, +retry-supervisor health, and the firehose reconnect loop remain deployment +blockers. The completed bootstrap durability, correctness-metadata, discovery, +and archive-recovery work does not admit an experiment while those independent +blockers remain. The detailed bootstrap audit is in [bootstrap-semantic-parity.md](bootstrap-semantic-parity.md). ## Audited checklist @@ -55,7 +55,7 @@ those independent blockers remain. The detailed bootstrap audit is in | Merge and post-bootstrap discovery | **blocked** | Source rows are filtered by the backfill revision; source cursor/latest-revision updates share a synced batch; merge cursor corruption/read errors fail closed; the restart guard treats only `FileNotFound` as cleanup-complete; discovery records active and inactive unknown repositories and follows arbitrary-length cursors with loop detection; successful cleanup is directory-synced. | The pending-repository pass still materializes all DIDs and does not use upstream's bounded retry runner. Latest-revision refresh and source-file failure coverage remain partial. | | Failed-repository healing | **blocked** | A real global/per-host worker gate exists, final redirect hosts are recorded, and retry state is stored. | Pass-level RocksDB/archive failures are logged and hidden from service health. Every failed DID and host gate is materialized in RAM before work. Persisted host parking is not upstream behavior, and park-read errors mean “not parked.” | | JSS sealed format interoperability | **verified** | Upstream-produced sealed fixtures are parsed; Stream-produced sealed files are consumed by the pinned Go reader; header/footer, block index, blooms, collections, compression, and checksums have reciprocal fixtures. | This verifies sealed-format compatibility only. It does not verify startup recovery or replay completeness. | -| Archive startup recovery | **partial** | Torn active tails with ordinary decode termination are truncated to the last complete block and sealed. A truncated sealed segment is subsequently rejected by manifest loading in the current composition. | `Archive.recover` swallows file-read/header errors and converts any active-iterator error, including allocation failure, into end-of-valid-data before rewriting the file. Upstream fails loud on these errors and checksum-verifies a sealed high-water mark. Add fault-specific recovery tests before claiming parity. | +| Archive startup recovery | **verified** | Startup validates the active header, walks complete frames, propagates read/decode/allocation failures, truncates and fsyncs only a framing-torn suffix, reconstructs block/event/sequence state, and resumes the same active file and segment index. Sealed high-water floors are accepted only through the checksum-verifying parser, including the empty-active/lower-sealed case. Physical tests prove repeated same-file restart, exact torn-tail truncation, byte-preserving failure on a complete corrupt frame, and rejection of a forged sealed `max_seq`. | This row establishes startup recovery. Cold replay and post-startup rewrite paths remain separate surfaces. | | `listSegments`, `getSegment`, `getBlock` | **verified** | The pinned Go client and direct HTTP tests cover normal responses, byte identity, checksums/ETags, conditional requests, ranges, cache headers, and storage-error responses through the production server. | The claim is limited to these three methods and the tested file states. | | `planBackfill` | **blocked** | Normal DID/collection/sequence filtering, whole-segment versus block mode, pagination, and configured limits match the pinned client fixtures. | `matchBlocks` errors are handled with `catch continue`, creating a successful plan with a false negative. Upstream's planner returns the error. Add an injected allocation/internal-error contract. | | Subscribe v1/v2 handshake and wire formats | **partial** | Cursor parsing, v1/v2 filter shapes, WebSocket framing, v1 deflate, v2 dictionary zstd, dictionary validators, size caps, pings, write deadlines, and shutdown close frames have production-path tests. | These tests do not establish gap-free cold replay or remote-data encoding isolation. | diff --git a/src/internal/archive.zig b/src/internal/archive.zig index e7ecc09..ac1d0be 100644 --- a/src/internal/archive.zig +++ b/src/internal/archive.zig @@ -272,7 +272,11 @@ pub const Archive = struct { .file = undefined, .stats = stats, }; - try a.recover(); + const resumed_active = try a.recover(); + errdefer if (resumed_active) { + a.writer.deinit(); + a.file.close(io); + }; // Recovery only advances next_seq from bytes that are already sealed // and fsynced. Publish that durable floor immediately; otherwise the // watermark (and /status) falsely reports zero until the first new @@ -280,7 +284,7 @@ pub const Archive = struct { 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(); + if (!resumed_active) try a.openNextSegment(); return a; } @@ -322,71 +326,104 @@ pub const Archive = struct { } } - /// scan the directory (the manifest), find the next index and seq, and - /// finish any interrupted active segment: walk its frames, truncate the - /// torn tail, seal in place. - fn recover(self: *Archive) !void { + /// Scan the directory, recover the durable seq floor, and resume the + /// highest active segment in place. Returns true when `writer` and `file` + /// were initialized from that active segment. + fn recover(self: *Archive) !bool { var max_index: ?u64 = null; var it = self.dir.iterate(); while (try it.next(self.io)) |entry| { const idx = parseSegmentIndex(entry.name) orelse continue; if (max_index == null or idx > max_index.?) max_index = idx; } - var candidate = max_index orelse return; - self.seg_index = candidate + 1; + const candidate = max_index orelse return false; + var name_buf: [64]u8 = undefined; + const name = formatSegmentName(&name_buf, candidate); + const bytes = try self.dir.readFileAlloc(self.io, name, self.allocator, .limited(1 << 31)); + defer self.allocator.free(bytes); + if (bytes.len < segment.header_size) return error.Truncated; + + if (segment.Header.decode(bytes)) |_| { + const max_seq = try self.sealedMaxSeq(candidate, bytes); + self.next_seq = @max(self.next_seq, max_seq + 1); + self.seg_index = candidate + 1; + return false; + } else |err| switch (err) { + error.ActiveSegment => {}, + else => return err, + } + + var valid_end: usize = segment.header_size; + var max_seq: u64 = 0; + var event_count: u64 = 0; + var block_count: u32 = 0; + var torn = false; + var walk: writer_mod.ActiveIterator = .{ .bytes = bytes }; while (true) { - var name_buf: [64]u8 = undefined; - const name = formatSegmentName(&name_buf, candidate); - const bytes = self.dir.readFileAlloc(self.io, name, self.allocator, .limited(1 << 31)) catch return; - defer self.allocator.free(bytes); - if (bytes.len < segment.header_size) return; - - if (segment.Header.decode(bytes)) |header| { - self.next_seq = @max(self.next_seq, header.max_seq + 1); - return; - } else |err| switch (err) { - error.ActiveSegment => {}, - else => { - log.warn("recovery: {s} unreadable ({s}), leaving as-is", .{ name, @errorName(err) }); - return; + var block = (walk.next(self.allocator) catch |err| switch (err) { + error.Truncated => { + torn = true; + break; }, - } + else => return err, + }) orelse break; + defer block.deinit(self.allocator); + if (block.events.len == 0) break; + valid_end = walk.offset; + event_count += block.events.len; + block_count += 1; + for (block.events) |event| max_seq = @max(max_seq, event.seq); + } + if (max_seq == 0 and candidate > 0) + max_seq = try self.sealedFloorBefore(candidate); + + var writer = try writer_mod.ActiveWriter.init(self.allocator); + errdefer writer.deinit(); + writer.buf.clearRetainingCapacity(); + try writer.buf.appendSlice(self.allocator, bytes[0..valid_end]); + writer.event_count = event_count; + writer.block_count = block_count; + + var file = try self.dir.openFile(self.io, name, .{ .mode = .read_write }); + errdefer file.close(self.io); + if (torn or valid_end != bytes.len) { + try file.setLength(self.io, valid_end); + try segment_io.sync(&file, self.io); + } - var valid_end: usize = segment.header_size; - var max_seq: u64 = 0; - var walk: writer_mod.ActiveIterator = .{ .bytes = bytes }; - while (true) { - var block = (walk.next(self.allocator) catch null) orelse break; - defer block.deinit(self.allocator); - if (block.events.len == 0) break; - valid_end = walk.offset; - max_seq = @max(max_seq, block.events[block.events.len - 1].seq); - } - if (valid_end == segment.header_size) { - // A normal restart can leave an empty newest active segment. - // Reuse its index, then inspect the preceding sealed segment - // to recover the durable seq floor instead of resetting to 1. - try self.dir.deleteFile(self.io, name); - self.seg_index = candidate; - if (candidate == 0) return; - candidate -= 1; - continue; - } + self.writer = writer; + self.file = file; + self.file_written = valid_end; + self.synced_block_count = block_count; + self.seg_index = candidate; + self.next_seq = @max(self.next_seq, max_seq + 1); + if (self.stats) |stats| + stats.ingest_active_segment_bytes.store(valid_end - segment.header_size, .monotonic); + log.info("recovery: resumed active segment {s} at {d} valid bytes", .{ name, valid_end }); + return true; + } - var w = try writer_mod.ActiveWriter.init(self.allocator); - defer w.deinit(); - w.buf.clearRetainingCapacity(); - try w.buf.appendSlice(self.allocator, bytes[0..valid_end]); - try w.seal(); + fn sealedMaxSeq(self: *Archive, candidate: u64, candidate_bytes: []const u8) !u64 { + var sealed = try segment.Sealed.parse(self.allocator, candidate_bytes); + defer sealed.deinit(self.allocator); + if (sealed.header.max_seq != 0) return sealed.header.max_seq; + if (candidate == 0) return 0; + return self.sealedFloorBefore(candidate); + } - var f = try self.dir.createFile(self.io, name, .{ .truncate = true }); - defer f.close(self.io); - try segment_io.writeStreamingAll(&f, self.io, w.bytes()); - try segment_io.sync(&f, self.io); - self.next_seq = @max(self.next_seq, max_seq + 1); - log.info("recovery: sealed interrupted segment {s} ({d} bytes valid)", .{ name, valid_end }); - return; + fn sealedFloorBefore(self: *Archive, before: u64) !u64 { + var candidate = before; + while (candidate > 0) { + candidate -= 1; + var name_buf: [64]u8 = undefined; + const name = formatSegmentName(&name_buf, candidate); + const bytes = try self.dir.readFileAlloc(self.io, name, self.allocator, .limited(1 << 31)); + defer self.allocator.free(bytes); + var sealed = try segment.Sealed.parse(self.allocator, bytes); + defer sealed.deinit(self.allocator); + if (sealed.header.max_seq != 0) return sealed.header.max_seq; } + return 0; } fn openNextSegment(self: *Archive) !void { @@ -1075,7 +1112,8 @@ test "archive: append, rotate, recover across restart" { { var a = try Archive.init(testing.allocator, io, data_dir); defer a.deinit(); - // recovery sealed the leftover; seqs continue with no reuse + // Recovery resumes the leftover active segment; seqs continue with no + // reuse and the segment index does not churn merely because of restart. const seq = try a.append(testEvent(999), 999); try testing.expectEqual(@as(u64, 101), seq); a.close(); @@ -1084,8 +1122,7 @@ test "archive: append, rotate, recover across restart" { { var a = try Archive.init(testing.allocator, io, data_dir); defer a.deinit(); - // Seals seq 101 and opens the next active segment, but deliberately - // appends nothing to it before another restart. + // The next restart resumes the same segment with seq 101 durable. try testing.expectEqual(@as(u64, 102), a.next_seq); a.close(); } @@ -1093,36 +1130,190 @@ test "archive: append, rotate, recover across restart" { { var a = try Archive.init(testing.allocator, io, data_dir); defer a.deinit(); - // The empty newest active segment must not hide the preceding sealed - // max seq and reset allocation to 1. + // Repeated recovery must retain the active segment's high-water mark. try testing.expectEqual(@as(u64, 102), a.next_seq); const seq = try a.append(testEvent(1000), 1000); try testing.expectEqual(@as(u64, 102), seq); } - // every sealed file in the dir must parse with our reader + // Every sealed file parses, and the one active file contains the durable + // suffix that was resumed across all three restarts. var dir = try Io.Dir.cwd().openDir(io, data_dir, .{}); defer dir.close(io); var seg_dir = try dir.openDir(io, "segments", .{ .iterate = true }); defer seg_dir.close(io); var it = seg_dir.iterate(); var sealed_count: usize = 0; + var active_count: usize = 0; var total_events: u64 = 0; while (try it.next(io)) |entry| { const bytes = try seg_dir.readFileAlloc(io, entry.name, testing.allocator, .limited(1 << 20)); defer testing.allocator.free(bytes); var sealed = segment.Sealed.parse(testing.allocator, bytes) catch |err| switch (err) { - error.ActiveSegment => continue, + error.ActiveSegment => { + active_count += 1; + var walk: writer_mod.ActiveIterator = .{ .bytes = bytes }; + while (try walk.next(testing.allocator)) |block_value| { + var block = block_value; + defer block.deinit(testing.allocator); + total_events += block.events.len; + } + continue; + }, else => return err, }; defer sealed.deinit(testing.allocator); sealed_count += 1; total_events += sealed.header.event_count; } - try testing.expect(sealed_count >= 2); + try testing.expect(sealed_count >= 1); + try testing.expectEqual(@as(usize, 1), active_count); try testing.expectEqual(@as(u64, 101), total_events); // event 102 is unflushed-active } +test "archive recovery fails closed on a complete corrupt active frame" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tmp = testing.tmpDir(.{}); + defer tmp.cleanup(); + var root_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const root_len = try tmp.dir.realPath(io, &root_buf); + var data_dir_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const data_dir = try std.fmt.bufPrint(&data_dir_buf, "{s}/archive", .{root_buf[0..root_len]}); + + { + var archive = try Archive.init(testing.allocator, io, data_dir); + defer archive.deinit(); + archive.max_events_per_block = 1; + archive.writer.max_events_per_block = 1; + _ = try archive.append(testEvent(1), 1); + archive.close(); + } + + var segment_path_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const segment_path = try std.fmt.bufPrint( + &segment_path_buf, + "{s}/segments/seg_0000000000.jss", + .{data_dir}, + ); + var file = try Io.Dir.cwd().openFile(io, segment_path, .{ .mode = .read_write }); + const stat = try file.stat(io); + const before = try testing.allocator.alloc(u8, @intCast(stat.size)); + defer testing.allocator.free(before); + try testing.expectEqual(before.len, try file.readPositionalAll(io, before, 0)); + // Preserve the complete length-prefixed frame while corrupting its zstd + // magic. This is not a torn tail and must never be truncated or sealed. + const frame_start = segment.header_size + 8; + before[frame_start] ^= 0xff; + try file.writePositionalAll(io, before[frame_start..][0..1], frame_start); + try file.sync(io); + file.close(io); + + try testing.expectError( + error.DecompressFailed, + Archive.init(testing.allocator, io, data_dir), + ); + + const after = try Io.Dir.cwd().readFileAlloc( + io, + segment_path, + testing.allocator, + .limited(1 << 20), + ); + defer testing.allocator.free(after); + try testing.expectEqualSlices(u8, before, after); +} + +test "archive recovery truncates only a torn suffix and resumes the same segment" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tmp = testing.tmpDir(.{}); + defer tmp.cleanup(); + var root_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const root_len = try tmp.dir.realPath(io, &root_buf); + var data_dir_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const data_dir = try std.fmt.bufPrint(&data_dir_buf, "{s}/archive", .{root_buf[0..root_len]}); + + { + var archive = try Archive.init(testing.allocator, io, data_dir); + defer archive.deinit(); + archive.max_events_per_block = 1; + archive.writer.max_events_per_block = 1; + _ = try archive.append(testEvent(1), 1); + archive.close(); + } + + var segment_path_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const segment_path = try std.fmt.bufPrint( + &segment_path_buf, + "{s}/segments/seg_0000000000.jss", + .{data_dir}, + ); + var file = try Io.Dir.cwd().openFile(io, segment_path, .{ .mode = .read_write }); + const durable_len = (try file.stat(io)).size; + var torn: [12]u8 = undefined; + std.mem.writeInt(u64, torn[0..8], 1 << 20, .little); + @memset(torn[8..], 0xff); + try file.writePositionalAll(io, &torn, durable_len); + try file.sync(io); + file.close(io); + + { + var recovered = try Archive.init(testing.allocator, io, data_dir); + defer recovered.deinit(); + try testing.expectEqual(@as(u64, 0), recovered.seg_index); + try testing.expectEqual(@as(usize, @intCast(durable_len)), recovered.file_written); + try testing.expectEqual(durable_len, (try recovered.file.stat(io)).size); + const seq = try recovered.append(testEvent(2), 2); + try testing.expectEqual(@as(u64, 2), seq); + recovered.close(); + } + var reopened = try Archive.init(testing.allocator, io, data_dir); + defer reopened.deinit(); + try testing.expectEqual(@as(u64, 3), reopened.next_seq); +} + +test "archive recovery checksum-verifies the sealed high-water floor" { + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tmp = testing.tmpDir(.{}); + defer tmp.cleanup(); + var root_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const root_len = try tmp.dir.realPath(io, &root_buf); + var data_dir_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const data_dir = try std.fmt.bufPrint(&data_dir_buf, "{s}/archive", .{root_buf[0..root_len]}); + + { + var archive = try Archive.init(testing.allocator, io, data_dir); + defer archive.deinit(); + archive.max_events_per_block = 1; + archive.writer.max_events_per_block = 1; + _ = try archive.append(testEvent(1), 1); + try archive.rotate(); + } + + var sealed_path_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const sealed_path = try std.fmt.bufPrint( + &sealed_path_buf, + "{s}/segments/seg_0000000000.jss", + .{data_dir}, + ); + var sealed_file = try Io.Dir.cwd().openFile(io, sealed_path, .{ .mode = .read_write }); + var corrupt_max_seq: [8]u8 = undefined; + std.mem.writeInt(u64, &corrupt_max_seq, 999_999, .little); + try sealed_file.writePositionalAll(io, &corrupt_max_seq, 34); + try sealed_file.sync(io); + sealed_file.close(io); + + try testing.expectError( + error.BadChecksum, + Archive.init(testing.allocator, io, data_dir), + ); +} + test "archive: metadata failure does not publish or strand durability" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); diff --git a/tests/upstream_oracle/stream_differential_test.go b/tests/upstream_oracle/stream_differential_test.go index 538b461..4b64663 100644 --- a/tests/upstream_oracle/stream_differential_test.go +++ b/tests/upstream_oracle/stream_differential_test.go @@ -7,6 +7,7 @@ package oracle import ( "bytes" "context" + "encoding/json" "fmt" "io" "log/slog" @@ -750,8 +751,10 @@ func TestStreamDifferentialOracle(t *testing.T) { require.Error(t, CompareEventLogMultiset(expected, mutated), "oracle must detect a missing intermediate update") - // Reopen from the durable steady-state metadata. Recovery seals the active - // file, and serving must begin without re-bootstrap or content changes. + // Reopen from durable steady-state metadata. Like upstream, recovery resumes + // the active segment in place: the archive XRPC frontier remains the last + // sealed segment, while a normal subscriber must replay both sealed and + // active durable rows before joining live delivery. publicThrough := initialMaxSeqFrom(all) cmd2 := exec.CommandContext(ctx, streamBin, cmd.Args[1:]...) cmd2.Stdout = logs @@ -767,17 +770,27 @@ func TestStreamDifferentialOracle(t *testing.T) { }() waitForStreamOracleServing(t, cmd2, baseURL, logs) + sealedThrough := streamSealedHighWater(t, baseURL) + require.Less(t, sealedThrough, publicThrough, + "restart oracle must retain an active segment beyond the sealed XRPC frontier") + sealedEvents := collectStreamOracleBackfill(t, baseURL, sealedThrough) + require.NoError(t, CompareEventLogMultiset( + NormalizeEventLog(observedEventsThrough(all, sealedThrough)), + NormalizeEventLog(sealedEvents), + ), "backfill-only public client must reproduce the sealed archive frontier") + // Exercise the pinned upstream public client against Stream's actual XRPC - // archive plan/getSegment/getBlock + /subscribe-v2 decode path. - publicEvents := collectStreamOracleBackfill(t, baseURL, publicThrough) + // archive plan/getSegment/getBlock plus cold active-segment replay and + // /subscribe-v2 decode path. + publicEvents := collectStreamOracleThrough(t, baseURL, publicThrough) publicModel, err := Reconstruct(EventsSortedBySeq(publicEvents)) require.NoError(t, err) require.NoError(t, Compare(ground, publicModel), "pinned upstream public client must reconstruct Stream to simulator ground truth") require.Zero(t, fan.TotalDrops(), "simulator fanout must not drop oracle frames") - t.Logf("semantic receipt: bootstrap=%d live_rows=%d total=%d public=%d upstream_seq=%d..%d pin=%s", - len(initial), len(expected), len(all), len(publicEvents), startTip+1, endTip, streamDifferentialPin) + t.Logf("semantic receipt: bootstrap=%d live_rows=%d total=%d sealed_public=%d public=%d sealed_through=%d upstream_seq=%d..%d pin=%s", + len(initial), len(expected), len(all), len(sealedEvents), len(publicEvents), sealedThrough, startTip+1, endTip, streamDifferentialPin) require.NoError(t, cmd2.Process.Signal(syscall.SIGTERM)) require.NoErrorf(t, cmd2.Wait(), "restarted Stream did not shut down cleanly:\n%s", logs.String()) @@ -930,6 +943,40 @@ func initialMaxSeqFrom(events []ObservedEvent) uint64 { return maxSeq } +func observedEventsThrough(events []ObservedEvent, through uint64) []ObservedEvent { + out := make([]ObservedEvent, 0, len(events)) + for _, event := range events { + if event.Seq <= through { + out = append(out, event) + } + } + return out +} + +func streamSealedHighWater(t *testing.T, baseURL string) uint64 { + t.Helper() + resp, err := (&http.Client{Timeout: 5 * time.Second}).Get( + baseURL + "/xrpc/network.bsky.jetstream.listSegments?limit=1000", + ) + require.NoError(t, err) + defer func() { require.NoError(t, resp.Body.Close()) }() + require.Equal(t, http.StatusOK, resp.StatusCode) + + var payload struct { + Segments []struct { + MaxSeq uint64 `json:"maxSeq"` + } `json:"segments"` + } + require.NoError(t, json.NewDecoder(resp.Body).Decode(&payload)) + require.NotEmpty(t, payload.Segments) + var highWater uint64 + for _, segment := range payload.Segments { + highWater = max(highWater, segment.MaxSeq) + } + require.NotZero(t, highWater) + return highWater +} + func collectStreamOracleBackfill(t *testing.T, baseURL string, through uint64) []ObservedEvent { t.Helper() client, err := jetstream.Subscribe(baseURL, @@ -955,3 +1002,33 @@ func collectStreamOracleBackfill(t *testing.T, baseURL string, through uint64) [ require.Equal(t, through, initialMaxSeqFrom(out), "public client must reach archive high-water") return out } + +func collectStreamOracleThrough(t *testing.T, baseURL string, through uint64) []ObservedEvent { + t.Helper() + client, err := jetstream.Subscribe(baseURL, + jetstream.WithAfterSeq(0), + jetstream.WithBatchSize(32), + ) + require.NoError(t, err) + defer func() { require.NoError(t, client.Close()) }() + + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + var out []ObservedEvent + for batch, err := range client.Events(ctx) { + require.NoError(t, err) + for _, event := range batch.Events() { + observed := observedEventFromClient(t, event) + if observed.Seq <= through { + out = append(out, observed) + } + } + if initialMaxSeqFrom(out) >= through { + break + } + } + require.NotEmpty(t, out) + require.Equal(t, through, initialMaxSeqFrom(out), + "normal public client must replay through the active durable frontier") + return out +}