atproto relay in zig zlay.waow.tech
relay zig atproto
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341//! upstream status checker — lazy re-check of stale account status//!//! a #commit/#sync arriving from an account's own (host-authority-confirmed)//! PDS while our record says upstream-inactive is self-contradicting: a PDS//! does not emit repo events for accounts that are inactive on it. the usual//! cause is a lost one-shot #account activation during a migration (the old//! host's deactivation was recorded, the new host's activation was dropped),//! which otherwise mutes the account forever.//!//! this worker resolves the contradiction the way indigo's EnsureAccountActive//! does — ask the PDS via com.atproto.sync.getRepoStatus and store the answer —//! but queued on a background thread instead of inline, so frame workers never//! block on network I/O. see docs/handoffs/HANDOFF-2026-08-06-stale-upstream-status.md.//!//! runs on pool_io — enqueue() is called from frame worker threads,//! background worker is a plain std.Thread. mirrors resync.zig.
const std = @import("std");const Io = std.Io;const http = std.http;const event_log_mod = @import("event_log.zig");const host_check = @import("atproto/host_check.zig");
const Allocator = std.mem.Allocator;const log = std.log.scoped(.status_check);
const queue_capacity = 4096;/// per-uid re-check floor — a stuck account triggers on every commit, so/// without this a busy account would hammer its PDS between queue drains.const recheck_interval_seconds: i64 = 60;/// prune the per-uid attempt map when it grows past thisconst attempt_map_prune_at: usize = 65_536;
const CheckItem = struct { uid: u64, did_buf: [128]u8, did_len: u8, host_buf: [256]u8, host_len: u16,
fn did(self: *const CheckItem) []const u8 { return self.did_buf[0..self.did_len]; }
fn hostname(self: *const CheckItem) []const u8 { return self.host_buf[0..self.host_len]; }};
pub const StatusChecker = struct { allocator: Allocator, /// pool_io (Threaded) — used for all synchronization AND the HTTP client. io: Io, persist: *event_log_mod.DiskPersist,
// bounded ring buffer queue queue: [queue_capacity]CheckItem = undefined, head: usize = 0, tail: usize = 0, len: usize = 0, mutex: Io.Mutex = Io.Mutex.init, cond: Io.Condition = Io.Condition.init,
// uid → epoch seconds of last enqueue, for dedup + rate limiting last_attempt: std.AutoHashMapUnmanaged(u64, i64) = .empty,
running: std.atomic.Value(bool) = .{ .raw = false }, thread: ?std.Thread = null,
// stats processed: std.atomic.Value(u64) = .{ .raw = 0 }, recovered: std.atomic.Value(u64) = .{ .raw = 0 }, failed: std.atomic.Value(u64) = .{ .raw = 0 }, dropped: std.atomic.Value(u64) = .{ .raw = 0 },
pub fn init(allocator: Allocator, io: Io, persist: *event_log_mod.DiskPersist) StatusChecker { return .{ .allocator = allocator, .io = io, .persist = persist }; }
pub fn start(self: *StatusChecker) !void { self.running.store(true, .release); self.thread = try std.Thread.spawn(.{}, run, .{self}); }
pub fn stop(self: *StatusChecker) void { self.running.store(false, .release); self.cond.signal(self.io); }
pub fn deinit(self: *StatusChecker) void { self.stop(); if (self.thread) |t| t.join(); self.thread = null; self.last_attempt.deinit(self.allocator); }
/// enqueue an account for an upstream status re-check. non-blocking; /// drops if the queue is full, inputs are oversized, or the uid was /// attempted within recheck_interval_seconds. pub fn enqueue(self: *StatusChecker, uid: u64, did: []const u8, hostname: []const u8) void { if (uid == 0 or did.len == 0 or did.len > 128 or hostname.len == 0 or hostname.len > 256) return;
const now = Io.Timestamp.now(self.io, .awake).toSeconds();
self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io);
if (self.last_attempt.get(uid)) |at| { if (now - at < recheck_interval_seconds) return; }
if (self.len >= queue_capacity) { _ = self.dropped.fetchAdd(1, .monotonic); return; }
if (self.last_attempt.count() >= attempt_map_prune_at) { self.pruneAttemptsLocked(now); } self.last_attempt.put(self.allocator, uid, now) catch {};
var item: CheckItem = .{ .uid = uid, .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); }
fn pruneAttemptsLocked(self: *StatusChecker, now: i64) void { var it = self.last_attempt.iterator(); var expired: std.ArrayListUnmanaged(u64) = .empty; defer expired.deinit(self.allocator); while (it.next()) |entry| { if (now - entry.value_ptr.* >= recheck_interval_seconds) { expired.append(self.allocator, entry.key_ptr.*) catch break; } } for (expired.items) |uid| _ = self.last_attempt.remove(uid); if (self.last_attempt.count() >= attempt_map_prune_at) { self.last_attempt.clearRetainingCapacity(); } }
fn dequeue(self: *StatusChecker) ?CheckItem { 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: *StatusChecker) void { log.info("status check 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("status check worker stopped (processed={d}, recovered={d}, failed={d}, dropped={d})", .{ self.processed.load(.monotonic), self.recovered.load(.monotonic), self.failed.load(.monotonic), self.dropped.load(.monotonic), }); }
fn processItem(self: *StatusChecker, client: *http.Client, item: *const CheckItem) void { const did = item.did(); const hostname = item.hostname();
host_check.rejectPrivateHost(self.allocator, hostname) catch { _ = self.failed.fetchAdd(1, .monotonic); return; };
var url_buf: [512]u8 = undefined; const url = std.fmt.bufPrint(&url_buf, "https://{s}/xrpc/com.atproto.sync.getRepoStatus?did={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("getRepoStatus failed for {s} on {s}: {s}", .{ did, hostname, @errorName(err) }); _ = self.failed.fetchAdd(1, .monotonic); return; };
if (result.status != .ok) { log.debug("getRepoStatus {d} for {s} on {s}", .{ @backingInt(result.status), did, hostname }); _ = self.failed.fetchAdd(1, .monotonic); return; }
const new_status = mapRepoStatus(self.allocator, aw.written()) orelse { log.debug("getRepoStatus parse failed for {s}", .{did}); _ = self.failed.fetchAdd(1, .monotonic); return; };
self.persist.updateAccountUpstreamStatus(item.uid, new_status) catch |err| { log.debug("upstream status update failed for uid={d}: {s}", .{ item.uid, @errorName(err) }); _ = self.failed.fetchAdd(1, .monotonic); return; };
_ = self.processed.fetchAdd(1, .monotonic); if (std.mem.eql(u8, new_status, "active")) { _ = self.recovered.fetchAdd(1, .monotonic); log.info("stale upstream status corrected: uid={d} did={s} host={s} -> active", .{ item.uid, did, hostname, }); } else { log.debug("upstream status confirmed non-active: uid={d} did={s} status={s}", .{ item.uid, did, new_status, }); } }
pub fn queueDepth(self: *StatusChecker) usize { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); return self.len; }};
/// map a getRepoStatus JSON body to zlay's upstream_status value./// active → "active"; inactive uses the reported status, "inactive" fallback./// mirrors the #account event mapping in frame_worker.zig / indigo's/// host_checker.go FetchAccountStatus. returns a static string.fn mapRepoStatus(allocator: Allocator, body: []const u8) ?[]const u8 { const Response = struct { active: bool, status: ?[]const u8 = null, }; const parsed = std.json.parseFromSlice(Response, allocator, body, .{ .ignore_unknown_fields = true, }) catch return null; defer parsed.deinit();
if (parsed.value.active) return "active"; const status = parsed.value.status orelse return "inactive"; // normalize to the known set so we store static strings const known = [_][]const u8{ "deactivated", "deleted", "takendown", "suspended", "desynchronized", "throttled" }; for (known) |k| { if (std.mem.eql(u8, status, k)) return k; } return "inactive";}
// --- tests ---
test "mapRepoStatus maps active, known, and unknown statuses" { const a = std.testing.allocator; try std.testing.expectEqualStrings("active", mapRepoStatus(a, "{\"did\":\"did:plc:x\",\"active\":true}").?); try std.testing.expectEqualStrings("deactivated", mapRepoStatus(a, "{\"active\":false,\"status\":\"deactivated\"}").?); try std.testing.expectEqualStrings("takendown", mapRepoStatus(a, "{\"active\":false,\"status\":\"takendown\"}").?); // unknown status normalizes to inactive try std.testing.expectEqualStrings("inactive", mapRepoStatus(a, "{\"active\":false,\"status\":\"weird\"}").?); // missing status on inactive account try std.testing.expectEqualStrings("inactive", mapRepoStatus(a, "{\"active\":false}").?); // garbage try std.testing.expect(mapRepoStatus(a, "not json") == null);}
test "StatusChecker enqueue dedups within recheck interval" { var sc: StatusChecker = .{ .allocator = std.testing.allocator, .io = std.testing.io, .persist = undefined, // not used in queue mechanics .running = .{ .raw = true }, }; defer sc.last_attempt.deinit(std.testing.allocator);
sc.enqueue(7, "did:plc:test123", "pds.example.com"); try std.testing.expectEqual(@as(usize, 1), sc.queueDepth());
// same uid again within the interval — dropped silently sc.enqueue(7, "did:plc:test123", "pds.example.com"); try std.testing.expectEqual(@as(usize, 1), sc.queueDepth());
// different uid still goes through sc.enqueue(8, "did:plc:other", "pds.example.com"); try std.testing.expectEqual(@as(usize, 2), sc.queueDepth());
const item = sc.dequeue().?; try std.testing.expectEqual(@as(u64, 7), item.uid); try std.testing.expectEqualStrings("did:plc:test123", item.did()); try std.testing.expectEqualStrings("pds.example.com", item.hostname());}
test "StatusChecker rejects invalid inputs and drops when full" { var sc: StatusChecker = .{ .allocator = std.testing.allocator, .io = std.testing.io, .persist = undefined, .running = .{ .raw = true }, }; defer sc.last_attempt.deinit(std.testing.allocator);
sc.enqueue(0, "did:plc:test", "pds.example.com"); // uid 0 sc.enqueue(1, "", "pds.example.com"); // empty did sc.enqueue(2, "did:plc:test", ""); // empty hostname try std.testing.expectEqual(@as(usize, 0), sc.queueDepth());
sc.len = queue_capacity; // pretend full sc.enqueue(3, "did:plc:test", "pds.example.com"); try std.testing.expectEqual(@as(u64, 1), sc.dropped.load(.monotonic)); sc.len = 0;}