from __future__ import annotations import json import shutil from abc import ABC, abstractmethod from dataclasses import dataclass from datetime import datetime, timezone from pathlib import Path from typing import Any from syncgames.errors import StoreError from syncgames.hashutil import sha256_bytes def utc_now() -> datetime: return datetime.now(timezone.utc) def iso_ts(dt: datetime | None = None) -> str: d = dt or utc_now() return d.strftime("%Y%m%dT%H%M%SZ") def parse_expires(s: str | None) -> datetime | None: if not s: return None if s.endswith("Z"): s = s[:-1] + "+00:00" return datetime.fromisoformat(s) @dataclass class Meta: raw: dict[str, Any] @property def live_hash(self) -> str | None: return self.raw.get("live_hash") @property def lease(self) -> dict[str, Any] | None: return self.raw.get("lease") @property def versions_to_keep(self) -> int: return int(self.raw.get("versions_to_keep", 5)) def to_json(self) -> bytes: return (json.dumps(self.raw, indent=2, sort_keys=True) + "\n").encode("utf-8") class ObjectStore(ABC): @abstractmethod def get_meta(self, game_id: str) -> Meta | None: ... @abstractmethod def put_meta(self, game_id: str, meta: Meta) -> None: ... @abstractmethod def list_live(self, game_id: str) -> list[str]: ... @abstractmethod def download_live(self, game_id: str, dest_dir: Path) -> dict[str, str]: ... @abstractmethod def upload_prefix( self, game_id: str, kind: str, prefix_extra: str, local_dir: Path ) -> dict[str, str]: """kind is 'live' or 'history'; prefix_extra for history is 'device/ts'.""" ... @abstractmethod def delete_live_not_in(self, game_id: str, keep_rels: set[str]) -> None: ... @abstractmethod def list_history(self, game_id: str) -> list[str]: """Return list of 'device/ts' prefixes newest-first-ish.""" ... @abstractmethod def copy_history_to_live(self, game_id: str, device_ts: str) -> None: ... @abstractmethod def prune_device_history(self, game_id: str, device_id: str, keep: int) -> int: ... @abstractmethod def retire_game(self, game_id: str) -> str: ... @abstractmethod def probe(self) -> str: ... @abstractmethod def export_ssot(self, out_dir: Path) -> None: ... def empty_meta(game_id: str, versions_to_keep: int = 5) -> Meta: return Meta( { "schema": 1, "game_id": game_id, "live_hash": None, "file_checksums": {}, "lease": None, "versions_to_keep": versions_to_keep, "updated_at": utc_now().strftime("%Y-%m-%dT%H:%M:%SZ"), "updated_by": None, } ) class MinioStore(ObjectStore): def __init__( self, *, endpoint_url: str, bucket: str, access_key: str, secret_key: str, region: str = "us-east-1", path_style: bool = True, ) -> None: import boto3 from botocore.client import Config if not endpoint_url or not access_key or not secret_key: raise StoreError( "MinIO endpoint_url, access_key, and secret_key are required", key="config_error", ) self.bucket = bucket cfg = Config( signature_version="s3v4", s3={"addressing_style": "path" if path_style else "auto"}, connect_timeout=10, read_timeout=60, retries={"max_attempts": 2, "mode": "standard"}, ) self.s3 = boto3.client( "s3", endpoint_url=endpoint_url.rstrip("/"), aws_access_key_id=access_key, aws_secret_access_key=secret_key, region_name=region, config=cfg, ) def _key(self, *parts: str) -> str: return "/".join(p.strip("/") for p in parts if p) def get_meta(self, game_id: str) -> Meta | None: key = self._key("games", game_id, "meta.json") try: obj = self.s3.get_object(Bucket=self.bucket, Key=key) data = json.loads(obj["Body"].read().decode("utf-8")) return Meta(data) except self.s3.exceptions.NoSuchKey: return None except Exception as e: # noqa: BLE001 # boto3 often raises ClientError code = getattr(e, "response", {}).get("Error", {}).get("Code", "") if code in {"404", "NoSuchKey", "NotFound"}: return None raise StoreError(f"get_meta failed: {e}") from e def put_meta(self, game_id: str, meta: Meta) -> None: key = self._key("games", game_id, "meta.json") try: self.s3.put_object( Bucket=self.bucket, Key=key, Body=meta.to_json(), ContentType="application/json", ) except Exception as e: # noqa: BLE001 raise StoreError(f"put_meta failed: {e}") from e def list_live(self, game_id: str) -> list[str]: prefix = self._key("games", game_id, "live") + "/" return self._list_rels(prefix) def _list_rels(self, prefix: str) -> list[str]: rels: list[str] = [] token = None while True: kw: dict[str, Any] = {"Bucket": self.bucket, "Prefix": prefix} if token: kw["ContinuationToken"] = token resp = self.s3.list_objects_v2(**kw) for item in resp.get("Contents") or []: key = item["Key"] if key.endswith("/"): continue rels.append(key[len(prefix) :]) if not resp.get("IsTruncated"): break token = resp.get("NextContinuationToken") return rels def download_live(self, game_id: str, dest_dir: Path) -> dict[str, str]: from syncgames.hashutil import sha256_file prefix = self._key("games", game_id, "live") + "/" dest_dir.mkdir(parents=True, exist_ok=True) checksums: dict[str, str] = {} for rel in self.list_live(game_id): target = dest_dir / rel target.parent.mkdir(parents=True, exist_ok=True) self.s3.download_file(self.bucket, prefix + rel, str(target)) checksums[rel] = sha256_file(target) return checksums def upload_prefix( self, game_id: str, kind: str, prefix_extra: str, local_dir: Path ) -> dict[str, str]: from syncgames.hashutil import sha256_file if kind == "live": base = self._key("games", game_id, "live") elif kind == "history": base = self._key("games", game_id, "history", prefix_extra) else: raise StoreError(f"Unknown kind {kind}") checksums: dict[str, str] = {} for path in sorted(p for p in local_dir.rglob("*") if p.is_file()): rel = path.relative_to(local_dir).as_posix() key = f"{base}/{rel}" self.s3.upload_file(str(path), self.bucket, key) checksums[rel] = sha256_file(path) return checksums def delete_live_not_in(self, game_id: str, keep_rels: set[str]) -> None: prefix = self._key("games", game_id, "live") + "/" for rel in self.list_live(game_id): if rel not in keep_rels: self.s3.delete_object(Bucket=self.bucket, Key=prefix + rel) def list_history(self, game_id: str) -> list[str]: prefix = self._key("games", game_id, "history") + "/" devices_ts: set[str] = set() token = None while True: kw: dict[str, Any] = {"Bucket": self.bucket, "Prefix": prefix} if token: kw["ContinuationToken"] = token resp = self.s3.list_objects_v2(**kw) for item in resp.get("Contents") or []: key = item["Key"][len(prefix) :] parts = key.split("/") if len(parts) >= 2: devices_ts.add(f"{parts[0]}/{parts[1]}") if not resp.get("IsTruncated"): break token = resp.get("NextContinuationToken") return sorted(devices_ts, reverse=True) def copy_history_to_live(self, game_id: str, device_ts: str) -> None: src_prefix = self._key("games", game_id, "history", device_ts) + "/" # Clear live then copy for rel in self.list_live(game_id): self.s3.delete_object( Bucket=self.bucket, Key=self._key("games", game_id, "live", rel) ) for rel in self._list_rels(src_prefix): src = src_prefix + rel dst = self._key("games", game_id, "live", rel) self.s3.copy_object( Bucket=self.bucket, CopySource={"Bucket": self.bucket, "Key": src}, Key=dst, ) def prune_device_history(self, game_id: str, device_id: str, keep: int) -> int: prefix = self._key("games", game_id, "history", device_id) + "/" stamps: set[str] = set() for rel in self._list_rels(prefix): stamps.add(rel.split("/", 1)[0]) ordered = sorted(stamps, reverse=True) removed = 0 for stamp in ordered[keep:]: for rel in self._list_rels(prefix + stamp + "/"): self.s3.delete_object(Bucket=self.bucket, Key=prefix + stamp + "/" + rel) removed += 1 return removed def retire_game(self, game_id: str) -> str: date = utc_now().strftime("%Y%m%d") dest = f"retired/{game_id}-{date}" src_prefix = self._key("games", game_id) + "/" for rel in self._list_rels(src_prefix): src = src_prefix + rel self.s3.copy_object( Bucket=self.bucket, CopySource={"Bucket": self.bucket, "Key": src}, Key=f"{dest}/{rel}", ) self.s3.delete_object(Bucket=self.bucket, Key=src) # meta.json meta_key = self._key("games", game_id, "meta.json") try: self.s3.copy_object( Bucket=self.bucket, CopySource={"Bucket": self.bucket, "Key": meta_key}, Key=f"{dest}/meta.json", ) self.s3.delete_object(Bucket=self.bucket, Key=meta_key) except Exception: # noqa: BLE001 pass return dest def probe(self) -> str: key = self._key("games", "_probe", f"ping-{iso_ts()}.txt") body = b"syncgames-probe\n" self.s3.put_object(Bucket=self.bucket, Key=key, Body=body) got = self.s3.get_object(Bucket=self.bucket, Key=key)["Body"].read() self.s3.delete_object(Bucket=self.bucket, Key=key) if got != body: raise StoreError("probe mismatch — possible NGINX buffering corruption") return sha256_bytes(body) def export_ssot(self, out_dir: Path) -> None: out_dir.mkdir(parents=True, exist_ok=True) token = None while True: kw: dict[str, Any] = {"Bucket": self.bucket, "Prefix": "games/"} if token: kw["ContinuationToken"] = token resp = self.s3.list_objects_v2(**kw) for item in resp.get("Contents") or []: key = item["Key"] target = out_dir / key target.parent.mkdir(parents=True, exist_ok=True) self.s3.download_file(self.bucket, key, str(target)) if not resp.get("IsTruncated"): break token = resp.get("NextContinuationToken") class FilesystemStore(ObjectStore): """Path 1 fallback: SSOT as a directory tree.""" def __init__(self, root: Path) -> None: if not root: raise StoreError("ssot_root required for filesystem store", key="config_error") self.root = root self.root.mkdir(parents=True, exist_ok=True) def _game(self, game_id: str) -> Path: return self.root / "games" / game_id def get_meta(self, game_id: str) -> Meta | None: p = self._game(game_id) / "meta.json" if not p.exists(): return None return Meta(json.loads(p.read_text(encoding="utf-8"))) def put_meta(self, game_id: str, meta: Meta) -> None: p = self._game(game_id) / "meta.json" p.parent.mkdir(parents=True, exist_ok=True) p.write_bytes(meta.to_json()) def list_live(self, game_id: str) -> list[str]: live = self._game(game_id) / "live" if not live.exists(): return [] return [p.relative_to(live).as_posix() for p in live.rglob("*") if p.is_file()] def download_live(self, game_id: str, dest_dir: Path) -> dict[str, str]: from syncgames.hashutil import sha256_file live = self._game(game_id) / "live" dest_dir.mkdir(parents=True, exist_ok=True) checksums: dict[str, str] = {} if not live.exists(): return checksums for path in live.rglob("*"): if not path.is_file(): continue rel = path.relative_to(live).as_posix() target = dest_dir / rel target.parent.mkdir(parents=True, exist_ok=True) shutil.copy2(path, target) checksums[rel] = sha256_file(target) return checksums def upload_prefix( self, game_id: str, kind: str, prefix_extra: str, local_dir: Path ) -> dict[str, str]: from syncgames.hashutil import sha256_file if kind == "live": dest = self._game(game_id) / "live" else: dest = self._game(game_id) / "history" / Path(prefix_extra) dest.mkdir(parents=True, exist_ok=True) checksums: dict[str, str] = {} for path in local_dir.rglob("*"): if not path.is_file(): continue rel = path.relative_to(local_dir).as_posix() target = dest / rel target.parent.mkdir(parents=True, exist_ok=True) shutil.copy2(path, target) checksums[rel] = sha256_file(path) return checksums def delete_live_not_in(self, game_id: str, keep_rels: set[str]) -> None: live = self._game(game_id) / "live" if not live.exists(): return for path in list(live.rglob("*")): if path.is_file(): rel = path.relative_to(live).as_posix() if rel not in keep_rels: path.unlink() def list_history(self, game_id: str) -> list[str]: hist = self._game(game_id) / "history" if not hist.exists(): return [] out: list[str] = [] for device in hist.iterdir(): if not device.is_dir(): continue for ts in device.iterdir(): if ts.is_dir(): out.append(f"{device.name}/{ts.name}") return sorted(out, reverse=True) def copy_history_to_live(self, game_id: str, device_ts: str) -> None: src = self._game(game_id) / "history" / Path(device_ts) live = self._game(game_id) / "live" if live.exists(): shutil.rmtree(live) shutil.copytree(src, live) def prune_device_history(self, game_id: str, device_id: str, keep: int) -> int: base = self._game(game_id) / "history" / device_id if not base.exists(): return 0 stamps = sorted([p for p in base.iterdir() if p.is_dir()], key=lambda p: p.name, reverse=True) removed = 0 for p in stamps[keep:]: shutil.rmtree(p) removed += 1 return removed def retire_game(self, game_id: str) -> str: date = utc_now().strftime("%Y%m%d") dest = self.root / "retired" / f"{game_id}-{date}" dest.parent.mkdir(parents=True, exist_ok=True) src = self._game(game_id) if src.exists(): shutil.move(str(src), str(dest)) return str(dest.relative_to(self.root)) def probe(self) -> str: p = self.root / "games" / "_probe" / "ping.txt" p.parent.mkdir(parents=True, exist_ok=True) body = b"syncgames-probe\n" p.write_bytes(body) got = p.read_bytes() p.unlink() if got != body: raise StoreError("filesystem probe mismatch") return sha256_bytes(body) def export_ssot(self, out_dir: Path) -> None: if out_dir.resolve() == self.root.resolve(): return if out_dir.exists(): shutil.rmtree(out_dir) shutil.copytree(self.root, out_dir) def build_store(agent) -> ObjectStore: # AgentConfig if agent.store == "filesystem": return FilesystemStore(agent.ssot_root) # type: ignore[arg-type] return MinioStore( endpoint_url=agent.endpoint_url, bucket=agent.bucket, access_key=agent.access_key, secret_key=agent.secret_key, region=agent.region, path_style=agent.path_style, )