//! 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 one pub 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); }