# retoor import asyncio import hashlib import io import shutil import sqlite3 import subprocess import tarfile from datetime import datetime, timedelta, timezone from pathlib import Path import brotli import zstandard from molodetz import config from molodetz.database import ( admin_uids_ordered, delete_backup_rows, enqueue_backup, list_backups, now_iso, parse_iso, pending_backups, shard_path, update_backup, ) from molodetz.services.base import BaseService, ConfigField from molodetz.utils.audit import record_system from molodetz.utils.notifications import create_notification TARGETS = ("database", "uploads", "keys", "full") def database_file(): prefix = "sqlite:///" if not config.DATABASE_URL.startswith(prefix): raise RuntimeError("backups support SQLite databases only") return Path(config.DATABASE_URL[len(prefix):]) def _snapshot_database(staging): target = staging / "molodetz.db" source = sqlite3.connect(str(database_file())) try: destination = sqlite3.connect(str(target)) try: source.backup(destination) finally: destination.close() finally: source.close() return target def build_archive(uid, target, codec="zstd"): if target not in TARGETS: raise ValueError(f"unknown backup target {target}") staging = config.DATA_PATHS["backup_staging"] / uid staging.mkdir(parents=True, exist_ok=True) try: buffer = io.BytesIO() with tarfile.open(fileobj=buffer, mode="w") as archive: if target in ("database", "full"): archive.add(_snapshot_database(staging), arcname="molodetz.db") if target in ("uploads", "full") and config.DATA_PATHS["uploads"].exists(): archive.add(config.DATA_PATHS["uploads"], arcname="uploads") if target in ("keys", "full") and config.DATA_PATHS["keys"].exists(): archive.add(config.DATA_PATHS["keys"], arcname="keys") raw = buffer.getvalue() if codec == "brotli": data = brotli.compress(raw, quality=6) suffix = "tar.br" else: data = zstandard.ZstdCompressor(level=10).compress(raw) suffix = "tar.zst" checksum = hashlib.sha256(data).hexdigest()[:16] directory = f"{shard_path(uid)}" stored_name = f"{target}-{uid.replace('-', '')}-{checksum}.{suffix}" path = config.DATA_PATHS["backups"] / directory / stored_name path.parent.mkdir(parents=True, exist_ok=True) path.write_bytes(data) return directory, stored_name, len(data) finally: shutil.rmtree(staging, ignore_errors=True) def backup_path(row): root = config.DATA_PATHS["backups"].resolve() path = (root / row["directory"] / row["stored_name"]).resolve() if root not in path.parents: raise ValueError("backup path escapes the backup root") return path def prune_backups(keep_last, dry_run=False): done = [row for row in list_backups(limit=10000) if row["status"] == "done"] victims = done[keep_last:] if not dry_run: for row in victims: try: backup_path(row).unlink(missing_ok=True) except ValueError: pass delete_backup_rows([row["uid"] for row in victims]) return len(victims) def clear_backups(dry_run=False): rows = list_backups(limit=100000) if not dry_run: for row in rows: if row.get("directory") and row.get("stored_name"): try: backup_path(row).unlink(missing_ok=True) except ValueError: pass delete_backup_rows([row["uid"] for row in rows]) return len(rows) def offload(path): remote = config.BACKUP_OFFLOAD_REMOTE if not remote or shutil.which(config.RCLONE_BIN) is None: return "skipped" result = subprocess.run( [config.RCLONE_BIN, "--config", config.RCLONE_CONFIG, "copy", str(path), remote], capture_output=True, text=True, timeout=900, ) return "ok" if result.returncode == 0 else f"failed: {result.stderr.strip()[:200]}" def finish_backup(row, archive=None, error=None): if error is not None: update_backup(row["uid"], status="failed", error=str(error)[:500], finished_at=now_iso()) record_system("backup.failed", actor_kind="service", origin="service", result="error", message=str(error)[:200]) return False directory, stored_name, size = archive update_backup(row["uid"], status="done", directory=directory, stored_name=stored_name, size_bytes=size, finished_at=now_iso()) record_system("backup.finished", actor_kind="service", origin="service", payload={"uid": row["uid"], "size": size}) return True def process_backup(row, codec): try: archive = build_archive(row["uid"], row["target"], codec) except Exception as exc: return finish_backup(row, error=exc) return finish_backup(row, archive=archive) async def process_backup_async(row, codec): try: archive = await asyncio.to_thread(build_archive, row["uid"], row["target"], codec) except Exception as exc: return finish_backup(row, error=exc) return finish_backup(row, archive=archive) class BackupService(BaseService): name = "backup" title = "Backups" description = "Works through the backup queue, schedules automatic backups and prunes old archives." default_interval = 10 min_interval = 5 metrics_interval = 30 config_fields = [ ConfigField("schedule_hours", "int", 24, "Automatic every N hours (0 = off)", minimum=0, maximum=720, group="Schedule"), ConfigField("schedule_target", "select", "database", "Automatic target", options=list(TARGETS), group="Schedule"), ConfigField("keep_last", "int", 14, "Keep last N", minimum=1, maximum=1000, group="Retention"), ConfigField("codec", "select", "zstd", "Compression", options=["zstd", "brotli"], group="Archive"), ConfigField("offload", "bool", "0", "Offload via rclone", group="Archive"), ] def __init__(self): super().__init__() self.processed = 0 def _schedule_due(self, hours): if hours <= 0: return False rows = list_backups(limit=1) if not rows: return True last = parse_iso(rows[0]["created_at"]) return last is None or datetime.now(timezone.utc) - last >= timedelta(hours=hours) async def run_once(self): cfg = self.get_config() if self._schedule_due(int(cfg["schedule_hours"])): enqueue_backup(cfg["schedule_target"], "system") self.log(f"automatic backup scheduled: {cfg['schedule_target']}") for row in pending_backups(): update_backup(row["uid"], status="running") ok = await process_backup_async(row, cfg["codec"]) self.processed += 1 self.log(f"backup {row['target']} {'finished' if ok else 'failed'}") primary = admin_uids_ordered()[:1] for uid in primary: create_notification(uid, "backup.finished" if ok else "backup.failed", f"Backup {row['target']} {'finished' if ok else 'failed'}", "/admin/backups") if ok and cfg["offload"] == "1": done = [item for item in list_backups(limit=5) if item["uid"] == row["uid"]] if done: self.log(f"offload: {offload(backup_path(done[0]))}") pruned = prune_backups(int(cfg["keep_last"])) if pruned: self.log(f"{pruned} old backups pruned") def collect_metrics(self): rows = list_backups(limit=200) return { "cards": [ {"label": "Archives", "value": sum(1 for row in rows if row["status"] == "done")}, {"label": "Queue", "value": sum(1 for row in rows if row["status"] == "pending")}, {"label": "Processed since start", "value": self.processed}, ], "table": None, }