# /// 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)") parser.add_argument( "--statuses", default="deactivated", help="comma-separated upstream_status values to sweep (default: deactivated — " "the only state the migration bug produces; 'all' sweeps every non-active status)", ) 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) if args.statuses == "all": rows = conn.execute( "SELECT uid, did, upstream_status FROM account" " WHERE status = 'active' AND upstream_status != 'active' ORDER BY uid" ).fetchall() else: statuses = [s.strip() for s in args.statuses.split(",") if s.strip()] rows = conn.execute( "SELECT uid, did, upstream_status FROM account" " WHERE status = 'active' AND upstream_status = ANY(%s) ORDER BY uid", (statuses,), ).fetchall() print(f"{len(rows)} candidate accounts (status=active, upstream_status in {args.statuses})") 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 try: pds = resolve_pds(client, did) except httpx.HTTPError as e: failed += 1 print(f" probe-failed uid={uid} {did}: resolve error: {type(e).__name__}") continue if not pds: failed += 1 print(f" probe-failed uid={uid} {did}: no PDS resolved") continue try: actual = fetch_status(client, pds, did) except httpx.HTTPError as e: failed += 1 print(f" probe-failed uid={uid} {did} @ {pds}: {type(e).__name__}") continue 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()