atproto relay in zig zlay.waow.tech
relay zig atproto
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341//! host ops — background thread for host bookkeeping (cursor flush, failure tracking)//!//! two mechanisms://!//! 1. CursorMap — coalescing cursor store. subscribers atomically write their//! latest seq (nanoseconds, no lock). the worker thread sweeps every 5s and//! batch-flushes only changed cursors to the DB. this replaces the old FIFO//! cursor queue that caused ~350 individual UPDATEs/sec and busy-spin on full.//!//! 2. HostOpsQueue — MPSC ring for rare ops (increment_failures, reset_failures,//! update_status). same atomic spinlock pattern as BroadcastQueue. these ops//! are infrequent (connect/disconnect/error), so the queue never fills.
const std = @import("std");const event_log_mod = @import("event_log.zig");const broadcaster = @import("broadcaster.zig");const util = @import("util/util.zig");
const Io = std.Io;const log = std.log.scoped(.relay);const timestamp = util.timestamp;
// ---------------------------------------------------------------------------// cursor coalescing map// ---------------------------------------------------------------------------
pub const CursorMap = struct { const MAX_SLOTS = 4096; /// sentinel: slot is inactive (freed or never used) — worker skips these const EMPTY: u64 = 0;
host_ids: [MAX_SLOTS]u64 = @splat(0), seqs: [MAX_SLOTS]std.atomic.Value(u64) = @splat(.{ .raw = EMPTY }), /// last value successfully written to DB (worker-thread only, no atomics) last_flushed: [MAX_SLOTS]u64 = @splat(0), /// high-water mark — slots [0..count) have been used at least once count: std.atomic.Value(u32) = .{ .raw = 0 },
/// free list for slot reuse (protected by free_lock spinlock) free_stack: [MAX_SLOTS]u32 = undefined, free_top: u32 = 0, free_lock: std.atomic.Value(u32) = .{ .raw = 0 },
/// register a host and return its slot index. called at spawn time (slurper). pub fn register(self: *CursorMap, host_id: u64, initial_seq: u64) ?u32 { // try reusing a freed slot first const slot = self.popFree() orelse blk: { const s = self.count.fetchAdd(1, .acq_rel); if (s >= MAX_SLOTS) { _ = self.count.fetchSub(1, .monotonic); log.warn("cursor_map: out of slots (max {d})", .{MAX_SLOTS}); return null; } break :blk s; }; // write host_id and last_flushed before making the slot visible via seqs self.host_ids[slot] = host_id; self.last_flushed[slot] = initial_seq; // release: makes host_ids write visible to worker thread's acquire load self.seqs[slot].store(initial_seq, .release); return slot; }
/// return a slot to the free list. called when a subscriber exits (slurper). pub fn unregister(self: *CursorMap, slot: u32) void { // mark inactive first — worker sees EMPTY on next sweep and skips self.seqs[slot].store(EMPTY, .release); self.host_ids[slot] = 0; self.pushFree(slot); }
/// update cursor (called from subscriber — any Io context). single atomic store. pub fn update(self: *CursorMap, slot: u32, seq: u64) void { self.seqs[slot].store(seq, .release); }
/// sweep all slots, flush changed cursors to DB. called from worker thread only. fn flush(self: *CursorMap, persist: *event_log_mod.DiskPersist) void { const n = self.count.load(.acquire); var flushed: u32 = 0; for (0..n) |i| { const current = self.seqs[i].load(.acquire); // EMPTY means freed or never-used slot — skip if (current == EMPTY) continue; if (current != self.last_flushed[i]) { persist.updateHostSeq(self.host_ids[i], current) catch |err| { log.debug("cursor flush failed for host_id={d}: {s}", .{ self.host_ids[i], @errorName(err) }); continue; }; self.last_flushed[i] = current; flushed += 1; } } if (flushed > 0) { log.debug("cursor_map: flushed {d}/{d} cursors", .{ flushed, n }); } }
fn popFree(self: *CursorMap) ?u32 { self.lockFree(); defer self.unlockFree(); if (self.free_top == 0) return null; self.free_top -= 1; return self.free_stack[self.free_top]; }
fn pushFree(self: *CursorMap, slot: u32) void { self.lockFree(); defer self.unlockFree(); self.free_stack[self.free_top] = slot; self.free_top += 1; }
fn lockFree(self: *CursorMap) void { while (self.free_lock.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { std.atomic.spinLoopHint(); } }
fn unlockFree(self: *CursorMap) void { self.free_lock.store(0, .release); }};
// ---------------------------------------------------------------------------// rare ops queue (failures, status)// ---------------------------------------------------------------------------
pub const HostOp = struct { host_id: u64, kind: Kind, payload: Payload,
pub const Kind = enum { increment_failures, reset_failures, update_status, takedown_user, };
pub const Payload = union { /// pointer to subscriber's host_shutdown atomic — set by worker on exhaustion host_shutdown: *std.atomic.Value(bool), none: void, status: Status, takedown: Takedown,
pub const Status = struct { buf: [16]u8 = @splat(0), len: u8 = 0,
pub fn init(s: []const u8) Status { var result: Status = .{}; const n: u8 = @intCast(@min(s.len, 16)); @memcpy(result.buf[0..n], s[0..n]); result.len = n; return result; }
pub fn slice(self: *const Status) []const u8 { return self.buf[0..self.len]; } };
/// inline buffer for takedown — #account CBOR frame is <200 bytes pub const Takedown = struct { uid: u64, frame_buf: [256]u8 = @splat(0), frame_len: u16 = 0,
pub fn frameSlice(self: *const Takedown) []const u8 { return self.frame_buf[0..self.frame_len]; } }; };};
pub const HostOpsQueue = struct { const CAPACITY = 4096; const CURSOR_FLUSH_INTERVAL_SEC: i64 = 5;
items: [CAPACITY]HostOp = undefined, head: std.atomic.Value(u32) = .{ .raw = 0 }, tail: std.atomic.Value(u32) = .{ .raw = 0 }, push_lock: std.atomic.Value(u32) = .{ .raw = 0 },
persist: *event_log_mod.DiskPersist, cursor_map: *CursorMap, shutdown: *std.atomic.Value(bool), max_consecutive_failures: u32, bc: ?*broadcaster.Broadcaster = null,
/// push a rare op (called from any Io context). spins until space is available. /// only used for failures/status — cursors go through CursorMap. pub fn push(self: *HostOpsQueue, op: HostOp) void { while (true) { // acquire spinlock while (self.push_lock.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { std.atomic.spinLoopHint(); }
const tail = self.tail.load(.monotonic); const next_tail = (tail + 1) % CAPACITY; if (next_tail == self.head.load(.acquire)) { // full — release lock, yield, retry self.push_lock.store(0, .release); std.atomic.spinLoopHint(); continue; }
self.items[tail] = op; self.tail.store(next_tail, .release); self.push_lock.store(0, .release); return; } }
/// pop an op (single-consumer — worker thread only). fn pop(self: *HostOpsQueue) ?HostOp { const head = self.head.load(.monotonic); if (head == self.tail.load(.acquire)) return null;
const item = self.items[head]; self.head.store((head + 1) % CAPACITY, .release); return item; }
/// worker thread entry point. runs on pool_io (Threaded). /// drains rare ops + playback requests immediately, sweeps cursor map every 5s. pub fn run(self: *HostOpsQueue, pool_io: Io) void { var last_cursor_flush: i64 = timestamp(pool_io);
while (!self.shutdown.load(.acquire)) { // drain rare ops (failures, status, takedowns) var drained: u32 = 0; while (self.pop()) |op| { self.execute(op); drained += 1; }
// drain playback requests drained += self.drainPlaybackRequests();
// periodic cursor sweep const now = timestamp(pool_io); if (now - last_cursor_flush >= CURSOR_FLUSH_INTERVAL_SEC) { self.cursor_map.flush(self.persist); last_cursor_flush = now; }
if (drained == 0) { // brief sleep — rare ops are infrequent, cursor sweep is timer-driven pool_io.sleep(Io.Duration.fromMilliseconds(100), .awake) catch return; } }
// final cursor flush + drain remaining ops on shutdown self.cursor_map.flush(self.persist); while (self.pop()) |op| { self.execute(op); } _ = self.drainPlaybackRequests(); }
/// drain all pending playback requests from the MPSC queue fn drainPlaybackRequests(self: *HostOpsQueue) u32 { var maybe_batch = self.persist.popPlaybackBatch(); var count: u32 = 0; // Treiber stack pops in LIFO order — fine for playback (each request is independent) while (maybe_batch) |req| { maybe_batch = req.next.load(.acquire); self.persist.playback(req.since, req.allocator, &req.entries, req.limits) catch |e| { req.err = e; }; req.done.store(true, .release); count += 1; } return count; }
fn execute(self: *HostOpsQueue, op: HostOp) void { switch (op.kind) { .increment_failures => { const failures = self.persist.incrementHostFailures(op.host_id) catch 0; if (failures >= self.max_consecutive_failures) { log.warn("host_ops: host_id={d} exhausted after {d} failures", .{ op.host_id, failures }); self.persist.updateHostStatus(op.host_id, "exhausted") catch {}; op.payload.host_shutdown.store(true, .release); } }, .reset_failures => { self.persist.resetHostFailures(op.host_id) catch |err| { log.debug("host_ops: reset failures failed for host_id={d}: {s}", .{ op.host_id, @errorName(err) }); }; }, .update_status => { self.persist.updateHostStatus(op.host_id, op.payload.status.slice()) catch |err| { log.debug("host_ops: update status failed for host_id={d}: {s}", .{ op.host_id, @errorName(err) }); }; }, .takedown_user => { self.executeTakedown(op.payload.takedown); }, } }
/// execute takedown on pool_io: takeDownUser + persist + broadcast fn executeTakedown(self: *HostOpsQueue, td: HostOp.Payload.Takedown) void { self.persist.takeDownUser(td.uid) catch |err| { log.warn("host_ops: takedown failed for uid={d}: {s}", .{ td.uid, @errorName(err) }); return; };
if (td.frame_len == 0) return; const frame_bytes = td.frameSlice(); const bc = self.bc orelse return;
// persist the #account event under ordering lock while (bc.persist_order.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { std.atomic.spinLoopHint(); }
if (self.persist.persist(.account, td.uid, frame_bytes)) |relay_seq| { bc.stats.relay_seq.store(relay_seq, .release);
const broadcast_data = broadcaster.resequenceFrame(self.persist.allocator, frame_bytes, relay_seq) orelse frame_bytes; const owned = self.persist.allocator.dupe(u8, broadcast_data) catch { bc.persist_order.store(0, .release); log.warn("host_ops: failed to alloc broadcast data for takedown uid={d}", .{td.uid}); return; }; bc.broadcast_queue.push(relay_seq, owned, &bc.stats); bc.persist_order.store(0, .release); log.info("host_ops: emitted #account takedown for uid={d} (seq={d})", .{ td.uid, relay_seq }); } else |err| { bc.persist_order.store(0, .release); log.warn("host_ops: persist #account takedown failed: {s}", .{@errorName(err)}); } }};