"""Inbound notification channel for the Riedel RRCS integration. After RegisterForAllEvents, RRCS reverses the roles: it becomes the XML-RPC client and posts method calls at us. We expose an unauthenticated view on a secret path to receive them, because RRCS has no way to present a bearer token. Answering matters. If the KeepAlive option is enabled in the RRCS route options it sends GetAlive every three seconds of inactivity and tears the channel down if we stay quiet (chapter 9.8). """ from __future__ import annotations import logging import xmlrpc.client from typing import Any, Callable from aiohttp import web from homeassistant.components.http import HomeAssistantView from homeassistant.core import HomeAssistant, callback from homeassistant.helpers.dispatcher import async_dispatcher_send from .const import ALARM_METHODS, DOMAIN, EVENT_RRCS_NOTIFICATION, PANEL_SPY_METHODS from .coordinator import RRCSCoordinator from .rrcs import GpioAddress, parse_xmlrpc _LOGGER = logging.getLogger(__name__) DATA_LISTENERS = f"{DOMAIN}_listeners" DATA_VIEW = f"{DOMAIN}_view_registered" NotificationHandler = Callable[[str, tuple[Any, ...]], None] def signal_event(entry_id: str) -> str: """Dispatcher signal carrying notifications for one config entry.""" return f"{DOMAIN}_{entry_id}_event" @callback def async_register_view(hass: HomeAssistant) -> None: """Register the notification view once for the whole integration.""" if hass.data.get(DATA_VIEW): return hass.http.register_view(RRCSNotificationView) hass.data[DATA_VIEW] = True @callback def async_register_listener( hass: HomeAssistant, token: str, handler: NotificationHandler ) -> None: """Attach a per-entry handler to a callback token.""" hass.data.setdefault(DATA_LISTENERS, {})[token] = handler @callback def async_remove_listener(hass: HomeAssistant, token: str) -> None: """Detach a per-entry handler.""" hass.data.get(DATA_LISTENERS, {}).pop(token, None) class RRCSNotificationView(HomeAssistantView): """Receives XML-RPC notifications pushed by RRCS.""" url = "/api/riedel_rrcs/{token}" name = "api:riedel_rrcs" requires_auth = False async def post(self, request: web.Request, token: str) -> web.Response: """Handle one inbound XML-RPC method call.""" hass: HomeAssistant = request.app["hass"] handler = hass.data.get(DATA_LISTENERS, {}).get(token) if handler is None: _LOGGER.debug("Notification received on unknown token %s", token[:6]) return web.Response(status=404) raw = await request.read() try: params, method = parse_xmlrpc(raw) except Exception as err: # noqa: BLE001 - never let a bad frame 500 _LOGGER.warning("Unparseable RRCS notification: %s", err) return web.Response(status=400) transkey = params[0] if params and isinstance(params[0], str) else "R0000000000" if method: try: handler(method, params) except Exception: # noqa: BLE001 - keep the channel alive regardless _LOGGER.exception("Error handling RRCS notification %s", method) body = xmlrpc.client.dumps( ([transkey, 0],), methodresponse=True, encoding="utf-8" ) return web.Response(body=body.encode("utf-8"), content_type="text/xml") class RRCSNotificationDispatcher: """Turns inbound RRCS method calls into state updates and HA events.""" def __init__( self, hass: HomeAssistant, entry_id: str, coordinator: RRCSCoordinator ) -> None: """Initialise the dispatcher.""" self.hass = hass self.entry_id = entry_id self.coordinator = coordinator @callback def handle(self, method: str, params: tuple[Any, ...]) -> None: """Route one notification.""" if method == "GetAlive": return args = list(params[1:]) if method == "LogicSourceChange" and len(args) >= 2: object_id, state = args[0], args[1] if isinstance(object_id, int): self.coordinator.apply_logic_source_change(object_id, bool(state)) elif method in ("GpInputChange", "GpOutputChange") and len(args) >= 6: net, node, port, slot, index, state = args[:6] if all(isinstance(value, int) for value in (net, node, port, slot, index)): address = GpioAddress( net=net, node=node, port=port, slot=slot, index=index, is_input=method == "GpInputChange", ) self.coordinator.apply_gpio_change(address, bool(state)) elif method == "ConnectArtistRestored": self.coordinator.apply_connection_change(True) elif method in ("ConnectArtistFailure", "GatewayShutdown"): self.coordinator.apply_connection_change(False) elif method == "ConfigurationChange": # The configuration changed underneath us; re-enumerate on the next # tick rather than trying to guess what moved. self.hass.async_create_task(self.coordinator.async_request_refresh()) payload = { "entry_id": self.entry_id, "method": method, "params": _sanitise(args), } self.hass.bus.async_fire(EVENT_RRCS_NOTIFICATION, payload) if method in ALARM_METHODS or method in PANEL_SPY_METHODS or method.startswith( "SendString" ): async_dispatcher_send(self.hass, signal_event(self.entry_id), payload) def _sanitise(value: Any) -> Any: """Coerce decoded XML-RPC values into something the event bus can carry.""" if isinstance(value, dict): return {str(key): _sanitise(item) for key, item in value.items()} if isinstance(value, (list, tuple)): return [_sanitise(item) for item in value] if isinstance(value, (str, int, float, bool)) or value is None: return value if isinstance(value, bytes): return value.decode("utf-8", errors="replace") return str(value)