Files

408 lines
14 KiB
Python

from typing import List, Optional
from datetime import datetime, timezone
from fastapi import APIRouter, Depends, HTTPException, status
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, func, desc
from app.core.database import get_db
from app.models.models import User, Client, BackupJob, BackupFile, JobRun
from app.schemas.schemas import (
JobCreate, JobUpdate, JobResponse, AgentJobRegister, PurgeOrphansRequest,
JobRunResponse, JobRunStartRequest, JobRunCompleteRequest
)
from app.api.deps import get_current_user, require_admin, get_current_client
from app.services.event_service import log_event
from app.ws.manager import ws_manager
import logging
logger = logging.getLogger("uvicorn.error")
router = APIRouter(prefix="/jobs", tags=["Backup Jobs"])
@router.get("", response_model=List[JobResponse])
async def list_jobs(
client_id: Optional[int] = None,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user)
):
query = select(BackupJob)
if client_id:
query = query.where(BackupJob.client_id == client_id)
query = query.order_by(desc(BackupJob.created_at))
result = await db.execute(query)
return result.scalars().all()
@router.get("/agent/assigned", response_model=List[JobResponse])
async def get_agent_jobs(
db: AsyncSession = Depends(get_db),
current_client: Client = Depends(get_current_client)
):
"""Called by the Windows Agent to query its active backup jobs."""
query = select(BackupJob).where(BackupJob.client_id == current_client.id, BackupJob.is_active == True)
result = await db.execute(query)
return result.scalars().all()
@router.post("/agent/register", response_model=JobResponse)
async def agent_register_job(
payload: AgentJobRegister,
db: AsyncSession = Depends(get_db),
current_client: Client = Depends(get_current_client)
):
"""Called by the Windows Agent to register a new local job/folder on the server."""
max_id_res = await db.execute(select(func.max(BackupJob.id)))
max_id = max_id_res.scalar() or 0
job_code = f"JOB-{max_id + 1:03d}"
job = BackupJob(
job_code=job_code,
client_id=current_client.id,
name=payload.name,
source_path=payload.source_path,
file_patterns=payload.file_patterns,
exclude_patterns=payload.exclude_patterns,
skip_hidden=payload.skip_hidden,
skip_system=payload.skip_system,
skip_readonly=payload.skip_readonly,
sync_deletions=payload.sync_deletions,
local_destination_path=payload.local_destination_path,
schedule_cron=payload.schedule_cron,
keep_daily=payload.keep_daily,
keep_weekly=payload.keep_weekly,
keep_monthly=payload.keep_monthly,
min_stable_time_seconds=payload.min_stable_time_seconds,
status="IDLE",
is_active=True
)
db.add(job)
await db.commit()
await db.refresh(job)
await log_event(
db=db,
event_type="JOB_CREATED",
message=f"Agent registered backup job '{job.name}' ({job.job_code}) from device.",
severity="INFO",
client_id=current_client.id,
job_id=job.id
)
return job
@router.delete("/agent/{job_id}")
async def agent_delete_job(
job_id: int,
db: AsyncSession = Depends(get_db),
current_client: Client = Depends(get_current_client)
):
"""Called by the Windows Agent to delete a backup job it owns."""
res = await db.execute(select(BackupJob).where(BackupJob.id == job_id, BackupJob.client_id == current_client.id))
job = res.scalar_one_or_none()
if not job:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="Job not found or not owned by this client"
)
await db.delete(job)
await db.commit()
await log_event(
db=db,
event_type="JOB_DELETED",
message=f"Agent deleted backup job {job.job_code} from device.",
severity="WARNING",
client_id=current_client.id
)
return {"message": "Job successfully deleted by agent"}
@router.post("", response_model=JobResponse)
async def create_job(
payload: JobCreate,
db: AsyncSession = Depends(get_db),
admin_user: User = Depends(require_admin)
):
# Verify client exists
client_res = await db.execute(select(Client).where(Client.id == payload.client_id))
client = client_res.scalar_one_or_none()
if not client:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Client not found")
max_id_res = await db.execute(select(func.max(BackupJob.id)))
max_id = max_id_res.scalar() or 0
job_code = f"JOB-{max_id + 1:03d}"
job = BackupJob(
job_code=job_code,
client_id=payload.client_id,
name=payload.name,
source_path=payload.source_path,
file_patterns=payload.file_patterns,
exclude_patterns=payload.exclude_patterns,
skip_hidden=payload.skip_hidden,
skip_system=payload.skip_system,
skip_readonly=payload.skip_readonly,
sync_deletions=payload.sync_deletions,
local_destination_path=payload.local_destination_path,
schedule_cron=payload.schedule_cron,
keep_daily=payload.keep_daily,
keep_weekly=payload.keep_weekly,
keep_monthly=payload.keep_monthly,
min_stable_time_seconds=payload.min_stable_time_seconds,
status="IDLE",
is_active=True
)
db.add(job)
await db.commit()
await db.refresh(job)
await log_event(
db=db,
event_type="JOB_CREATED",
message=f"Created backup job '{job.name}' ({job.job_code}) for client {client.name}.",
severity="INFO",
client_id=client.id,
job_id=job.id,
user_email=admin_user.email
)
return job
@router.get("/{job_id}", response_model=JobResponse)
async def get_job(
job_id: int,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user)
):
res = await db.execute(select(BackupJob).where(BackupJob.id == job_id))
job = res.scalar_one_or_none()
if not job:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Job not found")
return job
@router.put("/{job_id}", response_model=JobResponse)
async def update_job(
job_id: int,
payload: JobUpdate,
db: AsyncSession = Depends(get_db),
admin_user: User = Depends(require_admin)
):
res = await db.execute(select(BackupJob).where(BackupJob.id == job_id))
job = res.scalar_one_or_none()
if not job:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Job not found")
for field, value in payload.model_dump(exclude_unset=True).items():
setattr(job, field, value)
await db.commit()
await db.refresh(job)
await log_event(
db=db,
event_type="CONFIG_CHANGED",
message=f"Updated backup job '{job.name}' ({job.job_code}).",
severity="INFO",
client_id=job.client_id,
job_id=job.id,
user_email=admin_user.email
)
return job
@router.delete("/{job_id}")
async def delete_job(
job_id: int,
db: AsyncSession = Depends(get_db),
admin_user: User = Depends(require_admin)
):
res = await db.execute(select(BackupJob).where(BackupJob.id == job_id))
job = res.scalar_one_or_none()
if not job:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Job not found")
await db.delete(job)
await db.commit()
await log_event(
db=db,
event_type="JOB_DELETED",
message=f"Deleted backup job {job.job_code}.",
severity="WARNING",
client_id=job.client_id,
user_email=admin_user.email
)
return {"message": f"Job {job.job_code} deleted"}
@router.post("/{job_id}/trigger")
async def trigger_job(
job_id: int,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user)
):
"""Notifies the connected agent to start executing this backup job immediately."""
res = await db.execute(select(BackupJob).where(BackupJob.id == job_id))
job = res.scalar_one_or_none()
if not job:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Job not found")
job.status = "QUEUED"
await db.commit()
# Broadcast event so the agent / web UI knows the job was triggered
await ws_manager.broadcast("JOB_TRIGGERED", {
"job_id": job.id,
"job_code": job.job_code,
"client_id": job.client_id,
"timestamp": datetime.now(timezone.utc).isoformat()
})
return {"message": f"Job {job.job_code} triggered"}
@router.post("/{job_id}/purge-orphans")
async def purge_orphans(
job_id: int,
payload: PurgeOrphansRequest,
db: AsyncSession = Depends(get_db),
current_client: Client = Depends(get_current_client)
):
"""
Called by the Windows Agent to delete backup files on the server that no longer exist
in the client's local directory (orphans), enforcing replica mode.
"""
res = await db.execute(select(BackupJob).where(BackupJob.id == job_id, BackupJob.client_id == current_client.id))
job = res.scalar_one_or_none()
if not job:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail="Job not found or not owned by this client"
)
if not job.sync_deletions:
return {"purged_count": 0, "message": "Replica mode (sync_deletions) is not enabled for this job."}
files_res = await db.execute(
select(BackupFile)
.where(
BackupFile.job_id == job_id,
BackupFile.client_id == current_client.id,
BackupFile.is_active == True
)
)
db_files = files_res.scalars().all()
client_paths = {p.strip().replace("\\", "/").lower() for p in payload.active_relative_paths}
purged_count = 0
total_freed_bytes = 0
from app.storage.local import storage_provider
for db_file in db_files:
normalized_db_path = db_file.client_relative_path.strip().replace("\\", "/").lower()
if normalized_db_path not in client_paths:
try:
await storage_provider.delete_backup_file(db_file.relative_path)
db_file.is_active = False
total_freed_bytes += db_file.file_size
purged_count += 1
await log_event(
db=db,
event_type="FILE_PURGED",
message=f"Purged orphan file '{db_file.filename}' from server storage (Replica Mode).",
severity="WARNING",
client_id=current_client.id,
job_id=job.id
)
except Exception as e:
logger.error(f"Failed to purge orphan file {db_file.filename}: {e}")
if purged_count > 0:
current_client.storage_used_bytes = max(0, current_client.storage_used_bytes - total_freed_bytes)
await db.commit()
return {
"purged_count": purged_count,
"freed_bytes": total_freed_bytes,
"message": f"Successfully purged {purged_count} orphan files from server."
}
@router.post("/{job_id}/runs/start", response_model=JobRunResponse)
async def start_job_run(
job_id: int,
payload: JobRunStartRequest,
db: AsyncSession = Depends(get_db),
current_client: Client = Depends(get_current_client)
):
"""Called by the Windows Agent to register the start of a scheduled or manual job run."""
res = await db.execute(select(BackupJob).where(BackupJob.id == job_id, BackupJob.client_id == current_client.id))
job = res.scalar_one_or_none()
if not job:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Job not found or not owned by client")
job.status = "RUNNING"
job.last_run_at = datetime.now(timezone.utc)
run = JobRun(
job_id=job.id,
client_id=current_client.id,
started_at=datetime.now(timezone.utc),
status="RUNNING"
)
db.add(run)
await db.commit()
await db.refresh(run)
return run
@router.post("/{job_id}/runs/{run_id}/complete", response_model=JobRunResponse)
async def complete_job_run(
job_id: int,
run_id: int,
payload: JobRunCompleteRequest,
db: AsyncSession = Depends(get_db),
current_client: Client = Depends(get_current_client)
):
"""Called by the Windows Agent to record completion metrics for a run."""
res = await db.execute(select(JobRun).where(JobRun.id == run_id, JobRun.job_id == job_id, JobRun.client_id == current_client.id))
run = res.scalar_one_or_none()
if not run:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Job run not found")
job_res = await db.execute(select(BackupJob).where(BackupJob.id == job_id))
job = job_res.scalar_one_or_none()
run.completed_at = datetime.now(timezone.utc)
run.status = payload.status
run.files_scanned = payload.files_scanned
run.files_copied = payload.files_copied
run.files_skipped = payload.files_skipped
run.errors_count = payload.errors_count
run.bytes_transferred = payload.bytes_transferred
run.error_summary = payload.error_summary
if job:
job.status = payload.status
job.last_run_at = run.completed_at
await db.commit()
await db.refresh(run)
return run
@router.get("/{job_id}/runs", response_model=List[JobRunResponse])
async def list_job_runs(
job_id: int,
db: AsyncSession = Depends(get_db),
current_user: User = Depends(get_current_user)
):
"""Called by the Web UI to retrieve execution history for a job."""
result = await db.execute(
select(JobRun)
.where(JobRun.job_id == job_id)
.order_by(desc(JobRun.started_at))
.limit(50)
)
return result.scalars().all()