jetstream client
atproto jetstream client
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499//! subscribeEvents frame decoding — the live half of the v2 wire.//!//! ports upstream livedecode.go (bluesky-social/jetstream, pin 289b032)://! one xrpc.v1.json envelope per text frame, discriminated by $type//! ("message" | "error"); message payloads are the lexicon union//! network.bsky.jetstream.subscribeEvents#{commit,identity,account,sync,info}.//!//! decode policy, matching upstream exactly://! - #info advisories surface as .info (the session loop logs them; no//! seq, no cursor advance)//! - an unknown envelope or payload $type is a NEWER server's frame kind://! skip for forward compatibility//! - a MISSING $type (envelope or payload) is malformed, not future://! error, so a wrong endpoint or protocol revision cannot look healthy//! while delivering nothing//! - lexicon-required fields are enforced here (the JSON layer cannot)://! seq >= 1, commit did/rev/collection/rkey, and payload-presence DIDs//! on identity/account/sync//! - untrusted diagnostic strings (error codes, #info names/messages)//! are bounded before they can reach logs//!//! records match upstream exactly: Commit.record_cbor is the canonical//! DAG-CBOR (live: canonicalized via zat.cbor.fromJson + encode, mirroring//! upstream decodeLiveRecord; archive: the stored payload verbatim), and//! Commit.record is the generic form derived from those canonical bytes.
const std = @import("std");const zat = @import("zat");const cbor = zat.cbor;const multibase = zat.multibase;const json = std.json;
const Allocator = std.mem.Allocator;
pub const nsid = "network.bsky.jetstream.subscribeEvents";
/// bounds on untrusted server-supplied diagnostic strings before they enter/// error values and logs (upstream maxLiveDiag{Name,Message}Bytes)pub const max_diag_name_bytes = 128;pub const max_diag_message_bytes = 1024;
pub const Kind = enum { commit, identity, account, sync };pub const Operation = enum { create, update, delete };
pub const Commit = struct { operation: Operation, collection: []const u8, rkey: []const u8, rev: []const u8, cid: ?[]const u8 = null, /// parsed record for create/update; null on delete. derived from /// record_cbor (upstream decodeLiveRecord: the generic form comes /// from the canonical bytes, not the wire JSON) record: ?std.json.Value = null, /// the record's canonical DAG-CBOR encoding (upstream Event.RecordCBOR). /// live: canonicalized from the wire JSON (zat.cbor.fromJson + encode); /// archive: the stored payload verbatim. empty on deletes. lifetime: /// the decode arena (live) or the block buffer (archive). record_cbor: []const u8 = "",};
pub const Identity = struct { did: []const u8, handle: ?[]const u8 = null, seq: ?i64 = null, time: ?[]const u8 = null,};
pub const Account = struct { did: []const u8, active: bool, status: ?[]const u8 = null, seq: ?i64 = null, time: ?[]const u8 = null,};
pub const Sync = struct { did: []const u8, rev: []const u8, seq: ?i64 = null, time: ?[]const u8 = null,};
/// one decoded live event. slices reference the frame buffer and the/// decode arena; they are valid only during the handler callback.pub const Event = struct { seq: u64, time_us: i64, did: []const u8, payload: union(Kind) { commit: Commit, identity: Identity, account: Account, sync: Sync, },};
/// an #info advisory (e.g. OutdatedCursor on a clamped timestamp resume)pub const Info = struct { name: []const u8, message: []const u8,};
/// a terminal xrpc.v1.json error frame; the server closes right after/// sending onepub const StreamError = struct { code: []const u8, message: []const u8,};
pub const Decoded = union(enum) { event: Event, info: Info, stream_error: StreamError, /// valid frame from a newer protocol revision; advance without emitting skip,};
pub const DecodeError = error{ MalformedFrame, /// a create/update record that cannot canonicalize to DAG-CBOR MalformedRecord, MissingEnvelopeType, MissingPayload, MissingPayloadType, MissingRequiredField, InvalidSeq, InvalidTime, UnknownOperation, MissingRecord, OutOfMemory,};
fn bound(s: []const u8, limit: usize) []const u8 { if (s.len <= limit) return s; var cut = limit; // rune-aligned truncation, like upstream boundLiveString while (cut > 0 and (s[cut] & 0xC0) == 0x80) cut -= 1; return s[0..cut];}
/// decode one text frame. `arena` owns the parsed JSON structure; the/// caller keeps it alive for the duration of the handler callback.pub fn decodeFrame(arena: Allocator, data: []const u8) DecodeError!Decoded { const parsed = std.json.parseFromSliceLeaky(std.json.Value, arena, data, .{}) catch |err| switch (err) { error.OutOfMemory => return error.OutOfMemory, else => return error.MalformedFrame, }; if (parsed != .object) return error.MalformedFrame; const env = parsed.object;
const env_type = getString(env, "$type") orelse { // no $type at all is not a newer revision — it is a malformed frame // (a v1 /subscribe server, or not a subscribeEvents endpoint at all) return error.MissingEnvelopeType; }; if (std.mem.eql(u8, env_type, "error")) { const code = getString(env, "error") orelse return error.MalformedFrame; return .{ .stream_error = .{ .code = bound(code, max_diag_name_bytes), .message = bound(getString(env, "message") orelse "", max_diag_message_bytes), } }; } if (!std.mem.eql(u8, env_type, "message")) return .skip;
const payload_value = env.get("payload") orelse return error.MissingPayload; if (payload_value != .object) return error.MissingPayload; const payload = payload_value.object; const payload_type = getString(payload, "$type") orelse return error.MissingPayloadType;
if (std.mem.eql(u8, payload_type, nsid ++ "#info")) { return .{ .info = .{ .name = bound(getString(payload, "name") orelse "", max_diag_name_bytes), .message = bound(getString(payload, "message") orelse "", max_diag_message_bytes), } }; }
const kind: Kind = if (std.mem.eql(u8, payload_type, nsid ++ "#commit")) .commit else if (std.mem.eql(u8, payload_type, nsid ++ "#identity")) .identity else if (std.mem.eql(u8, payload_type, nsid ++ "#account")) .account else if (std.mem.eql(u8, payload_type, nsid ++ "#sync")) .sync else return .skip; // a newer server's message kind
// envelope fields shared by every message kind: 1-based seq (0 = the // required field was absent; accepting it would hand the dedup an event // it silently swallows) and the canonical datetime const seq_raw = getInt(payload, "seq") orelse return error.InvalidSeq; if (seq_raw <= 0) return error.InvalidSeq; const seq: u64 = @intCast(seq_raw); const time_str = getString(payload, "time") orelse return error.InvalidTime; const time_us = (zat.Datetime.parse(time_str) orelse return error.InvalidTime).micros; const did = getString(payload, "did") orelse return error.MissingRequiredField; if (did.len == 0) return error.MissingRequiredField;
switch (kind) { .commit => { const collection = getString(payload, "collection") orelse return error.MissingRequiredField; const rkey = getString(payload, "rkey") orelse return error.MissingRequiredField; const rev = getString(payload, "rev") orelse return error.MissingRequiredField; if (collection.len == 0 or rkey.len == 0 or rev.len == 0) return error.MissingRequiredField; const op_str = getString(payload, "operation") orelse return error.MissingRequiredField; const operation = std.meta.stringToEnum(Operation, op_str) orelse return error.UnknownOperation; var record: ?std.json.Value = null; var record_cbor: []const u8 = ""; switch (operation) { .create, .update => { const wire = payload.get("record") orelse return error.MissingRecord; if (wire != .object) return error.MissingRecord; // upstream decodeLiveRecord: canonicalize the atproto // JSON to DAG-CBOR, then derive the generic record from // those canonical bytes (deterministic encoding lets // callers verify the declared CID without a duplicate // byte string on the wire) const value = cbor.fromJson(arena, wire) catch |err| switch (err) { error.OutOfMemory => return error.OutOfMemory, else => return error.MalformedRecord, }; record_cbor = cbor.encodeAlloc(arena, value) catch |err| switch (err) { error.OutOfMemory => return error.OutOfMemory, error.WriteFailed => return error.OutOfMemory, }; const decoded_back = cbor.decode(arena, record_cbor) catch return error.MalformedRecord; record = try cborToJson(arena, decoded_back.value); }, .delete => {}, // no record payload on deletes } return .{ .event = .{ .seq = seq, .time_us = time_us, .did = did, .payload = .{ .commit = .{ .operation = operation, .collection = collection, .rkey = rkey, .rev = rev, .cid = getString(payload, "cid"), .record = record, .record_cbor = record_cbor, } } } }; }, .identity => { const inner_value = payload.get("identity") orelse return error.MissingRequiredField; if (inner_value != .object) return error.MissingRequiredField; const inner = inner_value.object; // payload-presence check (did is set by every producer); scalar // required fields are indistinguishable from their zero values const inner_did = getString(inner, "did") orelse return error.MissingRequiredField; if (inner_did.len == 0) return error.MissingRequiredField; return .{ .event = .{ .seq = seq, .time_us = time_us, .did = did, .payload = .{ .identity = .{ .did = inner_did, .handle = getString(inner, "handle"), .seq = getInt(inner, "seq"), .time = getString(inner, "time"), } } } }; }, .account => { const inner_value = payload.get("account") orelse return error.MissingRequiredField; if (inner_value != .object) return error.MissingRequiredField; const inner = inner_value.object; const inner_did = getString(inner, "did") orelse return error.MissingRequiredField; if (inner_did.len == 0) return error.MissingRequiredField; const active_value = inner.get("active") orelse return error.MissingRequiredField; if (active_value != .bool) return error.MissingRequiredField; return .{ .event = .{ .seq = seq, .time_us = time_us, .did = did, .payload = .{ .account = .{ .did = inner_did, .active = active_value.bool, .status = getString(inner, "status"), .seq = getInt(inner, "seq"), .time = getString(inner, "time"), } } } }; }, .sync => { const inner_value = payload.get("sync") orelse return error.MissingRequiredField; if (inner_value != .object) return error.MissingRequiredField; const inner = inner_value.object; // archived #sync payloads from an async resync legitimately // carry empty time/seq; did is the only reliable presence marker const inner_did = getString(inner, "did") orelse return error.MissingRequiredField; if (inner_did.len == 0) return error.MissingRequiredField; return .{ .event = .{ .seq = seq, .time_us = time_us, .did = did, .payload = .{ .sync = .{ .did = inner_did, .rev = getString(inner, "rev") orelse "", .seq = getInt(inner, "seq"), .time = getString(inner, "time"), } } } }; }, }}
fn getString(obj: std.json.ObjectMap, key: []const u8) ?[]const u8 { const v = obj.get(key) orelse return null; return if (v == .string) v.string else null;}
fn getInt(obj: std.json.ObjectMap, key: []const u8) ?i64 { const v = obj.get(key) orelse return null; return if (v == .integer) v.integer else null;}
// === DAG-CBOR record → the json.Value shape the live client delivers ===
/// bytes and links use the atproto lex-JSON conventions the live jetstream/// emits: {"$bytes": base64-no-pad} and {"$link": cid-base32-string}.pub fn cborToJson(arena: Allocator, value: cbor.Value) Allocator.Error!json.Value { return switch (value) { .null => .null, .boolean => |b| .{ .bool = b }, .unsigned => |u| if (std.math.cast(i64, u)) |i| .{ .integer = i } else .{ .float = @floatFromInt(u) }, .negative => |i| .{ .integer = i }, .float => |f| .{ .float = f }, .text => |s| .{ .string = s }, .bytes => |b| blk: { const b64 = std.base64.standard_no_pad.Encoder; const encoded = try arena.alloc(u8, b64.calcSize(b.len)); _ = b64.encode(encoded, b); var obj: json.ObjectMap = .empty; try obj.put(arena, "$bytes", .{ .string = encoded }); break :blk .{ .object = obj }; }, .cid => |c| blk: { const encoded = multibase.base32lower.encode(arena, c.raw) catch return error.OutOfMemory; var obj: json.ObjectMap = .empty; try obj.put(arena, "$link", .{ .string = encoded }); break :blk .{ .object = obj }; }, .array => |items| blk: { var arr = json.Array.init(arena); try arr.ensureTotalCapacityPrecise(items.len); for (items) |item| arr.appendAssumeCapacity(try cborToJson(arena, item)); break :blk .{ .array = arr }; }, .map => |entries| blk: { var obj: json.ObjectMap = .empty; try obj.ensureTotalCapacity(arena, entries.len); for (entries) |entry| obj.putAssumeCapacity(entry.key, try cborToJson(arena, entry.value)); break :blk .{ .object = obj }; }, };}
// === tests ===// fixtures mirror stream's wire.zig v2 encoder output byte-shape and// upstream livedecode_test.go's malformed-frame matrix
const testing = std.testing;
fn decodeTest(arena: Allocator, frame: []const u8) DecodeError!Decoded { return decodeFrame(arena, frame);}
test "commit frame decodes with envelope fields and flat commit payload" { var arena = std.heap.ArenaAllocator.init(testing.allocator); defer arena.deinit(); const frame = \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#commit","cid":"bafy123","collection":"app.bsky.feed.post","did":"did:plc:abc","operation":"create","record":{"$type":"app.bsky.feed.post","text":"hi"},"rev":"3k2a","rkey":"3k2abcdefghij","seq":42,"time":"2026-07-13T00:00:01.500000Z"}} ; const decoded = try decodeTest(arena.allocator(), frame); const ev = decoded.event; try testing.expectEqual(@as(u64, 42), ev.seq); try testing.expectEqual(@as(i64, (1783900800 + 1) * std.time.us_per_s + 500_000), ev.time_us); try testing.expectEqualStrings("did:plc:abc", ev.did); const commit = ev.payload.commit; try testing.expectEqual(Operation.create, commit.operation); try testing.expectEqualStrings("app.bsky.feed.post", commit.collection); try testing.expectEqualStrings("3k2abcdefghij", commit.rkey); try testing.expectEqualStrings("bafy123", commit.cid.?); try testing.expectEqualStrings("hi", commit.record.?.object.get("text").?.string);
// upstream Event.RecordCBOR parity: the canonical DAG-CBOR bytes of // the wire record, byte-equal to encoding the hand-built value var check_arena = std.heap.ArenaAllocator.init(testing.allocator); defer check_arena.deinit(); const ca = check_arena.allocator(); const expected_cbor = try cbor.encodeAlloc(ca, .{ .map = &.{ .{ .key = "$type", .value = .{ .text = "app.bsky.feed.post" } }, .{ .key = "text", .value = .{ .text = "hi" } }, } }); try testing.expectEqualSlices(u8, expected_cbor, commit.record_cbor);}
test "delete commit carries no record; create without one errors" { var arena = std.heap.ArenaAllocator.init(testing.allocator); defer arena.deinit(); const del = \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#commit","collection":"c","did":"did:plc:a","operation":"delete","rev":"r","rkey":"k","seq":1,"time":"2026-01-01T00:00:00.000000Z"}} ; const decoded = try decodeTest(arena.allocator(), del); try testing.expect(decoded.event.payload.commit.record == null); try testing.expectEqual(@as(usize, 0), decoded.event.payload.commit.record_cbor.len);
const create_missing = \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#commit","collection":"c","did":"did:plc:a","operation":"create","rev":"r","rkey":"k","seq":1,"time":"2026-01-01T00:00:00.000000Z"}} ; try testing.expectError(error.MissingRecord, decodeTest(arena.allocator(), create_missing));}
test "identity, account, and sync enforce payload-presence DIDs" { var arena = std.heap.ArenaAllocator.init(testing.allocator); defer arena.deinit(); const identity = \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#identity","did":"did:plc:a","identity":{"did":"did:plc:a","handle":"alice.test","seq":7,"time":"t"},"seq":9,"time":"2026-01-01T00:00:00.000000Z"}} ; const id_ev = (try decodeTest(arena.allocator(), identity)).event; try testing.expectEqualStrings("alice.test", id_ev.payload.identity.handle.?);
const identity_hollow = \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#identity","did":"did:plc:a","identity":{},"seq":9,"time":"2026-01-01T00:00:00.000000Z"}} ; try testing.expectError(error.MissingRequiredField, decodeTest(arena.allocator(), identity_hollow));
const account = \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#account","account":{"active":false,"did":"did:plc:a","status":"takendown"},"did":"did:plc:a","seq":10,"time":"2026-01-01T00:00:00.000000Z"}} ; const acct = (try decodeTest(arena.allocator(), account)).event.payload.account; try testing.expect(!acct.active); try testing.expectEqualStrings("takendown", acct.status.?);
const sync_frame = \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#sync","did":"did:plc:a","sync":{"did":"did:plc:a","rev":"3k2a"},"seq":11,"time":"2026-01-01T00:00:00.000000Z"}} ; const sync_ev = (try decodeTest(arena.allocator(), sync_frame)).event.payload.sync; try testing.expectEqualStrings("3k2a", sync_ev.rev);}
test "error and info frames surface typed and bounded" { var arena = std.heap.ArenaAllocator.init(testing.allocator); defer arena.deinit(); const err_frame = \\{"$type":"error","error":"FutureCursor","message":"cursor is ahead of the stream"} ; const stream_err = (try decodeTest(arena.allocator(), err_frame)).stream_error; try testing.expectEqualStrings("FutureCursor", stream_err.code);
const info_frame = \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#info","name":"OutdatedCursor","message":"resumed from seq 5"}} ; const info = (try decodeTest(arena.allocator(), info_frame)).info; try testing.expectEqualStrings("OutdatedCursor", info.name);}
test "malformed vs future frames: missing $type errors, unknown $type skips" { var arena = std.heap.ArenaAllocator.init(testing.allocator); defer arena.deinit(); // v1 /subscribe JSON has no envelope $type: malformed, never skipped try testing.expectError(error.MissingEnvelopeType, decodeTest(arena.allocator(), \\{"did":"did:plc:a","time_us":1,"kind":"commit"} )); // a newer revision's envelope kind skips try testing.expectEqual(Decoded.skip, try decodeTest(arena.allocator(), \\{"$type":"snapshot","payload":{}} )); // a newer server's payload kind skips try testing.expectEqual(Decoded.skip, try decodeTest(arena.allocator(), \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#hologram","seq":1}} )); // a payload with NO $type is malformed, not future try testing.expectError(error.MissingPayloadType, decodeTest(arena.allocator(), \\{"$type":"message","payload":{"seq":1}} )); // seq 0 means the required field was absent try testing.expectError(error.InvalidSeq, decodeTest(arena.allocator(), \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#commit","collection":"c","did":"did:plc:a","operation":"delete","rev":"r","rkey":"k","seq":0,"time":"2026-01-01T00:00:00.000000Z"}} )); try testing.expectError(error.MalformedFrame, decodeTest(arena.allocator(), "not json"));}
test "live record canonicalization: floats, $link, $bytes survive to canonical CBOR" { var arena = std.heap.ArenaAllocator.init(testing.allocator); defer arena.deinit(); const frame = \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#commit","collection":"c","did":"did:plc:abc","operation":"create","record":{"ratio":1.5,"ref":{"$link":"bafyreidfayvfuwqa7qlnopdjiqrxzs6blmoeu4rujcjtnci5beludirz2a"},"blob":{"$bytes":"3q2+7w"},"whole":42.0},"rev":"r","rkey":"k","seq":7,"time":"2026-07-13T00:00:01.000000Z"}} ; const decoded = try decodeTest(arena.allocator(), frame); const commit = decoded.event.payload.commit;
// the canonical bytes round-trip through the strict decoder const back = try cbor.decode(arena.allocator(), commit.record_cbor); var ratio: ?f64 = null; var whole: ?u64 = null; for (back.value.map) |entry| { if (std.mem.eql(u8, entry.key, "ratio")) ratio = entry.value.float; if (std.mem.eql(u8, entry.key, "whole")) whole = entry.value.unsigned; if (std.mem.eql(u8, entry.key, "ref")) try testing.expect(entry.value == .cid); if (std.mem.eql(u8, entry.key, "blob")) try testing.expectEqualSlices(u8, &.{ 0xde, 0xad, 0xbe, 0xef }, entry.value.bytes); } try testing.expectEqual(@as(f64, 1.5), ratio.?); // whole floats canonicalize to integers (atmos FromJSON) try testing.expectEqual(@as(u64, 42), whole.?);
// the generic record is DERIVED from the canonical bytes: the $link // survives as a sentinel object and the float as a JSON number try testing.expectEqual(@as(f64, 1.5), commit.record.?.object.get("ratio").?.float); try testing.expect(commit.record.?.object.get("ref").?.object.get("$link") != null);}