The firmware publishes status/heartbeat, system/alerts, system/info and
status/playback with retain=true. On every backend (re)connect - every
restart and every uvicorn --reload - the broker replays the last message on
each of those topics for every device that ever connected. We handled those
replays as if they had just happened:
- heartbeats: a row with received_at=now() for every device, so devices
that have been dead for months showed ONLINE for 90s after each restart
and got pinged. ~770k such rows exist locally.
- boot_report: the last boot logged again as a new reboot (the phantom
PANIC entries on the Health tab).
- alerts / other info events: logged again as new occurrences.
MQTT delivers retain=1 only for replays caused by a new subscription; live
publishes always arrive with retain=0. The flag is now passed through to the
handlers and the WS broadcast:
- heartbeat: replays are not stored. A live heartbeat is.
- {"state":"offline"} heartbeat (LWT / graceful disconnect) is no longer
stored as a sign of life. It marks the device offline immediately in
a small in-memory set (mqtt/presence.py) used by /mqtt/status and the ping
loop; a later live heartbeat clears it. Replayed offline markers also mark
offline, since a retained message is the device's last word.
- boot_report: live -> always a new boot. Replay -> stored only if it
differs from the device's latest boot row (i.e. we missed it while down).
- alerts: replay still syncs the current-alert row; history gets a row on a
live alert (even an identical repeat - faults recur) or on a replay that
changes state. Replaces the 98dd16b rule that dropped identical live alerts.
- other info events: replays are not logged.
- Frontend (DeviceList, DeviceDetail, LogsTab) ignores retained WS messages
for live updates, and flips a device offline on the offline marker instead
of marking it online.
Verified locally after a backend restart: only the 7 actually-live devices
got new heartbeat rows (none from the replays), and no boot/alert/info rows
were created.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
201 lines
7.5 KiB
Python
201 lines
7.5 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:])
|
|
# The broker sets retain=1 only on a message it replays from its
|
|
# retained store because we just (re)subscribed: a stale snapshot
|
|
# of the device's last state, not something that just happened.
|
|
# Live publishes always arrive with retain=0, even on topics the
|
|
# firmware publishes retained. Every backend restart (and every
|
|
# uvicorn --reload) replays these for every device.
|
|
retained = bool(msg.retain)
|
|
|
|
if self._loop and self._loop.is_running():
|
|
asyncio.run_coroutine_threadsafe(
|
|
self._process_message(serial, topic_type, payload, topic, retained),
|
|
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, retained: bool = False):
|
|
from mqtt.logger import handle_message
|
|
await handle_message(serial, topic_type, payload, retained=retained)
|
|
|
|
ws_data = {
|
|
"type": topic_type,
|
|
"device_serial": serial,
|
|
"payload": payload,
|
|
"topic": raw_topic,
|
|
"retained": retained,
|
|
}
|
|
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
|
|
from mqtt import presence
|
|
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 or presence.is_marked_offline(hb["device_serial"]):
|
|
continue
|
|
self.publish_command(
|
|
device_serial=hb["device_serial"],
|
|
cmd="ping",
|
|
contents={"ts": int(now * 1000)},
|
|
)
|
|
|
|
|
|
mqtt_manager = MqttManager()
|