atproto relay in zig zlay.waow.tech
relay zig atproto
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148# /// 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 andare stuck with status='active' AND upstream_status != 'active'. for each,this resolves the DID's current PDS from plc.directory (host_id in the dbmay itself be stale for accounts that never committed since migrating —do not trust it), asks that PDS via getRepoStatus, and writes the answerback. 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 argparseimport osimport sysimport time
import httpximport 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()