//! subscribeEvents live client — the v2 live tail. //! //! ports upstream live.go (bluesky-social/jetstream, pin 289b032): //! dial /xrpc/network.bsky.jetstream.subscribeEvents offering the //! xrpc.v1.json subprotocol, decode lexicon frames (livedecode.zig), //! deduplicate the at-least-once reconnect overlap by seq, and reconnect //! with bounded exponential backoff. reconnects resume at the highest //! delivered seq, so events produced while disconnected are replayed by //! the server instead of silently dropped. //! //! pre-upgrade rejections are typed exactly like upstream (the server //! answers the upgrade with an HTTP 400 xrpc error envelope, surfaced via //! websocket.zig's Client.handshake_failure): //! - CursorTooOld is TERMINAL: the seq will not become valid by //! retrying (the lookback floor only advances). returned to the //! caller so the cutover engine re-enters archive backfill. //! - InvalidRequest is PERMANENT: the same immutable filter cannot //! succeed on reconnect; returned as fatal. //! - UnknownZstdDictionary is RECOVERABLE in place: the server rotated //! its dictionary; refetch it (or degrade to uncompressed) and //! reconnect — never 400-loop on an ID the server keeps refusing. //! //! dict-zstd compression is negotiated at the application layer //! (?zstdDictionary=; event frames arrive as BINARY zstd frames), and //! the decompressed size is bounded by the read limit. v2 never //! negotiates permessage-deflate (upstream removed it server-side). const std = @import("std"); const zat = @import("zat"); const filter = @import("filter.zig"); const websocket = @import("websocket"); const livedecode = @import("livedecode.zig"); const zstd = @import("zstd.zig"); const Io = std.Io; const Allocator = std.mem.Allocator; const posix = std.posix; const log = std.log.scoped(.jetstream); pub const subscribe_nsid = livedecode.nsid; pub const subprotocol = "xrpc.v1.json"; pub const timestamp_cursor_min: u64 = 1_000_000_000_000_000; /// bounds a single websocket message; also caps a compressed frame's /// decompressed size (upstream defaultLiveReadLimit) pub const default_read_limit: usize = 32 << 20; const default_backoff_min_ns: u64 = 250 * std.time.ns_per_ms; const default_backoff_max_ns: u64 = 30 * std.time.ns_per_s; pub const Options = struct { /// base URL of the instance, e.g. "https://jetstream1.us-east.bsky.network" /// (plain http for loopback tests) host: []const u8, /// wire resume point. null = start at the live tip (the cursor param is /// omitted; upstream WithLiveCursor(0)); 0 = replay from the beginning /// of the server's retention; a positive seq resumes inclusively — the /// client's own seq dedup turns that into "> last delivered". /// Values >= timestamp_cursor_min are unix microseconds, resolved by the server. cursor: ?u64 = null, /// highest seq the caller already holds: the at-least-once re-delivery /// at or below it is dropped. 0 = nothing delivered yet, so the first /// real event (seq >= 1) always passes. dedup_floor: u64 = 0, kinds: []const livedecode.Kind = &.{}, collections: []const []const u8 = &.{}, dids: []const []const u8 = &.{}, read_limit: usize = default_read_limit, backoff_min_ns: u64 = default_backoff_min_ns, backoff_max_ns: u64 = default_backoff_max_ns, /// dictionary blob from getZstdDictionary; opts into dict-zstd frames. /// null = plain uncompressed text frames. zstd_dict: ?[]const u8 = null, /// re-fetches the CURRENT dictionary after the server rejects the /// pinned ID (rotation). null disables in-place recovery; the client /// 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 pub const RunError = error{ /// pre-upgrade 400 CursorTooOld: re-enter backfill from lastSeq() CursorTooOld, /// pre-upgrade 400 InvalidRequest: the filter can never succeed InvalidRequest, Canceled, OutOfMemory, }; pub const LiveClient = struct { io: Io, allocator: Allocator, options: Options, /// highest seq delivered; the dedup floor and the reconnect resume /// cursor. read after run() returns (the cutover engine resumes a /// re-backfill from it). last_seq: u64 = 0, seen_any: bool = false, decoder: ?zstd.DictDecoder = null, /// owned copy of a refetched dictionary blob (rotation recovery) owned_dict: ?[]u8 = null, pub fn init(io: Io, allocator: Allocator, options: Options) LiveClient { const cursor = options.cursor orelse 0; var self: LiveClient = .{ .io = io, .allocator = allocator, .options = options, .last_seq = @max(options.dedup_floor, if (cursor < timestamp_cursor_min) cursor else 0), }; if (options.zstd_dict) |dict| { self.decoder = zstd.DictDecoder.init(dict, options.read_limit) catch blk: { // the blob came from getZstdDictionary moments ago; a parse // failure is a server/transport fault, not a reason to // crash. degrade to uncompressed (documented, logged). log.warn("invalid zstd dictionary; live tail will be uncompressed", .{}); break :blk null; }; } return self; } pub fn deinit(self: *LiveClient) void { if (self.decoder) |*d| d.deinit(); if (self.owned_dict) |owned| self.allocator.free(owned); self.* = undefined; } pub fn lastSeq(self: *const LiveClient) u64 { return self.last_seq; } /// tail the live stream, invoking handler.onEvent for each decoded /// event in delivery order. handler contract: /// fn onEvent(*H, livedecode.Event) bool — false stops the tail /// optional fn onError(*H, anyerror) bool — recoverable read/decode/ /// reconnect errors; false stops (default: log and continue) /// optional fn onInfo(*H, livedecode.Info) void — #info advisories /// returns null on a clean stop (handler asked), typed RunError on a /// terminal pre-upgrade rejection or cancellation. pub fn run(self: *LiveClient, handler: anytype) RunError!void { var backoff = self.options.backoff_min_ns; while (true) { const seq_before = self.last_seq; const outcome = self.session(handler) catch |err| switch (err) { error.Canceled => return error.Canceled, error.OutOfMemory => return error.OutOfMemory, // transient transport/dial/read failure: reconnect below else => SessionEnd{ .transient = err }, }; switch (outcome) { .stopped => return, .cursor_too_old => return error.CursorTooOld, .invalid_request => return error.InvalidRequest, .dict_rejected => { self.refreshDict(); if (!self.emitError(handler, error.UnknownZstdDictionary)) return; }, .transient => |err| { if (!self.emitError(handler, err)) return; }, } // a session that delivered new events was healthy: reset the // backoff so a long-lived connection that finally drops // reconnects promptly, not at the accumulated max if (self.last_seq != seq_before) backoff = self.options.backoff_min_ns; self.io.sleep(Io.Duration.fromNanoseconds(@intCast(backoff)), .awake) catch return error.Canceled; backoff = @min(backoff * 2, self.options.backoff_max_ns); } } const SessionEnd = union(enum) { stopped, cursor_too_old, invalid_request, dict_rejected, transient: anyerror, }; fn emitError(self: *LiveClient, handler: anytype, err: anyerror) bool { _ = self; if (comptime @hasDecl(@TypeOf(handler.*), "onError")) return handler.onError(err); log.warn("live tail reconnecting: {s}", .{@errorName(err)}); return true; } fn session(self: *LiveClient, handler: anytype) !SessionEnd { const target = try parseHost(self.options.host); var path_buf: [4096]u8 = undefined; const path = try self.subscribePath(&path_buf); var client = try websocket.Client.init(self.io, self.allocator, .{ .host = target.host, .port = target.port, .tls = target.tls, .max_size = self.options.read_limit, }); defer client.deinit(); var headers_buf: [1024]u8 = undefined; const headers = try std.fmt.bufPrint( &headers_buf, "Host: {s}\r\nSec-WebSocket-Protocol: {s}\r\n", .{ target.host, subprotocol }, ); client.handshake(path, .{ .headers = headers }) catch |err| { if (err == error.InvalidHandshakeResponse) { if (client.handshake_failure) |failure| return classifyRejection(failure); } return err; }; configureKeepalive(&client); if (comptime @hasDecl(@TypeOf(handler.*), "onConnect")) handler.onConnect(target.host); var ws_handler: WsHandler(@TypeOf(handler.*)) = .{ .client = self, .handler = handler }; self.readFrames(&client, &ws_handler) catch |err| { if (err == error.Canceled) return error.Canceled; // handler-requested stop falls through to the flag checks // below (previously it hit the transient arm, so a live-phase // onEvent=false reconnect-looped instead of stopping) if (err != error.Stop) return .{ .transient = err }; }; if (ws_handler.stopped) return .stopped; if (ws_handler.failure) |err| switch (err) { error.OutOfMemory => return error.OutOfMemory, else => return .{ .transient = err }, }; // the server closed (possibly right after a terminal error frame, // which the handler surfaced through onError already) return .{ .transient = error.ConnectionClosed }; } /// frame pump. without a stall timeout this is the library's readLoop /// on the caller's thread; with one, the readLoop runs in a concurrent /// task while this thread watches EVENT progress (last_seq) — 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. socket read timeouts are NOT used: /// Io.Threaded treats the resulting EAGAIN as a programmer bug. 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); const Pump = struct { ws: *websocket.Client, handler: @TypeOf(ws_handler), done: std.atomic.Value(bool) = .init(false), err: ?anyerror = null, fn run(pump: *@This()) void { defer pump.done.store(true, .release); pump.ws.readLoop(pump.handler) catch |err| { pump.err = err; }; } }; var pump = Pump{ .ws = client, .handler = ws_handler }; var future = try self.io.concurrent(Pump.run, .{&pump}); var joined = false; defer if (!joined) future.cancel(self.io); const poll_ns = stallPollNs(stall_ns); var last_seq_seen = self.last_seq; var quiet_ns: u64 = 0; while (true) { try self.io.sleep(Io.Duration.fromNanoseconds(@intCast(poll_ns)), .awake); if (pump.done.load(.acquire)) { joined = true; future.await(self.io); if (pump.err) |err| return err; return; } if (self.last_seq != last_seq_seen) { last_seq_seen = self.last_seq; quiet_ns = 0; } else { quiet_ns += poll_ns; if (quiet_ns >= stall_ns) { joined = true; future.cancel(self.io); return error.LiveStalled; } } } } /// poll interval: a quarter of the stall window, clamped to [50ms, 5s] fn stallPollNs(stall_ns: u64) u64 { return @min(@max(stall_ns / 4, 50 * std.time.ns_per_ms), 5 * std.time.ns_per_s); } /// 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/503 during a deploy reconnect-loops like an abrupt close. fn classifyRejection(failure: websocket.Client.HandshakeFailure) SessionEnd { if (failure.status == 400) { var name_buf: [1024]u8 = undefined; const name = xrpcErrorName(failure.body(), &name_buf) orelse ""; if (std.mem.eql(u8, name, "CursorTooOld")) return .cursor_too_old; if (std.mem.eql(u8, name, "UnknownZstdDictionary")) return .dict_rejected; if (std.mem.eql(u8, name, "InvalidRequest")) return .invalid_request; } return .{ .transient = error.HandshakeRejected }; } fn xrpcErrorName(body: []const u8, buf: []u8) ?[]const u8 { var fba = std.heap.FixedBufferAllocator.init(buf); // bounded parse of {"error": name, ...}; allocation failure or // malformed JSON just means "no structured name" const parsed = std.json.parseFromSliceLeaky(std.json.Value, fba.allocator(), body, .{}) catch return null; if (parsed != .object) return null; const v = parsed.object.get("error") orelse return null; return if (v == .string) v.string else null; } /// recover from a server-side dictionary rotation: refetch the current /// dictionary and swap the decoder so the next dial negotiates the new /// ID. when the refetch is unavailable, fails, or returns the very ID /// the server just rejected (mixed-version fleet), degrade to an /// uncompressed tail — compression is an optimization; the tail must /// keep flowing. (upstream refreshDict) fn refreshDict(self: *LiveClient) void { var current = self.decoder orelse return; const rejected = current.id; if (self.options.refetch_dict) |refetch| { if (refetch(self.options.refetch_dict_ctx, self.allocator)) |blob| { if (zstd.DictDecoder.init(blob, self.options.read_limit)) |fresh| { if (fresh.id != rejected) { current.deinit(); if (self.owned_dict) |owned| self.allocator.free(owned); self.owned_dict = blob; self.decoder = fresh; log.info("live zstd dictionary rotated; refetched (rejected id {d}, new id {d})", .{ rejected, self.decoder.?.id }); return; } var discard = fresh; discard.deinit(); } else |_| {} self.allocator.free(blob); } } current.deinit(); self.decoder = null; log.warn("live zstd dictionary rejected and refetch unavailable; continuing uncompressed (rejected id {d})", .{rejected}); } fn subscribePath(self: *LiveClient, buf: []u8) ![]const u8 { var w: Io.Writer = .fixed(buf); try w.writeAll("/xrpc/" ++ subscribe_nsid); var sep: u8 = '?'; // wire cursor: once any event has been delivered, resume at the // highest delivered seq — this is what keeps a reconnect from // re-anchoring at the tip and silently dropping events produced // while disconnected. before any delivery, use the configured // start: omit when null (live from tip), else send it. if (self.seen_any) { try w.print("{c}cursor={d}", .{ sep, self.last_seq }); sep = '&'; } else if (self.options.cursor) |cursor| { try w.print("{c}cursor={d}", .{ sep, cursor }); sep = '&'; } for (self.options.kinds) |kind| { try w.print("{c}kinds={s}", .{ sep, @tagName(kind) }); sep = '&'; } for (self.options.collections) |collection| { try w.print("{c}collections={s}", .{ sep, collection }); sep = '&'; } for (self.options.dids) |did| { try w.print("{c}dids={s}", .{ sep, did }); sep = '&'; } if (self.decoder) |decoder| { try w.print("{c}zstdDictionary={d}", .{ sep, decoder.id }); sep = '&'; } return w.buffered(); } fn WsHandler(comptime H: type) type { return struct { client: *LiveClient, handler: *H, stopped: bool = false, failure: ?anyerror = null, const Self = @This(); pub fn serverMessage(self: *Self, data: []const u8, message_type: enum { text, binary }) !void { const client = self.client; var frame: []const u8 = data; var decompressed: ?[]u8 = null; defer if (decompressed) |owned| client.allocator.free(owned); switch (message_type) { .binary => { // dict-zstd connection: every event frame is a // BINARY zstd frame. on an uncompressed connection // stray binary is ignored (jetstream frames are // text JSON). const decoder = if (client.decoder) |*d| d else return; decompressed = decoder.decompressAlloc(client.allocator, data) catch |err| { // upstream input, never crash: surface and keep the tail if (!client.emitError(self.handler, err)) { self.stopped = true; return error.Stop; } return; }; frame = decompressed.?; }, .text => {}, } var arena = std.heap.ArenaAllocator.init(client.allocator); defer arena.deinit(); const decoded = livedecode.decodeFrame(arena.allocator(), frame) catch |err| { if (err == error.OutOfMemory) { self.failure = error.OutOfMemory; return error.Stop; } // a malformed data frame is upstream input; surface it // but keep the connection (one bad frame must not drop // the tail) if (!client.emitError(self.handler, err)) { self.stopped = true; return error.Stop; } return; }; switch (decoded) { .skip => {}, .info => |info| { if (comptime @hasDecl(H, "onInfo")) { self.handler.onInfo(info); } else { log.info("live stream info frame: {s}: {s}", .{ info.name, info.message }); } }, .stream_error => |stream_err| { // terminal error frame: the server closes right // after sending it. surface the typed reason; the // reconnect loop handles the close that follows. log.warn("live stream error frame: {s}: {s}", .{ stream_err.code, stream_err.message }); if (!client.emitError(self.handler, error.StreamErrorFrame)) { self.stopped = true; return error.Stop; } }, .event => |event| { // deduplicate the at-least-once reconnect overlap: // skip anything at or below the highest delivered // seq. last_seq 0 with nothing delivered means the // first real event (seq >= 1) passes. if (event.seq <= client.last_seq) return; client.last_seq = event.seq; client.seen_any = true; if (!filter.wants(event, client.options.kinds, client.options.dids, client.options.collections)) return; if (!self.handler.onEvent(event)) { self.stopped = true; return error.Stop; } }, } } pub fn close(_: *Self) void {} }; } }; const Target = struct { host: []const u8, port: u16, tls: bool, }; fn parseHost(base: []const u8) !Target { const uri = std.Uri.parse(base) catch return error.InvalidHost; const tls = std.ascii.eqlIgnoreCase(uri.scheme, "https") or std.ascii.eqlIgnoreCase(uri.scheme, "wss"); if (!tls and !std.ascii.eqlIgnoreCase(uri.scheme, "http") and !std.ascii.eqlIgnoreCase(uri.scheme, "ws")) return error.InvalidHost; const host_component = uri.host orelse return error.InvalidHost; const host = switch (host_component) { .raw => |raw| raw, .percent_encoded => |enc| enc, }; if (host.len == 0) return error.InvalidHost; return .{ .host = host, .port = uri.port orelse (if (tls) @as(u16, 443) else 80), .tls = tls }; } fn configureKeepalive(client: *websocket.Client) void { // TCP keepalive catches half-open sockets the read loop would block on // forever; mirrors zat.JetstreamClient (10s idle, 5s interval, 2 probes) const fd = client.stream.stream.socket.handle; const builtin = @import("builtin"); posix.setsockopt(fd, posix.SOL.SOCKET, posix.SO.KEEPALIVE, &std.mem.toBytes(@as(i32, 1))) catch return; const tcp: i32 = @intCast(posix.IPPROTO.TCP); if (builtin.os.tag == .linux) { posix.setsockopt(fd, tcp, posix.TCP.KEEPIDLE, &std.mem.toBytes(@as(i32, 10))) catch return; } else if (builtin.os.tag == .macos) { posix.setsockopt(fd, tcp, posix.TCP.KEEPALIVE, &std.mem.toBytes(@as(i32, 10))) catch return; } posix.setsockopt(fd, tcp, posix.TCP.KEEPINTVL, &std.mem.toBytes(@as(i32, 5))) catch return; posix.setsockopt(fd, tcp, posix.TCP.KEEPCNT, &std.mem.toBytes(@as(i32, 2))) catch return; } // === tests === const testing = std.testing; test "subscribe path: cursor semantics, filters, and dictionary id" { var client = LiveClient.init(std.Options.debug_io, testing.allocator, .{ .host = "https://example.test", .cursor = null, }); defer client.deinit(); var buf: [4096]u8 = undefined; // from tip: no cursor param at all try testing.expectEqualStrings("/xrpc/" ++ subscribe_nsid, try client.subscribePath(&buf)); // configured cursor before any delivery client.options.cursor = 0; try testing.expectEqualStrings("/xrpc/" ++ subscribe_nsid ++ "?cursor=0", try client.subscribePath(&buf)); // after delivery, reconnects anchor at the highest delivered seq even // when the start was from-tip client.options.cursor = null; client.last_seq = 41; client.seen_any = true; try testing.expectEqualStrings("/xrpc/" ++ subscribe_nsid ++ "?cursor=41", try client.subscribePath(&buf)); client.options.kinds = &.{ .commit, .identity }; client.options.collections = &.{"app.bsky.feed.post"}; client.options.dids = &.{"did:plc:abc"}; try testing.expectEqualStrings( "/xrpc/" ++ subscribe_nsid ++ "?cursor=41&kinds=commit&kinds=identity&collections=app.bsky.feed.post&dids=did:plc:abc", try client.subscribePath(&buf), ); } test "rejection classification matches the wire contract by error name" { const mk = struct { fn failure(status: u16, body: []const u8) websocket.Client.HandshakeFailure { var f: websocket.Client.HandshakeFailure = .{ .status = status }; @memcpy(f.body_buf[0..body.len], body); f.body_len = body.len; return f; } }; try testing.expectEqual( LiveClient.SessionEnd.cursor_too_old, LiveClient.classifyRejection(mk.failure(400, "{\"error\":\"CursorTooOld\",\"message\":\"below floor 9\"}")), ); try testing.expectEqual( LiveClient.SessionEnd.dict_rejected, LiveClient.classifyRejection(mk.failure(400, "{\"error\":\"UnknownZstdDictionary\"}")), ); try testing.expectEqual( LiveClient.SessionEnd.invalid_request, LiveClient.classifyRejection(mk.failure(400, "{\"error\":\"InvalidRequest\",\"message\":\"kinds\"}")), ); // name matching, never substrings: a 400 whose message MENTIONS // CursorTooOld is not a cursor rejection const mention = LiveClient.classifyRejection(mk.failure(400, "{\"error\":\"Nope\",\"message\":\"CursorTooOld\"}")); try testing.expect(mention == .transient); // proxy 503 during a deploy is transient, reconnect-loop territory const proxy = LiveClient.classifyRejection(mk.failure(503, "starting")); try testing.expect(proxy == .transient); } test "stallPollNs: quarter window clamped to [50ms, 5s]" { try std.testing.expectEqual(@as(u64, 50 * std.time.ns_per_ms), LiveClient.stallPollNs(100 * std.time.ns_per_ms)); // tiny window: floor try std.testing.expectEqual(@as(u64, 5 * std.time.ns_per_s), LiveClient.stallPollNs(600 * std.time.ns_per_s)); // huge window: ceiling try std.testing.expectEqual(@as(u64, 2_500 * std.time.ns_per_ms), LiveClient.stallPollNs(10 * std.time.ns_per_s)); // 10s window: 2.5s poll }