diff --git a/docs/upstream-subscribe-protocol.md b/docs/upstream-subscribe-protocol.md index 10e737c..00fbf30 100644 --- a/docs/upstream-subscribe-protocol.md +++ b/docs/upstream-subscribe-protocol.md @@ -64,6 +64,21 @@ endpoint: `GET /subscribe` → websocket upgrade. pre-upgrade errors are plain H `service not ready: cursor resolution failed`, metric mode `resolve_failed`, and internal detail only in server logs. It is not a permanent client-input 400. +- **hot-window timestamp cursors resolve against the hot tail, not the + archive** (Stream translation of upstream's active-segment replay walk). + Upstream's replay engine walks sealed segments and then the active segment + before handing to live; Stream's port replaced the active-segment walk with + the hot-tail deque, so a time cursor at or after the oldest retained hot + frame must start on the hot tail directly. Resolving it through + `resolveTimeToSeq` instead lands at the last sealed-segment boundary and + strands the subscriber on the cold reader, which delivers `read-batch` + bursts at block-seal cadence — the jerky-resume defect observed live + 2026-08-08 (coral and bufo-bot resumed with near-now cursors after a proxy + restart and trailed seals indefinitely). Known-open follow-up: a genuinely + cold resume (cursor older than hot retention) still trails seal cadence + after catching up — the cold→hot seam handoff in `coldBatch` does not fire + at the tip (reproduced with a 10s-rewind probe that stayed on 1024-row + bursts 7+ minutes after connect). ### compress - an ordinary `/subscribe` connection negotiates RFC 7692 diff --git a/src/internal/serve/server.zig b/src/internal/serve/server.zig index 8a2466e..5ba838e 100644 --- a/src/internal/serve/server.zig +++ b/src/internal/serve/server.zig @@ -1744,6 +1744,17 @@ fn resolveCursor(hub: *Hub, v2: bool, raw_param: ?[]const u8) !CursorPlan { } if (parsed > now_us) return .{ .mode = .clamped }; + // Hot-first: a time cursor within the retained hot window starts on the + // hot tail directly. Resolving it through the archive instead lands at + // the last sealed-segment boundary — up to a whole open segment behind — + // and strands the subscriber on the cold reader, which delivers in + // read-batch bursts at seal cadence (the "jerky resume" defect). + if (hub.tail.oldestTime()) |oldest_hot_time| { + if (parsed >= oldest_hot_time) return .{ + .start_idx = hub.tail.indexForCursor(parsed), + .mode = .time_us, + }; + } if (hub.archive) |archive| { const resolved = cold.resolveTimeToSeq(hub.allocator, hub.io, archive, parsed) catch |err| { log.err("timestamp cursor resolution failed for {d}: {s}", .{ parsed, @errorName(err) }); @@ -1867,6 +1878,46 @@ test "cursor resolution is pre-upgrade and preserves v1-v2 floor semantics" { try std.testing.expectError(error.InvalidCursor, resolveCursor(&hub, true, "-1")); } +test "near-now time cursor within hot retention starts hot, never cold" { + const testing = std.testing; + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tail = tail_mod.Tail.init(testing.allocator, io, 1 << 20); + defer tail.deinit(); + try tail.append(.commit, "did:plc:test", "app.bsky.feed.post", 1_750_000_000_000_000, 100, "{}", "{}", .none); + try tail.append(.commit, "did:plc:test", "app.bsky.feed.post", 1_750_000_000_500_000, 101, "{}", "{}", .none); + + // The archive is present but knows nothing recent: resolving a hot-window + // time cursor through it strands the subscriber on the cold reader (the + // jerky-resume defect). Hot rows must win. + var tmp = testing.tmpDir(.{}); + defer tmp.cleanup(); + var data_dir_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const data_dir_len = try tmp.dir.realPath(io, &data_dir_buf); + var archive = try archive_mod.Archive.init(testing.allocator, io, data_dir_buf[0..data_dir_len]); + defer archive.deinit(); + + var stats: metrics.Stats = .{}; + var hub: Hub = .{ + .allocator = testing.allocator, + .io = io, + .tail = &tail, + .stats = &stats, + .archive = &archive, + }; + + const plan = try resolveCursor(&hub, false, "1750000000200000"); + try testing.expectEqual(metrics.CursorMode.time_us, plan.mode); + try testing.expectEqual(@as(?u64, null), plan.cold_from_seq); + // resumes at the first retained frame witnessed at/after the cursor + try testing.expectEqual(tail.indexForSeq(101), plan.start_idx); + + // a cursor predating hot retention still consults the archive (cold path) + const old_plan = try resolveCursor(&hub, false, "1749999999999999"); + try testing.expect(old_plan.cold_from_seq != null); +} + test "XRPC readiness gate returns the upstream JSON service error" { const testing = std.testing; var threaded: Io.Threaded = .init(testing.allocator, .{});