Home Assistant 2026.9 deprecates the `via_device` (domain, identifier) tuple in favour of `via_device_id`, which takes the parent device's registry id; `via_device` stops working in 2027.8.0. Register the hub device explicitly in async_setup_entry before the platforms are forwarded, so the per-identity entities can look it up and pass its `.id`. The hub device info moves into a shared helper so both call sites describe the same device.
474 lines
16 KiB
Python
474 lines
16 KiB
Python
"""The Reticulum integration.
|
|
|
|
One config entry = the shared Reticulum stack (RNS + TCP interface).
|
|
Each config subentry (type ``identity``) = an independent LXMF identity.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from datetime import timedelta
|
|
from typing import Any
|
|
|
|
import voluptuous as vol
|
|
|
|
from homeassistant.config_entries import ConfigEntry, ConfigSubentry
|
|
from homeassistant.const import EVENT_HOMEASSISTANT_STOP, Platform
|
|
from homeassistant.core import (
|
|
HomeAssistant,
|
|
ServiceCall,
|
|
ServiceResponse,
|
|
SupportsResponse,
|
|
)
|
|
from homeassistant.exceptions import ConfigEntryNotReady, HomeAssistantError
|
|
from homeassistant.helpers import config_validation as cv
|
|
from homeassistant.helpers import device_registry as dr
|
|
from homeassistant.helpers.event import async_track_time_interval
|
|
|
|
from .const import (
|
|
ATTR_CONTENT,
|
|
ATTR_DESTINATION,
|
|
ATTR_FIELDS,
|
|
ATTR_IDENTITY,
|
|
ATTR_MAX_MESSAGES,
|
|
ATTR_METHOD,
|
|
ATTR_TITLE,
|
|
CONF_ALLOW_ALL,
|
|
CONF_ALLOWED_IDENTITIES,
|
|
CONF_ANNOUNCE_INTERVAL,
|
|
CONF_ASSIST_AGENT,
|
|
CONF_ASSIST_LANGUAGE,
|
|
CONF_DISPLAY_NAME,
|
|
CONF_ENABLE_ASSIST,
|
|
CONF_GREETING,
|
|
CONF_INTERFACE_NAME,
|
|
CONF_LOGLEVEL,
|
|
CONF_PROPAGATION_NODE,
|
|
CONF_SYNC_INTERVAL,
|
|
CONF_TARGET_HOST,
|
|
CONF_TARGET_PORT,
|
|
DEFAULT_ANNOUNCE_INTERVAL,
|
|
DEFAULT_DISPLAY_NAME,
|
|
DEFAULT_INTERFACE_NAME,
|
|
DEFAULT_LOGLEVEL,
|
|
DEFAULT_SYNC_INTERVAL,
|
|
DELIVERY_METHODS,
|
|
DOMAIN,
|
|
EVENT_MESSAGE_RECEIVED,
|
|
SERVICE_ANNOUNCE,
|
|
SERVICE_REQUEST_PATH,
|
|
SERVICE_SEND_MESSAGE,
|
|
SERVICE_SET_PROPAGATION_NODE,
|
|
SERVICE_SYNC_PROPAGATION,
|
|
STORAGE_SUBDIR,
|
|
SUBENTRY_TYPE_IDENTITY,
|
|
)
|
|
from .entity import hub_device_info
|
|
from .reticulum_client import (
|
|
IdentityManager,
|
|
ReticulumError,
|
|
ReticulumStack,
|
|
)
|
|
|
|
_LOGGER = logging.getLogger(__name__)
|
|
|
|
PLATFORMS: list[Platform] = [
|
|
Platform.BINARY_SENSOR,
|
|
Platform.BUTTON,
|
|
Platform.NOTIFY,
|
|
Platform.SENSOR,
|
|
]
|
|
|
|
INTERFACE_POLL = timedelta(seconds=15)
|
|
|
|
|
|
class RuntimeData:
|
|
"""Runtime objects for the hub entry."""
|
|
|
|
def __init__(self, stack: ReticulumStack) -> None:
|
|
self.stack = stack
|
|
self.identities: dict[str, IdentityManager] = {}
|
|
self.unsubs: list[Any] = []
|
|
# f"{subentry_id}:{source}" -> HA conversation_id
|
|
self.conversation_ids: dict[str, str] = {}
|
|
|
|
|
|
async def async_setup_entry(hass: HomeAssistant, entry: ConfigEntry) -> bool:
|
|
"""Set up the Reticulum stack and its identities."""
|
|
storage_dir = hass.config.path(STORAGE_SUBDIR)
|
|
stack = ReticulumStack(
|
|
hass,
|
|
storage_dir,
|
|
entry.entry_id,
|
|
target_host=entry.data[CONF_TARGET_HOST],
|
|
target_port=entry.data[CONF_TARGET_PORT],
|
|
interface_name=entry.data.get(CONF_INTERFACE_NAME, DEFAULT_INTERFACE_NAME),
|
|
loglevel=int(entry.options.get(CONF_LOGLEVEL, DEFAULT_LOGLEVEL)),
|
|
)
|
|
try:
|
|
await stack.async_start()
|
|
except Exception as err: # noqa: BLE001
|
|
raise ConfigEntryNotReady(f"Could not start Reticulum: {err}") from err
|
|
|
|
runtime = RuntimeData(stack)
|
|
hass.data.setdefault(DOMAIN, {})[entry.entry_id] = runtime
|
|
|
|
# Start each identity subentry.
|
|
for subentry in entry.subentries.values():
|
|
if subentry.subentry_type != SUBENTRY_TYPE_IDENTITY:
|
|
continue
|
|
display_name = subentry.data.get(CONF_DISPLAY_NAME) or subentry.title
|
|
manager = IdentityManager(
|
|
hass, stack, subentry.subentry_id, display_name=display_name
|
|
)
|
|
try:
|
|
await manager.async_start()
|
|
except Exception as err: # noqa: BLE001
|
|
raise ConfigEntryNotReady(
|
|
f"Could not start identity '{display_name}': {err}"
|
|
) from err
|
|
manager.incoming_handler = _make_incoming_handler(
|
|
hass, entry, subentry, runtime, manager
|
|
)
|
|
runtime.identities[subentry.subentry_id] = manager
|
|
|
|
# Register the hub device before the platforms, so the per-identity devices can
|
|
# link to it with `via_device_id`.
|
|
dr.async_get(hass).async_get_or_create(
|
|
config_entry_id=entry.entry_id, **hub_device_info(entry)
|
|
)
|
|
|
|
await hass.config_entries.async_forward_entry_setups(entry, PLATFORMS)
|
|
|
|
# Per-identity announce + periodic schedules.
|
|
for subentry in entry.subentries.values():
|
|
if subentry.subentry_type != SUBENTRY_TYPE_IDENTITY:
|
|
continue
|
|
manager = runtime.identities.get(subentry.subentry_id)
|
|
if manager is None:
|
|
continue
|
|
await _safe_announce(manager)
|
|
_schedule_identity(hass, runtime, manager, subentry)
|
|
|
|
# Hub interface-status poll.
|
|
runtime.unsubs.append(
|
|
async_track_time_interval(
|
|
hass, stack.refresh_interface_status, INTERFACE_POLL
|
|
)
|
|
)
|
|
|
|
async def _on_ha_stop(_event) -> None:
|
|
await stack.async_persist()
|
|
|
|
runtime.unsubs.append(
|
|
hass.bus.async_listen_once(EVENT_HOMEASSISTANT_STOP, _on_ha_stop)
|
|
)
|
|
entry.async_on_unload(entry.add_update_listener(_async_entry_updated))
|
|
_async_register_services(hass)
|
|
return True
|
|
|
|
|
|
async def async_unload_entry(hass: HomeAssistant, entry: ConfigEntry) -> bool:
|
|
"""Unload the hub entry."""
|
|
unload_ok = await hass.config_entries.async_unload_platforms(entry, PLATFORMS)
|
|
runtime: RuntimeData | None = hass.data.get(DOMAIN, {}).get(entry.entry_id)
|
|
if runtime is not None:
|
|
for unsub in runtime.unsubs:
|
|
unsub()
|
|
for manager in runtime.identities.values():
|
|
await manager.async_stop()
|
|
await runtime.stack.async_stop()
|
|
if unload_ok:
|
|
hass.data[DOMAIN].pop(entry.entry_id, None)
|
|
if not hass.data[DOMAIN]:
|
|
_async_unregister_services(hass)
|
|
return unload_ok
|
|
|
|
|
|
async def _async_entry_updated(hass: HomeAssistant, entry: ConfigEntry) -> None:
|
|
"""Reload when options or subentries change."""
|
|
await hass.config_entries.async_reload(entry.entry_id)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Periodic scheduling
|
|
# ---------------------------------------------------------------------------
|
|
def _schedule_identity(
|
|
hass: HomeAssistant,
|
|
runtime: RuntimeData,
|
|
manager: IdentityManager,
|
|
subentry: ConfigSubentry,
|
|
) -> None:
|
|
data = subentry.data
|
|
announce_interval = data.get(CONF_ANNOUNCE_INTERVAL, DEFAULT_ANNOUNCE_INTERVAL)
|
|
if announce_interval and announce_interval > 0:
|
|
|
|
async def _do_announce(_now, _mgr=manager) -> None:
|
|
await _safe_announce(_mgr)
|
|
|
|
runtime.unsubs.append(
|
|
async_track_time_interval(
|
|
hass, _do_announce, timedelta(seconds=announce_interval)
|
|
)
|
|
)
|
|
|
|
propagation_node = data.get(CONF_PROPAGATION_NODE)
|
|
if propagation_node:
|
|
hass.async_create_task(_set_propagation(manager, propagation_node))
|
|
sync_interval = data.get(CONF_SYNC_INTERVAL, DEFAULT_SYNC_INTERVAL)
|
|
if sync_interval and sync_interval > 0:
|
|
|
|
async def _do_sync(_now, _mgr=manager) -> None:
|
|
try:
|
|
await _mgr.async_sync_propagation()
|
|
except ReticulumError as err:
|
|
_LOGGER.warning("Propagation sync failed: %s", err)
|
|
|
|
runtime.unsubs.append(
|
|
async_track_time_interval(
|
|
hass, _do_sync, timedelta(seconds=sync_interval)
|
|
)
|
|
)
|
|
|
|
|
|
async def _set_propagation(manager: IdentityManager, node: str) -> None:
|
|
try:
|
|
await manager.async_set_propagation_node(node)
|
|
except ReticulumError as err:
|
|
_LOGGER.warning("Could not set propagation node: %s", err)
|
|
|
|
|
|
async def _safe_announce(manager: IdentityManager) -> None:
|
|
try:
|
|
await manager.async_announce()
|
|
except ReticulumError as err:
|
|
_LOGGER.debug("Announce skipped: %s", err)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Incoming message pipeline (per identity)
|
|
# ---------------------------------------------------------------------------
|
|
def _make_incoming_handler(
|
|
hass: HomeAssistant,
|
|
entry: ConfigEntry,
|
|
subentry: ConfigSubentry,
|
|
runtime: RuntimeData,
|
|
manager: IdentityManager,
|
|
):
|
|
async def _handle(payload: dict[str, Any]) -> None:
|
|
source = payload["source"]
|
|
hass.bus.async_fire(EVENT_MESSAGE_RECEIVED, payload)
|
|
|
|
data = subentry.data
|
|
allow_all = data.get(CONF_ALLOW_ALL, True)
|
|
allowed = {
|
|
a.strip().lower()
|
|
for a in data.get(CONF_ALLOWED_IDENTITIES, [])
|
|
if a and a.strip()
|
|
}
|
|
if not allow_all and source.lower() not in allowed:
|
|
_LOGGER.debug("Ignoring message from non-allowed sender %s", source)
|
|
return
|
|
|
|
# When not routing to Assist, optionally auto-reply with the greeting.
|
|
if not data.get(CONF_ENABLE_ASSIST, True):
|
|
greeting = (data.get(CONF_GREETING) or "").strip()
|
|
if greeting:
|
|
await _reply(manager, source, greeting, "Home Assistant")
|
|
return
|
|
|
|
text = (payload.get("content") or "").strip()
|
|
if not text:
|
|
return
|
|
await _run_assist(hass, subentry, runtime, manager, source, text)
|
|
|
|
return _handle
|
|
|
|
|
|
async def _run_assist(
|
|
hass: HomeAssistant,
|
|
subentry: ConfigSubentry,
|
|
runtime: RuntimeData,
|
|
manager: IdentityManager,
|
|
source: str,
|
|
text: str,
|
|
) -> None:
|
|
from homeassistant.components import conversation # noqa: PLC0415
|
|
from homeassistant.core import Context # noqa: PLC0415
|
|
|
|
data = subentry.data
|
|
agent_id = data.get(CONF_ASSIST_AGENT) or None
|
|
language = data.get(CONF_ASSIST_LANGUAGE) or hass.config.language
|
|
conv_key = f"{subentry.subentry_id}:{source}"
|
|
|
|
try:
|
|
result = await conversation.async_converse(
|
|
hass=hass,
|
|
text=text,
|
|
conversation_id=runtime.conversation_ids.get(conv_key),
|
|
context=Context(),
|
|
language=language,
|
|
agent_id=agent_id,
|
|
)
|
|
except Exception as err: # noqa: BLE001
|
|
_LOGGER.exception("Assist conversation failed")
|
|
await _reply(manager, source, f"⚠️ Assist error: {err}", "Error")
|
|
return
|
|
|
|
conv_id = getattr(result, "conversation_id", None)
|
|
if conv_id:
|
|
runtime.conversation_ids[conv_key] = conv_id
|
|
|
|
reply = _extract_speech(result)
|
|
if reply:
|
|
await _reply(manager, source, reply, "Assist")
|
|
|
|
|
|
def _extract_speech(result: Any) -> str | None:
|
|
try:
|
|
speech = result.response.speech
|
|
if isinstance(speech, dict):
|
|
return speech.get("plain", {}).get("speech")
|
|
except AttributeError:
|
|
return None
|
|
return None
|
|
|
|
|
|
async def _reply(
|
|
manager: IdentityManager, destination: str, content: str, title: str
|
|
) -> None:
|
|
try:
|
|
await manager.async_send_message(destination, content, title=title)
|
|
except ReticulumError as err:
|
|
_LOGGER.warning("Could not send reply to %s: %s", destination, err)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Services
|
|
# ---------------------------------------------------------------------------
|
|
SEND_MESSAGE_SCHEMA = vol.Schema(
|
|
{
|
|
vol.Required(ATTR_DESTINATION): cv.string,
|
|
vol.Required(ATTR_CONTENT): cv.string,
|
|
vol.Optional(ATTR_TITLE, default=""): cv.string,
|
|
vol.Optional(ATTR_METHOD, default="direct"): vol.In(DELIVERY_METHODS),
|
|
vol.Optional(ATTR_FIELDS): dict,
|
|
vol.Optional(ATTR_IDENTITY): cv.string,
|
|
}
|
|
)
|
|
REQUEST_PATH_SCHEMA = vol.Schema(
|
|
{vol.Required(ATTR_DESTINATION): cv.string, vol.Optional(ATTR_IDENTITY): cv.string}
|
|
)
|
|
SET_PROPAGATION_SCHEMA = vol.Schema(
|
|
{vol.Required(ATTR_DESTINATION): cv.string, vol.Optional(ATTR_IDENTITY): cv.string}
|
|
)
|
|
SYNC_SCHEMA = vol.Schema(
|
|
{
|
|
vol.Optional(ATTR_MAX_MESSAGES): vol.All(vol.Coerce(int), vol.Range(min=1)),
|
|
vol.Optional(ATTR_IDENTITY): cv.string,
|
|
}
|
|
)
|
|
ANNOUNCE_SCHEMA = vol.Schema({vol.Optional(ATTR_IDENTITY): cv.string})
|
|
|
|
|
|
def _resolve_manager(hass: HomeAssistant, identity: str | None) -> IdentityManager:
|
|
"""Pick an identity manager by name/address, or the first one."""
|
|
data: dict[str, RuntimeData] = hass.data.get(DOMAIN, {})
|
|
managers: list[IdentityManager] = []
|
|
for runtime in data.values():
|
|
managers.extend(runtime.identities.values())
|
|
if not managers:
|
|
raise HomeAssistantError("No Reticulum identity is configured")
|
|
if identity:
|
|
key = identity.strip().lower()
|
|
for manager in managers:
|
|
if (
|
|
manager.display_name.lower() == key
|
|
or (manager.state.lxmf_address or "").lower() == key
|
|
):
|
|
return manager
|
|
raise HomeAssistantError(f"No Reticulum identity matches '{identity}'")
|
|
return managers[0]
|
|
|
|
|
|
def _async_register_services(hass: HomeAssistant) -> None:
|
|
if hass.services.has_service(DOMAIN, SERVICE_SEND_MESSAGE):
|
|
return
|
|
|
|
async def _send(call: ServiceCall) -> ServiceResponse:
|
|
manager = _resolve_manager(hass, call.data.get(ATTR_IDENTITY))
|
|
try:
|
|
message_hash = await manager.async_send_message(
|
|
call.data[ATTR_DESTINATION],
|
|
call.data[ATTR_CONTENT],
|
|
title=call.data.get(ATTR_TITLE, ""),
|
|
method=call.data.get(ATTR_METHOD, "direct"),
|
|
fields=call.data.get(ATTR_FIELDS),
|
|
)
|
|
except ReticulumError as err:
|
|
raise HomeAssistantError(str(err)) from err
|
|
return {"message_hash": message_hash, "identity": manager.display_name}
|
|
|
|
async def _announce(call: ServiceCall) -> None:
|
|
manager = _resolve_manager(hass, call.data.get(ATTR_IDENTITY))
|
|
try:
|
|
await manager.async_announce()
|
|
except ReticulumError as err:
|
|
raise HomeAssistantError(str(err)) from err
|
|
|
|
async def _request_path(call: ServiceCall) -> None:
|
|
manager = _resolve_manager(hass, call.data.get(ATTR_IDENTITY))
|
|
try:
|
|
await manager.async_request_path(call.data[ATTR_DESTINATION])
|
|
except ReticulumError as err:
|
|
raise HomeAssistantError(str(err)) from err
|
|
|
|
async def _set_propagation(call: ServiceCall) -> None:
|
|
manager = _resolve_manager(hass, call.data.get(ATTR_IDENTITY))
|
|
try:
|
|
await manager.async_set_propagation_node(call.data[ATTR_DESTINATION])
|
|
except ReticulumError as err:
|
|
raise HomeAssistantError(str(err)) from err
|
|
|
|
async def _sync(call: ServiceCall) -> None:
|
|
manager = _resolve_manager(hass, call.data.get(ATTR_IDENTITY))
|
|
try:
|
|
await manager.async_sync_propagation(call.data.get(ATTR_MAX_MESSAGES))
|
|
except ReticulumError as err:
|
|
raise HomeAssistantError(str(err)) from err
|
|
|
|
hass.services.async_register(
|
|
DOMAIN,
|
|
SERVICE_SEND_MESSAGE,
|
|
_send,
|
|
schema=SEND_MESSAGE_SCHEMA,
|
|
supports_response=SupportsResponse.OPTIONAL,
|
|
)
|
|
hass.services.async_register(
|
|
DOMAIN, SERVICE_ANNOUNCE, _announce, schema=ANNOUNCE_SCHEMA
|
|
)
|
|
hass.services.async_register(
|
|
DOMAIN, SERVICE_REQUEST_PATH, _request_path, schema=REQUEST_PATH_SCHEMA
|
|
)
|
|
hass.services.async_register(
|
|
DOMAIN,
|
|
SERVICE_SET_PROPAGATION_NODE,
|
|
_set_propagation,
|
|
schema=SET_PROPAGATION_SCHEMA,
|
|
)
|
|
hass.services.async_register(
|
|
DOMAIN, SERVICE_SYNC_PROPAGATION, _sync, schema=SYNC_SCHEMA
|
|
)
|
|
|
|
|
|
def _async_unregister_services(hass: HomeAssistant) -> None:
|
|
for service in (
|
|
SERVICE_SEND_MESSAGE,
|
|
SERVICE_ANNOUNCE,
|
|
SERVICE_REQUEST_PATH,
|
|
SERVICE_SET_PROPAGATION_NODE,
|
|
SERVICE_SYNC_PROPAGATION,
|
|
):
|
|
if hass.services.has_service(DOMAIN, service):
|
|
hass.services.async_remove(DOMAIN, service)
|