//! 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 event pub 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 wire const 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 Event pub 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 true fn 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); }