"""Async OpenRGB SDK v5 adapter with reconnect and serialized writes.""" from __future__ import annotations import asyncio import contextlib import logging import struct import time from datetime import UTC, datetime from typing import Any from ...config import Settings from ...errors import ConnectorError from ..base import ( Connector, ConnectorDevice, ConnectorHealth, DeviceState, HealthStatus, InventoryCallback, RGBColor, ) from .protocol import ( HEADER, PROTOCOL_VERSION, Controller, Mode, ModeFlag, PacketHeader, PacketId, ProtocolError, ZoneFlag, pack_header, pack_mode, pack_update_leds, pack_update_single_led, pack_update_zone, parse_controller, parse_header, parse_profile_list, ) LOGGER = logging.getLogger(__name__) class OpenRGBAdapter(Connector): id = "openrgb-local" kind = "openrgb" def __init__(self, settings: Settings) -> None: self.settings = settings self._reader: asyncio.StreamReader | None = None self._writer: asyncio.StreamWriter | None = None self._reader_task: asyncio.Task[None] | None = None self._monitor_task: asyncio.Task[None] | None = None self._connect_lock = asyncio.Lock() self._request_lock = asyncio.Lock() self._write_lock = asyncio.Lock() self._response_queues: dict[int, asyncio.Queue[tuple[PacketHeader, bytes] | Exception]] = {} self._controllers: dict[str, Controller] = {} self._inventory_callback: InventoryCallback | None = None self._stopping = False self._connected_at: datetime | None = None self._last_success_at: datetime | None = None self._last_error: str | None = None self._protocol_version: int | None = None self._reconnect_attempt = 0 self._last_nonzero: dict[str, list[RGBColor]] = {} @property def connected(self) -> bool: return self._writer is not None and not self._writer.is_closing() and self._protocol_version == 5 async def start(self) -> None: if self._monitor_task and not self._monitor_task.done(): return self._stopping = False self._monitor_task = asyncio.create_task(self._monitor(), name="openrgb-monitor") async def stop(self) -> None: self._stopping = True if self._monitor_task: self._monitor_task.cancel() with contextlib.suppress(asyncio.CancelledError): await self._monitor_task await self._disconnect() def set_inventory_callback(self, callback: InventoryCallback | None) -> None: self._inventory_callback = callback async def _monitor(self) -> None: while not self._stopping: try: await self._ensure_connected() if not self._controllers: await self.inventory() await self._notify_inventory_callback() await asyncio.sleep(10) except asyncio.CancelledError: raise except Exception as exc: self._last_error = str(exc) LOGGER.warning( "OpenRGB connection unavailable: %s", exc, extra={"connector": self.id, "error_category": "connection"}, ) await self._disconnect() delay = min(30.0, 0.5 * (2 ** min(self._reconnect_attempt, 6))) self._reconnect_attempt += 1 await asyncio.sleep(delay) async def _ensure_connected(self) -> None: if self.connected: return async with self._connect_lock: if self.connected: return await self._disconnect() try: connect = asyncio.open_connection(self.settings.openrgb_host, self.settings.openrgb_port) self._reader, self._writer = await asyncio.wait_for( connect, timeout=self.settings.openrgb_connect_timeout ) self._reader_task = asyncio.create_task(self._reader_loop(), name="openrgb-reader") header, payload = await self._request_connected( PacketId.REQUEST_PROTOCOL_VERSION, struct.pack(" None: writer, self._writer = self._writer, None reader_task, self._reader_task = self._reader_task, None self._reader = None self._protocol_version = None self._controllers.clear() if reader_task and reader_task is not asyncio.current_task(): reader_task.cancel() with contextlib.suppress(asyncio.CancelledError): await reader_task if writer: writer.close() with contextlib.suppress(Exception): await writer.wait_closed() error = ConnectorError( "openrgb_disconnected", "De verbinding met OpenRGB is verbroken.", retryable=True, ) for queue in self._response_queues.values(): with contextlib.suppress(asyncio.QueueFull): queue.put_nowait(error) async def _reader_loop(self) -> None: assert self._reader is not None try: while not self._stopping: raw_header = await self._reader.readexactly(HEADER.size) header = parse_header(raw_header, self.settings.openrgb_max_packet_size) payload = await self._reader.readexactly(header.payload_size) if header.packet_id == PacketId.DEVICE_LIST_UPDATED: if payload: raise ProtocolError("Device-list-updated bevat onverwachte data.") self._controllers.clear() if self._inventory_callback: asyncio.create_task( self._notify_inventory_callback(), name="openrgb-inventory-change" ) continue queue = self._response_queues.setdefault(header.packet_id, asyncio.Queue(maxsize=1)) if queue.full(): raise ProtocolError( "Onverwacht dubbel OpenRGB-antwoord.", packet_id=header.packet_id ) queue.put_nowait((header, payload)) except asyncio.CancelledError: raise except (asyncio.IncompleteReadError, ConnectionError, OSError, ProtocolError) as exc: self._last_error = str(exc) error = ConnectorError( "openrgb_connection_lost", "De verbinding met OpenRGB is onverwacht verbroken.", retryable=True, details={"reason": str(exc)}, ) for queue in self._response_queues.values(): with contextlib.suppress(asyncio.QueueFull): queue.put_nowait(error) if self._writer: self._writer.close() async def _notify_inventory_callback(self) -> None: if self._inventory_callback: await self._inventory_callback(self.id) async def _send_connected( self, device_index: int, packet_id: int | PacketId, payload: bytes = b"", *, expect_response: bool = False, ) -> None: del expect_response if self._writer is None or self._writer.is_closing(): raise ConnectorError( "openrgb_disconnected", "OpenRGB is niet verbonden.", retryable=True ) if len(payload) > self.settings.openrgb_max_packet_size: raise ProtocolError("Uitgaand OpenRGB-pakket is te groot.", size=len(payload)) self._writer.write(pack_header(device_index, packet_id, len(payload)) + payload) await self._writer.drain() async def _request_connected( self, packet_id: int | PacketId, payload: bytes = b"", *, device_index: int = 0, timeout: float | None = None, ) -> tuple[PacketHeader, bytes]: packet_value = int(packet_id) async with self._request_lock: queue = self._response_queues.setdefault(packet_value, asyncio.Queue(maxsize=1)) while not queue.empty(): queue.get_nowait() await self._send_connected(device_index, packet_id, payload, expect_response=True) try: result = await asyncio.wait_for( queue.get(), timeout=timeout or self.settings.openrgb_command_timeout ) except TimeoutError as exc: raise ConnectorError( "openrgb_timeout", "OpenRGB antwoordde niet binnen de ingestelde tijd.", retryable=True, details={"packet_id": packet_value}, ) from exc if isinstance(result, Exception): raise result return result async def _request( self, packet_id: int | PacketId, payload: bytes = b"", *, device_index: int = 0, ) -> tuple[PacketHeader, bytes]: await self._ensure_connected() return await self._request_connected( packet_id, payload, device_index=device_index ) async def _write(self, device_index: int, packet_id: int | PacketId, payload: bytes = b"") -> None: await self._ensure_connected() async with self._write_lock: await asyncio.wait_for( self._send_connected(device_index, packet_id, payload), timeout=self.settings.openrgb_command_timeout, ) self._last_success_at = datetime.now(UTC) async def test_connection(self) -> ConnectorHealth: started = time.perf_counter() try: await self._ensure_connected() await self._request(PacketId.REQUEST_CONTROLLER_COUNT) return ConnectorHealth( status=HealthStatus.HEALTHY, message="OpenRGB SDK protocol 5 is bereikbaar.", connected=True, latency_ms=(time.perf_counter() - started) * 1000, last_success_at=datetime.now(UTC), details={"protocol_version": self._protocol_version}, ) except Exception as exc: return ConnectorHealth( status=HealthStatus.DEGRADED, message=str(exc), connected=False, latency_ms=(time.perf_counter() - started) * 1000, last_success_at=self._last_success_at, details={"host": self.settings.openrgb_host, "port": self.settings.openrgb_port}, ) async def health(self) -> ConnectorHealth: return ConnectorHealth( status=HealthStatus.HEALTHY if self.connected else HealthStatus.DEGRADED, message="OpenRGB is verbonden." if self.connected else (self._last_error or "OpenRGB is niet verbonden."), connected=self.connected, last_success_at=self._last_success_at, details={ "host": self.settings.openrgb_host, "port": self.settings.openrgb_port, "protocol_version": self._protocol_version, "controllers": len(self._controllers), "connected_at": self._connected_at.isoformat() if self._connected_at else None, }, ) async def discover(self) -> list[ConnectorDevice]: return await self.inventory() async def inventory(self) -> list[ConnectorDevice]: _header, payload = await self._request(PacketId.REQUEST_CONTROLLER_COUNT) if len(payload) != 4: raise ProtocolError("OpenRGB-controllertelling heeft geen vier bytes.") count = struct.unpack(" 4096: raise ProtocolError("OpenRGB rapporteert onredelijk veel controllers.", count=count) controllers: dict[str, Controller] = {} for index in range(count): _header, controller_payload = await self._request( PacketId.REQUEST_CONTROLLER_DATA, struct.pack(" ConnectorDevice: active_mode = next( (mode for mode in controller.modes if mode.index == controller.active_mode), None ) active_flags = ModeFlag(active_mode.flags) if active_mode else ModeFlag(0) colors = ( active_mode.colors if active_mode and ModeFlag.HAS_MODE_SPECIFIC_COLOR in active_flags else controller.colors ) direct_mode = not active_mode or active_mode.name.casefold() in {"direct", "custom"} if active_mode and active_mode.name.casefold() == "off": power = False elif direct_mode: power = any(color.red or color.green or color.blue for color in controller.colors) else: power = True brightness = None if active_mode and ModeFlag.HAS_BRIGHTNESS in ModeFlag(active_mode.flags): span = active_mode.brightness_max - active_mode.brightness_min brightness = ( round((active_mode.brightness - active_mode.brightness_min) * 100 / span) if span else 100 ) return ConnectorDevice( external_id=controller.fingerprint, fingerprint=controller.fingerprint, name=controller.name, vendor=controller.vendor or None, model=controller.description or controller.name, serial=controller.serial or None, location=controller.location or None, firmware_version=controller.firmware_version or None, controller_index=controller.index, source="openrgb", device_type=controller.device_type_name, capabilities=controller.capabilities, state=DeviceState( power=power, brightness=brightness, colors=colors, mode=active_mode.name if active_mode else None, mode_index=active_mode.index if active_mode else None, speed=active_mode.speed if active_mode and active_mode.flags & ModeFlag.HAS_SPEED else None, direction=( active_mode.direction if active_mode and active_mode.flags & (ModeFlag.HAS_DIRECTION_LR | ModeFlag.HAS_DIRECTION_UD | ModeFlag.HAS_DIRECTION_HV) else None ), ), zones=[ { "index": zone.index, "name": zone.name, "type": zone.zone_type, "led_count": zone.led_count, "leds_min": zone.leds_min, "leds_max": zone.leds_max, "start_index": zone.start_index, "resizable_effects_only": bool(zone.flags & ZoneFlag.RESIZE_EFFECTS_ONLY), "segments": [ { "name": segment.name, "type": segment.zone_type, "start_index": segment.start_index, "led_count": segment.led_count, } for segment in zone.segments ], } for zone in controller.zones ], modes=[ { "index": mode.index, "name": mode.name, "flags": mode.flags, "speed_min": mode.speed_min if mode.flags & ModeFlag.HAS_SPEED else None, "speed_max": mode.speed_max if mode.flags & ModeFlag.HAS_SPEED else None, "brightness": bool(mode.flags & ModeFlag.HAS_BRIGHTNESS), "colors_min": mode.colors_min, "colors_max": mode.colors_max, } for mode in controller.modes ], led_count=len(controller.leds), metadata={"controller_flags": controller.flags, "led_alt_names": controller.led_alt_names}, ) def _controller(self, external_id: str) -> Controller: try: return self._controllers[external_id] except KeyError as exc: raise ConnectorError( "openrgb_identity_changed", "De OpenRGB-controlleridentiteit is gewijzigd; voer eerst een rescan uit.", status_code=409, retryable=True, ) from exc async def get_state(self, external_id: str) -> DeviceState: controller = self._controller(external_id) _header, payload = await self._request( PacketId.REQUEST_CONTROLLER_DATA, struct.pack(" DeviceState: controller = self._controller(external_id) # Re-read identity immediately before a hardware write. await self.get_state(external_id) controller = self._controller(external_id) colors = list(state.colors or []) selected_mode: Mode | None = None mode_specific_colors = False if ( state.mode is not None or state.mode_index is not None or state.brightness is not None or state.speed is not None or state.direction is not None ): selected_mode = self._select_mode(controller, state) mode = selected_mode mode_flags = ModeFlag(mode.flags) mode_specific_colors = bool( colors and ModeFlag.HAS_MODE_SPECIFIC_COLOR in mode_flags and state.zone_index is None and state.led_index is None ) brightness_percent = state.brightness if ( brightness_percent is not None and ModeFlag.HAS_BRIGHTNESS not in ModeFlag(mode.flags) and colors ): LOGGER.info( "OpenRGB brightness omitted for color write: controller=%s mode=%s", controller.index, mode.name, ) brightness_percent = None payload = pack_mode( mode, mode.index, brightness_percent=brightness_percent, speed=state.speed, direction=state.direction, colors=colors if mode_specific_colors else None, ) await self._write(controller.index, PacketId.UPDATE_MODE, payload) controller.active_mode = mode.index if brightness_percent is not None: span = mode.brightness_max - mode.brightness_min mode.brightness = round(mode.brightness_min + span * brightness_percent / 100) if state.speed is not None: mode.speed = state.speed if state.direction is not None: mode.direction = state.direction if mode_specific_colors: mode.colors = list(colors) has_leds = bool(controller.leds) if state.power is False: if has_leds and controller.colors and any( c.red or c.green or c.blue for c in controller.colors ): self._last_nonzero[external_id] = list(controller.colors) colors = [RGBColor(red=0, green=0, blue=0)] if has_leds else [] mode_specific_colors = False elif state.power is True and not colors and selected_mode is None and has_leds: colors = self._last_nonzero.get( external_id, [RGBColor(red=255, green=255, blue=255)] ) if colors and not mode_specific_colors: active_mode = next( (mode for mode in controller.modes if mode.index == controller.active_mode), None, ) if ( selected_mode is not None or active_mode is None or active_mode.name.casefold() not in {"custom", "direct"} ): await self._write(controller.index, PacketId.SET_CUSTOM_MODE) custom_mode = next( ( mode for mode in controller.modes if mode.name.casefold() in {"custom", "direct"} ), None, ) if custom_mode is not None: controller.active_mode = custom_mode.index if state.led_index is not None: if len(colors) != 1: raise ProtocolError("Een ledopdracht vereist exact één kleur.") payload = pack_update_single_led(state.led_index, colors[0], len(controller.leds)) await self._write(controller.index, PacketId.UPDATE_SINGLE_LED, payload) controller.colors[state.led_index] = colors[0] elif state.zone_index is not None: if state.zone_index >= len(controller.zones): raise ProtocolError("OpenRGB-zoneindex valt buiten bereik.") zone = controller.zones[state.zone_index] expanded = expand_colors(colors, zone.led_count) payload = pack_update_zone(zone.index, expanded, zone.led_count) await self._write(controller.index, PacketId.UPDATE_ZONE_LEDS, payload) controller.colors[zone.start_index : zone.start_index + zone.led_count] = expanded else: expanded = expand_colors(colors, len(controller.leds)) payload = pack_update_leds(expanded, len(controller.leds)) await self._write(controller.index, PacketId.UPDATE_LEDS, payload) controller.colors = expanded if any(color.red or color.green or color.blue for color in controller.colors): self._last_nonzero[external_id] = list(controller.colors) # OpenRGB may map SET_CUSTOM_MODE to a controller-specific mode (for example, # Corsair DRAM reports Direct even when it also advertises Custom). Read the # controller back so persisted desired state reflects hardware truth and does # not cause an endless reconciliation loop after discovery or restart. return await self.get_state(external_id) def _select_mode(self, controller: Controller, state: DeviceState) -> Mode: if state.mode_index is not None: mode = next((item for item in controller.modes if item.index == state.mode_index), None) elif state.mode is not None: mode = next( (item for item in controller.modes if item.name.casefold() == state.mode.casefold()), None, ) else: mode = next( (item for item in controller.modes if item.index == controller.active_mode), None ) if mode is None: raise ProtocolError("De gevraagde OpenRGB-modus bestaat niet.") return mode async def rescan(self) -> None: await self._write(0, PacketId.REQUEST_RESCAN_DEVICES) self._controllers.clear() async def list_profiles(self) -> list[str]: _header, payload = await self._request(PacketId.REQUEST_PROFILE_LIST) return parse_profile_list(payload) async def save_profile(self, name: str) -> None: await self._profile_write(PacketId.REQUEST_SAVE_PROFILE, name) async def load_profile(self, name: str) -> None: await self._profile_write(PacketId.REQUEST_LOAD_PROFILE, name) devices = await self.inventory() for device in devices: controller = self._controller(device.external_id) mode = self._select_mode(controller, DeviceState()) await self._write(controller.index, PacketId.UPDATE_MODE, pack_mode(mode, mode.index)) async def delete_profile(self, name: str) -> None: await self._profile_write(PacketId.REQUEST_DELETE_PROFILE, name) async def _profile_write(self, packet_id: PacketId, name: str) -> None: if not name or "\0" in name or len(name.encode("utf-8")) > 255: raise ProtocolError("Ongeldige OpenRGB-profielnaam.") await self._write(0, packet_id, name.encode("utf-8") + b"\0") async def resize_zone(self, external_id: str, zone_index: int, new_size: int) -> None: controller = self._controller(external_id) if not 0 <= zone_index < len(controller.zones): raise ProtocolError("OpenRGB-zoneindex valt buiten bereik.") zone = controller.zones[zone_index] if not min(zone.leds_min, zone.leds_max) <= new_size <= max(zone.leds_min, zone.leds_max): raise ProtocolError("Nieuwe OpenRGB-zonegrootte valt buiten bereik.") await self._write(controller.index, PacketId.RESIZE_ZONE, struct.pack(" None: controller = self._controller(external_id) if not 0 <= zone_index < len(controller.zones): raise ProtocolError("OpenRGB-zoneindex valt buiten bereik.") await self._write(controller.index, PacketId.CLEAR_SEGMENTS, struct.pack(" dict[str, Any]: return { "type": "object", "properties": { "host": {"type": "string", "default": "127.0.0.1"}, "port": {"type": "integer", "minimum": 1, "maximum": 65535, "default": 6742}, }, "required": ["host", "port"], "secret_fields": [], "warning": "OpenRGB SDK heeft geen authenticatie of TLS; gebruik uitsluitend loopback of een beveiligde tunnel.", } def expand_colors(colors: list[RGBColor], expected: int) -> list[RGBColor]: if expected <= 0: raise ProtocolError("OpenRGB-controller heeft geen adresseerbare leds.") if colors and all(color == colors[0] for color in colors): return [colors[0]] * expected if len(colors) != expected: raise ProtocolError( "Geef één kleur of exact één kleur per led op.", colors=len(colors), expected=expected ) return colors