//! 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()); }