cc4c3fcecb
- Backend: Customer/Satellite Models, customer_id auf Server/Job/Audit - Satellite-API: heartbeat, poll (atomares Claiming), logs, result, scan-result, health-report - Auth via X-Api-Key (SHA-256 gehasht) - Job-Queue: pending/claimed/running/success/failed + Stale-Janitor - Batch-Trigger: ein Job pro Server, Satellite arbeitet sequenziell ab - Credentials bleiben lokal: nur symbolische credential_ref zentral - Neues Paket satellite/: Pull-Loop, WinRM/SSH/CAU/Scanner, PyInstaller-tauglich - Frontend: Kunden-Switcher, Satelliten-View, Polling statt WebSocket - Entfernt: WebSocket/Socket.io, Redis, zentrale Credentials, JobRunner - Docs: README/AGENTS/PROMPT auf neue Architektur aktualisiert
218 lines
7.0 KiB
Python
218 lines
7.0 KiB
Python
"""Update job routes: trigger (single + batch), list, logs, cancel.
|
|
|
|
Triggering only queues a job - a satellite of that customer picks it up
|
|
on its next poll and executes it locally.
|
|
"""
|
|
|
|
import json
|
|
|
|
from fastapi import APIRouter, Depends, Query, Request
|
|
from sqlalchemy import func, select
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from app.api.deps import client_ip, get_current_user
|
|
from app.core.database import get_db
|
|
from app.core.exceptions import JobNotCancellableError, NotFoundError, ValidationError
|
|
from app.models.customer import Customer
|
|
from app.models.server import Server
|
|
from app.models.update_job import JobStatus, JobType, UpdateJob, UpdateLog
|
|
from app.models.user import User
|
|
from app.schemas.update import (
|
|
BatchJobTriggerRequest,
|
|
BatchJobTriggerResponse,
|
|
JobTriggerRequest,
|
|
UpdateJobRead,
|
|
UpdateLogRead,
|
|
)
|
|
from app.services.audit import AuditService
|
|
|
|
router = APIRouter()
|
|
|
|
|
|
@router.post("/trigger", response_model=UpdateJobRead, status_code=201)
|
|
async def trigger_update(
|
|
payload: JobTriggerRequest,
|
|
request: Request,
|
|
db: AsyncSession = Depends(get_db),
|
|
user: User = Depends(get_current_user),
|
|
) -> UpdateJob:
|
|
job = await _create_job(db, payload, user.username)
|
|
await AuditService(db).log(
|
|
username=user.username,
|
|
action="update.trigger",
|
|
target=f"job:{job.id}",
|
|
customer_id=payload.customer_id,
|
|
details={"type": payload.type.value, "server_id": payload.server_id},
|
|
ip_address=client_ip(request),
|
|
)
|
|
return job
|
|
|
|
|
|
@router.post("/trigger-batch", response_model=BatchJobTriggerResponse, status_code=201)
|
|
async def trigger_batch(
|
|
payload: BatchJobTriggerRequest,
|
|
request: Request,
|
|
db: AsyncSession = Depends(get_db),
|
|
user: User = Depends(get_current_user),
|
|
) -> BatchJobTriggerResponse:
|
|
"""Queue one job per server - the satellite works through them in order."""
|
|
stmt = select(Server).where(Server.customer_id == payload.customer_id)
|
|
if payload.server_ids:
|
|
stmt = stmt.where(Server.id.in_(payload.server_ids))
|
|
result = await db.execute(stmt.order_by(Server.name))
|
|
servers = list(result.scalars().all())
|
|
if not servers:
|
|
raise NotFoundError("Keine Server für diesen Kunden gefunden")
|
|
|
|
job_ids: list[int] = []
|
|
for server in servers:
|
|
job = await _create_job(
|
|
db,
|
|
JobTriggerRequest(
|
|
customer_id=payload.customer_id,
|
|
type=payload.type,
|
|
server_id=server.id,
|
|
reboot_if_required=payload.reboot_if_required,
|
|
),
|
|
user.username,
|
|
)
|
|
job_ids.append(job.id)
|
|
|
|
await AuditService(db).log(
|
|
username=user.username,
|
|
action="update.trigger_batch",
|
|
customer_id=payload.customer_id,
|
|
details={"type": payload.type.value, "count": len(job_ids)},
|
|
ip_address=client_ip(request),
|
|
)
|
|
return BatchJobTriggerResponse(created=len(job_ids), job_ids=job_ids)
|
|
|
|
|
|
@router.get("", response_model=list[UpdateJobRead])
|
|
async def list_jobs(
|
|
customer_id: int | None = None,
|
|
status: JobStatus | None = None,
|
|
limit: int = Query(default=50, le=200),
|
|
db: AsyncSession = Depends(get_db),
|
|
_user: User = Depends(get_current_user),
|
|
) -> list[UpdateJob]:
|
|
stmt = select(UpdateJob).order_by(UpdateJob.id.desc()).limit(limit)
|
|
if customer_id is not None:
|
|
stmt = stmt.where(UpdateJob.customer_id == customer_id)
|
|
if status:
|
|
stmt = stmt.where(UpdateJob.status == status)
|
|
result = await db.execute(stmt)
|
|
return list(result.scalars().all())
|
|
|
|
|
|
@router.get("/{job_id}", response_model=UpdateJobRead)
|
|
async def get_job(
|
|
job_id: int,
|
|
db: AsyncSession = Depends(get_db),
|
|
_user: User = Depends(get_current_user),
|
|
) -> UpdateJob:
|
|
job = await db.get(UpdateJob, job_id)
|
|
if not job:
|
|
raise NotFoundError("Job nicht gefunden")
|
|
return job
|
|
|
|
|
|
@router.get("/{job_id}/logs", response_model=list[UpdateLogRead])
|
|
async def get_job_logs(
|
|
job_id: int,
|
|
after_id: int = 0,
|
|
limit: int = Query(default=500, le=2000),
|
|
db: AsyncSession = Depends(get_db),
|
|
_user: User = Depends(get_current_user),
|
|
) -> list[UpdateLog]:
|
|
stmt = (
|
|
select(UpdateLog)
|
|
.where(UpdateLog.job_id == job_id, UpdateLog.id > after_id)
|
|
.order_by(UpdateLog.id)
|
|
.limit(limit)
|
|
)
|
|
result = await db.execute(stmt)
|
|
return list(result.scalars().all())
|
|
|
|
|
|
@router.post("/{job_id}/cancel", response_model=UpdateJobRead)
|
|
async def cancel_job(
|
|
job_id: int,
|
|
request: Request,
|
|
db: AsyncSession = Depends(get_db),
|
|
user: User = Depends(get_current_user),
|
|
) -> UpdateJob:
|
|
job = await db.get(UpdateJob, job_id)
|
|
if not job:
|
|
raise NotFoundError("Job nicht gefunden")
|
|
# Only pending jobs can be cancelled centrally - a claimed/running job
|
|
# is already on the satellite and finishes there.
|
|
if job.status != JobStatus.PENDING:
|
|
raise JobNotCancellableError()
|
|
|
|
job.status = JobStatus.CANCELLED
|
|
await AuditService(db).log(
|
|
username=user.username,
|
|
action="update.cancel",
|
|
target=f"job:{job_id}",
|
|
customer_id=job.customer_id,
|
|
ip_address=client_ip(request),
|
|
)
|
|
return job
|
|
|
|
|
|
@router.get("/stats/summary")
|
|
async def job_stats(
|
|
customer_id: int | None = None,
|
|
db: AsyncSession = Depends(get_db),
|
|
_user: User = Depends(get_current_user),
|
|
) -> dict:
|
|
base = select(func.count(UpdateJob.id))
|
|
if customer_id is not None:
|
|
base = base.where(UpdateJob.customer_id == customer_id)
|
|
total = await db.scalar(base)
|
|
running = await db.scalar(
|
|
base.where(UpdateJob.status.in_([JobStatus.CLAIMED, JobStatus.RUNNING]))
|
|
)
|
|
failed = await db.scalar(base.where(UpdateJob.status == JobStatus.FAILED))
|
|
pending = await db.scalar(base.where(UpdateJob.status == JobStatus.PENDING))
|
|
return {
|
|
"total": total or 0,
|
|
"running": running or 0,
|
|
"failed": failed or 0,
|
|
"pending": pending or 0,
|
|
}
|
|
|
|
|
|
async def _create_job(
|
|
db: AsyncSession, payload: JobTriggerRequest, username: str
|
|
) -> UpdateJob:
|
|
customer = await db.get(Customer, payload.customer_id)
|
|
if not customer:
|
|
raise NotFoundError("Kunde nicht gefunden")
|
|
|
|
if payload.type == JobType.NETWORK_SCAN:
|
|
if payload.server_id is not None:
|
|
raise ValidationError("NETWORK_SCAN hat keinen Ziel-Server")
|
|
else:
|
|
if payload.server_id is None:
|
|
raise ValidationError("server_id erforderlich")
|
|
server = await db.get(Server, payload.server_id)
|
|
if not server or server.customer_id != payload.customer_id:
|
|
raise NotFoundError("Server nicht gefunden")
|
|
|
|
params: dict = {"reboot_if_required": payload.reboot_if_required}
|
|
if payload.scan_subnet:
|
|
params["scan_subnet"] = payload.scan_subnet
|
|
|
|
job = UpdateJob(
|
|
customer_id=payload.customer_id,
|
|
server_id=payload.server_id,
|
|
type=payload.type,
|
|
created_by=username,
|
|
params=json.dumps(params),
|
|
)
|
|
db.add(job)
|
|
await db.flush()
|
|
return job
|