From e44df7a377e1edfd90bd65cc5bcfcffed9ee1eb2 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Tue, 11 Aug 2026 16:55:43 -0500 Subject: [PATCH] cold: binary-search active-segment time resolution so mid-band connects stop timing out resolveTimeInActive decompressed every fsynced block from the front of the active segment inside the websocket handshake. a time cursor older than the hot tail but not yet sealed - a band hours wide - took 40s+ to resolve mid-fill and timed the connect out entirely: coral retried the same cursor against the same host forever (its zat 0.3.25 has no idle watchdog and never resets backoff). block boundaries now come from the 8-byte length prefixes alone and only log2(blocks) probes decompress, matching how sealed segments already resolve via their block index. regression test sweeps every microsecond boundary across a ten-block active region against the exhaustive definition. Co-Authored-By: Claude Opus 5 (1M context) --- src/internal/storage/cold.zig | 149 +++++++++++++++++++++++++++++----- 1 file changed, 128 insertions(+), 21 deletions(-) diff --git a/src/internal/storage/cold.zig b/src/internal/storage/cold.zig index 51002a7..91b468f 100644 --- a/src/internal/storage/cold.zig +++ b/src/internal/storage/cold.zig @@ -438,34 +438,76 @@ fn resolveTimeInActive( var file = archive.dir.openFile(io, name, .{}) catch return .{ .seq = sealed_fallback }; defer file.close(io); + // Walk only the 8-byte length prefixes to learn block boundaries (no + // decompression), then binary-search with one decompressed probe per + // step. The previous sequential decompress-everything walk took tens of + // seconds on a mid-fill active segment and ran inside the websocket + // handshake: any time cursor aged past the hot tail but not yet sealed + // (a band hours wide) timed the connect out entirely (2026-08-11 coral + // wedge). Sealed segments already resolve this way via their block index. + var offsets: std.ArrayList(u64) = .empty; + defer offsets.deinit(allocator); + { + var offset: usize = segment.header_size; + while (offset < snapshot.active_length) { + if (offset + 8 > snapshot.active_length) break; + var frame_len_bytes: [8]u8 = undefined; + _ = try file.readPositionalAll(io, &frame_len_bytes, offset); + const frame_len = std.mem.readInt(u64, &frame_len_bytes, .little); + if (frame_len > snapshot.active_length - offset - 8) break; + try offsets.append(allocator, offset); + offset += 8 + @as(usize, @intCast(frame_len)); + } + } + if (offsets.items.len == 0) return .{ .seq = @max(sealed_fallback, 1) }; + + const Probe = struct { + fn load(alloc: Allocator, io_: Io, f: *Io.File, off: u64, end: u64) !segment.Block { + var frame_len_bytes: [8]u8 = undefined; + _ = try f.readPositionalAll(io_, &frame_len_bytes, @intCast(off)); + const frame_len = std.mem.readInt(u64, &frame_len_bytes, .little); + if (frame_len > end - off - 8) return error.Truncated; + const compressed = try alloc.alloc(u8, @intCast(frame_len)); + _ = try f.readPositionalAll(io_, compressed, @intCast(off + 8)); + var frame_reader: std.Io.Reader = .fixed(compressed); + var zstd_stream: std.compress.zstd.Decompress = .init(&frame_reader, &.{}, .{}); + var decoded: std.Io.Writer.Allocating = .init(alloc); + _ = zstd_stream.reader.streamRemaining(&decoded.writer) catch return error.Truncated; + return segment.decodeBlock(alloc, try decoded.toOwnedSlice()); + } + }; + + // leftmost block whose last event is witnessed at/after time_us var next_after: u64 = sealed_fallback; - var offset: usize = segment.header_size; - while (offset < snapshot.active_length) { + var lo: usize = 0; + var hi: usize = offsets.items.len; + while (lo < hi) { + const mid = lo + (hi - lo) / 2; var block_arena = std.heap.ArenaAllocator.init(allocator); defer block_arena.deinit(); - if (offset + 8 > snapshot.active_length) break; - var frame_len_bytes: [8]u8 = undefined; - _ = try file.readPositionalAll(io, &frame_len_bytes, offset); - const frame_len = std.mem.readInt(u64, &frame_len_bytes, .little); - if (frame_len > snapshot.active_length - offset - 8) break; - const compressed = try block_arena.allocator().alloc(u8, @intCast(frame_len)); - _ = try file.readPositionalAll(io, compressed, offset + 8); - offset += 8 + @as(usize, @intCast(frame_len)); - var frame_reader: std.Io.Reader = .fixed(compressed); - var zstd_stream: std.compress.zstd.Decompress = .init(&frame_reader, &.{}, .{}); - var decoded: std.Io.Writer.Allocating = .init(block_arena.allocator()); - _ = zstd_stream.reader.streamRemaining(&decoded.writer) catch break; - const block = try segment.decodeBlock(block_arena.allocator(), try decoded.toOwnedSlice()); - if (block.events.len == 0) continue; - const last = block.events[block.events.len - 1]; - if (last.witnessed_at < time_us) { - next_after = last.seq +| 1; + const block = Probe.load(block_arena.allocator(), io, &file, offsets.items[mid], snapshot.active_length) catch break; + if (block.events.len == 0) { + // empty block cannot answer; treat as older so the search moves on + lo = mid + 1; continue; } - for (block.events) |event| { - if (event.witnessed_at >= time_us) return .{ .seq = @max(event.seq, 1) }; + const last = block.events[block.events.len - 1]; + if (last.witnessed_at >= time_us) { + hi = mid; + } else { + next_after = last.seq +| 1; + lo = mid + 1; } } + if (lo < offsets.items.len) { + var block_arena = std.heap.ArenaAllocator.init(allocator); + defer block_arena.deinit(); + if (Probe.load(block_arena.allocator(), io, &file, offsets.items[lo], snapshot.active_length)) |block| { + for (block.events) |event| { + if (event.witnessed_at >= time_us) return .{ .seq = @max(event.seq, 1) }; + } + } else |_| {} + } return .{ .seq = @max(next_after, 1) }; } @@ -1446,6 +1488,71 @@ test "timestamp cursor resolves into the active region, not the sealed boundary" try testing.expectEqual(@as(u64, 7), ahead.seq); } +test "active-region time resolution: binary search sweeps every boundary like the linear walk" { + // 2026-08-11 coral wedge regression: the sequential decompress walk this + // replaced was so slow mid-fill that connects timed out. correctness of + // the replacement: every time between/at/past rows must resolve to the + // same seq the exhaustive definition gives — first fsynced active row + // witnessed at/after the time, else the next not-yet-durable seq. + 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 = 2; + archive.writer.max_events_per_block = 2; + + const base: i64 = 1_700_000_000_000_000; + _ = try archive.append(.{ + .seq = 0, + .witnessed_at = base, + .indexed_at = 0, + .kind = .delete, + .did = "did:plc:sweepfixture", + .collection = "app.bsky.feed.post", + .rkey = "3k2abcdefghij", + .rev = "3k2abcdefghij", + .payload = "", + }, base); + try archive.rotate(); + // 21 rows -> ten flushed 2-row blocks plus one pending (not fsynced) row + var times: [21]i64 = undefined; + for (0..21) |i| { + times[i] = base + 1_000 + @as(i64, @intCast(i)) * 137; + _ = try archive.append(.{ + .seq = 0, + .witnessed_at = times[i], + .indexed_at = 0, + .kind = .delete, + .did = "did:plc:sweepfixture", + .collection = "app.bsky.feed.post", + .rkey = "3k2abcdefghij", + .rev = "3k2abcdefghij", + .payload = "", + }, times[i]); + } + + // rows 1..20 are seqs 2..21; the fsynced prefix covers seqs 2..21 (ten + // blocks); the 21st appended row (seq 22) is pending. exhaustive check: + var t = times[0] - 1; + while (t <= times[20] + 200) : (t += 1) { + var expected: u64 = 22; // next not-yet-durable + for (times[0..20], 0..) |w, i| { + if (w >= t) { + expected = 2 + @as(u64, @intCast(i)); + break; + } + } + const got = try resolveTimeToSeq(testing.allocator, io, &archive, t); + try testing.expectEqual(expected, got.seq); + } +} + test "timestamp cursor surfaces corrupt selected block" { const testing = std.testing; var threaded: Io.Threaded = .init(testing.allocator, .{}); -- 2.51.2