324 lines
12 KiB
Python
324 lines
12 KiB
Python
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import threading
|
||
from collections.abc import Mapping
|
||
from datetime import UTC, datetime
|
||
from pathlib import Path
|
||
from typing import Any
|
||
|
||
from bleak.exc import BleakError
|
||
from fastapi import FastAPI, HTTPException, WebSocket, WebSocketDisconnect
|
||
from fastapi.staticfiles import StaticFiles
|
||
from pydantic import BaseModel, Field
|
||
|
||
from k1link import __version__
|
||
from k1link.artifacts import write_json_atomic
|
||
from k1link.ble.scanner import scan
|
||
from k1link.ble.wifi_provisioning import AP_FALLBACK_IPV4, provision_wifi_once
|
||
from k1link.mqtt import validate_private_ipv4
|
||
from k1link.viewer.runtime import VisualizationRuntime, new_live_session_dir
|
||
|
||
REPOSITORY_ROOT = Path(__file__).resolve().parents[3]
|
||
|
||
|
||
class BleScanRequest(BaseModel):
|
||
duration_seconds: float = Field(default=6.0, ge=1.0, le=60.0)
|
||
|
||
|
||
class ConnectRequest(BaseModel):
|
||
device_id: str = Field(min_length=1, max_length=128)
|
||
ssid: str = Field(min_length=1, max_length=128)
|
||
password: str = Field(min_length=1, max_length=256)
|
||
|
||
|
||
class LiveRequest(BaseModel):
|
||
host: str | None = Field(default=None, max_length=15)
|
||
duration_seconds: float = Field(default=3600.0, ge=1.0, le=86_400.0)
|
||
|
||
|
||
class ReplayRequest(BaseModel):
|
||
path: str = Field(min_length=1, max_length=4096)
|
||
speed: float = Field(default=1.0, ge=0.0, le=100.0)
|
||
loop: bool = False
|
||
|
||
|
||
class ConsoleService:
|
||
def __init__(self, repository_root: Path) -> None:
|
||
self.repository_root = repository_root.resolve()
|
||
self._lock = threading.Lock()
|
||
self._devices: list[dict[str, Any]] = []
|
||
self._selected_device_id: str | None = None
|
||
self._k1_ip: str | None = None
|
||
self._operation_phase: str | None = None
|
||
self._operation_message: str | None = None
|
||
self.runtime = VisualizationRuntime()
|
||
|
||
def state(self) -> dict[str, Any]:
|
||
runtime = self.runtime.snapshot()
|
||
metrics = runtime["metrics"]
|
||
with self._lock:
|
||
operation_phase = self._operation_phase
|
||
operation_message = self._operation_message
|
||
devices = list(self._devices)
|
||
selected_device_id = self._selected_device_id
|
||
k1_ip = self._k1_ip
|
||
|
||
runtime_active = runtime["source_mode"] != "idle" or runtime["phase"] in {
|
||
"starting_live",
|
||
"stopping",
|
||
"error",
|
||
}
|
||
if operation_phase is not None:
|
||
phase = operation_phase
|
||
message = operation_message
|
||
elif runtime_active:
|
||
phase = runtime["phase"]
|
||
message = runtime["message"]
|
||
elif k1_ip is not None:
|
||
phase = "connected"
|
||
message = runtime["message"]
|
||
elif selected_device_id is not None:
|
||
phase = "device_selected"
|
||
message = "K1 выбран. Теперь введите название и пароль Wi-Fi."
|
||
else:
|
||
phase = "idle"
|
||
message = runtime["message"]
|
||
|
||
return {
|
||
"phase": phase,
|
||
"message": message,
|
||
"devices": devices,
|
||
"selected_device_id": selected_device_id,
|
||
"k1_ip": k1_ip,
|
||
"foxglove_ws_url": runtime["foxglove_ws_url"],
|
||
"foxglove_viewer_url": runtime["foxglove_viewer_url"],
|
||
"source_mode": runtime["source_mode"],
|
||
"metrics": {
|
||
"pipeline_ms": metrics["mqtt_to_publish_ms"],
|
||
"end_to_end_ms": metrics["mqtt_to_publish_ms"],
|
||
"decode_ms": metrics["decode_publish_ms"],
|
||
"frame_rate": metrics["pcl_fps"],
|
||
"frame_rate_hz": metrics["pcl_fps"],
|
||
"point_count": metrics["last_point_count"],
|
||
"dropped_preview_frames": metrics["preview_dropped"],
|
||
**metrics,
|
||
},
|
||
}
|
||
|
||
async def scan_ble(self, duration_seconds: float) -> dict[str, Any]:
|
||
self._set_operation("scanning", "Ищем K1 по Bluetooth (BLE)…")
|
||
try:
|
||
result = await scan(duration_seconds)
|
||
devices = [
|
||
{
|
||
"device_id": item["macos_uuid"],
|
||
"name": item["local_name"] or item["name"],
|
||
"rssi": item["rssi"],
|
||
"address": None,
|
||
"connectable": True,
|
||
}
|
||
for item in result["devices"]
|
||
if item["k1_name_candidate"]
|
||
]
|
||
with self._lock:
|
||
self._devices = devices
|
||
if len(devices) == 1:
|
||
self._selected_device_id = str(devices[0]["device_id"])
|
||
self._operation_message = (
|
||
f"Поиск завершён. Найдено устройств K1: {len(devices)}."
|
||
)
|
||
finally:
|
||
with self._lock:
|
||
self._operation_phase = None
|
||
return self.state()
|
||
|
||
async def connect(self, request: ConnectRequest) -> dict[str, Any]:
|
||
known_ids = {str(item["device_id"]) for item in self.state()["devices"]}
|
||
if request.device_id not in known_ids:
|
||
raise ValueError("сначала найдите и выберите K1 через Bluetooth")
|
||
self._set_operation(
|
||
"provisioning",
|
||
"Передаём в K1 настройки Wi-Fi одним подтверждённым запросом.",
|
||
)
|
||
session_dir = _new_operation_session_dir(
|
||
self.repository_root,
|
||
"viewer_wifi_provisioning",
|
||
)
|
||
session_dir.mkdir(parents=True, exist_ok=False)
|
||
password = request.password
|
||
try:
|
||
result = await provision_wifi_once(
|
||
request.device_id,
|
||
request.ssid,
|
||
password,
|
||
timeout_seconds=45.0,
|
||
write_mode="auto",
|
||
)
|
||
write_json_atomic(session_dir / "provisioning.sensitive.json", result)
|
||
ipv4 = _provisioned_ipv4(result)
|
||
write_json_atomic(
|
||
session_dir / "manifest.redacted.json",
|
||
{
|
||
"schema_version": 1,
|
||
"started_at_utc": result["started_at_utc"],
|
||
"completed_at_utc": result["completed_at_utc"],
|
||
"operation": "single_reviewed_wifi_provisioning_write",
|
||
"profile_id": result["profile_id"],
|
||
"outcome": result["outcome"],
|
||
"k1_lan_address_observed": ipv4 is not None,
|
||
"credentials_persisted_by_connector": False,
|
||
},
|
||
)
|
||
if ipv4 is None:
|
||
raise RuntimeError(
|
||
"K1 не сообщил адрес в локальной сети; автоматического повтора не было"
|
||
)
|
||
with self._lock:
|
||
self._selected_device_id = request.device_id
|
||
self._k1_ip = ipv4
|
||
self._operation_message = "K1 подключён к Wi-Fi и сообщил локальный адрес."
|
||
finally:
|
||
password = ""
|
||
with self._lock:
|
||
self._operation_phase = None
|
||
return self.state()
|
||
|
||
def start_live(self, host: str | None, duration_seconds: float) -> dict[str, Any]:
|
||
target = host or self.state()["k1_ip"]
|
||
if not isinstance(target, str) or not target:
|
||
raise ValueError(
|
||
"сначала подключите K1 к Wi-Fi или укажите его локальный адрес"
|
||
)
|
||
target = validate_private_ipv4(target)
|
||
out_dir = new_live_session_dir(self.repository_root)
|
||
self.runtime.start_live(target, out_dir, duration_seconds=duration_seconds)
|
||
return self.state()
|
||
|
||
def start_replay(self, path: str, speed: float, loop: bool) -> dict[str, Any]:
|
||
replay_path = Path(path).expanduser().resolve()
|
||
if not replay_path.is_relative_to(self.repository_root):
|
||
raise ValueError("файл записи должен находиться внутри репозитория")
|
||
self.runtime.start_replay(replay_path, speed=speed, loop=loop)
|
||
return self.state()
|
||
|
||
def stop(self) -> dict[str, Any]:
|
||
self.runtime.stop()
|
||
return self.state()
|
||
|
||
def _set_operation(self, phase: str, message: str) -> None:
|
||
with self._lock:
|
||
self._operation_phase = phase
|
||
self._operation_message = message
|
||
|
||
|
||
service = ConsoleService(REPOSITORY_ROOT)
|
||
app = FastAPI(
|
||
title="NODE.DC K1 Live Console API",
|
||
version=__version__,
|
||
docs_url="/api/docs",
|
||
redoc_url=None,
|
||
openapi_url="/api/openapi.json",
|
||
)
|
||
|
||
|
||
@app.get("/api/health")
|
||
def health() -> dict[str, Any]:
|
||
return {
|
||
"ok": True,
|
||
"status": "ok",
|
||
"service": "k1-live-console",
|
||
"version": __version__,
|
||
}
|
||
|
||
|
||
@app.get("/api/state")
|
||
def get_state() -> dict[str, Any]:
|
||
return service.state()
|
||
|
||
|
||
@app.post("/api/ble/scan")
|
||
async def scan_ble(request: BleScanRequest) -> dict[str, Any]:
|
||
try:
|
||
return await service.scan_ble(request.duration_seconds)
|
||
except (BleakError, OSError, RuntimeError, ValueError) as exc:
|
||
raise HTTPException(status_code=502, detail=f"Ошибка поиска BLE: {exc}") from exc
|
||
|
||
|
||
@app.post("/api/connect")
|
||
async def connect(request: ConnectRequest) -> dict[str, Any]:
|
||
try:
|
||
return await service.connect(request)
|
||
except (BleakError, OSError, TimeoutError, RuntimeError, ValueError) as exc:
|
||
raise HTTPException(
|
||
status_code=502,
|
||
detail=f"Ошибка подключения K1 к Wi-Fi: {exc}",
|
||
) from exc
|
||
|
||
|
||
@app.post("/api/session/live")
|
||
def start_live(request: LiveRequest) -> dict[str, Any]:
|
||
try:
|
||
return service.start_live(request.host, request.duration_seconds)
|
||
except (OSError, RuntimeError, ValueError) as exc:
|
||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||
|
||
|
||
@app.post("/api/session/replay")
|
||
def start_replay(request: ReplayRequest) -> dict[str, Any]:
|
||
try:
|
||
return service.start_replay(request.path, request.speed, request.loop)
|
||
except (OSError, RuntimeError, ValueError) as exc:
|
||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||
|
||
|
||
@app.post("/api/session/stop")
|
||
def stop_session() -> dict[str, Any]:
|
||
return service.stop()
|
||
|
||
|
||
@app.websocket("/api/events")
|
||
async def events(websocket: WebSocket) -> None:
|
||
await websocket.accept()
|
||
try:
|
||
while True:
|
||
await websocket.send_json({"state": service.state()})
|
||
await asyncio.sleep(0.5)
|
||
except WebSocketDisconnect:
|
||
return
|
||
|
||
|
||
frontend_dist = REPOSITORY_ROOT / "apps" / "k1-viewer" / "dist"
|
||
if frontend_dist.is_dir():
|
||
app.mount("/", StaticFiles(directory=frontend_dist, html=True), name="frontend")
|
||
|
||
|
||
def _new_operation_session_dir(repository_root: Path, suffix: str) -> Path:
|
||
stamp = datetime.now(UTC).strftime("%Y%m%dT%H%M%SZ")
|
||
base = repository_root / "sessions" / f"{stamp}_{suffix}"
|
||
candidate = base
|
||
serial = 1
|
||
while candidate.exists():
|
||
serial += 1
|
||
candidate = base.with_name(f"{base.name}_{serial:02d}")
|
||
return candidate
|
||
|
||
|
||
def _provisioned_ipv4(result: Mapping[str, Any]) -> str | None:
|
||
observations = result.get("observations")
|
||
if not isinstance(observations, list):
|
||
return None
|
||
for observation in reversed(observations):
|
||
if not isinstance(observation, dict):
|
||
continue
|
||
status = observation.get("status")
|
||
if not isinstance(status, dict):
|
||
continue
|
||
address = status.get("ipv4")
|
||
if isinstance(address, str) and address != AP_FALLBACK_IPV4:
|
||
try:
|
||
return validate_private_ipv4(address)
|
||
except ValueError:
|
||
continue
|
||
return None
|