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 <noreply@anthropic.com>
This commit is contained in:
2026-09-30 15:48:48 +03:00
co-authored by Claude Opus 5.5
parent 36f9540a89
commit 98dd16b597
2 changed files with 34 additions and 6 deletions
+29 -4
View File
@@ -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,
+5 -2
View File
@@ -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):