From 98dd16b5978febcf5598eb641ffe9de41e85af6e Mon Sep 17 00:00:00 2001 From: bonamin Date: Wed, 30 Sep 2026 15:48:48 +0300 Subject: [PATCH] fix(mqtt): stop logging retained boot_report/alerts as new events The firmware publishes boot_report on system/info and alerts on system/alerts with retain=true. Every time the backend (re)connects - which under uvicorn --reload means every backend file save - the broker redelivers the last retained message and we inserted it again with occurred_at=now(). Result: the Health tab showed fresh PANIC boots and "Device reset due to fault" alerts for a device that had been up for 4 days. - insert_boot_event skips the insert when a row with the same (device_serial, boot_count) already exists. boot_count is the firmware's lifetime counter, so it uniquely identifies a boot. - upsert_alert only writes when state/message actually changed and returns whether it did; the alert-event history row is only added on a change. A redelivered identical alert no longer bumps updated_at either. Existing duplicate rows are not touched by this commit. Co-Authored-By: Claude Opus 5.5 --- backend/database/pg_mqtt.py | 33 +++++++++++++++++++++++++++++---- backend/mqtt/logger.py | 7 +++++-- 2 files changed, 34 insertions(+), 6 deletions(-) diff --git a/backend/database/pg_mqtt.py b/backend/database/pg_mqtt.py index c95b1ae..0a9a51e 100644 --- a/backend/database/pg_mqtt.py +++ b/backend/database/pg_mqtt.py @@ -300,9 +300,13 @@ async def get_pending_command(device_serial: str) -> dict | None: # --------------------------------------------------------------------------- async def upsert_alert(device_serial: str, subsystem: str, state: str, - message: str | None = None): + message: str | None = None) -> bool: + """Insert or update the current alert. Returns False (and leaves the row, + including updated_at, untouched) when it already holds this exact state and + message — which is what a retained system/alerts message looks like when the + broker redelivers it on every backend (re)subscribe.""" async with AsyncSessionLocal() as session: - await session.execute( + result = await session.execute( text(""" INSERT INTO device_alerts (device_serial, subsystem, state, message, updated_at) VALUES (:serial, :subsystem, :state, :message, now()) @@ -311,10 +315,15 @@ async def upsert_alert(device_serial: str, subsystem: str, state: str, state = EXCLUDED.state, message = EXCLUDED.message, updated_at = EXCLUDED.updated_at + WHERE device_alerts.state IS DISTINCT FROM EXCLUDED.state + OR device_alerts.message IS DISTINCT FROM EXCLUDED.message + RETURNING id """), {"serial": device_serial, "subsystem": subsystem, "state": state, "message": message}, ) + changed = result.fetchone() is not None await session.commit() + return changed async def delete_alert(device_serial: str, subsystem: str): @@ -435,8 +444,24 @@ 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: + crash_exc_vaddr: int | None = None) -> int | None: + """Insert a boot event, unless one with the same boot_count already exists + for this device (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.""" async with AsyncSessionLocal() as session: + if boot_count is not None: + dup = await session.execute( + text(""" + SELECT 1 FROM device_boot_events + WHERE device_serial = :serial AND boot_count = :boot_count + LIMIT 1 + """), + {"serial": device_serial, "boot_count": boot_count}, + ) + if dup.first() is not None: + return None result = await session.execute( text(""" INSERT INTO device_boot_events @@ -461,7 +486,7 @@ async def insert_boot_event(device_serial: str, boot_count: int | None, ) row = result.fetchone() await session.commit() - return row[0] + return row[0] if row else None async def get_boot_events(device_serial: str, limit: int = 100, offset: int = 0, diff --git a/backend/mqtt/logger.py b/backend/mqtt/logger.py index deddcea..f61c08b 100644 --- a/backend/mqtt/logger.py +++ b/backend/mqtt/logger.py @@ -86,10 +86,13 @@ async def _handle_alerts(serial: str, payload: dict): if state == "CLEARED": await db.delete_alert(serial, subsystem) else: - await db.upsert_alert(serial, subsystem, state, payload.get("msg")) + 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. - await db.insert_alert_event(serial, subsystem, state, payload.get("msg")) + # Skipped when nothing changed: alerts are retained, so the broker + # re-sends every active one on each backend (re)subscribe. + if changed: + await db.insert_alert_event(serial, subsystem, state, payload.get("msg")) async def _handle_info(serial: str, payload: dict):