diff --git a/docs/handoffs/HANDOFF-2026-08-06-stale-upstream-status.md b/docs/handoffs/HANDOFF-2026-08-06-stale-upstream-status.md new file mode 100644 index 0000000..6a633bc --- /dev/null +++ b/docs/handoffs/HANDOFF-2026-08-06-stale-upstream-status.md @@ -0,0 +1,96 @@ +# handoff 2026-08-06 — stale upstream_status after migration: fixed, sweep ready + +for the operator, re: the coverage-gap investigation (zlay 98.7% vs 99.1%, +~70% of misses on migration-destination PDSes). your diagnosis was right: +migrated accounts stuck with `upstream_status='deactivated'` while active +on their PDS. root cause confirmed at code level, fix on branch, recovery +sweep script included. + +## root cause (one-shot event loss, no recovery path) + +during a migration the old PDS's `#account` deactivation is recorded +(legitimately — the stored host still matches, so no authority check runs). +the new PDS then emits `#account active=true` **exactly once**. that event +carries a host change, which triggers the synchronous DID-doc host-authority +check — and if it rejects (PLC propagation lag, a transient resolver +failure, or the 60s negative rejection cache primed by an earlier racing +event), the event is dropped wholesale and nothing ever retries it. + +subsequent commits from the new host do eventually pass authority and fix +the host pointer, but each is then dropped by the inactive-account check +(`status AND upstream_status` both required). permanent mute — the exact +db state you found: correct host, `status='active'`, +`upstream_status='deactivated'`. + +reviewed indigo's relay and hydrant before fixing. indigo has the same +event-loss window but never gets stuck: its `EnsureAccountActive` re-checks +`getRepoStatus` against the PDS whenever a commit arrives for an +upstream-inactive account. that lazy re-check is the mechanism zlay had +removed (the old TODO in frame_worker.zig). the key property that makes it +cheap: a PDS does not emit repo events for accounts inactive on it, so a +commit from the account's own confirmed host contradicting our record is +itself strong evidence the record is stale — the check almost only fires +when it should. + +## the fix + +`fix: recover accounts muted by stale upstream account status` + +1. **StatusChecker** (`src/internal/status_checker.zig`) — background + worker (indigo's EnsureAccountActive, but queued off the frame workers). + when a commit/sync from the account's own host is dropped with + `status='active'` but `upstream_status!='active'`, the account is queued; + the worker asks the PDS via `com.atproto.sync.getRepoStatus` and stores + the answer. per-uid 60s rate limit, 4096-deep queue, SSRF-guarded, + local takedowns never re-checked. +2. **`#account` events bypass the negative host-authority cache** — they + are rare and one-shot; a cached rejection from an event racing PLC + propagation must not swallow a migration's activation. (indigo/hydrant + instead purge their DID-doc caches and re-resolve on mismatch; zlay + resolves fresh every time, so the negative cache was our only stale + layer.) + +regression tests: status mapping vs the real getRepoStatus response shape, +queue dedup/rate-limit mechanics, and the local-vs-upstream activity split. +full suite passes. + +## expected signals after deploy + +- `status check worker started` at boot. +- `stale upstream status corrected: uid=... did=... host=... -> active` + at info level — that's the fix working. each such account resumes + flowing from its next commit (~one event lost per recovery, same as + indigo's behavior). +- coverage on eurosky/blacksky/northsky should start climbing organically + as stuck accounts post. + +## the recovery sweep (run once, after the fix is deployed) + +the fix only heals accounts that commit; the long tail of quiet stuck +accounts needs the one-time sweep. `scripts/sweep_stale_upstream.py`: + + DATABASE_URL=postgres://... uv run scripts/sweep_stale_upstream.py # dry run + DATABASE_URL=postgres://... uv run scripts/sweep_stale_upstream.py --apply # write + +it resolves each candidate's **current PDS from plc.directory** (the db +host pointer is itself stale for accounts that never committed since +migrating — don't trust it), probes getRepoStatus there, and only writes +what the PDS reports. genuinely deactivated accounts (your 12/12 mushroom +control) are confirmed and left alone. dry-run first and eyeball the +would-fix lines; safe to re-run. ~2.4k+ candidates × 2 small HTTP GETs, +50ms pacing — minutes, negligible bandwidth. + +host pointers are deliberately left alone: once unmuted, an account's next +commit migrates its pointer through the normal authority path. + +## deploy + +normal flow: `just zlay publish-remote ReleaseSafe zio` with a ref +at or past the fix commit. + +## still open + +- the collection index entries removed at deactivation repopulate lazily + from new commits, not at recovery time. acceptable; noting it. +- consumer connect/disconnect logs still lack remote addr + cursor (carried + over from the OOM handoff). diff --git a/scripts/sweep_stale_upstream.py b/scripts/sweep_stale_upstream.py new file mode 100644 index 0000000..036745e --- /dev/null +++ b/scripts/sweep_stale_upstream.py @@ -0,0 +1,123 @@ +# /// script +# requires-python = ">=3.12" +# dependencies = ["psycopg[binary]", "httpx"] +# /// +"""one-time recovery sweep for accounts muted by stale upstream_status. + +context: docs/handoffs/HANDOFF-2026-08-06-stale-upstream-status.md. +accounts that migrated PDSes lost their one-shot #account activation and +are stuck with status='active' AND upstream_status != 'active'. for each, +this resolves the DID's current PDS from plc.directory (host_id in the db +may itself be stale for accounts that never committed since migrating — +do not trust it), asks that PDS via getRepoStatus, and writes the answer +back. genuinely inactive accounts are left untouched. + +usage: + DATABASE_URL=postgres://... uv run scripts/sweep_stale_upstream.py [--apply] + +dry-run by default; --apply writes. safe to re-run. +""" + +import argparse +import os +import sys +import time + +import httpx +import psycopg + +KNOWN_STATUSES = {"deactivated", "deleted", "takendown", "suspended", "desynchronized", "throttled"} + + +def resolve_pds(client: httpx.Client, did: str) -> str | None: + if did.startswith("did:plc:"): + r = client.get(f"https://plc.directory/{did}") + if r.status_code != 200: + return None + doc = r.json() + elif did.startswith("did:web:"): + host = did.removeprefix("did:web:") + if ":" in host or "/" in host: + return None + r = client.get(f"https://{host}/.well-known/did.json") + if r.status_code != 200: + return None + doc = r.json() + else: + return None + for svc in doc.get("service", []): + if svc.get("id") == "#atproto_pds": + return svc.get("serviceEndpoint") + return None + + +def fetch_status(client: httpx.Client, pds: str, did: str) -> str | None: + """returns the upstream_status value to store, or None on probe failure.""" + r = client.get(f"{pds}/xrpc/com.atproto.sync.getRepoStatus", params={"did": did}) + if r.status_code != 200: + return None + body = r.json() + if body.get("active"): + return "active" + status = body.get("status") + return status if status in KNOWN_STATUSES else "inactive" + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--apply", action="store_true", help="write corrections (default: dry run)") + parser.add_argument("--limit", type=int, default=0, help="stop after N accounts (0 = all)") + args = parser.parse_args() + + database_url = os.environ.get("DATABASE_URL") + if not database_url: + sys.exit("DATABASE_URL is required") + + conn = psycopg.connect(database_url) + rows = conn.execute( + "SELECT uid, did, upstream_status FROM account" + " WHERE status = 'active' AND upstream_status != 'active' ORDER BY uid" + ).fetchall() + print(f"{len(rows)} candidate accounts (status=active, upstream_status!=active)") + + client = httpx.Client(timeout=10, headers={"user-agent": "zlay-sweep (atproto-relay)"}) + checked = corrected = confirmed = failed = 0 + + for uid, did, stored in rows: + if args.limit and checked >= args.limit: + break + checked += 1 + + pds = resolve_pds(client, did) + if not pds: + failed += 1 + print(f" probe-failed uid={uid} {did}: no PDS resolved") + continue + + actual = fetch_status(client, pds, did) + if actual is None: + failed += 1 + print(f" probe-failed uid={uid} {did} @ {pds}") + continue + + if actual == stored: + confirmed += 1 + else: + corrected += 1 + print(f" {'FIX' if args.apply else 'would-fix'} uid={uid} {did}: {stored} -> {actual} (pds={pds})") + if args.apply: + conn.execute( + "UPDATE account SET upstream_status = %s WHERE uid = %s", (actual, uid) + ) + conn.commit() + + time.sleep(0.05) # be polite to plc.directory and the PDSes + + print( + f"done: checked={checked} corrected={corrected} confirmed-inactive={confirmed} probe-failed={failed}" + + ("" if args.apply else " (dry run — rerun with --apply)") + ) + + +if __name__ == "__main__": + main() diff --git a/src/internal/event_log.zig b/src/internal/event_log.zig index 2f8656f..ef54d4c 100644 --- a/src/internal/event_log.zig +++ b/src/internal/event_log.zig @@ -688,18 +688,35 @@ pub const DiskPersist = struct { /// check if an account is active (both local status and upstream status). /// Go relay: models.Account.IsActive() - pub fn isAccountActive(self: *DiskPersist, uid: u64) !bool { + pub const AccountActivity = struct { + local_ok: bool, + upstream_ok: bool, + + pub fn active(self: AccountActivity) bool { + return self.local_ok and self.upstream_ok; + } + }; + + /// account activity split by source: local (relay takedown) vs upstream + /// (PDS-reported via #account events). callers that drop events for + /// upstream-inactive accounts use the distinction to trigger a + /// StatusChecker re-check — local takedowns must never be re-checked away. + pub fn accountActivity(self: *DiskPersist, uid: u64) !AccountActivity { var row = (try self.db.rowUnsafe( "SELECT status, upstream_status FROM account WHERE uid = $1", .{@as(i64, @intCast(uid))}, - )) orelse return false; + )) orelse return .{ .local_ok = false, .upstream_ok = false }; defer row.deinit() catch {}; const status = row.get([]const u8, 0); const upstream = row.get([]const u8, 1); - // active if local is active AND upstream is active - const local_ok = std.mem.eql(u8, status, "active"); - const upstream_ok = std.mem.eql(u8, upstream, "active"); - return local_ok and upstream_ok; + return .{ + .local_ok = std.mem.eql(u8, status, "active"), + .upstream_ok = std.mem.eql(u8, upstream, "active"), + }; + } + + pub fn isAccountActive(self: *DiskPersist, uid: u64) !bool { + return (try self.accountActivity(uid)).active(); } // --- host management --- @@ -1964,6 +1981,53 @@ test "uidForDid assigns and caches UIDs" { try std.testing.expect(uid1 != uid2); } +test "accountActivity distinguishes local takedown from stale upstream status" { + const base_url = try requireDatabaseUrl(); + var tdb = try TestDb.create(base_url); + defer tdb.destroy(); + + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + const dir_path = try tmpDirRealPath(std.testing.allocator, tmp); + defer std.testing.allocator.free(dir_path); + + var dp = try DiskPersist.init(std.testing.allocator, dir_path, tdb.url, 5, std.testing.io); + defer dp.deinit(); + + const uid = try dp.uidForDid("did:plc:migrant"); + + // fresh account: fully active + var activity = try dp.accountActivity(uid); + try std.testing.expect(activity.local_ok and activity.upstream_ok); + try std.testing.expect(activity.active()); + + // the migration-stale shape: locally active, upstream deactivated. + // this is the only combination that may trigger a StatusChecker re-check. + try dp.updateAccountUpstreamStatus(uid, "deactivated"); + activity = try dp.accountActivity(uid); + try std.testing.expect(activity.local_ok); + try std.testing.expect(!activity.upstream_ok); + try std.testing.expect(!activity.active()); + try std.testing.expect(!try dp.isAccountActive(uid)); + + // recovery: the checker writes the PDS's answer back + try dp.updateAccountUpstreamStatus(uid, "active"); + activity = try dp.accountActivity(uid); + try std.testing.expect(activity.active()); + try std.testing.expect(try dp.isAccountActive(uid)); + + // local takedown: local_ok false — must never look like the re-check shape + _ = try dp.db.exec("UPDATE account SET status = 'takendown' WHERE uid = $1", .{@as(i64, @intCast(uid))}); + activity = try dp.accountActivity(uid); + try std.testing.expect(!activity.local_ok); + try std.testing.expect(activity.upstream_ok); + try std.testing.expect(!activity.active()); + + // unknown uid: inactive on both axes + activity = try dp.accountActivity(999_999); + try std.testing.expect(!activity.local_ok and !activity.upstream_ok); +} + test "getAccountState retains rev-only legacy state" { const base_url = try requireDatabaseUrl(); var tdb = try TestDb.create(base_url); diff --git a/src/internal/frame_worker.zig b/src/internal/frame_worker.zig index 4a8d2ae..0cbbde4 100644 --- a/src/internal/frame_worker.zig +++ b/src/internal/frame_worker.zig @@ -15,6 +15,7 @@ const validator_mod = @import("validator.zig"); const event_log_mod = @import("event_log.zig"); const collection_index_mod = @import("collection_index/index.zig"); const resync_mod = @import("collection_index/resync.zig"); +const status_checker_mod = @import("status_checker.zig"); const thread_pool = @import("util/thread_pool.zig"); const util = @import("util/util.zig"); @@ -34,6 +35,7 @@ pub const FrameWork = struct { persist: ?*event_log_mod.DiskPersist, collection_index: ?*collection_index_mod.CollectionIndex, resyncer: ?*resync_mod.Resyncer, + status_checker: ?*status_checker_mod.StatusChecker = null, }; pub fn processFrame(work: *FrameWork) void { @@ -102,7 +104,9 @@ pub fn processFrame(work: *FrameWork) void { const elapsed: u64 = @intCast(@max(0, microTimestamp(work.io) - ha_t0)); _ = work.bc.stats.host_authority_time_us.fetchAdd(elapsed, .monotonic); } - switch (work.validator.resolveHostAuthority(d, work.host_id, work.hostname)) { + switch (work.validator.resolveHostAuthority(d, work.host_id, work.hostname, .{ + .skip_reject_cache = is_account, + })) { .migrate => { // DID doc confirms the host — update and continue dp.setAccountHostId(result.uid, work.host_id) catch {}; @@ -164,13 +168,21 @@ pub fn processFrame(work: *FrameWork) void { var commit_rev: ?[]const u8 = null; if (is_commit or is_sync) { // drop events for inactive accounts. - // status self-corrects when the next #account event arrives. - // TODO: reintroduce PDS re-check via background queue with shared - // long-lived http client + timeout (not inline on frame workers). if (work.persist) |dp| { if (uid > 0) { - const active = dp.isAccountActive(uid) catch true; - if (!active) { + const activity = dp.accountActivity(uid) catch + event_log_mod.DiskPersist.AccountActivity{ .local_ok = true, .upstream_ok = true }; + if (!activity.active()) { + // a repo event from the account's own confirmed host while + // upstream_status is non-active is self-contradicting — the + // PDS wouldn't emit it for an account inactive on it. usual + // cause: a migration's #account activation was lost. queue a + // getRepoStatus re-check (never for local takedowns). + if (activity.local_ok and !activity.upstream_ok) { + if (work.status_checker) |sc| { + if (did) |d| sc.enqueue(uid, d, work.hostname); + } + } _ = work.bc.stats.skipped.fetchAdd(1, .monotonic); return; } diff --git a/src/internal/slurper.zig b/src/internal/slurper.zig index abea5e8..0642b46 100644 --- a/src/internal/slurper.zig +++ b/src/internal/slurper.zig @@ -19,6 +19,7 @@ const event_log_mod = @import("event_log.zig"); const subscriber_mod = @import("subscriber.zig"); const collection_index_mod = @import("collection_index/index.zig"); const resync_mod = @import("collection_index/resync.zig"); +const status_checker_mod = @import("status_checker.zig"); const frame_worker_mod = @import("frame_worker.zig"); const host_ops_mod = @import("host_ops.zig"); const atproto = @import("atproto/main.zig"); @@ -53,6 +54,7 @@ pub const Slurper = struct { persist: *event_log_mod.DiskPersist, collection_index: ?*collection_index_mod.CollectionIndex = null, resyncer: ?*resync_mod.Resyncer = null, + status_checker: ?*status_checker_mod.StatusChecker = null, host_ops: ?*host_ops_mod.HostOpsQueue = null, cursor_map: ?*host_ops_mod.CursorMap = null, db_queue: ?*event_log_mod.DbRequestQueue = null, @@ -405,6 +407,7 @@ pub const Slurper = struct { ); sub.collection_index = self.collection_index; sub.resyncer = self.resyncer; + sub.status_checker = self.status_checker; if (self.frame_pool) |*fp| sub.pool = fp; sub.pool_io = self.pool_io; sub.host_ops = self.host_ops; diff --git a/src/internal/status_checker.zig b/src/internal/status_checker.zig new file mode 100644 index 0000000..217796b --- /dev/null +++ b/src/internal/status_checker.zig @@ -0,0 +1,340 @@ +//! 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}", .{ @intFromEnum(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; +} diff --git a/src/internal/subscriber.zig b/src/internal/subscriber.zig index 944f426..1a70101 100644 --- a/src/internal/subscriber.zig +++ b/src/internal/subscriber.zig @@ -14,6 +14,7 @@ const validator_mod = @import("validator.zig"); const event_log_mod = @import("event_log.zig"); const collection_index_mod = @import("collection_index/index.zig"); const resync_mod = @import("collection_index/resync.zig"); +const status_checker_mod = @import("status_checker.zig"); const frame_worker_mod = @import("frame_worker.zig"); const host_ops_mod = @import("host_ops.zig"); const util = @import("util/util.zig"); @@ -198,6 +199,7 @@ pub const Subscriber = struct { persist: ?*event_log_mod.DiskPersist, collection_index: ?*collection_index_mod.CollectionIndex = null, resyncer: ?*resync_mod.Resyncer = null, + status_checker: ?*status_checker_mod.StatusChecker = null, pool: ?*frame_worker_mod.FramePool = null, /// dedicated Threaded io for frame workers — safe from plain OS threads pool_io: ?Io = null, @@ -544,6 +546,7 @@ const FrameHandler = struct { .persist = sub.persist, .collection_index = sub.collection_index, .resyncer = sub.resyncer, + .status_checker = sub.status_checker, }, sub.shutdown)) { // pool accepted — advance cursor past this frame _ = sub.bc.stats.pool_queued_bytes.fetchAdd(duped.len, .monotonic); @@ -598,7 +601,9 @@ const FrameHandler = struct { // host authority enforcement (mirrors frame_worker path) if ((result.is_new or result.host_changed) and !is_identity) { - switch (sub.validator.resolveHostAuthority(d, sub.options.host_id, sub.options.hostname)) { + switch (sub.validator.resolveHostAuthority(d, sub.options.host_id, sub.options.hostname, .{ + .skip_reject_cache = is_account, + })) { .migrate => { dp.setAccountHostId(result.uid, sub.options.host_id) catch {}; if (result.host_changed) { @@ -657,11 +662,17 @@ const FrameHandler = struct { var commit_data_cid: ?[]const u8 = null; var commit_rev: ?[]const u8 = null; if (is_commit or is_sync) { - // drop for inactive accounts + // drop for inactive accounts (see frame_worker for the re-check rationale) if (sub.persist) |dp| { if (uid > 0) { - const active = dp.isAccountActive(uid) catch true; - if (!active) { + const activity = dp.accountActivity(uid) catch + event_log_mod.DiskPersist.AccountActivity{ .local_ok = true, .upstream_ok = true }; + if (!activity.active()) { + if (activity.local_ok and !activity.upstream_ok) { + if (sub.status_checker) |sc| { + if (did) |d| sc.enqueue(uid, d, sub.options.hostname); + } + } _ = sub.bc.stats.skipped.fetchAdd(1, .monotonic); return; } diff --git a/src/internal/validator.zig b/src/internal/validator.zig index dba00b7..6e798b1 100644 --- a/src/internal/validator.zig +++ b/src/internal/validator.zig @@ -832,11 +832,21 @@ pub const Validator = struct { /// .accept — should not happen (caller should only call on new/mismatch) /// .migrate — DID doc confirms this host, caller should update host_id /// .reject — DID doc does not confirm, caller should drop the event + pub const HostAuthorityOptions = struct { + /// skip the negative rejection cache and always do a live DID doc + /// check. used for #account events: they are rare and one-shot — a + /// migration's activation swallowed by a cached rejection (from an + /// earlier event racing PLC propagation) mutes the account until a + /// StatusChecker re-check, so they always deserve a fresh look. + skip_reject_cache: bool = false, + }; + pub fn resolveHostAuthority( self: *Validator, did: []const u8, incoming_host_id: u64, incoming_host: []const u8, + options: HostAuthorityOptions, ) HostAuthority { const persist = self.persist orelse return .migrate; // no DB — can't check const parsed = zat.Did.parse(did) orelse { @@ -844,7 +854,7 @@ pub const Validator = struct { return .reject; }; - if (self.host_authority_cache.recentlyRejected(did, incoming_host_id)) { + if (!options.skip_reject_cache and self.host_authority_cache.recentlyRejected(did, incoming_host_id)) { _ = self.stats.host_authority_cache_hits.fetchAdd(1, .monotonic); return .reject; } diff --git a/src/main.zig b/src/main.zig index ba1591d..47d1ba3 100644 --- a/src/main.zig +++ b/src/main.zig @@ -35,6 +35,7 @@ const collection_index_mod = @import("internal/collection_index/index.zig"); const backfill_mod = @import("internal/collection_index/backfill.zig"); const cleaner_mod = @import("internal/collection_index/cleaner.zig"); const resync_mod = @import("internal/collection_index/resync.zig"); +const status_checker_mod = @import("internal/status_checker.zig"); const host_ops_mod = @import("internal/host_ops.zig"); const api = @import("internal/api/main.zig"); const util = @import("internal/util/util.zig"); @@ -394,6 +395,12 @@ fn runRelay() !void { try resyncer.start(); defer resyncer.deinit(); + // init status checker (lazy getRepoStatus re-check for stale upstream + // account status — recovers accounts muted by a lost migration activation) + var status_checker = status_checker_mod.StatusChecker.init(allocator, pool_io, &dp); + try status_checker.start(); + defer status_checker.deinit(); + // init cursor map + host ops queue — subscriber threads write cursor seqs // to the coalescing map (atomic store, no lock) and push rare ops (failures, // status) to the MPSC queue. a background thread sweeps cursors every 5s @@ -431,6 +438,7 @@ fn runRelay() !void { defer slurper.deinit(); slurper.collection_index = &ci; slurper.resyncer = &resyncer; + slurper.status_checker = &status_checker; slurper.host_ops = &host_ops_queue; slurper.cursor_map = &cursor_map; slurper.db_queue = &db_queue; diff --git a/src/tests.zig b/src/tests.zig index 4101b70..54c7dae 100644 --- a/src/tests.zig +++ b/src/tests.zig @@ -15,6 +15,7 @@ test { _ = @import("internal/event_log.zig"); _ = @import("internal/slurper.zig"); _ = @import("internal/frame_worker.zig"); + _ = @import("internal/status_checker.zig"); _ = @import("internal/collection_index/index.zig"); _ = @import("internal/collection_index/backfill.zig"); _ = @import("main.zig");