Something went wrong. Try again.
atproto relay implementation in zig zlay.waow.tech
Something went wrong. Try again.
Zig
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015//! relay frame validator — DID key resolution + real signature verification//!//! validates firehose commit frames by verifying the commit signature against//! the pre-resolved signing key for the DID. accepts pre-decoded CBOR payload//! from the subscriber (decoded via zat SDK). on cache miss, skips validation//! and queues background resolution. no frame is ever blocked on network I/O.
const std = @import("std");const Io = std.Io;const zat = @import("zat");const broadcaster = @import("broadcaster.zig");const event_log_mod = @import("event_log.zig");const lru = @import("lru.zig");
const Allocator = std.mem.Allocator;const log = std.log.scoped(.relay);
/// decoded and cached signing key for a DIDconst CachedKey = struct { key_type: zat.multicodec.KeyType, raw: [33]u8, // compressed public key (secp256k1 or p256) len: u8, resolve_time: i64 = 0, // epoch seconds when resolved};
pub const ValidationResult = struct { valid: bool, skipped: bool, data_cid: ?[]const u8 = null, // MST root CID from verified commit commit_rev: ?[]const u8 = null, // rev from verified commit};
/// configuration for commit validation checkspub const ValidatorConfig = struct { /// verify MST structure during signature verification verify_mst: bool = false, // off by default for relay throughput /// verify commit diffs via MST inversion (sync 1.1) verify_commit_diff: bool = false, /// max allowed operations per commit max_ops: usize = 200, /// max clock skew for rev timestamps (seconds) rev_clock_skew: i64 = 300, // 5 minutes};
pub const Validator = struct { allocator: Allocator, stats: *broadcaster.Stats, config: ValidatorConfig, persist: ?*event_log_mod.DiskPersist = null, // DID → signing key cache (decoded, ready for verification) cache: lru.LruCache(CachedKey), // background resolve queue queue: std.ArrayListUnmanaged([]const u8) = .empty, // in-flight set — prevents duplicate DID entries in the queue queued_set: std.StringHashMapUnmanaged(void) = .empty, queue_mutex: Io.Mutex = Io.Mutex.init, queue_cond: Io.Condition = Io.Condition.init, resolver_futures: [max_resolver_threads]?Io.Future(void) = .{null} ** max_resolver_threads, alive: std.atomic.Value(bool) = .{ .raw = true }, max_cache_size: u32 = 250_000, io: Io, // pool of reusable resolvers for inline host authority checks. // frame workers acquire/release via atomic flag to avoid creating // a fresh resolver (and fresh TLS handshake) per call. host_resolvers: [host_resolver_pool_size]zat.DidResolver = undefined, host_resolver_available: [host_resolver_pool_size]std.atomic.Value(bool) = .{std.atomic.Value(bool){ .raw = false }} ** host_resolver_pool_size, host_resolver_inited: bool = false,
const max_resolver_threads = 8; const default_resolver_threads = 4; const max_queue_size: usize = 100_000; const host_resolver_pool_size: usize = 4;
pub fn init(allocator: Allocator, stats: *broadcaster.Stats, io: Io) Validator { return initWithConfig(allocator, stats, .{}, io); }
pub fn initWithConfig(allocator: Allocator, stats: *broadcaster.Stats, config: ValidatorConfig, io: Io) Validator { return .{ .allocator = allocator, .stats = stats, .config = config, .cache = lru.LruCache(CachedKey).init(allocator, 250_000, io), .io = io, }; }
pub fn deinit(self: *Validator) void { self.alive.store(false, .release); self.queue_cond.broadcast(self.io); for (&self.resolver_futures) |*slot| { if (slot.*) |*f| { f.cancel(self.io); } slot.* = null; }
if (self.host_resolver_inited) { for (&self.host_resolvers) |*r| { r.deinit(); } self.host_resolver_inited = false; }
self.cache.deinit();
// free queued DIDs for (self.queue.items) |did| { self.allocator.free(did); } self.queue.deinit(self.allocator); self.queued_set.deinit(self.allocator); }
/// start background resolver threads and host authority resolver pool pub fn start(self: *Validator) !void { self.max_cache_size = parseEnvInt(u32, "VALIDATOR_CACHE_SIZE", self.max_cache_size); self.cache.capacity = self.max_cache_size; const n = parseEnvInt(u8, "RESOLVER_THREADS", default_resolver_threads); const count = @min(n, max_resolver_threads); for (self.resolver_futures[0..count]) |*slot| { slot.* = try self.io.concurrent(resolveLoop, .{self}); }
// init host authority resolver pool (reused across calls) for (&self.host_resolvers) |*r| { r.* = zat.DidResolver.initWithOptions(self.io, self.allocator, .{}); } for (&self.host_resolver_available) |*a| { a.store(true, .release); } self.host_resolver_inited = true; }
/// validate a #sync frame: signature verification only (no ops, no MST). /// #sync resets a repo to a new commit state — used for recovery from broken streams. /// on cache miss, queues background resolution and skips. pub fn validateSync(self: *Validator, payload: zat.cbor.Value) ValidationResult { const did = payload.getString("did") orelse { _ = self.stats.skipped.fetchAdd(1, .monotonic); return .{ .valid = true, .skipped = true }; };
if (zat.Did.parse(did) == null) { _ = self.stats.failed.fetchAdd(1, .monotonic); _ = self.stats.failed_bad_did.fetchAdd(1, .monotonic); return .{ .valid = false, .skipped = false }; }
// check rev is valid TID (if present) if (payload.getString("rev")) |rev| { if (zat.Tid.parse(rev) == null) { _ = self.stats.failed.fetchAdd(1, .monotonic); _ = self.stats.failed_bad_rev.fetchAdd(1, .monotonic); return .{ .valid = false, .skipped = false }; } }
const blocks = payload.getBytes("blocks") orelse { _ = self.stats.failed.fetchAdd(1, .monotonic); _ = self.stats.failed_missing_blocks.fetchAdd(1, .monotonic); return .{ .valid = false, .skipped = false }; };
// #sync CAR should be small (just the signed commit block) // lexicon maxLength: 10000 if (blocks.len > 10_000) { _ = self.stats.failed.fetchAdd(1, .monotonic); _ = self.stats.failed_oversized_blocks.fetchAdd(1, .monotonic); return .{ .valid = false, .skipped = false }; }
// cache lookup const cached_key: ?CachedKey = self.cache.get(did);
if (cached_key == null) { _ = self.stats.cache_misses.fetchAdd(1, .monotonic); _ = self.stats.skipped.fetchAdd(1, .monotonic); self.queueResolve(did); return .{ .valid = true, .skipped = true }; }
_ = self.stats.cache_hits.fetchAdd(1, .monotonic);
// verify signature (no MST, no ops) const public_key = zat.multicodec.PublicKey{ .key_type = cached_key.?.key_type, .raw = cached_key.?.raw[0..cached_key.?.len], };
var arena = std.heap.ArenaAllocator.init(self.allocator); defer arena.deinit();
const result = zat.verifyCommitCar(arena.allocator(), blocks, public_key, .{ .verify_mst = false, .expected_did = did, .max_car_size = 10 * 1024, }) catch |err| { log.debug("sync verification failed for {s}: {s}", .{ did, @errorName(err) }); // sync spec: on signature failure, key may have rotated. // evict cached key and queue re-resolution. skip this frame. self.evictKey(did); self.queueResolve(did); _ = self.stats.skipped.fetchAdd(1, .monotonic); return .{ .valid = true, .skipped = true }; };
_ = self.stats.validated.fetchAdd(1, .monotonic); return .{ .valid = true, .skipped = false, .data_cid = result.commit_cid, .commit_rev = result.commit_rev, }; }
/// validate a commit frame using a pre-decoded CBOR payload (from SDK decoder). /// on cache miss, queues background resolution and skips. pub fn validateCommit(self: *Validator, payload: zat.cbor.Value) ValidationResult { // extract DID from decoded payload const did = payload.getString("repo") orelse { _ = self.stats.skipped.fetchAdd(1, .monotonic); return .{ .valid = true, .skipped = true }; };
// check cache for pre-resolved signing key const cached_key: ?CachedKey = self.cache.get(did);
if (cached_key == null) { // cache miss — queue for background resolution, skip validation _ = self.stats.cache_misses.fetchAdd(1, .monotonic); _ = self.stats.skipped.fetchAdd(1, .monotonic); self.queueResolve(did); return .{ .valid = true, .skipped = true }; }
_ = self.stats.cache_hits.fetchAdd(1, .monotonic);
// cache hit — do structure checks + signature verification if (self.verifyCommit(payload, did, cached_key.?)) |vr| { _ = self.stats.validated.fetchAdd(1, .monotonic); return vr; } else |err| { log.debug("commit verification failed for {s}: {s}", .{ did, @errorName(err) }); // sync spec: on signature failure, key may have rotated. // evict cached key and queue re-resolution. skip this frame // (treat as cache miss). next commit will use the refreshed key. self.evictKey(did); self.queueResolve(did); _ = self.stats.skipped.fetchAdd(1, .monotonic); return .{ .valid = true, .skipped = true }; } }
fn verifyCommit(self: *Validator, payload: zat.cbor.Value, expected_did: []const u8, cached_key: CachedKey) !ValidationResult { // commit structure checks first (cheap, no allocation) self.checkCommitStructure(payload) catch { _ = self.stats.failed_bad_structure.fetchAdd(1, .monotonic); return error.InvalidFrame; };
// extract blocks (raw CAR bytes) from the pre-decoded payload const blocks = payload.getBytes("blocks") orelse return error.InvalidFrame;
// blocks size check — lexicon maxLength: 2000000 if (blocks.len > 2_000_000) return error.InvalidFrame;
// build public key for verification const public_key = zat.multicodec.PublicKey{ .key_type = cached_key.key_type, .raw = cached_key.raw[0..cached_key.len], };
// run real signature verification (needs its own arena for CAR/MST temporaries) var arena = std.heap.ArenaAllocator.init(self.allocator); defer arena.deinit(); const alloc = arena.allocator();
// try sync 1.1 path: extract ops and use verifyCommitDiff if (self.config.verify_commit_diff) { if (self.extractOps(alloc, payload)) |msg_ops| { // get stored prev_data from payload const prev_data: ?[]const u8 = if (payload.get("prevData")) |pd| switch (pd) { .cid => |c| c.raw, .null => null, else => null, } else null;
const diff_result = zat.verifyCommitDiff(alloc, blocks, msg_ops, prev_data, public_key, .{ .expected_did = expected_did, .skip_inversion = prev_data == null, }) catch |err| { return err; };
return .{ .valid = true, .skipped = false, .data_cid = diff_result.data_cid, .commit_rev = diff_result.commit_rev, }; } }
// fallback: legacy verification (signature + optional MST walk) const result = zat.verifyCommitCar(alloc, blocks, public_key, .{ .verify_mst = self.config.verify_mst, .expected_did = expected_did, }) catch |err| { return err; };
return .{ .valid = true, .skipped = false, .data_cid = result.commit_cid, .commit_rev = result.commit_rev, }; }
/// extract ops from payload and convert to mst.Operation array. /// the firehose format uses a single "path" field ("collection/rkey"), /// not separate "collection"/"rkey" fields. fn extractOps(self: *Validator, alloc: Allocator, payload: zat.cbor.Value) ?[]const zat.MstOperation { _ = self; const ops_array = payload.getArray("ops") orelse return null; var ops: std.ArrayListUnmanaged(zat.MstOperation) = .empty; for (ops_array) |op| { const action = op.getString("action") orelse continue; const path = op.getString("path") orelse continue;
// validate path contains "/" (collection/rkey) if (std.mem.indexOfScalar(u8, path, '/') == null) continue;
// extract CID values const cid_value: ?[]const u8 = if (op.get("cid")) |v| switch (v) { .cid => |c| c.raw, else => null, } else null;
var value: ?[]const u8 = null; var prev: ?[]const u8 = null;
if (std.mem.eql(u8, action, "create")) { value = cid_value; } else if (std.mem.eql(u8, action, "update")) { value = cid_value; prev = if (op.get("prev")) |v| switch (v) { .cid => |c| c.raw, else => null, } else null; } else if (std.mem.eql(u8, action, "delete")) { prev = if (op.get("prev")) |v| switch (v) { .cid => |c| c.raw, else => null, } else null; } else continue;
ops.append(alloc, .{ .path = path, .value = value, .prev = prev, }) catch return null; }
if (ops.items.len == 0) return null; return ops.items; }
fn checkCommitStructure(self: *Validator, payload: zat.cbor.Value) !void { // check repo field is a valid DID const repo = payload.getString("repo") orelse return error.InvalidFrame; if (zat.Did.parse(repo) == null) return error.InvalidFrame;
// check rev is a valid TID if (payload.getString("rev")) |rev| { if (zat.Tid.parse(rev) == null) return error.InvalidFrame; }
// check ops count if (payload.get("ops")) |ops_value| { switch (ops_value) { .array => |ops| { if (ops.len > self.config.max_ops) return error.InvalidFrame; // validate each op has valid path (collection/rkey) for (ops) |op| { if (op.getString("path")) |path| { if (std.mem.indexOfScalar(u8, path, '/')) |sep| { const collection = path[0..sep]; const rkey = path[sep + 1 ..]; if (zat.Nsid.parse(collection) == null) return error.InvalidFrame; if (rkey.len > 0) { if (zat.Rkey.parse(rkey) == null) return error.InvalidFrame; } } else return error.InvalidFrame; // path must contain '/' } } }, else => return error.InvalidFrame, } } }
fn queueResolve(self: *Validator, did: []const u8) void { // check if already cached (race between validate and resolver) if (self.cache.contains(did)) return;
const duped = self.allocator.dupe(u8, did) catch return;
self.queue_mutex.lockUncancelable(self.io); defer self.queue_mutex.unlock(self.io);
// skip if already queued (prevents unbounded queue growth) if (self.queued_set.contains(duped)) { self.allocator.free(duped); return; }
// cap queue size — drop DID without adding to queued_set so it can be re-queued later if (self.queue.items.len >= max_queue_size) { self.allocator.free(duped); return; }
self.queue.append(self.allocator, duped) catch { self.allocator.free(duped); return; }; self.queued_set.put(self.allocator, duped, {}) catch {}; self.queue_cond.signal(self.io); }
fn resolveLoop(self: *Validator) void { var resolver = zat.DidResolver.initWithOptions(self.io, self.allocator, .{ .keep_alive = true }); defer resolver.deinit();
while (self.alive.load(.acquire)) { var did: ?[]const u8 = null; { self.queue_mutex.lockUncancelable(self.io); defer self.queue_mutex.unlock(self.io); while (self.queue.items.len == 0 and self.alive.load(.acquire)) { self.queue_cond.waitUncancelable(self.io, &self.queue_mutex); } if (self.queue.items.len > 0) { did = self.queue.orderedRemove(0); _ = self.queued_set.remove(did.?); } }
const d = did orelse continue; defer self.allocator.free(d);
// skip if already cached (resolved while queued) if (self.cache.contains(d)) continue;
// resolve DID → signing key const parsed = zat.Did.parse(d) orelse continue; var doc = resolver.resolve(parsed) catch |err| { log.debug("DID resolve failed for {s}: {s}", .{ d, @errorName(err) }); continue; }; defer doc.deinit();
// extract and decode signing key const vm = doc.signingKey() orelse continue; const key_bytes = zat.multibase.decode(self.allocator, vm.public_key_multibase) catch continue; defer self.allocator.free(key_bytes); const public_key = zat.multicodec.parsePublicKey(key_bytes) catch continue;
// store decoded key in cache (fixed-size, no pointer chasing) var cached = CachedKey{ .key_type = public_key.key_type, .raw = undefined, .len = @intCast(public_key.raw.len), .resolve_time = timestamp(self.io), }; @memcpy(cached.raw[0..public_key.raw.len], public_key.raw);
self.cache.put(d, cached) catch continue;
// --- host validation (merged from migration queue) --- // while we have the DID doc, check PDS endpoint and update host if needed. // best-effort: failures don't prevent signing key caching. if (self.persist) |persist| { if (doc.pdsEndpoint()) |pds_endpoint| { if (extractHostFromUrl(pds_endpoint)) |pds_host| { const pds_host_id = (persist.getHostIdForHostname(pds_host) catch null) orelse continue; const uid = persist.uidForDid(d) catch continue; const current_host = persist.getAccountHostId(uid) catch continue; if (current_host != 0 and current_host != pds_host_id) { persist.setAccountHostId(uid, pds_host_id) catch {}; log.info("host updated via DID doc: {s} -> host {d}", .{ d, pds_host_id }); } } } } } }
/// evict a DID's cached signing key (e.g. on #identity event). /// the next commit from this DID will trigger a fresh resolution. pub fn evictKey(self: *Validator, did: []const u8) void { _ = self.cache.remove(did); }
/// cache size (for diagnostics) pub fn cacheSize(self: *Validator) u32 { return self.cache.count(); }
/// resolve queue length (for diagnostics — non-blocking) pub fn resolveQueueLen(self: *Validator) usize { if (!self.queue_mutex.tryLock()) return 0; defer self.queue_mutex.unlock(self.io); return self.queue.items.len; }
/// resolve dedup set size (for diagnostics — non-blocking) pub fn resolveQueuedSetCount(self: *Validator) u32 { if (!self.queue_mutex.tryLock()) return 0; defer self.queue_mutex.unlock(self.io); return self.queued_set.count(); }
/// signing key cache hashmap backing capacity (for memory attribution) pub fn cacheMapCapacity(self: *Validator) u32 { return self.cache.mapCapacity(); }
/// resolver dedup set hashmap backing capacity (for memory attribution — non-blocking) pub fn resolveQueuedSetCapacity(self: *Validator) u32 { if (!self.queue_mutex.tryLock()) return 0; defer self.queue_mutex.unlock(self.io); return self.queued_set.capacity(); }
pub const HostAuthority = enum { accept, migrate, reject };
/// synchronous host authority check. called on first-seen DIDs (is_new) /// and host migrations (host_changed). resolves the DID doc to verify the /// PDS endpoint matches the incoming host. retries once on failure to /// handle transient network errors. /// /// uses a pooled resolver to avoid creating a fresh resolver (and fresh /// TLS handshake) per call. blocks briefly if all pool slots are in use. /// /// returns: /// .accept — should not happen (caller should only call on new/mismatch) /// .migrate — DID doc confirms this host, caller should update host_id /// .reject — DID doc does not confirm, caller should drop the event pub fn resolveHostAuthority(self: *Validator, did: []const u8, incoming_host_id: u64) HostAuthority { const persist = self.persist orelse return .migrate; // no DB — can't check const parsed = zat.Did.parse(did) orelse return .reject;
const idx = self.acquireHostResolver(); defer self.releaseHostResolver(idx);
var resolver = &self.host_resolvers[idx];
// first resolve attempt var doc = resolver.resolve(parsed) catch { // retry once on network failure var doc2 = resolver.resolve(parsed) catch return .reject; defer doc2.deinit(); return self.checkPdsHost(&doc2, persist, incoming_host_id); }; defer doc.deinit(); return self.checkPdsHost(&doc, persist, incoming_host_id); }
/// acquire a resolver from the pool. spins until one is available. fn acquireHostResolver(self: *Validator) usize { while (self.alive.load(.acquire)) { for (0..host_resolver_pool_size) |i| { if (self.host_resolver_available[i].cmpxchgStrong(true, false, .acquire, .monotonic) == null) { return i; } } self.io.sleep(Io.Duration.fromMilliseconds(1), .awake) catch {}; } return 0; // shutdown path — caller will exit soon }
fn releaseHostResolver(self: *Validator, idx: usize) void { self.host_resolver_available[idx].store(true, .release); }
fn checkPdsHost(self: *Validator, doc: *zat.DidDocument, persist: *event_log_mod.DiskPersist, incoming_host_id: u64) HostAuthority { _ = self; const pds_endpoint = doc.pdsEndpoint() orelse return .reject; const pds_host = extractHostFromUrl(pds_endpoint) orelse return .reject; const pds_host_id = (persist.getHostIdForHostname(pds_host) catch null) orelse return .reject; if (pds_host_id == incoming_host_id) return .migrate; return .reject; }};
/// extract hostname from a URL like "https://pds.example.com" or "https://pds.example.com:443/path"pub fn extractHostFromUrl(url: []const u8) ?[]const u8 { // strip scheme var rest = url; if (std.mem.startsWith(u8, rest, "https://")) { rest = rest["https://".len..]; } else if (std.mem.startsWith(u8, rest, "http://")) { rest = rest["http://".len..]; } // strip path if (std.mem.indexOfScalar(u8, rest, '/')) |i| { rest = rest[0..i]; } // strip port if (std.mem.indexOfScalar(u8, rest, ':')) |i| { rest = rest[0..i]; } if (rest.len == 0) return null; return rest;}
fn getenv(key: [*:0]const u8) ?[]const u8 { const ptr = std.c.getenv(key) orelse return null; return std.mem.sliceTo(ptr, 0);}
fn parseEnvInt(comptime T: type, key: [*:0]const u8, default: T) T { const val = getenv(key) orelse return default; return std.fmt.parseInt(T, val, 10) catch default;}
fn timestamp(io: Io) i64 { return @intCast(@divFloor(Io.Timestamp.now(io, .real).nanoseconds, std.time.ns_per_s));}
// --- tests ---
test "validateCommit skips on cache miss" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
// build a commit payload using SDK const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = "did:plc:test123" } }, .{ .key = "seq", .value = .{ .unsigned = 42 } }, .{ .key = "rev", .value = .{ .text = "3k2abc000000" } }, .{ .key = "time", .value = .{ .text = "2024-01-15T10:30:00Z" } }, } };
const result = v.validateCommit(payload); try std.testing.expect(result.valid); try std.testing.expect(result.skipped); try std.testing.expectEqual(@as(u64, 1), stats.cache_misses.load(.acquire));}
test "validateCommit skips when no repo field" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
// payload without "repo" field const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "seq", .value = .{ .unsigned = 42 } }, } };
const result = v.validateCommit(payload); try std.testing.expect(result.valid); try std.testing.expect(result.skipped); try std.testing.expectEqual(@as(u64, 1), stats.skipped.load(.acquire));}
test "checkCommitStructure rejects invalid DID" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = "not-a-did" } }, } };
try std.testing.expectError(error.InvalidFrame, v.checkCommitStructure(payload));}
test "checkCommitStructure accepts valid commit" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = "did:plc:test123" } }, .{ .key = "rev", .value = .{ .text = "3k2abcdefghij" } }, } };
try v.checkCommitStructure(payload);}
test "validateSync skips on cache miss" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "did", .value = .{ .text = "did:plc:test123" } }, .{ .key = "seq", .value = .{ .unsigned = 42 } }, .{ .key = "rev", .value = .{ .text = "3k2abcdefghij" } }, .{ .key = "blocks", .value = .{ .bytes = "deadbeef" } }, } };
const result = v.validateSync(payload); try std.testing.expect(result.valid); try std.testing.expect(result.skipped); try std.testing.expectEqual(@as(u64, 1), stats.cache_misses.load(.acquire));}
test "validateSync rejects invalid DID" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "did", .value = .{ .text = "not-a-did" } }, .{ .key = "blocks", .value = .{ .bytes = "deadbeef" } }, } };
const result = v.validateSync(payload); try std.testing.expect(!result.valid); try std.testing.expect(!result.skipped); try std.testing.expectEqual(@as(u64, 1), stats.failed.load(.acquire));}
test "validateSync rejects missing blocks" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "did", .value = .{ .text = "did:plc:test123" } }, .{ .key = "rev", .value = .{ .text = "3k2abcdefghij" } }, } };
const result = v.validateSync(payload); try std.testing.expect(!result.valid); try std.testing.expect(!result.skipped); try std.testing.expectEqual(@as(u64, 1), stats.failed.load(.acquire));}
test "validateSync skips when no did field" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "seq", .value = .{ .unsigned = 42 } }, } };
const result = v.validateSync(payload); try std.testing.expect(result.valid); try std.testing.expect(result.skipped); try std.testing.expectEqual(@as(u64, 1), stats.skipped.load(.acquire));}
test "LRU cache evicts least recently used" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); v.cache.capacity = 3; defer v.deinit();
const mk = CachedKey{ .key_type = .p256, .raw = .{0} ** 33, .len = 33 };
try v.cache.put("did:plc:aaa", mk); try v.cache.put("did:plc:bbb", mk); try v.cache.put("did:plc:ccc", mk);
// access "aaa" to promote it _ = v.cache.get("did:plc:aaa");
// insert "ddd" — should evict "bbb" (LRU) try v.cache.put("did:plc:ddd", mk);
try std.testing.expect(v.cache.get("did:plc:bbb") == null); try std.testing.expect(v.cache.get("did:plc:aaa") != null); try std.testing.expect(v.cache.get("did:plc:ccc") != null); try std.testing.expect(v.cache.get("did:plc:ddd") != null); try std.testing.expectEqual(@as(u32, 3), v.cache.count());}
test "checkCommitStructure rejects too many ops" { var stats = broadcaster.Stats{}; var v = Validator.initWithConfig(std.testing.allocator, &stats, .{ .max_ops = 2 }, std.testing.io); defer v.deinit();
// build ops array with 3 items (over limit of 2) const ops = [_]zat.cbor.Value{ .{ .map = &.{.{ .key = "action", .value = .{ .text = "create" } }} }, .{ .map = &.{.{ .key = "action", .value = .{ .text = "create" } }} }, .{ .map = &.{.{ .key = "action", .value = .{ .text = "create" } }} }, };
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = "did:plc:test123" } }, .{ .key = "ops", .value = .{ .array = &ops } }, } };
try std.testing.expectError(error.InvalidFrame, v.checkCommitStructure(payload));}
// --- spec conformance tests ---
test "spec: #commit blocks > 2,000,000 bytes rejected" { // lexicon maxLength for #commit blocks: 2,000,000 var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
// insert a fake cached key so we reach the blocks size check const did = "did:plc:test123"; try v.cache.put(did, .{ .key_type = .p256, .raw = .{0} ** 33, .len = 33, .resolve_time = 100, });
// blocks with 2,000,001 bytes (1 byte over limit) const oversized_blocks = try std.testing.allocator.alloc(u8, 2_000_001); defer std.testing.allocator.free(oversized_blocks); @memset(oversized_blocks, 0);
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = did } }, .{ .key = "rev", .value = .{ .text = "3k2abcdefghij" } }, .{ .key = "blocks", .value = .{ .bytes = oversized_blocks } }, } };
const result = v.validateCommit(payload); try std.testing.expect(!result.valid or result.skipped);}
test "spec: #commit blocks = 2,000,000 bytes accepted (boundary)" { // lexicon maxLength for #commit blocks: 2,000,000 — exactly at limit should pass size check var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
const did = "did:plc:test123"; try v.cache.put(did, .{ .key_type = .p256, .raw = .{0} ** 33, .len = 33, .resolve_time = 100, });
// exactly 2,000,000 bytes — should pass size check (may fail signature verify, that's ok) const exact_blocks = try std.testing.allocator.alloc(u8, 2_000_000); defer std.testing.allocator.free(exact_blocks); @memset(exact_blocks, 0);
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = did } }, .{ .key = "rev", .value = .{ .text = "3k2abcdefghij" } }, .{ .key = "blocks", .value = .{ .bytes = exact_blocks } }, } };
const result = v.validateCommit(payload); // should not be rejected for size — may fail signature verification (that's fine, // it means we passed the size check). with P1.1c, sig failure → skipped=true. try std.testing.expect(result.valid or result.skipped);}
test "spec: #sync blocks > 10,000 bytes rejected" { // lexicon maxLength for #sync blocks: 10,000 var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "did", .value = .{ .text = "did:plc:test123" } }, .{ .key = "rev", .value = .{ .text = "3k2abcdefghij" } }, .{ .key = "blocks", .value = .{ .bytes = &([_]u8{0} ** 10_001) } }, } };
const result = v.validateSync(payload); try std.testing.expect(!result.valid); try std.testing.expect(!result.skipped);}
test "spec: #sync blocks = 10,000 bytes accepted (boundary)" { // lexicon maxLength for #sync blocks: 10,000 — exactly at limit should pass size check var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "did", .value = .{ .text = "did:plc:test123" } }, .{ .key = "rev", .value = .{ .text = "3k2abcdefghij" } }, .{ .key = "blocks", .value = .{ .bytes = &([_]u8{0} ** 10_000) } }, } };
const result = v.validateSync(payload); // should pass size check — will be a cache miss → skipped (no cached key) try std.testing.expect(result.valid); try std.testing.expect(result.skipped);}
test "extractOps reads path field from firehose format" { var stats = broadcaster.Stats{}; var v = Validator.initWithConfig(std.testing.allocator, &stats, .{ .verify_commit_diff = true }, std.testing.io); defer v.deinit();
// use arena since extractOps allocates an ArrayList internally var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit();
const ops = [_]zat.cbor.Value{ .{ .map = &.{ .{ .key = "action", .value = .{ .text = "create" } }, .{ .key = "path", .value = .{ .text = "app.bsky.feed.post/3k2abc000000" } }, .{ .key = "cid", .value = .{ .cid = .{ .raw = "fakecid12345" } } }, } }, .{ .map = &.{ .{ .key = "action", .value = .{ .text = "delete" } }, .{ .key = "path", .value = .{ .text = "app.bsky.feed.like/3k2def000000" } }, } }, };
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = "did:plc:test123" } }, .{ .key = "ops", .value = .{ .array = &ops } }, } };
const result = v.extractOps(arena.allocator(), payload); try std.testing.expect(result != null); try std.testing.expectEqual(@as(usize, 2), result.?.len); try std.testing.expectEqualStrings("app.bsky.feed.post/3k2abc000000", result.?[0].path); try std.testing.expectEqualStrings("app.bsky.feed.like/3k2def000000", result.?[1].path); try std.testing.expect(result.?[0].value != null); // create has cid try std.testing.expect(result.?[1].value == null); // delete has no cid}
test "extractOps rejects malformed path without slash" { var stats = broadcaster.Stats{}; var v = Validator.initWithConfig(std.testing.allocator, &stats, .{ .verify_commit_diff = true }, std.testing.io); defer v.deinit();
var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit();
const ops = [_]zat.cbor.Value{ .{ .map = &.{ .{ .key = "action", .value = .{ .text = "create" } }, .{ .key = "path", .value = .{ .text = "noslash" } }, .{ .key = "cid", .value = .{ .cid = .{ .raw = "fakecid12345" } } }, } }, };
const payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = "did:plc:test123" } }, .{ .key = "ops", .value = .{ .array = &ops } }, } };
// malformed path (no slash) → all ops skipped → returns null const result = v.extractOps(arena.allocator(), payload); try std.testing.expect(result == null);}
test "checkCommitStructure validates path field" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
// valid path const valid_ops = [_]zat.cbor.Value{ .{ .map = &.{ .{ .key = "action", .value = .{ .text = "create" } }, .{ .key = "path", .value = .{ .text = "app.bsky.feed.post/3k2abc000000" } }, } }, };
const valid_payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = "did:plc:test123" } }, .{ .key = "ops", .value = .{ .array = &valid_ops } }, } };
try v.checkCommitStructure(valid_payload);
// invalid collection in path const invalid_ops = [_]zat.cbor.Value{ .{ .map = &.{ .{ .key = "action", .value = .{ .text = "create" } }, .{ .key = "path", .value = .{ .text = "not-an-nsid/3k2abc000000" } }, } }, };
const invalid_payload: zat.cbor.Value = .{ .map = &.{ .{ .key = "repo", .value = .{ .text = "did:plc:test123" } }, .{ .key = "ops", .value = .{ .array = &invalid_ops } }, } };
try std.testing.expectError(error.InvalidFrame, v.checkCommitStructure(invalid_payload));}
test "queueResolve deduplicates repeated DIDs" { var stats = broadcaster.Stats{}; var v = Validator.init(std.testing.allocator, &stats, std.testing.io); defer v.deinit();
// queue the same DID 100 times for (0..100) |_| { v.queueResolve("did:plc:duplicate"); }
// should have exactly 1 entry, not 100 try std.testing.expectEqual(@as(usize, 1), v.queue.items.len); try std.testing.expectEqual(@as(u32, 1), v.queued_set.count());}