Session-gated MinIO save sync with AppImage GUI, CLI edit/session flow, and Gitea release helper. Co-authored-by: Cursor <[email protected]>
486 lines
17 KiB
Python
486 lines
17 KiB
Python
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,
|
|
)
|