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:
@@ -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,
|
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:
|
async with AsyncSessionLocal() as session:
|
||||||
await session.execute(
|
result = await session.execute(
|
||||||
text("""
|
text("""
|
||||||
INSERT INTO device_alerts (device_serial, subsystem, state, message, updated_at)
|
INSERT INTO device_alerts (device_serial, subsystem, state, message, updated_at)
|
||||||
VALUES (:serial, :subsystem, :state, :message, now())
|
VALUES (:serial, :subsystem, :state, :message, now())
|
||||||
@@ -311,10 +315,15 @@ async def upsert_alert(device_serial: str, subsystem: str, state: str,
|
|||||||
state = EXCLUDED.state,
|
state = EXCLUDED.state,
|
||||||
message = EXCLUDED.message,
|
message = EXCLUDED.message,
|
||||||
updated_at = EXCLUDED.updated_at
|
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},
|
{"serial": device_serial, "subsystem": subsystem, "state": state, "message": message},
|
||||||
)
|
)
|
||||||
|
changed = result.fetchone() is not None
|
||||||
await session.commit()
|
await session.commit()
|
||||||
|
return changed
|
||||||
|
|
||||||
|
|
||||||
async def delete_alert(device_serial: str, subsystem: str):
|
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_task: str | None = None,
|
||||||
crash_pc: int | None = None,
|
crash_pc: int | None = None,
|
||||||
crash_exc_cause: 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:
|
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(
|
result = await session.execute(
|
||||||
text("""
|
text("""
|
||||||
INSERT INTO device_boot_events
|
INSERT INTO device_boot_events
|
||||||
@@ -461,7 +486,7 @@ async def insert_boot_event(device_serial: str, boot_count: int | None,
|
|||||||
)
|
)
|
||||||
row = result.fetchone()
|
row = result.fetchone()
|
||||||
await session.commit()
|
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,
|
async def get_boot_events(device_serial: str, limit: int = 100, offset: int = 0,
|
||||||
|
|||||||
@@ -86,9 +86,12 @@ async def _handle_alerts(serial: str, payload: dict):
|
|||||||
if state == "CLEARED":
|
if state == "CLEARED":
|
||||||
await db.delete_alert(serial, subsystem)
|
await db.delete_alert(serial, subsystem)
|
||||||
else:
|
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
|
# Append-only history — survives past the alert being resolved, used to
|
||||||
# answer "when was the most recent issue" even once it's cleared.
|
# 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:
|
||||||
await db.insert_alert_event(serial, subsystem, state, payload.get("msg"))
|
await db.insert_alert_event(serial, subsystem, state, payload.get("msg"))
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user