diff --git a/README.md b/README.md index 4ea7896..5876c8f 100644 --- a/README.md +++ b/README.md @@ -27,22 +27,10 @@ filtered archive replay, and a gapless replay→live cutover. typed pre-upgrade errors (`CursorTooOld` / `UnknownZstdDictionary` / `InvalidRequest`), reconnect backoff that resumes at the highest delivered seq, TCP keepalive. -- **multi-host failover** (extension beyond the upstream Go client, which - binds a single `service`): `fallback_hosts` rotates to another instance - after `failover_after_errors` consecutive live failures with no events - between (or a sweep whose bounded retries exhaust). v2 seqs are - per-instance, so the switch re-anchors by witnessed time: the new host's - `listSegments` bounds place the sweep (`anchorForFloor`), and a delivery - time floor at `last_time_us − failover_rewind_us` trims the - segment-granularity over-fetch, bounding duplicates to the rewind margin - (default 10s, the v1 client's number). failback retries the primary after - any rotation that followed delivered events, mirroring zat's v1 - round-robin (`nextHostIndex`). cross-host semantics are at-least-once for - the live flow; a caller-supplied `after_seq`/`live_cursor` is bound to - the primary's seq space, so a primary that dies before delivering - anything fails loudly rather than silently restarting elsewhere. note: - re-anchoring reads the new host's `listSegments`, so token-gated - archives need `api_key` for failover, not just for sweeps. +- **multi-host failover** — `fallback_hosts` rotates to another instance + when the current one is unhealthy, re-anchoring across per-instance seq + spaces by witnessed time. an extension beyond the upstream Go client; + mechanics, triggers, and caveats in [docs/failover.md](docs/failover.md). - **`ArchiveBackfill`** — the archive half standalone: `planSnapshot` (collections/dids, seq bounds, `api_key` for token-gated instances) → `getBlock`/`getSegment` → jss v1 columnar decode. two handler contracts: diff --git a/docs/failover.md b/docs/failover.md new file mode 100644 index 0000000..41be224 --- /dev/null +++ b/docs/failover.md @@ -0,0 +1,78 @@ +# multi-host failover + +an extension beyond the upstream Go client, which binds a single `service`. +motivated by the 2026-08-16 `bsky.network` front outage: every consumer of +one hostname went dark together, and v2's per-instance seq cursors meant +none of them could hop to a healthy sibling the way zat's v1 client +round-robins (v1 cursors are wall-clock `time_us`; see zat +`src/internal/streaming/jetstream.zig`). + +## the problem + +jetstream v2 assigns each event a per-instance monotonic seq, and the +subscribe cursor is that seq. this buys exact at-least-once delivery and +gapless archive replay against one instance — and makes the cursor +meaningless on any other instance. a client that wants to survive its +instance dying needs a portable anchor. + +## the bridge: witnessed time + +every delivered event carries `time_us` (the instance's `witnessed_at`). +witnessed time transfers between instances approximately — each instance +witnessed the same firehose independently, skewed by connection lag. so on +failover the client re-enters the NEXT host's seq space through that +host's own archive metadata: + +1. **anchor** — fetch the new host's `listSegments`; each sealed segment + row carries `minSeq/maxSeq/minWitnessedAt/maxWitnessedAt`. + `anchorForFloor(spans, floor_us)` (archive_backfill.zig) reduces them: + - floor inside the sealed range → sweep from one seq below the earliest + intersecting segment; + - floor past every sealed segment → the window lives in the unsealed + tail; anchor at the sealed tip and let the live cutover's server-side + cold replay cover the rest; + - no sealed archive at all → no anchor; live from the tip. +2. **trim** — segment granularity is coarse (a segment spans hours), so + the seq bound decides what to FETCH and a delivery time floor at + `last_time_us − failover_rewind_us` decides what to DELIVER. rows + witnessed before the floor are skipped in both the sweep and live + halves, bounding duplicates to the rewind margin (default 10s, the v1 + client's number). +3. **resume** — the normal sweep→cutover loop runs on the new host from + the anchor. all machinery past the anchor point is the standard + single-host path. + +## when rotation triggers + +- **live half**: `failover_after_errors` consecutive transient failures + with no events between (an event between failures resets the count — + a limping host is not a dead one). +- **archive half**: a sweep whose bounded internal retries exhaust + (`fetchWithRetry`, 4 attempts) is a host-health verdict when there is + somewhere else to go. + +`failover_after_errors = 0` or an empty `fallback_hosts` disables +rotation entirely; the single-host error surface is exactly the +pre-failover one. + +## failback + +mirrors the v1 client's `nextHostIndex`: a rotation that follows delivered +events retries the primary first; failed attempts keep rotating through +the list with exponential backoff between full cycles. + +## semantics and caveats + +- cross-host flow is **at-least-once**; duplicates are bounded by the + rewind margin. consumers should already be idempotent (the same + contract as same-host reconnects). +- instances witness independently: the margin covers connection skew, not + rows the new host witnessed at very different times (e.g. its own + resyncs). a switch can miss such rows. +- a caller-supplied `after_seq`/`live_cursor` is bound to the PRIMARY's + seq space. if the primary dies before any event is delivered there is + no time anchor yet, and the client fails with `error.HostUnhealthy` + rather than silently restarting somewhere else with a broken resume + contract. +- re-anchoring reads the new host's `listSegments`, so token-gated + archives need `api_key` for failover, not just for sweeps. diff --git a/src/client.zig b/src/client.zig index 09ba342..4d58774 100644 --- a/src/client.zig +++ b/src/client.zig @@ -78,6 +78,12 @@ pub const Options = struct { /// witnessed-time rewind margin when re-anchoring on a new host, /// covering inter-instance witness skew (v1 client rewinds 10s) failover_rewind_us: i64 = 10 * std.time.us_per_s, + /// treat a live session with no EVENTS for this long as failed + /// (error.LiveStalled → counts toward failover_after_errors). catches + /// the connected-but-silent failure mode that produces no transport + /// errors at all (the 2026-08-16 bsky.network front wedge). opt-in: + /// a heavily filtered stream can be legitimately quiet for hours. + failover_stall_ns: u64 = 0, /// replay start, exclusive: sweep the sealed archive from here and cut /// over to live. null = no replay — a pure live tail from the tip /// (set live_cursor to resume a live tail at a known seq instead). @@ -322,6 +328,7 @@ fn liveOptions(options: Options, host: []const u8, cursor: ?u64, dedup_floor: u6 .collections = options.collections, .dids = options.dids, .zstd_dict = dict, + .stall_timeout_ns = options.failover_stall_ns, }; if (options.backoff_min_ns > 0) out.backoff_min_ns = options.backoff_min_ns; if (options.backoff_max_ns > 0) out.backoff_max_ns = options.backoff_max_ns; diff --git a/src/live.zig b/src/live.zig index ef0c230..bca36c5 100644 --- a/src/live.zig +++ b/src/live.zig @@ -73,6 +73,13 @@ pub const Options = struct { /// then degrades to an uncompressed tail. refetch_dict: ?*const fn (ctx: ?*anyopaque, allocator: Allocator) ?[]u8 = null, refetch_dict_ctx: ?*anyopaque = null, + /// declare the session stalled when no EVENT arrives for this long + /// (surfaced as transient error.LiveStalled → reconnect/failover). + /// event-based, not frame-based: a connected-but-silent server that + /// still answers pings must not count as alive (the 2026-08-16 + /// bsky.network front wedge). 0 disables — a heavily filtered stream + /// can be legitimately quiet for hours, so this is opt-in. + stall_timeout_ns: u64 = 0, }; /// terminal outcomes of run(); everything transient reconnects internally @@ -208,8 +215,12 @@ pub const LiveClient = struct { if (comptime @hasDecl(@TypeOf(handler.*), "onConnect")) handler.onConnect(target.host); var ws_handler: WsHandler(@TypeOf(handler.*)) = .{ .client = self, .handler = handler }; - client.readLoop(&ws_handler) catch |err| switch (err) { + self.readFrames(&client, &ws_handler) catch |err| switch (err) { error.Canceled => return error.Canceled, + // handler-requested stop: the flags below carry the reason + // (previously fell into the transient arm, so a live-phase + // onEvent=false reconnect-looped instead of stopping) + error.Stop => {}, else => return .{ .transient = err }, }; if (ws_handler.stopped) return .stopped; @@ -222,6 +233,59 @@ pub const LiveClient = struct { return .{ .transient = error.ConnectionClosed }; } + /// frame pump. without a stall timeout this is the library's readLoop; + /// with one, a poll-timeout read loop that watches EVENT progress + /// (last_seq), so a connected-but-silent server — even one answering + /// pings, which resets any frame-level liveness check — surfaces as + /// error.LiveStalled within one poll of the deadline. + fn readFrames(self: *LiveClient, client: *websocket.Client, ws_handler: anytype) !void { + const stall_ns = self.options.stall_timeout_ns; + if (stall_ns == 0) return client.readLoop(ws_handler); + + try client.readTimeout(stallPollMs(stall_ns)); + var last_seq_seen = self.last_seq; + var quiet_since: u64 = @intCast(Io.Timestamp.now(self.io, .awake).nanoseconds); + while (true) { + const message = client.read() catch |err| switch (err) { + error.Closed => return, + else => return err, + } orelse { + // poll timeout: frames may or may not be flowing — only + // event progress proves the stream alive + const now: u64 = @intCast(Io.Timestamp.now(self.io, .awake).nanoseconds); + if (self.last_seq != last_seq_seen) { + last_seq_seen = self.last_seq; + quiet_since = now; + } else if (now - quiet_since >= stall_ns) { + return error.LiveStalled; + } + continue; + }; + defer client.done(message); + switch (message.type) { + .text, .binary => { + try ws_handler.serverMessage(message.data, if (message.type == .text) .text else .binary); + if (self.last_seq != last_seq_seen) { + last_seq_seen = self.last_seq; + quiet_since = @intCast(Io.Timestamp.now(self.io, .awake).nanoseconds); + } + }, + .ping => try client.writeFrame(.pong, @constCast(message.data)), + .close => { + client.close(.{}) catch {}; + return; + }, + .pong => {}, + } + } + } + + /// poll interval: a quarter of the stall window, clamped to [250ms, 5s] + fn stallPollMs(stall_ns: u64) u32 { + const quarter_ms = stall_ns / std.time.ns_per_ms / 4; + return @intCast(@min(@max(quarter_ms, 250), 5_000)); + } + /// pre-upgrade rejection classification: match the structured xrpc /// envelope's error NAME, never body substrings (the wire contract; /// upstream dialWebsocket). any other status is transient — a proxy @@ -502,3 +566,9 @@ test "rejection classification matches the wire contract by error name" { const proxy = LiveClient.classifyRejection(mk.failure(503, "starting")); try testing.expect(proxy == .transient); } + +test "stallPollMs: quarter window clamped to [250ms, 5s]" { + try std.testing.expectEqual(@as(u32, 250), LiveClient.stallPollMs(100 * std.time.ns_per_ms)); // tiny window: floor + try std.testing.expectEqual(@as(u32, 5_000), LiveClient.stallPollMs(600 * std.time.ns_per_s)); // huge window: ceiling + try std.testing.expectEqual(@as(u32, 2_500), LiveClient.stallPollMs(10 * std.time.ns_per_s)); // 10s window: 2.5s poll +}