diff --git a/src/internal/streaming/firehose.zig b/src/internal/streaming/firehose.zig index 1b62625..72d47ab 100644 --- a/src/internal/streaming/firehose.zig +++ b/src/internal/streaming/firehose.zig @@ -715,6 +715,11 @@ pub const FirehoseClient = struct { allocator: Allocator, options: Options, last_seq: ?i64 = null, + /// set once the websocket handshake for the current attempt succeeds. + /// An endpoint that accepted a connection is a working endpoint, however + /// the connection later ended, so the reconnect backoff starts over + /// rather than compounding. + connected_this_attempt: bool = false, pub fn init(io: Io, allocator: Allocator, options: Options) FirehoseClient { return .{ @@ -766,6 +771,23 @@ pub const FirehoseClient = struct { if (host_index > 0 and effective_index != prev_host_index) { backoff = 1; } + // ...and after a connection that was actually established. Backoff + // exists to spare an endpoint that is refusing us; an endpoint that + // accepted a connection and later closed it is not that. Without + // this a single-host consumer never resets — effective_index is + // always 0, so the host-switch branch above can never fire — and + // the delay climbs to max_backoff and stays there for the life of + // the process. Observed on a relay that closed every ~113s: every + // reconnect paid the full 60s, leaving the tail down 33% of the + // time while each individual connection was perfectly healthy. + // + // Keyed on the handshake, matching jcalabro/atmos + // (streaming/client.go: "Reset backoff after successful + // connection"). Keying it on delivered frames instead would leave + // a connection that is open but idle compounding its backoff — + // which for a filtered jetstream consumer is an ordinary state, + // not a fault. + if (self.connected_this_attempt) backoff = 1; log.info("connecting to host {d}/{d}: {s}", .{ effective_index + 1, self.options.hosts.len, host }); @@ -798,6 +820,7 @@ pub const FirehoseClient = struct { } fn connectAndRead(self: *FirehoseClient, host: []const u8, handler: anytype) !void { + self.connected_this_attempt = false; const ep = try Endpoint.parse(host); var path_buf: [256]u8 = undefined; @@ -829,6 +852,7 @@ pub const FirehoseClient = struct { try client.handshake(path, .{ .headers = host_header }); configureKeepalive(&client); + self.connected_this_attempt = true; log.info("firehose connected to {s}", .{ep.host}); if (comptime @hasDecl(@TypeOf(handler.*), "onConnect")) { diff --git a/src/internal/streaming/jetstream.zig b/src/internal/streaming/jetstream.zig index efe0e62..6881ee2 100644 --- a/src/internal/streaming/jetstream.zig +++ b/src/internal/streaming/jetstream.zig @@ -197,6 +197,12 @@ pub const JetstreamClient = struct { } }; + // A connection that delivered events was a working endpoint; the + // same signal that sends us back to the primary host also clears + // the backoff, so a single-host consumer does not compound its + // delay to max_backoff and stay there. See firehose.zig. + if (self.events_this_connection > 0) backoff = 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);