File size: 5,220 Bytes
8003382
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
7d3393c
8003382
7d3393c
8003382
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
7d3393c
8003382
 
 
7d3393c
8003382
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
"""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]")