- BackupSchedule model + schedule_service: next-run calc (hourly/daily/weekly, UTC), background scheduler loop (lifespan), run-one with retention pruning (keep newest N per stack on the destination), backup_failed notify event. - routers/schedules.py: schedules CRUD + run-now; registered in main.py. - Frontend: api/schedules.ts + Settings → Scheduled backups (list with next/last run + status, enable/disable, run-now, delete; add form with stack/destination/ frequency/time/weekday/retention/volumes). Rough-verified only (per request): py_compile, frontend tsc build, app import (95 routes), next-run math sanity. Full live run to be tested after deploy. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
157 lines
5.4 KiB
Python
157 lines
5.4 KiB
Python
"""Scheduled (recurring) stack backups to a destination.
|
|
|
|
A background loop wakes every minute, runs any schedules whose ``next_run`` has
|
|
passed, uploads the backup to the destination, prunes old backups per the
|
|
retention setting, and records status + the next run time.
|
|
|
|
Times are stored as naive UTC to match SQLite's datetime handling.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
import os
|
|
from datetime import datetime, timedelta, timezone
|
|
|
|
from sqlmodel import Session, select
|
|
|
|
from database import engine
|
|
from models.backup_destination import BackupDestination
|
|
from models.backup_schedule import BackupSchedule
|
|
from models.setting import EVENT_BACKUP_FAILED
|
|
from models.stack import Stack
|
|
from services import (
|
|
backup_destination_service as dest_service,
|
|
backup_service,
|
|
compose_service,
|
|
notify_service,
|
|
)
|
|
|
|
logger = logging.getLogger("stackpilot.schedule")
|
|
|
|
|
|
def _now() -> datetime:
|
|
return datetime.now(timezone.utc).replace(tzinfo=None)
|
|
|
|
|
|
def compute_next_run(
|
|
frequency: str, hour: int, minute: int, weekday: int, after: datetime
|
|
) -> datetime:
|
|
"""Next run strictly after ``after`` (naive UTC)."""
|
|
if frequency == "hourly":
|
|
nxt = after.replace(minute=minute % 60, second=0, microsecond=0)
|
|
if nxt <= after:
|
|
nxt += timedelta(hours=1)
|
|
return nxt
|
|
|
|
base = after.replace(hour=hour % 24, minute=minute % 60, second=0, microsecond=0)
|
|
if frequency == "weekly":
|
|
days_ahead = (weekday - base.weekday()) % 7
|
|
cand = base + timedelta(days=days_ahead)
|
|
if cand <= after:
|
|
cand += timedelta(days=7)
|
|
return cand
|
|
|
|
# daily (default)
|
|
if base <= after:
|
|
base += timedelta(days=1)
|
|
return base
|
|
|
|
|
|
def _backup_filename(stack_id: str, include_volumes: bool) -> str:
|
|
date = compose_service.now().strftime("%Y%m%d-%H%M%S")
|
|
suffix = "full" if include_volumes else "config"
|
|
return f"backup-{stack_id}-{suffix}-{date}.tar.gz"
|
|
|
|
|
|
def _prune(dest: BackupDestination, stack_id: str, keep: int) -> int:
|
|
if keep <= 0:
|
|
return 0
|
|
items = dest_service.list_backups(dest)
|
|
mine = [i for i in items if i["name"].startswith(f"backup-{stack_id}-")]
|
|
mine.sort(key=lambda x: x.get("modified") or 0, reverse=True)
|
|
removed = 0
|
|
for old in mine[keep:]:
|
|
try:
|
|
dest_service.delete(dest, old["name"])
|
|
removed += 1
|
|
except dest_service.DestinationError as exc:
|
|
logger.warning("retention delete failed for %s: %s", old["name"], exc)
|
|
return removed
|
|
|
|
|
|
async def run_schedule(session: Session, schedule: BackupSchedule) -> dict:
|
|
"""Execute one schedule now. Updates status + next_run. Returns a summary."""
|
|
now = _now()
|
|
stack = session.get(Stack, schedule.stack_id)
|
|
dest = session.get(BackupDestination, schedule.destination_id)
|
|
result: dict = {"ok": False}
|
|
path = None
|
|
try:
|
|
if not stack:
|
|
raise RuntimeError(f"stack '{schedule.stack_id}' not found")
|
|
if not dest:
|
|
raise RuntimeError(f"destination {schedule.destination_id} not found")
|
|
path = await backup_service.create_backup(
|
|
schedule.stack_id, stack.name,
|
|
include_volumes=schedule.include_volumes, stop_first=schedule.stop_first,
|
|
)
|
|
filename = _backup_filename(schedule.stack_id, schedule.include_volumes)
|
|
await asyncio.to_thread(dest_service.upload, dest, path, filename)
|
|
pruned = await asyncio.to_thread(_prune, dest, schedule.stack_id, schedule.keep)
|
|
schedule.last_status = "ok"
|
|
result = {"ok": True, "name": filename, "destination": dest.name, "pruned": pruned}
|
|
logger.info("Scheduled backup %s → %s ok (pruned %d)", schedule.stack_id, dest.name, pruned)
|
|
except Exception as exc: # noqa: BLE001
|
|
schedule.last_status = f"error: {exc}"[:300]
|
|
result = {"ok": False, "error": str(exc)}
|
|
logger.warning("Scheduled backup %s failed: %s", schedule.stack_id, exc)
|
|
try:
|
|
await notify_service.notify(
|
|
EVENT_BACKUP_FAILED,
|
|
f"Scheduled backup of '{schedule.stack_id}' failed",
|
|
str(exc),
|
|
session,
|
|
)
|
|
except Exception: # noqa: BLE001
|
|
pass
|
|
finally:
|
|
if path and os.path.exists(path):
|
|
os.unlink(path)
|
|
schedule.last_run = now
|
|
schedule.next_run = compute_next_run(
|
|
schedule.frequency, schedule.hour, schedule.minute, schedule.weekday, now
|
|
)
|
|
session.add(schedule)
|
|
session.commit()
|
|
session.refresh(schedule)
|
|
return result
|
|
|
|
|
|
async def run_due() -> None:
|
|
now = _now()
|
|
with Session(engine) as session:
|
|
schedules = session.exec(
|
|
select(BackupSchedule).where(BackupSchedule.enabled == True) # noqa: E712
|
|
).all()
|
|
due = []
|
|
for s in schedules:
|
|
if s.next_run is None:
|
|
s.next_run = compute_next_run(s.frequency, s.hour, s.minute, s.weekday, now)
|
|
session.add(s)
|
|
elif s.next_run <= now:
|
|
due.append(s)
|
|
session.commit()
|
|
for s in due:
|
|
await run_schedule(session, s)
|
|
|
|
|
|
async def scheduler_loop() -> None:
|
|
await asyncio.sleep(20) # let startup settle
|
|
while True:
|
|
try:
|
|
await run_due()
|
|
except Exception as exc: # noqa: BLE001
|
|
logger.warning("scheduler tick failed: %s", exc)
|
|
await asyncio.sleep(60)
|