jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188#!/usr/bin/env python3"""Prune atcr.io/zat.dev/stream via the PDS records atcr indexes.atcr has no registry-API manifest DELETE (every DELETE returns 405, evencorrectly scoped — recorded 2026-08-09). Images are really deleted byremoving their `io.atcr.tag` / `io.atcr.manifest` records from our own PDS,which is what atcr indexes. That path is not quota-blocked. atcr'sserver-side GC recomputes quota roughly daily, so usage lags deletion by upto ~24h.plan:- inventory: public listRecords over io.atcr.tag + io.atcr.manifest- scope: ONLY repository == "stream" or "stream/cache" — the same collections hold other projects' images (zds), which must never be touched- keep: every tag named by a receipt in receipts/*.json, plus 'latest'; manifests stay if reachable from a kept tag (index children included, resolved through the raw manifest blob)- delete: everything else stream-scoped, via com.atproto.repo.applyWritesauth: ATCR_USER (default zat.dev) + ATCR_APP_PASSWORD, an atproto apppassword — createSession against the account's PDS. Reads are public.usage: scripts/registry-prune [--dry-run]"""import jsonimport osimport pathlibimport sysimport urllib.errorimport urllib.requestHANDLE = os.environ.get("ATCR_USER", "zat.dev")REPOSITORIES = {"stream", "stream/cache"}DRY = "--dry-run" in sys.argvBATCH = 200 # applyWrites capdef get_json(url: str) -> dict: with urllib.request.urlopen(url, timeout=30) as resp: return json.load(resp)def resolve_pds(handle: str) -> tuple[str, str]: did = get_json( "https://public.api.bsky.app/xrpc/com.atproto.identity.resolveHandle" f"?handle={handle}" )["did"] doc = get_json(f"https://plc.directory/{did}") pds = next( s["serviceEndpoint"] for s in doc["service"] if s["id"] == "#atproto_pds" ) return did, pdsdef list_records(pds: str, did: str, collection: str) -> list[dict]: records, cursor = [], None while True: url = ( f"{pds}/xrpc/com.atproto.repo.listRecords" f"?repo={did}&collection={collection}&limit=100" ) if cursor: url += f"&cursor={cursor}" page = get_json(url) records += page["records"] cursor = page.get("cursor") if not cursor or not page["records"]: return recordsdef rkey(uri: str) -> str: return uri.rsplit("/", 1)[1]def bare_digest(d: str) -> str: return d.removeprefix("sha256:")def index_children(pds: str, did: str, record: dict) -> set[str]: """Child manifest digests of an OCI index, from the raw manifest blob.""" blob = record["value"].get("manifestBlob") if not blob: return set() cid = blob["ref"]["$link"] with urllib.request.urlopen( f"{pds}/xrpc/com.atproto.sync.getBlob?did={did}&cid={cid}", timeout=30 ) as resp: manifest = json.load(resp) return {bare_digest(m["digest"]) for m in manifest.get("manifests", [])}def create_session(pds: str) -> dict: password = os.environ.get("ATCR_APP_PASSWORD") or sys.exit( "no ATCR_APP_PASSWORD in env" ) req = urllib.request.Request( f"{pds}/xrpc/com.atproto.server.createSession", data=json.dumps({"identifier": HANDLE, "password": password}).encode(), headers={"Content-Type": "application/json"}, ) with urllib.request.urlopen(req, timeout=30) as resp: return json.load(resp)def apply_deletes(pds: str, session: dict, doomed: list[tuple[str, str]]) -> None: for start in range(0, len(doomed), BATCH): writes = [ { "$type": "com.atproto.repo.applyWrites#delete", "collection": collection, "rkey": key, } for collection, key in doomed[start : start + BATCH] ] req = urllib.request.Request( f"{pds}/xrpc/com.atproto.repo.applyWrites", data=json.dumps({"repo": session["did"], "writes": writes}).encode(), headers={ "Content-Type": "application/json", "Authorization": f"Bearer {session['accessJwt']}", }, ) urllib.request.urlopen(req, timeout=60).close() print(f" deleted {len(writes)} records")def main() -> None: keep_tags = {"latest"} for receipt in pathlib.Path("receipts").glob("*.json"): keep_tags.add(receipt.stem) print(f"keeping tags: {sorted(keep_tags)}") did, pds = resolve_pds(HANDLE) print(f"{HANDLE} = {did} @ {pds}") tags = [ t for t in list_records(pds, did, "io.atcr.tag") if t["value"].get("repository") in REPOSITORIES ] manifests = { rkey(m["uri"]): m for m in list_records(pds, did, "io.atcr.manifest") if m["value"].get("repository") in REPOSITORIES } print(f"stream-scoped: {len(tags)} tags, {len(manifests)} manifests") # closure of manifests reachable from kept tags (cache tags never kept) keep_digests: set[str] = set() doomed_tags: list[dict] = [] for tag in tags: value = tag["value"] kept = value["repository"] == "stream" and value["tag"] in keep_tags if not kept: doomed_tags.append(tag) continue digest = rkey(value["manifest"]) keep_digests.add(digest) if digest in manifests: keep_digests |= index_children(pds, did, manifests[digest]) doomed: list[tuple[str, str]] = [] for tag in doomed_tags: print(f" doomed tag {rkey(tag['uri'])} ({tag['value']['tag']})") doomed.append(("io.atcr.tag", rkey(tag["uri"]))) for digest, record in sorted(manifests.items()): if digest in keep_digests: continue print(f" doomed manifest {digest[:19]}… ({record['value']['repository']})") doomed.append(("io.atcr.manifest", digest)) kept = len(tags) - len(doomed_tags) print(f"kept {kept} tags / {len(keep_digests)} reachable manifests; " f"{len(doomed)} records doomed") if not doomed: print("nothing to prune.") return if DRY: print("dry run: no deletes issued.") return apply_deletes(pds, create_session(pds), doomed) print("done. atcr's server-side GC recomputes quota within ~24h.")if __name__ == "__main__": main()