Files
B0rbor4d cc4c3fcecb Hub-and-Spoke Umbau: Multi-Tenant Zentrale + Satellite-Agent
- 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
2026-08-07 03:42:06 +00:00

229 lines
7.2 KiB
Python

"""Satellite agent API - polled by remote satellites, authenticated via X-Api-Key.
Pull model: satellites poll for jobs, execute them locally in the customer
network, push log batches and final results back here.
"""
import json
from datetime import UTC, datetime
from fastapi import APIRouter, Depends
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.api.deps import get_current_satellite
from app.core.database import get_db
from app.core.exceptions import ForbiddenError, NotFoundError
from app.core.logging import get_logger
from app.models.satellite import Satellite
from app.models.server import Server, ServerType
from app.models.update_job import JobStatus, JobType, UpdateJob, UpdateLog
from app.schemas.satellite_api import (
HealthReportRequest,
HeartbeatRequest,
JobResultRequest,
LogBatchRequest,
PollResponse,
SatelliteJob,
ScanResultRequest,
)
from app.services.audit import AuditService
logger = get_logger(__name__)
router = APIRouter()
@router.post("/heartbeat")
async def heartbeat(
payload: HeartbeatRequest,
satellite: Satellite = Depends(get_current_satellite),
db: AsyncSession = Depends(get_db),
) -> dict:
satellite.last_seen_at = datetime.now(UTC)
satellite.version = payload.version
satellite.hostname = payload.hostname
return {"ok": True, "server_time": datetime.now(UTC).isoformat()}
@router.get("/poll", response_model=PollResponse)
async def poll_jobs(
satellite: Satellite = Depends(get_current_satellite),
db: AsyncSession = Depends(get_db),
) -> PollResponse:
"""Claim and return pending jobs for this satellite's customer.
Claiming is atomic-ish: status flips pending -> claimed in the same
transaction, so two satellites of one customer do not get the same job.
"""
result = await db.execute(
select(UpdateJob)
.where(
UpdateJob.customer_id == satellite.customer_id,
UpdateJob.status == JobStatus.PENDING,
)
.order_by(UpdateJob.id)
.limit(5)
.with_for_update()
)
jobs = list(result.scalars().all())
now = datetime.now(UTC)
out: list[SatelliteJob] = []
for job in jobs:
job.status = JobStatus.CLAIMED
job.satellite_id = satellite.id
job.claimed_at = now
job.last_report_at = now
params = json.loads(job.params) if job.params else {}
server = job.server
out.append(
SatelliteJob(
job_id=job.id,
type=job.type,
server_id=server.id if server else None,
server_name=server.name if server else None,
hostname=server.hostname if server else None,
port=server.port if server else None,
server_type=server.type.value if server else None,
credential_ref=server.credential_ref if server else None,
reboot_if_required=bool(params.get("reboot_if_required", False)),
scan_subnet=params.get("scan_subnet"),
)
)
if out:
logger.info(
"satellite.jobs_claimed",
satellite=satellite.name,
customer_id=satellite.customer_id,
count=len(out),
)
return PollResponse(jobs=out)
@router.post("/logs")
async def push_logs(
payload: LogBatchRequest,
satellite: Satellite = Depends(get_current_satellite),
db: AsyncSession = Depends(get_db),
) -> dict:
job = await _get_own_job(db, payload.job_id, satellite)
if job.status == JobStatus.CLAIMED:
job.status = JobStatus.RUNNING
job.started_at = datetime.now(UTC)
for line in payload.lines:
db.add(
UpdateLog(
job_id=job.id,
timestamp=line.timestamp,
level=line.level,
line=line.line,
)
)
if payload.progress_percent is not None:
job.progress_percent = payload.progress_percent
if payload.current_phase is not None:
job.current_phase = payload.current_phase
job.last_report_at = datetime.now(UTC)
return {"ok": True, "accepted": len(payload.lines)}
@router.post("/result")
async def push_result(
payload: JobResultRequest,
satellite: Satellite = Depends(get_current_satellite),
db: AsyncSession = Depends(get_db),
) -> dict:
job = await _get_own_job(db, payload.job_id, satellite)
job.status = JobStatus.SUCCESS if payload.status == "success" else JobStatus.FAILED
job.error = payload.error
job.finished_at = datetime.now(UTC)
job.last_report_at = job.finished_at
if job.status == JobStatus.SUCCESS:
job.progress_percent = 100
await AuditService(db).log(
username=f"satellite:{satellite.name}",
action="job.result",
target=f"job:{job.id}",
result="success" if payload.status == "success" else "failure",
customer_id=satellite.customer_id,
details={"type": job.type.value, "error": payload.error},
)
return {"ok": True}
@router.post("/scan-result")
async def push_scan_result(
payload: ScanResultRequest,
satellite: Satellite = Depends(get_current_satellite),
db: AsyncSession = Depends(get_db),
) -> dict:
"""Ingest discovered hosts from a NETWORK_SCAN job as server candidates."""
job = await _get_own_job(db, payload.job_id, satellite)
if job.type != JobType.NETWORK_SCAN:
raise ForbiddenError("Scan-Ergebnisse nur für NETWORK_SCAN Jobs")
created = 0
for host in payload.hosts:
existing = await db.execute(
select(Server).where(
Server.customer_id == satellite.customer_id,
Server.hostname.in_([host.hostname, host.ip]),
)
)
if existing.scalar_one_or_none():
continue
if host.winrm_open:
stype, port = ServerType.WINDOWS, 5985
elif host.ssh_open:
stype, port = ServerType.LINUX, 22
else:
continue # not manageable - skip
db.add(
Server(
customer_id=satellite.customer_id,
name=host.hostname,
hostname=host.ip,
port=port,
type=stype,
description=f"Auto-Discovery via Scan (Job #{job.id})",
discovered_by_scan=True,
)
)
created += 1
return {"ok": True, "created": created}
@router.post("/health-report")
async def push_health_report(
payload: HealthReportRequest,
satellite: Satellite = Depends(get_current_satellite),
db: AsyncSession = Depends(get_db),
) -> dict:
server = await db.get(Server, payload.server_id)
if not server or server.customer_id != satellite.customer_id:
raise NotFoundError("Server nicht gefunden")
server.last_health_at = datetime.now(UTC)
server.last_health_ok = payload.ok
server.last_health_message = payload.message
return {"ok": True}
async def _get_own_job(
db: AsyncSession, job_id: int, satellite: Satellite
) -> UpdateJob:
job = await db.get(UpdateJob, job_id)
if not job or job.customer_id != satellite.customer_id:
raise NotFoundError("Job nicht gefunden")
return job