diff --git a/CHANGELOG.md b/CHANGELOG.md index 4bdce14..d5e152a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,29 @@ # changelog +## 0.3.29 + +Idle-connection failover works on both streaming clients: a connection that +stays open but delivers nothing is now a fault the consumer can bound. + +- **feat**: `FirehoseClient` gains `idle_timeout_ms`, mirroring + `JetstreamClient`'s option. A relay that holds the socket open but goes + silent — TCP-healthy, stream stalled — is detected and failed over instead + of waited out. Observed live 2026-08-11: bsky.network went quiet for ~2.5 + minutes before closing the connection itself; consumers ate the whole gap. + Off by default (reconnect policy belongs to the consumer; a filtered + jetstream consumer may be legitimately quiet). jcalabro/atmos has no + equivalent — this is a deliberate, opt-in divergence. +- **fix**: the idle-timeout mechanism no longer uses `SO_RCVTIMEO`, which + panics under `std.Io.Threaded` in debug builds (the EAGAIN from a receive + timeout on a blocking socket is treated as a programmer bug). Both clients + now run the blocking read loop as a concurrent task with a watchdog that + compares a frame counter across quiet windows and cancels the read when a + full window passes with no frames. Fires between 1x and 2x the timeout. + `JetstreamClient.idle_timeout_ms` / `replay_idle_timeout_ms` were affected; + verified live against a silent local server (fires, control blocks + forever), the live firehose (13,322 events in 30s, zero false fires), and a + DID-filtered truly-quiet jetstream subscription (fires on schedule). + ## 0.3.28 OAuth requests can be pinned to validated addresses, DPoP nonces are optional, diff --git a/build.zig.zon b/build.zig.zon index 4704036..01cb27f 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -1,6 +1,6 @@ .{ .name = .zat, - .version = "0.3.28", + .version = "0.3.29", .fingerprint = 0x8da9db57ee82fbe4, .minimum_zig_version = "0.16.0-dev.3070+b22eb176b", .dependencies = .{ diff --git a/src/internal/streaming/firehose.zig b/src/internal/streaming/firehose.zig index 72d47ab..534a1c3 100644 --- a/src/internal/streaming/firehose.zig +++ b/src/internal/streaming/firehose.zig @@ -43,6 +43,14 @@ pub const Options = struct { /// path is fixed by the protocol. bare hostnames are implicit wss. hosts: []const []const u8 = &default_hosts, cursor: ?i64 = null, + /// Fail over after this many milliseconds with no WebSocket frames. + /// Catches a host that holds the socket open but goes silent, which TCP + /// keepalive never fires on (the socket is healthy; the stream behind it + /// stalled). Mirrors JetstreamClient's option of the same name. A live + /// relay emits continuously, so even a few seconds of total silence is + /// anomalous — but the default stays off: reconnect policy belongs to + /// the consumer. + idle_timeout_ms: ?u32 = null, max_message_size: usize = 5 * 1024 * 1024, // 5MB — firehose frames can be large }; @@ -720,6 +728,9 @@ pub const FirehoseClient = struct { /// the connection later ended, so the reconnect backoff starts over /// rather than compounding. connected_this_attempt: bool = false, + /// frames delivered on the current connection; written by the reader + /// task and read by the idle watchdog, hence atomic. + frames_this_connection: std.atomic.Value(u64) = .init(0), pub fn init(io: Io, allocator: Allocator, options: Options) FirehoseClient { return .{ @@ -864,7 +875,54 @@ pub const FirehoseClient = struct { .handler = handler, .client_state = self, }; - try client.readLoop(&ws_handler); + if (self.options.idle_timeout_ms) |timeout_ms| { + // A socket-level receive timeout (SO_RCVTIMEO) is off the table: + // std.Io.Threaded treats the resulting EAGAIN on a blocking + // socket as a programmer bug and panics in debug builds. Instead + // the blocking readLoop runs as a concurrent task and a watchdog + // here compares the frame counter across quiet windows, + // cancelling the read when a full window passes with no frames. + // Fires between timeout and 2x timeout after the last frame. + self.frames_this_connection.store(0, .monotonic); + var done: Io.Event = .unset; + var read_err: ?anyerror = null; + const Reader = struct { + fn run( + io: Io, + ws_client: *websocket.Client, + h: *WsHandler(@TypeOf(handler.*)), + done_ev: *Io.Event, + err_out: *?anyerror, + ) void { + ws_client.readLoop(h) catch |err| { + err_out.* = err; + }; + done_ev.set(io); + } + }; + var future = try self.io.concurrent(Reader.run, .{ self.io, &client, &ws_handler, &done, &read_err }); + var joined = false; + defer if (!joined) future.cancel(self.io); + var last_frames: u64 = 0; + while (true) { + done.waitTimeout(self.io, .{ .duration = .{ .raw = Io.Duration.fromMilliseconds(@intCast(timeout_ms)), .clock = .awake } }) catch |err| switch (err) { + error.Timeout => { + const frames = self.frames_this_connection.load(.monotonic); + if (frames == last_frames) return error.IdleTimeout; + last_frames = frames; + continue; + }, + error.Canceled => return error.Canceled, + }; + // readLoop finished on its own (close frame, socket error, stop) + future.await(self.io); + joined = true; + if (read_err) |err| return err; + return; + } + } else { + try client.readLoop(&ws_handler); + } } }; @@ -885,6 +943,7 @@ fn WsHandler(comptime H: type) type { const Self = @This(); pub fn serverMessage(self: *Self, data: []const u8) !void { + _ = self.client_state.frames_this_connection.fetchAdd(1, .monotonic); if (comptime @hasDecl(H, "shouldStop")) { // breaks the read loop; subscribe sees the stop and returns if (self.handler.shouldStop()) return error.SubscriptionStopped; diff --git a/src/internal/streaming/jetstream.zig b/src/internal/streaming/jetstream.zig index 6881ee2..f4b2399 100644 --- a/src/internal/streaming/jetstream.zig +++ b/src/internal/streaming/jetstream.zig @@ -149,6 +149,10 @@ pub const JetstreamClient = struct { options: Options, last_time_us: ?i64 = null, events_this_connection: u64 = 0, + /// raw frames on the current connection; written by the reader task and + /// read by the idle watchdog, hence atomic (events_this_connection stays + /// plain: it is only read after the reader task is joined). + frames_this_connection: std.atomic.Value(u64) = .init(0), pub fn init(io: Io, allocator: Allocator, options: Options) JetstreamClient { return .{ @@ -251,20 +255,53 @@ pub const JetstreamClient = struct { .client_state = self, }; if (self.options.replay_idle_timeout_ms orelse self.options.idle_timeout_ms) |timeout_ms| { - defer ws_handler.close(); - try client.readTimeout(timeout_ms); + // A socket-level receive timeout (SO_RCVTIMEO) is off the table: + // std.Io.Threaded treats the resulting EAGAIN on a blocking + // socket as a programmer bug and panics in debug builds. Instead + // the blocking readLoop runs as a concurrent task and a watchdog + // here compares the frame counter across quiet windows, + // cancelling the read when a full window passes with no frames. + // Fires between timeout and 2x timeout after the last frame. + // Mirrors firehose.zig. + self.frames_this_connection.store(0, .monotonic); + var done: Io.Event = .unset; + var read_err: ?anyerror = null; + const Reader = struct { + fn run( + io: Io, + ws_client: *websocket.Client, + h: *WsHandler(@TypeOf(handler.*)), + done_ev: *Io.Event, + err_out: *?anyerror, + ) void { + ws_client.readLoop(h) catch |err| { + err_out.* = err; + }; + done_ev.set(io); + } + }; + var future = try self.io.concurrent(Reader.run, .{ self.io, &client, &ws_handler, &done, &read_err }); + var joined = false; + defer if (!joined) future.cancel(self.io); + var last_frames: u64 = 0; while (true) { - const message = try client.read() orelse { - if (self.options.replay_idle_timeout_ms != null) return error.ReplayIdle; - return error.IdleTimeout; + done.waitTimeout(self.io, .{ .duration = .{ .raw = Io.Duration.fromMilliseconds(@intCast(timeout_ms)), .clock = .awake } }) catch |err| switch (err) { + error.Timeout => { + const frames = self.frames_this_connection.load(.monotonic); + if (frames == last_frames) { + if (self.options.replay_idle_timeout_ms != null) return error.ReplayIdle; + return error.IdleTimeout; + } + last_frames = frames; + continue; + }, + error.Canceled => return error.Canceled, }; - defer client.done(message); - switch (message.type) { - .text, .binary => try ws_handler.serverMessage(message.data), - .ping => try client.writePong(@constCast(message.data)), - .close => return, - .pong => {}, - } + // readLoop finished on its own (close frame, socket error, stop) + future.await(self.io); + joined = true; + if (read_err) |err| return err; + return; } } else { try client.readLoop(&ws_handler); @@ -310,6 +347,7 @@ fn WsHandler(comptime H: type) type { const Self = @This(); pub fn serverMessage(self: *Self, data: []const u8) !void { + _ = self.client_state.frames_this_connection.fetchAdd(1, .monotonic); var arena = std.heap.ArenaAllocator.init(self.allocator); defer arena.deinit();