feat: punkt 6 (backup) — ARQ cron job, backup config in settings, backup-now trigger, backup history, migration 0130
This commit is contained in:
@@ -0,0 +1,238 @@
|
||||
"""ARQ backup job — scheduled database backup via scripts/backup.py.
|
||||
|
||||
Reads backup configuration from system settings, executes backup.py as a
|
||||
subprocess, logs the result to audit_log, and notifies admins on failure.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
import uuid
|
||||
from typing import Any
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Path to the backup script relative to project root
|
||||
_BACKUP_SCRIPT = os.path.join(
|
||||
os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))),
|
||||
"scripts",
|
||||
"backup.py",
|
||||
)
|
||||
|
||||
|
||||
async def _get_backup_config() -> dict[str, Any]:
|
||||
"""Read backup configuration from system settings for all tenants.
|
||||
|
||||
Returns the first tenant's settings that have backup_enabled=True,
|
||||
or defaults if no settings exist.
|
||||
"""
|
||||
from sqlalchemy import select as sa_select
|
||||
|
||||
from app.core.db import get_worker_session_factory
|
||||
from app.models.system_settings import SystemSettings
|
||||
|
||||
factory = get_worker_session_factory()
|
||||
async with factory() as db:
|
||||
result = await db.execute(
|
||||
sa_select(SystemSettings).where(
|
||||
SystemSettings.backup_enabled.is_(True),
|
||||
SystemSettings.deleted_at.is_(None),
|
||||
).limit(1)
|
||||
)
|
||||
settings = result.scalar_one_or_none()
|
||||
|
||||
if settings is None:
|
||||
return {
|
||||
"backup_enabled": False,
|
||||
"backup_interval": "daily",
|
||||
"backup_retention_days": 7,
|
||||
"backup_destination": "local",
|
||||
"tenant_id": None,
|
||||
}
|
||||
|
||||
return {
|
||||
"backup_enabled": True,
|
||||
"backup_interval": settings.backup_interval,
|
||||
"backup_retention_days": settings.backup_retention_days,
|
||||
"backup_destination": settings.backup_destination,
|
||||
"tenant_id": settings.tenant_id,
|
||||
}
|
||||
|
||||
|
||||
async def run_backup_job(ctx: dict[str, Any]) -> dict[str, Any]:
|
||||
"""Execute a scheduled backup by calling scripts/backup.py as a subprocess.
|
||||
|
||||
Reads backup configuration from system settings. If backup_enabled is
|
||||
False, the job is silently skipped.
|
||||
|
||||
Returns a dict with keys: success (bool), message (str), backup_id (str|None).
|
||||
"""
|
||||
config = await _get_backup_config()
|
||||
|
||||
if not config["backup_enabled"]:
|
||||
logger.debug("Backup job skipped — backup_enabled is False")
|
||||
return {"success": False, "message": "Backup disabled", "backup_id": None}
|
||||
|
||||
tenant_id = config.get("tenant_id")
|
||||
retention_days = config.get("backup_retention_days", 7)
|
||||
destination = config.get("backup_destination", "local")
|
||||
|
||||
logger.info(
|
||||
"Starting scheduled backup: destination=%s, retention=%dd",
|
||||
destination,
|
||||
retention_days,
|
||||
)
|
||||
|
||||
# Build subprocess command
|
||||
cmd = [
|
||||
sys.executable,
|
||||
_BACKUP_SCRIPT,
|
||||
"--destination",
|
||||
destination,
|
||||
"--retention-days",
|
||||
str(retention_days),
|
||||
]
|
||||
|
||||
# Pass environment with DATABASE_URL
|
||||
env = os.environ.copy()
|
||||
|
||||
try:
|
||||
process = await asyncio.create_subprocess_exec(
|
||||
*cmd,
|
||||
env=env,
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE,
|
||||
)
|
||||
stdout, stderr = await process.communicate()
|
||||
|
||||
success = process.returncode == 0
|
||||
output = stdout.decode() if stdout else ""
|
||||
error = stderr.decode() if stderr else ""
|
||||
|
||||
if success:
|
||||
logger.info("Scheduled backup completed successfully: %s", output[-500:] if output else "")
|
||||
else:
|
||||
logger.error("Scheduled backup failed (exit %d): %s", process.returncode, error)
|
||||
|
||||
# Log to audit_log
|
||||
await _log_backup_result(
|
||||
tenant_id=tenant_id,
|
||||
success=success,
|
||||
output=output,
|
||||
error=error,
|
||||
destination=destination,
|
||||
retention_days=retention_days,
|
||||
)
|
||||
|
||||
# Notify admin on failure
|
||||
if not success and tenant_id:
|
||||
await _notify_admin_failure(tenant_id, error)
|
||||
|
||||
return {
|
||||
"success": success,
|
||||
"message": "Backup completed" if success else f"Backup failed: {error[:200]}",
|
||||
"backup_id": None,
|
||||
}
|
||||
|
||||
except Exception as exc:
|
||||
logger.exception("Backup job encountered an exception")
|
||||
if tenant_id:
|
||||
await _log_backup_result(
|
||||
tenant_id=tenant_id,
|
||||
success=False,
|
||||
output="",
|
||||
error=str(exc),
|
||||
destination=destination,
|
||||
retention_days=retention_days,
|
||||
)
|
||||
await _notify_admin_failure(tenant_id, str(exc))
|
||||
return {"success": False, "message": str(exc), "backup_id": None}
|
||||
|
||||
|
||||
async def _log_backup_result(
|
||||
tenant_id: uuid.UUID | None,
|
||||
success: bool,
|
||||
output: str,
|
||||
error: str,
|
||||
destination: str,
|
||||
retention_days: int,
|
||||
) -> None:
|
||||
"""Write backup result to audit_log."""
|
||||
from app.core.audit import log_audit
|
||||
from app.core.db import get_worker_session_factory
|
||||
|
||||
if tenant_id is None:
|
||||
return
|
||||
|
||||
factory = get_worker_session_factory()
|
||||
async with factory() as db:
|
||||
try:
|
||||
await log_audit(
|
||||
db,
|
||||
tenant_id=tenant_id,
|
||||
user_id=None,
|
||||
action="backup_success" if success else "backup_failed",
|
||||
entity_type="backup",
|
||||
entity_id=None,
|
||||
changes={
|
||||
"success": success,
|
||||
"destination": destination,
|
||||
"retention_days": retention_days,
|
||||
"output": output[-1000:] if output else "",
|
||||
"error": error[-1000:] if error else "",
|
||||
},
|
||||
)
|
||||
await db.commit()
|
||||
except Exception:
|
||||
logger.exception("Failed to write backup audit log")
|
||||
await db.rollback()
|
||||
|
||||
|
||||
async def _notify_admin_failure(tenant_id: uuid.UUID, error: str) -> None:
|
||||
"""Send a notification to admin users about backup failure."""
|
||||
from sqlalchemy import select as sa_select
|
||||
|
||||
from app.core.db import get_worker_session_factory
|
||||
from app.core.notifications import post_system_message
|
||||
from app.models.user import User
|
||||
|
||||
factory = get_worker_session_factory()
|
||||
async with factory() as db:
|
||||
try:
|
||||
# Find system admin users for this tenant
|
||||
result = await db.execute(
|
||||
sa_select(User).where(
|
||||
User.tenant_id == tenant_id,
|
||||
User.is_system_admin.is_(True),
|
||||
User.deleted_at.is_(None),
|
||||
).limit(1)
|
||||
)
|
||||
admin_user = result.scalar_one_or_none()
|
||||
|
||||
if admin_user is None:
|
||||
logger.warning("No admin user found to notify about backup failure")
|
||||
return
|
||||
|
||||
await post_system_message(
|
||||
db,
|
||||
tenant_id=tenant_id,
|
||||
user_id=admin_user.id,
|
||||
message_type="backup_failed",
|
||||
title="Backup fehlgeschlagen",
|
||||
body=f"Das geplante Backup ist fehlgeschlagen: {error[:500]}",
|
||||
entity_type="backup",
|
||||
severity="error",
|
||||
)
|
||||
await db.commit()
|
||||
except Exception:
|
||||
logger.exception("Failed to notify admin about backup failure")
|
||||
await db.rollback()
|
||||
|
||||
|
||||
# Register with the job registry
|
||||
from app.core.job_registry import register_job # noqa: E402
|
||||
|
||||
register_job("run_backup", run_backup_job)
|
||||
+6
-1
@@ -236,7 +236,7 @@ def _lazy_register_plugin_jobs() -> None:
|
||||
if not registry.list_discovered():
|
||||
registry.discover_builtins()
|
||||
|
||||
job_modules: list[str] = ["app.core.jobs", "app.services.import_export_jobs"]
|
||||
job_modules: list[str] = ["app.core.jobs", "app.core.backup_job", "app.services.import_export_jobs"]
|
||||
for plugin_name in registry.list_discovered():
|
||||
plugin = registry.get_plugin(plugin_name)
|
||||
if plugin is None:
|
||||
@@ -371,4 +371,9 @@ class WorkerSettings:
|
||||
_wrap_cron_with_lock("cleanup_outbox", cleanup_outbox_job, ttl_seconds=300),
|
||||
minute=0,
|
||||
),
|
||||
# Scheduled backup — daily at 02:00 (guarded by distributed lock)
|
||||
cron(
|
||||
_wrap_cron_with_lock("run_backup", get_job("run_backup"), ttl_seconds=600),
|
||||
hour=2, minute=0,
|
||||
),
|
||||
]
|
||||
|
||||
Reference in New Issue
Block a user