atproto relay in zig zlay.waow.tech
relay zig atproto
Something went wrong. Try again.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273//! relay subscriber — per-host firehose worker//!//! connects to a single PDS or relay upstream, decodes firehose frames//! using the zat SDK's CBOR codec, validates commit frames, and persists//! all events to disk with relay-assigned sequence numbers before broadcast.//!//! managed by the Slurper, which spawns one Subscriber per tracked host.
const std = @import("std");const backend_config = @import("backend_config");const websocket = @import("websocket");const zat = @import("zat");const broadcaster = @import("broadcaster.zig");const validator_mod = @import("validator.zig");const event_log_mod = @import("event_log.zig");const collection_index_mod = @import("collection_index/index.zig");const resync_mod = @import("collection_index/resync.zig");const status_checker_mod = @import("status_checker.zig");const frame_worker_mod = @import("frame_worker.zig");const host_ops_mod = @import("host_ops.zig");const ws_ping_mod = @import("ws_ping.zig");const util = @import("util/util.zig");
const Allocator = std.mem.Allocator;const Io = std.Io;const log = std.log.scoped(.relay);
const timestamp = util.timestamp;const milliTimestamp = util.milliTimestamp;const microTimestamp = util.microTimestamp;const nanoTimestamp = util.nanoTimestamp;
const max_consecutive_failures = 15;const cursor_flush_interval_sec = 4; // flush cursor to DB every N seconds (Go relay: 4s)const ping_interval_sec = 30; // keepalive ping interval (Go relay: 30s)const max_ping_failures = 4; // close connection after N consecutive ping failures (Go relay: 4)
// per-host rate limits — backpressure thresholds (matches indigo: 50/s, 2500/hr, 20k/day)// blocking instead of dropping means these trigger TCP backpressure on the upstream PDSconst default_per_second_limit: u64 = 50;const default_per_hour_limit: u64 = 2_500;const default_per_day_limit: u64 = 20_000;
// trusted hosts get much higher limits (Go relay: 5000/s, 50M/hr, 500M/day)const trusted_per_second_limit: u64 = 5_000;const trusted_per_hour_limit: u64 = 50_000_000;const trusted_per_day_limit: u64 = 500_000_000;
// Go relay: TrustedDomains config — hosts matching these suffixes get trusted limitsconst trusted_suffixes: []const []const u8 = &.{".host.bsky.network"};
pub fn isTrustedHost(hostname: []const u8) bool { for (trusted_suffixes) |suffix| { if (std.mem.endsWith(u8, hostname, suffix)) return true; } return false;}
/// compute rate limits scaled by account count (matches Go relay: slurper.go ComputeLimiterCounts)pub fn computeLimits(trusted: bool, account_count: u64) struct { sec: u64, hour: u64, day: u64 } { if (trusted) return .{ .sec = trusted_per_second_limit, .hour = trusted_per_hour_limit, .day = trusted_per_day_limit, }; return .{ .sec = default_per_second_limit + account_count / 1000, .hour = default_per_hour_limit + account_count, .day = default_per_day_limit + account_count * 10, };}
pub const Options = struct { hostname: []const u8 = "bsky.network", max_message_size: usize = 5 * 1024 * 1024, host_id: u64 = 0, account_count: u64 = 0, ca_bundle: ?std.crypto.Certificate.Bundle = null,};
/// simple sliding window rate limiter — tracks event counts per second/hour/day./// Sliding window rate limiter (same algorithm as Go relay's github.com/RussellLuo/slidingwindow)./// Uses millisecond timestamps for sub-second precision (critical for the 1-second window)./// Effective count = weight * prev_count + curr_count/// where weight = (window_size - elapsed) / window_size.const RateLimiter = struct { sec: SlidingWindow = .{ .size_ms = 1_000 }, hour: SlidingWindow = .{ .size_ms = 3_600_000 }, day: SlidingWindow = .{ .size_ms = 86_400_000 },
sec_limit: std.atomic.Value(u64) = .{ .raw = default_per_second_limit }, hour_limit: std.atomic.Value(u64) = .{ .raw = default_per_hour_limit }, day_limit: std.atomic.Value(u64) = .{ .raw = default_per_day_limit },
const SlidingWindow = struct { size_ms: i64, curr_start: i64 = 0, curr_count: u64 = 0, prev_count: u64 = 0,
fn advance(self: *SlidingWindow, now_ms: i64) void { const new_start = now_ms - @mod(now_ms, self.size_ms); if (new_start <= self.curr_start) return; const diff = @divTrunc(new_start - self.curr_start, self.size_ms); self.prev_count = if (diff == 1) self.curr_count else 0; self.curr_start = new_start; self.curr_count = 0; }
fn effectiveCount(self: *const SlidingWindow, now_ms: i64) u64 { const elapsed = now_ms - self.curr_start; const remaining = self.size_ms - elapsed; if (remaining <= 0 or self.prev_count == 0) return self.curr_count; const weighted_prev = self.prev_count * @as(u64, @intCast(remaining)) / @as(u64, @intCast(self.size_ms)); return weighted_prev + self.curr_count; } };
const Result = enum { allowed, sec, hour, day };
fn allow(self: *RateLimiter, now_ms: i64) Result { self.sec.advance(now_ms); self.hour.advance(now_ms); self.day.advance(now_ms);
if (self.sec.effectiveCount(now_ms) >= self.sec_limit.load(.monotonic)) return .sec; if (self.hour.effectiveCount(now_ms) >= self.hour_limit.load(.monotonic)) return .hour; if (self.day.effectiveCount(now_ms) >= self.day_limit.load(.monotonic)) return .day;
self.sec.curr_count += 1; self.hour.curr_count += 1; self.day.curr_count += 1; return .allowed; }
/// Block until all rate limit windows allow the event. /// Checks every 100ms, matching indigo's waitForLimiter behavior. /// Returns which tier (if any) caused a wait, for metrics. fn waitForAllow(self: *RateLimiter, shutdown: *std.atomic.Value(bool), io: Io) Result { // fast path: no waiting needed const now_ms = milliTimestamp(io); self.sec.advance(now_ms); self.hour.advance(now_ms); self.day.advance(now_ms);
if (self.sec.effectiveCount(now_ms) < self.sec_limit.load(.monotonic) and self.hour.effectiveCount(now_ms) < self.hour_limit.load(.monotonic) and self.day.effectiveCount(now_ms) < self.day_limit.load(.monotonic)) { self.sec.curr_count += 1; self.hour.curr_count += 1; self.day.curr_count += 1; return .allowed; }
// slow path: poll every 100ms until allowed (creates TCP backpressure) var waited: Result = .sec; while (!shutdown.load(.acquire)) { io.sleep(Io.Duration.fromMilliseconds(100), .awake) catch {};
const t = milliTimestamp(io); self.sec.advance(t); self.hour.advance(t); self.day.advance(t);
if (self.sec.effectiveCount(t) >= self.sec_limit.load(.monotonic)) { waited = .sec; continue; } if (self.hour.effectiveCount(t) >= self.hour_limit.load(.monotonic)) { waited = .hour; continue; } if (self.day.effectiveCount(t) >= self.day_limit.load(.monotonic)) { waited = .day; continue; }
self.sec.curr_count += 1; self.hour.curr_count += 1; self.day.curr_count += 1; return waited; } return waited; }
/// update rate limits from another thread (e.g. admin API). pub fn updateLimits(self: *RateLimiter, sec: u64, hour: u64, day: u64) void { self.sec_limit.store(sec, .monotonic); self.hour_limit.store(hour, .monotonic); self.day_limit.store(day, .monotonic); }};
pub const Subscriber = struct { allocator: Allocator, io: Io, options: Options, bc: *broadcaster.Broadcaster, validator: *validator_mod.Validator, persist: ?*event_log_mod.DiskPersist, collection_index: ?*collection_index_mod.CollectionIndex = null, resyncer: ?*resync_mod.Resyncer = null, status_checker: ?*status_checker_mod.StatusChecker = null, pool: ?*frame_worker_mod.FramePool = null, /// dedicated Threaded io for frame workers — safe from plain OS threads pool_io: ?Io = null, /// host ops queue — pushes rare DB ops to a background thread host_ops: ?*host_ops_mod.HostOpsQueue = null, /// coalescing cursor map — subscriber writes latest seq atomically, worker sweeps every 5s cursor_map: ?*host_ops_mod.CursorMap = null, cursor_slot: ?u32 = null, shutdown: *std.atomic.Value(bool), last_upstream_seq: ?u64 = null, last_cursor_flush: i64 = 0, /// ms timestamp of the last sign of life from this worker: a frame, a /// completed handshake, or the start of a connect attempt. Written by the /// subscriber fiber, read by the admin reconnect fiber, hence atomic. /// /// It stamps connect *attempts*, not just delivered frames, because a /// worker can wedge inside connectAndRead -- DNS, TLS and the websocket /// handshake all park on the same read path that lost wakes in the /// 2026-08-18 incident. connectAndRead never returning means run()'s /// reconnect loop never iterates either, so such a worker has no backoff /// to fall back on. A genuinely unreachable host keeps attempting, so it /// keeps stamping and stays inside the idle guard. last_progress_ms: std.atomic.Value(i64) = .{ .raw = 0 }, /// keepalive pinger knobs, shared across all subscribers and settable at /// runtime via /admin/ws-ping. null (tests, standalone) disables the /// zio-side pinger; the Threaded path is unaffected either way. ws_ping: ?*ws_ping_mod.WsPing = null, rate_limiter: RateLimiter = .{},
// per-host shutdown (e.g. FutureCursor — stops only this subscriber) host_shutdown: std.atomic.Value(bool) = .{ .raw = false },
pub fn init( allocator: Allocator, io: Io, bc: *broadcaster.Broadcaster, val: *validator_mod.Validator, persist: ?*event_log_mod.DiskPersist, shutdown: *std.atomic.Value(bool), options: Options, ) Subscriber { const trusted = isTrustedHost(options.hostname); const limits = computeLimits(trusted, options.account_count); return .{ .allocator = allocator, .io = io, .options = options, .bc = bc, .validator = val, .persist = persist, .shutdown = shutdown, .rate_limiter = .{ .sec_limit = .{ .raw = limits.sec }, .hour_limit = .{ .raw = limits.hour }, .day_limit = .{ .raw = limits.day }, }, }; }
/// check if this subscriber should stop (global or per-host shutdown) fn shouldStop(self: *Subscriber) bool { return self.shutdown.load(.acquire) or self.host_shutdown.load(.acquire); }
/// run the subscriber loop. reconnects with exponential backoff. /// blocks until shutdown flag is set or host is exhausted. pub fn run(self: *Subscriber) void { var backoff: u64 = 1; const max_backoff: u64 = 60;
// cursor is set at spawn time by slurper if (self.last_upstream_seq) |seq| { log.info("host {s}: resuming from cursor {d}", .{ self.options.hostname, seq }); }
while (!self.shouldStop()) { log.info("host {s}: connecting...", .{self.options.hostname}); self.last_progress_ms.store(milliTimestamp(self.io), .release);
self.connectAndRead() catch |err| { if (self.shouldStop()) return; log.err("host {s}: error: {s}, reconnecting in {d}s...", .{ self.options.hostname, @errorName(err), backoff }); };
if (self.shouldStop()) return;
// track failures for this host (pushed to background thread via host_ops queue) if (self.options.host_id > 0) { if (self.host_ops) |hq| { hq.push(.{ .host_id = self.options.host_id, .kind = .increment_failures, .payload = .{ .host_shutdown = &self.host_shutdown }, }); } }
// backoff sleep in small increments so we can check shutdown var remaining: u64 = backoff; while (remaining > 0 and !self.shouldStop()) { const chunk = @min(remaining, 1); self.io.sleep(Io.Duration.fromSeconds(@intCast(chunk)), .awake) catch {}; remaining -= chunk; } backoff = @min(backoff * 2, max_backoff); } }
/// flush cursor position (atomic store to coalescing map — worker sweeps every 5s) fn flushCursor(self: *Subscriber) void { if (self.options.host_id == 0) return; const seq = self.last_upstream_seq orelse return; if (self.cursor_map) |cm| { if (self.cursor_slot) |slot| { cm.update(slot, seq); } } }
/// resolve and connect trying addresses one at a time, instead of /// HostName.connect's racing per-address attempts (Io.Group + loser /// cancellation). a relay reconnecting to thousands of hosts wants the /// smallest possible concurrency surface per connect, not per-connect /// latency: sequential attempts keep the subscriber fiber the only /// actor, with no group spawns and no cancel storms in the hot path. fn connectSequential(io: Io, host_name: Io.net.HostName) !Io.net.Stream { var canonical_name_buffer: [Io.net.HostName.max_len]u8 = undefined; var lookup_buffer: [32]Io.net.HostName.LookupResult = undefined; var lookup_queue: Io.Queue(Io.net.HostName.LookupResult) = .init(&lookup_buffer); try host_name.lookup(io, &lookup_queue, .{ .port = 443, .canonical_name_buffer = &canonical_name_buffer, }); var last_err: ?anyerror = null; while (lookup_queue.getOne(io)) |result| switch (result) { .address => |address| { const stream = address.connect(io, .{ .mode = .stream }) catch |e| { last_err = e; continue; }; return stream; }, .canonical_name => continue, } else |err| switch (err) { error.Closed => {}, else => |e| return e, } return if (last_err) |e| e else error.ConnectionRefused; }
fn connectAndRead(self: *Subscriber) !void { 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_upstream_seq) |cursor| { try w.print("?cursor={d}", .{cursor}); } const path = w.buffered();
// Connect on the same io that will read the socket. The old split — // DNS + connect on pool_io, then TLS + WebSocket on self.io — existed // because Io.Uring never implemented netLookup. Under Threaded both // are the same backend, so it was a no-op; under any evented backend // it hands a foreign-domain fd to the reader, which is a defect (a // blocking fd read from a fiber). Backends that lack netLookup are // expected to delegate it, not to have callers route around them. // // each phase labels its catch site so operators can see which layer // is failing across the fleet via relay_subscriber_disconnect_total. const host_name = Io.net.HostName.init(self.options.hostname) catch |e| { _ = self.bc.stats.subscriber_disconnect_dns_connect.fetchAdd(1, .monotonic); return e; }; const net_stream = connectSequential(self.io, host_name) catch |e| { _ = self.bc.stats.subscriber_disconnect_dns_connect.fetchAdd(1, .monotonic); return e; };
var client = websocket.Client.initWithStream(self.io, self.allocator, net_stream, .{ .host = self.options.hostname, .port = 443, .tls = true, .max_size = self.options.max_message_size, .ca_bundle = self.options.ca_bundle, }) catch |e| { _ = self.bc.stats.subscriber_disconnect_tls_init.fetchAdd(1, .monotonic); return e; }; defer client.deinit();
var host_header_buf: [256]u8 = undefined; const host_header = std.fmt.bufPrint( &host_header_buf, "Host: {s}\r\n", .{self.options.hostname}, ) catch self.options.hostname;
client.handshake(path, .{ .headers = host_header }) catch |e| { _ = self.bc.stats.subscriber_disconnect_ws_handshake.fetchAdd(1, .monotonic); return e; }; log.info("host {s}: connected", .{self.options.hostname}); self.last_progress_ms.store(milliTimestamp(self.io), .release);
// reset failures on successful connect (via host_ops queue) if (self.options.host_id > 0) { if (self.host_ops) |hq| { hq.push(.{ .host_id = self.options.host_id, .kind = .reset_failures, .payload = .{ .none = {} }, }); } }
var conn = ConnState{ .activity_ms = .{ .raw = milliTimestamp(self.io) }, }; var handler = FrameHandler{ .subscriber = self, .conn = &conn, };
const heartbeat: websocket.Client.HeartbeatConfig = .{ .interval_ms = @intCast(@max(1, if (self.ws_ping) |p| p.interval_sec.load(.acquire) else ping_interval_sec) * 1000), .max_failures = @intCast(if (self.ws_ping) |p| p.max_failures.load(.acquire) else max_ping_failures), };
// read_loop disconnects only count toward `failed_attempts` if the // *next* handshake also fails — reset_failures was already pushed // above on successful handshake. if (comptime backend_config.use_zio) { if (self.ws_ping != null) { // under zio, readLoopWithHeartbeat's SO_RCVTIMEO mechanism is // inert: netRead parks the fiber in epoll and the socket // timeout never fires (the 2026-08-18/27 half-open wedges). // run the read loop as a cancelable task and drive liveness // from this fiber instead. cancellation (not a cross-fiber // socket close) is the teardown: the reader observes // cancel_requested before ever touching the fd again, so // there is no closed-fd-number-reuse race, and the fd is // closed exactly once by client.deinit after the task ends. var read_future = self.io.concurrent(readLoopTask, .{ self, &client, &handler, heartbeat, &conn }) catch |e| { _ = self.bc.stats.subscriber_disconnect_read_loop.fetchAdd(1, .monotonic); return e; }; self.pingLoop(&client, &conn); return read_future.cancel(self.io); } }
// Threaded backend (and pinger-less standalone/test use): SO_RCVTIMEO // works, so the merged keepalive + read loop covers liveness alone. client.readLoopWithHeartbeat(&handler, heartbeat) catch |e| { _ = self.bc.stats.subscriber_disconnect_read_loop.fetchAdd(1, .monotonic); return e; }; }
fn readLoopTask(self: *Subscriber, client: *websocket.Client, handler: *FrameHandler, heartbeat: websocket.Client.HeartbeatConfig, conn: *ConnState) anyerror!void { defer conn.read_done.store(true, .release); client.readLoopWithHeartbeat(handler, heartbeat) catch |e| { _ = self.bc.stats.subscriber_disconnect_read_loop.fetchAdd(1, .monotonic); return e; }; }
/// zio keepalive: runs in the subscriber fiber while readLoopTask reads. /// pings when the connection has been silent past the interval, tears it /// down after max_failures unanswered pings. see ws_ping.zig for why /// silence alone never kills and how the reader-in-handler guard works. /// /// a ping write that parks on a full send buffer stalls only this /// connection's pinger, and the kernel's retransmission limit /// (tcp_retries2, ~15-25 min) errors the socket underneath it — the /// unanswered bytes themselves are the backstop. fn pingLoop(self: *Subscriber, client: *websocket.Client, conn: *ConnState) void { const cfg = self.ws_ping.?; var pending: u64 = 0; var empty: [0]u8 = .{};
// stagger by host so a cold start's simultaneous connects don't // produce a synchronized fleet-wide ping wave every interval var jitter: u64 = self.options.host_id % @max(1, cfg.interval_sec.load(.acquire)); while (jitter > 0 and !conn.read_done.load(.acquire) and !self.shouldStop()) : (jitter -= 1) { self.io.sleep(Io.Duration.fromSeconds(1), .awake) catch {}; }
var last_ping_ms: i64 = 0; outer: while (true) { // sleep the interval in 1s chunks so read-loop exit and shutdown // are noticed promptly (bounds reconnect latency to ~1s) var slept: u64 = 0; while (slept < @max(1, cfg.interval_sec.load(.acquire))) : (slept += 1) { if (conn.read_done.load(.acquire) or self.shouldStop()) break :outer; self.io.sleep(Io.Duration.fromSeconds(1), .awake) catch {}; } // a ping is answered iff ANY frame arrived after it was sent. // judging by silence-vs-interval alone instead is subtly wrong: a // tick's spacing is interval + scheduler oversleep while the pong // lands at interval + rtt, so a healthy quiet host with rtt below // the oversleep would accumulate "unanswered" pings it answered. if (pending > 0 and conn.activity_ms.load(.acquire) > last_ping_ms) pending = 0; switch (ws_ping_mod.decide( cfg, milliTimestamp(self.io), conn.activity_ms.load(.acquire), conn.in_handler.load(.acquire), pending, )) { .idle => pending = 0, .ping => { client.writePing(&empty) catch |e| { // error.Closed = the reader closed the connection // under us (websocket.zig v0.1.13 lock-scoped // teardown) — a normal disconnect, not a failure if (e != error.Closed) { _ = self.bc.stats.ws_ping_write_failures.fetchAdd(1, .monotonic); log.info("host {s}: keepalive ping write failed ({s}), closing", .{ self.options.hostname, @errorName(e) }); } break; }; last_ping_ms = milliTimestamp(self.io); _ = self.bc.stats.ws_pings_sent.fetchAdd(1, .monotonic); pending += 1; }, .kill => { _ = self.bc.stats.ws_ping_timeout_closes.fetchAdd(1, .monotonic); log.warn("host {s}: {d} keepalive pings unanswered, closing half-open connection", .{ self.options.hostname, pending }); break; }, } } }};
/// per-connection liveness state shared between the read-loop task and the/// pinger. all fields are atomics for the Threaded test path; under zio both/// sides are fibers on the one loop thread.const ConnState = struct { /// ms timestamp of the last frame received on THIS connection (any type — /// data, pong). distinct from Subscriber.last_progress_ms, which spans /// reconnects and stamps connect attempts for the admin reconnect fiber. activity_ms: std.atomic.Value(i64), /// reader is inside serverMessage (possibly blocked in rate-limit or /// pool backpressure for a long time, by design) — pongs cannot be read, /// so the pinger must not judge the connection. in_handler: std.atomic.Value(bool) = .{ .raw = false }, read_done: std.atomic.Value(bool) = .{ .raw = false },};
const FrameHandler = struct { subscriber: *Subscriber, conn: *ConnState,
/// any frame from the peer proves the link — pongs answer our keepalive /// pings, so this is what resets the pinger's failure count. pub fn serverPong(self: *FrameHandler, data: []u8) !void { _ = data; self.conn.activity_ms.store(milliTimestamp(self.subscriber.io), .release); }
pub fn serverMessage(self: *FrameHandler, data: []const u8) !void { const sub = self.subscriber; const io = sub.io;
self.conn.activity_ms.store(milliTimestamp(io), .release); self.conn.in_handler.store(true, .release); defer { // time blocked inside the handler must not count as silence self.conn.in_handler.store(false, .release); self.conn.activity_ms.store(milliTimestamp(io), .release); }
// lightweight header decode for cursor tracking + routing var arena = std.heap.ArenaAllocator.init(sub.allocator); defer arena.deinit(); const alloc = arena.allocator();
const header_result = zat.cbor.decode(alloc, data) catch |err| { log.debug("frame header decode failed: {s} (len={d})", .{ @errorName(err), data.len }); _ = sub.bc.stats.decode_errors.fetchAdd(1, .monotonic); return; }; const header = header_result.value; const payload_data = data[header_result.consumed..];
// check op field (1 = message, -1 = error) const op = header.getInt("op") orelse return; if (op == -1) { // error frame from upstream — check for FutureCursor // Go relay: slurper.go — sets host to idle and disconnects if (zat.cbor.decodeAll(alloc, payload_data) catch null) |err_payload| { const err_name = err_payload.getString("error") orelse "unknown"; const err_msg = err_payload.getString("message") orelse ""; log.warn("host {s}: error frame: {s}: {s}", .{ sub.options.hostname, err_name, err_msg }); if (std.mem.eql(u8, err_name, "FutureCursor")) { // our cursor is ahead of the PDS — set host to idle, stop this subscriber only if (sub.options.host_id > 0) { if (sub.host_ops) |hq| { hq.push(.{ .host_id = sub.options.host_id, .kind = .update_status, .payload = .{ .status = host_ops_mod.HostOp.Payload.Status.init("idle") }, }); } } sub.host_shutdown.store(true, .release); } } return; }
const frame_type = header.getString("t") orelse return; const payload = zat.cbor.decodeAll(alloc, payload_data) catch |err| { log.debug("frame payload decode failed: {s} (type={s})", .{ @errorName(err), frame_type }); _ = sub.bc.stats.decode_errors.fetchAdd(1, .monotonic); return; };
// count every successfully decoded event (matches Go relay's events_received_counter) _ = sub.bc.stats.frames_in.fetchAdd(1, .monotonic); const now_ms = milliTimestamp(io); sub.last_progress_ms.store(now_ms, .release);
// extract seq for cursor tracking (deferred until after pool accepts) const upstream_seq = payload.getUint("seq");
// time-based cursor flush (Go relay: every 4 seconds) { const now = @divFloor(now_ms, std.time.ms_per_s); if (now - sub.last_cursor_flush >= cursor_flush_interval_sec) { sub.flushCursor(); sub.last_cursor_flush = now; } }
// per-host rate limiting — block until window opens (matches indigo's waitForLimiter) // blocking here stalls the websocket reader → TCP backpressure → PDS slows down switch (sub.rate_limiter.waitForAllow(sub.shutdown, io)) { .allowed => {}, .sec => { _ = sub.bc.stats.rate_limited.fetchAdd(1, .monotonic); _ = sub.bc.stats.rate_limited_sec.fetchAdd(1, .monotonic); }, .hour => { _ = sub.bc.stats.rate_limited.fetchAdd(1, .monotonic); _ = sub.bc.stats.rate_limited_hour.fetchAdd(1, .monotonic); }, .day => { _ = sub.bc.stats.rate_limited.fetchAdd(1, .monotonic); _ = sub.bc.stats.rate_limited_day.fetchAdd(1, .monotonic); }, }
// filter unknown frame types before submitting to pool (forward-compat) const is_commit = std.mem.eql(u8, frame_type, "#commit"); const is_sync = std.mem.eql(u8, frame_type, "#sync"); const is_account = std.mem.eql(u8, frame_type, "#account"); const is_identity = std.mem.eql(u8, frame_type, "#identity");
if (!is_commit and !is_sync and !is_account and !is_identity) { // advance cursor for intentionally skipped frames — // we won't process these on reconnect either if (upstream_seq) |s| sub.last_upstream_seq = s; return; }
// submit to frame pool for heavy processing (CBOR re-decode, validation, persist, broadcast) // route by DID so commits for the same account serialize (prevents chain break races) // blocks when queue is full — TCP backpressure propagates to upstream PDS (matches indigo) if (sub.pool) |pool| { const did_key: u64 = blk: { const d = payload.getString("repo") orelse payload.getString("did"); break :blk if (d) |s| std.hash.Wyhash.hash(0, s) else sub.options.host_id; }; const duped = sub.allocator.dupe(u8, data) catch return; // dupe hostname per-frame: subscriber teardown (slurper.runWorker) // frees sub.options.hostname after sub.run() returns, but FrameWorks // can still be queued in the pool. borrowing the slice would be a // use-after-free (corrupt hostnames in chain-break logs, etc.). const hostname_dup = sub.allocator.dupe(u8, sub.options.hostname) catch { sub.allocator.free(duped); return; }; const t0 = nanoTimestamp(io); if (pool.submit(did_key, .{ .data = duped, .host_id = sub.options.host_id, .hostname = hostname_dup, .allocator = sub.allocator, .io = sub.pool_io orelse sub.io, // pool_io (Threaded) for worker-safe ops .bc = sub.bc, .validator = sub.validator, .persist = sub.persist, .collection_index = sub.collection_index, .resyncer = sub.resyncer, .status_checker = sub.status_checker, }, sub.shutdown)) { // pool accepted — advance cursor past this frame _ = sub.bc.stats.pool_queued_bytes.fetchAdd(duped.len, .monotonic); if (upstream_seq) |s| sub.last_upstream_seq = s; if (nanoTimestamp(io) - t0 > 1_000_000) { // >1ms = had to wait _ = sub.bc.stats.pool_backpressure.fetchAdd(1, .monotonic); } } else { // shutdown requested — don't advance cursor so reconnect replays this frame sub.allocator.free(duped); sub.allocator.free(hostname_dup); } return; }
// fallback: no pool, process inline (original path for tests / standalone use) self.processInline(sub, alloc, data, payload, upstream_seq, frame_type, is_commit, is_sync, is_account, is_identity); }
/// inline processing path — used when no frame pool is configured (tests, standalone). /// this is the original serverMessage heavy processing logic. fn processInline( _: *FrameHandler, sub: *Subscriber, alloc: Allocator, data: []const u8, payload: zat.cbor.Value, upstream_seq: ?u64, _: []const u8, is_commit: bool, is_sync: bool, is_account: bool, is_identity: bool, ) void { const io = sub.io;
// extract DID: "repo" for commits, "did" for identity/account const did: ?[]const u8 = if (is_commit) payload.getString("repo") else payload.getString("did");
// on #identity event, evict cached signing key so next commit re-resolves if (is_identity) { if (did) |d| sub.validator.evictKey(d); }
// resolve DID → numeric UID for event header (host-aware) const uid: u64 = if (sub.persist) |dp| blk: { if (did) |d| { const result = dp.uidForDidFromHost(d, sub.options.host_id) catch break :blk @as(u64, 0);
// host authority enforcement (mirrors frame_worker path) if ((result.is_new or result.host_changed) and !is_identity) { switch (sub.validator.resolveHostAuthority(d, sub.options.host_id, sub.options.hostname, .{ .skip_reject_cache = is_account, })) { .migrate => { dp.setAccountHostId(result.uid, sub.options.host_id) catch {}; if (result.host_changed) { log.info("host {s}: account migrated uid={d} did={s}", .{ sub.options.hostname, result.uid, d, }); } }, .reject => { log.debug("host {s}: dropping event, host authority failed uid={d} did={s}", .{ sub.options.hostname, result.uid, d, }); _ = sub.bc.stats.failed.fetchAdd(1, .monotonic); _ = sub.bc.stats.failed_host_authority.fetchAdd(1, .monotonic); return; }, .accept => {}, } }
break :blk result.uid; } else break :blk @as(u64, 0); } else 0;
// process #account events: update upstream status // Go relay: ingest.go processAccountEvent if (is_account) { if (sub.persist) |dp| { if (uid > 0) { const is_active = payload.getBool("active") orelse false; const status_str = payload.getString("status"); const new_status: []const u8 = if (is_active) "active" else (status_str orelse "inactive"); dp.updateAccountUpstreamStatus(uid, new_status) catch |err| { log.debug("upstream status update failed: {s}", .{@errorName(err)}); };
// on any non-active status, remove all collection index entries // (covers deleted, takendown, suspended, deactivated, and future unknown statuses) if (!is_active) { if (sub.collection_index) |ci| { if (did) |d| { ci.removeAll(d) catch |err| { log.debug("collection removeAll failed: {s}", .{@errorName(err)}); }; } } } } } }
// for commits and syncs: check account is active, validate, extract state var commit_data_cid: ?[]const u8 = null; var commit_rev: ?[]const u8 = null; if (is_commit or is_sync) { // drop for inactive accounts (see frame_worker for the re-check rationale) if (sub.persist) |dp| { if (uid > 0) { const activity = dp.accountActivity(uid) catch event_log_mod.DiskPersist.AccountActivity{ .local_ok = true, .upstream_ok = true }; if (!activity.active()) { if (activity.local_ok and !activity.upstream_ok) { if (sub.status_checker) |sc| { if (did) |d| sc.enqueue(uid, d, sub.options.hostname); } } _ = sub.bc.stats.skipped.fetchAdd(1, .monotonic); return; } } }
// rev checks: stale, future, chain continuity if (is_commit and uid > 0) { if (payload.getString("rev")) |incoming_rev| { // future-rev rejection if (zat.Tid.parse(incoming_rev)) |tid| { const rev_us: i64 = @intCast(tid.timestamp()); const now_us = microTimestamp(io); const skew_us: i64 = sub.validator.config.rev_clock_skew * 1_000_000; if (rev_us > now_us + skew_us) { log.info("host {s}: dropping future rev uid={d} rev={s}", .{ sub.options.hostname, uid, incoming_rev, }); _ = sub.bc.stats.failed.fetchAdd(1, .monotonic); _ = sub.bc.stats.failed_future_rev.fetchAdd(1, .monotonic); return; } }
if (sub.persist) |dp| { if (dp.getAccountState(uid, alloc) catch null) |prev| { // stale rev check if (std.mem.order(u8, incoming_rev, prev.rev) != .gt) { log.debug("host {s}: dropping stale commit uid={d} rev={s} <= {s}", .{ sub.options.hostname, uid, incoming_rev, prev.rev, }); _ = sub.bc.stats.skipped.fetchAdd(1, .monotonic); return; }
// chain continuity: since should match stored rev if (payload.getString("since")) |since| { if (!std.mem.eql(u8, since, prev.rev)) { log.info("host {s}: chain break uid={d} since={s} stored_rev={s}", .{ sub.options.hostname, uid, since, prev.rev, }); _ = sub.bc.stats.chain_breaks_since.fetchAdd(1, .monotonic); } }
// With prior repo state, prevData is required for an // inductive Sync 1.1 proof. Missing/null cannot be // repaired by MST inversion, so drop the commit. if (!sub.validator.validatePrevDataPresence(payload)) { log.info("host {s}: dropping commit with missing prevData uid={d}", .{ sub.options.hostname, uid, }); return; }
// chain continuity: prevData CID should match stored data_cid. // only meaningful with verify_commit_diff — that's the only path // that stores a real MST root in data_cid. if (sub.validator.config.verify_commit_diff) { if (payload.get("prevData")) |pd| { if (pd == .cid) { const prev_data_encoded = zat.multibase.encode(alloc, .base32lower, pd.cid.raw) catch ""; if (prev_data_encoded.len > 0 and prev.data_cid.len > 0 and !std.mem.eql(u8, prev_data_encoded, prev.data_cid)) { log.info("host {s}: chain break uid={d} prevData mismatch", .{ sub.options.hostname, uid, }); _ = sub.bc.stats.chain_breaks_prev_data.fetchAdd(1, .monotonic); } } } } } } } }
if (is_commit) { const result = sub.validator.validateCommit(payload); if (!result.valid) return; commit_data_cid = result.data_cid; commit_rev = result.commit_rev;
// track collections from commit ops (phase 1 live indexing) if (sub.collection_index) |ci| { if (did) |d| { if (payload.get("ops")) |ops| { ci.trackCommitOps(d, ops); } } } } else { // #sync: signature verification only, no ops/MST const result = sub.validator.validateSync(payload); if (!result.valid) return; commit_data_cid = result.data_cid; commit_rev = result.commit_rev;
// enqueue collection index resync — #sync means repo discontinuity if (sub.resyncer) |r| { if (did) |d| r.enqueue(d, sub.options.hostname); } } }
// determine event kind for persistence const kind: event_log_mod.EvtKind = if (is_commit) .commit else if (is_sync) .sync else if (is_account) .account else // is_identity (unknown types already filtered above) .identity;
// persist + resequence + enqueue under ordering lock to guarantee // broadcast_queue insertion order matches seq assignment order. if (sub.persist) |dp| { var spins: u64 = 0; while (sub.bc.persist_order.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { spins += 1; std.atomic.spinLoopHint(); } if (spins > 0) { _ = sub.bc.stats.persist_order_spins.fetchAdd(spins, .monotonic); }
const relay_seq = dp.persist(kind, uid, data) catch |err| { sub.bc.persist_order.store(0, .release); log.warn("persist failed: {s}", .{@errorName(err)}); return; }; sub.bc.stats.relay_seq.store(relay_seq, .release);
const broadcast_data = broadcaster.resequenceFrame(alloc, data, relay_seq) orelse data; const owned = sub.allocator.dupe(u8, broadcast_data) catch { sub.bc.persist_order.store(0, .release); return; }; sub.bc.broadcast_queue.push(relay_seq, owned, &sub.bc.stats); sub.bc.persist_order.store(0, .release);
// update per-DID state outside the ordering lock (Postgres round-trip) if ((is_commit or is_sync) and uid > 0) { if (commit_rev) |rev| { const cid_str: []const u8 = if (commit_data_cid) |cid_raw| zat.multibase.encode(alloc, .base32lower, cid_raw) catch "" else ""; if (dp.updateAccountState(uid, rev, cid_str)) |updated| { if (!updated) { log.debug("host {s}: stale state update lost race uid={d} rev={s}", .{ sub.options.hostname, uid, rev, }); } } else |err| { log.debug("account state update failed: {s}", .{@errorName(err)}); } } } } else { const seq = upstream_seq orelse 0; const owned = sub.allocator.dupe(u8, data) catch return; sub.bc.broadcast_queue.push(seq, owned, &sub.bc.stats); } }
pub fn close(self: *FrameHandler) void { log.info("host {s}: connection closed", .{self.subscriber.options.hostname}); // flush cursor on disconnect self.subscriber.flushCursor(); }};
// --- tests ---
test "decode frame via SDK and extract fields" { const cbor = zat.cbor;
var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
// build a commit frame using SDK encoder const header: cbor.Value = .{ .map = &.{ .{ .key = "op", .value = .{ .unsigned = 1 } }, .{ .key = "t", .value = .{ .text = "#commit" } }, } }; const payload: cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = "did:plc:test123" } }, .{ .key = "seq", .value = .{ .unsigned = 12345 } }, .{ .key = "rev", .value = .{ .text = "3k2abc000000" } }, .{ .key = "time", .value = .{ .text = "2024-01-15T10:30:00Z" } }, } };
const header_bytes = try cbor.encodeAlloc(alloc, header); const payload_bytes = try cbor.encodeAlloc(alloc, payload);
var frame = try alloc.alloc(u8, header_bytes.len + payload_bytes.len); @memcpy(frame[0..header_bytes.len], header_bytes); @memcpy(frame[header_bytes.len..], payload_bytes);
// decode using SDK (same path as FrameHandler.serverMessage) const h_result = try cbor.decode(alloc, frame); const h = h_result.value; const p_data = frame[h_result.consumed..]; const p = try cbor.decodeAll(alloc, p_data);
try std.testing.expectEqualStrings("#commit", h.getString("t").?); try std.testing.expectEqual(@as(i64, 1), h.getInt("op").?); try std.testing.expectEqual(@as(i64, 12345), p.getInt("seq").?); try std.testing.expectEqualStrings("did:plc:test123", p.getString("repo").?);}
test "decode identity frame via SDK" { const cbor = zat.cbor;
var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
const header: cbor.Value = .{ .map = &.{ .{ .key = "op", .value = .{ .unsigned = 1 } }, .{ .key = "t", .value = .{ .text = "#identity" } }, } }; const payload: cbor.Value = .{ .map = &.{ .{ .key = "did", .value = .{ .text = "did:plc:alice" } }, .{ .key = "seq", .value = .{ .unsigned = 99 } }, .{ .key = "time", .value = .{ .text = "2024-01-15T10:30:00Z" } }, } };
const header_bytes = try cbor.encodeAlloc(alloc, header); const payload_bytes = try cbor.encodeAlloc(alloc, payload);
var frame = try alloc.alloc(u8, header_bytes.len + payload_bytes.len); @memcpy(frame[0..header_bytes.len], header_bytes); @memcpy(frame[header_bytes.len..], payload_bytes);
const h_result = try cbor.decode(alloc, frame); const h = h_result.value; const p_data = frame[h_result.consumed..]; const p = try cbor.decodeAll(alloc, p_data);
try std.testing.expectEqualStrings("#identity", h.getString("t").?); try std.testing.expectEqualStrings("did:plc:alice", p.getString("did").?); try std.testing.expectEqual(@as(i64, 99), p.getInt("seq").?);}
test "rate limiter enforces per-second limit" { var rl: RateLimiter = .{ .sec_limit = .{ .raw = 3 }, .hour_limit = .{ .raw = 1000 }, .day_limit = .{ .raw = 10000 } };
const now: i64 = 1_000_000_000; try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(now)); try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(now)); try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(now)); try std.testing.expectEqual(RateLimiter.Result.sec, rl.allow(now));
// at start of next second, prev carries over fully → still blocked try std.testing.expectEqual(RateLimiter.Result.sec, rl.allow(now + 1000));
// two seconds later prev is forgotten try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(now + 2000));}
test "rate limiter enforces per-hour limit" { var rl: RateLimiter = .{ .sec_limit = .{ .raw = 1000 }, .hour_limit = .{ .raw = 5 }, .day_limit = .{ .raw = 10000 } };
const now: i64 = 3_600_000 * 100; for (0..5) |_| { try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(now)); } try std.testing.expectEqual(RateLimiter.Result.hour, rl.allow(now + 100));
// next hour, prev carries over → still blocked try std.testing.expectEqual(RateLimiter.Result.hour, rl.allow(now + 3_600_000));
// two hours later prev is forgotten try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(now + 7_200_000));}
test "sliding window interpolates previous count by elapsed time" { var rl: RateLimiter = .{ .sec_limit = .{ .raw = 10 }, .hour_limit = .{ .raw = 1_000_000 }, .day_limit = .{ .raw = 1_000_000 } };
const base: i64 = 1_000_000_000; for (0..8) |_| { try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(base)); }
// next second: prev=8, curr=0. effective=8, room for 2 try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(base + 1000)); try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(base + 1000)); try std.testing.expectEqual(RateLimiter.Result.sec, rl.allow(base + 1000));
// halfway through: prev weight = 500/1000 = 0.5, so weighted_prev = 4. curr=2. eff=6, room for 4 try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(base + 1500)); try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(base + 1500)); try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(base + 1500)); try std.testing.expectEqual(RateLimiter.Result.allowed, rl.allow(base + 1500)); try std.testing.expectEqual(RateLimiter.Result.sec, rl.allow(base + 1500));}
test "waitForAllow blocks then allows after window advances" { const io = std.testing.io;
// verify that waitForAllow returns a non-.allowed result when the limit was hit, // indicating it had to wait. We use a tiny limit so the fast path is exhausted. var rl: RateLimiter = .{ .sec_limit = .{ .raw = 1 }, .hour_limit = .{ .raw = 1000 }, .day_limit = .{ .raw = 10000 } }; var shutdown = std.atomic.Value(bool){ .raw = false };
// first call takes the fast path try std.testing.expectEqual(RateLimiter.Result.allowed, rl.waitForAllow(&shutdown, io));
// second call must block (sec limit = 1), then return .sec after the window advances // this will sleep ~100ms+ until the sliding window allows it const before = milliTimestamp(io); const result = rl.waitForAllow(&shutdown, io); const elapsed = milliTimestamp(io) - before;
try std.testing.expectEqual(RateLimiter.Result.sec, result); try std.testing.expect(elapsed >= 100); // must have slept at least one 100ms poll}
test "waitForAllow respects shutdown" { const io = std.testing.io;
var rl: RateLimiter = .{ .sec_limit = .{ .raw = 1 }, .hour_limit = .{ .raw = 1000 }, .day_limit = .{ .raw = 10000 } }; var shutdown = std.atomic.Value(bool){ .raw = false };
// exhaust the limit _ = rl.waitForAllow(&shutdown, io);
// set shutdown before the next call shutdown.store(true, .release);
// should return immediately without blocking const before = milliTimestamp(io); _ = rl.waitForAllow(&shutdown, io); const elapsed = milliTimestamp(io) - before;
try std.testing.expect(elapsed < 50); // should not have slept}
test "trusted host detection" { try std.testing.expect(isTrustedHost("pds-123.host.bsky.network")); try std.testing.expect(isTrustedHost("abc.host.bsky.network")); try std.testing.expect(!isTrustedHost("bsky.network")); try std.testing.expect(!isTrustedHost("evil.bsky.network")); try std.testing.expect(!isTrustedHost("pds.example.com"));}
test "error frame (op=-1) is detected" { const cbor = zat.cbor;
var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
const header: cbor.Value = .{ .map = &.{ .{ .key = "op", .value = .{ .negative = -1 } }, .{ .key = "t", .value = .{ .text = "#info" } }, } };
const header_bytes = try cbor.encodeAlloc(alloc, header); const h_result = try cbor.decode(alloc, header_bytes); const h = h_result.value;
try std.testing.expectEqual(@as(i64, -1), h.getInt("op").?);}
// --- spec conformance tests ---
test "spec: unknown frame type (op=1, t=#unknown) is ignored" { // event stream spec: unknown t values must be ignored for forward-compat const cbor = zat.cbor;
var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
const header: cbor.Value = .{ .map = &.{ .{ .key = "op", .value = .{ .unsigned = 1 } }, .{ .key = "t", .value = .{ .text = "#unknown" } }, } }; const payload: cbor.Value = .{ .map = &.{ .{ .key = "did", .value = .{ .text = "did:plc:test123" } }, .{ .key = "seq", .value = .{ .unsigned = 1 } }, } };
const header_bytes = try cbor.encodeAlloc(alloc, header); const payload_bytes = try cbor.encodeAlloc(alloc, payload);
// decode header — verify it's a valid message with unknown type const h_result = try cbor.decode(alloc, header_bytes); const h = h_result.value;
try std.testing.expectEqual(@as(i64, 1), h.getInt("op").?); const frame_type = h.getString("t").?; try std.testing.expectEqualStrings("#unknown", frame_type);
// verify unknown type is NOT one of the known types (this is the filter logic) const is_commit = std.mem.eql(u8, frame_type, "#commit"); const is_sync = std.mem.eql(u8, frame_type, "#sync"); const is_account = std.mem.eql(u8, frame_type, "#account"); const is_identity = std.mem.eql(u8, frame_type, "#identity"); try std.testing.expect(!is_commit and !is_sync and !is_account and !is_identity);
// verify payload still decodes (frame is valid, just ignored) const p = try cbor.decodeAll(alloc, payload_bytes); try std.testing.expectEqualStrings("did:plc:test123", p.getString("did").?);}
test "spec: error frame (op=-1) is handled, not persisted" { // event stream spec: op=-1 frames are error notifications from upstream const cbor = zat.cbor;
var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator();
const header: cbor.Value = .{ .map = &.{ .{ .key = "op", .value = .{ .negative = -1 } }, .{ .key = "t", .value = .{ .text = "#error" } }, } }; const err_payload: cbor.Value = .{ .map = &.{ .{ .key = "error", .value = .{ .text = "FutureCursor" } }, .{ .key = "message", .value = .{ .text = "cursor is ahead of server" } }, } };
const header_bytes = try cbor.encodeAlloc(alloc, header); const h_result = try cbor.decode(alloc, header_bytes); const h = h_result.value;
// verify op=-1 is detected const op = h.getInt("op").?; try std.testing.expectEqual(@as(i64, -1), op);
// verify error payload decodes correctly const payload_bytes = try cbor.encodeAlloc(alloc, err_payload); const p = try cbor.decodeAll(alloc, payload_bytes); try std.testing.expectEqualStrings("FutureCursor", p.getString("error").?); try std.testing.expectEqualStrings("cursor is ahead of server", p.getString("message").?);}