Spaces:
Paused
Paused
| #!/usr/bin/env python3 | |
| """Interior Edit Study Space: gated annotation app with a global work queue. | |
| Assignment model: the 1,100 pairs form one shared queue. GET /next?rater=R | |
| returns the first pair that (a) has no active judgment from anyone and (b) is | |
| not currently leased to another rater (10-minute leases prevent two | |
| simultaneous raters from getting the same pair). POST /submit records a | |
| judgment (marking the pair done) or a retract from undo (reopening it). | |
| State survives restarts: at boot the app replays all judgment files from the | |
| dataset repo plus local ones; a huggingface_hub CommitScheduler pushes the | |
| local judgments folder to the dataset every 2 minutes. Files are suffixed | |
| with the boot id so restarts never overwrite earlier uploads. | |
| Access is gated by the STUDY_CODE secret (cookie set via POST /gate). | |
| Set DISABLE_SYNC=1 for local testing without the dataset. | |
| """ | |
| from __future__ import annotations | |
| import json | |
| import os | |
| import re | |
| import threading | |
| import time | |
| from functools import partial | |
| from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer | |
| from pathlib import Path | |
| from urllib.parse import parse_qs, urlparse | |
| PORT = int(os.environ.get("PORT", "7860")) | |
| STUDY_CODE = os.environ.get("STUDY_CODE", "") # empty = no gate | |
| DISABLE_SYNC = os.environ.get("DISABLE_SYNC", "") == "1" | |
| BUNDLE = Path(os.environ.get("BUNDLE_DIR", "bundle")) | |
| BOOT = time.strftime("%Y%m%d-%H%M%S") | |
| JUDGMENTS = Path("judgments") | |
| JUDGMENTS.mkdir(exist_ok=True) | |
| LEASE_SECONDS = 600 | |
| LOCK = threading.Lock() | |
| MANIFEST = json.loads((BUNDLE / "manifest.json").read_text()) | |
| PAIRS = MANIFEST["pairs"] | |
| PAIR_BY_ID = {p["pair_id"]: p for p in PAIRS} | |
| judged: set[str] = set() # pair_ids with an active judgment (any rater) | |
| leases: dict[str, tuple[str, float]] = {} # pair_id -> (rater, expiry) | |
| scheduler = None | |
| if not DISABLE_SYNC: | |
| from huggingface_hub import CommitScheduler | |
| scheduler = CommitScheduler( | |
| repo_id=os.environ["JUDGMENTS_DATASET"], | |
| repo_type="dataset", | |
| folder_path=str(JUDGMENTS), | |
| path_in_repo="judgments", | |
| every=2, # minutes | |
| private=True, | |
| ) | |
| def replay_history() -> None: | |
| """Rebuild the judged set from dataset history + local files.""" | |
| files: list[Path] = [] | |
| if not DISABLE_SYNC: | |
| try: | |
| from huggingface_hub import snapshot_download | |
| seed = snapshot_download( | |
| os.environ["JUDGMENTS_DATASET"], repo_type="dataset", | |
| allow_patterns=["judgments/*"], local_dir="seed", | |
| ) | |
| files += sorted(Path(seed).glob("judgments/*.jsonl")) | |
| except Exception as e: # noqa: BLE001 - empty dataset on first boot | |
| print(f"no history to replay ({type(e).__name__})") | |
| files += sorted(JUDGMENTS.glob("*.jsonl")) | |
| records = [] | |
| for path in files: | |
| for line in path.read_text().splitlines(): | |
| if line.strip(): | |
| records.append(json.loads(line)) | |
| records.sort(key=lambda r: str(r.get("ts", ""))) | |
| active: dict[tuple, bool] = {} | |
| for rec in records: | |
| key = (rec.get("rater"), rec.get("pair_id")) | |
| active[key] = not rec.get("retract") | |
| for (_, pair_id), is_active in active.items(): | |
| if is_active and pair_id in PAIR_BY_ID: | |
| judged.add(pair_id) | |
| print(f"replayed {len(files)} files: {len(judged)}/{len(PAIRS)} pairs already judged") | |
| def next_for(rater: str, lease: bool = True) -> dict: | |
| """Serve the first available pair; also lease the following one to the same | |
| rater and return it as `prefetch` so the client can preload its images.""" | |
| now = time.time() | |
| with LOCK: | |
| for pid in list(leases): | |
| if leases[pid][1] < now: | |
| del leases[pid] | |
| picked = [] | |
| for p in PAIRS: | |
| pid = p["pair_id"] | |
| if pid in judged: | |
| continue | |
| held = leases.get(pid) | |
| if held and held[0] != rater: | |
| continue | |
| if not lease: | |
| return {"judged": len(judged), "total": len(PAIRS)} | |
| leases[pid] = (rater, now + LEASE_SECONDS) | |
| picked.append(p) | |
| if len(picked) == 2: | |
| break | |
| counts = {"judged": len(judged), "total": len(PAIRS)} | |
| if not picked: | |
| return {"done": True, **counts} | |
| out = {"pair": picked[0], **counts} | |
| if len(picked) == 2: | |
| out["prefetch"] = picked[1] | |
| return out | |
| def record_judgment(rec: dict) -> None: | |
| rater = re.sub(r"[^A-Za-z0-9_-]", "_", str(rec.get("rater", "anon")))[:64] or "anon" | |
| line = json.dumps(rec, ensure_ascii=False) | |
| pid = rec.get("pair_id") | |
| with LOCK: | |
| if rec.get("retract"): | |
| judged.discard(pid) | |
| if pid in PAIR_BY_ID: # hand it straight back to the undoing rater | |
| leases[pid] = (rec.get("rater", ""), time.time() + LEASE_SECONDS) | |
| elif pid in PAIR_BY_ID: | |
| judged.add(pid) | |
| leases.pop(pid, None) | |
| if scheduler: | |
| with scheduler.lock: | |
| with open(JUDGMENTS / f"{rater}-{BOOT}.jsonl", "a") as f: | |
| f.write(line + "\n") | |
| else: | |
| with open(JUDGMENTS / f"{rater}-{BOOT}.jsonl", "a") as f: | |
| f.write(line + "\n") | |
| class Handler(SimpleHTTPRequestHandler): | |
| """Static files are open; /next (non-peek) and /submit require ?code=. | |
| The gate is deliberately cookie-free: the Space runs inside an iframe on | |
| huggingface.co, where third-party cookies are blocked, so any Set-Cookie | |
| flow silently loops. The annotation tool sends the code with every API | |
| call instead. | |
| """ | |
| def _send_json(self, obj: dict, status: int = 200): | |
| body = json.dumps(obj, ensure_ascii=False).encode() | |
| self.send_response(status) | |
| self.send_header("Content-Type", "application/json") | |
| self.send_header("Cache-Control", "no-store") | |
| self.send_header("Content-Length", str(len(body))) | |
| self.end_headers() | |
| self.wfile.write(body) | |
| def _code_ok(query: dict) -> bool: | |
| return not STUDY_CODE or query.get("code", [""])[0].strip() == STUDY_CODE | |
| def end_headers(self): | |
| if self.path.startswith("/images/"): # content-addressed enough: never edited in place | |
| self.send_header("Cache-Control", "public, max-age=604800, immutable") | |
| super().end_headers() | |
| def do_GET(self): | |
| parsed = urlparse(self.path) | |
| if parsed.path.rstrip("/") == "/next": | |
| query = parse_qs(parsed.query) | |
| if "peek" in query: # open probe: counts only, no lease, no pair | |
| self._send_json(next_for("", lease=False) | {"gated": bool(STUDY_CODE)}) | |
| return | |
| if not self._code_ok(query): | |
| self._send_json({"error": "wrong code"}, status=403) | |
| return | |
| rater = query.get("rater", ["anon"])[0] | |
| self._send_json(next_for(rater)) | |
| return | |
| super().do_GET() | |
| def do_POST(self): | |
| parsed = urlparse(self.path) | |
| if parsed.path.rstrip("/") != "/submit": | |
| self.send_error(404) | |
| return | |
| if not self._code_ok(parse_qs(parsed.query)): | |
| self._send_json({"error": "wrong code"}, status=403) | |
| return | |
| try: | |
| length = int(self.headers.get("Content-Length", 0)) | |
| record_judgment(json.loads(self.rfile.read(length))) | |
| self._send_json({"ok": True}) | |
| except Exception as e: # noqa: BLE001 - report to the client, keep serving | |
| self.send_error(400, str(e)) | |
| def log_message(self, fmt, *fmt_args): | |
| if self.command == "POST": | |
| super().log_message(fmt, *fmt_args) | |
| def main() -> None: | |
| replay_history() | |
| handler = partial(Handler, directory=str(BUNDLE)) | |
| server = ThreadingHTTPServer(("0.0.0.0", PORT), handler) | |
| print(f"serving on :{PORT} (boot {BOOT}, sync={'off' if DISABLE_SYNC else 'on'})") | |
| server.serve_forever() | |
| if __name__ == "__main__": | |
| main() | |