Files
bellsystems-cp/backend/mqtt/client.py
T
bonaminandClaude Sonnet 5 7c533b9245 feat(mqtt): add device health telemetry and migrate to v2 topic spec
Two efforts that landed together because the v2 topic work extends
tables the health-telemetry effort added days earlier in the same
files/functions, making them impractical to separate cleanly:

Health/diagnostics telemetry (schema, Jul 13-17):
- New Postgres tables: device_alert_events, device_boot_events,
  device_ping_samples, device_diagnostics_reports, plus a `source`
  column on device_logs to distinguish log origins
- Query/service layer in pg_mqtt.py and database/__init__.py for
  inserting and listing this history, plus a "latest metrics" endpoint
  combining most-recent diagnostics + ping RTT per device
- mqtt/router.py gains list endpoints for alert/boot/ping/diagnostics
  history, consumed by the upcoming Health tab

MQTT v2 topic migration (Sep 21):
- Heartbeat payload flattened per vesper_mqtt_topic_spec_v2.md, adding
  rssi/free_heap/state/ok fields
- Command replies move to control/ack, device-initiated events to
  control/reports; mqtt/client.py subscribes to the new topic set and
  runs a ping_loop (wired up in main.py) for RTT sampling
- mqtt/logger.py and pg_mqtt.py updated to parse and persist the new
  payload shape alongside the legacy fields for backwards compatibility

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-21 18:24:36 +03:00

192 lines
6.9 KiB
Python

import json
import logging
import asyncio
import time
from typing import Set
import paho.mqtt.client as paho_mqtt
from config import settings
logger = logging.getLogger("mqtt.client")
PING_INTERVAL_SECONDS = 60
# Only devices heard from within this window get pinged — no point spending
# broker traffic/RTT samples on a device that's already known offline; its
# heartbeat-derived "online" state will already reflect that on the console.
PING_ONLINE_WINDOW_SECONDS = 90
class MqttManager:
"""Singleton MQTT client manager."""
def __init__(self):
self._client: paho_mqtt.Client | None = None
self._connected = False
self._loop: asyncio.AbstractEventLoop | None = None
self._ws_subscribers: Set = set()
@property
def connected(self) -> bool:
return self._connected
def start(self, loop: asyncio.AbstractEventLoop):
self._loop = loop
self._client = paho_mqtt.Client(
callback_api_version=paho_mqtt.CallbackAPIVersion.VERSION2,
client_id=settings.mqtt_client_id,
clean_session=True,
)
if settings.mqtt_admin_username and settings.mqtt_admin_password:
self._client.username_pw_set(
settings.mqtt_admin_username,
settings.mqtt_admin_password,
)
self._client.on_connect = self._on_connect
self._client.on_disconnect = self._on_disconnect
self._client.on_message = self._on_message
try:
self._client.connect_async(
settings.mqtt_broker_host,
settings.mqtt_broker_port,
)
self._client.loop_start()
logger.info(f"MQTT client connecting to {settings.mqtt_broker_host}:{settings.mqtt_broker_port}")
except Exception as e:
logger.warning(f"MQTT client failed to start: {e}")
def stop(self):
if self._client:
self._client.loop_stop()
self._client.disconnect()
self._connected = False
logger.info("MQTT client disconnected")
def _on_connect(self, client, userdata, flags, reason_code, properties):
if reason_code == 0:
self._connected = True
logger.info("MQTT connected, subscribing to topics")
# v2 topic set — see vesper_mqtt_topic_spec_v2.md in the firmware repo.
# control/command is inbound-to-device only; the console never subscribes to it.
client.subscribe([
("vesper/+/control/ack", 1),
("vesper/+/control/reports", 1),
("vesper/+/status/heartbeat", 1),
("vesper/+/status/playback", 1),
("vesper/+/system/alerts", 1),
("vesper/+/system/info", 1),
("vesper/+/system/logs", 0),
("vesper/+/system/metrics", 0),
])
else:
logger.error(f"MQTT connection failed: {reason_code}")
def _on_disconnect(self, client, userdata, flags, reason_code, properties):
self._connected = False
if reason_code != 0:
logger.warning(f"MQTT disconnected unexpectedly: {reason_code}")
def _on_message(self, client, userdata, msg):
"""Called from paho thread — bridge to asyncio."""
try:
topic = msg.topic
payload = json.loads(msg.payload.decode("utf-8"))
parts = topic.split("/")
if len(parts) < 3 or parts[0] != "vesper":
return
serial = parts[1]
topic_type = "/".join(parts[2:])
if self._loop and self._loop.is_running():
asyncio.run_coroutine_threadsafe(
self._process_message(serial, topic_type, payload, topic),
self._loop,
)
except json.JSONDecodeError:
logger.warning(f"Invalid JSON on topic {msg.topic}")
except Exception as e:
logger.error(f"Error processing MQTT message: {e}")
async def _process_message(self, serial: str, topic_type: str,
payload: dict, raw_topic: str):
from mqtt.logger import handle_message
await handle_message(serial, topic_type, payload)
ws_data = {
"type": topic_type,
"device_serial": serial,
"payload": payload,
"topic": raw_topic,
}
await self._broadcast_ws(ws_data)
async def _broadcast_ws(self, data: dict):
msg = json.dumps(data)
dead = set()
for ws in self._ws_subscribers:
try:
await ws.send_text(msg)
except Exception:
dead.add(ws)
self._ws_subscribers -= dead
def add_ws_subscriber(self, websocket):
self._ws_subscribers.add(websocket)
def remove_ws_subscriber(self, websocket):
self._ws_subscribers.discard(websocket)
def publish_command(self, device_serial: str, cmd: str,
contents: dict) -> bool:
if not self._client or not self._connected:
return False
topic = f"vesper/{device_serial}/control/command"
payload = json.dumps({"v": 2, "cmd": cmd, "contents": contents})
result = self._client.publish(topic, payload, qos=1)
return result.rc == paho_mqtt.MQTT_ERR_SUCCESS
async def ping_loop(self):
"""Periodically pings every recently-online device with a client
timestamp so mqtt/logger.py::_handle_data_response can compute RTT
from the echoed pong. Deliberately bypasses db.insert_command — this
is a background health check, not a user-initiated command, and
shouldn't clutter the Control tab's command history.
Only available on RTC-equipped firmware builds (see API Reference —
ping is compiled out on agnus/agnus-mini). A pong simply never
arrives for those devices, so no ping-latency samples accumulate for
them; the Health tab handles an empty series as "unsupported".
"""
while True:
await asyncio.sleep(PING_INTERVAL_SECONDS)
try:
await self._ping_online_devices()
except Exception as e:
logger.error(f"Ping loop error: {e}")
async def _ping_online_devices(self):
import database as db
heartbeats = await db.get_latest_heartbeats()
now = time.time()
for hb in heartbeats:
try:
from datetime import datetime
received = datetime.fromisoformat(hb["received_at"])
age = now - received.timestamp()
except (ValueError, TypeError, KeyError):
continue
if age > PING_ONLINE_WINDOW_SECONDS:
continue
self.publish_command(
device_serial=hb["device_serial"],
cmd="ping",
contents={"ts": int(now * 1000)},
)
mqtt_manager = MqttManager()