139 lines
5.1 KiB
Python
139 lines
5.1 KiB
Python
import os
|
|
import hashlib
|
|
import shutil
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Optional
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
from sqlalchemy import select
|
|
from app.models import Job, JobExecution, Backup
|
|
from backup.adapters import get_adapter
|
|
from backup.retention import RetentionPolicy
|
|
|
|
class BackupEngine:
|
|
def __init__(self, db: AsyncSession):
|
|
self.db = db
|
|
|
|
async def execute_job(self, job_id: int, triggered_by: str = "manual") -> JobExecution:
|
|
# Create execution record
|
|
execution = JobExecution(
|
|
job_id=job_id,
|
|
status="pending",
|
|
triggered_by=triggered_by
|
|
)
|
|
self.db.add(execution)
|
|
await self.db.commit()
|
|
await self.db.refresh(execution)
|
|
|
|
try:
|
|
# Load job with source
|
|
result = await self.db.execute(
|
|
select(Job).where(Job.id == job_id)
|
|
)
|
|
job = result.scalar_one()
|
|
|
|
# Update status to running
|
|
execution.status = "running"
|
|
execution.started_at = datetime.now(timezone.utc)
|
|
await self.db.commit()
|
|
|
|
# Determine strategy
|
|
strategy = job.strategy
|
|
parent_backup_id = None
|
|
|
|
if strategy == "incremental":
|
|
# Find last successful full backup
|
|
result = await self.db.execute(
|
|
select(Backup)
|
|
.join(JobExecution)
|
|
.where(
|
|
JobExecution.job_id == job_id,
|
|
JobExecution.status == "success",
|
|
Backup.type == "full"
|
|
)
|
|
.order_by(Backup.created_at.desc())
|
|
)
|
|
last_full = result.scalar_one_or_none()
|
|
|
|
if last_full:
|
|
parent_backup_id = last_full.id
|
|
else:
|
|
# No full backup exists, do full instead
|
|
strategy = "full"
|
|
|
|
# Create backup directory
|
|
timestamp = datetime.now(timezone.utc).strftime("%Y-%m-%d_%H%M%S")
|
|
backup_dir = Path(job.destination_path) / str(job_id) / f"{timestamp}_{strategy}"
|
|
backup_dir.mkdir(parents=True, exist_ok=True)
|
|
|
|
# Get adapter and connect
|
|
adapter = get_adapter(job.source.type, job.source.config)
|
|
await adapter.connect()
|
|
|
|
try:
|
|
# Copy files
|
|
total_processed = 0
|
|
total_backed_up = 0
|
|
|
|
source_path = Path(job.source.config.get("path", "."))
|
|
|
|
for item in source_path.rglob("*"):
|
|
if item.is_file():
|
|
rel_path = item.relative_to(source_path)
|
|
dest_path = backup_dir / "data" / rel_path
|
|
dest_path.parent.mkdir(parents=True, exist_ok=True)
|
|
|
|
# Copy file
|
|
shutil.copy2(item, dest_path)
|
|
|
|
size = item.stat().st_size
|
|
total_processed += size
|
|
total_backed_up += size
|
|
|
|
# Calculate checksum
|
|
checksum = await self._calculate_checksum(backup_dir)
|
|
|
|
# Create backup record
|
|
backup = Backup(
|
|
execution_id=execution.id,
|
|
storage_path=str(backup_dir),
|
|
size_bytes=total_backed_up,
|
|
checksum=checksum,
|
|
type=strategy,
|
|
parent_backup_id=parent_backup_id
|
|
)
|
|
self.db.add(backup)
|
|
|
|
# Update execution
|
|
execution.status = "success"
|
|
execution.completed_at = datetime.now(timezone.utc)
|
|
execution.bytes_processed = total_processed
|
|
execution.bytes_backed_up = total_backed_up
|
|
|
|
# Apply retention policy
|
|
retention = RetentionPolicy(self.db)
|
|
keep_count = getattr(job, 'retention_count', None)
|
|
keep_days = getattr(job, 'retention_days', None)
|
|
if keep_count or keep_days:
|
|
await retention.apply_retention_for_job(job_id, keep_count, keep_days)
|
|
|
|
finally:
|
|
await adapter.disconnect()
|
|
|
|
except Exception as e:
|
|
execution.status = "failed"
|
|
execution.completed_at = datetime.now(timezone.utc)
|
|
execution.error_message = str(e)
|
|
|
|
await self.db.commit()
|
|
return execution
|
|
|
|
async def _calculate_checksum(self, path: Path) -> str:
|
|
hasher = hashlib.sha256()
|
|
for item in sorted(path.rglob("*")):
|
|
if item.is_file():
|
|
with open(item, "rb") as f:
|
|
while chunk := f.read(8192):
|
|
hasher.update(chunk)
|
|
return hasher.hexdigest()
|