Update docs and screenshots
CI / backend (pull_request) Canceled after 0s
CI / shell (pull_request) Canceled after 0s
CI / frontend (pull_request) Canceled after 0s
CI / arm64-smoke (pull_request) Canceled after 0s
CI / backend (push) Canceled after 0s
CI / shell (push) Canceled after 0s
CI / frontend (push) Canceled after 0s
CI / arm64-smoke (push) Canceled after 0s
CI / backend (pull_request) Canceled after 0s
CI / shell (pull_request) Canceled after 0s
CI / frontend (pull_request) Canceled after 0s
CI / arm64-smoke (pull_request) Canceled after 0s
CI / backend (push) Canceled after 0s
CI / shell (push) Canceled after 0s
CI / frontend (push) Canceled after 0s
CI / arm64-smoke (push) Canceled after 0s
This commit is contained in:
@@ -0,0 +1,2 @@
|
||||
"""DMX runtime services."""
|
||||
|
||||
@@ -0,0 +1,496 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import socket
|
||||
import struct
|
||||
import threading
|
||||
from array import array
|
||||
from collections import deque
|
||||
from collections.abc import Callable
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import UTC, datetime
|
||||
from typing import Any, Protocol
|
||||
|
||||
from app.dmx.frame import DmxFrame
|
||||
|
||||
|
||||
def utc_now() -> datetime:
|
||||
return datetime.now(UTC)
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class BackendStatus:
|
||||
connected: bool = True
|
||||
degraded: bool = False
|
||||
last_error: str | None = None
|
||||
reconnect_count: int = 0
|
||||
latency_ms: int = 5
|
||||
device_name: str = "Simulator universe 1"
|
||||
connected_since: datetime = field(default_factory=utc_now)
|
||||
backend_name: str = "simulator"
|
||||
last_successful_frame: datetime | None = None
|
||||
frames_sent: int = 0
|
||||
send_errors: int = 0
|
||||
selected_universe: int = 1
|
||||
selected_output_port: str | None = None
|
||||
|
||||
|
||||
class DmxBackend:
|
||||
async def startup(self) -> None:
|
||||
return None
|
||||
|
||||
async def shutdown(self) -> None:
|
||||
return None
|
||||
|
||||
async def send_frame(self, frame: DmxFrame) -> None: # pragma: no cover - interface
|
||||
raise NotImplementedError
|
||||
|
||||
def get_status(self) -> BackendStatus: # pragma: no cover - interface
|
||||
raise NotImplementedError
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class ArtNetNode:
|
||||
ip: str
|
||||
short_name: str
|
||||
long_name: str
|
||||
net: int
|
||||
sub_switch: int
|
||||
port_count: int
|
||||
raw_port_address: int
|
||||
|
||||
@property
|
||||
def label(self) -> str:
|
||||
name = self.long_name or self.short_name or self.ip
|
||||
return f"{name} ({self.ip})"
|
||||
|
||||
|
||||
class OlaClientAdapter(Protocol):
|
||||
def open(self) -> None: ...
|
||||
|
||||
def close(self) -> None: ...
|
||||
|
||||
def send_frame(self, universe: int, values: list[int]) -> None: ...
|
||||
|
||||
|
||||
class PythonOlaClientAdapter:
|
||||
def __init__(self, timeout_s: float = 1.0) -> None:
|
||||
self.timeout_s = timeout_s
|
||||
self._wrapper_class: type[Any] | None = None
|
||||
|
||||
def open(self) -> None:
|
||||
if self._wrapper_class is not None:
|
||||
return
|
||||
try:
|
||||
from ola.ClientWrapper import ClientWrapper
|
||||
except ImportError as exc: # pragma: no cover - depends on Pi runtime
|
||||
raise RuntimeError(
|
||||
"Python OLA bindings blev ikke fundet. "
|
||||
"Installer ola-python eller OLA's Python-modul."
|
||||
) from exc
|
||||
self._wrapper_class = ClientWrapper
|
||||
|
||||
def close(self) -> None:
|
||||
self._wrapper_class = None
|
||||
|
||||
def send_frame(self, universe: int, values: list[int]) -> None:
|
||||
self.open()
|
||||
assert self._wrapper_class is not None
|
||||
|
||||
wrapper = self._wrapper_class()
|
||||
client = wrapper.Client()
|
||||
dmx_buffer = array("B", values)
|
||||
completed = threading.Event()
|
||||
result: dict[str, object] = {
|
||||
"success": False,
|
||||
"error": f"Timeout under afsendelse til olad efter {self.timeout_s:.2f}s",
|
||||
}
|
||||
|
||||
def stop_wrapper() -> None:
|
||||
try:
|
||||
wrapper.Stop()
|
||||
except Exception:
|
||||
return
|
||||
|
||||
def on_timeout() -> None:
|
||||
completed.set()
|
||||
stop_wrapper()
|
||||
|
||||
def callback(state: object | None = None) -> None:
|
||||
if completed.is_set():
|
||||
return
|
||||
result["success"] = self._state_succeeded(state)
|
||||
if not result["success"]:
|
||||
result["error"] = self._state_message(state)
|
||||
completed.set()
|
||||
stop_wrapper()
|
||||
|
||||
timer = threading.Timer(self.timeout_s, on_timeout)
|
||||
timer.daemon = True
|
||||
timer.start()
|
||||
try:
|
||||
client.SendDmx(universe, dmx_buffer, callback)
|
||||
wrapper.Run()
|
||||
except Exception as exc:
|
||||
raise RuntimeError(f"OLA SendDmx fejlede: {exc}") from exc
|
||||
finally:
|
||||
timer.cancel()
|
||||
|
||||
if not bool(result["success"]):
|
||||
raise RuntimeError(str(result["error"]))
|
||||
|
||||
def _state_succeeded(self, state: object | None) -> bool:
|
||||
if state is None:
|
||||
return True
|
||||
if isinstance(state, bool):
|
||||
return state
|
||||
for attribute in ("Succeeded", "succeeded", "Ok", "ok", "success"):
|
||||
member = getattr(state, attribute, None)
|
||||
if callable(member):
|
||||
try:
|
||||
return bool(member())
|
||||
except Exception:
|
||||
continue
|
||||
if member is not None:
|
||||
return bool(member)
|
||||
return True
|
||||
|
||||
def _state_message(self, state: object | None) -> str:
|
||||
if state is None:
|
||||
return "Ukendt OLA-fejl"
|
||||
for attribute in ("message", "error", "status"):
|
||||
member = getattr(state, attribute, None)
|
||||
if callable(member):
|
||||
try:
|
||||
value = member()
|
||||
except Exception:
|
||||
continue
|
||||
if value:
|
||||
return str(value)
|
||||
elif member:
|
||||
return str(member)
|
||||
return str(state)
|
||||
|
||||
|
||||
class OlaDmxBackend(DmxBackend):
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
universe: int = 1,
|
||||
output_port: str | None = None,
|
||||
send_timeout_s: float = 1.0,
|
||||
adapter_factory: Callable[[], OlaClientAdapter] | None = None,
|
||||
) -> None:
|
||||
self._adapter_factory = adapter_factory or (
|
||||
lambda: PythonOlaClientAdapter(timeout_s=send_timeout_s)
|
||||
)
|
||||
self._adapter: OlaClientAdapter | None = None
|
||||
self._has_connected_once = False
|
||||
self._status = BackendStatus(
|
||||
connected=False,
|
||||
degraded=True,
|
||||
last_error="Afventer forbindelse til olad",
|
||||
backend_name="ola",
|
||||
device_name="OLA DMX output",
|
||||
selected_universe=universe,
|
||||
selected_output_port=output_port,
|
||||
)
|
||||
|
||||
async def startup(self) -> None:
|
||||
return None
|
||||
|
||||
async def shutdown(self) -> None:
|
||||
await self._disconnect()
|
||||
|
||||
async def send_frame(self, frame: DmxFrame) -> None:
|
||||
await self._ensure_connected()
|
||||
assert self._adapter is not None
|
||||
|
||||
try:
|
||||
await asyncio.to_thread(
|
||||
self._adapter.send_frame,
|
||||
self._status.selected_universe,
|
||||
frame.values.copy(),
|
||||
)
|
||||
except Exception as exc:
|
||||
self._status.connected = False
|
||||
self._status.degraded = True
|
||||
self._status.last_error = str(exc)
|
||||
self._status.send_errors += 1
|
||||
await self._disconnect()
|
||||
raise
|
||||
|
||||
self._status.connected = True
|
||||
self._status.degraded = False
|
||||
self._status.last_error = None
|
||||
self._status.frames_sent += 1
|
||||
self._status.last_successful_frame = utc_now()
|
||||
|
||||
def get_status(self) -> BackendStatus:
|
||||
return self._status
|
||||
|
||||
async def _ensure_connected(self) -> None:
|
||||
if self._adapter is not None:
|
||||
return
|
||||
adapter = self._adapter_factory()
|
||||
await asyncio.to_thread(adapter.open)
|
||||
self._adapter = adapter
|
||||
self._status.connected = True
|
||||
self._status.degraded = False
|
||||
self._status.last_error = None
|
||||
if self._has_connected_once:
|
||||
self._status.reconnect_count += 1
|
||||
else:
|
||||
self._has_connected_once = True
|
||||
self._status.connected_since = utc_now()
|
||||
|
||||
async def _disconnect(self) -> None:
|
||||
if self._adapter is None:
|
||||
return
|
||||
adapter = self._adapter
|
||||
self._adapter = None
|
||||
await asyncio.to_thread(adapter.close)
|
||||
|
||||
|
||||
ARTNET_PORT = 6454
|
||||
ARTNET_HEADER = b"Art-Net\x00"
|
||||
ARTNET_OPCODE_POLL = 0x2000
|
||||
ARTNET_OPCODE_POLL_REPLY = 0x2100
|
||||
ARTNET_OPCODE_DMX = 0x5000
|
||||
|
||||
|
||||
def build_artnet_dmx_packet(universe: int, values: list[int], sequence: int = 1) -> bytes:
|
||||
dmx_values = bytes(array("B", values[:512] + [0] * max(0, 512 - len(values))))
|
||||
port_address = max(0, universe - 1)
|
||||
sub_uni = port_address & 0xFF
|
||||
net = (port_address >> 8) & 0x7F
|
||||
length = len(dmx_values)
|
||||
return b"".join(
|
||||
[
|
||||
ARTNET_HEADER,
|
||||
struct.pack("<H", ARTNET_OPCODE_DMX),
|
||||
bytes([0x00, 0x0E]),
|
||||
bytes([sequence & 0xFF, 0x00]),
|
||||
bytes([sub_uni, net]),
|
||||
struct.pack(">H", length),
|
||||
dmx_values,
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
def build_artnet_poll_packet() -> bytes:
|
||||
return b"".join(
|
||||
[
|
||||
ARTNET_HEADER,
|
||||
struct.pack("<H", ARTNET_OPCODE_POLL),
|
||||
bytes([0x00, 0x0E]),
|
||||
bytes([0x00, 0x00]),
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
def parse_artnet_poll_reply(packet: bytes) -> ArtNetNode | None:
|
||||
if len(packet) < 190 or not packet.startswith(ARTNET_HEADER):
|
||||
return None
|
||||
opcode = struct.unpack_from("<H", packet, 8)[0]
|
||||
if opcode != ARTNET_OPCODE_POLL_REPLY:
|
||||
return None
|
||||
|
||||
ip = ".".join(str(part) for part in packet[10:14])
|
||||
net = packet[18]
|
||||
sub_switch = packet[19]
|
||||
short_name = packet[26:44].split(b"\x00", 1)[0].decode("utf-8", errors="ignore").strip()
|
||||
long_name = packet[44:108].split(b"\x00", 1)[0].decode("utf-8", errors="ignore").strip()
|
||||
port_count = struct.unpack_from(">H", packet, 172)[0]
|
||||
sw_out_0 = packet[190] if len(packet) > 190 else 0
|
||||
raw_port_address = (net << 8) | sw_out_0
|
||||
return ArtNetNode(
|
||||
ip=ip,
|
||||
short_name=short_name,
|
||||
long_name=long_name,
|
||||
net=net,
|
||||
sub_switch=sub_switch,
|
||||
port_count=port_count,
|
||||
raw_port_address=raw_port_address,
|
||||
)
|
||||
|
||||
|
||||
def discover_artnet_nodes(timeout_s: float = 1.0) -> list[ArtNetNode]:
|
||||
nodes: dict[str, ArtNetNode] = {}
|
||||
poll_packet = build_artnet_poll_packet()
|
||||
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM, socket.IPPROTO_UDP)
|
||||
try:
|
||||
sock.setsockopt(socket.SOL_SOCKET, socket.SO_BROADCAST, 1)
|
||||
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
sock.bind(("", 0))
|
||||
sock.settimeout(max(0.05, timeout_s / 5))
|
||||
sock.sendto(poll_packet, ("255.255.255.255", ARTNET_PORT))
|
||||
deadline = datetime.now(UTC).timestamp() + timeout_s
|
||||
while datetime.now(UTC).timestamp() < deadline:
|
||||
try:
|
||||
packet, _addr = sock.recvfrom(1024)
|
||||
except socket.timeout:
|
||||
continue
|
||||
node = parse_artnet_poll_reply(packet)
|
||||
if node is not None:
|
||||
nodes[node.ip] = node
|
||||
finally:
|
||||
sock.close()
|
||||
return sorted(nodes.values(), key=lambda item: (item.long_name or item.short_name or item.ip, item.ip))
|
||||
|
||||
|
||||
class ArtNetDmxBackend(DmxBackend):
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
universe: int = 1,
|
||||
target_host: str = "255.255.255.255",
|
||||
output_port: str | None = None,
|
||||
) -> None:
|
||||
self._target_host = target_host
|
||||
self._sequence = 0
|
||||
self._socket: socket.socket | None = None
|
||||
self._status = BackendStatus(
|
||||
connected=False,
|
||||
degraded=True,
|
||||
last_error="Afventer Artnet socket",
|
||||
backend_name="artnet",
|
||||
device_name=f"Art-Net output {target_host}",
|
||||
selected_universe=universe,
|
||||
selected_output_port=output_port or target_host,
|
||||
)
|
||||
|
||||
async def startup(self) -> None:
|
||||
if self._socket is not None:
|
||||
return
|
||||
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM, socket.IPPROTO_UDP)
|
||||
sock.setsockopt(socket.SOL_SOCKET, socket.SO_BROADCAST, 1)
|
||||
self._socket = sock
|
||||
self._status.connected = True
|
||||
self._status.degraded = False
|
||||
self._status.last_error = None
|
||||
self._status.connected_since = utc_now()
|
||||
|
||||
async def shutdown(self) -> None:
|
||||
if self._socket is not None:
|
||||
self._socket.close()
|
||||
self._socket = None
|
||||
self._status.connected = False
|
||||
|
||||
async def send_frame(self, frame: DmxFrame) -> None:
|
||||
if self._socket is None:
|
||||
await self.startup()
|
||||
assert self._socket is not None
|
||||
self._sequence = (self._sequence + 1) % 256
|
||||
packet = build_artnet_dmx_packet(self._status.selected_universe, frame.values.copy(), self._sequence)
|
||||
try:
|
||||
self._socket.sendto(packet, (self._target_host, ARTNET_PORT))
|
||||
except OSError as exc:
|
||||
self._status.connected = False
|
||||
self._status.degraded = True
|
||||
self._status.last_error = str(exc)
|
||||
self._status.send_errors += 1
|
||||
raise RuntimeError(str(exc)) from exc
|
||||
|
||||
self._status.connected = True
|
||||
self._status.degraded = False
|
||||
self._status.last_error = None
|
||||
self._status.frames_sent += 1
|
||||
self._status.last_successful_frame = utc_now()
|
||||
|
||||
def get_status(self) -> BackendStatus:
|
||||
return self._status
|
||||
|
||||
|
||||
class SimulatorDmxBackend(DmxBackend):
|
||||
def __init__(self, universe: int = 1) -> None:
|
||||
self._status = BackendStatus(
|
||||
selected_universe=universe,
|
||||
selected_output_port="simulator",
|
||||
device_name=f"Simulator universe {universe}",
|
||||
)
|
||||
self._history: deque[DmxFrame] = deque(maxlen=120)
|
||||
self._faults: set[str] = set()
|
||||
|
||||
def inject_fault(self, fault: str, enabled: bool) -> None:
|
||||
if enabled:
|
||||
self._faults.add(fault)
|
||||
else:
|
||||
self._faults.discard(fault)
|
||||
|
||||
def history(self) -> list[DmxFrame]:
|
||||
return list(self._history)
|
||||
|
||||
async def send_frame(self, frame: DmxFrame) -> None:
|
||||
if "ola_unavailable" in self._faults:
|
||||
self._record_error("OLA utilgaengelig i simulator")
|
||||
raise RuntimeError(self._status.last_error or "Simulatorfejl")
|
||||
if "usb_disconnected" in self._faults:
|
||||
self._record_error("USB-DMX er frakoblet i simulator")
|
||||
raise RuntimeError(self._status.last_error or "Simulatorfejl")
|
||||
if "send_timeout" in self._faults:
|
||||
await asyncio.sleep(0.25)
|
||||
self._record_error("Simuleret send timeout")
|
||||
raise TimeoutError(self._status.last_error or "Simuleret send timeout")
|
||||
|
||||
await asyncio.sleep(self._status.latency_ms / 1000)
|
||||
was_connected = self._status.connected
|
||||
self._status.connected = True
|
||||
self._status.degraded = "slow_backend" in self._faults
|
||||
self._status.last_error = None
|
||||
self._status.frames_sent += 1
|
||||
self._status.last_successful_frame = utc_now()
|
||||
if not was_connected:
|
||||
self._status.connected_since = utc_now()
|
||||
self._status.reconnect_count += 1
|
||||
self._history.append(frame.copy())
|
||||
|
||||
def get_status(self) -> BackendStatus:
|
||||
return self._status
|
||||
|
||||
def _record_error(self, message: str) -> None:
|
||||
self._status.connected = False
|
||||
self._status.degraded = True
|
||||
self._status.last_error = message
|
||||
self._status.send_errors += 1
|
||||
|
||||
|
||||
class FakeOlaClientAdapter:
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
output_port: str = "mock-port-1",
|
||||
open_results: list[Exception | None] | None = None,
|
||||
send_results: list[Exception | None] | None = None,
|
||||
) -> None:
|
||||
self.output_port = output_port
|
||||
self.open_results: deque[Exception | None] = deque(open_results or [])
|
||||
self.send_results: deque[Exception | None] = deque(send_results or [])
|
||||
self.is_open = False
|
||||
self.open_calls = 0
|
||||
self.close_calls = 0
|
||||
self.sent_frames: list[tuple[int, list[int]]] = []
|
||||
|
||||
def open(self) -> None:
|
||||
self.open_calls += 1
|
||||
if self.open_results:
|
||||
outcome = self.open_results.popleft()
|
||||
if outcome is not None:
|
||||
raise outcome
|
||||
self.is_open = True
|
||||
|
||||
def close(self) -> None:
|
||||
self.close_calls += 1
|
||||
self.is_open = False
|
||||
|
||||
def send_frame(self, universe: int, values: list[int]) -> None:
|
||||
if not self.is_open:
|
||||
raise RuntimeError("Mock OLA-adapter er ikke aaben")
|
||||
self.sent_frames.append((universe, values.copy()))
|
||||
if self.send_results:
|
||||
outcome = self.send_results.popleft()
|
||||
if outcome is not None:
|
||||
raise outcome
|
||||
|
||||
|
||||
FakeOlaAdapter = FakeOlaClientAdapter
|
||||
@@ -0,0 +1,137 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import time
|
||||
from collections import deque
|
||||
from datetime import UTC, datetime
|
||||
|
||||
from app.core.config import get_settings
|
||||
from app.dmx.backends import DmxBackend, SimulatorDmxBackend
|
||||
from app.dmx.frame import DmxFrame, FrameLayer, merge_layers
|
||||
from app.telemetry.service import TelemetryService
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class DmxEngine:
|
||||
def __init__(self, telemetry: TelemetryService, backend: DmxBackend | None = None) -> None:
|
||||
self._settings = get_settings()
|
||||
self.telemetry = telemetry
|
||||
self.backend = backend or SimulatorDmxBackend()
|
||||
self.layers: dict[str, FrameLayer] = {}
|
||||
self.blackout = False
|
||||
self.freeze = False
|
||||
self.master = 255
|
||||
self.current_frame = DmxFrame()
|
||||
self.current_frames: dict[int, DmxFrame] = {1: DmxFrame()}
|
||||
self._task: asyncio.Task[None] | None = None
|
||||
self._running = False
|
||||
self._frame_times: deque[float] = deque(maxlen=120)
|
||||
|
||||
@property
|
||||
def is_running(self) -> bool:
|
||||
return self._running
|
||||
|
||||
async def start(self) -> None:
|
||||
if self._running:
|
||||
return
|
||||
await self.backend.startup()
|
||||
self._running = True
|
||||
self._task = asyncio.create_task(self._loop())
|
||||
|
||||
async def stop(self) -> None:
|
||||
self._running = False
|
||||
if self._task is not None:
|
||||
self._task.cancel()
|
||||
try:
|
||||
await self._task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
self._task = None
|
||||
await self.backend.shutdown()
|
||||
|
||||
def configure_backend(self, backend: DmxBackend) -> None:
|
||||
self.backend = backend
|
||||
|
||||
async def replace_backend(self, backend: DmxBackend) -> None:
|
||||
was_running = self._running
|
||||
if was_running:
|
||||
await self.stop()
|
||||
self.backend = backend
|
||||
selected_universe = self.backend.get_status().selected_universe
|
||||
self.current_frame = DmxFrame(universe=selected_universe)
|
||||
self.current_frames = {selected_universe: self.current_frame.copy()}
|
||||
if was_running:
|
||||
await self.start()
|
||||
|
||||
async def _loop(self) -> None:
|
||||
interval = 1 / self._settings.target_fps
|
||||
while self._running:
|
||||
started = time.perf_counter()
|
||||
try:
|
||||
if not self.freeze:
|
||||
frames = merge_layers(list(self.layers.values()), blackout=self.blackout)
|
||||
if self.master < 255:
|
||||
for frame in frames.values():
|
||||
frame.values = [int((value / 255) * self.master) for value in frame.values]
|
||||
self.current_frames = {universe: frame.copy() for universe, frame in frames.items()}
|
||||
selected_universe = self.backend.get_status().selected_universe
|
||||
frame = self.current_frames.get(selected_universe, DmxFrame(universe=selected_universe))
|
||||
self.current_frame = frame.copy()
|
||||
await self.backend.send_frame(frame)
|
||||
self.telemetry.record_send_success(self.backend.get_status())
|
||||
self._frame_times.append(time.perf_counter() - started)
|
||||
except Exception as exc:
|
||||
logger.exception("DMX send failed")
|
||||
self.telemetry.record_send_failure(str(exc), self.backend.get_status())
|
||||
elapsed = time.perf_counter() - started
|
||||
await asyncio.sleep(max(0, interval - elapsed))
|
||||
|
||||
def set_layer(self, layer: FrameLayer) -> None:
|
||||
self.layers[layer.name] = layer
|
||||
|
||||
def remove_layer(self, name: str) -> None:
|
||||
self.layers.pop(name, None)
|
||||
|
||||
def trigger_blackout(self) -> None:
|
||||
self.blackout = True
|
||||
|
||||
def release_blackout(self) -> None:
|
||||
self.blackout = False
|
||||
|
||||
def get_frame(self, universe: int) -> DmxFrame:
|
||||
frame = self.current_frames.get(universe)
|
||||
if frame is None:
|
||||
return DmxFrame(universe=universe)
|
||||
return frame.copy()
|
||||
|
||||
def snapshot(self) -> dict[str, object]:
|
||||
status = self.backend.get_status()
|
||||
fps = 0.0
|
||||
if self._frame_times:
|
||||
average = sum(self._frame_times) / len(self._frame_times)
|
||||
if average > 0:
|
||||
fps = 1 / average
|
||||
return {
|
||||
"backend": status.backend_name,
|
||||
"connected": status.connected,
|
||||
"degraded": status.degraded,
|
||||
"last_error": status.last_error,
|
||||
"last_successful_frame": status.last_successful_frame.isoformat()
|
||||
if status.last_successful_frame is not None
|
||||
else None,
|
||||
"frames_sent": status.frames_sent,
|
||||
"send_errors": status.send_errors,
|
||||
"reconnect_count": status.reconnect_count,
|
||||
"selected_universe": status.selected_universe,
|
||||
"selected_output_port": status.selected_output_port,
|
||||
"available_universes": sorted(self.current_frames),
|
||||
"blackout": self.blackout,
|
||||
"freeze": self.freeze,
|
||||
"master": self.master,
|
||||
"fps": round(fps, 2),
|
||||
"frame": self.current_frame.values,
|
||||
"source_map": self.current_frame.source_map,
|
||||
"updated_at": datetime.now(UTC).isoformat(),
|
||||
}
|
||||
@@ -0,0 +1,91 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import UTC, datetime
|
||||
|
||||
|
||||
def clamp_channel(value: int) -> int:
|
||||
return max(0, min(255, int(value)))
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class FrameLayer:
|
||||
name: str
|
||||
priority: int
|
||||
values_by_universe: dict[int, dict[int, int]]
|
||||
precedence_map_by_universe: dict[int, dict[int, str]] = field(default_factory=dict)
|
||||
created_at: datetime = field(default_factory=lambda: datetime.now(UTC))
|
||||
|
||||
@property
|
||||
def values(self) -> dict[int, int]:
|
||||
return self.values_by_universe.get(1, {})
|
||||
|
||||
@property
|
||||
def precedence_map(self) -> dict[int, str]:
|
||||
return self.precedence_map_by_universe.get(1, {})
|
||||
|
||||
@classmethod
|
||||
def from_channel_values(
|
||||
cls,
|
||||
name: str,
|
||||
priority: int,
|
||||
values: dict[int, int],
|
||||
precedence_map: dict[int, str] | None = None,
|
||||
*,
|
||||
universe: int = 1,
|
||||
) -> "FrameLayer":
|
||||
return cls(
|
||||
name=name,
|
||||
priority=priority,
|
||||
values_by_universe={universe: values},
|
||||
precedence_map_by_universe={universe: precedence_map or {}},
|
||||
)
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class DmxFrame:
|
||||
universe: int = 1
|
||||
values: list[int] = field(default_factory=lambda: [0] * 512)
|
||||
source_map: list[str] = field(default_factory=lambda: ["idle"] * 512)
|
||||
|
||||
def set_channel(self, channel: int, value: int, source: str = "manual") -> None:
|
||||
if not 1 <= channel <= 512:
|
||||
raise ValueError("Channel must be in range 1..512")
|
||||
self.values[channel - 1] = clamp_channel(value)
|
||||
self.source_map[channel - 1] = source
|
||||
|
||||
def copy(self) -> DmxFrame:
|
||||
return DmxFrame(self.universe, self.values.copy(), self.source_map.copy())
|
||||
|
||||
|
||||
def merge_layers(layers: list[FrameLayer], blackout: bool = False) -> dict[int, DmxFrame]:
|
||||
universes = sorted(
|
||||
{
|
||||
int(universe)
|
||||
for layer in layers
|
||||
for universe in layer.values_by_universe
|
||||
}
|
||||
) or [1]
|
||||
if blackout:
|
||||
return {universe: DmxFrame(universe=universe) for universe in universes}
|
||||
|
||||
ordered = sorted(layers, key=lambda layer: (layer.priority, layer.created_at))
|
||||
frames = {universe: DmxFrame(universe=universe) for universe in universes}
|
||||
for universe, frame in frames.items():
|
||||
for channel in range(1, 513):
|
||||
channel_candidates: list[tuple[FrameLayer, int]] = []
|
||||
for layer in ordered:
|
||||
values = layer.values_by_universe.get(universe, {})
|
||||
if channel in values:
|
||||
channel_candidates.append((layer, clamp_channel(values[channel])))
|
||||
if not channel_candidates:
|
||||
continue
|
||||
|
||||
precedence_map = channel_candidates[-1][0].precedence_map_by_universe.get(universe, {})
|
||||
precedence = precedence_map.get(channel, "ltp").lower()
|
||||
if precedence == "htp":
|
||||
chosen_layer, chosen_value = max(channel_candidates, key=lambda pair: pair[1])
|
||||
else:
|
||||
chosen_layer, chosen_value = channel_candidates[-1]
|
||||
frame.set_channel(channel, chosen_value, chosen_layer.name)
|
||||
return frames
|
||||
@@ -0,0 +1,152 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import UTC, datetime
|
||||
from typing import Literal
|
||||
|
||||
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
|
||||
|
||||
from app.core.config import get_settings
|
||||
from app.core.database import SessionLocal
|
||||
from app.dmx.backends import (
|
||||
ArtNetDmxBackend,
|
||||
ArtNetNode,
|
||||
OlaDmxBackend,
|
||||
SimulatorDmxBackend,
|
||||
discover_artnet_nodes,
|
||||
)
|
||||
from app.dmx.engine import DmxEngine
|
||||
from app.models.entities import Setting
|
||||
|
||||
|
||||
DMX_OUTPUT_SETTING_KEY = "dmx_output"
|
||||
BackendKind = Literal["simulator", "ola", "artnet"]
|
||||
|
||||
|
||||
class DmxOutputService:
|
||||
def __init__(
|
||||
self,
|
||||
engine: DmxEngine,
|
||||
session_factory: async_sessionmaker[AsyncSession] = SessionLocal,
|
||||
) -> None:
|
||||
self.engine = engine
|
||||
self._session_factory = session_factory
|
||||
self._settings = get_settings()
|
||||
self._config = self._default_config()
|
||||
self._artnet_nodes: list[ArtNetNode] = []
|
||||
|
||||
async def startup(self) -> None:
|
||||
await self._load_config()
|
||||
self.engine.configure_backend(self._build_backend(self._config))
|
||||
|
||||
async def get_config(self) -> dict[str, object]:
|
||||
return {
|
||||
**self._config,
|
||||
"artnet_nodes": [self._serialize_node(node) for node in self._artnet_nodes],
|
||||
}
|
||||
|
||||
async def save_config(self, payload: dict[str, object]) -> dict[str, object]:
|
||||
normalized = self._normalize_config(payload)
|
||||
self._config = normalized
|
||||
await self._persist_config(normalized)
|
||||
backend = self._build_backend(normalized)
|
||||
if self.engine.is_running:
|
||||
await self.engine.replace_backend(backend)
|
||||
else:
|
||||
self.engine.configure_backend(backend)
|
||||
return await self.get_config()
|
||||
|
||||
async def discover_artnet(self, timeout_s: float = 1.0) -> dict[str, object]:
|
||||
nodes = await self._discover(timeout_s)
|
||||
self._artnet_nodes = nodes
|
||||
return {
|
||||
"items": [self._serialize_node(node) for node in nodes],
|
||||
"count": len(nodes),
|
||||
}
|
||||
|
||||
async def _discover(self, timeout_s: float) -> list[ArtNetNode]:
|
||||
return await self._to_thread_discovery(timeout_s)
|
||||
|
||||
async def _to_thread_discovery(self, timeout_s: float) -> list[ArtNetNode]:
|
||||
import asyncio
|
||||
|
||||
return await asyncio.to_thread(discover_artnet_nodes, timeout_s)
|
||||
|
||||
async def _load_config(self) -> None:
|
||||
async with self._session_factory() as session:
|
||||
setting = await session.get(Setting, DMX_OUTPUT_SETTING_KEY)
|
||||
if setting is None or not isinstance(setting.value, dict):
|
||||
return
|
||||
self._config = self._normalize_config(setting.value)
|
||||
|
||||
async def _persist_config(self, config: dict[str, object]) -> None:
|
||||
async with self._session_factory() as session:
|
||||
setting = await session.get(Setting, DMX_OUTPUT_SETTING_KEY)
|
||||
if setting is None:
|
||||
setting = Setting(
|
||||
key=DMX_OUTPUT_SETTING_KEY,
|
||||
value=config,
|
||||
updated_at=datetime.now(UTC),
|
||||
)
|
||||
session.add(setting)
|
||||
else:
|
||||
setting.value = config
|
||||
setting.updated_at = datetime.now(UTC)
|
||||
await session.commit()
|
||||
|
||||
def _default_config(self) -> dict[str, object]:
|
||||
backend: BackendKind = "simulator" if self._settings.simulator_enabled else "ola"
|
||||
return {
|
||||
"backend": backend,
|
||||
"universe": self._settings.ola_universe,
|
||||
"output_port": self._settings.ola_output_port or "",
|
||||
"target_host": "",
|
||||
}
|
||||
|
||||
def _normalize_config(self, payload: dict[str, object]) -> dict[str, object]:
|
||||
backend = str(payload.get("backend", self._default_config()["backend"])).lower()
|
||||
if backend not in {"simulator", "ola", "artnet"}:
|
||||
backend = "simulator"
|
||||
universe = int(payload.get("universe", self._settings.ola_universe) or self._settings.ola_universe)
|
||||
output_port = str(payload.get("output_port", "") or "").strip()
|
||||
target_host = str(payload.get("target_host", "") or "").strip()
|
||||
if backend == "artnet" and not target_host:
|
||||
target_host = "255.255.255.255"
|
||||
if backend == "simulator":
|
||||
output_port = "simulator"
|
||||
return {
|
||||
"backend": backend,
|
||||
"universe": max(1, min(63999, universe)),
|
||||
"output_port": output_port,
|
||||
"target_host": target_host,
|
||||
}
|
||||
|
||||
def _build_backend(self, config: dict[str, object]):
|
||||
backend = str(config["backend"])
|
||||
universe = int(config["universe"])
|
||||
output_port = str(config.get("output_port", "") or "") or None
|
||||
if backend == "ola":
|
||||
return OlaDmxBackend(
|
||||
universe=universe,
|
||||
output_port=output_port,
|
||||
send_timeout_s=self._settings.ola_send_timeout_ms / 1000,
|
||||
)
|
||||
if backend == "artnet":
|
||||
target_host = str(config.get("target_host", "") or "255.255.255.255")
|
||||
return ArtNetDmxBackend(
|
||||
universe=universe,
|
||||
target_host=target_host,
|
||||
output_port=output_port or target_host,
|
||||
)
|
||||
return SimulatorDmxBackend(universe=universe)
|
||||
|
||||
def _serialize_node(self, node: ArtNetNode) -> dict[str, object]:
|
||||
return {
|
||||
"ip": node.ip,
|
||||
"short_name": node.short_name,
|
||||
"long_name": node.long_name,
|
||||
"label": node.label,
|
||||
"net": node.net,
|
||||
"sub_switch": node.sub_switch,
|
||||
"port_count": node.port_count,
|
||||
"raw_port_address": node.raw_port_address,
|
||||
}
|
||||
Reference in New Issue
Block a user