diff --git a/src/internal/cold.zig b/src/internal/cold.zig index 318a68b..8407a70 100644 --- a/src/internal/cold.zig +++ b/src/internal/cold.zig @@ -98,7 +98,13 @@ pub const ColdReader = struct { 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; + var file = archive.dir.openFile(io, name, .{}) catch |err| { + // A manifest-listed segment we cannot open is a hole, not an + // empty range. Skipping it would advance the cursor past + // durable events that were never delivered. + log.warn("cold: open sealed segment {s}: {s}", .{ name, @errorName(err) }); + return err; + }; 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; @@ -434,7 +440,7 @@ pub fn replay( 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; + return err; }; defer file.close(io); const fresh_header = try segment.readHeaderFile(io, &file); @@ -516,7 +522,10 @@ pub fn replaySeq( 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; + var file = archive.dir.openFile(io, name, .{}) catch |err| { + log.warn("cold: open sealed segment {s}: {s}", .{ name, @errorName(err) }); + return err; + }; 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; @@ -898,6 +907,138 @@ test "cold batches include the active segment fsynced prefix" { try testing.expectEqualSlices(u64, &.{ 1, 2, 3, 4 }, sink.seqs[0..sink.count]); } +// Upstream's walkSealedSegment returns `open seg %d: %w` and leaves the +// replay cursor where it was. Skipping an unreadable manifest-listed segment +// would emit the following segment's events and advance the cursor past +// durable events the subscriber never received -- an undetectable hole. +test "cold replay refuses to skip an unreadable sealed segment" { + 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(); + + for (0..4) |i| { + _ = try archive.append(.{ + .seq = 0, + .witnessed_at = 1_700_000_000_000_000 + @as(i64, @intCast(i)), + .indexed_at = 0, + .kind = .create, + .did = "did:plc:coldholefixture", + .collection = "app.bsky.feed.post", + .rkey = "3k2abcdefghij", + .rev = "3k2abcdefghij", + .payload = "\xa1\x64text\x62hi", + }, 1_700_000_000_000_000 + @as(i64, @intCast(i))); + if (i == 1) try archive.rotate(); + } + try archive.rotate(); + + const summaries = try archive.manifest.all(testing.allocator); + defer testing.allocator.free(summaries); + try testing.expectEqual(@as(usize, 2), summaries.len); + + // The manifest still lists segment 0; the file is gone underneath it. + try archive.dir.deleteFile(io, "seg_0000000000.jss"); + + var reader = ColdReader.init(testing.allocator, io, 1 << 20); + defer reader.deinit(); + const Sink = struct { + delivered: usize = 0, + + fn wants(_: *@This(), _: wire.Kind, _: []const u8, _: []const u8, _: tail.SkipV1) bool { + return true; + } + fn write(_: *@This(), _: CachedFrame) anyerror!WriteResult { + return .sent; + } + fn observe(self: *@This(), _: u64, result: ReplayResult) anyerror!void { + if (result == .sent) self.delivered += 1; + } + }; + var sink: Sink = .{}; + try testing.expectError(error.FileNotFound, reader.readSeqBatch( + &archive, + &sink, + Sink.wants, + &sink, + Sink.write, + Sink.observe, + false, + false, + 1, + 5, + 16, + )); + // Nothing from the intact later segment may be emitted over the hole. + try testing.expectEqual(@as(usize, 0), sink.delivered); +} + +// Same invariant on the second cold walker: replaySeq must not step over an +// unreadable manifest-listed segment to reach the intact one behind it. +test "replaySeq refuses to skip an unreadable sealed segment" { + 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(); + + for (0..4) |i| { + _ = try archive.append(.{ + .seq = 0, + .witnessed_at = 1_700_000_000_000_000 + @as(i64, @intCast(i)), + .indexed_at = 0, + .kind = .create, + .did = "did:plc:replayseqhole", + .collection = "app.bsky.feed.post", + .rkey = "3k2abcdefghij", + .rev = "3k2abcdefghij", + .payload = "\xa1\x64text\x62hi", + }, 1_700_000_000_000_000 + @as(i64, @intCast(i))); + if (i == 1) try archive.rotate(); + } + try archive.rotate(); + try archive.dir.deleteFile(io, "seg_0000000000.jss"); + + const Sink = struct { + delivered: usize = 0, + + fn wants(_: *@This(), _: wire.Kind, _: []const u8, _: []const u8, _: tail.SkipV1) bool { + return true; + } + fn write(_: *@This(), _: []const u8) anyerror!WriteResult { + return .sent; + } + fn observe(self: *@This(), _: u64, result: ReplayResult) anyerror!void { + if (result == .sent) self.delivered += 1; + } + }; + var sink: Sink = .{}; + try testing.expectError(error.FileNotFound, replaySeq( + testing.allocator, + io, + &archive, + &sink, + Sink.wants, + &sink, + Sink.write, + Sink.observe, + false, + 1, + )); + try testing.expectEqual(@as(usize, 0), sink.delivered); +} + test "cold cache single-flights real block decode and shares wire memos" { const testing = std.testing; var threaded: Io.Threaded = .init(testing.allocator, .{});