"""Read atomic public result revisions from a Dataset without rebuilding the Space.""" import json import logging import os import shutil from pathlib import Path import threading import time class ResultFeed: def __init__(self, fallback, repo_id=None, cache_dir=None): self.current = str(fallback) self.repo_id = repo_id or os.getenv("HF_RESULTS_REPO") self.cache_dir = Path(cache_dir or Path(fallback).parent / ".result-cache") self.revision = None self.checked_at = 0.0 self.lock = threading.Lock() def refresh(self): if not self.repo_id: return self.current with self.lock: if time.monotonic() - self.checked_at < 30: return self.current self.checked_at = time.monotonic() try: from huggingface_hub import HfApi, snapshot_download revision = HfApi(token=False).dataset_info(self.repo_id, timeout=15).sha if revision == self.revision: return self.current destination = self.cache_dir / revision snapshot_download(self.repo_id, repo_type="dataset", revision=revision, allow_patterns=["results/**"], local_dir=destination, token=False, max_workers=4) results = destination / "results" status = json.loads((results / "online_status.json").read_text()) if status.get("aggregate_status") != "ok" or not list(results.glob("*/all_results.csv")): raise ValueError("Incomplete result revision") # Readers receive a complete immutable directory, never half an upload. self.current = str(results) self.revision = revision # Leave recent revisions long enough for in-flight UI callbacks. for old in self.cache_dir.iterdir(): if old.is_dir() and old != destination and time.time() - old.stat().st_mtime > 600: shutil.rmtree(old) except Exception: logging.exception("Result feed unavailable; retaining the last complete revision") return self.current