atproto utils for zig zat.dev
atproto sdk zig
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638//! firehose codec - com.atproto.sync.subscribeRepos//!//! encode and decode AT Protocol firehose events over WebSocket. messages are//! DAG-CBOR encoded (unlike jetstream, which is JSON). includes frame encoding///! decoding, CAR block packing, and CID creation for records.//!//! wire format per frame://! [DAG-CBOR header: {op, t}] [DAG-CBOR payload: {seq, repo, ops, blocks, ...}]//!//! see: https://atproto.com/specs/event-stream
const std = @import("std");const cbor = @import("../repo/cbor.zig");const car = @import("../repo/car.zig");const mst = @import("../repo/mst.zig");const sync = @import("sync.zig");const Did = @import("../syntax/did.zig").Did;const Nsid = @import("../syntax/nsid.zig").Nsid;const Rkey = @import("../syntax/rkey.zig").Rkey;const Tid = @import("../syntax/tid.zig").Tid;
const mem = std.mem;const Allocator = mem.Allocator;const posix = std.posix;const Io = std.Io;const log = std.log.scoped(.zat);
pub const CommitAction = sync.CommitAction;pub const AccountStatus = sync.AccountStatus;
pub const default_hosts = [_][]const u8{ "bsky.network", "northamerica.firehose.network", "europe.firehose.network", "asia.firehose.network",};
pub const Options = struct { /// relay endpoints, as full URLs following ecosystem norms (atmos /// streaming.Options.URL, jetstream --relay-url): the scheme carries /// tls and the port rides in the URL, e.g. "wss://bsky.network" or /// "ws://localhost:7777". any path is ignored — the subscribeRepos /// path is fixed by the protocol. bare hostnames are implicit wss. hosts: []const []const u8 = &default_hosts, cursor: ?i64 = null, /// Fail over after this many milliseconds with no WebSocket frames. /// Catches a host that holds the socket open but goes silent, which TCP /// keepalive never fires on (the socket is healthy; the stream behind it /// stalled). Mirrors JetstreamClient's option of the same name. A live /// relay emits continuously, so even a few seconds of total silence is /// anomalous — but the default stays off: reconnect policy belongs to /// the consumer. idle_timeout_ms: ?u32 = null, max_message_size: usize = 5 * 1024 * 1024, // 5MB — firehose frames can be large};
/// a parsed relay endpoint (see Options.hosts for the accepted forms)pub const Endpoint = struct { host: []const u8, port: u16, tls: bool,
pub fn parse(url: []const u8) error{InvalidEndpoint}!Endpoint { var rest = url; var tls = true; if (mem.startsWith(u8, rest, "wss://")) { rest = rest["wss://".len..]; } else if (mem.startsWith(u8, rest, "ws://")) { tls = false; rest = rest["ws://".len..]; } else if (mem.indexOf(u8, rest, "://") != null) { return error.InvalidEndpoint; } if (mem.indexOfScalar(u8, rest, '/')) |i| rest = rest[0..i]; if (mem.indexOfScalar(u8, rest, '?')) |i| rest = rest[0..i]; var port: u16 = if (tls) 443 else 80; var host = rest; if (mem.indexOfScalar(u8, rest, ':')) |i| { host = rest[0..i]; port = std.fmt.parseInt(u16, rest[i + 1 ..], 10) catch return error.InvalidEndpoint; } if (host.len == 0) return error.InvalidEndpoint; return .{ .host = host, .port = port, .tls = tls }; }
fn isDefaultPort(self: Endpoint) bool { return self.port == if (self.tls) @as(u16, 443) else 80; }};
/// decoded firehose eventpub const Event = union(enum) { commit: CommitEvent, sync: SyncEvent, identity: IdentityEvent, account: AccountEvent, info: InfoEvent,
pub fn seq(self: Event) ?i64 { return switch (self) { .commit => |c| c.seq, .sync => |s| s.seq, .identity => |i| i.seq, .account => |a| a.seq, .info => null, }; }};
pub const CommitEvent = struct { seq: i64, repo: []const u8, // DID rev: []const u8, // TID — revision of the commit time: []const u8, // datetime — when event was received since: ?[]const u8 = null, // TID — rev of preceding commit (null = full repo export) commit: ?cbor.Cid = null, // CID of the commit object blocks: []const u8 = &.{}, // raw CAR diff bytes ops: []const RepoOp, prev_data: ?cbor.Cid = null, // MST root CID of the previous revision blobs: []const cbor.Cid = &.{}, // new blobs referenced by records in this commit rebase: bool = false, too_big: bool = false,
pub fn toMstOperations(self: CommitEvent, allocator: Allocator) Allocator.Error![]mst.Operation { var ops: std.ArrayList(mst.Operation) = .empty; errdefer ops.deinit(allocator); for (self.ops) |op| { const path = if (op.path.len != 0) try allocator.dupe(u8, op.path) else try std.fmt.allocPrint(allocator, "{s}/{s}", .{ op.collection, op.rkey }); try ops.append(allocator, .{ .path = path, .value = if (op.cid) |cid| cid.raw else null, .prev = if (op.prev) |prev| prev.raw else null, }); } return ops.toOwnedSlice(allocator); }};
pub const RepoOp = struct { action: CommitAction, path: []const u8 = "", collection: []const u8, rkey: []const u8, cid: ?cbor.Cid = null, // CID of the record (null for deletes) prev: ?cbor.Cid = null, // CID of the previous record for updates/deletes record: ?cbor.Value = null, // decoded DAG-CBOR record from CAR block};
pub const CommitEventOp = struct { action: CommitAction, collection: []const u8, rkey: []const u8, cid: ?cbor.Cid = null, prev: ?cbor.Cid = null,};
pub const CommitEventParams = struct { seq: i64, repo_did: []const u8, commit_cid: cbor.Cid, rev: []const u8, since_rev: ?[]const u8 = null, prev_data: ?cbor.Cid = null, blocks: []const u8, ops: []const CommitEventOp, blobs: []const cbor.Cid = &.{}, time: []const u8, rebase: bool = false, too_big: bool = false,};
pub const SyncEvent = struct { seq: i64, did: []const u8, rev: []const u8, time: []const u8, blocks: []const u8, // raw CAR bytes containing the current commit object};
pub const IdentityEvent = struct { seq: i64, did: []const u8, time: []const u8, // datetime — when event was received handle: ?[]const u8 = null,};
pub const AccountEvent = struct { seq: i64, did: []const u8, time: []const u8, // datetime — when event was received active: bool = true, status: ?AccountStatus = null,};
pub const InfoEvent = struct { name: ?[]const u8 = null, message: ?[]const u8 = null,};
/// frame header from the wireconst FrameHeader = struct { op: i64, t: ?[]const u8 = null,};
const FrameOp = enum(i64) { message = 1, err = -1,};
pub const DecodeError = error{ InvalidFrame, InvalidHeader, UnexpectedEof, MissingField, InvalidRepoPath, UnknownOp, UnknownEventType,} || cbor.DecodeError || car.CarError;
pub const ValidateCommitError = error{ MissingCommitCid, InvalidRepoDid, InvalidRev, InvalidSince, InvalidRepoPath, MissingRecordCid, UnexpectedRecordCid, UnexpectedPrevRecordCid,};
/// cheaply extract the seq from a raw frame without CAR hydration./// returns null for frames without a seq (#info, errors, malformed).pub fn peekSeq(allocator: Allocator, data: []const u8) ?i64 { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const header_result = cbor.decode(arena.allocator(), data) catch return null; if ((header_result.value.getInt("op") orelse return null) != 1) return null; const payload = cbor.decodeAll(arena.allocator(), data[header_result.consumed..]) catch return null; return payload.getInt("seq");}
/// decode a raw WebSocket binary frame into a firehose Eventpub fn decodeFrame(allocator: Allocator, data: []const u8) DecodeError!Event { // frame = [CBOR header] [CBOR payload] concatenated const header_result = try cbor.decode(allocator, data); const header_val = header_result.value; const payload_data = data[header_result.consumed..];
// parse header const op = header_val.getInt("op") orelse return error.InvalidHeader; if (op == -1) return error.UnknownOp; // error frame
const t = header_val.getString("t") orelse return error.InvalidHeader;
// decode payload const payload = try cbor.decodeAll(allocator, payload_data);
if (mem.eql(u8, t, "#commit")) { return try decodeCommit(allocator, payload); } else if (mem.eql(u8, t, "#sync")) { return decodeSync(payload); } else if (mem.eql(u8, t, "#identity")) { return decodeIdentity(payload); } else if (mem.eql(u8, t, "#account")) { return decodeAccount(payload); } else if (mem.eql(u8, t, "#info")) { return .{ .info = .{ .name = payload.getString("name"), .message = payload.getString("message"), } }; }
return error.UnknownEventType;}
fn decodeCommit(allocator: Allocator, payload: cbor.Value) DecodeError!Event { const seq_val = payload.getInt("seq") orelse return error.MissingField; const repo = payload.getString("repo") orelse return error.MissingField; const rev = payload.getString("rev") orelse return error.MissingField; const time = payload.getString("time") orelse return error.MissingField;
const commit_cid = payload.getCid("commit") orelse return error.MissingField; const blocks_bytes = payload.getBytes("blocks") orelse return error.MissingField; const prev_data = try optionalCid(payload, "prevData");
// parse blobs array (array of CID links) var blobs: std.ArrayList(cbor.Cid) = .empty; if (payload.getArray("blobs")) |blob_values| { for (blob_values) |blob_val| { switch (blob_val) { .cid => |c| try blobs.append(allocator, c), else => {}, } } }
// parse CAR blocks for record hydration. CAR verification failure should not // hide the wire event; consumers can reject by verifying `blocks` explicitly. const parsed_car: ?car.Car = car.read(allocator, blocks_bytes) catch null;
// parse ops const ops_array = payload.getArray("ops"); var ops: std.ArrayList(RepoOp) = .empty;
if (ops_array) |op_values| { for (op_values) |op_val| { const action_str = op_val.getString("action") orelse continue; const action = CommitAction.parse(action_str) orelse continue; const path = op_val.getString("path") orelse continue;
// Decoding preserves hostile/noncanonical paths so downstream // ingest policy can drop one bad op without losing its valid // siblings. Structural validation remains available through // validateCommitEvent and encoding still calls it. const repo_path = splitRepoPath(path);
// extract CID from op and look up record from CAR blocks var op_cid: ?cbor.Cid = null; var op_prev: ?cbor.Cid = null; var record: ?cbor.Value = null; if (op_val.get("cid")) |cid_val| { switch (cid_val) { .cid => |cid| { op_cid = cid; if (parsed_car) |c| { if (car.findBlock(c, cid.raw)) |block_data| { record = cbor.decodeAll(allocator, block_data) catch null; } } }, else => {}, } } if (op_val.get("prev")) |prev_val| { switch (prev_val) { .cid => |cid| op_prev = cid, else => {}, } }
try ops.append(allocator, .{ .action = action, .path = path, .collection = repo_path.collection, .rkey = repo_path.rkey, .cid = op_cid, .prev = op_prev, .record = record, }); } }
return .{ .commit = .{ .seq = seq_val, .repo = repo, .rev = rev, .time = time, .since = payload.getString("since"), .commit = commit_cid, .blocks = blocks_bytes, .ops = try ops.toOwnedSlice(allocator), .prev_data = prev_data, .blobs = try blobs.toOwnedSlice(allocator), .rebase = payload.getBool("rebase") orelse false, .too_big = payload.getBool("tooBig") orelse false, } };}
/// Validate the structural shape of a decoded commit event.////// This is intentionally lighter than repo verification: it checks event/// vocabulary, identifier syntax, repo op paths, and op CID nullability. It/// does not parse CAR blocks, verify commit signatures, or prove MST roots.pub fn validateCommitEvent(commit: CommitEvent) ValidateCommitError!void { if (Did.parse(commit.repo) == null) return error.InvalidRepoDid; if (Tid.parse(commit.rev) == null) return error.InvalidRev; if (commit.since) |since| { if (Tid.parse(since) == null) return error.InvalidSince; } if (commit.commit == null) return error.MissingCommitCid;
for (commit.ops) |op| { try validateRepoOpPath(op); switch (op.action) { .create => { if (op.cid == null) return error.MissingRecordCid; if (op.prev != null) return error.UnexpectedPrevRecordCid; }, .update => { if (op.cid == null) return error.MissingRecordCid; }, .delete => { if (op.cid != null) return error.UnexpectedRecordCid; }, } }}
fn commitEvent(allocator: Allocator, params: CommitEventParams) (Allocator.Error || ValidateCommitError)!CommitEvent { var ops: std.ArrayList(RepoOp) = .empty; errdefer ops.deinit(allocator);
for (params.ops) |op| { try ops.append(allocator, .{ .action = op.action, .path = "", .collection = op.collection, .rkey = op.rkey, .cid = op.cid, .prev = op.prev, }); }
const event: CommitEvent = .{ .seq = params.seq, .repo = params.repo_did, .rev = params.rev, .time = params.time, .since = params.since_rev, .commit = params.commit_cid, .blocks = params.blocks, .ops = try ops.toOwnedSlice(allocator), .prev_data = params.prev_data, .blobs = params.blobs, .rebase = params.rebase, .too_big = params.too_big, }; errdefer allocator.free(event.ops);
try validateCommitEvent(event); return event;}
pub fn encodeCommitEvent(allocator: Allocator, params: CommitEventParams) ![]u8 { var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit();
const event = try commitEvent(arena.allocator(), params); return encodeFrame(allocator, .{ .commit = event });}
const RepoPath = struct { collection: []const u8, rkey: []const u8,};
fn splitRepoPath(path: []const u8) RepoPath { const slash = mem.indexOfScalar(u8, path, '/') orelse return .{ .collection = path, .rkey = "", }; return .{ .collection = path[0..slash], .rkey = path[slash + 1 ..], };}
fn parseRepoPath(path: []const u8) error{InvalidRepoPath}!RepoPath { const slash = mem.indexOfScalar(u8, path, '/') orelse return error.InvalidRepoPath; if (mem.indexOfScalarPos(u8, path, slash + 1, '/') != null) return error.InvalidRepoPath;
const result = splitRepoPath(path); if (Nsid.parse(result.collection) == null) return error.InvalidRepoPath; if (Rkey.parse(result.rkey) == null) return error.InvalidRepoPath;
return result;}
fn validateRepoOpPath(op: RepoOp) error{InvalidRepoPath}!void { if (op.path.len != 0) { _ = try parseRepoPath(op.path); return; } if (Nsid.parse(op.collection) == null) return error.InvalidRepoPath; if (Rkey.parse(op.rkey) == null) return error.InvalidRepoPath;}
fn optionalCid(value: cbor.Value, key: []const u8) DecodeError!?cbor.Cid { const found = value.get(key) orelse return null; return switch (found) { .cid => |cid| cid, .null => null, else => error.MissingField, };}
test "optional CID accepts omitted and explicit null fields" { try std.testing.expect(try optionalCid(.{ .map = &.{} }, "prevData") == null); try std.testing.expect(try optionalCid(.{ .map = &.{.{ .key = "prevData", .value = .null, }} }, "prevData") == null); try std.testing.expectError(error.MissingField, optionalCid(.{ .map = &.{.{ .key = "prevData", .value = .{ .text = "not-a-cid" }, }} }, "prevData"));}
fn decodeSync(payload: cbor.Value) DecodeError!Event { return .{ .sync = .{ .seq = payload.getInt("seq") orelse return error.MissingField, .did = payload.getString("did") orelse return error.MissingField, .rev = payload.getString("rev") orelse return error.MissingField, .time = payload.getString("time") orelse return error.MissingField, .blocks = payload.getBytes("blocks") orelse return error.MissingField, } };}
fn decodeIdentity(payload: cbor.Value) DecodeError!Event { return .{ .identity = .{ .seq = payload.getInt("seq") orelse return error.MissingField, .did = payload.getString("did") orelse return error.MissingField, .time = payload.getString("time") orelse return error.MissingField, .handle = payload.getString("handle"), } };}
fn decodeAccount(payload: cbor.Value) DecodeError!Event { const status_str = payload.getString("status"); return .{ .account = .{ .seq = payload.getInt("seq") orelse return error.MissingField, .did = payload.getString("did") orelse return error.MissingField, .time = payload.getString("time") orelse return error.MissingField, .active = payload.getBool("active") orelse true, .status = if (status_str) |s| AccountStatus.parse(s) else null, } };}
// === encoder ===
/// encode a firehose Event into a wire frame: [DAG-CBOR header] [DAG-CBOR payload]fn encodeFrame(allocator: Allocator, event: Event) ![]u8 { var aw: std.Io.Writer.Allocating = .init(allocator); errdefer aw.deinit();
const tag = switch (event) { .commit => "#commit", .sync => "#sync", .identity => "#identity", .account => "#account", .info => "#info", };
// encode header: {op: 1, t: "#..."} const header: cbor.Value = .{ .map = &.{ .{ .key = "op", .value = .{ .unsigned = 1 } }, .{ .key = "t", .value = .{ .text = tag } }, } }; try cbor.encode(allocator, &aw.writer, header);
// encode payload based on event type switch (event) { .commit => |commit| try encodeCommitPayload(allocator, &aw.writer, commit), .sync => |sync_event| try encodeSyncPayload(allocator, &aw.writer, sync_event), .identity => |id| try encodeIdentityPayload(allocator, &aw.writer, id), .account => |acct| try encodeAccountPayload(allocator, &aw.writer, acct), .info => |inf| try encodeInfoPayload(allocator, &aw.writer, inf), }
return try aw.toOwnedSlice();}
fn encodeCommitPayload(allocator: Allocator, writer: anytype, commit: CommitEvent) !void { // build ops array and CAR blocks simultaneously var op_values: std.ArrayList(cbor.Value) = .empty; defer op_values.deinit(allocator); var car_blocks: std.ArrayList(car.Block) = .empty; defer car_blocks.deinit(allocator); var root_cids: std.ArrayList(cbor.Cid) = .empty; defer root_cids.deinit(allocator);
for (commit.ops) |op| { const action_str: []const u8 = @tagName(op.action); const path = if (op.path.len != 0) op.path else try std.fmt.allocPrint(allocator, "{s}/{s}", .{ op.collection, op.rkey }); _ = try parseRepoPath(path);
if (op.record) |record| { // encode record, create CID, add to CAR blocks const record_bytes = try cbor.encodeAlloc(allocator, record); const cid = try cbor.Cid.forDagCbor(allocator, record_bytes);
try car_blocks.append(allocator, .{ .cid_raw = cid.raw, .data = record_bytes, });
if (root_cids.items.len == 0) { try root_cids.append(allocator, cid); }
var op_entries: std.ArrayList(cbor.Value.MapEntry) = .empty; defer op_entries.deinit(allocator); try op_entries.append(allocator, .{ .key = "action", .value = .{ .text = action_str } }); try op_entries.append(allocator, .{ .key = "cid", .value = .{ .cid = cid } }); try op_entries.append(allocator, .{ .key = "path", .value = .{ .text = path } }); if (op.prev) |prev| { try op_entries.append(allocator, .{ .key = "prev", .value = .{ .cid = prev } }); } try op_values.append(allocator, .{ .map = try op_entries.toOwnedSlice(allocator) }); } else { var op_entries: std.ArrayList(cbor.Value.MapEntry) = .empty; defer op_entries.deinit(allocator); try op_entries.append(allocator, .{ .key = "action", .value = .{ .text = action_str } }); try op_entries.append(allocator, .{ .key = "cid", .value = if (op.cid) |cid| .{ .cid = cid } else .null }); try op_entries.append(allocator, .{ .key = "path", .value = .{ .text = path } }); if (op.prev) |prev| { try op_entries.append(allocator, .{ .key = "prev", .value = .{ .cid = prev } }); } try op_values.append(allocator, .{ .map = try op_entries.toOwnedSlice(allocator) }); } }
const blocks_bytes = if (commit.blocks.len != 0) commit.blocks else blk: { const car_data = car.Car{ .roots = root_cids.items, .blocks = car_blocks.items, }; break :blk try car.writeAlloc(allocator, car_data); };
// build blobs array var blob_values: std.ArrayList(cbor.Value) = .empty; defer blob_values.deinit(allocator); for (commit.blobs) |blob| { try blob_values.append(allocator, .{ .cid = blob }); }
// build payload entries var entries: std.ArrayList(cbor.Value.MapEntry) = .empty; defer entries.deinit(allocator);
try entries.append(allocator, .{ .key = "blocks", .value = .{ .bytes = blocks_bytes } }); if (commit.commit) |commit_cid| { try entries.append(allocator, .{ .key = "commit", .value = .{ .cid = commit_cid } }); } try entries.append(allocator, .{ .key = "blobs", .value = .{ .array = blob_values.items } }); try entries.append(allocator, .{ .key = "ops", .value = .{ .array = op_values.items } }); try entries.append(allocator, .{ .key = "prevData", .value = if (commit.prev_data) |prev_data| .{ .cid = prev_data } else .null }); if (commit.rebase) { try entries.append(allocator, .{ .key = "rebase", .value = .{ .boolean = true } }); } try entries.append(allocator, .{ .key = "repo", .value = .{ .text = commit.repo } }); try entries.append(allocator, .{ .key = "rev", .value = .{ .text = commit.rev } }); try entries.append(allocator, .{ .key = "seq", .value = .{ .unsigned = @intCast(commit.seq) } }); if (commit.since) |s| { try entries.append(allocator, .{ .key = "since", .value = .{ .text = s } }); } try entries.append(allocator, .{ .key = "time", .value = .{ .text = commit.time } }); if (commit.too_big) { try entries.append(allocator, .{ .key = "tooBig", .value = .{ .boolean = true } }); }
try cbor.encode(allocator, writer, .{ .map = entries.items });}
fn encodeSyncPayload(allocator: Allocator, writer: anytype, sync_event: SyncEvent) !void { var entries: std.ArrayList(cbor.Value.MapEntry) = .empty; defer entries.deinit(allocator);
try entries.append(allocator, .{ .key = "blocks", .value = .{ .bytes = sync_event.blocks } }); try entries.append(allocator, .{ .key = "did", .value = .{ .text = sync_event.did } }); try entries.append(allocator, .{ .key = "rev", .value = .{ .text = sync_event.rev } }); try entries.append(allocator, .{ .key = "seq", .value = .{ .unsigned = @intCast(sync_event.seq) } }); try entries.append(allocator, .{ .key = "time", .value = .{ .text = sync_event.time } });
try cbor.encode(allocator, writer, .{ .map = entries.items });}
fn encodeIdentityPayload(allocator: Allocator, writer: anytype, identity: IdentityEvent) !void { var entries: std.ArrayList(cbor.Value.MapEntry) = .empty; defer entries.deinit(allocator);
try entries.append(allocator, .{ .key = "did", .value = .{ .text = identity.did } }); if (identity.handle) |h| { try entries.append(allocator, .{ .key = "handle", .value = .{ .text = h } }); } try entries.append(allocator, .{ .key = "seq", .value = .{ .unsigned = @intCast(identity.seq) } }); try entries.append(allocator, .{ .key = "time", .value = .{ .text = identity.time } });
try cbor.encode(allocator, writer, .{ .map = entries.items });}
fn encodeAccountPayload(allocator: Allocator, writer: anytype, account: AccountEvent) !void { var entries: std.ArrayList(cbor.Value.MapEntry) = .empty; defer entries.deinit(allocator);
if (!account.active) { try entries.append(allocator, .{ .key = "active", .value = .{ .boolean = false } }); } try entries.append(allocator, .{ .key = "did", .value = .{ .text = account.did } }); try entries.append(allocator, .{ .key = "seq", .value = .{ .unsigned = @intCast(account.seq) } }); if (account.status) |s| { try entries.append(allocator, .{ .key = "status", .value = .{ .text = @tagName(s) } }); } try entries.append(allocator, .{ .key = "time", .value = .{ .text = account.time } });
try cbor.encode(allocator, writer, .{ .map = entries.items });}
fn encodeInfoPayload(allocator: Allocator, writer: anytype, info: InfoEvent) !void { var entries: std.ArrayList(cbor.Value.MapEntry) = .empty; defer entries.deinit(allocator);
if (info.message) |m| { try entries.append(allocator, .{ .key = "message", .value = .{ .text = m } }); } if (info.name) |n| { try entries.append(allocator, .{ .key = "name", .value = .{ .text = n } }); }
try cbor.encode(allocator, writer, .{ .map = entries.items });}
pub const FirehoseClient = struct { io: Io, allocator: Allocator, options: Options, last_seq: ?i64 = null, /// set once the websocket handshake for the current attempt succeeds. /// An endpoint that accepted a connection is a working endpoint, however /// the connection later ended, so the reconnect backoff starts over /// rather than compounding. connected_this_attempt: bool = false, /// frames delivered on the current connection; written by the reader /// task and read by the idle watchdog, hence atomic. frames_this_connection: std.atomic.Value(usize) = .init(0),
pub fn init(io: Io, allocator: Allocator, options: Options) FirehoseClient { return .{ .io = io, .allocator = allocator, .options = options, .last_seq = if (options.cursor) |c| c else null, }; }
pub fn deinit(_: *FirehoseClient) void {}
/// subscribe with a user-provided handler. /// handler must implement: fn onEvent(*@TypeOf(handler), Event) void /// optional: fn onError(*@TypeOf(handler), anyerror) void /// optional: fn onConnect(*@TypeOf(handler), []const u8) void — called /// after the websocket handshake succeeds and before frame delivery /// optional: fn onReconnect(*@TypeOf(handler)) void — called immediately /// before every connection attempt after the initial attempt. /// optional: fn onRawFrame(*@TypeOf(handler), []const u8) void or !void — /// when declared, raw websocket frames are delivered INSTEAD of decoded /// events (the caller owns decoding). Cursor tracking advances only /// after the callback accepts the frame; a fallible callback can reject /// work without making a reconnect skip it. This matters for bounded /// downstream pipelines which may close or fail while the socket lives. /// optional: fn shouldStop(*@TypeOf(handler)) bool — checked before /// every frame delivery and around every reconnect. returning true /// makes subscribe return cleanly (bounded consumers: bootstrap /// capture, tests, graceful shutdown). without it, an interrupted /// read is indistinguishable from a connection error and the /// reconnect loop absorbs it forever; pair a stop flag with /// future.cancel to unblock an idle read. /// blocks until shouldStop (forever without it) — reconnects with /// exponential backoff on disconnect. rotates through hosts on each /// reconnect attempt. pub fn subscribe(self: *FirehoseClient, handler: anytype) Io.Cancelable!void { var backoff: u64 = 1; var host_index: usize = 0; const max_backoff: u64 = 60; var prev_host_index: usize = 0;
while (!stopRequested(handler)) { if (host_index > 0 and comptime @hasDecl(@TypeOf(handler.*), "onReconnect")) handler.onReconnect(); const host = self.options.hosts[host_index % self.options.hosts.len]; const effective_index = host_index % self.options.hosts.len;
// reset backoff on host switch (fresh host deserves a fresh chance) if (host_index > 0 and effective_index != prev_host_index) { backoff = 1; } // ...and after a connection that was actually established. Backoff // exists to spare an endpoint that is refusing us; an endpoint that // accepted a connection and later closed it is not that. Without // this a single-host consumer never resets — effective_index is // always 0, so the host-switch branch above can never fire — and // the delay climbs to max_backoff and stays there for the life of // the process. Observed on a relay that closed every ~113s: every // reconnect paid the full 60s, leaving the tail down 33% of the // time while each individual connection was perfectly healthy. // // Keyed on the handshake, matching jcalabro/atmos // (streaming/client.go: "Reset backoff after successful // connection"). Keying it on delivered frames instead would leave // a connection that is open but idle compounding its backoff — // which for a filtered jetstream consumer is an ordinary state, // not a fault. if (self.connected_this_attempt) backoff = 1;
log.info("connecting to host {d}/{d}: {s}", .{ effective_index + 1, self.options.hosts.len, host });
self.connectAndRead(host, handler) catch |err| { // a stop-interrupted read surfaces as a connection error; // don't report it as one if (stopRequested(handler)) return; if (comptime @hasDecl(@TypeOf(handler.*), "onError")) { handler.onError(err); } else { log.err("firehose error: {s}, reconnecting in {d}s...", .{ @errorName(err), backoff }); } }; if (stopRequested(handler)) return;
prev_host_index = effective_index; host_index += 1; self.io.sleep(Io.Duration.fromSeconds(@intCast(backoff)), .awake) catch |err| switch (err) { // Closing the socket can race the next cancellable operation // on some Io backends. That cancellation belongs to the dead // connection, not to the reconnect loop; an explicit handler // stop remains the only clean termination condition. error.Canceled => { if (stopRequested(handler)) return; continue; }, }; backoff = @min(backoff * 2, max_backoff); } }
fn connectAndRead(self: *FirehoseClient, host: []const u8, handler: anytype) !void { self.connected_this_attempt = false; const ep = try Endpoint.parse(host);
var path_buf: [256]u8 = undefined; var w: std.Io.Writer = .fixed(&path_buf);
try w.writeAll("/xrpc/com.atproto.sync.subscribeRepos"); if (self.last_seq) |cursor| { try w.print("?cursor={d}", .{cursor}); } const path = w.buffered();
log.info("connecting to {s}://{s}:{d}{s}", .{ if (ep.tls) "wss" else "ws", ep.host, ep.port, path });
const websocket = @import("websocket"); var client = try websocket.Client.init(self.io, self.allocator, .{ .host = ep.host, .port = ep.port, .tls = ep.tls, .max_size = self.options.max_message_size, }); defer client.deinit();
// Host header carries the port only when non-default (RFC 9110 §7.2) var host_header_buf: [256]u8 = undefined; const host_header = if (ep.isDefaultPort()) std.fmt.bufPrint(&host_header_buf, "Host: {s}\r\n", .{ep.host}) catch ep.host else std.fmt.bufPrint(&host_header_buf, "Host: {s}:{d}\r\n", .{ ep.host, ep.port }) catch ep.host;
try client.handshake(path, .{ .headers = host_header }); configureKeepalive(&client); self.connected_this_attempt = true;
log.info("firehose connected to {s}", .{ep.host}); if (comptime @hasDecl(@TypeOf(handler.*), "onConnect")) { handler.onConnect(host); }
var ws_handler = WsHandler(@TypeOf(handler.*)){ .allocator = self.allocator, .handler = handler, .client_state = self, }; if (self.options.idle_timeout_ms) |timeout_ms| { // A socket-level receive timeout (SO_RCVTIMEO) is off the table: // std.Io.Threaded treats the resulting EAGAIN on a blocking // socket as a programmer bug and panics in debug builds. Instead // the blocking readLoop runs as a concurrent task and a watchdog // here compares the frame counter across quiet windows, // cancelling the read when a full window passes with no frames. // Fires between timeout and 2x timeout after the last frame. self.frames_this_connection.store(0, .monotonic); var done: Io.Event = .unset; var read_err: ?anyerror = null; const Reader = struct { fn run( io: Io, ws_client: *websocket.Client, h: *WsHandler(@TypeOf(handler.*)), done_ev: *Io.Event, err_out: *?anyerror, ) void { ws_client.readLoop(h) catch |err| { err_out.* = err; }; done_ev.set(io); } }; var future = try self.io.concurrent(Reader.run, .{ self.io, &client, &ws_handler, &done, &read_err }); var joined = false; defer if (!joined) future.cancel(self.io); var last_frames: usize = 0; while (true) { done.waitTimeout(self.io, .{ .duration = .{ .raw = Io.Duration.fromMilliseconds(@intCast(timeout_ms)), .clock = .awake } }) catch |err| switch (err) { error.Timeout => { const frames = self.frames_this_connection.load(.monotonic); if (frames == last_frames) return error.IdleTimeout; last_frames = frames; continue; }, error.Canceled => return error.Canceled, }; // readLoop finished on its own (close frame, socket error, stop) future.await(self.io); joined = true; if (read_err) |err| return err; return; } } else { try client.readLoop(&ws_handler); } }};
/// true iff the handler declares shouldStop and it currently returns truefn stopRequested(handler: anytype) bool { if (comptime @hasDecl(@TypeOf(handler.*), "shouldStop")) { return handler.shouldStop(); } return false;}
fn WsHandler(comptime H: type) type { return struct { allocator: Allocator, handler: *H, client_state: *FirehoseClient,
const Self = @This();
pub fn serverMessage(self: *Self, data: []const u8) !void { _ = self.client_state.frames_this_connection.fetchAdd(1, .monotonic); if (comptime @hasDecl(H, "shouldStop")) { // breaks the read loop; subscribe sees the stop and returns if (self.handler.shouldStop()) return error.SubscriptionStopped; } if (comptime @hasDecl(H, "onRawFrame")) { // A frame is reconnect-safe only after the downstream accepts // it. Peeking before the callback is cheap, but publishing the // cursor before a bounded/fallible pipeline accepts the frame // can turn a local failure into silent upstream data loss. const seq = peekSeq(self.allocator, data); const result = self.handler.onRawFrame(data); if (comptime @TypeOf(result) != void) try result; if (seq) |s| self.client_state.last_seq = s; return; } var arena = std.heap.ArenaAllocator.init(self.allocator); defer arena.deinit();
const event = decodeFrame(arena.allocator(), data) catch |err| { log.debug("frame decode error: {s}", .{@errorName(err)}); return; };
if (event.seq()) |s| { self.client_state.last_seq = s; }
self.handler.onEvent(event); }
pub fn close(_: *Self) void { log.info("firehose connection closed", .{}); } };}
/// enable TCP keepalive so reads don't block forever when a peer/// disappears without FIN/RST (network partition, crash, power loss)./// detection time: 10s idle + 5s × 2 probes = 20s.fn configureKeepalive(client: anytype) void { 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 ===
test "decode frame header" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
// simulate a frame: header {op: 1, t: "#info"} + payload {name: "OutdatedCursor"} const header_bytes = [_]u8{ 0xa2, // map(2) 0x61, 't', 0x65, '#', 'i', 'n', 'f', 'o', // "t": "#info" 0x62, 'o', 'p', 0x01, // "op": 1 }; const payload_bytes = [_]u8{ 0xa1, // map(1) 0x64, 'n', 'a', 'm', 'e', // "name" 0x6e, 'O', 'u', 't', 'd', 'a', 't', 'e', 'd', 'C', 'u', 'r', 's', 'o', 'r', // "OutdatedCursor" };
var frame: [header_bytes.len + payload_bytes.len]u8 = undefined; @memcpy(frame[0..header_bytes.len], &header_bytes); @memcpy(frame[header_bytes.len..], &payload_bytes);
const event = try decodeFrame(alloc, &frame); const info = event.info; try std.testing.expectEqualStrings("OutdatedCursor", info.name.?);}
test "decode identity frame" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
// build frame via encoder for cleaner test const original = Event{ .identity = .{ .seq = 42, .did = "did:plc:test", .time = "2024-01-15T10:30:00Z", } }; const frame = try encodeFrame(alloc, original);
const event = try decodeFrame(alloc, frame); const identity = event.identity; try std.testing.expectEqual(@as(i64, 42), identity.seq); try std.testing.expectEqualStrings("did:plc:test", identity.did); try std.testing.expectEqualStrings("2024-01-15T10:30:00Z", identity.time);}
test "Event.seq works" { const info_event = Event{ .info = .{ .name = "test" } }; try std.testing.expect(info_event.seq() == null);
const identity_event = Event{ .identity = .{ .seq = 42, .did = "did:plc:test", .time = "2024-01-15T10:30:00Z", } }; try std.testing.expectEqual(@as(i64, 42), identity_event.seq().?);}
// === encoder tests ===
test "encode → decode info frame" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
const original = Event{ .info = .{ .name = "OutdatedCursor", .message = "cursor is behind", } };
const frame = try encodeFrame(alloc, original); const decoded = try decodeFrame(alloc, frame);
try std.testing.expectEqualStrings("OutdatedCursor", decoded.info.name.?); try std.testing.expectEqualStrings("cursor is behind", decoded.info.message.?);}
test "encode → decode identity frame" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
const original = Event{ .identity = .{ .seq = 42, .did = "did:plc:test123", .time = "2024-01-15T10:30:00Z", .handle = "alice.bsky.social", } };
const frame = try encodeFrame(alloc, original); const decoded = try decodeFrame(alloc, frame);
const id = decoded.identity; try std.testing.expectEqual(@as(i64, 42), id.seq); try std.testing.expectEqualStrings("did:plc:test123", id.did); try std.testing.expectEqualStrings("2024-01-15T10:30:00Z", id.time); try std.testing.expectEqualStrings("alice.bsky.social", id.handle.?);}
test "encode → decode account frame" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
const original = Event{ .account = .{ .seq = 100, .did = "did:plc:suspended", .time = "2024-01-15T10:30:00Z", .active = false, .status = .suspended, } };
const frame = try encodeFrame(alloc, original); const decoded = try decodeFrame(alloc, frame);
const acct = decoded.account; try std.testing.expectEqual(@as(i64, 100), acct.seq); try std.testing.expectEqualStrings("did:plc:suspended", acct.did); try std.testing.expectEqualStrings("2024-01-15T10:30:00Z", acct.time); try std.testing.expectEqual(false, acct.active); try std.testing.expectEqual(AccountStatus.suspended, acct.status.?);}
test "encode → decode commit frame with record" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
const record: cbor.Value = .{ .map = &.{ .{ .key = "$type", .value = .{ .text = "app.bsky.feed.post" } }, .{ .key = "text", .value = .{ .text = "hello firehose" } }, } }; const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit"); const prev_data = try cbor.Cid.forDagCbor(alloc, "prev-data");
const original = Event{ .commit = .{ .seq = 999, .repo = "did:plc:poster", .rev = "3k2abc000000", .time = "2024-01-15T10:30:00Z", .since = "3k2abd000000", .commit = commit_cid, .prev_data = prev_data, .ops = &.{.{ .action = .create, .collection = "app.bsky.feed.post", .rkey = "3k2abc", .record = record, }}, } };
const frame = try encodeFrame(alloc, original); const decoded = try decodeFrame(alloc, frame);
const commit = decoded.commit; try std.testing.expectEqual(@as(i64, 999), commit.seq); try std.testing.expectEqualStrings("did:plc:poster", commit.repo); try std.testing.expectEqualStrings("3k2abc000000", commit.rev); try std.testing.expectEqualStrings("2024-01-15T10:30:00Z", commit.time); try std.testing.expectEqualStrings("3k2abd000000", commit.since.?); try std.testing.expectEqualSlices(u8, commit_cid.raw, commit.commit.?.raw); try std.testing.expect(commit.blocks.len > 0); try std.testing.expectEqualSlices(u8, prev_data.raw, commit.prev_data.?.raw); try std.testing.expectEqual(@as(usize, 0), commit.blobs.len); try std.testing.expectEqual(@as(usize, 1), commit.ops.len);
const op = commit.ops[0]; try std.testing.expectEqual(CommitAction.create, op.action); try std.testing.expectEqualStrings("app.bsky.feed.post/3k2abc", op.path); try std.testing.expectEqualStrings("app.bsky.feed.post", op.collection); try std.testing.expectEqualStrings("3k2abc", op.rkey); try std.testing.expect(op.cid != null);
// record should be decoded from the CAR blocks const rec = op.record.?; try std.testing.expectEqualStrings("hello firehose", rec.getString("text").?); try std.testing.expectEqualStrings("app.bsky.feed.post", rec.getString("$type").?);
const mst_ops = try commit.toMstOperations(alloc); try std.testing.expectEqual(@as(usize, 1), mst_ops.len); try std.testing.expectEqualStrings("app.bsky.feed.post/3k2abc", mst_ops[0].path); try std.testing.expectEqualSlices(u8, op.cid.?.raw, mst_ops[0].value.?); try std.testing.expect(mst_ops[0].prev == null);}
test "encode → decode commit with delete (no record)" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator(); const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit"); const prev_data = try cbor.Cid.forDagCbor(alloc, "prev-data"); const prev_record = try cbor.Cid.forDagCbor(alloc, "prev-record");
const original = Event{ .commit = .{ .seq = 500, .repo = "did:plc:deleter", .rev = "3k2xyz000000", .time = "2024-01-15T10:30:00Z", .commit = commit_cid, .prev_data = prev_data, .ops = &.{.{ .action = .delete, .collection = "app.bsky.feed.post", .rkey = "abc123", .prev = prev_record, .record = null, }}, } };
const frame = try encodeFrame(alloc, original); const decoded = try decodeFrame(alloc, frame);
try std.testing.expectEqual(@as(i64, 500), decoded.commit.seq); try std.testing.expectEqualStrings("3k2xyz000000", decoded.commit.rev); try std.testing.expectEqualStrings("2024-01-15T10:30:00Z", decoded.commit.time); try std.testing.expectEqual(@as(usize, 1), decoded.commit.ops.len); try std.testing.expectEqual(CommitAction.delete, decoded.commit.ops[0].action); try std.testing.expect(decoded.commit.ops[0].cid == null); try std.testing.expectEqualSlices(u8, prev_record.raw, decoded.commit.ops[0].prev.?.raw); try std.testing.expect(decoded.commit.ops[0].record == null);}
test "encode → decode commit with null prevData" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator(); const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit");
const original = Event{ .commit = .{ .seq = 501, .repo = "did:plc:initial", .rev = "3k2xyz000001", .time = "2024-01-15T10:30:00Z", .since = null, .commit = commit_cid, .blocks = "car bytes", .prev_data = null, .ops = &.{}, } };
const frame = try encodeFrame(alloc, original); const decoded = try decodeFrame(alloc, frame);
try std.testing.expectEqual(@as(i64, 501), decoded.commit.seq); try std.testing.expectEqualStrings("did:plc:initial", decoded.commit.repo); try std.testing.expectEqualStrings("3k2xyz000001", decoded.commit.rev); try std.testing.expect(decoded.commit.since == null); try std.testing.expect(decoded.commit.prev_data == null); try std.testing.expectEqualStrings("car bytes", decoded.commit.blocks);}
test "commit op paths are validated" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator(); const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit");
const event = Event{ .commit = .{ .seq = 502, .repo = "did:plc:badpath", .rev = "3k2xyz000002", .time = "2024-01-15T10:30:00Z", .commit = commit_cid, .blocks = "car bytes", .prev_data = null, .ops = &.{.{ .action = .create, .path = "app.bsky.feed.post/not/one/rkey", .collection = "app.bsky.feed.post", .rkey = "unused", }}, } };
try std.testing.expectError(error.InvalidRepoPath, encodeFrame(alloc, event));}
test "repo path decoding preserves invalid paths for downstream policy" { const no_slash = splitRepoPath("nosslashatall"); try std.testing.expectEqualStrings("nosslashatall", no_slash.collection); try std.testing.expectEqualStrings("", no_slash.rkey);
const extra_slash = splitRepoPath("app.bsky.feed.post/bad/key"); try std.testing.expectEqualStrings("app.bsky.feed.post", extra_slash.collection); try std.testing.expectEqualStrings("bad/key", extra_slash.rkey);
try std.testing.expectError(error.InvalidRepoPath, parseRepoPath("nosslashatall")); try std.testing.expectError(error.InvalidRepoPath, parseRepoPath("app.bsky.feed.post/bad/key"));}
test "validate commit event accepts structurally valid event" { const alloc = std.testing.allocator; const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit"); defer alloc.free(commit_cid.raw); const record_cid = try cbor.Cid.forDagCbor(alloc, "record"); defer alloc.free(record_cid.raw); const prev_record_cid = try cbor.Cid.forDagCbor(alloc, "prev-record"); defer alloc.free(prev_record_cid.raw); const rev = Tid.fromTimestamp(1704067200000000, 1); const since = Tid.fromTimestamp(1704067100000000, 1);
try validateCommitEvent(.{ .seq = 600, .repo = "did:plc:validator", .rev = rev.str(), .time = "2024-01-15T10:30:00Z", .since = since.str(), .commit = commit_cid, .blocks = "car bytes", .ops = &.{ .{ .action = .create, .path = "app.bsky.feed.post/3k2valid", .collection = "unused.invalid.value", .rkey = "unused", .cid = record_cid, }, .{ .action = .update, .collection = "app.bsky.feed.post", .rkey = "3k2valid", .cid = record_cid, .prev = prev_record_cid, }, .{ .action = .delete, .collection = "app.bsky.feed.post", .rkey = "3k2delete", .prev = prev_record_cid, }, }, });}
test "validate commit event rejects commit CID in since" { const alloc = std.testing.allocator; const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit"); defer alloc.free(commit_cid.raw); const rev = Tid.fromTimestamp(1704067200000000, 1);
try std.testing.expectError(error.InvalidSince, validateCommitEvent(.{ .seq = 601, .repo = "did:plc:validator", .rev = rev.str(), .time = "2024-01-15T10:30:00Z", .since = "bafyreihyrpefhacm2x43w6c5ylz6dibjtnfubn2noldubqefbzzrskc6sy", .commit = commit_cid, .blocks = "car bytes", .ops = &.{}, }));}
test "validate commit event enforces op CID nullability" { const alloc = std.testing.allocator; const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit"); defer alloc.free(commit_cid.raw); const record_cid = try cbor.Cid.forDagCbor(alloc, "record"); defer alloc.free(record_cid.raw); const prev_record_cid = try cbor.Cid.forDagCbor(alloc, "prev-record"); defer alloc.free(prev_record_cid.raw); const rev = Tid.fromTimestamp(1704067200000000, 1);
try std.testing.expectError(error.MissingRecordCid, validateCommitEvent(.{ .seq = 602, .repo = "did:plc:validator", .rev = rev.str(), .time = "2024-01-15T10:30:00Z", .commit = commit_cid, .ops = &.{.{ .action = .create, .collection = "app.bsky.feed.post", .rkey = "3k2missing", }}, }));
try std.testing.expectError(error.UnexpectedRecordCid, validateCommitEvent(.{ .seq = 603, .repo = "did:plc:validator", .rev = rev.str(), .time = "2024-01-15T10:30:00Z", .commit = commit_cid, .ops = &.{.{ .action = .delete, .collection = "app.bsky.feed.post", .rkey = "3k2delete", .cid = record_cid, }}, }));
try std.testing.expectError(error.UnexpectedPrevRecordCid, validateCommitEvent(.{ .seq = 604, .repo = "did:plc:validator", .rev = rev.str(), .time = "2024-01-15T10:30:00Z", .commit = commit_cid, .ops = &.{.{ .action = .create, .collection = "app.bsky.feed.post", .rkey = "3k2create", .cid = record_cid, .prev = prev_record_cid, }}, }));}
test "commit event builder names since as since_rev" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator(); const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit"); const record_cid = try cbor.Cid.forDagCbor(alloc, "record"); const rev = Tid.fromTimestamp(1704067200000000, 1); const since = Tid.fromTimestamp(1704067100000000, 1);
const frame = try encodeCommitEvent(alloc, .{ .seq = 700, .repo_did = "did:plc:builder", .commit_cid = commit_cid, .rev = rev.str(), .since_rev = since.str(), .prev_data = null, .blocks = "car bytes", .ops = &.{.{ .action = .create, .collection = "app.bsky.feed.post", .rkey = "3k2builder", .cid = record_cid, }}, .time = "2024-01-15T10:30:00Z", .rebase = true, });
const decoded = try decodeFrame(alloc, frame); try std.testing.expectEqual(@as(i64, 700), decoded.commit.seq); try std.testing.expectEqualStrings("did:plc:builder", decoded.commit.repo); try std.testing.expectEqualStrings(rev.str(), decoded.commit.rev); try std.testing.expectEqualStrings(since.str(), decoded.commit.since.?); try std.testing.expect(decoded.commit.prev_data == null); try std.testing.expect(decoded.commit.rebase); try std.testing.expectEqual(@as(usize, 1), decoded.commit.ops.len); try std.testing.expectEqualStrings("app.bsky.feed.post/3k2builder", decoded.commit.ops[0].path); try std.testing.expectEqualSlices(u8, record_cid.raw, decoded.commit.ops[0].cid.?.raw);}
test "commit event builder rejects commit CID as since_rev" { const alloc = std.testing.allocator; const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit"); defer alloc.free(commit_cid.raw); const rev = Tid.fromTimestamp(1704067200000000, 1);
try std.testing.expectError(error.InvalidSince, encodeCommitEvent(alloc, .{ .seq = 701, .repo_did = "did:plc:builder", .commit_cid = commit_cid, .rev = rev.str(), .since_rev = "bafyreihyrpefhacm2x43w6c5ylz6dibjtnfubn2noldubqefbzzrskc6sy", .prev_data = null, .blocks = "car bytes", .ops = &.{}, .time = "2024-01-15T10:30:00Z", }));}
test "encode → decode sync frame" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
const original = Event{ .sync = .{ .seq = 777, .did = "did:plc:sync", .rev = "3k2sync00000", .time = "2024-01-15T10:30:00Z", .blocks = "car bytes", } };
const frame = try encodeFrame(alloc, original); const decoded = try decodeFrame(alloc, frame);
try std.testing.expectEqual(@as(i64, 777), decoded.sync.seq); try std.testing.expectEqualStrings("did:plc:sync", decoded.sync.did); try std.testing.expectEqualStrings("3k2sync00000", decoded.sync.rev); try std.testing.expectEqualStrings("car bytes", decoded.sync.blocks);}
test "Endpoint.parse accepts URLs and bare hostnames" { const bare = try Endpoint.parse("bsky.network"); try std.testing.expectEqualStrings("bsky.network", bare.host); try std.testing.expectEqual(@as(u16, 443), bare.port); try std.testing.expect(bare.tls);
const wss = try Endpoint.parse("wss://bsky.network/xrpc/whatever?x=1"); try std.testing.expectEqualStrings("bsky.network", wss.host); try std.testing.expectEqual(@as(u16, 443), wss.port); try std.testing.expect(wss.tls);
const ws = try Endpoint.parse("ws://localhost:7777"); try std.testing.expectEqualStrings("localhost", ws.host); try std.testing.expectEqual(@as(u16, 7777), ws.port); try std.testing.expect(!ws.tls);
const ws_default = try Endpoint.parse("ws://localhost"); try std.testing.expectEqual(@as(u16, 80), ws_default.port);
try std.testing.expectError(error.InvalidEndpoint, Endpoint.parse("https://bsky.network")); try std.testing.expectError(error.InvalidEndpoint, Endpoint.parse("wss://")); try std.testing.expectError(error.InvalidEndpoint, Endpoint.parse("ws://host:notaport"));}
test "peekSeq extracts seq without full decode" { // header {op:1, t:"#commit"} + payload {seq: 42} const frame = [_]u8{ 0xa2, 0x62, 'o', 'p', 0x01, 0x61, 't', 0x67, '#', 'c', 'o', 'm', 'm', 'i', 't', 0xa1, 0x63, 's', 'e', 'q', 0x18, 0x2a, }; try std.testing.expectEqual(@as(?i64, 42), peekSeq(std.testing.allocator, &frame)); try std.testing.expectEqual(@as(?i64, null), peekSeq(std.testing.allocator, "junk"));}
test "raw frame cursor advances only after downstream acceptance" { const frame = [_]u8{ 0xa2, 0x62, 'o', 'p', 0x01, 0x61, 't', 0x67, '#', 'c', 'o', 'm', 'm', 'i', 't', 0xa1, 0x63, 's', 'e', 'q', 0x18, 0x2a, }; var client = FirehoseClient.init(std.testing.io, std.testing.allocator, .{}); defer client.deinit();
const Rejecting = struct { calls: usize = 0, fn onRawFrame(self: *@This(), _: []const u8) !void { self.calls += 1; return error.PipelineClosed; } }; var rejecting: Rejecting = .{}; var rejecting_ws: WsHandler(Rejecting) = .{ .allocator = std.testing.allocator, .handler = &rejecting, .client_state = &client, }; try std.testing.expectError(error.PipelineClosed, rejecting_ws.serverMessage(&frame)); try std.testing.expectEqual(@as(usize, 1), rejecting.calls); try std.testing.expectEqual(@as(?i64, null), client.last_seq);
const Accepting = struct { client: *FirehoseClient, cursor_during_callback: ?i64 = null, fn onRawFrame(self: *@This(), _: []const u8) !void { self.cursor_during_callback = self.client.last_seq; } }; var accepting: Accepting = .{ .client = &client }; var accepting_ws: WsHandler(Accepting) = .{ .allocator = std.testing.allocator, .handler = &accepting, .client_state = &client, }; try accepting_ws.serverMessage(&frame); try std.testing.expectEqual(@as(?i64, null), accepting.cursor_during_callback); try std.testing.expectEqual(@as(?i64, 42), client.last_seq);}
test "shouldStop hook: subscribe returns cleanly instead of reconnecting" { var threaded: Io.Threaded = .init(std.testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io();
const StopHandler = struct { stop: bool = false, errors: usize = 0, fn onEvent(_: *@This(), _: Event) void {} fn shouldStop(self: *@This()) bool { return self.stop; } fn onError(self: *@This(), _: anyerror) void { // first failed connect (port 9, nothing listening) requests stop self.errors += 1; self.stop = true; } };
var client = FirehoseClient.init(io, std.testing.allocator, .{ .hosts = &.{"ws://127.0.0.1:9"}, }); defer client.deinit();
// stop set before the first connect: returns without dialing var pre_stopped: StopHandler = .{ .stop = true }; try client.subscribe(&pre_stopped); try std.testing.expectEqual(@as(usize, 0), pre_stopped.errors);
// stop set from onError: returns after the first failed connect, // without sleeping into the reconnect backoff var stops_on_error: StopHandler = .{}; try client.subscribe(&stops_on_error); try std.testing.expectEqual(@as(usize, 1), stops_on_error.errors);}
test "handlers without shouldStop never stop-request" { const Plain = struct { fn onEvent(_: *@This(), _: Event) void {} }; var h: Plain = .{}; try std.testing.expect(!stopRequested(&h));}
test "onReconnect runs at the retry attempt boundary" { var threaded: Io.Threaded = .init(std.testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io();
const Handler = struct { stop: bool = false, errors: usize = 0, reconnects: usize = 0, fn onEvent(_: *@This(), _: Event) void {} fn shouldStop(self: *@This()) bool { return self.stop; } fn onError(self: *@This(), _: anyerror) void { self.errors += 1; if (self.errors == 2) self.stop = true; } fn onReconnect(self: *@This()) void { self.reconnects += 1; } };
var client = FirehoseClient.init(io, std.testing.allocator, .{ .hosts = &.{"ws://127.0.0.1:9"}, }); defer client.deinit(); var handler: Handler = .{}; try client.subscribe(&handler); try std.testing.expectEqual(@as(usize, 2), handler.errors); try std.testing.expectEqual(@as(usize, 1), handler.reconnects);}