//! 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 this const 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; }