From 471e1795bfdadf0f3418cc7d0a926b01cee99aed Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 8 Jun 2026 13:46:34 -0400 Subject: [PATCH] feat(report): optional MQTT publisher and dashboard state (V2-P1-006) Add opt-in REPORTS.MQTT with retained shared state and live voice_event topics, topology static TG in v2 payloads, dashboard_state builder, and SIGHUP reload. --- adn-server.example.yaml | 12 + docs/en/server/protocols/report-v2.md | 55 ++++ docs/es/server/protocols/report-v2.md | 55 ++++ pyproject.toml | 3 +- schemas/report-v2.json | 22 +- src/adn_server/application/ports.py | 29 ++ src/adn_server/application/report/__init__.py | 2 + .../application/report/dashboard_state.py | 137 ++++++++++ src/adn_server/application/report/payloads.py | 69 ++++- .../infrastructure/config_validator.py | 24 ++ .../twisted_adapters/report/mqtt_config.py | 154 +++++++++++ .../twisted_adapters/report/mqtt_publisher.py | 256 ++++++++++++++++++ .../twisted_adapters/report/mqtt_topics.py | 62 +++++ .../twisted_adapters/report_server.py | 43 ++- src/adn_server/main.py | 22 +- tests/application/test_dashboard_state.py | 87 ++++++ tests/application/test_report_payloads.py | 42 +++ tests/infrastructure/test_mqtt_auth.py | 73 +++++ tests/infrastructure/test_mqtt_config.py | 134 +++++++++ .../test_mqtt_publish_filter.py | 39 +++ tests/infrastructure/test_mqtt_publisher.py | 174 ++++++++++++ tests/infrastructure/test_mqtt_reload.py | 123 +++++++++ tests/infrastructure/test_mqtt_topics.py | 40 +++ .../infrastructure/test_report_server_mqtt.py | 62 +++++ 24 files changed, 1705 insertions(+), 14 deletions(-) create mode 100644 src/adn_server/application/report/dashboard_state.py create mode 100644 src/adn_server/infrastructure/twisted_adapters/report/mqtt_config.py create mode 100644 src/adn_server/infrastructure/twisted_adapters/report/mqtt_publisher.py create mode 100644 src/adn_server/infrastructure/twisted_adapters/report/mqtt_topics.py create mode 100644 tests/application/test_dashboard_state.py create mode 100644 tests/infrastructure/test_mqtt_auth.py create mode 100644 tests/infrastructure/test_mqtt_config.py create mode 100644 tests/infrastructure/test_mqtt_publish_filter.py create mode 100644 tests/infrastructure/test_mqtt_publisher.py create mode 100644 tests/infrastructure/test_mqtt_reload.py create mode 100644 tests/infrastructure/test_mqtt_topics.py create mode 100644 tests/infrastructure/test_report_server_mqtt.py diff --git a/adn-server.example.yaml b/adn-server.example.yaml index 2d547dc..d0b4558 100644 --- a/adn-server.example.yaml +++ b/adn-server.example.yaml @@ -29,6 +29,18 @@ REPORTS: REPORT_INTERVAL: 60 REPORT_PORT: 4321 REPORT_CLIENTS: "127.0.0.1" + # Optional MQTT (disabled unless explicitly enabled). + # MQTT: + # ENABLED: true + # URL: mqtt://127.0.0.1:1883 + # # Auth: USERNAME/PASSWORD here, or mqtt://user:pass@host:1883 in URL (YAML overrides URL). + # USERNAME: "" + # PASSWORD: "" + # # TLS (mqtts://): optional CA bundle for broker verification + # CAFILE: "" + # TOPIC_PREFIX: adn/73010 + # QOS: 0 + # Wire: state (retained), voice_event (live) LOGGER: # ENABLED: false # omit or true = normal; false = no console/file logs (NullHandler) diff --git a/docs/en/server/protocols/report-v2.md b/docs/en/server/protocols/report-v2.md index a7942b0..3313339 100644 --- a/docs/en/server/protocols/report-v2.md +++ b/docs/en/server/protocols/report-v2.md @@ -252,6 +252,61 @@ REPORTS: No `PROTOCOL` switch on **2.x** — wire is always JSON (`infrastructure/twisted_adapters/report/wire.py`). Payload mapping: **`application/report/`**. +### Optional MQTT mirror (V2-P1-006) + +By default the server sends reports **only** over TCP netstring (adn-monitor and other TCP clients). MQTT is **disabled** unless you explicitly enable it. + +**Enable** only when both are set: + +1. `REPORTS.MQTT.ENABLED: true` (boolean `true`, not merely present) +2. `REPORTS.MQTT.URL` — broker URL (`mqtt://host:1883` or `mqtts://host:8883`) + +```yaml +REPORTS: + REPORT: true + REPORT_PORT: 4321 + MQTT: + ENABLED: true + URL: mqtt://127.0.0.1:1883 + TOPIC_PREFIX: adn/73010 # optional; default adn/{GLOBAL.SERVER_ID} + USERNAME: my-mqtt-user # optional; overrides user in URL + PASSWORD: my-mqtt-secret # optional; overrides password in URL + CAFILE: /path/to/ca.pem # optional; broker TLS trust store (mqtts://) + QOS: 0 # optional, 0–2 +``` + +**Client ID:** auto-generated at startup as `adn-server-{GLOBAL.SERVER_ID}-{random}` (not configurable; random suffix avoids broker session collisions on restart). + +**Authentication:** username/password via `USERNAME` and `PASSWORD`, or embedded in the URL (`mqtt://user:pass@host:1883`). YAML credentials override URL userinfo. Password may be empty if the broker allows it. With `mqtts://`, set `CAFILE` when the broker uses a private CA. + +Requires optional dependency: `pip install 'adn-server[mqtt]'` (paho-mqtt). If MQTT is enabled but the library is missing, the server logs an error and continues with TCP only. + +**MQTT wire (fixed, not configurable):** only **`voice_event`** (telemetry) and **`state`** (snapshot). TCP-only types (`topology`, `routing_table`, `delta`, `hello`) are **not** published on MQTT. + +**Topic convention** (shared under `{prefix}`): + +| Topic | Direction | JSON `type` | Retain | +|-------|-----------|-------------|--------| +| `voice_event` | server → broker | `voice_event` | no | +| `state` | server → broker | `dashboard_state` | yes | + +`state` carries masters with connected peers, homebrew peers, and openbridges (monitor WebSocket `conf,lnksys` + `conf,opb` intent). It is **retained** so new subscribers receive the last snapshot without requesting it. **Topology-driven refreshes** republish `{prefix}/state` when the dashboard changes (dedup). + +**Triggers:** live `voice_event`; retained `{prefix}/state` on topology changes and after MQTT connect (dedup). + +**Example** (`SERVER_ID` 7302): + +```bash +mosquitto_sub -h BROKER -p 1883 -u USER -P PASS -t 'adn/7302/state' -v +mosquitto_sub -h BROKER -p 1883 -u USER -P PASS -t 'adn/7302/voice_event' -v +``` + +**Broker ACL:** consumers need **subscribe** on `adn/7302/state` and `adn/7302/voice_event`; the server `client_id` needs **publish** on those two topics only (no server-side subscribe). + +**Reload (`systemctl reload` / SIGHUP):** when `REPORTS.MQTT` changes (enable/disable, URL, credentials, TLS, `TOPIC_PREFIX`, `QOS`), the server disconnects the old MQTT client and connects with the new settings, or stays offline if `ENABLED` becomes false. + +**QoS:** configurable via `MQTT.QOS` (default `0`). + **`rule_timer`**, **`stat_trimmer`**, and **`bridgeDebug`** use **routing deltas** when only part of `BRIDGES` changed; connect, `REPORT_INTERVAL`, reload, and client `CONFIG_REQ` / `BRIDGE_REQ` send **full** snapshots. ## Version pairing (supported combinations) diff --git a/docs/es/server/protocols/report-v2.md b/docs/es/server/protocols/report-v2.md index bd642c0..d7133c5 100644 --- a/docs/es/server/protocols/report-v2.md +++ b/docs/es/server/protocols/report-v2.md @@ -249,6 +249,61 @@ REPORTS: Sin `PROTOCOL` en **2.x** — wire siempre JSON (`wire.py`). Mapeo: **`application/report/`**. +### Espejo MQTT opcional (V2-P1-006) + +Por defecto los informes salen **solo** por TCP netstring (adn-monitor y otros clientes TCP). MQTT queda **deshabilitado** salvo activación explícita. + +**Activar** solo cuando se cumplen ambas condiciones: + +1. `REPORTS.MQTT.ENABLED: true` (booleano `true`, no basta con existir la clave) +2. `REPORTS.MQTT.URL` — URL del broker (`mqtt://host:1883` o `mqtts://host:8883`) + +```yaml +REPORTS: + REPORT: true + REPORT_PORT: 4321 + MQTT: + ENABLED: true + URL: mqtt://127.0.0.1:1883 + TOPIC_PREFIX: adn/73010 # opcional; por defecto adn/{GLOBAL.SERVER_ID} + USERNAME: my-mqtt-user # opcional; sustituye al usuario en la URL + PASSWORD: my-mqtt-secret # opcional; sustituye a la contraseña en la URL + CAFILE: /path/to/ca.pem # opcional; CA del broker con mqtts:// + QOS: 0 # opcional, 0–2 +``` + +**Client ID:** generado al arrancar como `adn-server-{GLOBAL.SERVER_ID}-{random}` (no configurable; el sufijo aleatorio evita colisiones de sesión en el broker al reiniciar). + +**Autenticación:** usuario/contraseña con `USERNAME` y `PASSWORD`, o embebidos en la URL (`mqtt://user:pass@host:1883`). Las credenciales YAML tienen prioridad sobre la URL. La contraseña puede ir vacía si el broker lo permite. Con `mqtts://`, indique `CAFILE` si el broker usa una CA propia. + +Dependencia opcional: `pip install 'adn-server[mqtt]'` (paho-mqtt). Si MQTT está habilitado pero falta la librería, el servidor registra error y sigue solo con TCP. + +**Wire MQTT (fijo, no configurable):** solo **`voice_event`** (telemetría) y **`state`** (snapshot). Los tipos solo-TCP (`topology`, `routing_table`, `delta`, `hello`) **no** se publican por MQTT. + +**Convención de topics** (compartidos bajo `{prefix}`): + +| Topic | Dirección | `type` JSON | Retain | +|-------|-----------|-------------|--------| +| `voice_event` | servidor → broker | `voice_event` | no | +| `state` | servidor → broker | `dashboard_state` | sí | + +`state` incluye masters con peers conectados, peers sueltos y openbridges (equivalente WebSocket `conf,lnksys` + `conf,opb`). Va con **retain** para que nuevos suscriptores reciban el último snapshot sin pedirlo. Los **refrescos por topología** republican `{prefix}/state` cuando cambia el dashboard (dedup). + +**Disparadores:** `voice_event` en vivo; `{prefix}/state` retenido al cambiar topología y tras conectar MQTT (dedup). + +**Ejemplo** (`SERVER_ID` 7302): + +```bash +mosquitto_sub -h BROKER -p 1883 -u USER -P PASS -t 'adn/7302/state' -v +mosquitto_sub -h BROKER -p 1883 -u USER -P PASS -t 'adn/7302/voice_event' -v +``` + +**ACL TBMQ:** consumidores con **subscribe** en `adn/7302/state` y `adn/7302/voice_event`; el `client_id` del servidor solo necesita **publish** en esos dos topics (sin subscribe en el servidor). + +**Recarga (`systemctl reload` / SIGHUP):** si cambia `REPORTS.MQTT` (activar/desactivar, URL, credenciales, TLS, `TOPIC_PREFIX`, `QOS`), el servidor desconecta el cliente MQTT anterior y conecta con la nueva config, o se queda sin MQTT si `ENABLED` pasa a false. + +**QoS:** configurable con `MQTT.QOS` (por defecto `0`). + Los timers usan **deltas** cuando solo cambia parte de `BRIDGES`; conexión, `REPORT_INTERVAL`, reload y `CONFIG_REQ` / `BRIDGE_REQ` envían instantáneas **completas**. ## Acoplamiento de versiones (combinaciones soportadas) diff --git a/pyproject.toml b/pyproject.toml index 135c891..5749921 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -21,7 +21,8 @@ dependencies = [ ] [project.optional-dependencies] -dev = ["pytest>=7.0", "jsonschema>=4.0"] +mqtt = ["paho-mqtt>=2.0"] +dev = ["pytest>=7.0", "jsonschema>=4.0", "paho-mqtt>=2.0"] docs = ["mkdocs>=1.6", "mkdocs-material>=9.5", "pymdown-extensions>=10.3"] [tool.setuptools.packages.find] diff --git a/schemas/report-v2.json b/schemas/report-v2.json index e8e8a82..da0edfd 100644 --- a/schemas/report-v2.json +++ b/schemas/report-v2.json @@ -80,6 +80,16 @@ "$ref": "#/$defs/dmr_id", "description": "OPENBRIDGE NETWORK_ID (legacy CONFIG bytes_4 value as integer)." }, + "ts1_static": { + "type": "array", + "items": { "type": "string" }, + "description": "MASTER slot-1 static talkgroups (from YAML and peer RPTO OPTIONS)." + }, + "ts2_static": { + "type": "array", + "items": { "type": "string" }, + "description": "MASTER slot-2 static talkgroups (from YAML and peer RPTO OPTIONS)." + }, "peers": { "type": "array", "items": { @@ -101,7 +111,17 @@ "package_id": { "type": "string" }, "software_id": { "type": "string" }, "colorcode": { "type": "string" }, - "tx_power": { "type": "string" } + "tx_power": { "type": "string" }, + "ts1_static": { + "type": "array", + "items": { "type": "string" }, + "description": "Slot-1 static TGs from this peer's RPTO OPTIONS." + }, + "ts2_static": { + "type": "array", + "items": { "type": "string" }, + "description": "Slot-2 static TGs from this peer's RPTO OPTIONS." + } } } } diff --git a/src/adn_server/application/ports.py b/src/adn_server/application/ports.py index 92327e1..919d2ad 100644 --- a/src/adn_server/application/ports.py +++ b/src/adn_server/application/ports.py @@ -112,6 +112,35 @@ class ReportWireEncoder(ABC): ... +class ReportMqttPublisher(ABC): + """Optional second sink: publish the same report v2 JSON payloads to an MQTT broker.""" + + @abstractmethod + def start( + self, + wire: ReportWireEncoder, + get_systems: Any, + get_bridges: Any, + ) -> None: + """Connect to broker and publish bootstrap snapshots (hello + full topology + routing).""" + ... + + @abstractmethod + def publish_frames(self, frames: tuple[bytes, ...]) -> None: + """Publish zero or more wire frames (opcode + JSON) to MQTT topics.""" + ... + + @abstractmethod + def publish_dashboard(self, systems: dict[str, Any]) -> None: + """Publish slim ``dashboard_state`` (linked systems only).""" + ... + + @abstractmethod + def stop(self) -> None: + """Disconnect from broker.""" + ... + + class ReportSender(ABC): """Send config and bridge state to report TCP clients (CONFIG_SND, BRIDGE_SND, BRDG_EVENT).""" diff --git a/src/adn_server/application/report/__init__.py b/src/adn_server/application/report/__init__.py index 9e03ac5..95ea7d5 100644 --- a/src/adn_server/application/report/__init__.py +++ b/src/adn_server/application/report/__init__.py @@ -6,6 +6,7 @@ from .queue import ( BoundedReportQueue, QueuedReportSender, ) +from .dashboard_state import build_dashboard_state from .payloads import ( REPORT_FEATURES, REPORT_PROTOCOL, @@ -24,6 +25,7 @@ __all__ = [ "QueuedReportSender", "REPORT_FEATURES", "REPORT_PROTOCOL", + "build_dashboard_state", "build_routing_table", "build_topology", "hello_connected_system_names", diff --git a/src/adn_server/application/report/dashboard_state.py b/src/adn_server/application/report/dashboard_state.py new file mode 100644 index 0000000..6275dcc --- /dev/null +++ b/src/adn_server/application/report/dashboard_state.py @@ -0,0 +1,137 @@ +"""Minimal dashboard snapshot for external MQTT consumers (not full topology / routing).""" + +from __future__ import annotations + +import time +from typing import Any + +from adn_server.domain import int_id +from .payloads import _peer_field_json, build_topology, static_tg_list + + +def _connected_topology_peers(system: dict[str, Any]) -> list[dict[str, Any]]: + return [p for p in system.get("peers", []) if isinstance(p, dict) and p.get("connected")] + + +def _upstream_peer_connected(cfg: dict[str, Any]) -> bool: + mode = cfg.get("MODE", "MASTER") + if mode not in ("PEER", "XLXPEER"): + return False + for key in ("XLXSTATS", "STATS"): + block = cfg.get(key) + if isinstance(block, dict) and block.get("CONNECTION") == "YES": + return True + return False + + +def _upstream_peer_block(name: str, cfg: dict[str, Any]) -> dict[str, Any]: + """Homebrew / XLX upstream (``CTABLE.PEERS``), not hotspots under a MASTER.""" + mode = cfg.get("MODE", "PEER") + block: dict[str, Any] = {"mode": mode, "connected": True} + for legacy_key, json_key in ( + ("CALLSIGN", "callsign"), + ("LOCATION", "location"), + ("DESCRIPTION", "description"), + ("URL", "url"), + ("MASTER_IP", "master_ip"), + ("MASTER_PORT", "master_port"), + ): + text = _peer_field_json(cfg.get(legacy_key)) + if text is not None: + block[json_key] = text + radio_id = cfg.get("RADIO_ID") + if radio_id is not None: + block["radio_id"] = int_id(radio_id) + stats_key = "XLXSTATS" if mode == "XLXPEER" else "STATS" + stats = cfg.get(stats_key) + if isinstance(stats, dict) and stats.get("CONNECTED"): + try: + block["connected_at"] = int(float(stats["CONNECTED"])) + except (TypeError, ValueError): + pass + return block + + +def _openbridge_block(name: str, cfg: dict[str, Any], topology_row: dict[str, Any] | None) -> dict[str, Any]: + """Enabled OPENBRIDGE legs (``CTABLE.OPENBRIDGES``); STREAMS stay empty here (live chips = monitor/voice).""" + block: dict[str, Any] = {"mode": "OPENBRIDGE", "streams": {}} + network_id = cfg.get("NETWORK_ID") + if network_id is not None: + block["network_id"] = int_id(network_id) + row = topology_row or {} + if row.get("ip"): + block["ip"] = row["ip"] + if row.get("port") is not None: + block["port"] = int(row["port"]) + if row.get("enhanced_obp") or cfg.get("ENHANCED_OBP"): + block["enhanced_obp"] = True + return block + + +def build_dashboard_state( + systems: dict[str, Any], + *, + server_id: str | None = None, + ts: float | None = None, +) -> dict[str, Any]: + """Slim linked-systems view (masters with peers, homebrew peers, openbridges). + + Mirrors adn-monitor WebSocket ``ctable_for_lnksys`` + ``ctable_for_opb`` intent: + no routing_table, no idle masters, no secrets. + """ + epoch = time.time() if ts is None else ts + topology = build_topology(systems, seq=0, ts=epoch) + topology_by_name = { + s["name"]: s + for s in topology.get("systems", []) + if isinstance(s, dict) and s.get("name") + } + masters: dict[str, Any] = {} + peers: dict[str, Any] = {} + openbridges: dict[str, Any] = {} + + for name, cfg in systems.items(): + if not isinstance(cfg, dict) or not cfg.get("ENABLED", True): + continue + mode = cfg.get("MODE", "MASTER") + topo = topology_by_name.get(name) + + if mode == "MASTER": + if topo is None: + continue + live = _connected_topology_peers(topo) + if not live: + continue + block: dict[str, Any] = { + "mode": "MASTER", + "peers": {int(p["id"]): p for p in live if "id" in p}, + } + if topo.get("ip"): + block["ip"] = topo["ip"] + if topo.get("port") is not None: + block["port"] = int(topo["port"]) + master_ts1 = static_tg_list(cfg.get("TS1_STATIC")) + master_ts2 = static_tg_list(cfg.get("TS2_STATIC")) + for peer_row in block["peers"].values(): + if master_ts1 and "ts1_static" not in peer_row: + peer_row["ts1_static"] = master_ts1 + if master_ts2 and "ts2_static" not in peer_row: + peer_row["ts2_static"] = master_ts2 + masters[name] = block + elif mode == "OPENBRIDGE": + openbridges[name] = _openbridge_block(name, cfg, topo) + elif mode in ("PEER", "XLXPEER") and _upstream_peer_connected(cfg): + peers[name] = _upstream_peer_block(name, cfg) + + payload: dict[str, Any] = { + "type": "dashboard_state", + "ts": float(epoch), + "ctable": { + "MASTERS": masters, + "PEERS": peers, + "OPENBRIDGES": openbridges, + }, + } + if server_id is not None: + payload["server_id"] = server_id + return payload diff --git a/src/adn_server/application/report/payloads.py b/src/adn_server/application/report/payloads.py index dc1afb7..9ad8b42 100644 --- a/src/adn_server/application/report/payloads.py +++ b/src/adn_server/application/report/payloads.py @@ -19,7 +19,7 @@ REPORT_FEATURES = ( "DELTA_UPDATES", ) -_SYSTEM_MODES = frozenset({"MASTER", "PEER", "OPENBRIDGE"}) +_SYSTEM_MODES = frozenset({"MASTER", "PEER", "XLXPEER", "OPENBRIDGE"}) _TO_TYPES = frozenset({"ON", "OFF", "STAT", "NONE"}) _CSV_FAMILIES = { "GROUP VOICE": "GROUP", @@ -38,6 +38,60 @@ def _peer_connected(peer: dict[str, Any]) -> bool: return peer.get("CONNECTION") == "YES" +def static_tg_list(value: Any) -> list[str]: + """Normalize legacy TS1_STATIC / TS2_STATIC (comma string or list) to string TG ids.""" + if value is None: + return [] + if isinstance(value, str): + return [x.strip() for x in value.split(",") if x.strip()] + if isinstance(value, list): + return [str(x).strip() for x in value if str(x).strip()] + text = str(value).strip() + return [text] if text else [] + + +def parse_peer_options_static(options: Any) -> tuple[list[str], list[str]]: + """Parse hotspot RPTO OPTIONS (``TS1=…;TS2=…;``) into static TG id lists.""" + if options is None: + return [], [] + if isinstance(options, bytes): + text = options.decode("utf-8", errors="replace") + else: + text = str(options) + text = text.rstrip("\x00").strip() + if not text: + return [], [] + parsed: dict[str, str] = {} + for part in text.split(";"): + part = part.strip() + if "=" not in part: + continue + key, value = part.split("=", 1) + parsed[key.strip().upper()] = value.strip() + for old, new in (("TS1", "TS1_STATIC"), ("TS2", "TS2_STATIC")): + if old in parsed and new not in parsed: + parsed[new] = parsed[old] + ts1_parts: list[str] = [] + if "TS1_1" in parsed: + ts1_parts.append(parsed["TS1_1"]) + for i in range(2, 10): + p = parsed.get(f"TS1_{i}") + if p: + ts1_parts.append(p) + elif parsed.get("TS1_STATIC"): + ts1_parts = [x.strip() for x in parsed["TS1_STATIC"].split(",") if x.strip()] + ts2_parts: list[str] = [] + if "TS2_1" in parsed: + ts2_parts.append(parsed["TS2_1"]) + for i in range(2, 10): + p = parsed.get(f"TS2_{i}") + if p: + ts2_parts.append(p) + elif parsed.get("TS2_STATIC"): + ts2_parts = [x.strip() for x in parsed["TS2_STATIC"].split(",") if x.strip()] + return ts1_parts, ts2_parts + + def _peer_field_json(value: Any) -> str | None: """Sanitize a legacy peer field for JSON (no secrets).""" if value is None: @@ -95,6 +149,12 @@ def _topology_peer_row(peer_key: Any, peer: dict[str, Any]) -> dict[str, Any]: text = _peer_field_json(peer[legacy_key]) if text is not None: row[json_key] = text + if "OPTIONS" in peer: + ts1, ts2 = parse_peer_options_static(peer.get("OPTIONS")) + if ts1: + row["ts1_static"] = ts1 + if ts2: + row["ts2_static"] = ts2 return row @@ -159,6 +219,13 @@ def build_topology(systems: dict[str, Any], *, seq: int, ts: float | None = None entry["enhanced_obp"] = True if mode == "OPENBRIDGE" and cfg.get("NETWORK_ID") is not None: entry["network_id"] = _dmr_id(cfg["NETWORK_ID"]) + if mode == "MASTER": + ts1 = static_tg_list(cfg.get("TS1_STATIC")) + ts2 = static_tg_list(cfg.get("TS2_STATIC")) + if ts1: + entry["ts1_static"] = ts1 + if ts2: + entry["ts2_static"] = ts2 peers_out: list[dict[str, Any]] = [] for peer_key, peer in cfg.get("PEERS", {}).items(): if not isinstance(peer, dict): diff --git a/src/adn_server/infrastructure/config_validator.py b/src/adn_server/infrastructure/config_validator.py index 3d14560..9f1aa03 100644 --- a/src/adn_server/infrastructure/config_validator.py +++ b/src/adn_server/infrastructure/config_validator.py @@ -157,6 +157,29 @@ def _validate_global(global_cfg: dict[str, Any], errors: list[str]) -> None: errors.append("GLOBAL.PASS_SECURITY: required when GLOBAL.URL_SECURITY is set.") +def _validate_mqtt(reports_cfg: dict[str, Any], errors: list[str]) -> None: + mqtt_block = reports_cfg.get("MQTT") + if isinstance(mqtt_block, dict): + if "ENABLED" in mqtt_block: + _expect_bool("REPORTS.MQTT.ENABLED", mqtt_block["ENABLED"], errors) + if mqtt_block.get("ENABLED") is True: + url = mqtt_block.get("URL") or reports_cfg.get("MQTT_URL") + if _is_empty(url): + errors.append("REPORTS.MQTT.URL: required when REPORTS.MQTT.ENABLED is true.") + if "QOS" in mqtt_block: + qos = mqtt_block["QOS"] + if not isinstance(qos, int) or qos < 0 or qos > 2: + errors.append(f"REPORTS.MQTT.QOS: expected integer 0–2, got {qos!r}.") + if "MQTT_ENABLED" in reports_cfg: + _expect_bool("REPORTS.MQTT_ENABLED", reports_cfg["MQTT_ENABLED"], errors) + if reports_cfg.get("MQTT_ENABLED") is True and _is_empty(reports_cfg.get("MQTT_URL")): + errors.append("REPORTS.MQTT_URL: required when REPORTS.MQTT_ENABLED is true.") + if "MQTT_QOS" in reports_cfg: + qos = reports_cfg["MQTT_QOS"] + if not isinstance(qos, int) or qos < 0 or qos > 2: + errors.append(f"REPORTS.MQTT_QOS: expected integer 0–2, got {qos!r}.") + + def _validate_reports(reports_cfg: dict[str, Any], errors: list[str]) -> None: if "REPORT" in reports_cfg: _expect_bool("REPORTS.REPORT", reports_cfg["REPORT"], errors) @@ -169,6 +192,7 @@ def _validate_reports(reports_cfg: dict[str, Any], errors: list[str]) -> None: errors.append( f"REPORTS.REPORT_CLIENTS: expected string or list, got {type(clients).__name__} ({clients!r})." ) + _validate_mqtt(reports_cfg, errors) def _validate_logger(logger_cfg: dict[str, Any], errors: list[str]) -> None: diff --git a/src/adn_server/infrastructure/twisted_adapters/report/mqtt_config.py b/src/adn_server/infrastructure/twisted_adapters/report/mqtt_config.py new file mode 100644 index 0000000..101ff12 --- /dev/null +++ b/src/adn_server/infrastructure/twisted_adapters/report/mqtt_config.py @@ -0,0 +1,154 @@ +"""Parse optional REPORTS.MQTT settings (disabled unless explicitly enabled).""" + +from __future__ import annotations + +import secrets +from dataclasses import dataclass +from typing import Any +from urllib.parse import unquote, urlparse + + +@dataclass(frozen=True) +class MqttBroker: + host: str + port: int + use_tls: bool + display_url: str + + +# Fixed MQTT wire (not configurable): live voice + retained shared ``{prefix}/state`` only. +MQTT_PUBLISH_VOICE_EVENT = "voice_event" +MQTT_PUBLISH_STATE = "state" + + +@dataclass(frozen=True) +class MqttSettings: + broker: MqttBroker + topic_prefix: str + client_id: str + username: str | None + password: str | None + qos: int + cafile: str | None = None + + +def _non_empty(value: Any) -> str | None: + if value is None: + return None + text = str(value).strip() + return text or None + + +def _password_value(value: Any) -> str | None: + """Return password string; empty YAML password is a valid (empty) secret.""" + if value is None: + return None + return str(value) + + +def parse_mqtt_broker(url: str) -> tuple[MqttBroker, str | None, str | None]: + """Parse broker URL; return endpoint and optional username/password from userinfo.""" + normalized = url if "://" in url else f"mqtt://{url}" + parsed = urlparse(normalized) + scheme = (parsed.scheme or "mqtt").lower() + host = parsed.hostname or "127.0.0.1" + if parsed.port is not None: + port = parsed.port + elif scheme in ("mqtts", "ssl", "tls"): + port = 8883 + else: + port = 1883 + use_tls = scheme in ("mqtts", "ssl", "tls") + display_url = f"{'mqtts' if use_tls else 'mqtt'}://{host}:{port}" + + url_user = unquote(parsed.username) if parsed.username is not None else None + url_pass = unquote(parsed.password) if parsed.password is not None else None + if url_user is not None: + url_user = url_user.strip() or None + return MqttBroker(host=host, port=port, use_tls=use_tls, display_url=display_url), url_user, url_pass + + +def _server_id(config: dict[str, Any]) -> str: + """Display server id for MQTT topics/client_id (GLOBAL.SERVER_ID is bytes after normalize).""" + sid = config.get("GLOBAL", {}).get("SERVER_ID", "server") + if isinstance(sid, bytes): + raw = sid[:4].ljust(4, b"\x00")[:4] + return str(int.from_bytes(raw, "big")) + if isinstance(sid, int): + return str(sid & 0xFFFFFFFF) + text = str(sid).strip() + return text or "server" + + +def default_topic_prefix(config: dict[str, Any]) -> str: + return f"adn/{_server_id(config)}" + + +def default_mqtt_client_id(config: dict[str, Any]) -> str: + """Derived client id: adn-server-{SERVER_ID}-{random} (not configurable).""" + return f"adn-server-{_server_id(config)}-{secrets.token_hex(4)}" + + +def _resolve_credentials( + mqtt_block: dict[str, Any] | None, + reports: dict[str, Any], + url_user: str | None, + url_pass: str | None, +) -> tuple[str | None, str | None]: + """Explicit YAML USERNAME/PASSWORD override credentials embedded in the broker URL.""" + username = url_user + password = url_pass + if isinstance(mqtt_block, dict): + if "USERNAME" in mqtt_block: + username = _non_empty(mqtt_block.get("USERNAME")) + elif "MQTT_USERNAME" in reports: + username = _non_empty(reports.get("MQTT_USERNAME")) + if "PASSWORD" in mqtt_block: + password = _password_value(mqtt_block.get("PASSWORD")) + elif "MQTT_PASSWORD" in reports: + password = _password_value(reports.get("MQTT_PASSWORD")) + else: + if "MQTT_USERNAME" in reports: + username = _non_empty(reports.get("MQTT_USERNAME")) + if "MQTT_PASSWORD" in reports: + password = _password_value(reports.get("MQTT_PASSWORD")) + return username, password + + +def mqtt_settings_from_config(config: dict[str, Any]) -> MqttSettings | None: + """Return settings only when MQTT is explicitly enabled and a broker URL is set.""" + reports = config.get("REPORTS", {}) + mqtt_block = reports.get("MQTT") + if isinstance(mqtt_block, dict): + enabled = mqtt_block.get("ENABLED") is True + raw_url = _non_empty(mqtt_block.get("URL")) or _non_empty(reports.get("MQTT_URL")) + topic_prefix = _non_empty(mqtt_block.get("TOPIC_PREFIX")) or _non_empty( + reports.get("MQTT_TOPIC_PREFIX") + ) + cafile = _non_empty(mqtt_block.get("CAFILE")) or _non_empty(reports.get("MQTT_CAFILE")) + qos_raw = mqtt_block.get("QOS", reports.get("MQTT_QOS", 0)) + qos = int(qos_raw) if qos_raw is not None else 0 + else: + enabled = reports.get("MQTT_ENABLED") is True + raw_url = _non_empty(reports.get("MQTT_URL")) + topic_prefix = _non_empty(reports.get("MQTT_TOPIC_PREFIX")) + cafile = _non_empty(reports.get("MQTT_CAFILE")) + qos_raw = reports.get("MQTT_QOS", 0) + qos = int(qos_raw) if qos_raw is not None else 0 + mqtt_block = None + + if not enabled or not raw_url: + return None + + broker, url_user, url_pass = parse_mqtt_broker(raw_url) + username, password = _resolve_credentials(mqtt_block, reports, url_user, url_pass) + + return MqttSettings( + broker=broker, + topic_prefix=topic_prefix or default_topic_prefix(config), + client_id=default_mqtt_client_id(config), + username=username, + password=password, + qos=max(0, min(qos, 2)), + cafile=cafile, + ) diff --git a/src/adn_server/infrastructure/twisted_adapters/report/mqtt_publisher.py b/src/adn_server/infrastructure/twisted_adapters/report/mqtt_publisher.py new file mode 100644 index 0000000..99d85a9 --- /dev/null +++ b/src/adn_server/infrastructure/twisted_adapters/report/mqtt_publisher.py @@ -0,0 +1,256 @@ +"""Optional MQTT publisher for report v2 JSON (same payloads as TCP wire).""" + +from __future__ import annotations + +import json +import logging +from collections.abc import Callable +from typing import Any + +from adn_server.application.report.dashboard_state import build_dashboard_state +from adn_server.application.ports import ReportMqttPublisher, ReportWireEncoder + +from .mqtt_config import MQTT_PUBLISH_VOICE_EVENT, MqttSettings, mqtt_settings_from_config +from .mqtt_topics import frame_message_type, mqtt_shared_state_topic, topic_for_frame + +logger = logging.getLogger(__name__) + +_MQTT_WIRE_LABEL = "state,voice_event" +_SHARED_STATE_DEDUP_KEY = "__shared__" + +try: + import paho.mqtt.client as mqtt +except ImportError: # pragma: no cover - exercised via create_report_mqtt_publisher + mqtt = None # type: ignore[assignment] + + +class NullReportMqttPublisher(ReportMqttPublisher): + """No-op sink when MQTT is disabled.""" + + def start(self, wire: ReportWireEncoder, get_systems: Any, get_bridges: Any) -> None: + del wire, get_systems, get_bridges + + def publish_frames(self, frames: tuple[bytes, ...]) -> None: + del frames + + def publish_dashboard(self, systems: dict[str, Any]) -> None: + del systems + + def stop(self) -> None: + pass + + +class PahoReportMqttPublisher(ReportMqttPublisher): + """Publish retained shared ``state`` and live ``voice_event`` only.""" + + def __init__(self, settings: MqttSettings) -> None: + self._settings = settings + self._client: Any = None + self._connected = False + self._get_systems: Callable[[], dict[str, Any]] | None = None + self._last_state_dedup: dict[str, bytes] = {} + + def start( + self, + wire: ReportWireEncoder, + get_systems: Callable[[], dict[str, Any]], + get_bridges: Callable[[], dict[str, Any]], + ) -> None: + del wire, get_bridges + if mqtt is None: + logger.error("(REPORT) MQTT enabled but paho-mqtt is not installed; publisher disabled") + return + self._get_systems = get_systems + + client = mqtt.Client( + callback_api_version=mqtt.CallbackAPIVersion.VERSION2, + client_id=self._settings.client_id, + protocol=mqtt.MQTTv311, + ) + _apply_mqtt_auth(client, self._settings) + client.on_connect = self._on_connect + client.on_disconnect = self._on_disconnect + broker = self._settings.broker + if broker.use_tls: + if self._settings.cafile: + client.tls_set(ca_certs=self._settings.cafile) + else: + client.tls_set() + self._client = client + try: + client.connect(broker.host, broker.port, keepalive=60) + client.loop_start() + auth_user = self._settings.username + logger.info( + "(REPORT) MQTT publisher started broker=%s client_id=%s prefix=%s qos=%s auth=%s wire=%s", + broker.display_url, + self._settings.client_id, + self._settings.topic_prefix, + self._settings.qos, + auth_user or "no", + _MQTT_WIRE_LABEL, + ) + except Exception as e: + logger.warning("(REPORT) MQTT connect failed: %s", e) + self._client = None + + def publish_dashboard(self, systems: dict[str, Any]) -> None: + self._emit_shared_dashboard(systems, force=False) + + def _emit_shared_dashboard(self, systems: dict[str, Any], *, force: bool = False) -> None: + topic = mqtt_shared_state_topic(self._settings.topic_prefix) + self._publish_state(systems, topic=topic, dedup_key=_SHARED_STATE_DEDUP_KEY, force=force) + + def _publish_state( + self, + systems: dict[str, Any], + *, + topic: str, + dedup_key: str, + force: bool, + ) -> None: + if not self._connected or self._client is None: + return + payload = build_dashboard_state(systems, server_id=_server_id_from_settings(self._settings)) + body = json.dumps(payload, separators=(",", ":")).encode("utf-8") + content_key = json.dumps( + {"ctable": payload.get("ctable"), "server_id": payload.get("server_id")}, + separators=(",", ":"), + sort_keys=True, + ).encode("utf-8") + if not force and self._last_state_dedup.get(dedup_key) == content_key: + logger.debug("(REPORT) MQTT state unchanged, skip publish topic=%s", topic) + return + self._last_state_dedup[dedup_key] = content_key + try: + info = self._client.publish( + topic, + payload=body, + qos=self._settings.qos, + retain=True, + ) + if info.rc != mqtt.MQTT_ERR_SUCCESS: + logger.warning("(REPORT) MQTT state publish rc=%s topic=%s", info.rc, topic) + else: + logger.info("(REPORT) MQTT state published topic=%s bytes=%s", topic, len(body)) + except Exception as e: + logger.warning("(REPORT) MQTT state publish failed topic=%s: %s", topic, e) + + def publish_frames(self, frames: tuple[bytes, ...]) -> None: + if not self._connected or self._client is None: + return + prefix = self._settings.topic_prefix + qos = self._settings.qos + for frame in frames: + if frame_message_type(frame) != MQTT_PUBLISH_VOICE_EVENT: + continue + topic = topic_for_frame(frame, prefix) + if topic is None: + continue + payload = frame[1:] + try: + info = self._client.publish(topic, payload=payload, qos=qos, retain=False) + if info.rc != mqtt.MQTT_ERR_SUCCESS: + logger.debug("(REPORT) MQTT publish rc=%s topic=%s", info.rc, topic) + except Exception as e: + logger.debug("(REPORT) MQTT publish failed topic=%s: %s", topic, e) + + def stop(self) -> None: + client = self._client + self._client = None + self._connected = False + self._last_state_dedup.clear() + if client is None: + return + try: + client.loop_stop() + client.disconnect() + except Exception as e: + logger.debug("(REPORT) MQTT disconnect: %s", e) + + def _on_connect(self, client: Any, userdata: Any, connect_flags: Any, reason_code: Any, properties: Any) -> None: + del client, userdata, connect_flags, properties + if getattr(reason_code, "value", reason_code) != 0: + logger.warning("(REPORT) MQTT broker rejected connection: %s", reason_code) + return + self._connected = True + systems = self._get_systems() if self._get_systems else {} + if systems: + self._emit_shared_dashboard(systems, force=False) + + def _on_disconnect( + self, + client: Any, + userdata: Any, + disconnect_flags: Any, + reason_code: Any, + properties: Any, + ) -> None: + del client, userdata, disconnect_flags, properties + self._connected = False + if getattr(reason_code, "value", reason_code) != 0: + logger.debug("(REPORT) MQTT disconnected: %s", reason_code) + + +def _server_id_from_settings(settings: MqttSettings) -> str: + prefix = settings.topic_prefix + if prefix.startswith("adn/"): + return prefix[4:] + return prefix + + +def _apply_mqtt_auth(client: Any, settings: MqttSettings) -> None: + if settings.username is None and settings.password is None: + return + client.username_pw_set(settings.username or "", settings.password) + + +def reconcile_mqtt_publisher( + factory: Any, + current: ReportMqttPublisher | None, + before: MqttSettings | None, + after: MqttSettings | None, + *, + report_enabled: bool, +) -> ReportMqttPublisher | None: + """Stop/start MQTT client after SIGHUP when REPORTS.MQTT settings change.""" + if before == after and (current is not None) == (after is not None): + return current + if current is not None: + try: + current.stop() + except Exception as e: + logger.debug("(REPORT) MQTT stop on reload: %s", e) + if after is None: + factory.set_mqtt(None) + if before is not None: + logger.info("(REPORT) MQTT disconnected (disabled in config reload)") + return None + publisher = create_report_mqtt_publisher_from_settings(after) + factory.set_mqtt(publisher) + if publisher is None: + logger.warning("(REPORT) MQTT enabled in config but publisher could not start") + return None + if report_enabled: + factory.start_mqtt() + action = "reconnected" if before is not None else "connected" + logger.info("(REPORT) MQTT %s after config reload broker=%s", action, after.broker.display_url) + return publisher + + +def create_report_mqtt_publisher_from_settings(settings: MqttSettings) -> ReportMqttPublisher | None: + if mqtt is None: + logger.error( + "(REPORT) MQTT enabled but paho-mqtt is missing; " + "install with: pip install 'adn-server[mqtt]'" + ) + return None + return PahoReportMqttPublisher(settings) + + +def create_report_mqtt_publisher(config: dict[str, Any]) -> ReportMqttPublisher | None: + """Return a publisher when MQTT is explicitly enabled; otherwise None.""" + settings = mqtt_settings_from_config(config) + if settings is None: + return None + return create_report_mqtt_publisher_from_settings(settings) diff --git a/src/adn_server/infrastructure/twisted_adapters/report/mqtt_topics.py b/src/adn_server/infrastructure/twisted_adapters/report/mqtt_topics.py new file mode 100644 index 0000000..ca3a548 --- /dev/null +++ b/src/adn_server/infrastructure/twisted_adapters/report/mqtt_topics.py @@ -0,0 +1,62 @@ +"""Map report v2 wire frames to MQTT topic suffixes.""" + +from __future__ import annotations + +import json +from typing import Any + +from .opcodes import REPORT_OPCODES + + +def mqtt_shared_state_topic(prefix: str) -> str: + """Retained ``dashboard_state`` snapshot for all consumers (topology-driven refreshes).""" + return f"{prefix}/state" + + +def frame_message_type(frame: bytes) -> str | None: + """Report v2 message kind for MQTT publish filtering.""" + if len(frame) < 2: + return None + opcode = frame[0:1] + if opcode == REPORT_OPCODES["HELLO"]: + return "hello" + if opcode == REPORT_OPCODES["TOPOLOGY_SND"]: + return "topology" + if opcode == REPORT_OPCODES["ROUTING_TABLE_SND"]: + return "routing_table" + if opcode == REPORT_OPCODES["VOICE_EVENT_SND"]: + return "voice_event" + if opcode == REPORT_OPCODES["DELTA_SND"]: + return "delta" + return None + + +def topic_for_frame(frame: bytes, prefix: str) -> str | None: + """Return full MQTT topic for a TCP report frame (opcode + JSON), or None if unknown.""" + if len(frame) < 2: + return None + opcode = frame[0:1] + if opcode == REPORT_OPCODES["HELLO"]: + return f"{prefix}/hello" + if opcode == REPORT_OPCODES["TOPOLOGY_SND"]: + return f"{prefix}/topology" + if opcode == REPORT_OPCODES["ROUTING_TABLE_SND"]: + return f"{prefix}/routing_table" + if opcode == REPORT_OPCODES["VOICE_EVENT_SND"]: + return f"{prefix}/voice_event" + if opcode == REPORT_OPCODES["DELTA_SND"]: + return _delta_topic(frame[1:], prefix) + return None + + +def _delta_topic(payload: bytes, prefix: str) -> str: + try: + doc: dict[str, Any] = json.loads(payload.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError): + return f"{prefix}/delta" + patch_type = doc.get("patch", {}).get("type") + if patch_type == "topology": + return f"{prefix}/delta/topology" + if patch_type == "routing_table": + return f"{prefix}/delta/routing_table" + return f"{prefix}/delta" diff --git a/src/adn_server/infrastructure/twisted_adapters/report_server.py b/src/adn_server/infrastructure/twisted_adapters/report_server.py index d412637..e0c0df8 100644 --- a/src/adn_server/infrastructure/twisted_adapters/report_server.py +++ b/src/adn_server/infrastructure/twisted_adapters/report_server.py @@ -31,7 +31,7 @@ from typing import Any from twisted.internet.protocol import Factory from twisted.protocols.basic import NetstringReceiver -from adn_server.application.ports import ReportWireEncoder +from adn_server.application.ports import ReportMqttPublisher, ReportWireEncoder from .report import REPORT_OPCODES, create_report_wire @@ -70,9 +70,15 @@ class ReportProtocol(NetstringReceiver): class ReportServerFactory(Factory): """Twisted factory: ACL, client list, broadcast via injected ``ReportWireEncoder``.""" - def __init__(self, config: dict[str, Any]) -> None: + def __init__( + self, + config: dict[str, Any], + *, + mqtt: ReportMqttPublisher | None = None, + ) -> None: self._config = config self._wire: ReportWireEncoder = create_report_wire(config) + self._mqtt = mqtt self.clients: list[ReportProtocol] = [] self._systems: dict[str, Any] = {} self._bridges: dict[str, Any] = {} @@ -85,6 +91,12 @@ class ReportServerFactory(Factory): return ReportProtocol(self) return None + def set_config(self, config: dict[str, Any]) -> None: + self._config = config + + def set_mqtt(self, mqtt: ReportMqttPublisher | None) -> None: + self._mqtt = mqtt + def set_systems(self, systems: dict[str, Any]) -> None: self._systems = systems @@ -95,6 +107,16 @@ class ReportServerFactory(Factory): for frame in frames: client.sendString(frame) + def _broadcast_frames(self, frames: tuple[bytes, ...]) -> None: + if not frames: + return + for client in self.clients: + self._send_frames(client, frames) + + def start_mqtt(self) -> None: + if self._mqtt is not None: + self._mqtt.start(self._wire, lambda: self._systems, lambda: self._bridges) + def _send_hello_to(self, client: ReportProtocol) -> None: try: self._send_frames(client, self._wire.hello_frames(self._systems)) @@ -108,16 +130,17 @@ class ReportServerFactory(Factory): self._send_frames(client, self._wire.bridge_frames(self._bridges, full_snapshot=full_snapshot)) def send_config(self, *, incremental: bool = False) -> None: - full = not incremental - for client in self.clients: - self._send_config_to(client, full_snapshot=full) + frames = self._wire.config_frames(self._systems, full_snapshot=not incremental) + self._broadcast_frames(frames) + if self._mqtt is not None: + self._mqtt.publish_dashboard(self._systems) def send_bridge(self, *, incremental: bool = False) -> None: - full = not incremental - for client in self.clients: - self._send_bridge_to(client, full_snapshot=full) + frames = self._wire.bridge_frames(self._bridges, full_snapshot=not incremental) + self._broadcast_frames(frames) def send_bridge_event(self, event: str) -> None: frames = self._wire.bridge_event_frames(event) - for client in self.clients: - self._send_frames(client, frames) + self._broadcast_frames(frames) + if self._mqtt is not None: + self._mqtt.publish_frames(frames) diff --git a/src/adn_server/main.py b/src/adn_server/main.py index b56edc7..55545a2 100644 --- a/src/adn_server/main.py +++ b/src/adn_server/main.py @@ -89,6 +89,11 @@ from .infrastructure.security.password_download import DefaultSecurityDownloader from .infrastructure.security.user_passwords_loader import UserPasswordsLoader from .infrastructure.talker_alias_emblc import default_ta_emblc_encoder from .application.report.queue import BoundedReportQueue, QueuedReportSender +from .infrastructure.twisted_adapters.report.mqtt_config import mqtt_settings_from_config +from .infrastructure.twisted_adapters.report.mqtt_publisher import ( + create_report_mqtt_publisher, + reconcile_mqtt_publisher, +) from .infrastructure.twisted_adapters.report.worker import start_report_queue_worker from .infrastructure.twisted_adapters.report_server import ReportServerFactory from .infrastructure.twisted_adapters.udp_hbp import HBPProtocolFactory @@ -256,7 +261,8 @@ def main() -> None: # Protocol registry for send_to_system (legacy: systems[name].send_system(packet)) protocols: dict[str, Any] = {} udp_ports: dict[str, Any] = {} - report_factory = ReportServerFactory(config) + report_mqtt = create_report_mqtt_publisher(config) + report_factory = ReportServerFactory(config, mqtt=report_mqtt) def send_to_system(system_name: str, packet: bytes, **kwargs: Any) -> None: p = protocols.get(system_name) @@ -366,6 +372,9 @@ def main() -> None: report_inner, on_errback=lambda e: _looping_errback(e, logger), ) + if report_mqtt is not None: + report_factory.start_mqtt() + reactor.addSystemEventTrigger("during", "shutdown", report_mqtt.stop) # Reporting loop (REPORT_INTERVAL) — same logs as legacy after send_config/send_bridge def reporting_loop(): @@ -535,6 +544,8 @@ def main() -> None: user_passwords_loader.load(config) def _do_config_reload() -> None: + nonlocal report_mqtt + mqtt_before = mqtt_settings_from_config(config) new_config = prepare_reload_config(runtime_holder) try: result = reload_server_config( @@ -551,6 +562,15 @@ def main() -> None: log=logger, ) swap_runtime_config(runtime_holder, new_config, config_path=config_path) + report_factory.set_config(config) + mqtt_after = mqtt_settings_from_config(config) + report_mqtt = reconcile_mqtt_publisher( + report_factory, + report_mqtt, + mqtt_before, + mqtt_after, + report_enabled=config.get("REPORTS", {}).get("REPORT", True), + ) if result.added or result.removed or result.updated or result.rebound: _on_config_systems_changed() except Exception as e: diff --git a/tests/application/test_dashboard_state.py b/tests/application/test_dashboard_state.py new file mode 100644 index 0000000..aac3035 --- /dev/null +++ b/tests/application/test_dashboard_state.py @@ -0,0 +1,87 @@ +"""Minimal dashboard_state payload.""" + +from __future__ import annotations + +from adn_server.application.report.dashboard_state import build_dashboard_state + + +def test_dashboard_state_omits_idle_masters(): + systems = { + "SYSTEM-0": { + "MODE": "MASTER", + "ENABLED": True, + "PEERS": { + b"\x00\x00\x00\x01": {"CONNECTION": "NO", "CONNECTED": 0}, + }, + }, + "ECHO": { + "MODE": "MASTER", + "ENABLED": True, + "PEERS": { + b"\x00\x00\x00\x02": {"CONNECTION": "YES", "CONNECTED": 1000, "IP": "10.0.0.2"}, + }, + }, + } + state = build_dashboard_state(systems, server_id="7302") + assert state["type"] == "dashboard_state" + assert state["server_id"] == "7302" + assert "SYSTEM-0" not in state["ctable"]["MASTERS"] + assert "ECHO" in state["ctable"]["MASTERS"] + assert 2 in state["ctable"]["MASTERS"]["ECHO"]["peers"] + + +def test_dashboard_state_includes_master_static_tgs_on_peers(): + systems = { + "MASTER-A": { + "MODE": "MASTER", + "ENABLED": True, + "TS1_STATIC": "91,92", + "TS2_STATIC": "730", + "PEERS": { + b"\x00\x2f\xd0\x31": {"CONNECTION": "YES", "CONNECTED": 1000}, + }, + }, + } + state = build_dashboard_state(systems) + peers = state["ctable"]["MASTERS"]["MASTER-A"]["peers"] + assert len(peers) == 1 + peer = next(iter(peers.values())) + assert peer["ts1_static"] == ["91", "92"] + assert peer["ts2_static"] == ["730"] + + +def test_dashboard_state_includes_enabled_openbridge(): + systems = { + "OBP-CL": { + "MODE": "OPENBRIDGE", + "ENABLED": True, + "IP": "44.31.61.68", + "PORT": 62999, + "NETWORK_ID": 73010, + "ENHANCED_OBP": True, + "PEERS": {}, + }, + } + state = build_dashboard_state(systems) + obp = state["ctable"]["OPENBRIDGES"]["OBP-CL"] + assert obp["mode"] == "OPENBRIDGE" + assert obp["network_id"] == 73010 + assert obp["enhanced_obp"] is True + assert obp["streams"] == {} + + +def test_dashboard_state_includes_connected_upstream_peer(): + systems = { + "XLX-730": { + "MODE": "XLXPEER", + "ENABLED": True, + "CALLSIGN": "XLX730", + "RADIO_ID": 730, + "XLXSTATS": {"CONNECTION": "YES", "CONNECTED": 1700000000}, + "PEERS": {}, + }, + } + state = build_dashboard_state(systems) + assert "XLX-730" in state["ctable"]["PEERS"] + assert state["ctable"]["PEERS"]["XLX-730"]["connected"] is True + assert state["ctable"]["MASTERS"] == {} diff --git a/tests/application/test_report_payloads.py b/tests/application/test_report_payloads.py index 119136b..e9e8cc6 100644 --- a/tests/application/test_report_payloads.py +++ b/tests/application/test_report_payloads.py @@ -15,6 +15,7 @@ from adn_server.application.report import ( parse_bridge_event_csv, routing_table_delta, ) +from adn_server.application.report.payloads import parse_peer_options_static from adn_server.domain import bytes_3 _SCHEMA_PATH = Path(__file__).resolve().parents[2] / "schemas" / "report-v2.json" @@ -28,6 +29,47 @@ def validator() -> jsonschema.Draft202012Validator: return jsonschema.Draft202012Validator(schema) +def test_parse_peer_options_static_ts2(): + ts1, ts2 = parse_peer_options_static(b"TS2=730444;TIMER=15;") + assert ts1 == [] + assert ts2 == ["730444"] + + +def test_build_topology_exports_peer_options_static() -> None: + systems = { + "SYSTEM-10": { + "MODE": "MASTER", + "ENABLED": True, + "PEERS": { + bytes_3(7301896): { + "CONNECTION": "YES", + "CONNECTED": 1000, + "OPTIONS": b"TS2=730444;TIMER=15;", + }, + }, + }, + } + doc = build_topology(systems, seq=1) + peer = doc["systems"][0]["peers"][0] + assert peer["ts2_static"] == ["730444"] + + +def test_build_topology_exports_master_static_tgs() -> None: + systems = { + "MASTER-A": { + "MODE": "MASTER", + "ENABLED": True, + "TS1_STATIC": "91,92", + "TS2_STATIC": "730", + "PEERS": {}, + }, + } + doc = build_topology(systems, seq=1) + master = doc["systems"][0] + assert master["ts1_static"] == ["91", "92"] + assert master["ts2_static"] == ["730"] + + def test_build_topology_exports_peer_connected_at() -> None: login_ts = 1717555100.0 systems = { diff --git a/tests/infrastructure/test_mqtt_auth.py b/tests/infrastructure/test_mqtt_auth.py new file mode 100644 index 0000000..0e57217 --- /dev/null +++ b/tests/infrastructure/test_mqtt_auth.py @@ -0,0 +1,73 @@ +"""MQTT username/password authentication wiring.""" + +from __future__ import annotations + +from typing import Any + +from adn_server.infrastructure.twisted_adapters.report.mqtt_config import MqttBroker, MqttSettings +from adn_server.infrastructure.twisted_adapters.report.mqtt_publisher import PahoReportMqttPublisher + + +class _AuthCapturingClient: + def __init__(self, **kwargs: Any) -> None: + self.kwargs = kwargs + self.auth: tuple[str, str | None] | None = None + self.tls: dict[str, Any] | None = None + + def username_pw_set(self, username: str, password: str | None) -> None: + self.auth = (username, password) + + def tls_set(self, **kwargs: Any) -> None: + self.tls = kwargs + + def connect(self, host: str, port: int, keepalive: int) -> None: + self.connect_args = (host, port, keepalive) + + def loop_start(self) -> None: + pass + + +class _FakeMqttModule: + MQTTv311 = 4 + CallbackAPIVersion = type("CallbackAPIVersion", (), {"VERSION2": 2}) + + def Client(self, **kwargs: Any) -> _AuthCapturingClient: + return _AuthCapturingClient(**kwargs) + + +def _settings(**kwargs: Any) -> MqttSettings: + defaults = { + "broker": MqttBroker(host="broker", port=1883, use_tls=False, display_url="mqtt://broker:1883"), + "topic_prefix": "adn/1", + "client_id": "test", + "username": "ops", + "password": "secret", + "qos": 0, + } + defaults.update(kwargs) + return MqttSettings(**defaults) + + +def test_start_applies_username_password(monkeypatch): + fake = _FakeMqttModule() + monkeypatch.setattr( + "adn_server.infrastructure.twisted_adapters.report.mqtt_publisher.mqtt", + fake, + ) + pub = PahoReportMqttPublisher(_settings()) + pub.start(wire=object(), get_systems=lambda: {}, get_bridges=lambda: {}) + assert pub._client is not None + assert pub._client.auth == ("ops", "secret") + + +def test_start_tls_with_cafile(monkeypatch): + fake = _FakeMqttModule() + monkeypatch.setattr( + "adn_server.infrastructure.twisted_adapters.report.mqtt_publisher.mqtt", + fake, + ) + broker = MqttBroker(host="secure", port=8883, use_tls=True, display_url="mqtts://secure:8883") + pub = PahoReportMqttPublisher(_settings(broker=broker, cafile="/etc/ssl/certs/ca.pem")) + pub.start(wire=object(), get_systems=lambda: {}, get_bridges=lambda: {}) + assert pub._client is not None + assert pub._client.tls == {"ca_certs": "/etc/ssl/certs/ca.pem"} diff --git a/tests/infrastructure/test_mqtt_config.py b/tests/infrastructure/test_mqtt_config.py new file mode 100644 index 0000000..585f672 --- /dev/null +++ b/tests/infrastructure/test_mqtt_config.py @@ -0,0 +1,134 @@ +"""Tests for optional REPORTS.MQTT configuration.""" + +from __future__ import annotations + +from adn_server.infrastructure.config_validator import validate_config +from adn_server.infrastructure.twisted_adapters.report.mqtt_config import ( + mqtt_settings_from_config, + parse_mqtt_broker, +) + + +def test_mqtt_disabled_by_default(): + config = {"REPORTS": {"REPORT": True}, "GLOBAL": {"SERVER_ID": 73010}} + assert mqtt_settings_from_config(config) is None + + +def test_mqtt_disabled_when_url_without_enabled(): + config = { + "REPORTS": {"MQTT_URL": "mqtt://127.0.0.1:1883"}, + "GLOBAL": {"SERVER_ID": 73010}, + } + assert mqtt_settings_from_config(config) is None + + +def test_mqtt_enabled_with_nested_block(): + config = { + "GLOBAL": {"SERVER_ID": 73010}, + "REPORTS": { + "MQTT": { + "ENABLED": True, + "URL": "mqtt://broker.example:1883", + "TOPIC_PREFIX": "lab/adn", + "QOS": 1, + } + }, + } + settings = mqtt_settings_from_config(config) + assert settings is not None + assert settings.broker.display_url == "mqtt://broker.example:1883" + assert settings.topic_prefix == "lab/adn" + assert settings.qos == 1 + + +def test_mqtt_default_topic_prefix_from_server_id(): + config = { + "GLOBAL": {"SERVER_ID": 99999}, + "REPORTS": {"MQTT": {"ENABLED": True, "URL": "mqtt://127.0.0.1:1883"}}, + } + settings = mqtt_settings_from_config(config) + assert settings is not None + assert settings.topic_prefix == "adn/99999" + + +def test_validate_mqtt_enabled_requires_url(): + config = { + "GLOBAL": {"SERVER_ID": 1}, + "REPORTS": {"MQTT": {"ENABLED": True}}, + "SYSTEMS": {}, + } + try: + validate_config(config) + raised = False + except Exception as exc: + raised = True + assert "REPORTS.MQTT.URL" in str(exc) + assert raised + + +def test_mqtt_credentials_from_url(): + config = { + "GLOBAL": {"SERVER_ID": 1}, + "REPORTS": { + "MQTT": { + "ENABLED": True, + "URL": "mqtt://reportuser:secr%40et@mqtt.example:1883", + } + }, + } + settings = mqtt_settings_from_config(config) + assert settings is not None + assert settings.username == "reportuser" + assert settings.password == "secr@et" + assert settings.broker.host == "mqtt.example" + assert "reportuser" not in settings.broker.display_url + + +def test_mqtt_yaml_credentials_override_url(): + config = { + "GLOBAL": {"SERVER_ID": 1}, + "REPORTS": { + "MQTT": { + "ENABLED": True, + "URL": "mqtt://urluser:urlpass@127.0.0.1:1883", + "USERNAME": "yamluser", + "PASSWORD": "yamlpass", + } + }, + } + settings = mqtt_settings_from_config(config) + assert settings is not None + assert settings.username == "yamluser" + assert settings.password == "yamlpass" + + +def test_mqtt_client_id_derived_from_server_id(): + config = { + "GLOBAL": {"SERVER_ID": 7302}, + "REPORTS": {"MQTT": {"ENABLED": True, "URL": "mqtt://127.0.0.1:1883"}}, + } + settings = mqtt_settings_from_config(config) + assert settings is not None + assert settings.client_id.startswith("adn-server-7302-") + suffix = settings.client_id.removeprefix("adn-server-7302-") + assert len(suffix) == 8 + assert all(c in "0123456789abcdef" for c in suffix) + + +def test_mqtt_topic_prefix_from_normalized_server_id_bytes(): + """After config_normalizer, SERVER_ID is 4-byte big-endian (7302 -> b'\\x00\\x00\\x1c\\x86').""" + config = { + "GLOBAL": {"SERVER_ID": (7302).to_bytes(4, "big")}, + "REPORTS": {"MQTT": {"ENABLED": True, "URL": "mqtt://127.0.0.1:1883"}}, + } + settings = mqtt_settings_from_config(config) + assert settings is not None + assert settings.topic_prefix == "adn/7302" + assert settings.client_id.startswith("adn-server-7302-") + + +def test_parse_mqtts_default_port(): + broker, user, passwd = parse_mqtt_broker("mqtts://secure.example") + assert broker.port == 8883 + assert broker.use_tls is True + assert user is None and passwd is None diff --git a/tests/infrastructure/test_mqtt_publish_filter.py b/tests/infrastructure/test_mqtt_publish_filter.py new file mode 100644 index 0000000..fba1528 --- /dev/null +++ b/tests/infrastructure/test_mqtt_publish_filter.py @@ -0,0 +1,39 @@ +"""MQTT publish type filtering.""" + +from __future__ import annotations + +import json + +from adn_server.infrastructure.twisted_adapters.report.mqtt_config import MqttBroker, MqttSettings +from adn_server.infrastructure.twisted_adapters.report.mqtt_publisher import PahoReportMqttPublisher +from adn_server.infrastructure.twisted_adapters.report.opcodes import REPORT_OPCODES + + +class _FakeClient: + def __init__(self) -> None: + self.published: list[tuple[str, bytes]] = [] + + def publish(self, topic: str, payload: bytes, qos: int, retain: bool) -> object: + del qos, retain + self.published.append((topic, payload)) + return type("Info", (), {"rc": 0})() + + +def test_publish_frames_skips_routing_by_default(): + pub = PahoReportMqttPublisher( + MqttSettings( + broker=MqttBroker(host="h", port=1883, use_tls=False, display_url="mqtt://h:1883"), + topic_prefix="adn/1", + client_id="c", + username=None, + password=None, + qos=0, + ) + ) + pub._connected = True + pub._client = _FakeClient() + routing = REPORT_OPCODES["ROUTING_TABLE_SND"] + b'{"type":"routing_table","seq":1,"routes":[]}' + voice = REPORT_OPCODES["VOICE_EVENT_SND"] + json.dumps({"type": "voice_event", "phase": "start"}).encode() + pub.publish_frames((routing, voice)) + assert len(pub._client.published) == 1 + assert pub._client.published[0][0] == "adn/1/voice_event" diff --git a/tests/infrastructure/test_mqtt_publisher.py b/tests/infrastructure/test_mqtt_publisher.py new file mode 100644 index 0000000..a261ff7 --- /dev/null +++ b/tests/infrastructure/test_mqtt_publisher.py @@ -0,0 +1,174 @@ +"""Tests for MQTT publisher frame fan-out.""" + +from __future__ import annotations + +import json +from typing import Any + +from adn_server.infrastructure.twisted_adapters.report.mqtt_config import MqttBroker, MqttSettings +from adn_server.infrastructure.twisted_adapters.report.mqtt_publisher import PahoReportMqttPublisher +from adn_server.infrastructure.twisted_adapters.report.opcodes import REPORT_OPCODES + + +class _FakeMqttModule: + MQTTv311 = 4 + MQTT_ERR_SUCCESS = 0 + CallbackAPIVersion = type("CallbackAPIVersion", (), {"VERSION2": 2}) + + class Client: + def __init__(self, **kwargs: Any) -> None: + self.kwargs = kwargs + self.published: list[tuple[str, bytes, int, bool]] = [] + + def username_pw_set(self, username: str, password: str | None) -> None: + self.auth = (username, password) + + def connect(self, host: str, port: int, keepalive: int) -> None: + self.host = (host, port, keepalive) + + def loop_start(self) -> None: + pass + + def publish(self, topic: str, payload: bytes, qos: int, retain: bool) -> Any: + self.published.append((topic, payload, qos, retain)) + return type("Info", (), {"rc": 0})() + + def loop_stop(self) -> None: + pass + + def disconnect(self) -> None: + pass + + +def test_publish_frames_maps_topics(monkeypatch): + fake = _FakeMqttModule() + monkeypatch.setattr( + "adn_server.infrastructure.twisted_adapters.report.mqtt_publisher.mqtt", + fake, + ) + pub = PahoReportMqttPublisher( + MqttSettings( + broker=MqttBroker(host="127.0.0.1", port=1883, use_tls=False, display_url="mqtt://127.0.0.1:1883"), + topic_prefix="adn/test", + client_id="test", + username=None, + password=None, + qos=0, + ) + ) + pub._connected = True + pub._client = fake.Client() + voice = {"type": "voice_event", "phase": "start"} + frame = REPORT_OPCODES["VOICE_EVENT_SND"] + json.dumps(voice).encode("utf-8") + pub.publish_frames((frame,)) + assert pub._client.published == [ + ("adn/test/voice_event", json.dumps(voice).encode("utf-8"), 0, False), + ] + + +def test_on_connect_publishes_shared_state(monkeypatch): + fake = _FakeMqttModule() + monkeypatch.setattr( + "adn_server.infrastructure.twisted_adapters.report.mqtt_publisher.mqtt", + fake, + ) + pub = PahoReportMqttPublisher( + MqttSettings( + broker=MqttBroker(host="127.0.0.1", port=1883, use_tls=False, display_url="mqtt://127.0.0.1:1883"), + topic_prefix="adn/7302", + client_id="test", + username=None, + password=None, + qos=0, + ) + ) + client = fake.Client() + pub._client = client + pub._get_systems = lambda: {"MASTER1": {"MODE": "MASTER", "ENABLED": True, "PEERS": {}}} + pub._on_connect(client, None, None, type("RC", (), {"value": 0})(), None) + assert len(client.published) == 1 + topic, body, qos, retain = client.published[0] + assert topic == "adn/7302/state" + assert qos == 0 + assert retain is True + assert json.loads(body.decode("utf-8"))["type"] == "dashboard_state" + + +def test_emit_dashboard_skips_unchanged_state(monkeypatch): + fake = _FakeMqttModule() + monkeypatch.setattr( + "adn_server.infrastructure.twisted_adapters.report.mqtt_publisher.mqtt", + fake, + ) + pub = PahoReportMqttPublisher( + MqttSettings( + broker=MqttBroker(host="127.0.0.1", port=1883, use_tls=False, display_url="mqtt://127.0.0.1:1883"), + topic_prefix="adn/test", + client_id="test", + username=None, + password=None, + qos=0, + ) + ) + pub._connected = True + pub._client = fake.Client() + systems = { + "ECHO": { + "MODE": "MASTER", + "ENABLED": True, + "PEERS": {b"\x00\x00\x00\x01": {"CONNECTION": "YES", "CONNECTED": 1}}, + }, + } + pub._emit_shared_dashboard(systems) + pub._emit_shared_dashboard(systems) + state_published = [ + t for t, b, _q, _r in pub._client.published if b and json.loads(b.decode())["type"] == "dashboard_state" + ] + assert state_published == ["adn/test/state"] + + +def test_publish_dashboard_publishes_shared_state(monkeypatch): + fake = _FakeMqttModule() + monkeypatch.setattr( + "adn_server.infrastructure.twisted_adapters.report.mqtt_publisher.mqtt", + fake, + ) + pub = PahoReportMqttPublisher( + MqttSettings( + broker=MqttBroker(host="127.0.0.1", port=1883, use_tls=False, display_url="mqtt://127.0.0.1:1883"), + topic_prefix="adn/7302", + client_id="test", + username=None, + password=None, + qos=0, + ) + ) + pub._connected = True + pub._client = fake.Client() + pub.publish_dashboard({}) + assert len(pub._client.published) == 1 + assert pub._client.published[0][0] == "adn/7302/state" + assert pub._client.published[0][3] is True + pub._client.published.clear() + pub.publish_dashboard( + { + "ECHO": { + "MODE": "MASTER", + "ENABLED": True, + "PEERS": {b"\x00\x00\x00\x01": {"CONNECTION": "YES", "CONNECTED": 1}}, + }, + } + ) + topic, _body, _qos, retain = pub._client.published[0] + assert topic == "adn/7302/state" + assert retain is True + pub.publish_dashboard( + { + "ECHO": { + "MODE": "MASTER", + "ENABLED": True, + "PEERS": {b"\x00\x00\x00\x01": {"CONNECTION": "YES", "CONNECTED": 1}}, + }, + } + ) + assert len(pub._client.published) == 1 diff --git a/tests/infrastructure/test_mqtt_reload.py b/tests/infrastructure/test_mqtt_reload.py new file mode 100644 index 0000000..46f03bf --- /dev/null +++ b/tests/infrastructure/test_mqtt_reload.py @@ -0,0 +1,123 @@ +"""MQTT reconnect on config reload.""" + +from __future__ import annotations + +from typing import Any + +from adn_server.infrastructure.twisted_adapters.report.mqtt_config import ( + MqttBroker, + MqttSettings, + mqtt_settings_from_config, +) +from adn_server.infrastructure.twisted_adapters.report.mqtt_publisher import reconcile_mqtt_publisher + + +class _RecordingMqtt: + def __init__(self) -> None: + self.stopped = False + self.started = False + + def start(self, wire: Any, get_systems: Any, get_bridges: Any) -> None: + del wire, get_systems, get_bridges + self.started = True + + def publish_frames(self, frames: tuple[bytes, ...]) -> None: + del frames + + def publish_dashboard(self, systems: dict[str, Any]) -> None: + del systems + + def stop(self) -> None: + self.stopped = True + + +class _FactoryStub: + def __init__(self) -> None: + self._mqtt = None + self.start_mqtt_calls = 0 + + def set_mqtt(self, mqtt: Any) -> None: + self._mqtt = mqtt + + def start_mqtt(self) -> None: + self.start_mqtt_calls += 1 + if self._mqtt is not None: + self._mqtt.start(None, lambda: {}, lambda: {}) + + +def _settings(url: str = "mqtt://127.0.0.1:1883") -> MqttSettings: + return MqttSettings( + broker=MqttBroker(host="127.0.0.1", port=1883, use_tls=False, display_url=url), + topic_prefix="adn/1", + client_id="adn-server-1-deadbeef", + username=None, + password=None, + qos=0, + ) + + +def test_reload_noop_when_mqtt_unchanged(): + factory = _FactoryStub() + current = _RecordingMqtt() + same = _settings() + result = reconcile_mqtt_publisher(factory, current, same, same, report_enabled=True) + assert result is current + assert not current.stopped + + +def test_reload_disconnects_when_mqtt_disabled(): + factory = _FactoryStub() + current = _RecordingMqtt() + result = reconcile_mqtt_publisher(factory, current, _settings(), None, report_enabled=True) + assert result is None + assert current.stopped + assert factory._mqtt is None + + +def test_reload_enables_mqtt(monkeypatch): + factory = _FactoryStub() + + class _NewPub(_RecordingMqtt): + pass + + monkeypatch.setattr( + "adn_server.infrastructure.twisted_adapters.report.mqtt_publisher.create_report_mqtt_publisher_from_settings", + lambda _s: _NewPub(), + ) + result = reconcile_mqtt_publisher(factory, None, None, _settings(), report_enabled=True) + assert isinstance(result, _NewPub) + assert result.started + assert factory.start_mqtt_calls == 1 + + +def test_reload_restarts_when_broker_changes(monkeypatch): + factory = _FactoryStub() + old = _RecordingMqtt() + + class _NewPub(_RecordingMqtt): + pass + + monkeypatch.setattr( + "adn_server.infrastructure.twisted_adapters.report.mqtt_publisher.create_report_mqtt_publisher_from_settings", + lambda _s: _NewPub(), + ) + result = reconcile_mqtt_publisher( + factory, + old, + _settings("mqtt://127.0.0.1:1883"), + _settings("mqtt://other:1883"), + report_enabled=True, + ) + assert old.stopped + assert isinstance(result, _NewPub) + assert result.started + + +def test_mqtt_settings_detect_enabled_toggle(): + off = {"GLOBAL": {"SERVER_ID": 1}, "REPORTS": {}} + on = { + "GLOBAL": {"SERVER_ID": 1}, + "REPORTS": {"MQTT": {"ENABLED": True, "URL": "mqtt://b:1883"}}, + } + assert mqtt_settings_from_config(off) is None + assert mqtt_settings_from_config(on) is not None diff --git a/tests/infrastructure/test_mqtt_topics.py b/tests/infrastructure/test_mqtt_topics.py new file mode 100644 index 0000000..82f2134 --- /dev/null +++ b/tests/infrastructure/test_mqtt_topics.py @@ -0,0 +1,40 @@ +"""Tests for MQTT topic mapping from report wire frames.""" + +from __future__ import annotations + +import json + +from adn_server.infrastructure.twisted_adapters.report.mqtt_topics import ( + mqtt_shared_state_topic, + topic_for_frame, +) +from adn_server.infrastructure.twisted_adapters.report.opcodes import REPORT_OPCODES + + +def _frame(opcode: bytes, payload: dict) -> bytes: + return opcode + json.dumps(payload, separators=(",", ":")).encode("utf-8") + + +def test_shared_state_topic(): + assert mqtt_shared_state_topic("adn/7302") == "adn/7302/state" + + +def test_topology_topic(): + frame = _frame(REPORT_OPCODES["TOPOLOGY_SND"], {"type": "topology", "seq": 1}) + assert topic_for_frame(frame, "adn/73010") == "adn/73010/topology" + + +def test_delta_topology_topic(): + frame = _frame( + REPORT_OPCODES["DELTA_SND"], + {"type": "delta", "since_seq": 1, "patch": {"type": "topology", "systems": []}}, + ) + assert topic_for_frame(frame, "pfx") == "pfx/delta/topology" + + +def test_delta_routing_topic(): + frame = _frame( + REPORT_OPCODES["DELTA_SND"], + {"type": "delta", "since_seq": 2, "patch": {"type": "routing_table", "bridges": {}}}, + ) + assert topic_for_frame(frame, "pfx") == "pfx/delta/routing_table" diff --git a/tests/infrastructure/test_report_server_mqtt.py b/tests/infrastructure/test_report_server_mqtt.py new file mode 100644 index 0000000..0c24b16 --- /dev/null +++ b/tests/infrastructure/test_report_server_mqtt.py @@ -0,0 +1,62 @@ +"""ReportServerFactory MQTT fan-out.""" + +from __future__ import annotations + +import json +from typing import Any + +from adn_server.application.ports import ReportMqttPublisher +from adn_server.infrastructure.twisted_adapters.report.opcodes import REPORT_OPCODES +from adn_server.infrastructure.twisted_adapters.report_server import ReportServerFactory + + +class RecordingMqtt(ReportMqttPublisher): + def __init__(self) -> None: + self.frames: list[tuple[bytes, ...]] = [] + self.dashboard_calls = 0 + self.started = False + + def start(self, wire: Any, get_systems: Any, get_bridges: Any) -> None: + del wire, get_systems, get_bridges + self.started = True + + def publish_frames(self, frames: tuple[bytes, ...]) -> None: + if frames: + self.frames.append(frames) + + def publish_dashboard(self, systems: dict[str, Any]) -> None: + del systems + self.dashboard_calls += 1 + + def stop(self) -> None: + pass + + +def test_send_config_publishes_state_not_tcp_wire(): + mqtt = RecordingMqtt() + factory = ReportServerFactory({"REPORTS": {}}, mqtt=mqtt) + factory.set_systems({"SYS": {"MODE": "MASTER", "ENABLED": True}}) + factory.send_config() + assert mqtt.frames == [] + assert mqtt.dashboard_calls == 1 + + +def test_send_bridge_does_not_mirror_tcp_wire_to_mqtt(): + mqtt = RecordingMqtt() + factory = ReportServerFactory({"REPORTS": {}}, mqtt=mqtt) + factory.set_bridges({"73010": {"ACTIVE": True}}) + factory.send_bridge() + assert mqtt.frames == [] + + +def test_send_bridge_event_publishes_voice_event_to_mqtt(): + mqtt = RecordingMqtt() + factory = ReportServerFactory({"REPORTS": {"PROTOCOL": "v2"}}, mqtt=mqtt) + factory.set_systems({}) + factory.set_bridges({}) + factory.send_bridge_event("GROUP VOICE,START,RX,MASTER-A,2155905152,1001,3120001,2,52090") + assert len(mqtt.frames) == 1 + frame = mqtt.frames[0][0] + assert frame[:1] == REPORT_OPCODES["VOICE_EVENT_SND"] + payload = json.loads(frame[1:].decode("utf-8")) + assert payload["type"] == "voice_event"