153 lines
5.2 KiB
Python
153 lines
5.2 KiB
Python
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"]),
|
|
)
|