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", length), dmx_values, ] ) def build_artnet_poll_packet() -> bytes: return b"".join( [ ARTNET_HEADER, struct.pack(" ArtNetNode | None: if len(packet) < 190 or not packet.startswith(ARTNET_HEADER): return None opcode = 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