remote-postgres / app /backup.py
imkrish's picture
Whole-cluster backups + on-demand backup/restore endpoints
7d3393c verified
Raw History Blame Contribute Delete
5.22 kB
"""Backup / restore for the in-container Postgres.
Two durable targets, either or both:
* BACKUP_DIR — local dir (durable only if it's on HF persistent /data)
* HF_BACKUP_REPO — a private HF Dataset repo (durable even on the free tier)
A fresh boot calls `restore()` to pull the latest dump back; a background scheduler
calls `backup()` every BACKUP_INTERVAL_MIN minutes.
"""
import os
import sys
import glob
import gzip
import time
import threading
import subprocess
from datetime import datetime
PG_USER = os.environ.get("POSTGRES_USER", "demo")
PG_DB = os.environ.get("POSTGRES_DB", "demo")
PG_PORT = os.environ.get("PGPORT", "5432")
PGHOST = "/tmp"
BACKUP_DIR = os.environ.get("BACKUP_DIR", "/home/appuser/backups")
KEEP = int(os.environ.get("BACKUP_KEEP", "24"))
INTERVAL_MIN = int(os.environ.get("BACKUP_INTERVAL_MIN", "30"))
HF_REPO = os.environ.get("HF_BACKUP_REPO", "").strip() # e.g. "user/pg-backups"
HF_TOKEN = os.environ.get("HF_TOKEN", "").strip()
HF_LATEST = "backups/latest.sql.gz"
# Surfaced on the dashboard.
STATE = {"count": 0, "last": None, "last_name": None, "error": None,
"dir": BACKUP_DIR, "hf": bool(HF_REPO and HF_TOKEN)}
def _ts():
return datetime.now().strftime("%Y%m%d-%H%M%S")
def _hf_api():
from huggingface_hub import HfApi
return HfApi(token=HF_TOKEN)
def ensure_repo():
if not (HF_REPO and HF_TOKEN):
return
try:
_hf_api().create_repo(repo_id=HF_REPO, repo_type="dataset",
private=True, exist_ok=True)
print(f"[backup] HF dataset ready: {HF_REPO}", flush=True)
except Exception as e:
print(f"[backup] create_repo failed: {e}", flush=True)
def backup():
"""Dump the DB to a gzip file locally, rotate, and mirror to HF (latest)."""
os.makedirs(BACKUP_DIR, exist_ok=True)
name = f"backup-{_ts()}.sql.gz"
path = os.path.join(BACKUP_DIR, name)
try:
# whole-cluster dump: every database + roles (not just one DB)
dump = subprocess.run(
["pg_dumpall", "-h", PGHOST, "-p", PG_PORT, "-U", PG_USER],
stdout=subprocess.PIPE, check=True,
).stdout
with gzip.open(path, "wb") as f:
f.write(dump)
# keep only the newest KEEP local dumps
for old in sorted(glob.glob(os.path.join(BACKUP_DIR, "backup-*.sql.gz")))[:-KEEP]:
try:
os.remove(old)
except OSError:
pass
if HF_REPO and HF_TOKEN:
try:
_hf_api().upload_file(
path_or_fileobj=path, path_in_repo=HF_LATEST,
repo_id=HF_REPO, repo_type="dataset",
commit_message=f"backup {name}",
)
except Exception as e:
print(f"[backup] HF upload failed: {e}", flush=True)
STATE.update(count=STATE["count"] + 1, last=datetime.now(),
last_name=name, error=None)
print(f"[backup] wrote {name} ({os.path.getsize(path)} bytes)", flush=True)
return path
except Exception as e:
STATE["error"] = str(e)
print(f"[backup] FAILED: {e}", flush=True)
raise
def restore():
"""Restore the most recent dump (prefers HF, falls back to local). Returns bool."""
src = None
if HF_REPO and HF_TOKEN:
try:
from huggingface_hub import hf_hub_download
src = hf_hub_download(HF_REPO, HF_LATEST, repo_type="dataset",
token=HF_TOKEN)
print(f"[backup] fetched latest dump from HF dataset {HF_REPO}", flush=True)
except Exception as e:
print(f"[backup] no HF backup to restore ({e})", flush=True)
if src is None:
local = sorted(glob.glob(os.path.join(BACKUP_DIR, "backup-*.sql.gz")))
src = local[-1] if local else None
if not src:
print("[backup] no backup found — starting empty", flush=True)
return False
print(f"[backup] restoring from {src}", flush=True)
# connect to 'postgres'; the dumpall script \connects into each database itself
with gzip.open(src, "rb") as f:
subprocess.run(
["psql", "-h", PGHOST, "-p", PG_PORT, "-U", PG_USER,
"-d", "postgres", "-v", "ON_ERROR_STOP=0", "-q"],
input=f.read(), check=True,
)
print("[backup] restore complete", flush=True)
return True
def _loop():
while True:
time.sleep(INTERVAL_MIN * 60)
try:
backup()
except Exception:
pass
def start_scheduler():
"""Called from the web app: ensure the HF repo exists and start periodic backups."""
ensure_repo()
threading.Thread(target=_loop, daemon=True).start()
print(f"[backup] scheduler started — every {INTERVAL_MIN} min, "
f"dir={BACKUP_DIR}, hf={'on' if STATE['hf'] else 'off'}", flush=True)
if __name__ == "__main__":
cmd = sys.argv[1] if len(sys.argv) > 1 else ""
if cmd == "backup":
backup()
elif cmd == "restore":
restore()
elif cmd == "ensure-repo":
ensure_repo()
else:
print("usage: backup.py [backup|restore|ensure-repo]")