Spaces:
Running
Running
File size: 2,275 Bytes
e317359 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 | """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
|