LiveHouse-TS / space /src /result_feed.py
ziyuzhou02's picture
Deploy GitHub main 3feb6cda1511
e317359 verified
Raw History Blame Contribute Delete
2.28 kB
"""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