commit de0975673149ebd251db0e544075316852b67434 Author: Ben Nicholson Date: Wed Aug 12 15:31:50 2026 +1000 Upload files to "/" diff --git a/README.md b/README.md new file mode 100644 index 0000000..ef9f970 --- /dev/null +++ b/README.md @@ -0,0 +1,302 @@ +# Riedel RRCS for Home Assistant + +A custom integration that talks to **Riedel Router Control Software (RRCS)**, the +third-party gateway to an Artist intercom system. Built against the RRCS +8.9.1.Rev1 interface specification. + +It gives you logic sources and GPIOs as real Home Assistant entities, alarms and +panel events on the event bus, and services for crosspoints, keys, labels and +gains. + +> **Integration, not add-on.** Home Assistant add-ons are Supervisor-only Docker +> containers and cannot create entities. This is a custom integration, so it +> works the same on HAOS, Home Assistant Container and Core. + +--- + +## What you get + +### Entities + +| Platform | Entity | Source | +| --- | --- | --- | +| `switch` | One per logic source | `GetAllLogicSources_v2` / `SetLogicSourceState` | +| `switch` | One per GP output | `GetGpOutputState` / `SetGpOutput` | +| `binary_sensor` | One per GP input | `GetGpInputState` | +| `binary_sensor` | Artist connection | `IsConnectedToArtist` | +| `sensor` | Gateway state (Working / Standby) | `GetState` | +| `sensor` | Active crosspoints | `GetAllActiveXpsCount` | +| `sensor` | RRCS version, logic source count | diagnostic, disabled by default | +| `event` | Alarm | ring, node, client and port alarms (ch. 9.7) | +| `event` | Send string | `SendString` / `SendStringOff` | +| `event` | Panel spy | panel spy events, disabled by default | + +### Services + +`set_xp`, `kill_xp`, `set_xp_volume`, `set_gp_output`, `set_logic_source`, +`press_key`, `set_key_label`, `clear_key_label`, `set_port_alias`, +`set_input_gain`, `set_output_gain`, and `call_method` as an escape hatch to any +RRCS method that has no dedicated service. + +### Events + +Every inbound notification is fired on the bus as `riedel_rrcs_notification` +with `entry_id`, `method` and `params`. + +--- + +## Installation + +### HACS + +Add this repository as a custom repository of type *Integration*, install it, +restart Home Assistant, then add **Riedel RRCS** from +*Settings → Devices & Services → Add Integration*. + +### Manual + +Copy `custom_components/riedel_rrcs/` into your Home Assistant `config/custom_components/` +directory and restart. + +--- + +## Configuration + +Everything is done in the UI. The initial dialog asks for: + +| Field | Notes | +| --- | --- | +| Host | The machine running RRCS | +| Port | 8193 by default | +| RPC path | `/` unless the gateway is behind a reverse proxy | +| Request timeout | Seconds | +| Transaction key prefix | A single character. Avoid `R`, which RRCS reserves for its own requests | + +The connection is tested with `GetVersion` before the entry is created. + +### Options + +| Option | Default | Notes | +| --- | --- | --- | +| Polling interval | 30 s | Raise it if push notifications are working | +| Receive push notifications | on | See below | +| Callback port | HA's HTTP port | Override if HA serves HTTPS | +| Discover GPIOs automatically | on | Requires RRCS 8.8 or later | +| Poll GPIO states | on | Turn off to reduce gateway load once push is confirmed | +| Manual GPIO list | empty | JSON, see below | + +--- + +## Push notifications + +With notifications enabled, the integration calls `RegisterForAllEvents` and +serves an endpoint at `/api/riedel_rrcs/`. RRCS then reverses the +roles and posts XML-RPC method calls at Home Assistant, so logic source, GPIO, +crosspoint and alarm changes arrive immediately rather than at the next poll. + +Three things to know: + +1. **The endpoint is unauthenticated.** RRCS has no way to present a bearer + token, so the path carries a random secret instead. It is only reachable from + whatever can already reach your Home Assistant HTTP port. +2. **RRCS derives the destination host from the source address of the + registration request.** The gateway has to be able to reach Home Assistant + back on that same IP, on the callback port, over **plain HTTP**. If Home + Assistant serves HTTPS you must set a plain-HTTP callback port in the + options; a warning is logged if this looks wrong. +3. **Answer or be dropped.** If KeepAlive is enabled in the RRCS route options, + RRCS sends `GetAlive` every three seconds of inactivity and tears the channel + down if it goes unanswered. This integration answers it. It also re-checks + the registration on every poll with `IsRegisteredForAllEvents` and + re-registers if RRCS has forgotten us, which it does after a restart. + +--- + +## GPIOs + +### Automatic discovery + +`GetAllGpIns` and `GetAllGpOuts` arrived in RRCS 8.8, but the specification +documents them by name only — the return struct is literally `...` on page 40. +Discovery here is a tolerant walk over whatever comes back, mapping the three +`TGPIOAddress` variants from chapter 6.7 onto the Net/Node/Port/Slot/Index +addressing that `Set`/`GetGpOutput` actually wants. + +That means the Bay-to-Slot mapping in particular is an educated guess. If the +discovered entities look wrong, download the diagnostics for the config entry: +the raw `GetAllGpIns` / `GetAllGpOuts` payloads are included verbatim. + +### Manual list + +Anything you define manually wins over discovery. The options flow takes a JSON +array: + +```json +[ + {"name": "Studio A tally", "direction": "in", "node": 2, "port": 128, "slot": 0, "index": 3}, + {"name": "Red light", "direction": "out", "node": 2, "port": 128, "slot": 0, "index": 5}, + {"name": "Panel GPI", "direction": "in", "node": 4, "port": 9, "index": 0} +] +``` + +`node` and `index` are required. Defaults are `net` 1, `port` 128, `slot` 0 and +`direction` `in`. + +Addressing follows chapter 6.9: + +- **GPIO client card** — `port` is 128 and `slot` selects the bay: 0–15 for bays + 1–16, 16 for bay X, 17 for bay Y, 20 for bay B. +- **Panel GPIO** — `port` is the panel's port address and `slot` is ignored. +- Port address for Artist 32/64/128 is `((slot - 1) * 8) + port - 1`, so port 5 + in slot 2 is address 12. + +--- + +## Examples + +### Toggle a logic source from a physical button + +This replaces `toggle_v5.vbs` — the logic source is just a switch now. + +```yaml +automation: + - alias: Toggle busy from desk button + triggers: + - trigger: state + entity_id: binary_sensor.desk_button + to: "on" + actions: + - action: switch.toggle + target: + entity_id: switch.rrcs_ben_busy +``` + +### Press a panel key, primary trigger + +The XML-RPC equivalent of `ras_pi_RRCS_BusyLight.py`, press and release: + +```yaml +script: + ben_busy_pulse: + sequence: + - action: riedel_rrcs.press_key + data: + node: 64 + port: 46 + page: 2 + expansion_panel: 0 + key_number: 1 + is_virt_key: true + trigger: 1 + press: true + - delay: "00:00:00.25" + - action: riedel_rrcs.press_key + data: + node: 64 + port: 46 + page: 2 + expansion_panel: 0 + key_number: 1 + is_virt_key: true + trigger: 1 + press: false +``` + +### Route audio on an event + +```yaml +- action: riedel_rrcs.set_xp + data: + source_node: 2 + source_port: 12 + dest_node: 3 + dest_port: 4 + priority: "2" +``` + +Set `destructive: true` to drop any existing route to the destination first. +Passing net, node and port all as 0 clears every route for a source or +destination. + +### React to a ring alarm + +```yaml +automation: + - alias: Notify on upstream failure + triggers: + - trigger: state + entity_id: event.rrcs_alarm + conditions: + - condition: template + value_template: >- + {{ trigger.to_state.attributes.event_type in + ['UpstreamFailed', 'DownstreamFailed', 'NodeControllerFailed'] }} + actions: + - action: notify.ntfy + data: + message: >- + RRCS {{ trigger.to_state.attributes.event_type }}: + {{ trigger.to_state.attributes.params }} +``` + +Or catch everything on the bus: + +```yaml +triggers: + - trigger: event + event_type: riedel_rrcs_notification + event_data: + method: PortInactive +``` + +### Call anything else + +```yaml +- action: riedel_rrcs.call_method + data: + method: GetAllPorts + response_variable: ports +- action: notify.persistent_notification + data: + message: "{{ ports.result[1] | count }} ports" +``` + +`call_method` prepends a transaction key for you unless you turn that off. The +decoded response comes back under `result`. + +--- + +## Notes and limitations + +- **RRCS only undoes its own work.** `SetLogicSourceState`, `SetGpOutput` and + `KillXp` only affect state that RRCS itself set. A logic source driven from a + panel key will not respond to "off", and `GetXpStatus` can still report true + after a `KillXp`. +- **One request at a time.** Requests are serialised behind a lock. RRCS is a + single Windows service in front of a live ring, and chapter 12's timing + figures (10 ms for queries, 50 ms for state changes, seconds for + configuration changes) all assume an otherwise idle gateway. Don't set the + polling interval aggressively low on a large system. +- **No configuration changes.** `ConfigurationChange` and friends are not + wrapped as services deliberately — they rewrite the Artist configuration and + are better done in Director. `call_method` will still let you if you mean it. +- **Standby gateways.** In a redundant pair, a standby gateway answers only the + redundancy command set, so most entities will be unavailable against one. + `GetState` and the connection sensor still work. +- **Crosspoint volume** does not work for destinations on an Artist 1024. +- Encoding: some gateways declare UTF-8 but emit single-byte port labels. The + parser strips the declaration and decodes leniently rather than dropping the + whole response. + +## Troubleshooting + +Turn on debug logging to see every RPC in both directions: + +```yaml +logger: + logs: + custom_components.riedel_rrcs: debug +``` + +Config entry diagnostics include the coordinator state, the resolved GPIO +addresses and the raw GPIO list payloads. diff --git a/__init__.py b/__init__.py new file mode 100644 index 0000000..001ecac --- /dev/null +++ b/__init__.py @@ -0,0 +1,594 @@ +"""The Riedel RRCS integration.""" + +from __future__ import annotations + +import logging +import secrets +from typing import Any + +import voluptuous as vol +from homeassistant.config_entries import ConfigEntry, ConfigEntryState +from homeassistant.const import CONF_HOST, CONF_PORT, CONF_SCAN_INTERVAL, Platform +from homeassistant.core import HomeAssistant, ServiceCall, ServiceResponse, SupportsResponse +from homeassistant.exceptions import ConfigEntryNotReady, HomeAssistantError, ServiceValidationError +from homeassistant.helpers import config_validation as cv +from homeassistant.helpers.aiohttp_client import async_get_clientsession + +from .const import ( + ATTR_ENTRY_ID, + CONF_CALLBACK_PORT, + CONF_CALLBACK_TOKEN, + CONF_DISCOVER_GPIO, + CONF_GPIO_ENTITIES, + CONF_NOTIFICATIONS, + CONF_POLL_GPIO, + CONF_RPC_PATH, + CONF_TIMEOUT, + CONF_TRANSKEY_PREFIX, + DEFAULT_DISCOVER_GPIO, + DEFAULT_NOTIFICATIONS, + DEFAULT_PATH, + DEFAULT_POLL_GPIO, + DEFAULT_PORT, + DEFAULT_SCAN_INTERVAL, + DEFAULT_TIMEOUT, + DEFAULT_TRANSKEY_PREFIX, + DOMAIN, + NOTIFICATION_URL_TEMPLATE, + SERVICE_CALL_METHOD, + SERVICE_CLEAR_KEY_LABEL, + SERVICE_KILL_XP, + SERVICE_PRESS_KEY, + SERVICE_SET_GP_OUTPUT, + SERVICE_SET_INPUT_GAIN, + SERVICE_SET_KEY_LABEL, + SERVICE_SET_LOGIC_SOURCE, + SERVICE_SET_OUTPUT_GAIN, + SERVICE_SET_PORT_ALIAS, + SERVICE_SET_XP, + SERVICE_SET_XP_VOLUME, +) +from .coordinator import RRCSCoordinator +from .models import RRCSConfigEntry, RRCSRuntimeData +from .notification import ( + RRCSNotificationDispatcher, + async_register_listener, + async_register_view, + async_remove_listener, +) +from .rrcs import ( + GpioAddress, + RRCSClient, + RRCSConnectionError, + RRCSError, + parse_gpio_config, +) + +_LOGGER = logging.getLogger(__name__) + +PLATFORMS: list[Platform] = [ + Platform.BINARY_SENSOR, + Platform.EVENT, + Platform.SENSOR, + Platform.SWITCH, +] + +CONFIG_SCHEMA = cv.config_entry_only_config_schema(DOMAIN) + + +async def async_setup(hass: HomeAssistant, config: dict[str, Any]) -> bool: + """Register the integration's services.""" + _async_register_services(hass) + return True + + +async def async_setup_entry(hass: HomeAssistant, entry: RRCSConfigEntry) -> bool: + """Set up Riedel RRCS from a config entry.""" + options = {**entry.data, **entry.options} + + client = RRCSClient( + session=async_get_clientsession(hass), + host=entry.data[CONF_HOST], + port=entry.data.get(CONF_PORT, DEFAULT_PORT), + path=options.get(CONF_RPC_PATH, DEFAULT_PATH), + timeout=options.get(CONF_TIMEOUT, DEFAULT_TIMEOUT), + transkey_prefix=options.get(CONF_TRANSKEY_PREFIX, DEFAULT_TRANSKEY_PREFIX), + ) + + # Fail fast and let HA retry rather than creating half a device. + try: + await client.get_version() + except RRCSError as err: + raise ConfigEntryNotReady(f"Cannot reach RRCS at {client.url}: {err}") from err + + coordinator = RRCSCoordinator( + hass, + entry, + client, + scan_interval=options.get(CONF_SCAN_INTERVAL, DEFAULT_SCAN_INTERVAL), + poll_gpio=options.get(CONF_POLL_GPIO, DEFAULT_POLL_GPIO), + ) + + gpio_inputs, gpio_outputs, gpio_names = await _async_collect_gpios(client, options) + coordinator.set_gpios(gpio_inputs, gpio_outputs, gpio_names) + + runtime = RRCSRuntimeData( + client=client, + coordinator=coordinator, + gpio_inputs=gpio_inputs, + gpio_outputs=gpio_outputs, + gpio_names=gpio_names, + ) + entry.runtime_data = runtime + + if options.get(CONF_NOTIFICATIONS, DEFAULT_NOTIFICATIONS): + await _async_start_notifications(hass, entry, runtime, options) + + await coordinator.async_config_entry_first_refresh() + await hass.config_entries.async_forward_entry_setups(entry, PLATFORMS) + entry.async_on_unload(entry.add_update_listener(_async_update_listener)) + return True + + +async def async_unload_entry(hass: HomeAssistant, entry: RRCSConfigEntry) -> bool: + """Unload a config entry.""" + runtime = entry.runtime_data + + if runtime.callback_token: + async_remove_listener(hass, runtime.callback_token) + if runtime.callback_port and runtime.callback_path: + try: + await runtime.client.unregister_for_all_events( + runtime.callback_port, runtime.callback_path + ) + except RRCSError as err: + _LOGGER.debug("Could not unregister notifications cleanly: %s", err) + + return await hass.config_entries.async_unload_platforms(entry, PLATFORMS) + + +async def _async_update_listener(hass: HomeAssistant, entry: RRCSConfigEntry) -> None: + """Reload the entry when its options change.""" + await hass.config_entries.async_reload(entry.entry_id) + + +async def _async_start_notifications( + hass: HomeAssistant, + entry: RRCSConfigEntry, + runtime: RRCSRuntimeData, + options: dict[str, Any], +) -> None: + """Stand up the callback endpoint and register it with the gateway.""" + token = entry.data.get(CONF_CALLBACK_TOKEN) + if not token: + token = secrets.token_hex(8) + hass.config_entries.async_update_entry( + entry, data={**entry.data, CONF_CALLBACK_TOKEN: token} + ) + + port = options.get(CONF_CALLBACK_PORT) or hass.http.server_port + path = NOTIFICATION_URL_TEMPLATE.format(token=token) + + if getattr(hass.config.api, "use_ssl", False) and not options.get(CONF_CALLBACK_PORT): + _LOGGER.warning( + "Home Assistant is serving HTTPS on port %s but RRCS pushes plain HTTP. " + "Set a plain-HTTP callback port in the integration options or " + "notifications will not arrive", + port, + ) + + async_register_view(hass) + dispatcher = RRCSNotificationDispatcher(hass, entry.entry_id, runtime.coordinator) + async_register_listener(hass, token, dispatcher.handle) + + try: + await runtime.client.register_for_all_events(port, path) + except RRCSError as err: + # Not fatal: polling still works, and the coordinator retries the + # registration on every refresh. + _LOGGER.warning("Could not register for RRCS notifications: %s", err) + + runtime.callback_token = token + runtime.callback_port = port + runtime.callback_path = path + runtime.coordinator.set_registration(port, path) + + +async def _async_collect_gpios( + client: RRCSClient, options: dict[str, Any] +) -> tuple[list[GpioAddress], list[GpioAddress], dict[str, str]]: + """Build the GPIO entity list from discovery plus manual configuration.""" + addresses: dict[str, GpioAddress] = {} + names: dict[str, str] = {} + + if options.get(CONF_DISCOVER_GPIO, DEFAULT_DISCOVER_GPIO): + try: + for address in await client.discover_gpios(): + addresses[address.key] = address + except RRCSError as err: + _LOGGER.debug("GPIO discovery failed: %s", err) + + for address, name in parse_gpio_config(options.get(CONF_GPIO_ENTITIES)): + addresses[address.key] = address + if name: + names[address.key] = name + + inputs = [address for address in addresses.values() if address.is_input] + outputs = [address for address in addresses.values() if not address.is_input] + return inputs, outputs, names + + +# --- Services ---------------------------------------------------------------------- + +_ENTRY_FIELD = {vol.Optional(ATTR_ENTRY_ID): cv.string} + +_XP_FIELDS = { + vol.Required("source_net", default=1): vol.Coerce(int), + vol.Required("source_node"): vol.Coerce(int), + vol.Required("source_port"): vol.Coerce(int), + vol.Required("dest_net", default=1): vol.Coerce(int), + vol.Required("dest_node"): vol.Coerce(int), + vol.Required("dest_port"): vol.Coerce(int), +} + +_KEY_FIELDS = { + vol.Required("node"): vol.Coerce(int), + vol.Required("port"): vol.Coerce(int), + vol.Optional("is_input", default=False): cv.boolean, + vol.Optional("page", default=1): vol.Coerce(int), + vol.Optional("expansion_panel", default=0): vol.Coerce(int), + vol.Required("key_number"): vol.Coerce(int), + vol.Optional("is_virt_key", default=False): cv.boolean, +} + +SET_XP_SCHEMA = vol.Schema( + { + **_ENTRY_FIELD, + **_XP_FIELDS, + vol.Optional("priority"): vol.All(vol.Coerce(int), vol.Range(min=0, max=4)), + vol.Optional("destructive", default=False): cv.boolean, + } +) + +KILL_XP_SCHEMA = vol.Schema({**_ENTRY_FIELD, **_XP_FIELDS}) + +SET_XP_VOLUME_SCHEMA = vol.Schema( + { + **_ENTRY_FIELD, + **_XP_FIELDS, + vol.Optional("single", default=True): cv.boolean, + vol.Optional("conference", default=False): cv.boolean, + vol.Required("volume"): vol.All(vol.Coerce(int), vol.Range(min=0, max=256)), + } +) + +SET_GP_OUTPUT_SCHEMA = vol.Schema( + { + **_ENTRY_FIELD, + vol.Optional("net", default=1): vol.Coerce(int), + vol.Required("node"): vol.Coerce(int), + vol.Optional("port", default=128): vol.Coerce(int), + vol.Optional("slot", default=0): vol.Coerce(int), + vol.Required("index"): vol.Coerce(int), + vol.Required("state"): cv.boolean, + } +) + +SET_LOGIC_SOURCE_SCHEMA = vol.Schema( + { + **_ENTRY_FIELD, + vol.Required("object_id"): vol.Coerce(int), + vol.Required("state"): cv.boolean, + } +) + +PRESS_KEY_SCHEMA = vol.Schema( + { + **_ENTRY_FIELD, + **_KEY_FIELDS, + vol.Optional("press", default=True): cv.boolean, + vol.Optional("trigger"): vol.All(vol.Coerce(int), vol.Range(min=1, max=2)), + vol.Optional("pool_port", default=0): vol.Coerce(int), + } +) + +SET_KEY_LABEL_SCHEMA = vol.Schema( + { + **_ENTRY_FIELD, + **_KEY_FIELDS, + vol.Required("label"): vol.All(cv.string, vol.Length(max=8)), + vol.Optional("marker"): vol.Coerce(int), + } +) + +CLEAR_KEY_LABEL_SCHEMA = vol.Schema( + {**_ENTRY_FIELD, **_KEY_FIELDS, vol.Optional("clear_marker", default=False): cv.boolean} +) + +SET_PORT_ALIAS_SCHEMA = vol.Schema( + { + **_ENTRY_FIELD, + vol.Optional("net", default=1): vol.Coerce(int), + vol.Required("node"): vol.Coerce(int), + vol.Required("port"): vol.Coerce(int), + vol.Required("alias"): vol.All(cv.string, vol.Length(max=8)), + vol.Optional("is_input", default=False): cv.boolean, + } +) + +_GAIN_SCHEMA = vol.Schema( + { + **_ENTRY_FIELD, + vol.Optional("net", default=1): vol.Coerce(int), + vol.Required("node"): vol.Coerce(int), + vol.Required("port"): vol.Coerce(int), + # Half-decibel steps; -128 mutes. + vol.Required("gain"): vol.All(vol.Coerce(int), vol.Range(min=-128, max=36)), + } +) + +CALL_METHOD_SCHEMA = vol.Schema( + { + **_ENTRY_FIELD, + vol.Required("method"): cv.string, + vol.Optional("params", default=list): vol.Any(list, dict), + vol.Optional("include_transkey", default=True): cv.boolean, + } +) + + +def _resolve_client(hass: HomeAssistant, call: ServiceCall) -> RRCSClient: + """Find the client a service call is aimed at.""" + entries = [ + entry + for entry in hass.config_entries.async_entries(DOMAIN) + if entry.state is ConfigEntryState.LOADED + ] + entry_id = call.data.get(ATTR_ENTRY_ID) + if entry_id: + for entry in entries: + if entry.entry_id == entry_id: + return entry.runtime_data.client + raise ServiceValidationError(f"No loaded RRCS config entry with id {entry_id}") + if not entries: + raise ServiceValidationError("No RRCS gateway is currently loaded") + if len(entries) > 1: + raise ServiceValidationError( + "Several RRCS gateways are configured; pass entry_id to pick one" + ) + return entries[0].runtime_data.client + + +def _key_args(data: dict[str, Any]) -> list[Any]: + """Build the shared key addressing arguments.""" + return [ + data["node"], + data["port"], + data["is_input"], + data["page"], + data["expansion_panel"], + data["key_number"], + data["is_virt_key"], + ] + + +def _async_register_services(hass: HomeAssistant) -> None: + """Register every RRCS service, once.""" + if hass.services.has_service(DOMAIN, SERVICE_SET_XP): + return + + async def _guard(coro) -> Any: + try: + return await coro + except RRCSConnectionError as err: + raise HomeAssistantError(f"RRCS unreachable: {err}") from err + except RRCSError as err: + raise HomeAssistantError(str(err)) from err + + async def set_xp(call: ServiceCall) -> None: + client = _resolve_client(hass, call) + data = call.data + args = [ + data["source_net"], + data["source_node"], + data["source_port"], + data["dest_net"], + data["dest_node"], + data["dest_port"], + ] + priority = data.get("priority") + if data.get("destructive"): + await _guard( + client.call("SetXpDestructive", *args, priority or 1, check=True) + ) + elif priority is not None: + await _guard(client.call("SetXpPrio", *args, priority, check=True)) + else: + await _guard(client.call("SetXp", *args, check=True)) + + async def kill_xp(call: ServiceCall) -> None: + client = _resolve_client(hass, call) + data = call.data + await _guard( + client.call( + "KillXp", + data["source_net"], + data["source_node"], + data["source_port"], + data["dest_net"], + data["dest_node"], + data["dest_port"], + check=True, + ) + ) + + async def set_xp_volume(call: ServiceCall) -> None: + client = _resolve_client(hass, call) + data = call.data + await _guard( + client.call( + "SetXpVolume", + data["source_net"], + data["source_node"], + data["source_port"], + data["dest_net"], + data["dest_node"], + data["dest_port"], + data["single"], + data["conference"], + data["volume"], + check=True, + ) + ) + + async def set_gp_output(call: ServiceCall) -> None: + client = _resolve_client(hass, call) + data = call.data + address = GpioAddress( + net=data["net"], + node=data["node"], + port=data["port"], + slot=data["slot"], + index=data["index"], + is_input=False, + ) + await _guard(client.set_gp_output(address, data["state"])) + + async def set_logic_source(call: ServiceCall) -> None: + client = _resolve_client(hass, call) + await _guard( + client.set_logic_source(call.data["object_id"], call.data["state"]) + ) + + async def press_key(call: ServiceCall) -> None: + client = _resolve_client(hass, call) + data = call.data + args = _key_args(data) + if data.get("trigger") is not None: + await _guard( + client.call( + "PressKeyEx", + *args, + data["press"], + data["trigger"], + data["pool_port"], + ) + ) + else: + await _guard( + client.call( + "PressKey", *args, data["press"], data["pool_port"], check=True + ) + ) + + async def set_key_label(call: ServiceCall) -> None: + client = _resolve_client(hass, call) + data = call.data + args = _key_args(data) + if data.get("marker") is not None: + await _guard( + client.call( + "SetKeyLabelAndMarker", + *args, + data["label"], + data["marker"], + check=True, + ) + ) + else: + await _guard(client.call("SetKeyLabel", *args, data["label"], check=True)) + + async def clear_key_label(call: ServiceCall) -> None: + client = _resolve_client(hass, call) + method = ( + "ClearKeyLabelAndMarker" if call.data["clear_marker"] else "ClearKeyLabel" + ) + await _guard(client.call(method, *_key_args(call.data), check=True)) + + async def set_port_alias(call: ServiceCall) -> None: + client = _resolve_client(hass, call) + data = call.data + await _guard( + client.call( + "SetPortAlias", + data["net"], + data["node"], + data["port"], + data["alias"], + data["is_input"], + check=True, + ) + ) + + def _gain_handler(method: str): + async def handler(call: ServiceCall) -> None: + client = _resolve_client(hass, call) + data = call.data + await _guard( + client.call( + method, data["net"], data["node"], data["port"], data["gain"], check=True + ) + ) + + return handler + + async def call_method(call: ServiceCall) -> ServiceResponse: + client = _resolve_client(hass, call) + params = call.data["params"] + if isinstance(params, dict): + params = list(params.values()) + if call.data["include_transkey"]: + result = await _guard(client.call(call.data["method"], *params)) + else: + result = await _guard( + client.call_raw(call.data["method"], tuple(params)) + ) + return {"result": _jsonable(result)} + + hass.services.async_register(DOMAIN, SERVICE_SET_XP, set_xp, SET_XP_SCHEMA) + hass.services.async_register(DOMAIN, SERVICE_KILL_XP, kill_xp, KILL_XP_SCHEMA) + hass.services.async_register( + DOMAIN, SERVICE_SET_XP_VOLUME, set_xp_volume, SET_XP_VOLUME_SCHEMA + ) + hass.services.async_register( + DOMAIN, SERVICE_SET_GP_OUTPUT, set_gp_output, SET_GP_OUTPUT_SCHEMA + ) + hass.services.async_register( + DOMAIN, SERVICE_SET_LOGIC_SOURCE, set_logic_source, SET_LOGIC_SOURCE_SCHEMA + ) + hass.services.async_register(DOMAIN, SERVICE_PRESS_KEY, press_key, PRESS_KEY_SCHEMA) + hass.services.async_register( + DOMAIN, SERVICE_SET_KEY_LABEL, set_key_label, SET_KEY_LABEL_SCHEMA + ) + hass.services.async_register( + DOMAIN, SERVICE_CLEAR_KEY_LABEL, clear_key_label, CLEAR_KEY_LABEL_SCHEMA + ) + hass.services.async_register( + DOMAIN, SERVICE_SET_PORT_ALIAS, set_port_alias, SET_PORT_ALIAS_SCHEMA + ) + hass.services.async_register( + DOMAIN, SERVICE_SET_INPUT_GAIN, _gain_handler("SetInputGain"), _GAIN_SCHEMA + ) + hass.services.async_register( + DOMAIN, SERVICE_SET_OUTPUT_GAIN, _gain_handler("SetOutputGain"), _GAIN_SCHEMA + ) + hass.services.async_register( + DOMAIN, + SERVICE_CALL_METHOD, + call_method, + CALL_METHOD_SCHEMA, + supports_response=SupportsResponse.OPTIONAL, + ) + + +def _jsonable(value: Any) -> Any: + """Make a decoded XML-RPC value safe to hand back as a service response.""" + if isinstance(value, dict): + return {str(key): _jsonable(item) for key, item in value.items()} + if isinstance(value, (list, tuple)): + return [_jsonable(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) diff --git a/config_flow.py b/config_flow.py new file mode 100644 index 0000000..e218361 --- /dev/null +++ b/config_flow.py @@ -0,0 +1,165 @@ +"""Config flow for the Riedel RRCS integration.""" + +from __future__ import annotations + +import logging +from typing import Any + +import voluptuous as vol +from homeassistant.config_entries import ( + ConfigEntry, + ConfigFlow, + ConfigFlowResult, + OptionsFlow, +) +from homeassistant.const import CONF_HOST, CONF_PORT, CONF_SCAN_INTERVAL +from homeassistant.core import callback +from homeassistant.helpers.aiohttp_client import async_get_clientsession +from homeassistant.helpers import selector + +from .const import ( + CONF_CALLBACK_PORT, + CONF_DISCOVER_GPIO, + CONF_GPIO_ENTITIES, + CONF_NOTIFICATIONS, + CONF_POLL_GPIO, + CONF_RPC_PATH, + CONF_TIMEOUT, + CONF_TRANSKEY_PREFIX, + DEFAULT_DISCOVER_GPIO, + DEFAULT_NOTIFICATIONS, + DEFAULT_PATH, + DEFAULT_POLL_GPIO, + DEFAULT_PORT, + DEFAULT_SCAN_INTERVAL, + DEFAULT_TIMEOUT, + DEFAULT_TRANSKEY_PREFIX, + DOMAIN, +) +from .rrcs import RRCSClient, RRCSError, parse_gpio_config + +_LOGGER = logging.getLogger(__name__) + +STEP_USER_SCHEMA = vol.Schema( + { + vol.Required(CONF_HOST): str, + vol.Required(CONF_PORT, default=DEFAULT_PORT): vol.Coerce(int), + vol.Optional(CONF_RPC_PATH, default=DEFAULT_PATH): str, + vol.Optional(CONF_TIMEOUT, default=DEFAULT_TIMEOUT): vol.Coerce(int), + vol.Optional(CONF_TRANSKEY_PREFIX, default=DEFAULT_TRANSKEY_PREFIX): vol.All( + str, vol.Length(min=1, max=1) + ), + } +) + + +class RRCSConfigFlow(ConfigFlow, domain=DOMAIN): + """Handle the config flow.""" + + VERSION = 1 + + async def async_step_user( + self, user_input: dict[str, Any] | None = None + ) -> ConfigFlowResult: + """Handle the initial step.""" + errors: dict[str, str] = {} + + if user_input is not None: + host = user_input[CONF_HOST] + port = user_input[CONF_PORT] + await self.async_set_unique_id(f"{host}:{port}") + self._abort_if_unique_id_configured() + + client = RRCSClient( + session=async_get_clientsession(self.hass), + host=host, + port=port, + path=user_input.get(CONF_RPC_PATH, DEFAULT_PATH), + timeout=user_input.get(CONF_TIMEOUT, DEFAULT_TIMEOUT), + transkey_prefix=user_input.get( + CONF_TRANSKEY_PREFIX, DEFAULT_TRANSKEY_PREFIX + ), + ) + try: + version = await client.get_version() + except RRCSError as err: + _LOGGER.debug("RRCS connection test failed: %s", err) + errors["base"] = "cannot_connect" + else: + title = f"RRCS {host}" + if version: + title = f"{title} ({version})" + return self.async_create_entry(title=title, data=user_input) + + return self.async_show_form( + step_id="user", data_schema=STEP_USER_SCHEMA, errors=errors + ) + + @staticmethod + @callback + def async_get_options_flow(config_entry: ConfigEntry) -> RRCSOptionsFlow: + """Return the options flow.""" + return RRCSOptionsFlow(config_entry) + + +class RRCSOptionsFlow(OptionsFlow): + """Handle the options flow.""" + + def __init__(self, config_entry: ConfigEntry) -> None: + """Initialise the options flow.""" + self._entry = config_entry + + async def async_step_init( + self, user_input: dict[str, Any] | None = None + ) -> ConfigFlowResult: + """Manage the options.""" + errors: dict[str, str] = {} + + if user_input is not None: + try: + parse_gpio_config(user_input.get(CONF_GPIO_ENTITIES)) + except ValueError as err: + _LOGGER.debug("Invalid GPIO list: %s", err) + errors[CONF_GPIO_ENTITIES] = "invalid_gpio_list" + else: + return self.async_create_entry(data=user_input) + + current = {**self._entry.data, **self._entry.options} + + schema = vol.Schema( + { + vol.Optional( + CONF_SCAN_INTERVAL, + default=current.get(CONF_SCAN_INTERVAL, DEFAULT_SCAN_INTERVAL), + ): vol.All(vol.Coerce(int), vol.Range(min=5, max=3600)), + vol.Optional( + CONF_NOTIFICATIONS, + default=current.get(CONF_NOTIFICATIONS, DEFAULT_NOTIFICATIONS), + ): bool, + vol.Optional( + CONF_CALLBACK_PORT, + description={"suggested_value": current.get(CONF_CALLBACK_PORT)}, + ): vol.Any(None, vol.Coerce(int)), + vol.Optional( + CONF_DISCOVER_GPIO, + default=current.get(CONF_DISCOVER_GPIO, DEFAULT_DISCOVER_GPIO), + ): bool, + vol.Optional( + CONF_POLL_GPIO, + default=current.get(CONF_POLL_GPIO, DEFAULT_POLL_GPIO), + ): bool, + vol.Optional( + CONF_GPIO_ENTITIES, + description={ + "suggested_value": current.get(CONF_GPIO_ENTITIES, "") + }, + ): selector.TextSelector( + selector.TextSelectorConfig(multiline=True) + ), + vol.Optional( + CONF_TIMEOUT, default=current.get(CONF_TIMEOUT, DEFAULT_TIMEOUT) + ): vol.All(vol.Coerce(int), vol.Range(min=1, max=120)), + } + ) + + return self.async_show_form(step_id="init", data_schema=schema, errors=errors) diff --git a/coordinator.py b/coordinator.py new file mode 100644 index 0000000..16ec8a7 --- /dev/null +++ b/coordinator.py @@ -0,0 +1,163 @@ +"""Polling coordinator for the Riedel RRCS integration.""" + +from __future__ import annotations + +import logging +from dataclasses import dataclass, field +from datetime import timedelta + +from homeassistant.config_entries import ConfigEntry +from homeassistant.core import HomeAssistant +from homeassistant.helpers.update_coordinator import DataUpdateCoordinator, UpdateFailed + +from .const import DOMAIN +from .rrcs import GpioAddress, LogicSource, RRCSClient, RRCSError + +_LOGGER = logging.getLogger(__name__) + + +@dataclass +class RRCSData: + """Everything the entities read from.""" + + connected: bool = False + gateway_state: str | None = None + version: str | None = None + active_xps: int | None = None + logic_sources: dict[int, LogicSource] = field(default_factory=dict) + gpio_inputs: dict[str, bool] = field(default_factory=dict) + gpio_outputs: dict[str, bool] = field(default_factory=dict) + + +class RRCSCoordinator(DataUpdateCoordinator[RRCSData]): + """Polls the gateway and folds in pushed notifications between polls.""" + + def __init__( + self, + hass: HomeAssistant, + entry: ConfigEntry, + client: RRCSClient, + scan_interval: int, + poll_gpio: bool, + ) -> None: + """Initialise the coordinator.""" + super().__init__( + hass, + _LOGGER, + name=f"{DOMAIN} {client.host}", + update_interval=timedelta(seconds=scan_interval), + config_entry=entry, + ) + self.client = client + self.poll_gpio = poll_gpio + self.gpio_inputs: list[GpioAddress] = [] + self.gpio_outputs: list[GpioAddress] = [] + self.gpio_names: dict[str, str] = {} + self._registration: tuple[int, str] | None = None + + def set_gpios( + self, + inputs: list[GpioAddress], + outputs: list[GpioAddress], + names: dict[str, str], + ) -> None: + """Record the GPIOs this entry exposes as entities.""" + self.gpio_inputs = inputs + self.gpio_outputs = outputs + self.gpio_names = names + + def set_registration(self, tcp_port: int, url_path: str) -> None: + """Remember the notification registration so it can be re-asserted.""" + self._registration = (tcp_port, url_path) + + async def _async_update_data(self) -> RRCSData: + """Fetch the current gateway and Artist state.""" + data = RRCSData() + try: + data.connected = await self.client.is_connected_to_artist() + data.gateway_state = await self.client.get_state() + + # The version never changes at runtime, so only ask once. + previous = self.data + data.version = previous.version if previous else None + if data.version is None: + data.version = await self.client.get_version() + + if data.connected: + data.logic_sources = await self.client.get_logic_sources() + data.active_xps = await self.client.get_active_xp_count() + + if self.poll_gpio: + for address in self.gpio_inputs: + data.gpio_inputs[address.key] = await self.client.get_gp_input_state( + address + ) + for address in self.gpio_outputs: + data.gpio_outputs[address.key] = ( + await self.client.get_gp_output_state(address) + ) + elif previous is not None: + data.gpio_inputs = dict(previous.gpio_inputs) + data.gpio_outputs = dict(previous.gpio_outputs) + elif previous is not None: + # Keep the last known picture rather than blanking every entity + # while the gateway is disconnected from the ring. + data.logic_sources = dict(previous.logic_sources) + data.gpio_inputs = dict(previous.gpio_inputs) + data.gpio_outputs = dict(previous.gpio_outputs) + except RRCSError as err: + raise UpdateFailed(str(err)) from err + + await self._async_check_registration() + return data + + async def _async_check_registration(self) -> None: + """Re-register for notifications if RRCS has forgotten us. + + RRCS drops a notification channel when it restarts, when the Artist + connection is re-established, or when a GetAlive goes unanswered, and it + does not tell us that it has. + """ + if self._registration is None: + return + tcp_port, url_path = self._registration + try: + if await self.client.is_registered_for_all_events(tcp_port, url_path): + return + _LOGGER.info("RRCS notification registration lost; re-registering") + await self.client.register_for_all_events(tcp_port, url_path) + except RRCSError as err: + _LOGGER.warning("Could not refresh RRCS notification registration: %s", err) + + # --- Push updates -------------------------------------------------------------- + + def apply_logic_source_change(self, object_id: int, state: bool) -> None: + """Fold a LogicSourceChange notification into the current data.""" + if self.data is None: + return + source = self.data.logic_sources.get(object_id) + if source is None: + # An object we have not enumerated yet; the next poll will pick it up. + return + self.data.logic_sources[object_id] = LogicSource( + object_id=source.object_id, + long_name=source.long_name, + label=source.label, + state=state, + ) + self.async_set_updated_data(self.data) + + def apply_gpio_change(self, address: GpioAddress, state: bool) -> None: + """Fold a GpInputChange / GpOutputChange notification into the data.""" + if self.data is None: + return + target = self.data.gpio_inputs if address.is_input else self.data.gpio_outputs + target[address.key] = state + self.async_set_updated_data(self.data) + + def apply_connection_change(self, connected: bool) -> None: + """Fold a ConnectArtistFailure / ConnectArtistRestored notification in.""" + if self.data is None: + return + self.data.connected = connected + self.async_set_updated_data(self.data) diff --git a/notification.py b/notification.py new file mode 100644 index 0000000..71ed01d --- /dev/null +++ b/notification.py @@ -0,0 +1,170 @@ +"""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)