124 lines
3.1 KiB
Python
Executable File
124 lines
3.1 KiB
Python
Executable File
#!/usr/bin/env python3
|
|
"""Run a PostgreSQL source->dest sync on a cron schedule."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import os
|
|
import subprocess
|
|
import sys
|
|
import time
|
|
from datetime import datetime
|
|
from zoneinfo import ZoneInfo
|
|
|
|
from croniter import croniter
|
|
|
|
CRON_MACROS = {
|
|
"@yearly": "0 0 1 1 *",
|
|
"@annually": "0 0 1 1 *",
|
|
"@monthly": "0 0 1 * *",
|
|
"@weekly": "0 0 * * 0",
|
|
"@daily": "0 0 * * *",
|
|
"@midnight": "0 0 * * *",
|
|
"@hourly": "0 * * * *",
|
|
}
|
|
|
|
SYNC_SCRIPT = os.environ.get("SYNC_SCRIPT", "/usr/local/bin/sync.sh")
|
|
|
|
|
|
def log(msg: str) -> None:
|
|
print(f"[{datetime.now().astimezone().isoformat(timespec='seconds')}] {msg}", file=sys.stderr)
|
|
|
|
|
|
def zone() -> ZoneInfo:
|
|
name = os.environ.get("TZ") or "UTC"
|
|
try:
|
|
return ZoneInfo(name)
|
|
except Exception as exc:
|
|
raise SystemExit(f"Invalid TZ '{name}': {exc}") from exc
|
|
|
|
|
|
def resolve_cron_expr() -> str:
|
|
cron = os.environ.get("CRON", "").strip()
|
|
if not cron:
|
|
raise SystemExit("CRON must be set (e.g. '0 3 * * *')")
|
|
|
|
expr = CRON_MACROS.get(cron.lower(), cron)
|
|
if not croniter.is_valid(expr):
|
|
raise SystemExit(f"Invalid CRON '{expr}' (expected a 5-field cron expression)")
|
|
return expr
|
|
|
|
|
|
def require_env(name: str) -> str:
|
|
value = os.environ.get(name, "").strip()
|
|
if not value:
|
|
raise SystemExit(f"{name} must be set")
|
|
return value
|
|
|
|
|
|
def run_on_startup() -> bool:
|
|
return os.environ.get("ON_STARTUP", "true").lower() in ("1", "true", "yes", "on")
|
|
|
|
|
|
def next_run(expr: str, after: datetime) -> datetime:
|
|
return croniter(expr, after).get_next(datetime)
|
|
|
|
|
|
def redacted(url: str) -> str:
|
|
"""Hide password in postgresql://user:pass@host/db style URLs."""
|
|
if "://" not in url:
|
|
return url
|
|
scheme, rest = url.split("://", 1)
|
|
if "@" not in rest or ":" not in rest.split("@", 1)[0]:
|
|
return url
|
|
creds, hostpart = rest.split("@", 1)
|
|
user = creds.split(":", 1)[0]
|
|
return f"{scheme}://{user}:***@{hostpart}"
|
|
|
|
|
|
def do_sync() -> None:
|
|
log("Running sync...")
|
|
result = subprocess.run([SYNC_SCRIPT], check=False)
|
|
if result.returncode != 0:
|
|
log(f"Sync failed with exit code {result.returncode}")
|
|
else:
|
|
log("Sync finished OK")
|
|
|
|
|
|
def sleep_until(target: datetime, tz: ZoneInfo) -> None:
|
|
while True:
|
|
remaining = (target - datetime.now(tz)).total_seconds()
|
|
if remaining <= 0:
|
|
return
|
|
time.sleep(min(remaining, 60.0))
|
|
|
|
|
|
def main() -> None:
|
|
source = require_env("SOURCE_DB_URL")
|
|
dest = require_env("DEST_DB_URL")
|
|
expr = resolve_cron_expr()
|
|
tz = zone()
|
|
on_startup = run_on_startup()
|
|
|
|
log(
|
|
"Starting database-syncer "
|
|
f"(TZ={tz.key}, CRON='{expr}', "
|
|
f"SOURCE='{redacted(source)}', DEST='{redacted(dest)}')"
|
|
)
|
|
|
|
if on_startup:
|
|
do_sync()
|
|
|
|
now = datetime.now(tz)
|
|
nxt = next_run(expr, now)
|
|
log(f"Next scheduled sync at {nxt.isoformat(timespec='seconds')}")
|
|
|
|
while True:
|
|
sleep_until(nxt, tz)
|
|
do_sync()
|
|
nxt = next_run(expr, datetime.now(tz))
|
|
log(f"Next scheduled sync at {nxt.isoformat(timespec='seconds')}")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|