atproto relay in zig zlay.waow.tech
relay zig atproto
Something went wrong. Try again.
10 kB · 322 lines
Zig
at zig-0.17
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323//! collection index resyncer — updates collection index on #sync events//!//! a #sync event means a repo discontinuity (migration, rebase, bulk import).//! the collection index may be stale for that DID. this worker fetches//! describeRepo from the PDS to get the current collection list, then//! replaces the index entries for that DID.//!//! runs on pool_io — enqueue() is called from frame worker threads,//! background worker is a plain std.Thread.
const std = @import("std");const Io = std.Io;const http = std.http;const collection_index_mod = @import("index.zig");
const Allocator = std.mem.Allocator;const log = std.log.scoped(.resync);
const queue_capacity = 4096;
const ResyncItem = struct { did_buf: [128]u8, did_len: u8, host_buf: [256]u8, host_len: u16,
fn did(self: *const ResyncItem) []const u8 { return self.did_buf[0..self.did_len]; }
fn hostname(self: *const ResyncItem) []const u8 { return self.host_buf[0..self.host_len]; }};
pub const Resyncer = struct { allocator: Allocator, /// pool_io (Threaded) — used for all synchronization AND the HTTP client. /// this struct must never touch the Evented io. io: Io, collection_index: *collection_index_mod.CollectionIndex,
// bounded ring buffer queue queue: [queue_capacity]ResyncItem, head: usize, tail: usize, len: usize, mutex: Io.Mutex, cond: Io.Condition,
running: std.atomic.Value(bool), thread: ?std.Thread,
// stats processed: std.atomic.Value(u64), failed: std.atomic.Value(u64), dropped: std.atomic.Value(u64),
pub fn init( allocator: Allocator, io: Io, collection_index: *collection_index_mod.CollectionIndex, ) Resyncer { return .{ .allocator = allocator, .io = io, .collection_index = collection_index, .queue = undefined, .head = 0, .tail = 0, .len = 0, .mutex = Io.Mutex.init, .cond = Io.Condition.init, .running = .{ .raw = false }, .thread = null, .processed = .{ .raw = 0 }, .failed = .{ .raw = 0 }, .dropped = .{ .raw = 0 }, }; }
/// start the background worker thread. pub fn start(self: *Resyncer) !void { self.running.store(true, .release); self.thread = try std.Thread.spawn(.{}, run, .{self}); }
/// enqueue a DID for resync. non-blocking, drops if queue full or inputs too long. /// called from frame worker threads (plain std.Thread). pub fn enqueue(self: *Resyncer, did: []const u8, hostname: []const u8) void { if (did.len == 0 or did.len > 128 or hostname.len == 0 or hostname.len > 256) return;
self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io);
if (self.len >= queue_capacity) { _ = self.dropped.fetchAdd(1, .monotonic); return; }
var item: ResyncItem = .{ .did_buf = undefined, .did_len = @intCast(did.len), .host_buf = undefined, .host_len = @intCast(hostname.len), }; @memcpy(item.did_buf[0..did.len], did); @memcpy(item.host_buf[0..hostname.len], hostname);
self.queue[self.tail] = item; self.tail = (self.tail + 1) % queue_capacity; self.len += 1; self.cond.signal(self.io); }
/// dequeue one item. blocks until available or shutdown. fn dequeue(self: *Resyncer) ?ResyncItem { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io);
while (self.len == 0 and self.running.load(.acquire)) { self.cond.waitUncancelable(self.io, &self.mutex); }
if (self.len == 0) return null;
const item = self.queue[self.head]; self.head = (self.head + 1) % queue_capacity; self.len -= 1; return item; }
fn run(self: *Resyncer) void { log.info("resync worker started", .{});
var client: http.Client = .{ .allocator = self.allocator, .io = self.io }; defer client.deinit();
while (self.running.load(.acquire)) { const item = self.dequeue() orelse continue; self.processItem(&client, &item);
// brief pause between items self.io.sleep(Io.Duration.fromMilliseconds(50), .awake) catch {}; }
log.info("resync worker stopped (processed={d}, failed={d}, dropped={d})", .{ self.processed.load(.monotonic), self.failed.load(.monotonic), self.dropped.load(.monotonic), }); }
fn processItem(self: *Resyncer, client: *http.Client, item: *const ResyncItem) void { const did = item.did(); const hostname = item.hostname();
// fetch describeRepo from PDS var url_buf: [512]u8 = undefined; const url = std.fmt.bufPrint(&url_buf, "https://{s}/xrpc/com.atproto.repo.describeRepo?repo={s}", .{ hostname, did }) catch { _ = self.failed.fetchAdd(1, .monotonic); return; };
var aw: std.Io.Writer.Allocating = .init(self.allocator); defer aw.deinit();
const result = client.fetch(.{ .location = .{ .url = url }, .response_writer = &aw.writer, .method = .GET, }) catch |err| { log.debug("resync: describeRepo failed for {s} on {s}: {s}", .{ did, hostname, @errorName(err) }); _ = self.failed.fetchAdd(1, .monotonic); return; };
if (result.status != .ok) { log.debug("resync: describeRepo {d} for {s} on {s}", .{ @backingInt(result.status), did, hostname }); _ = self.failed.fetchAdd(1, .monotonic); return; }
const body = aw.written();
// parse {"collections":["app.bsky.feed.post",...], ...} const parsed = std.json.parseFromSlice( DescribeRepoResponse, self.allocator, body, .{ .ignore_unknown_fields = true }, ) catch { log.debug("resync: JSON parse failed for {s}", .{did}); _ = self.failed.fetchAdd(1, .monotonic); return; }; defer parsed.deinit();
const collections = parsed.value.collections orelse { // no collections field — just remove all self.collection_index.removeAll(did) catch {}; _ = self.processed.fetchAdd(1, .monotonic); log.debug("resync: {s} — no collections, removed all", .{did}); return; };
// removeAll then re-add each collection self.collection_index.removeAll(did) catch |err| { log.debug("resync: removeAll failed for {s}: {s}", .{ did, @errorName(err) }); _ = self.failed.fetchAdd(1, .monotonic); return; };
for (collections) |collection| { self.collection_index.addCollection(did, collection) catch {}; }
_ = self.processed.fetchAdd(1, .monotonic); log.debug("resync: {s} — {d} collections indexed", .{ did, collections.len }); }
pub fn queueDepth(self: *Resyncer) usize { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); return self.len; }
pub fn stop(self: *Resyncer) void { self.running.store(false, .release); self.cond.signal(self.io); }
pub fn deinit(self: *Resyncer) void { self.stop(); if (self.thread) |t| t.join(); self.thread = null; }
const DescribeRepoResponse = struct { collections: ?[]const []const u8 = null, };};
// --- tests ---
test "Resyncer init and enqueue" { // just test the queue mechanics, not the HTTP/RocksDB parts var r: Resyncer = .{ .allocator = std.testing.allocator, .io = std.testing.io, .collection_index = undefined, // not used in this test .queue = undefined, .head = 0, .tail = 0, .len = 0, .mutex = Io.Mutex.init, .cond = Io.Condition.init, .running = .{ .raw = true }, .thread = null, .processed = .{ .raw = 0 }, .failed = .{ .raw = 0 }, .dropped = .{ .raw = 0 }, };
// enqueue and dequeue r.enqueue("did:plc:test123", "pds.example.com"); try std.testing.expectEqual(@as(usize, 1), r.queueDepth());
const item = r.dequeue().?; try std.testing.expectEqualStrings("did:plc:test123", item.did()); try std.testing.expectEqualStrings("pds.example.com", item.hostname()); try std.testing.expectEqual(@as(usize, 0), r.queueDepth());}
test "Resyncer drops when full" { var r: Resyncer = .{ .allocator = std.testing.allocator, .io = std.testing.io, .collection_index = undefined, .queue = undefined, .head = 0, .tail = 0, .len = queue_capacity, // pretend full .mutex = Io.Mutex.init, .cond = Io.Condition.init, .running = .{ .raw = true }, .thread = null, .processed = .{ .raw = 0 }, .failed = .{ .raw = 0 }, .dropped = .{ .raw = 0 }, };
r.enqueue("did:plc:test", "pds.example.com"); try std.testing.expectEqual(@as(u64, 1), r.dropped.load(.monotonic));}
test "Resyncer rejects oversized inputs" { var r: Resyncer = .{ .allocator = std.testing.allocator, .io = std.testing.io, .collection_index = undefined, .queue = undefined, .head = 0, .tail = 0, .len = 0, .mutex = Io.Mutex.init, .cond = Io.Condition.init, .running = .{ .raw = true }, .thread = null, .processed = .{ .raw = 0 }, .failed = .{ .raw = 0 }, .dropped = .{ .raw = 0 }, };
// empty DID r.enqueue("", "pds.example.com"); try std.testing.expectEqual(@as(usize, 0), r.queueDepth());
// empty hostname r.enqueue("did:plc:test", ""); try std.testing.expectEqual(@as(usize, 0), r.queueDepth());}