diff --git a/src/internal/streaming/jetstream.zig b/src/internal/streaming/jetstream.zig index eca06c3..efe0e62 100644 --- a/src/internal/streaming/jetstream.zig +++ b/src/internal/streaming/jetstream.zig @@ -48,6 +48,11 @@ pub const Options = struct { /// frames. Jetstream has no server-side end cursor, so an idle boundary is /// needed when no matching event exists at or beyond `end_cursor`. replay_idle_timeout_ms: ?u32 = null, + /// Fail over after this many milliseconds with no WebSocket frames on a + /// live subscription. Catches a host that holds the socket open but goes + /// silent, which TCP keepalive never fires on. Ignored during finite + /// replay (`replay_idle_timeout_ms` governs there). + idle_timeout_ms: ?u32 = null, max_message_size: usize = 1024 * 1024, }; @@ -143,6 +148,7 @@ pub const JetstreamClient = struct { allocator: Allocator, options: Options, last_time_us: ?i64 = null, + events_this_connection: u64 = 0, pub fn init(io: Io, allocator: Allocator, options: Options) JetstreamClient { return .{ @@ -160,26 +166,27 @@ pub const JetstreamClient = struct { /// optional: fn onError(*@TypeOf(handler), anyerror) void /// optional: fn onConnect(*@TypeOf(handler), []const u8) void — called with host on connect /// blocks forever — reconnects with exponential backoff on disconnect. - /// rotates through hosts on each reconnect attempt. + /// rotates through hosts on failed attempts; after a connection that + /// delivered events drops, retries from hosts[0] so consumers return to + /// the primary instead of sticking to a fallback. pub fn subscribe(self: *JetstreamClient, handler: anytype) Io.Cancelable!void { var backoff: u64 = 1; - var host_index: usize = 0; const max_backoff: u64 = 60; - var prev_host_index: usize = 0; + var current: usize = 0; + var prev: ?usize = null; while (true) { - const host = self.options.hosts[host_index % self.options.hosts.len]; - const effective_index = host_index % self.options.hosts.len; + const host = self.options.hosts[current]; // rewind cursor by 10s on host switch (different instances may lag) - if (host_index > 0 and effective_index != prev_host_index) { + if (prev != null and current != prev.?) { if (self.last_time_us) |t| { self.last_time_us = t - 10_000_000; } backoff = 1; } - log.info("connecting to host {d}/{d}: {s}", .{ effective_index + 1, self.options.hosts.len, host }); + log.info("connecting to host {d}/{d}: {s}", .{ current + 1, self.options.hosts.len, host }); self.connectAndRead(host, handler) catch |err| { if (err == error.EndCursorReached or err == error.ReplayIdle) return; @@ -190,13 +197,20 @@ pub const JetstreamClient = struct { } }; - prev_host_index = effective_index; - host_index += 1; + prev = current; + current = nextHostIndex(current, self.events_this_connection > 0, self.options.hosts.len); try self.io.sleep(Io.Duration.fromSeconds(@intCast(backoff)), .awake); backoff = @min(backoff * 2, max_backoff); } } + /// after a working connection drops, go back to the primary; a failed + /// attempt keeps rotating through the fallbacks. + fn nextHostIndex(current: usize, delivered_events: bool, host_count: usize) usize { + if (delivered_events) return 0; + return (current + 1) % host_count; + } + fn connectAndRead(self: *JetstreamClient, host: []const u8, handler: anytype) !void { var path_buf: [2048]u8 = undefined; const path = try self.buildSubscribePath(&path_buf); @@ -223,16 +237,21 @@ pub const JetstreamClient = struct { handler.onConnect(host); } + self.events_this_connection = 0; + var ws_handler = WsHandler(@TypeOf(handler.*)){ .allocator = self.allocator, .handler = handler, .client_state = self, }; - if (self.options.replay_idle_timeout_ms) |timeout_ms| { + if (self.options.replay_idle_timeout_ms orelse self.options.idle_timeout_ms) |timeout_ms| { defer ws_handler.close(); try client.readTimeout(timeout_ms); while (true) { - const message = try client.read() orelse return error.ReplayIdle; + const message = try client.read() orelse { + if (self.options.replay_idle_timeout_ms != null) return error.ReplayIdle; + return error.IdleTimeout; + }; defer client.done(message); switch (message.type) { .text, .binary => try ws_handler.serverMessage(message.data), @@ -298,6 +317,7 @@ fn WsHandler(comptime H: type) type { } self.client_state.last_time_us = event.timeUs(); + self.client_state.events_this_connection += 1; self.handler.onEvent(event); } @@ -636,6 +656,33 @@ test "round-robin cycles through hosts" { } } +test "failback: a working connection's drop retries the primary first" { + // fallback delivered events → next attempt is hosts[0] + try std.testing.expectEqual(@as(usize, 0), JetstreamClient.nextHostIndex(2, true, 3)); + // failed attempt keeps rotating + try std.testing.expectEqual(@as(usize, 1), JetstreamClient.nextHostIndex(0, false, 3)); + try std.testing.expectEqual(@as(usize, 0), JetstreamClient.nextHostIndex(2, false, 3)); + // primary itself dropping reconnects to primary + try std.testing.expectEqual(@as(usize, 0), JetstreamClient.nextHostIndex(0, true, 3)); +} + +test "serverMessage counts events per connection" { + const Handler = struct { + pub fn onEvent(_: *@This(), _: Event) void {} + }; + var client = JetstreamClient.init(std.Options.debug_io, std.testing.allocator, .{}); + var handler = Handler{}; + var ws_handler = WsHandler(Handler){ + .allocator = std.testing.allocator, + .handler = &handler, + .client_state = &client, + }; + try ws_handler.serverMessage( + \\{"did":"did:plc:a","time_us":100,"kind":"commit","commit":{"operation":"create","collection":"x","rkey":"1"}} + ); + try std.testing.expectEqual(@as(u64, 1), client.events_this_connection); +} + test "options default hosts are used" { const opts = Options{}; try std.testing.expectEqual(@as(usize, 12), opts.hosts.len);