e4e76a7ede
- NodeIdentity/NodeRole in hms_domain: persistente node_id aus userdata/identity, Rollen RENDER_NODE/COORDINATOR/CONTROL_DESK, plausible Kombinationen validiert, keine Auto-Leader-Wahl (§6.3) - Control Core: create_app(identity) injizierbar; Dev-Modus ephemeral; Registry registriert eigene Node; neue Endpunkte /system/identity (ohne Secrets, §27.1) und /cluster/nodes (UI-Kategorien §6.3) - 16 neue/aktualisierte Integrationstests; Gesamtsuite 199 gruen
218 lines
7.6 KiB
Python
218 lines
7.6 KiB
Python
"""FastAPI-Anwendung des Control Core.
|
|
|
|
Endpunkte:
|
|
- GET /api/v1/system/health
|
|
- GET /api/v1/system/identity (node_id, display_name, roles; nicht vertraulich)
|
|
- GET /api/v1/system/capabilities
|
|
- GET /api/v1/parameters
|
|
- POST /api/v1/commands (parameter.set mit Revision-Prüfung und Idempotenz)
|
|
- POST /api/v1/commands/{command_id}/release
|
|
- GET /api/v1/diagnostics
|
|
- GET /api/v1/cluster/nodes
|
|
- WS /ws (State-Snapshot + Updates)
|
|
|
|
Commands folgen §23.2: command_id, type, expected_revision, actor, payload.
|
|
Node-Identität: persistente node_id aus userdata/identity (§3.6, §6.3);
|
|
im Dev-Modus ohne App-Root wird eine ephemeral-Identität erzeugt.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import uuid
|
|
|
|
from fastapi import FastAPI, HTTPException, WebSocket, WebSocketDisconnect
|
|
from hms_capabilities.probe import CapabilityReport
|
|
from hms_cluster.registry import NodeRegistry
|
|
from hms_domain.identity import NodeIdentity, NodeRole
|
|
from hms_parameter.engine import (
|
|
ControlSource,
|
|
ParameterEngine,
|
|
RevisionConflict,
|
|
)
|
|
from hms_protocol.idempotency import IdempotencyRegistry
|
|
from pydantic import BaseModel, Field
|
|
|
|
|
|
class SetParameterCommand(BaseModel):
|
|
"""parameter.set-Command (§23.2)."""
|
|
|
|
command_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
|
|
type: str = "parameter.set"
|
|
expected_revision: int | None = None
|
|
actor: dict = Field(default_factory=lambda: {"type": "web", "id": "operator-session"})
|
|
payload: dict
|
|
|
|
|
|
class ReleaseCommand(BaseModel):
|
|
source: str = "web"
|
|
|
|
|
|
class _State:
|
|
def __init__(self, identity: NodeIdentity) -> None:
|
|
self.identity = identity
|
|
self.engine = ParameterEngine()
|
|
self.registry = IdempotencyRegistry()
|
|
self.report = CapabilityReport()
|
|
self.nodes = NodeRegistry()
|
|
self.subscribers: list[asyncio.Queue] = []
|
|
|
|
|
|
def create_app(identity: NodeIdentity | None = None) -> FastAPI:
|
|
"""Erzeugt die Control-Core-App.
|
|
|
|
identity: produktiv vom Launcher geladene persistente Identität
|
|
(userdata/identity). Ohne Angabe gilt Dev-Modus mit ephemeral-Identität
|
|
(jede App-Instanz erhält eine eigene ID; für Showbetrieb unzulässig).
|
|
"""
|
|
if identity is None:
|
|
identity = NodeIdentity.ephemeral(
|
|
"Dev Node", frozenset({NodeRole.RENDER_NODE, NodeRole.COORDINATOR})
|
|
)
|
|
app = FastAPI(title="HMS MediaEngine Control Core", version="0.1.0")
|
|
state = _State(identity)
|
|
# Eigene Node in die Registry eintragen (§6.3)
|
|
state.nodes.register(
|
|
node_id=identity.node_id,
|
|
display_name=identity.display_name,
|
|
roles=tuple(r.value for r in identity.roles),
|
|
)
|
|
|
|
def _broadcast(event: dict) -> None:
|
|
for queue in list(state.subscribers):
|
|
queue.put_nowait(event)
|
|
|
|
@app.get("/api/v1/system/health")
|
|
async def health() -> dict:
|
|
return {
|
|
"status": "ok",
|
|
"phase": 1,
|
|
"node_id": state.identity.node_id,
|
|
"revision": state.engine.revision,
|
|
}
|
|
|
|
@app.get("/api/v1/system/identity")
|
|
async def system_identity() -> dict:
|
|
"""Nicht vertrauliche Selbstauskunft (§27.1: keine Tokens/Secrets)."""
|
|
return {
|
|
"node_id": state.identity.node_id,
|
|
"display_name": state.identity.display_name,
|
|
"roles": sorted(r.value for r in state.identity.roles),
|
|
"renders_locally": state.identity.renders_locally,
|
|
"is_coordinator": state.identity.is_coordinator,
|
|
}
|
|
|
|
@app.get("/api/v1/system/capabilities")
|
|
async def capabilities() -> dict:
|
|
return state.report.as_dict()
|
|
|
|
@app.get("/api/v1/parameters")
|
|
async def parameters() -> dict:
|
|
snap = state.engine.snapshot()
|
|
return {"revision": snap.revision, "values": snap.as_dict()}
|
|
|
|
@app.post("/api/v1/commands")
|
|
async def post_command(cmd: SetParameterCommand) -> dict:
|
|
if cmd.type != "parameter.set":
|
|
raise HTTPException(status_code=400, detail=f"unknown command type {cmd.type!r}")
|
|
if not state.registry.register(cmd.command_id):
|
|
prior = state.registry.result(cmd.command_id)
|
|
if prior is not None:
|
|
return {"status": "ack", "duplicate": True, "result": prior}
|
|
raise HTTPException(status_code=409, detail="command already in flight")
|
|
path = cmd.payload.get("parameter_path")
|
|
value = cmd.payload.get("value")
|
|
if not path or value is None:
|
|
raise HTTPException(status_code=400, detail="payload requires parameter_path and value")
|
|
try:
|
|
revision = state.engine.set_value(
|
|
path=path,
|
|
value=float(value),
|
|
source=ControlSource.WEB,
|
|
expected_revision=cmd.expected_revision,
|
|
)
|
|
except RevisionConflict as exc:
|
|
raise HTTPException(
|
|
status_code=409,
|
|
detail={
|
|
"error": "REVISION_CONFLICT",
|
|
"current": exc.current,
|
|
"expected": exc.expected,
|
|
},
|
|
) from exc
|
|
except ValueError as exc:
|
|
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
|
result = {
|
|
"status": "ack",
|
|
"command_id": cmd.command_id,
|
|
"revision": revision,
|
|
"effective": state.engine.effective_value(path),
|
|
}
|
|
state.registry.complete(cmd.command_id, result)
|
|
_broadcast(
|
|
{
|
|
"type": "parameter.update",
|
|
"parameter_path": path,
|
|
"value": value,
|
|
"revision": revision,
|
|
}
|
|
)
|
|
return result
|
|
|
|
@app.post("/api/v1/commands/{command_id}/release")
|
|
async def release_override(command_id: str) -> dict:
|
|
# Release nach §11.3; command_id referenziert den ursprünglichen Command.
|
|
return {"status": "not_implemented_in_phase0"}
|
|
|
|
@app.get("/api/v1/diagnostics")
|
|
async def diagnostics() -> dict:
|
|
return {
|
|
"renderer": "not_connected", # IPC-Handshake folgt (ADR-0003)
|
|
"artnet": "not_started",
|
|
"revision": state.engine.revision,
|
|
"node_id": state.identity.node_id,
|
|
}
|
|
|
|
@app.get("/api/v1/cluster/nodes")
|
|
async def cluster_nodes() -> dict:
|
|
"""Node-Übersicht nach UI-Kategorien (§6.3): getrennt aufgeführt."""
|
|
nodes = []
|
|
for entry in state.nodes.all():
|
|
state.nodes.evaluate_health(entry.node_id)
|
|
nodes.append(
|
|
{
|
|
"node_id": entry.node_id,
|
|
"display_name": entry.display_name,
|
|
"roles": list(entry.roles),
|
|
"health": entry.health.value,
|
|
"category": entry.category.value,
|
|
}
|
|
)
|
|
return {"self": state.identity.node_id, "nodes": nodes}
|
|
|
|
@app.websocket("/ws")
|
|
async def websocket_endpoint(ws: WebSocket) -> None:
|
|
await ws.accept()
|
|
queue: asyncio.Queue = asyncio.Queue(maxsize=256)
|
|
state.subscribers.append(queue)
|
|
try:
|
|
snap = state.engine.snapshot()
|
|
await ws.send_json(
|
|
{"type": "snapshot", "revision": snap.revision, "values": snap.as_dict()}
|
|
)
|
|
while True:
|
|
try:
|
|
event = await asyncio.wait_for(queue.get(), timeout=15.0)
|
|
await ws.send_json(event)
|
|
except TimeoutError:
|
|
await ws.send_json({"type": "heartbeat"})
|
|
except WebSocketDisconnect:
|
|
pass
|
|
finally:
|
|
state.subscribers.remove(queue)
|
|
|
|
return app
|
|
|
|
|
|
app = create_app()
|