Files
homelable/backend/app/services/zigbee_service.py
T
pranjal-joshi 103e24e5fa feat: add Zigbee2MQTT network map importer
- Backend: async MQTT service (aiomqtt) to fetch Z2M networkmap via bridge API
- Backend: FastAPI router at /api/v1/zigbee with /import and /test-connection
- Backend: Pydantic v2 schemas for request/response validation
- Backend: coordinator → router → end-device parent_id hierarchy builder
- Frontend: ZigbeeImportModal with MQTT config form, Test Connection, Fetch Devices
- Frontend: device list grouped by type (coordinator/router/enddevice) with checkboxes
- Frontend: ZigbeeCoordinatorNode, ZigbeeRouterNode, ZigbeeEndDeviceNode canvas nodes
- Frontend: Zigbee Import button in sidebar alongside Scan Network
- Frontend: handleZigbeeAddToCanvas wires selected devices + edges onto canvas
- Tests: full unit test suite for parser, hierarchy builder, MQTT mocks
- Tests: API endpoint tests for /zigbee/import and /zigbee/test-connection
- Tests: Vitest component tests for ZigbeeImportModal
- Docs: docs/zigbee-import.md with full usage, MQTT config, troubleshooting guide
- Docs: README.md Zigbee2MQTT Import section

Co-authored-by: CyberKeys <noreply@openclaw.ai>
2026-05-04 13:58:58 +00:00

233 lines
8.3 KiB
Python

"""Zigbee2MQTT service: connects to MQTT broker and fetches the network map."""
from __future__ import annotations
import asyncio
import json
import logging
logger = logging.getLogger(__name__)
_NETWORKMAP_REQUEST_TOPIC = "{base_topic}/bridge/request/networkmap"
_NETWORKMAP_RESPONSE_TOPIC = "{base_topic}/bridge/response/networkmap"
_CONNECTION_TIMEOUT = 5.0 # seconds to verify broker reachability
_NETWORKMAP_TIMEOUT = 10.0 # seconds to wait for the networkmap response
def _z2m_type_to_homelable(device_type: str) -> str:
"""Map a Z2M device type string to a homelable node type."""
mapping = {
"Coordinator": "zigbee_coordinator",
"Router": "zigbee_router",
"EndDevice": "zigbee_enddevice",
}
return mapping.get(device_type, "zigbee_enddevice")
def parse_networkmap(payload: dict) -> tuple[list[dict], list[dict]]:
"""Parse a Z2M networkmap response payload into node + edge lists.
Returns:
(nodes, edges) where each node/edge is a plain dict with the fields
expected by ZigbeeNodeOut / ZigbeeEdgeOut.
"""
data = payload.get("data", {})
routes = data.get("routes", [])
nodes_list: list[dict] = []
edges_list: list[dict] = []
seen_ids: set[str] = set()
# Coordinator is always present; find it first so we can wire the hierarchy
coordinator_id: str | None = None
for route in routes:
source = route.get("source", {})
if not source:
continue
ieee = source.get("ieeeAddr") or source.get("ieee_address") or ""
if not ieee:
continue
device_type: str = source.get("type", "EndDevice")
friendly_name: str = source.get("friendlyName") or source.get("friendly_name") or ieee
model: str | None = source.get("modelID") or source.get("model")
vendor: str | None = source.get("vendor")
description: str | None = source.get("description")
if ieee not in seen_ids:
seen_ids.add(ieee)
node_type = _z2m_type_to_homelable(device_type)
node: dict = {
"id": ieee,
"label": friendly_name,
"type": node_type,
"ieee_address": ieee,
"friendly_name": friendly_name,
"device_type": device_type,
"model": model,
"vendor": vendor,
"lqi": None,
"parent_id": None,
}
nodes_list.append(node)
if device_type == "Coordinator":
coordinator_id = ieee
# Walk the route targets to build edges and collect additional nodes
targets = route.get("routes", [])
for target_entry in targets:
target_ieee = target_entry.get("target", {}).get("ieeeAddr") or target_entry.get("target", {}).get("ieee_address") or ""
lqi: int | None = target_entry.get("lqi")
if not target_ieee:
continue
if target_ieee not in seen_ids:
seen_ids.add(target_ieee)
t_source = target_entry.get("target", {})
t_type: str = t_source.get("type", "EndDevice")
t_fn: str = t_source.get("friendlyName") or t_source.get("friendly_name") or target_ieee
t_model: str | None = t_source.get("modelID") or t_source.get("model")
t_vendor: str | None = t_source.get("vendor")
t_node: dict = {
"id": target_ieee,
"label": t_fn,
"type": _z2m_type_to_homelable(t_type),
"ieee_address": target_ieee,
"friendly_name": t_fn,
"device_type": t_type,
"model": t_model,
"vendor": t_vendor,
"lqi": lqi,
"parent_id": None,
}
nodes_list.append(t_node)
edges_list.append({"source": ieee, "target": target_ieee})
# Build parent_id hierarchy: coordinator → routers → end devices
if coordinator_id:
router_ids = {n["id"] for n in nodes_list if n["device_type"] == "Router"}
for node in nodes_list:
if node["device_type"] == "Router":
node["parent_id"] = coordinator_id
elif node["device_type"] == "EndDevice":
# Try to find the nearest router from the edge list
parent = _find_parent_router(node["id"], router_ids, edges_list)
node["parent_id"] = parent or coordinator_id
return nodes_list, edges_list
def _find_parent_router(
device_id: str,
router_ids: set[str],
edges: list[dict],
) -> str | None:
"""Return the first router that has a direct edge to device_id."""
for edge in edges:
if edge["target"] == device_id and edge["source"] in router_ids:
return edge["source"]
if edge["source"] == device_id and edge["target"] in router_ids:
return edge["target"]
return None
async def fetch_networkmap(
mqtt_host: str,
mqtt_port: int,
base_topic: str,
username: str | None = None,
password: str | None = None,
) -> tuple[list[dict], list[dict]]:
"""Connect to the MQTT broker, request the Z2M networkmap, and return (nodes, edges).
Raises:
TimeoutError: if the broker does not respond in time.
ConnectionError: if the broker cannot be reached.
ValueError: if the response payload is malformed.
"""
try:
import aiomqtt # type: ignore[import]
except ImportError as exc: # pragma: no cover
raise ImportError(
"aiomqtt is required for Zigbee import. "
"Install it with: pip install aiomqtt"
) from exc
request_topic = _NETWORKMAP_REQUEST_TOPIC.format(base_topic=base_topic)
response_topic = _NETWORKMAP_RESPONSE_TOPIC.format(base_topic=base_topic)
result_event: asyncio.Event = asyncio.Event()
response_payload: dict = {}
try:
async with aiomqtt.Client(
hostname=mqtt_host,
port=mqtt_port,
username=username,
password=password,
timeout=_CONNECTION_TIMEOUT,
) as client:
await client.subscribe(response_topic)
await client.publish(
request_topic,
json.dumps({"type": "raw", "routes": False}),
)
async def _wait_for_response() -> None:
async for message in client.messages:
if str(message.topic) == response_topic:
try:
response_payload.update(json.loads(message.payload))
except (json.JSONDecodeError, TypeError) as exc:
raise ValueError(f"Malformed networkmap response: {exc}") from exc
result_event.set()
break
await asyncio.wait_for(_wait_for_response(), timeout=_NETWORKMAP_TIMEOUT)
except aiomqtt.MqttError as exc:
raise ConnectionError(f"MQTT connection failed: {exc}") from exc
except asyncio.TimeoutError as exc:
raise TimeoutError(
f"Timed out waiting for networkmap response from {mqtt_host}:{mqtt_port}"
) from exc
if not response_payload:
raise ValueError("Empty networkmap response received")
return parse_networkmap(response_payload)
async def test_mqtt_connection(
mqtt_host: str,
mqtt_port: int,
username: str | None = None,
password: str | None = None,
) -> bool:
"""Attempt a quick MQTT connection to verify broker reachability.
Returns True on success, raises ConnectionError on failure.
"""
try:
import aiomqtt # type: ignore[import]
except ImportError as exc: # pragma: no cover
raise ImportError("aiomqtt is required") from exc
try:
async with aiomqtt.Client(
hostname=mqtt_host,
port=mqtt_port,
username=username,
password=password,
timeout=_CONNECTION_TIMEOUT,
):
return True
except aiomqtt.MqttError as exc:
raise ConnectionError(f"MQTT connection failed: {exc}") from exc
except asyncio.TimeoutError as exc:
raise TimeoutError(f"Connection to {mqtt_host}:{mqtt_port} timed out") from exc