From 46c5c0a84659fdd0d298f1da0709e8a2814ad3ef Mon Sep 17 00:00:00 2001 From: bonamin Date: Wed, 30 Sep 2026 16:01:41 +0300 Subject: [PATCH] fix(mqtt): treat broker-replayed retained messages as stale state 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 --- backend/database/pg_mqtt.py | 17 +++--- backend/mqtt/client.py | 17 ++++-- backend/mqtt/logger.py | 56 +++++++++++++++---- backend/mqtt/presence.py | 26 +++++++++ backend/mqtt/router.py | 3 +- .../pages/bellcloud/devices/DeviceDetail.jsx | 8 +++ .../pages/bellcloud/devices/DeviceList.jsx | 8 +++ .../pages/bellcloud/devices/tabs/LogsTab.jsx | 3 + 8 files changed, 111 insertions(+), 27 deletions(-) create mode 100644 backend/mqtt/presence.py diff --git a/backend/database/pg_mqtt.py b/backend/database/pg_mqtt.py index d2caffc..132d199 100644 --- a/backend/database/pg_mqtt.py +++ b/backend/database/pg_mqtt.py @@ -444,20 +444,17 @@ async def insert_boot_event(device_serial: str, boot_count: int | None, crash_task: str | None = None, crash_pc: int | None = None, crash_exc_cause: int | None = None, - crash_exc_vaddr: int | None = None) -> int | None: - """Insert a boot event, unless it is a repeat of the device's most recent - one (returns None then). boot_report is published retained on system/info, - so the broker redelivers the last one every time the backend (re)subscribes - — without this check each backend restart logged the device's last boot - (crash detail included) again as a brand-new event. + crash_exc_vaddr: int | None = None, + skip_if_latest: bool = False) -> int | None: + """Insert a boot event. With skip_if_latest=True (a retained replay of the + device's last boot_report), skip it and return None when it matches the + device's most recent row: that boot is already recorded. Only the LATEST row is compared, never "any row with this boot_count": the firmware's lifetime counter gets reset (reflash / telemetry reset), so the - same boot_count legitimately appears again for a later, different boot. A - genuine new boot always differs from the latest row — the counter either - moved forward or was reset to a lower number.""" + same boot_count legitimately appears again for a later, different boot.""" async with AsyncSessionLocal() as session: - if boot_count is not None: + if skip_if_latest and boot_count is not None: latest = await session.execute( text(""" SELECT boot_count, reset_reason FROM device_boot_events diff --git a/backend/mqtt/client.py b/backend/mqtt/client.py index f88a8d9..e7691d9 100644 --- a/backend/mqtt/client.py +++ b/backend/mqtt/client.py @@ -99,10 +99,17 @@ class MqttManager: 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), + self._process_message(serial, topic_type, payload, topic, retained), self._loop, ) except json.JSONDecodeError: @@ -111,15 +118,16 @@ class MqttManager: logger.error(f"Error processing MQTT message: {e}") async def _process_message(self, serial: str, topic_type: str, - payload: dict, raw_topic: str): + payload: dict, raw_topic: str, retained: bool = False): from mqtt.logger import handle_message - await handle_message(serial, topic_type, payload) + 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) @@ -170,6 +178,7 @@ class MqttManager: 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: @@ -179,7 +188,7 @@ class MqttManager: age = now - received.timestamp() except (ValueError, TypeError, KeyError): continue - if age > PING_ONLINE_WINDOW_SECONDS: + if age > PING_ONLINE_WINDOW_SECONDS or presence.is_marked_offline(hb["device_serial"]): continue self.publish_command( device_serial=hb["device_serial"], diff --git a/backend/mqtt/logger.py b/backend/mqtt/logger.py index f61c08b..bef1adc 100644 --- a/backend/mqtt/logger.py +++ b/backend/mqtt/logger.py @@ -1,6 +1,7 @@ import logging import time import database as db +from mqtt import presence logger = logging.getLogger("mqtt.logger") @@ -15,15 +16,22 @@ LEVEL_MAP = { } -async def handle_message(serial: str, topic_type: str, payload: dict): +async def handle_message(serial: str, topic_type: str, payload: dict, + retained: bool = False): + """`retained` is True when the broker replayed this from its retained store + on (re)subscribe: a stale copy of the device's last message, NOT a new + event. The firmware publishes status/heartbeat, system/alerts, system/info + and status/playback retained, so every backend restart replays them for + every device that ever connected. Handlers for those topics must never + record a retained replay as something that just happened.""" try: # v2 topic set — see project-vesper's vesper_mqtt_topic_spec_v2.md. if topic_type == "status/heartbeat": - await _handle_heartbeat(serial, payload) + await _handle_heartbeat(serial, payload, retained) elif topic_type == "system/alerts": - await _handle_alerts(serial, payload) + await _handle_alerts(serial, payload, retained) elif topic_type == "system/info": - await _handle_info(serial, payload) + await _handle_info(serial, payload, retained) elif topic_type == "system/logs": await _handle_log(serial, payload) elif topic_type == "system/metrics": @@ -40,7 +48,21 @@ async def handle_message(serial: str, topic_type: str, payload: dict): logger.error(f"Error handling {topic_type} for {serial}: {e}") -async def _handle_heartbeat(serial: str, payload: dict): +async def _handle_heartbeat(serial: str, payload: dict, retained: bool = False): + # Online status is "a heartbeat row newer than 90s", so only a live sign of + # life may be stored: + # - retained replay: the device's last heartbeat from whenever. Storing it + # with received_at=now() marked every device ever seen (even long-dead + # ones) online for 90s after each backend restart. + # - state "offline": the LWT / graceful-disconnect marker published when + # the device goes AWAY. Not a heartbeat, but it IS current state even + # when replayed (a retained message is the device's last word). + if payload.get("state") == "offline": + presence.mark_offline(serial) + return + if retained: + return + presence.mark_alive(serial) # Store silently — do not log as a visible event. # The console surfaces an alert only when the device goes silent (no heartbeat for 90s). # v2 heartbeat payload is FLAT — no {"status","type","payload"} wrapper, @@ -76,7 +98,7 @@ async def _handle_log(serial: str, payload: dict): ) -async def _handle_alerts(serial: str, payload: dict): +async def _handle_alerts(serial: str, payload: dict, retained: bool = False): subsystem = payload.get("subsystem", "") state = payload.get("state", "") if not subsystem or not state: @@ -86,16 +108,18 @@ async def _handle_alerts(serial: str, payload: dict): if state == "CLEARED": await db.delete_alert(serial, subsystem) else: + # Retained replays still sync the current-alert row, so a fresh DB + # learns about alerts raised while the backend was down. changed = await db.upsert_alert(serial, subsystem, state, payload.get("msg")) # Append-only history — survives past the alert being resolved, used to # answer "when was the most recent issue" even once it's cleared. - # Skipped when nothing changed: alerts are retained, so the broker - # re-sends every active one on each backend (re)subscribe. - if changed: + # A live alert is always a new occurrence (the same fault can recur); + # a retained replay only counts if it isn't already the current state. + if changed or not retained: await db.insert_alert_event(serial, subsystem, state, payload.get("msg")) -async def _handle_info(serial: str, payload: dict): +async def _handle_info(serial: str, payload: dict, retained: bool = False): event_type = payload.get("type", "") # boot_report carries structured fields (boot_count, crash detail, etc.) — @@ -106,7 +130,12 @@ async def _handle_info(serial: str, payload: dict): # Note: diagnostics_report used to also arrive here (F-056) but moved to # its own system/metrics topic in the v2 spec — see _handle_metrics below. if event_type == "boot_report": - await _handle_boot_report(serial, payload) + await _handle_boot_report(serial, payload, retained) + return + + # Any other info event replayed from the retained store (e.g. an old + # playback_started) already happened in the past — don't log it again. + if retained: return data = payload.get("payload", {}) @@ -126,7 +155,7 @@ async def _handle_info(serial: str, payload: dict): ) -async def _handle_boot_report(serial: str, payload: dict): +async def _handle_boot_report(serial: str, payload: dict, retained: bool = False): """Parses the firmware's consolidated boot_report event (see CommunicationRouter::reportBootOnce() / publishInfo()) into device_boot_events. Like all status/info events, the report fields are @@ -145,6 +174,9 @@ async def _handle_boot_report(serial: str, payload: dict): crash_pc=crash.get("pc"), crash_exc_cause=crash.get("exc_cause"), crash_exc_vaddr=crash.get("exc_vaddr"), + # A live boot_report is always a new boot. A retained replay is the + # device's last boot, recorded only if we missed it while down. + skip_if_latest=retained, ) diff --git a/backend/mqtt/presence.py b/backend/mqtt/presence.py new file mode 100644 index 0000000..60c1cec --- /dev/null +++ b/backend/mqtt/presence.py @@ -0,0 +1,26 @@ +"""Devices that announced they went offline. + +The firmware's LWT (unclean disconnect) and its graceful-disconnect message both +publish {"state": "offline", "ok": false} to status/heartbeat, retained. Online +status is otherwise "a heartbeat newer than 90s", which would keep showing a +device that just dropped as online for up to 90s. This remembers the marker so +/mqtt/status and the ping loop treat the device as offline right away. + +In-memory only, which is enough: on backend start the broker replays each +device's retained heartbeat, so a device whose last word was "offline" gets +re-marked, and any live heartbeat clears the mark. +""" + +_offline: set[str] = set() + + +def mark_offline(serial: str) -> None: + _offline.add(serial) + + +def mark_alive(serial: str) -> None: + _offline.discard(serial) + + +def is_marked_offline(serial: str) -> bool: + return serial in _offline diff --git a/backend/mqtt/router.py b/backend/mqtt/router.py index 79a083e..4e1ebd2 100644 --- a/backend/mqtt/router.py +++ b/backend/mqtt/router.py @@ -11,6 +11,7 @@ from mqtt.models import ( LatestDiagnosticsEntry, LatestPingEntry, DeviceReportListResponse, ) from mqtt.client import mqtt_manager +from mqtt import presence import database as db from datetime import datetime, timezone @@ -40,7 +41,7 @@ async def get_all_device_status( devices.append(DeviceMqttStatus( device_serial=hb["device_serial"], - online=seconds_ago < 90, + online=seconds_ago < 90 and not presence.is_marked_offline(hb["device_serial"]), last_heartbeat=HeartbeatEntry(**hb), seconds_since_heartbeat=seconds_ago, last_alert_event=AlertEventEntry(**alert_event) if alert_event else None, diff --git a/frontend/src/pages/bellcloud/devices/DeviceDetail.jsx b/frontend/src/pages/bellcloud/devices/DeviceDetail.jsx index 920361f..e036b5d 100644 --- a/frontend/src/pages/bellcloud/devices/DeviceDetail.jsx +++ b/frontend/src/pages/bellcloud/devices/DeviceDetail.jsx @@ -194,7 +194,15 @@ export default function DeviceDetail() { // v2 heartbeat payload is FLAT (no nested .payload wrapper) — see // vesper_mqtt_topic_spec_v2.md. field names: fw_version, uptime_human. if (msg.type !== 'status/heartbeat') return + // Retained replay = the device's LAST heartbeat, re-sent by the broker when + // the backend reconnects — not proof the device is alive now. + if (msg.retained) return const hb = msg.payload || {} + // LWT / graceful-disconnect marker: the device just went away. + if (hb.state === 'offline') { + setMqttStatus(prev => prev ? { ...prev, online: false } : prev) + return + } setMqttStatus(prev => ({ device_serial: deviceSerial, online: true, diff --git a/frontend/src/pages/bellcloud/devices/DeviceList.jsx b/frontend/src/pages/bellcloud/devices/DeviceList.jsx index 91a79d7..adc1df0 100644 --- a/frontend/src/pages/bellcloud/devices/DeviceList.jsx +++ b/frontend/src/pages/bellcloud/devices/DeviceList.jsx @@ -436,9 +436,17 @@ export default function DeviceList() { if (msg?.type !== 'status/heartbeat') return // v2 heartbeat payload is FLAT (no nested .payload wrapper) — see // vesper_mqtt_topic_spec_v2.md. field names: fw_version, uptime_human. + // Retained replay = the device's LAST heartbeat, re-sent by the broker + // when the backend reconnects — not proof the device is alive now. + if (msg.retained) return const hb = msg.payload || {} const serial = msg.device_serial if (!serial) return + // LWT / graceful-disconnect marker: the device just went away. + if (hb.state === 'offline') { + setMqttStatusMap(prev => prev[serial] ? { ...prev, [serial]: { ...prev[serial], online: false } } : prev) + return + } setMqttStatusMap(prev => ({ ...prev, [serial]: { diff --git a/frontend/src/pages/bellcloud/devices/tabs/LogsTab.jsx b/frontend/src/pages/bellcloud/devices/tabs/LogsTab.jsx index f682155..d530387 100644 --- a/frontend/src/pages/bellcloud/devices/tabs/LogsTab.jsx +++ b/frontend/src/pages/bellcloud/devices/tabs/LogsTab.jsx @@ -523,6 +523,9 @@ function UnifiedPanel({ sn, channel, onChannelChange, range }) { const handleWsMessage = useCallback((data) => { if (data.device_serial !== sn) return + // Retained replays (broker re-sending the last message after the backend + // reconnects) are old events, not live ones. + if (data.retained) return if (data.type === 'system/logs' || data.type === 'system/info') { const entry = { _kind: 'log',