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.pull/1/head
parent
99169e77ca
commit
471e1795bf
@ -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
|
||||
@ -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,
|
||||
)
|
||||
@ -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)
|
||||
@ -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"
|
||||
@ -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"] == {}
|
||||
@ -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"}
|
||||
@ -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
|
||||
@ -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"
|
||||
@ -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
|
||||
@ -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
|
||||
@ -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"
|
||||
@ -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"
|
||||
Loading…
Reference in new issue