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 from app.schemas.schemas import JobCreate, JobUpdate, JobResponse, AgentJobRegister, PurgeOrphansRequest 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.""" count_res = await db.execute(select(func.count(BackupJob.id))) job_count = count_res.scalar() or 0 job_code = f"JOB-{job_count + 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, 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") count_res = await db.execute(select(func.count(BackupJob.id))) job_count = count_res.scalar() or 0 job_code = f"JOB-{job_count + 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, 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 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." }