220 lines
8.0 KiB
Python
220 lines
8.0 KiB
Python
# retoor <retoor@molodetz.nl>
|
|
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,
|
|
}
|