diff --git a/backend/alembic/versions/c9d0e1f2a3b4_boot_event_crash_detail.py b/backend/alembic/versions/c9d0e1f2a3b4_boot_event_crash_detail.py new file mode 100644 index 0000000..f4876d1 --- /dev/null +++ b/backend/alembic/versions/c9d0e1f2a3b4_boot_event_crash_detail.py @@ -0,0 +1,32 @@ +"""boot event crash detail (firmware F-070) + +Firmware F-070 enriches boot_report (and telemetry.get_boot_history entries) +with a fuller `crash` object (backtrace, backtrace_corrupted, elf_sha256) and a +new, independent `pre_crash` object (uptime/heap snapshot + optional +abort_msg). Both are stored verbatim as JSONB so future additive fields don't +need another migration. The existing flat crash_* columns stay and are still +filled, so older rows and older firmware keep working unchanged. + +Revision ID: c9d0e1f2a3b4 +Revises: b8c9d0e1f2a3 +Create Date: 2026-09-30 00:00:00.000000 +""" +from typing import Sequence, Union +import sqlalchemy as sa +from sqlalchemy.dialects import postgresql +from alembic import op + +revision: str = "c9d0e1f2a3b4" +down_revision: Union[str, None] = "b8c9d0e1f2a3" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.add_column("device_boot_events", sa.Column("crash", postgresql.JSONB(), nullable=True)) + op.add_column("device_boot_events", sa.Column("pre_crash", postgresql.JSONB(), nullable=True)) + + +def downgrade() -> None: + op.drop_column("device_boot_events", "pre_crash") + op.drop_column("device_boot_events", "crash") diff --git a/backend/database/__init__.py b/backend/database/__init__.py index c609dab..bd6f48f 100644 --- a/backend/database/__init__.py +++ b/backend/database/__init__.py @@ -23,6 +23,8 @@ from database.pg_mqtt import ( insert_boot_event, get_boot_events, get_latest_boot_event, + merge_device_boot_history, + get_crash_events, insert_ping_sample, get_ping_samples, get_latest_ping_samples, @@ -74,6 +76,8 @@ __all__ = [ "insert_boot_event", "get_boot_events", "get_latest_boot_event", + "merge_device_boot_history", + "get_crash_events", "insert_ping_sample", "get_ping_samples", "get_latest_ping_samples", diff --git a/backend/database/pg_mqtt.py b/backend/database/pg_mqtt.py index 132d199..c5f2bdd 100644 --- a/backend/database/pg_mqtt.py +++ b/backend/database/pg_mqtt.py @@ -438,13 +438,56 @@ async def get_latest_alert_event(device_serial: str) -> dict | None: # Health tab's boot timeline + restart chart. # --------------------------------------------------------------------------- +_BOOT_COLUMNS = """ + id, device_serial, boot_count, reset_reason, is_fault, free_heap, + crash_task, crash_pc, crash_exc_cause, crash_exc_vaddr, + crash::text AS crash, pre_crash::text AS pre_crash, + occurred_at AT TIME ZONE 'UTC' AS occurred_at +""" + +# How far apart a device-reported boot timestamp (telemetry.get_boot_history +# "ts", device clock) and our own occurred_at (broker receive time of the live +# boot_report) may be and still count as the same boot. The live report goes +# out on first MQTT connect, so it normally lands within a minute of boot. +_BOOT_MATCH_WINDOW = timedelta(minutes=30) + + +def _json_or_none(obj: dict | None) -> str | None: + """Encode for a CAST(CAST(:x AS TEXT) AS JSONB) bind. Passed as text on + purpose: a bare CAST(:x AS JSONB) makes asyncpg type the parameter as jsonb + and JSON-encode our already-encoded string a second time.""" + return json.dumps(obj) if obj else None + + +def _boot_row_to_dict(row) -> dict: + d = _row_to_dict(row) + for key in ("crash", "pre_crash"): + val = d.get(key) + if isinstance(val, str): + try: + d[key] = json.loads(val) + except ValueError: + d[key] = None + return d + + +def _crash_columns(crash: dict | None) -> dict: + """The legacy flat crash_* columns, still filled so older readers and the + fleet grouping query keep working without unpacking JSON.""" + crash = crash or {} + return { + "crash_task": crash.get("task"), + "crash_pc": crash.get("pc"), + "crash_exc_cause": crash.get("exc_cause"), + "crash_exc_vaddr": crash.get("exc_vaddr"), + } + + async def insert_boot_event(device_serial: str, boot_count: int | None, reset_reason: str | None, is_fault: bool, free_heap: int | None, - crash_task: str | None = None, - crash_pc: int | None = None, - crash_exc_cause: int | None = None, - crash_exc_vaddr: int | None = None, + crash: dict | None = None, + pre_crash: dict | 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 @@ -452,7 +495,10 @@ async def insert_boot_event(device_serial: str, boot_count: int | None, 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.""" + same boot_count legitimately appears again for a later, different boot. + + `crash` / `pre_crash` are the firmware's objects stored verbatim (F-070); + both are omitted by older firmware and simply stay NULL.""" async with AsyncSessionLocal() as session: if skip_if_latest and boot_count is not None: latest = await session.execute( @@ -471,10 +517,12 @@ async def insert_boot_event(device_serial: str, boot_count: int | None, text(""" INSERT INTO device_boot_events (device_serial, boot_count, reset_reason, is_fault, free_heap, - crash_task, crash_pc, crash_exc_cause, crash_exc_vaddr, occurred_at) + crash_task, crash_pc, crash_exc_cause, crash_exc_vaddr, + crash, pre_crash, occurred_at) VALUES (:serial, :boot_count, :reset_reason, :is_fault, :free_heap, - :crash_task, :crash_pc, :crash_exc_cause, :crash_exc_vaddr, now()) + :crash_task, :crash_pc, :crash_exc_cause, :crash_exc_vaddr, + CAST(CAST(:crash AS TEXT) AS JSONB), CAST(CAST(:pre_crash AS TEXT) AS JSONB), now()) RETURNING id """), { @@ -483,10 +531,9 @@ async def insert_boot_event(device_serial: str, boot_count: int | None, "reset_reason": reset_reason, "is_fault": is_fault, "free_heap": free_heap, - "crash_task": crash_task, - "crash_pc": crash_pc, - "crash_exc_cause": crash_exc_cause, - "crash_exc_vaddr": crash_exc_vaddr, + **_crash_columns(crash), + "crash": _json_or_none(crash), + "pre_crash": _json_or_none(pre_crash), }, ) row = result.fetchone() @@ -494,6 +541,102 @@ async def insert_boot_event(device_serial: str, boot_count: int | None, return row[0] if row else None +async def merge_device_boot_history(device_serial: str, boots: list[dict]) -> dict: + """Merge the device's own SD boot log (telemetry.get_boot_history reply) + into device_boot_events. + + Each entry is {"ts": epoch s, "boot": n, "reason": str, "crash"?, "pre_crash"?}. + An entry matches an existing row when boot_count and reset_reason are equal + and the times are within _BOOT_MATCH_WINDOW (boot_count alone is not unique, + see insert_boot_event). A match only fills in crash/pre_crash the row is + missing; nothing already stored is overwritten. An entry with no match is a + boot we never saw live and is inserted at the device's timestamp. Entries + without a usable ts can't be placed in time and are skipped. + """ + inserted = updated = skipped = 0 + async with AsyncSessionLocal() as session: + for entry in boots: + if not isinstance(entry, dict): + skipped += 1 + continue + ts = entry.get("ts") + boot_count = entry.get("boot") + reason = entry.get("reason") + if not isinstance(ts, (int, float)) or ts <= 0: + skipped += 1 + continue + at = datetime.fromtimestamp(ts, tz=timezone.utc) + crash = entry.get("crash") if isinstance(entry.get("crash"), dict) else None + pre_crash = entry.get("pre_crash") if isinstance(entry.get("pre_crash"), dict) else None + + match = (await session.execute( + text(""" + SELECT id, crash IS NULL AS no_crash, pre_crash IS NULL AS no_pre_crash + FROM device_boot_events + WHERE device_serial = :serial + AND boot_count IS NOT DISTINCT FROM :boot_count + AND reset_reason IS NOT DISTINCT FROM :reason + AND occurred_at BETWEEN :lo AND :hi + ORDER BY abs(extract(epoch FROM occurred_at - :at)) + LIMIT 1 + """), + {"serial": device_serial, "boot_count": boot_count, "reason": reason, + "lo": at - _BOOT_MATCH_WINDOW, "hi": at + _BOOT_MATCH_WINDOW, "at": at}, + )).first() + + if match is None: + await session.execute( + text(""" + INSERT INTO device_boot_events + (device_serial, boot_count, reset_reason, is_fault, free_heap, + crash_task, crash_pc, crash_exc_cause, crash_exc_vaddr, + crash, pre_crash, occurred_at) + VALUES + (:serial, :boot_count, :reason, :is_fault, NULL, + :crash_task, :crash_pc, :crash_exc_cause, :crash_exc_vaddr, + CAST(CAST(:crash AS TEXT) AS JSONB), CAST(CAST(:pre_crash AS TEXT) AS JSONB), :at) + """), + { + "serial": device_serial, "boot_count": boot_count, "reason": reason, + "is_fault": bool(entry.get("is_fault", reason in FAULT_RESET_REASONS)), + **_crash_columns(crash), + "crash": _json_or_none(crash), "pre_crash": _json_or_none(pre_crash), + "at": at, + }, + ) + inserted += 1 + continue + + fill_crash = crash is not None and match.no_crash + fill_pre = pre_crash is not None and match.no_pre_crash + if not (fill_crash or fill_pre): + continue + sets, params = [], {"id": match.id} + if fill_crash: + sets.append("crash = CAST(CAST(:crash AS TEXT) AS JSONB)") + sets.append("crash_task = COALESCE(crash_task, :crash_task)") + sets.append("crash_pc = COALESCE(crash_pc, :crash_pc)") + sets.append("crash_exc_cause = COALESCE(crash_exc_cause, :crash_exc_cause)") + sets.append("crash_exc_vaddr = COALESCE(crash_exc_vaddr, :crash_exc_vaddr)") + params.update(_crash_columns(crash)) + params["crash"] = _json_or_none(crash) + if fill_pre: + sets.append("pre_crash = CAST(CAST(:pre_crash AS TEXT) AS JSONB)") + params["pre_crash"] = _json_or_none(pre_crash) + await session.execute( + text(f"UPDATE device_boot_events SET {', '.join(sets)} WHERE id = :id"), + params, + ) + updated += 1 + await session.commit() + return {"inserted": inserted, "updated": updated, "skipped": skipped} + + +# Mirrors the firmware's isFaultReset() — used only when a history entry +# doesn't say is_fault itself. +FAULT_RESET_REASONS = {"PANIC", "TASK_WATCHDOG", "INTERRUPT_WATCHDOG", "OTHER_WATCHDOG", "BROWNOUT"} + + async def get_boot_events(device_serial: str, limit: int = 100, offset: int = 0, since: datetime | None = None, until: datetime | None = None) -> tuple[list, int]: @@ -509,9 +652,7 @@ async def get_boot_events(device_serial: str, limit: int = 100, offset: int = 0, rows_result = await session.execute( text(f""" - SELECT id, device_serial, boot_count, reset_reason, is_fault, free_heap, - crash_task, crash_pc, crash_exc_cause, crash_exc_vaddr, - occurred_at AT TIME ZONE 'UTC' AS occurred_at + SELECT {_BOOT_COLUMNS} FROM device_boot_events WHERE device_serial = :serial{range_clause} ORDER BY occurred_at DESC @@ -521,16 +662,14 @@ async def get_boot_events(device_serial: str, limit: int = 100, offset: int = 0, ) rows = rows_result.mappings().all() - return [_row_to_dict(r) for r in rows], total + return [_boot_row_to_dict(r) for r in rows], total async def get_latest_boot_event(device_serial: str) -> dict | None: async with AsyncSessionLocal() as session: result = await session.execute( - text(""" - SELECT id, device_serial, boot_count, reset_reason, is_fault, free_heap, - crash_task, crash_pc, crash_exc_cause, crash_exc_vaddr, - occurred_at AT TIME ZONE 'UTC' AS occurred_at + text(f""" + SELECT {_BOOT_COLUMNS} FROM device_boot_events WHERE device_serial = :serial ORDER BY occurred_at DESC @@ -540,7 +679,32 @@ async def get_latest_boot_event(device_serial: str) -> dict | None: ) row = result.mappings().first() - return _row_to_dict(row) if row else None + return _boot_row_to_dict(row) if row else None + + +async def get_crash_events(since: datetime | None = None, + until: datetime | None = None, + limit: int = 5000) -> list: + """Fleet-wide fault boots that carry any crash detail (coredump summary + and/or pre-crash snapshot), newest first. Grouped in mqtt/crash_groups.py.""" + async with AsyncSessionLocal() as session: + params: dict = {"limit": limit} + range_clause = _time_range_clause("occurred_at", since, until, params) + rows_result = await session.execute( + text(f""" + SELECT {_BOOT_COLUMNS} + FROM device_boot_events + WHERE is_fault + AND (crash IS NOT NULL OR pre_crash IS NOT NULL OR crash_task IS NOT NULL) + {range_clause} + ORDER BY occurred_at DESC + LIMIT :limit + """), + params, + ) + rows = rows_result.mappings().all() + + return [_boot_row_to_dict(r) for r in rows] # --------------------------------------------------------------------------- diff --git a/backend/mqtt/crash_groups.py b/backend/mqtt/crash_groups.py new file mode 100644 index 0000000..7fc9aba --- /dev/null +++ b/backend/mqtt/crash_groups.py @@ -0,0 +1,90 @@ +"""Fleet-wide crash grouping (firmware F-070). + +Groups fault boots so recurring crash types are visible across devices: + - with an abort_msg: by the message with hex addresses stripped, since the + same assert/abort reports different addresses per build; + - otherwise: by faulting task + Xtensa exception cause; + - fault resets with no coredump and no abort_msg: by reset reason alone. + +When abort_msg is present, pc/exc_cause/exc_vaddr describe the abort() +mechanism (always StoreProhibited at 0x0), so they must NOT drive the grouping. +""" + +import re + +# Xtensa EXCCAUSE values the ESP32 actually produces. +EXC_CAUSE_NAMES = { + 0: "IllegalInstruction", + 2: "InstructionFetchError", + 3: "LoadStoreError", + 6: "IntegerDivideByZero", + 9: "LoadStoreAlignment", + 28: "LoadProhibited", + 29: "StoreProhibited", +} + +_HEX = re.compile(r"0x[0-9a-fA-F]+") +_SPACES = re.compile(r"\s+") + + +def normalize_abort_msg(msg: str) -> str: + return _SPACES.sub(" ", _HEX.sub("0x…", msg)).strip() + + +def exc_cause_name(cause) -> str: + if cause is None: + return "Unknown exception" + return EXC_CAUSE_NAMES.get(cause, f"Exception cause {cause}") + + +def group_key(event: dict) -> tuple[str, str, str]: + """(key, kind, title) for one boot event.""" + crash = event.get("crash") or {} + pre = event.get("pre_crash") or {} + abort_msg = pre.get("abort_msg") + if abort_msg: + norm = normalize_abort_msg(abort_msg) + return f"abort:{norm}", "abort", norm + + task = crash.get("task", event.get("crash_task")) + cause = crash.get("exc_cause", event.get("crash_exc_cause")) + if task is not None or cause is not None: + title = f"{exc_cause_name(cause)} in {task or 'unknown task'}" + return f"exc:{task}:{cause}", "exception", title + + reason = event.get("reset_reason") or "UNKNOWN" + return f"reason:{reason}", "reason", f"{reason.replace('_', ' ')} (no coredump)" + + +def group_crashes(events: list[dict]) -> list[dict]: + """events must be newest-first (as get_crash_events returns them).""" + groups: dict[str, dict] = {} + for ev in events: + key, kind, title = group_key(ev) + g = groups.get(key) + if g is None: + g = groups[key] = { + "key": key, + "kind": kind, + "title": title, + "count": 0, + "first_at": ev["occurred_at"], + "last_at": ev["occurred_at"], + "latest": ev, + "devices": {}, + } + g["count"] += 1 + g["first_at"] = ev["occurred_at"] # newest-first input, so the last seen is the oldest + d = g["devices"].setdefault(ev["device_serial"], { + "device_serial": ev["device_serial"], "count": 0, "last_at": ev["occurred_at"], + }) + d["count"] += 1 + + result = [] + for g in groups.values(): + devices = sorted(g["devices"].values(), key=lambda d: (-d["count"], d["device_serial"])) + result.append({**g, "devices": devices, "device_count": len(devices)}) + # Most frequent first; ties broken by most recent (two stable sorts). + result.sort(key=lambda g: g["last_at"], reverse=True) + result.sort(key=lambda g: g["count"], reverse=True) + return result diff --git a/backend/mqtt/logger.py b/backend/mqtt/logger.py index bef1adc..c4e11c3 100644 --- a/backend/mqtt/logger.py +++ b/backend/mqtt/logger.py @@ -163,17 +163,19 @@ async def _handle_boot_report(serial: str, payload: dict, retained: bool = False info event as { "type": ..., "payload": {...} }. """ data = payload.get("payload", {}) - crash = data.get("crash") or {} + # F-070: `crash` (coredump summary, now incl. backtrace/elf_sha256) and + # `pre_crash` (heap/uptime snapshot + optional abort_msg) are stored + # verbatim. Older firmware omits them / sends the 4-field crash only. + crash = data.get("crash") if isinstance(data.get("crash"), dict) else None + pre_crash = data.get("pre_crash") if isinstance(data.get("pre_crash"), dict) else None await db.insert_boot_event( device_serial=serial, boot_count=data.get("boot_count"), reset_reason=data.get("reset_reason"), is_fault=bool(data.get("is_fault", False)), free_heap=data.get("free_heap"), - crash_task=crash.get("task"), - crash_pc=crash.get("pc"), - crash_exc_cause=crash.get("exc_cause"), - crash_exc_vaddr=crash.get("exc_vaddr"), + crash=crash, + pre_crash=pre_crash, # 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, @@ -252,6 +254,15 @@ async def _handle_ack(serial: str, payload: dict): await db.insert_ping_sample(device_serial=serial, rtt_ms=rtt_ms) return + # The device's own SD boot log — merged into device_boot_events whoever + # asked for it (Health tab "Sync from device", API reference, Control tab), + # so crash detail for boots we never saw live isn't lost. + if payload.get("type") == "telemetry.get_boot_history" and status == "SUCCESS": + boots = (payload.get("data") or {}).get("boots") + if isinstance(boots, list): + result = await db.merge_device_boot_history(serial, boots) + logger.info(f"Merged device boot history for {serial}: {result}") + pending = await db.get_pending_command(serial) if pending: cmd_status = "success" if status == "SUCCESS" else "error" diff --git a/backend/mqtt/models.py b/backend/mqtt/models.py index 6922aa2..8653b5a 100644 --- a/backend/mqtt/models.py +++ b/backend/mqtt/models.py @@ -129,6 +129,9 @@ class BootEventEntry(BaseModel): crash_pc: Optional[int] = None crash_exc_cause: Optional[int] = None crash_exc_vaddr: Optional[int] = None + # Firmware F-070 objects, verbatim — None for older firmware/rows. + crash: Optional[Dict[str, Any]] = None + pre_crash: Optional[Dict[str, Any]] = None occurred_at: str @@ -137,6 +140,29 @@ class BootEventListResponse(BaseModel): total: int +class CrashGroupDevice(BaseModel): + device_serial: str + count: int + last_at: str + + +class CrashGroup(BaseModel): + key: str + kind: str # "abort" | "exception" | "reason" + title: str + count: int + device_count: int + first_at: str + last_at: str + latest: BootEventEntry + devices: List[CrashGroupDevice] + + +class CrashGroupListResponse(BaseModel): + groups: List[CrashGroup] + total_crashes: int + + class PingSampleEntry(BaseModel): id: int device_serial: str diff --git a/backend/mqtt/router.py b/backend/mqtt/router.py index 4e1ebd2..c0c020b 100644 --- a/backend/mqtt/router.py +++ b/backend/mqtt/router.py @@ -9,7 +9,9 @@ from mqtt.models import ( AlertEventListResponse, BootEventListResponse, PingSampleListResponse, DiagnosticsReportListResponse, LatestMetricsResponse, LatestDiagnosticsEntry, LatestPingEntry, DeviceReportListResponse, + CrashGroupListResponse, ) +from mqtt.crash_groups import group_crashes from mqtt.client import mqtt_manager from mqtt import presence import database as db @@ -177,6 +179,17 @@ async def get_device_boot_events( return BootEventListResponse(events=events, total=total) +@router.get("/crash-groups", response_model=CrashGroupListResponse) +async def get_crash_groups( + since: Optional[datetime] = Query(None, description="ISO timestamp — only crashes at/after this time"), + until: Optional[datetime] = Query(None, description="ISO timestamp — only crashes at/before this time"), + _user: TokenPayload = Depends(require_permission("mqtt", "view")), +): + """Fault boots across the whole fleet, grouped by crash signature.""" + events = await db.get_crash_events(since=since, until=until) + return CrashGroupListResponse(groups=group_crashes(events), total_crashes=len(events)) + + @router.get("/ping-samples/{device_serial}", response_model=PingSampleListResponse) async def get_device_ping_samples( device_serial: str,