jetstream client
atproto jetstream client
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589//! 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=<id>; 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 internallypub 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}