| """Hugging Face bucket relay for finished Small Cuts scenes. |
| |
| The Space uses this as a read-only scene source. The private engine or a local |
| publisher writes finished scene manifests + media into an HF bucket; the Space |
| downloads those files into a temp cache and serves them through Gradio. |
| """ |
|
|
| from __future__ import annotations |
|
|
| import copy |
| import json |
| import os |
| import tempfile |
| import threading |
| import time |
| from dataclasses import dataclass |
| from datetime import datetime, timezone |
| from pathlib import Path |
| from typing import Any, Protocol |
| from urllib.parse import quote, urlparse |
|
|
| import httpx |
|
|
| RELAY_BUCKET_ENV = "SMALL_CUTS_RELAY_BUCKET" |
| RELAY_PREFIX_ENV = "SMALL_CUTS_RELAY_PREFIX" |
| DEFAULT_RELAY_PREFIX = "relay" |
| RELAY_MANIFEST = "manifest.json" |
| RELAY_CACHE_DIR = Path(tempfile.gettempdir()) / "small-cuts-hf-relay" |
| GRADIO_FILE_ROUTE = "/gradio_api/file=" |
| DEFAULT_SCENE_LIMIT = 60 |
| MEDIA_KEYS = ("frame_url", "card_url", "audio_url", "clip_url") |
| PUBLISH_VISIBILITIES = frozenset({"shared", "public"}) |
| HTTP_TIMEOUT_S = 20.0 |
| MANIFEST_CACHE_TTL_S = 5.0 |
| RELAY_CACHE_MAX_BYTES = 512 * 1024 * 1024 |
|
|
|
|
| class BucketFileSystem(Protocol): |
| def cat(self, path: str) -> bytes: ... |
|
|
|
|
| class BucketRelayError(RuntimeError): |
| """Raised when the bucket relay cannot read or hydrate its manifest.""" |
|
|
|
|
| @dataclass(frozen=True) |
| class RelaySnapshot: |
| path: Path |
| scene_count: int |
| manifest_path: Path |
|
|
|
|
| def gradio_file_url(path: str | Path) -> str: |
| return f"{GRADIO_FILE_ROUTE}{quote(str(path))}" |
|
|
|
|
| def _normalize_prefix(prefix: str) -> str: |
| return prefix.strip().strip("/") |
|
|
|
|
| def _safe_bucket_slug(bucket_id: str) -> str: |
| return bucket_id.replace("/", "__") |
|
|
|
|
| class BucketSceneClient: |
| """Read finished NarratedScene payloads from a Hugging Face bucket manifest.""" |
|
|
| base_url = "" |
| readonly = True |
|
|
| def __init__( |
| self, |
| bucket_id: str, |
| *, |
| prefix: str = DEFAULT_RELAY_PREFIX, |
| fs: BucketFileSystem | None = None, |
| cache_dir: str | Path | None = None, |
| register_static_paths: Any | None = None, |
| manifest_cache_ttl_s: float = MANIFEST_CACHE_TTL_S, |
| cache_max_bytes: int = RELAY_CACHE_MAX_BYTES, |
| ) -> None: |
| self.bucket_id = bucket_id.strip() |
| if not self.bucket_id: |
| raise ValueError("bucket_id is required") |
| self.prefix = _normalize_prefix(prefix) |
| self.root = f"hf://buckets/{self.bucket_id}" |
| if self.prefix: |
| self.root = f"{self.root}/{self.prefix}" |
| self._fs = fs |
| self.cache_dir = ( |
| Path(cache_dir) |
| if cache_dir is not None |
| else (RELAY_CACHE_DIR / _safe_bucket_slug(self.bucket_id)) |
| ) |
| if cache_dir is None and self.prefix: |
| self.cache_dir = self.cache_dir / self.prefix |
| self.cache_dir.mkdir(parents=True, exist_ok=True) |
| if register_static_paths is not None: |
| register_static_paths([self.cache_dir]) |
| self.manifest_cache_ttl_s = manifest_cache_ttl_s |
| self.cache_max_bytes = cache_max_bytes |
| self._manifest_lock = threading.Lock() |
| self._media_lock = threading.Lock() |
| self._manifest_cache: tuple[float, list[dict[str, Any]]] | None = None |
| self._prune_cache() |
|
|
| @property |
| def fs(self) -> BucketFileSystem: |
| if self._fs is None: |
| from huggingface_hub import HfFileSystem |
|
|
| self._fs = HfFileSystem() |
| return self._fs |
|
|
| def list_scenes(self, limit: int = DEFAULT_SCENE_LIMIT) -> list[dict[str, Any]]: |
| with self._manifest_lock: |
| now = time.monotonic() |
| if ( |
| self._manifest_cache is not None |
| and now - self._manifest_cache[0] < self.manifest_cache_ttl_s |
| ): |
| return copy.deepcopy(self._manifest_cache[1][-limit:]) |
| try: |
| raw = self.fs.cat(f"{self.root}/{RELAY_MANIFEST}") |
| except FileNotFoundError: |
| self._manifest_cache = (now, []) |
| return [] |
| try: |
| manifest = json.loads(raw.decode("utf-8")) |
| scenes = manifest.get("scenes", []) |
| if not isinstance(scenes, list): |
| raise ValueError("relay manifest scenes must be a list") |
| hydrated = [] |
| for scene in scenes: |
| try: |
| hydrated.append(self._hydrate_scene(scene)) |
| except FileNotFoundError: |
| continue |
| self._manifest_cache = (now, hydrated) |
| return copy.deepcopy(hydrated[-limit:]) |
| except (OSError, json.JSONDecodeError, KeyError, TypeError, ValueError) as exc: |
| raise BucketRelayError( |
| f"could not read relay bucket {self.bucket_id}: {exc}" |
| ) from exc |
|
|
| def media_url(self, path: str | None) -> str | None: |
| if not path: |
| return None |
| if path.startswith(("http://", "https://", "data:", GRADIO_FILE_ROUTE)): |
| return path |
| relative = self._relative_media_path(path) |
| target = self.cache_dir / relative |
| with self._media_lock: |
| if not target.exists(): |
| target.parent.mkdir(parents=True, exist_ok=True) |
| tmp = target.with_name(f".{target.name}.{os.getpid()}.{threading.get_ident()}.tmp") |
| try: |
| tmp.write_bytes(self.fs.cat(f"{self.root}/{relative.as_posix()}")) |
| tmp.replace(target) |
| finally: |
| tmp.unlink(missing_ok=True) |
| self._prune_cache(protected=target) |
| return gradio_file_url(target) |
|
|
| def _hydrate_scene(self, scene: dict[str, Any]) -> dict[str, Any]: |
| hydrated = copy.deepcopy(scene) |
| media = hydrated.get("media") |
| if not isinstance(media, dict): |
| hydrated["media"] = {} |
| return hydrated |
| for key in MEDIA_KEYS: |
| media[key] = self.media_url(media.get(key)) |
| return hydrated |
|
|
| def _relative_media_path(self, path: str) -> Path: |
| value = path.strip().lstrip("/") |
| if self.prefix and value.startswith(f"{self.prefix}/"): |
| value = value[len(self.prefix) + 1 :] |
| relative = Path(value) |
| if relative.is_absolute() or ".." in relative.parts: |
| raise ValueError(f"unsafe bucket media path: {path}") |
| return relative |
|
|
| def _prune_cache(self, protected: Path | None = None) -> None: |
| if self.cache_max_bytes <= 0 or not self.cache_dir.exists(): |
| return |
| protected_resolved = protected.resolve() if protected is not None else None |
| files = [path for path in self.cache_dir.rglob("*") if path.is_file()] |
| total = sum(path.stat().st_size for path in files) |
| if total <= self.cache_max_bytes: |
| return |
| for path in sorted(files, key=lambda item: item.stat().st_mtime): |
| if protected_resolved is not None and path.resolve() == protected_resolved: |
| continue |
| size = path.stat().st_size |
| path.unlink(missing_ok=True) |
| total -= size |
| if total <= self.cache_max_bytes: |
| break |
|
|
|
|
| def prepare_relay_snapshot( |
| engine_url: str, |
| output_dir: str | Path, |
| *, |
| limit: int = DEFAULT_SCENE_LIMIT, |
| include_private: bool = False, |
| source: str | None = None, |
| client: httpx.Client | None = None, |
| ) -> RelaySnapshot: |
| """Stage a bucket-ready manifest + media snapshot from the private engine.""" |
| base_url = engine_url.rstrip("/") |
| output = Path(output_dir) |
| media_root = output / "media" |
| output.mkdir(parents=True, exist_ok=True) |
| media_root.mkdir(parents=True, exist_ok=True) |
| close_client = client is None |
| http = client or httpx.Client(timeout=HTTP_TIMEOUT_S) |
| try: |
| response = http.get(f"{base_url}/v1/scenes") |
| response.raise_for_status() |
| scenes = response.json().get("scenes", [])[-limit:] |
| published = [ |
| _stage_scene_media(base_url, output, scene, http, source=source) |
| for scene in scenes |
| if _should_publish_scene(scene, include_private=include_private) |
| ] |
| finally: |
| if close_client: |
| http.close() |
| manifest = { |
| "contract_version": "1.1.0", |
| "published_at": datetime.now(timezone.utc).isoformat(), |
| "source_engine": base_url, |
| "scenes": published, |
| } |
| manifest_path = output / RELAY_MANIFEST |
| manifest_path.write_text(json.dumps(manifest, indent=2, sort_keys=True) + "\n") |
| return RelaySnapshot(output, len(published), manifest_path) |
|
|
|
|
| def _should_publish_scene(scene: dict[str, Any], *, include_private: bool) -> bool: |
| if include_private: |
| return True |
| return scene.get("visibility") in PUBLISH_VISIBILITIES |
|
|
|
|
| def _stage_scene_media( |
| engine_url: str, |
| output_dir: Path, |
| scene: dict[str, Any], |
| client: httpx.Client, |
| *, |
| source: str | None = None, |
| ) -> dict[str, Any]: |
| staged = copy.deepcopy(scene) |
| if source: |
| staged["source"] = source |
| staged["source_icon"] = source |
| media = staged.get("media") |
| if not isinstance(media, dict): |
| staged["media"] = {} |
| return staged |
| scene_dir = _safe_path_segment(str(staged.get("scene_id") or "scene")) |
| for key in MEDIA_KEYS: |
| media[key] = _download_media(engine_url, output_dir, scene_dir, media.get(key), client) |
| return staged |
|
|
|
|
| def _download_media( |
| engine_url: str, |
| output_dir: Path, |
| scene_dir: str, |
| url: str | None, |
| client: httpx.Client, |
| ) -> str | None: |
| if not url: |
| return None |
| if url.startswith(("data:", GRADIO_FILE_ROUTE)): |
| return None |
| absolute = url if url.startswith(("http://", "https://")) else f"{engine_url}/{url.lstrip('/')}" |
| relative = _relay_media_path(url, scene_dir) |
| target = output_dir / relative |
| target.parent.mkdir(parents=True, exist_ok=True) |
| response = client.get(absolute) |
| response.raise_for_status() |
| target.write_bytes(response.content) |
| return relative.as_posix() |
|
|
|
|
| def _relay_media_path(url: str, scene_dir: str) -> Path: |
| parsed = urlparse(url) |
| source_path = (parsed.path if parsed.scheme else url.split("?", 1)[0]).lstrip("/") |
| if source_path.startswith("media/"): |
| relative = Path(source_path) |
| else: |
| filename = _safe_path_segment(Path(source_path).name or "media.bin") |
| relative = Path("media") / scene_dir / filename |
| if relative.is_absolute() or ".." in relative.parts: |
| raise ValueError(f"unsafe relay media path: {url}") |
| return relative |
|
|
|
|
| def _safe_path_segment(value: str) -> str: |
| cleaned = "".join(ch if ch.isalnum() or ch in ("-", "_", ".") else "-" for ch in value) |
| return cleaned.strip(".-") or "item" |
|
|