jetstream client
atproto jetstream client
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529//! the unified Jetstream v2 client — live tail, filtered archive replay,//! and gapless replay→live cutover through one handler.//!//! ports upstream client_core.go runBackfillThenLive (bluesky-social///! jetstream, pin 289b032, design §11/§13/§14)://!//! 1. sweep the sealed archive: page planSnapshot and deliver the whole//! sealed range (cursor, plannedThroughSeq] in seq order//! (ArchiveBackfill with the extended onRow handler — every matching//! row including identity/account/sync markers, each with its seq);//! 2. connect subscribeEvents ONCE at cursor = the sealed tip — no rewind//! margin, no client buffer. the live client's seq dedup makes the//! seam at-least-once with no gap; segments sealed during the download//! are covered by the server's cold replay.//!//! a pre-upgrade 400 CursorTooOld at connect (slow handoff / fell off//! live) is NOT fatal: re-enter archive pagination from the last durably//! processed seq. cutover = max(sealed tip, cursor) keeps the dedup floor//! and resume monotonic non-decreasing (a live tail routinely delivers//! events past the sealed tip; a re-learned tip below the cursor must not//! regress it). re-backfill cycles that fail to advance the cursor are//! bounded (upstream maxRebackfillStalls = 5) — a non-advancing cycle is a//! pathological loop, not a real fall-behind.//!//! === multi-host failover (live tail only) ===//!//! v2 seqs are per-instance. Live failover sends the last witnessed time//! minus the rewind margin directly as the next host's timestamp cursor.//! The server resolves it; no archive requests or API key are needed.//! Sequence deduplication resets on host rotation. Subsequent reconnects//! use the receiving host's delivered sequences.//!//! The overlap tolerates witness-time skew, not arbitrary coverage gaps.//! Retention clamps surface through onInfo (OutdatedCursor), or logging//! without a callback. A caller's sequence cannot transfer before any//! event supplies a witnessed-time anchor; that fails with HostUnhealthy.
const std = @import("std");const zat = @import("zat");const archive = @import("archive_backfill.zig");const live_mod = @import("live.zig");const livedecode = @import("livedecode.zig");
const Io = std.Io;const Allocator = std.mem.Allocator;const log = std.log.scoped(.jetstream);
/// upstream maxRebackfillStalls: consecutive too-old cycles that made no/// cursor progress before the stream is declared pathologicalpub const max_rebackfill_stalls = 5;
pub const Options = struct { /// instance base URLs, e.g. "https://stream.waow.tech". the first is /// the primary; the rest are rotated through when the current host is /// declared unhealthy, and failback returns to hosts[0] — the same /// list-with-preferred-primary contract as zat's v1 JetstreamClient /// and FirehoseClient. cross-instance resume re-anchors by witnessed /// time (see module doc); a single host keeps the exact single-host /// error surface. hosts: []const []const u8, /// consecutive live transient failures WITH NO EVENTS BETWEEN before /// the current host is declared unhealthy and rotation kicks in. /// only meaningful with multiple hosts; 0 disables rotation. failover_after_errors: u8 = 5, /// witnessed-time rewind margin when re-anchoring on a new host, /// covering inter-instance witness skew (v1 client rewinds 10s) failover_rewind_us: i64 = 10 * std.time.us_per_s, /// treat a live session with no EVENTS for this long as failed /// (error.LiveStalled → counts toward failover_after_errors). catches /// the connected-but-silent failure mode that produces no transport /// errors at all (the 2026-08-16 bsky.network front wedge). opt-in: /// a heavily filtered stream can be legitimately quiet for hours. failover_stall_ns: u64 = 0, /// replay start, exclusive: sweep the sealed archive from here and cut /// over to live. null = no replay — a pure live tail from the tip /// (set live_cursor to resume a live tail at a known seq instead). after_seq: ?u64 = null, /// stop after the sealed range instead of cutting over to live /// (upstream WithSnapshotOnly). only meaningful with after_seq. snapshot_only: bool = false, /// live-only sequence or unix-microsecond timestamp; null = from the tip live_cursor: ?u64 = null, kinds: []const livedecode.Kind = &.{}, collections: []const []const u8 = &.{}, dids: []const []const u8 = &.{}, /// Bearer credential for archive requests only. Live connections are public. api_key: ?[]const u8 = null, /// block fetches in flight during the archive sweep (see ArchiveBackfill) concurrency: usize = 4, /// Dict-zstd live frames are enabled by default: the current dictionary is fetched /// via getZstdDictionary; a fetch failure degrades to an uncompressed /// tail (compression is an optimization, never a stream failure) zstd_compression: bool = true, /// live reconnect backoff overrides (tests); zero = defaults backoff_min_ns: u64 = 0, backoff_max_ns: u64 = 0,};
/// cross-host delivery state: the time anchor rotation re-enters from,/// and the failback signalconst Progress = struct { /// witnessed time of the newest delivered event, 0 = nothing yet last_time_us: i64 = 0, /// events delivered since the last rotation (failback: retry the /// primary first, matching the v1 client) delivered_since_rotate: bool = false,};
/// v1 client's rotation rule (zat jetstream.zig nextHostIndex): a healthy/// connection's drop retries the primary first; a failed attempt keeps/// rotatingpub fn nextHostIndex(current: usize, delivered_events: bool, host_count: usize) usize { if (delivered_events) return 0; return (current + 1) % host_count;}
const HostOutcome = enum { stopped, rotate };
/// run the unified stream through `handler`:/// fn onEvent(*H, livedecode.Event) bool — every event, archive and live,/// in seq order across the seam; false stops cleanly/// optional fn onError(*H, anyerror) bool — recoverable archive/live errors;/// true continues, false stops cleanly. Without it, archive errors return./// optional fn onInfo(*H, livedecode.Info) void/// blocks until the handler stops it, the snapshot completes/// (snapshot_only), or a terminal error.pub fn subscribe(io: Io, allocator: Allocator, options: Options, handler: anytype) !void { if (options.hosts.len == 0) return error.NoHosts; if (options.failover_rewind_us < 0) return error.InvalidRequest; const host_count = options.hosts.len; const can_rotate = host_count > 1 and options.failover_after_errors > 0;
var progress = Progress{}; var host_idx: usize = 0; var after = options.after_seq; var live_cur = options.live_cursor; var time_floor: i64 = 0;
const rot_min: u64 = if (options.backoff_min_ns > 0) options.backoff_min_ns else 250 * std.time.ns_per_ms; const rot_max: u64 = if (options.backoff_max_ns > 0) options.backoff_max_ns else 30 * std.time.ns_per_s; var rot_backoff = rot_min;
while (true) { const host = options.hosts[host_idx]; const outcome = try runOnHost(io, allocator, options, host, after, live_cur, time_floor, can_rotate, &progress, handler); switch (outcome) { .stopped => return, .rotate => { // no time anchor yet + a caller-supplied seq anchor: that // anchor is bound to this host's seq space and cannot be // mapped — fail loudly instead of silently restarting // somewhere else (module doc) if (progress.last_time_us == 0 and (after != null or (live_cur != null and live_cur.? < live_mod.timestamp_cursor_min))) return error.HostUnhealthy;
const delivered = progress.delivered_since_rotate; const prev = host_idx; host_idx = nextHostIndex(host_idx, delivered, host_count); progress.delivered_since_rotate = false; rot_backoff = if (delivered) rot_min else @min(rot_backoff * 2, rot_max); log.warn("host {s} unhealthy; failing over to {s}", .{ options.hosts[prev], options.hosts[host_idx] });
if (progress.last_time_us > 0) { const floor = progress.last_time_us -| options.failover_rewind_us; if (floor < live_mod.timestamp_cursor_min) return error.InvalidRequest; after = null; live_cur = @intCast(floor); time_floor = floor; } io.sleep(Io.Duration.fromNanoseconds(@intCast(rot_backoff)), .awake) catch return error.Canceled; }, } }}
/// the pre-failover single-host flow, verbatim, parameterized by host and/// anchor. can_rotate=false propagates every error exactly as the/// original did; can_rotate=true converts live-tail host-health errors/// (transient threshold, stall) into .rotate. Archive errors never request/// rotation: accepted entry/decode errors continue on the same host, while/// plan and resource failures propagate. Archives are not interchangeable.fn runOnHost( io: Io, allocator: Allocator, options: Options, host: []const u8, after_seq: ?u64, live_cursor: ?u64, time_floor_us: i64, can_rotate: bool, progress: *Progress, handler: anytype,) !HostOutcome { var dict_buf: ?[]u8 = null; defer if (dict_buf) |owned| allocator.free(owned); if (options.zstd_compression) dict_buf = fetchZstdDict(io, allocator, host);
var dict_source = DictSource{ .io = io, .host = host }; const H = @TypeOf(handler.*);
if (after_seq == null) { // pure live tail (upstream runLiveOnly) var client = live_mod.LiveClient.init(io, allocator, liveOptions(options, host, live_cursor, 0, dict_buf)); defer client.deinit(); client.options.refetch_dict = DictSource.fetch; client.options.refetch_dict_ctx = &dict_source; var filtering = FilteringHandler(H){ .inner = handler, .kinds = options.kinds, .time_floor_us = time_floor_us, .progress = progress, .rotate_threshold = if (can_rotate) options.failover_after_errors else 0, }; try client.run(&filtering); return if (filtering.rotate_pending) .rotate else .stopped; }
var cursor = after_seq.?; var stalls: u8 = 0; while (true) { var sweeper = SweepHandler(H){ .inner = handler, .kinds = options.kinds, .progress = progress, }; const result = try archive.run(io, allocator, .{ .host = host, .collections = options.collections, .dids = options.dids, .kinds = options.kinds, .api_key = options.api_key, .after_seq = cursor, .concurrency = options.concurrency, }, &sweeper); if (result.stopped) return .stopped; if (options.snapshot_only) return .stopped;
// cut over at the HIGHER of the freshly-learned sealed tip and the // cursor already processed through (monotonic non-decreasing; see // module doc) const cutover = @max(result.planned_through_seq, cursor); var client = live_mod.LiveClient.init(io, allocator, liveOptions(options, host, cutover, cutover, dict_buf)); defer client.deinit(); client.options.refetch_dict = DictSource.fetch; client.options.refetch_dict_ctx = &dict_source; var filtering = FilteringHandler(H){ .inner = handler, .kinds = options.kinds, .time_floor_us = time_floor_us, .progress = progress, .rotate_threshold = if (can_rotate) options.failover_after_errors else 0, }; client.run(&filtering) catch |err| switch (err) { error.CursorTooOld => { // §14: re-backfill from the last durably-processed seq, // requiring strict cursor progress within the stall bound const resume_seq = @max(client.lastSeq(), cutover); if (resume_seq <= cursor) { stalls += 1; if (stalls >= max_rebackfill_stalls) { log.err("re-backfill made no progress after {d} cursor-too-old cycles at seq {d}", .{ stalls, resume_seq }); return error.RebackfillStalled; } } else { stalls = 0; } cursor = resume_seq; continue; }, else => |e| return e, }; return if (filtering.rotate_pending) .rotate else .stopped; }}
fn liveOptions(options: Options, host: []const u8, cursor: ?u64, dedup_floor: u64, dict: ?[]const u8) live_mod.Options { var out: live_mod.Options = .{ .host = host, .cursor = cursor, .dedup_floor = dedup_floor, .kinds = options.kinds, .collections = options.collections, .dids = options.dids, .zstd_dict = dict, .stall_timeout_ns = options.failover_stall_ns, }; if (options.backoff_min_ns > 0) out.backoff_min_ns = options.backoff_min_ns; if (options.backoff_max_ns > 0) out.backoff_max_ns = options.backoff_max_ns; return out;}
/// archive-half adapter: receives every row (with seq) from the sweep and/// applies the kinds filter — the plan prunes server-side where it can,/// but the client-side matcher remains the correctness backstop.fn SweepHandler(comptime H: type) type { return struct { inner: *H, kinds: []const livedecode.Kind, progress: *Progress,
pub fn onError(self: *@This(), err: anyerror) !bool { if (comptime @hasDecl(H, "onError")) return self.inner.onError(err); return err; }
pub fn onRow(self: *@This(), event: livedecode.Event) bool { if (!wantsKind(self.kinds, event)) return true; noteDelivery(self.progress, event.time_us); return self.inner.onEvent(event); } };}
/// live-half adapter: same kinds backstop over the tail (the server also/// filters; upstream wantsLive). counts consecutive transports failures/// with no events between toward the failover threshold.fn FilteringHandler(comptime H: type) type { return struct { inner: *H, kinds: []const livedecode.Kind, time_floor_us: i64 = 0, progress: *Progress, /// 0 = never rotate (single-host / rotation disabled) rotate_threshold: u8 = 0, consecutive_errors: u8 = 0, rotate_pending: bool = false,
pub fn onEvent(self: *@This(), event: livedecode.Event) bool { self.consecutive_errors = 0; if (self.time_floor_us > 0 and event.time_us < self.time_floor_us) return true; if (!wantsKind(self.kinds, event)) return true; noteDelivery(self.progress, event.time_us); return self.inner.onEvent(event); }
pub fn onError(self: *@This(), err: anyerror) bool { if (self.rotate_threshold > 0) { self.consecutive_errors += 1; if (self.consecutive_errors >= self.rotate_threshold) { self.rotate_pending = true; log.warn("live tail: {d} consecutive failures without progress (last: {s}); requesting failover", .{ self.consecutive_errors, @errorName(err) }); return false; } } if (comptime @hasDecl(H, "onError")) return self.inner.onError(err); log.warn("live tail: {s} (reconnecting)", .{@errorName(err)}); return true; }
pub fn onInfo(self: *@This(), info: livedecode.Info) void { if (comptime @hasDecl(H, "onInfo")) return self.inner.onInfo(info); log.info("live stream info: {s}: {s}", .{ info.name, info.message }); }
pub fn onConnect(self: *@This(), host: []const u8) void { if (comptime @hasDecl(H, "onConnect")) return self.inner.onConnect(host); } };}
fn noteDelivery(progress: *Progress, time_us: i64) void { if (time_us > progress.last_time_us) progress.last_time_us = time_us; progress.delivered_since_rotate = true;}
fn wantsKind(kinds: []const livedecode.Kind, event: livedecode.Event) bool { if (kinds.len == 0) return true; const kind: livedecode.Kind = event.payload; for (kinds) |k| if (k == kind) return true; return false;}
const DictSource = struct { io: Io, host: []const u8,
fn fetch(ctx: ?*anyopaque, allocator: Allocator) ?[]u8 { const self: *DictSource = @ptrCast(@alignCast(ctx.?)); return fetchZstdDict(self.io, allocator, self.host); }};
/// fetch the server's current live-tail compression dictionary. nil on any/// failure: compression is an optimization, so a failed fetch degrades to/// an uncompressed tail (logged) rather than failing the stream.fn fetchZstdDict(io: Io, allocator: Allocator, host: []const u8) ?[]u8 { var transport = zat.HttpTransport.init(io, allocator); defer transport.deinit(); var url_buf: [512]u8 = undefined; const url = std.fmt.bufPrint(&url_buf, "{s}/xrpc/network.bsky.jetstream.getZstdDictionary", .{host}) catch return null; var response = transport.fetch(.{ .url = url, .max_response_size = 16 << 20 }) catch |err| { log.warn("getZstdDictionary failed; live tail will be uncompressed: {s}", .{@errorName(err)}); return null; }; defer response.deinit(allocator); if (response.status != .ok) { log.warn("getZstdDictionary returned {d}; live tail will be uncompressed", .{@backingInt(response.status)}); return null; } return allocator.dupe(u8, response.body) catch null;}
// === tests ===
const testing = std.testing;
test "cutover and resume stay monotonic non-decreasing" { // the max() that prevents a re-learned sealed tip BELOW the cursor from // regressing the dedup floor (upstream's §14 anti-regression comment) try testing.expectEqual(@as(u64, 100), @max(@as(u64, 80), @as(u64, 100))); // tip behind cursor try testing.expectEqual(@as(u64, 120), @max(@as(u64, 120), @as(u64, 100))); // tip ahead}
test "kinds matcher: empty admits all, otherwise exact" { const ev_commit = livedecode.Event{ .seq = 1, .time_us = 0, .did = "did:plc:a", .payload = .{ .commit = .{ .operation = .delete, .collection = "c", .rkey = "k", .rev = "r", } } }; const ev_sync = livedecode.Event{ .seq = 2, .time_us = 0, .did = "did:plc:a", .payload = .{ .sync = .{ .did = "did:plc:a", .rev = "r", } } }; try testing.expect(wantsKind(&.{}, ev_commit)); try testing.expect(wantsKind(&.{.commit}, ev_commit)); try testing.expect(!wantsKind(&.{.commit}, ev_sync)); try testing.expect(wantsKind(&.{ .commit, .sync }, ev_sync));}
test "nextHostIndex mirrors the v1 client: failback to primary after progress" { // a connection that delivered events, then dropped: primary first try testing.expectEqual(@as(usize, 0), nextHostIndex(2, true, 3)); // failed attempts keep rotating, wrapping try testing.expectEqual(@as(usize, 1), nextHostIndex(0, false, 3)); try testing.expectEqual(@as(usize, 2), nextHostIndex(1, false, 3)); try testing.expectEqual(@as(usize, 0), nextHostIndex(2, false, 3));}
test "hosts: first is primary, order preserved" { const opts = Options{ .hosts = &.{ "https://a.example", "https://b.example", "https://c.example" } }; try testing.expectEqualStrings("https://a.example", opts.hosts[0]); try testing.expectEqual(@as(usize, 3), opts.hosts.len);}
const CountingHandler = struct { events: usize = 0, last_seq: u64 = 0,
pub fn onEvent(self: *@This(), event: livedecode.Event) bool { self.events += 1; self.last_seq = event.seq; return true; }};
fn commitEvent(seq: u64, time_us: i64) livedecode.Event { return .{ .seq = seq, .time_us = time_us, .did = "did:plc:a", .payload = .{ .commit = .{ .operation = .create, .collection = "c", .rkey = "k", .rev = "r", } } };}
test "archive delivery tracks progress and live replay applies its time floor" { var inner = CountingHandler{}; var progress = Progress{};
var sweep = SweepHandler(CountingHandler){ .inner = &inner, .kinds = &.{}, .progress = &progress, }; try testing.expect(sweep.onRow(commitEvent(2, 1500))); try testing.expectEqual(@as(usize, 1), inner.events); try testing.expectEqual(@as(i64, 1500), progress.last_time_us); try testing.expect(progress.delivered_since_rotate);
var live = FilteringHandler(CountingHandler){ .inner = &inner, .kinds = &.{}, .time_floor_us = 1000, .progress = &progress, }; try testing.expect(live.onEvent(commitEvent(3, 900))); try testing.expectEqual(@as(usize, 1), inner.events); try testing.expect(live.onEvent(commitEvent(4, 2000))); try testing.expectEqual(@as(usize, 2), inner.events); try testing.expectEqual(@as(i64, 2000), progress.last_time_us);}
test "failover threshold: consecutive errors rotate, any event resets" { var inner = CountingHandler{}; var progress = Progress{}; var live = FilteringHandler(CountingHandler){ .inner = &inner, .kinds = &.{}, .progress = &progress, .rotate_threshold = 3, };
try testing.expect(live.onError(error.ReadFailed)); try testing.expect(live.onError(error.ReadFailed)); // an event between failures resets the count: the host is limping, // not dead try testing.expect(live.onEvent(commitEvent(1, 100))); try testing.expect(live.onError(error.ReadFailed)); try testing.expect(live.onError(error.ReadFailed)); try testing.expect(!live.rotate_pending); // third consecutive: rotation requested, run() stops try testing.expect(!live.onError(error.ReadFailed)); try testing.expect(live.rotate_pending);}
test "failover threshold disabled: errors fall through to default handling" { var inner = CountingHandler{}; var progress = Progress{}; var live = FilteringHandler(CountingHandler){ .inner = &inner, .kinds = &.{}, .progress = &progress, .rotate_threshold = 0, }; // never rotates, never stops on the default path for (0..20) |_| try testing.expect(live.onError(error.ReadFailed)); try testing.expect(!live.rotate_pending);}