from __future__ import annotations from datetime import datetime, timezone from zoneinfo import ZoneInfo from croniter import croniter from .db import db, row_to_dict, utc_now from .settings_store import get_settings def validate_cron(expr: str) -> None: if not croniter.is_valid(expr): raise ValueError("Invalid cron expression") def _scheduler_timezone(timezone_name: str | None = None) -> ZoneInfo: return ZoneInfo(timezone_name or get_settings().get("timezone") or "UTC") def next_run(expr: str, base: datetime | None = None, timezone_name: str | None = None) -> str: base = base or datetime.now(timezone.utc) if base.tzinfo is None: base = base.replace(tzinfo=timezone.utc) scheduler_tz = _scheduler_timezone(timezone_name) local_base = base.astimezone(scheduler_tz) local_next = croniter(expr, local_base).get_next(datetime) if local_next.tzinfo is None: local_next = local_next.replace(tzinfo=scheduler_tz) return local_next.astimezone(timezone.utc).replace(microsecond=0).isoformat() def list_jobs() -> list[dict]: with db() as conn: return [dict(row) for row in conn.execute("SELECT * FROM backup_jobs ORDER BY id DESC")] def get_job(job_id: int) -> dict: with db() as conn: row = conn.execute("SELECT * FROM backup_jobs WHERE id = ?", (job_id,)).fetchone() if not row: raise KeyError(f"Backup job {job_id} not found") return dict(row) def create_job(payload: dict) -> dict: validate_cron(payload["cron_schedule"]) now = utc_now() with db() as conn: cur = conn.execute( """ INSERT INTO backup_jobs( guest_vmid, guest_name, guest_type, node, enabled, cron_schedule, proxmox_storage, backup_mode, compression, retention_type, retention_value, next_run_at, created_at, updated_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( payload["guest_vmid"], payload["guest_name"], payload["guest_type"], payload["node"], int(payload.get("enabled", True)), payload["cron_schedule"], payload["proxmox_storage"], payload["backup_mode"], payload["compression"], payload["retention_type"], int(payload["retention_value"]), next_run(payload["cron_schedule"]), now, now, ), ) row = conn.execute("SELECT * FROM backup_jobs WHERE id = ?", (cur.lastrowid,)).fetchone() return dict(row) def update_job(job_id: int, payload: dict) -> dict: existing = get_job(job_id) merged = {**existing, **payload} validate_cron(merged["cron_schedule"]) merged["next_run_at"] = next_run(merged["cron_schedule"]) if payload.get("cron_schedule") else existing["next_run_at"] merged["updated_at"] = utc_now() with db() as conn: conn.execute( """ UPDATE backup_jobs SET guest_vmid = ?, guest_name = ?, guest_type = ?, node = ?, enabled = ?, cron_schedule = ?, proxmox_storage = ?, backup_mode = ?, compression = ?, retention_type = ?, retention_value = ?, next_run_at = ?, updated_at = ? WHERE id = ? """, ( merged["guest_vmid"], merged["guest_name"], merged["guest_type"], merged["node"], int(merged["enabled"]), merged["cron_schedule"], merged["proxmox_storage"], merged["backup_mode"], merged["compression"], merged["retention_type"], int(merged["retention_value"]), merged["next_run_at"], merged["updated_at"], job_id, ), ) row = conn.execute("SELECT * FROM backup_jobs WHERE id = ?", (job_id,)).fetchone() return dict(row) def delete_job(job_id: int) -> None: with db() as conn: conn.execute("DELETE FROM backup_jobs WHERE id = ?", (job_id,)) def recalculate_next_runs() -> None: now = utc_now() with db() as conn: rows = conn.execute("SELECT id, cron_schedule FROM backup_jobs").fetchall() for row in rows: conn.execute( "UPDATE backup_jobs SET next_run_at = ?, updated_at = ? WHERE id = ?", (next_run(row["cron_schedule"]), now, row["id"]), ) def due_jobs(limit: int) -> list[dict]: now = utc_now() with db() as conn: rows = conn.execute( """ SELECT * FROM backup_jobs WHERE enabled = 1 AND next_run_at IS NOT NULL AND next_run_at <= ? ORDER BY next_run_at ASC LIMIT ? """, (now, limit), ).fetchall() return [dict(row) for row in rows] def advance_job(job: dict) -> None: with db() as conn: conn.execute( "UPDATE backup_jobs SET next_run_at = ?, updated_at = ? WHERE id = ?", (next_run(job["cron_schedule"]), utc_now(), job["id"]), )