"""Operator-only fleet admission; the separate private mTLS listener is in fleet.""" from __future__ import annotations import asyncio import ipaddress import json from contextlib import suppress from typing import Annotated from urllib.parse import urlsplit from fastapi import APIRouter, Depends, HTTPException, Request, Response from fastapi.responses import StreamingResponse from pydantic import BaseModel, ConfigDict, Field from k1link.fleet.registry import FleetRegistry from k1link.fleet.trust import PairingError def local_operator(request: Request) -> FleetRegistry: try: peer = ipaddress.ip_address(request.client.host) host = urlsplit(f"http://{request.headers.get('host', '')}") if not peer.is_loopback or host.hostname not in ("127.0.0.1", "localhost", "::1"): raise ValueError origin = request.headers.get("origin") if origin and origin != f"http://{request.headers['host']}": raise ValueError if request.headers.get("sec-fetch-site") == "cross-site": raise ValueError except (ValueError, AttributeError, KeyError): raise HTTPException(403, "Откройте Mission Core на компьютере оператора.") from None registry = getattr(request.app.state, "fleet", None) if registry is None: raise HTTPException(503, "Реестр аппаратов недоступен. Повторите подключение.") return registry class PreviewRequest(BaseModel): model_config = ConfigDict(extra="forbid") code: str = Field(min_length=1, max_length=4096) class AddRequest(BaseModel): model_config = ConfigDict(extra="forbid") preview_id: str = Field(min_length=43, max_length=43) name: str = Field(min_length=1, max_length=80) platform: str = Field(pattern="^(ugv|uav|stationary|other)$") router = APIRouter(prefix="/api/v1/fleet", tags=["fleet"]) @router.get("/{vehicle_id}/monitor") def board_monitor(vehicle_id: str, response: Response, fleet: Annotated[FleetRegistry, Depends(local_operator)], metric: str = "cpu.usage", window: int = 900, end: float | None = None): response.headers["Cache-Control"] = "no-store" try: with fleet.lock: node_id = fleet.find(vehicle_id)["node_id"] return fleet.monitor.query(node_id, metric, window, end) except PairingError as error: raise HTTPException(404, str(error)) from None except ValueError: raise HTTPException(400, "Выберите доступный период и показатель.") from None @router.get("") def fleet_list(response: Response, fleet: Annotated[FleetRegistry, Depends(local_operator)]): response.headers["Cache-Control"] = "no-store" return fleet.listing() @router.get("/events") async def fleet_events(fleet: Annotated[FleetRegistry, Depends(local_operator)]): async def stream(): queue, close = fleet.events.subscribe() try: while not fleet.stop.is_set(): # Full replacement snapshots recover missed events/reconnections. # The timeout updates link expiry locally; it never queries a Node. value = await asyncio.to_thread(fleet.listing) yield "retry: 3000\ndata: " + json.dumps(value) + "\n\n" with suppress(TimeoutError): await asyncio.wait_for(queue.get(), timeout=5) finally: close() return StreamingResponse( stream(), media_type="text/event-stream", headers={ "Cache-Control": "no-store", "X-Accel-Buffering": "no", }, ) @router.post("/preview") def fleet_preview( body: PreviewRequest, response: Response, fleet: Annotated[FleetRegistry, Depends(local_operator)], ): response.headers["Cache-Control"] = "no-store" try: return fleet.preview(body.code) except PairingError as error: raise HTTPException(409, str(error)) from None @router.post("") def fleet_add( body: AddRequest, response: Response, fleet: Annotated[FleetRegistry, Depends(local_operator)] ): response.headers["Cache-Control"] = "no-store" try: return fleet.add(body.preview_id, body.name, body.platform) except PairingError as error: raise HTTPException(409, str(error)) from None @router.delete("/{vehicle_id}") def fleet_revoke(vehicle_id: str, fleet: Annotated[FleetRegistry, Depends(local_operator)]): try: return fleet.revoke(vehicle_id) except PairingError as error: raise HTTPException(404, str(error)) from None @router.post("/{vehicle_id}/devices/operations") def sensor_command( vehicle_id: str, body: dict, fleet: Annotated[FleetRegistry, Depends(local_operator)] ): from k1link.fleet.sensors import submit try: return submit(fleet, vehicle_id, body) except (PairingError, ValueError) as error: raise HTTPException(409, str(error)) from None @router.get("/{vehicle_id}/devices/operations/{operation_id}") def sensor_operation( vehicle_id: str, operation_id: str, fleet: Annotated[FleetRegistry, Depends(local_operator)] ): from k1link.fleet.sensors import operation try: return operation(fleet, vehicle_id, operation_id) except PairingError as error: raise HTTPException(404, str(error)) from None @router.get("/{vehicle_id}/devices/enrollment") def enrollment_state( vehicle_id: str, response: Response, fleet: Annotated[FleetRegistry, Depends(local_operator)] ): response.headers["Cache-Control"] = "no-store" try: with fleet.lock: row = fleet.find(vehicle_id) return { **row.get("device_enrollment", {"available": False}), "node_id": row["node_id"], "name": row["name"], "fresh": fleet.public(row)["connectivity"] == "online", } except PairingError as error: raise HTTPException(404, str(error)) from None @router.post("/{vehicle_id}/devices/enrollment/operations") def enrollment_submit( vehicle_id: str, body: dict, response: Response, fleet: Annotated[FleetRegistry, Depends(local_operator)], ): response.headers["Cache-Control"] = "no-store" try: return fleet.device_enrollment.submit(fleet, vehicle_id, body) except (PairingError, ValueError, TypeError): # Validation diagnostics must never echo a supplied credential. raise HTTPException( 409, "Запрос подключения не принят. Обновите БК и проверьте параметры сети." ) from None @router.get("/{vehicle_id}/devices/enrollment/operations/{operation_id}") def enrollment_operation( vehicle_id: str, operation_id: str, response: Response, fleet: Annotated[FleetRegistry, Depends(local_operator)], ): response.headers["Cache-Control"] = "no-store" try: return fleet.device_enrollment.operation(fleet, vehicle_id, operation_id) except PairingError as error: raise HTTPException(404, str(error)) from None