diff --git a/backend/alembic/versions/a7b8c9d0e1f2_device_diagnostics_reports.py b/backend/alembic/versions/a7b8c9d0e1f2_device_diagnostics_reports.py new file mode 100644 index 0000000..da86419 --- /dev/null +++ b/backend/alembic/versions/a7b8c9d0e1f2_device_diagnostics_reports.py @@ -0,0 +1,73 @@ +"""device diagnostics reports + +Adds device_diagnostics_reports — structured storage for the firmware's +diagnostics_report MQTT event (vesper/{uid}/status/info, type="diagnostics_report"), +published every 5 minutes. Cleanly separate from device_boot_events (one row per +boot, event-driven) and heartbeats (one row every 30s, transport/liveness facts): +this table is periodic health telemetry — CPU temperature (min/max/avg over the +5-minute window), WiFi reconnect count + last disconnect reason, OTA/firmware +check state, and per-task stack high-water marks (bytes free). + +stack_high_water is stored as a JSON-encoded TEXT column rather than flattened +columns — unlike the other three groups (fixed field sets), the set of +monitored tasks is open-ended on the firmware side (see project-vesper's +Telemetry::registerTaskForStackMonitoring), so a fixed column per task would +need a migration every time a task is added. TEXT (not JSONB) matches this +codebase's existing convention for JSON blobs stored via raw SQL — see +commands.command_payload / response_payload. + +Revision ID: a7b8c9d0e1f2 +Revises: f6a7b8c9d0e1 +Create Date: 2026-07-17 00:00:00.000000 +""" +from typing import Sequence, Union +import sqlalchemy as sa +from alembic import op + +revision: str = "a7b8c9d0e1f2" +down_revision: Union[str, None] = "f6a7b8c9d0e1" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.create_table( + "device_diagnostics_reports", + sa.Column("id", sa.BigInteger(), primary_key=True, autoincrement=True), + sa.Column("device_serial", sa.String(128), nullable=False), + + # cpu_temp — omitted by firmware (all null here) if no samples were + # taken yet this window (e.g. very first report shortly after boot). + sa.Column("cpu_temp_avg", sa.Float(), nullable=True), + sa.Column("cpu_temp_min", sa.Float(), nullable=True), + sa.Column("cpu_temp_max", sa.Float(), nullable=True), + sa.Column("cpu_temp_samples", sa.Integer(), nullable=True), + + # wifi_reconnects — this-boot-only lifetime count, not device lifetime. + sa.Column("wifi_reconnect_count", sa.Integer(), nullable=True), + sa.Column("wifi_last_disconnect_reason", sa.String(64), nullable=True), + sa.Column("wifi_last_disconnect_uptime_ms", sa.BigInteger(), nullable=True), + + # ota + sa.Column("ota_current_version", sa.String(32), nullable=True), + sa.Column("ota_update_available", sa.Boolean(), nullable=True), + sa.Column("ota_available_version", sa.String(32), nullable=True), + sa.Column("ota_last_check_uptime_ms", sa.BigInteger(), nullable=True), + sa.Column("ota_last_error", sa.String(32), nullable=True), + + # stack_high_water — open-ended task set, see module docstring. + sa.Column("stack_high_water", sa.Text(), nullable=True), + + sa.Column("received_at", sa.DateTime(timezone=True), nullable=False, + server_default=sa.func.now()), + ) + op.create_index( + "idx_device_diagnostics_reports_serial_received", + "device_diagnostics_reports", + ["device_serial", sa.text("received_at DESC")], + ) + + +def downgrade() -> None: + op.drop_index("idx_device_diagnostics_reports_serial_received", table_name="device_diagnostics_reports") + op.drop_table("device_diagnostics_reports") diff --git a/backend/alembic/versions/b8c9d0e1f2a3_mqtt_v2_topic_migration.py b/backend/alembic/versions/b8c9d0e1f2a3_mqtt_v2_topic_migration.py new file mode 100644 index 0000000..0e42c42 --- /dev/null +++ b/backend/alembic/versions/b8c9d0e1f2a3_mqtt_v2_topic_migration.py @@ -0,0 +1,78 @@ +"""mqtt v2 topic migration + +Firmware moved to a new MQTT topic spec (project-vesper's +docs/reference/vesper_mqtt_topic_spec_v2.md, feature-catalog F-062). Three +schema changes needed to keep ingestion correct against the new payload shapes: + +1. heartbeats.state / heartbeats.ok — the v2 heartbeat payload is flat (no more + {"status":"INFO","type":"heartbeat","payload":{...}} wrapper) and adds two + new fields: state ("idle"/"playing"/"paused"/"error"/"booting") and ok + (overall health), both independently derived on the firmware side — not + mirrored from status/playback or system/alerts. Nullable: older rows + (pre-migration) and any device still reporting have no value for these. + +2. device_diagnostics_reports gains bell_strikes / bell_loads / cooling_active. + The 5-min diagnostics report moved from status/info to system/metrics and + picked up per-bell strike/heat data along the way (the topic spec calls out + "bell heat ratings" explicitly for this topic). Stored as JSON-encoded TEXT, + matching this table's existing stack_high_water column — same reasoning: + 16 fixed bell slots, but keeping the encoding consistent with the sibling + column is simpler than mixing raw JSONB and TEXT in one row shape. + +3. device_reports — new table for the new control/reports topic. Board- + initiated, unsolicited, critical events (currently only bell_overload). + Distinct from device_alert_events (subsystem health transitions) and + device_logs (routine log lines) — this is a narrow, insert-only table for + a different kind of signal: urgent, bell-mechanism-specific events the + physical tablets act on in real time. The console side is deliberately + light for now — store + list for historical/audit purposes; the tablets + are the real-time consumer of control/reports, not this console. + +Revision ID: b8c9d0e1f2a3 +Revises: a7b8c9d0e1f2 +Create Date: 2026-09-21 00:00:00.000000 +""" +from typing import Sequence, Union +import sqlalchemy as sa +from alembic import op + +revision: str = "b8c9d0e1f2a3" +down_revision: Union[str, None] = "a7b8c9d0e1f2" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.add_column("heartbeats", sa.Column("state", sa.String(16), nullable=True)) + op.add_column("heartbeats", sa.Column("ok", sa.Boolean(), nullable=True)) + + op.add_column("device_diagnostics_reports", sa.Column("bell_strikes", sa.Text(), nullable=True)) + op.add_column("device_diagnostics_reports", sa.Column("bell_loads", sa.Text(), nullable=True)) + op.add_column("device_diagnostics_reports", sa.Column("cooling_active", sa.Boolean(), nullable=True)) + + op.create_table( + "device_reports", + sa.Column("id", sa.BigInteger(), primary_key=True, autoincrement=True), + sa.Column("device_serial", sa.String(128), nullable=False), + sa.Column("report_type", sa.String(64), nullable=False), + sa.Column("payload", sa.Text(), nullable=True), + sa.Column("occurred_at", sa.DateTime(timezone=True), nullable=False, + server_default=sa.func.now()), + ) + op.create_index( + "idx_device_reports_serial_occurred", + "device_reports", + ["device_serial", sa.text("occurred_at DESC")], + ) + + +def downgrade() -> None: + op.drop_index("idx_device_reports_serial_occurred", table_name="device_reports") + op.drop_table("device_reports") + + op.drop_column("device_diagnostics_reports", "cooling_active") + op.drop_column("device_diagnostics_reports", "bell_loads") + op.drop_column("device_diagnostics_reports", "bell_strikes") + + op.drop_column("heartbeats", "ok") + op.drop_column("heartbeats", "state") diff --git a/backend/alembic/versions/d4e5f6a7b8c9_device_alert_events.py b/backend/alembic/versions/d4e5f6a7b8c9_device_alert_events.py new file mode 100644 index 0000000..5e00072 --- /dev/null +++ b/backend/alembic/versions/d4e5f6a7b8c9_device_alert_events.py @@ -0,0 +1,48 @@ +"""device_alert_events + +Adds an insert-only history log of device alert transitions (WARNING/CRITICAL/ +FAILED only — CLEARED is not logged here). device_alerts remains the current- +state table and is untouched; this is purely additive, used to answer "when +was the most recent issue on this device" even after it has been resolved. + +Also adds an optional rssi column to heartbeats, for a signal-strength +indicator once firmware starts reporting it on the heartbeat payload. + +Revision ID: d4e5f6a7b8c9 +Revises: c3d4e5f6a7b8 +Create Date: 2026-07-13 00:00:00.000000 +""" +from typing import Sequence, Union +import sqlalchemy as sa +from alembic import op + +revision: str = "d4e5f6a7b8c9" +down_revision: Union[str, None] = "c3d4e5f6a7b8" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.create_table( + "device_alert_events", + sa.Column("id", sa.BigInteger(), primary_key=True, autoincrement=True), + sa.Column("device_serial", sa.String(128), nullable=False), + sa.Column("subsystem", sa.String(128), nullable=False), + sa.Column("state", sa.String(64), nullable=False), + sa.Column("message", sa.Text(), nullable=True), + sa.Column("occurred_at", sa.DateTime(timezone=True), nullable=False, + server_default=sa.func.now()), + ) + op.create_index( + "idx_device_alert_events_serial_occurred", + "device_alert_events", + ["device_serial", sa.text("occurred_at DESC")], + ) + + op.add_column("heartbeats", sa.Column("rssi", sa.Integer(), nullable=True)) + + +def downgrade() -> None: + op.drop_column("heartbeats", "rssi") + op.drop_index("idx_device_alert_events_serial_occurred", table_name="device_alert_events") + op.drop_table("device_alert_events") diff --git a/backend/alembic/versions/e5f6a7b8c9d0_device_logs_source.py b/backend/alembic/versions/e5f6a7b8c9d0_device_logs_source.py new file mode 100644 index 0000000..e26ca09 --- /dev/null +++ b/backend/alembic/versions/e5f6a7b8c9d0_device_logs_source.py @@ -0,0 +1,39 @@ +"""device_logs source column + +Adds a 'source' column to device_logs distinguishing the debug log stream +(vesper/{id}/logs, source='log') from the general info stream +(vesper/{id}/status/info, source='info'), which was previously discarded +after only being written to the server debug log. Existing rows default to +'log' since that's the only source ever persisted before this migration. + +Adding a column to a partitioned parent table applies it to all existing +and future partitions automatically. + +Revision ID: e5f6a7b8c9d0 +Revises: d4e5f6a7b8c9 +Create Date: 2026-07-14 00:00:00.000000 +""" +from typing import Sequence, Union +import sqlalchemy as sa +from alembic import op + +revision: str = "e5f6a7b8c9d0" +down_revision: Union[str, None] = "d4e5f6a7b8c9" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.execute(""" + ALTER TABLE device_logs + ADD COLUMN source TEXT NOT NULL DEFAULT 'log' + """) + op.execute(""" + CREATE INDEX idx_device_logs_source + ON device_logs(device_serial, source, received_at DESC) + """) + + +def downgrade() -> None: + op.execute("DROP INDEX IF EXISTS idx_device_logs_source") + op.execute("ALTER TABLE device_logs DROP COLUMN source") diff --git a/backend/alembic/versions/f6a7b8c9d0e1_device_health_tab.py b/backend/alembic/versions/f6a7b8c9d0e1_device_health_tab.py new file mode 100644 index 0000000..1222d9a --- /dev/null +++ b/backend/alembic/versions/f6a7b8c9d0e1_device_health_tab.py @@ -0,0 +1,82 @@ +"""device health tab — boot events, heartbeat free_heap, ping samples + +Adds three pieces of schema needed for the device Health tab: + +1. device_boot_events — structured, insert-only history of the firmware's + boot_report MQTT event (vesper/{uid}/status/info, type="boot_report"). + Previously this payload was flattened into a single device_logs text line + with all structured fields (boot_count, crash detail) discarded on arrival. + This table gives the console a real timeline to query and chart against, + instead of parsing log strings. + +2. heartbeats.free_heap — firmware now includes free_heap on every 30s + heartbeat (previously only available once per boot via boot_report), so + the console can chart heap trend instead of seeing one point per boot. + +3. device_ping_samples — backend-computed RTT samples. The firmware's ping + command now echoes back a caller-supplied timestamp; the backend pings + each online device on an interval and records (now - echoed_ts) here. + A dedicated table rather than reusing `commands` because `commands` isn't + shaped for time-series charting (mixed command types, no fast per-device + time-range query path) and pruning ping history independently of other + command history is desirable. + +Revision ID: f6a7b8c9d0e1 +Revises: e5f6a7b8c9d0 +Create Date: 2026-07-16 00:00:00.000000 +""" +from typing import Sequence, Union +import sqlalchemy as sa +from alembic import op + +revision: str = "f6a7b8c9d0e1" +down_revision: Union[str, None] = "e5f6a7b8c9d0" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.create_table( + "device_boot_events", + sa.Column("id", sa.BigInteger(), primary_key=True, autoincrement=True), + sa.Column("device_serial", sa.String(128), nullable=False), + sa.Column("boot_count", sa.Integer(), nullable=True), + sa.Column("reset_reason", sa.String(64), nullable=True), + sa.Column("is_fault", sa.Boolean(), nullable=False, server_default=sa.false()), + sa.Column("free_heap", sa.Integer(), nullable=True), + sa.Column("crash_task", sa.String(64), nullable=True), + sa.Column("crash_pc", sa.BigInteger(), nullable=True), + sa.Column("crash_exc_cause", sa.Integer(), nullable=True), + sa.Column("crash_exc_vaddr", sa.BigInteger(), nullable=True), + sa.Column("occurred_at", sa.DateTime(timezone=True), nullable=False, + server_default=sa.func.now()), + ) + op.create_index( + "idx_device_boot_events_serial_occurred", + "device_boot_events", + ["device_serial", sa.text("occurred_at DESC")], + ) + + op.add_column("heartbeats", sa.Column("free_heap", sa.Integer(), nullable=True)) + + op.create_table( + "device_ping_samples", + sa.Column("id", sa.BigInteger(), primary_key=True, autoincrement=True), + sa.Column("device_serial", sa.String(128), nullable=False), + sa.Column("rtt_ms", sa.Integer(), nullable=False), + sa.Column("sampled_at", sa.DateTime(timezone=True), nullable=False, + server_default=sa.func.now()), + ) + op.create_index( + "idx_device_ping_samples_serial_sampled", + "device_ping_samples", + ["device_serial", sa.text("sampled_at DESC")], + ) + + +def downgrade() -> None: + op.drop_index("idx_device_ping_samples_serial_sampled", table_name="device_ping_samples") + op.drop_table("device_ping_samples") + op.drop_column("heartbeats", "free_heap") + op.drop_index("idx_device_boot_events_serial_occurred", table_name="device_boot_events") + op.drop_table("device_boot_events") diff --git a/backend/database/__init__.py b/backend/database/__init__.py index 18451c5..c609dab 100644 --- a/backend/database/__init__.py +++ b/backend/database/__init__.py @@ -16,6 +16,31 @@ from database.pg_mqtt import ( upsert_alert, delete_alert, get_alerts, + insert_alert_event, + get_latest_alert_events, + get_latest_alert_event, + get_alert_events, + insert_boot_event, + get_boot_events, + get_latest_boot_event, + insert_ping_sample, + get_ping_samples, + get_latest_ping_samples, + insert_diagnostics_report, + get_diagnostics_reports, + get_latest_diagnostics_report, + get_latest_diagnostics_reports, + insert_report, + get_reports, + delete_device_logs, + delete_device_heartbeats, + delete_device_commands, + delete_device_alert_events, + delete_device_current_alerts, + delete_device_boot_events, + delete_device_ping_samples, + delete_device_diagnostics_reports, + delete_device_reports, partition_manager_loop, ensure_current_partitions, ) @@ -42,6 +67,31 @@ __all__ = [ "upsert_alert", "delete_alert", "get_alerts", + "insert_alert_event", + "get_latest_alert_events", + "get_latest_alert_event", + "get_alert_events", + "insert_boot_event", + "get_boot_events", + "get_latest_boot_event", + "insert_ping_sample", + "get_ping_samples", + "get_latest_ping_samples", + "insert_diagnostics_report", + "get_diagnostics_reports", + "get_latest_diagnostics_report", + "get_latest_diagnostics_reports", + "insert_report", + "get_reports", + "delete_device_logs", + "delete_device_heartbeats", + "delete_device_commands", + "delete_device_alert_events", + "delete_device_current_alerts", + "delete_device_boot_events", + "delete_device_ping_samples", + "delete_device_diagnostics_reports", + "delete_device_reports", "partition_manager_loop", "ensure_current_partitions", ] diff --git a/backend/database/pg_mqtt.py b/backend/database/pg_mqtt.py index 0e35f64..c95b1ae 100644 --- a/backend/database/pg_mqtt.py +++ b/backend/database/pg_mqtt.py @@ -26,15 +26,15 @@ logger = logging.getLogger("database.pg_mqtt") # --------------------------------------------------------------------------- async def insert_log(device_serial: str, level: str, message: str, - device_timestamp: int | None = None) -> int: + device_timestamp: int | None = None, source: str = "log") -> int: async with AsyncSessionLocal() as session: result = await session.execute( text(""" - INSERT INTO device_logs (device_serial, level, message, device_timestamp, received_at) - VALUES (:serial, :level, :message, :ts, now()) + INSERT INTO device_logs (device_serial, level, message, device_timestamp, received_at, source) + VALUES (:serial, :level, :message, :ts, now(), :source) RETURNING id """), - {"serial": device_serial, "level": level, "message": message, "ts": device_timestamp}, + {"serial": device_serial, "level": level, "message": message, "ts": device_timestamp, "source": source}, ) row = result.fetchone() await session.commit() @@ -43,15 +43,17 @@ async def insert_log(device_serial: str, level: str, message: str, async def insert_heartbeat(device_serial: str, device_id: str, firmware_version: str, ip_address: str, - gateway: str, uptime_ms: int, uptime_display: str) -> int: + gateway: str, uptime_ms: int, uptime_display: str, + rssi: int | None = None, free_heap: int | None = None, + state: str | None = None, ok: bool | None = None) -> int: async with AsyncSessionLocal() as session: result = await session.execute( text(""" INSERT INTO heartbeats (device_serial, device_id, firmware_version, ip_address, - gateway, uptime_ms, uptime_display, received_at) + gateway, uptime_ms, uptime_display, rssi, free_heap, state, ok, received_at) VALUES - (:serial, :device_id, :fw, :ip, :gw, :uptime_ms, :uptime_display, now()) + (:serial, :device_id, :fw, :ip, :gw, :uptime_ms, :uptime_display, :rssi, :free_heap, :state, :ok, now()) RETURNING id """), { @@ -62,6 +64,10 @@ async def insert_heartbeat(device_serial: str, device_id: str, "gw": gateway, "uptime_ms": uptime_ms, "uptime_display": uptime_display, + "rssi": rssi, + "free_heap": free_heap, + "state": state, + "ok": ok, }, ) row = result.fetchone() @@ -113,18 +119,48 @@ async def update_command_response(command_id: int, status: str, # Query operations # --------------------------------------------------------------------------- +LEVEL_ORDER = ["INFO", "WARN", "ERROR"] + + async def get_logs(device_serial: str, level: str | None = None, - search: str | None = None, - limit: int = 100, offset: int = 0) -> tuple[list, int]: + search: str | None = None, source: str | list[str] | None = None, + min_level: bool = False, + limit: int = 100, offset: int = 0, + since: datetime | None = None, until: datetime | None = None) -> tuple[list, int]: + """ + level: exact level match (legacy behaviour). + min_level: if True, `level` is treated as a floor — also returns higher-severity + levels (e.g. min_level=True, level='WARN' returns WARN and ERROR). + source: single source string, or a list of sources (e.g. ['log', 'info']) to + include several channels in one query — used by the unified feed. + since/until: same per-device time-range filter as the other history + endpoints (heartbeats, boot events, etc.) — lets the console's + TimeRangeSelect scope the Logs sub-tab too. + """ where = "device_serial = :serial" params: dict = {"serial": device_serial, "limit": limit, "offset": offset} if level: - where += " AND level = :level" - params["level"] = level + if min_level and level in LEVEL_ORDER: + allowed = LEVEL_ORDER[LEVEL_ORDER.index(level):] + where += " AND level = ANY(:levels)" + params["levels"] = allowed + else: + where += " AND level = :level" + params["level"] = level if search: where += " AND message ILIKE :search" params["search"] = f"%{search}%" + if source: + sources = [source] if isinstance(source, str) else list(source) + where += " AND source = ANY(:sources)" + params["sources"] = sources + if since is not None: + where += " AND received_at >= :since" + params["since"] = since + if until is not None: + where += " AND received_at <= :until" + params["until"] = until async with AsyncSessionLocal() as session: count_result = await session.execute( @@ -134,7 +170,7 @@ async def get_logs(device_serial: str, level: str | None = None, rows_result = await session.execute( text(f""" - SELECT id, device_serial, level, message, device_timestamp, + SELECT id, device_serial, level, message, device_timestamp, source, received_at AT TIME ZONE 'UTC' AS received_at FROM device_logs WHERE {where} @@ -148,26 +184,47 @@ async def get_logs(device_serial: str, level: str | None = None, return [_row_to_dict(r) for r in rows], total -async def get_heartbeats(device_serial: str, limit: int = 100, - offset: int = 0) -> tuple[list, int]: +def _time_range_clause(column: str, since: datetime | None, until: datetime | None, + params: dict) -> str: + """Appends `since`/`until` bind params to `params` in place and returns an + ` AND {column} >= :since AND {column} <= :until`-style SQL fragment (only + the bounds actually provided). Shared by every history/telemetry query + that the console's per-device time-range filter needs to scope — see + TimeRangeSelect on the frontend.""" + clause = "" + if since is not None: + clause += f" AND {column} >= :since" + params["since"] = since + if until is not None: + clause += f" AND {column} <= :until" + params["until"] = until + return clause + + +async def get_heartbeats(device_serial: str, limit: int = 100, offset: int = 0, + since: datetime | None = None, + until: datetime | None = None) -> tuple[list, int]: async with AsyncSessionLocal() as session: + params: dict = {"serial": device_serial} + range_clause = _time_range_clause("received_at", since, until, params) + count_result = await session.execute( - text("SELECT COUNT(*) FROM heartbeats WHERE device_serial = :serial"), - {"serial": device_serial}, + text(f"SELECT COUNT(*) FROM heartbeats WHERE device_serial = :serial{range_clause}"), + params, ) total = count_result.scalar() rows_result = await session.execute( - text(""" + text(f""" SELECT id, device_serial, device_id, firmware_version, ip_address, - gateway, uptime_ms, uptime_display, + gateway, uptime_ms, uptime_display, rssi, free_heap, state, ok, received_at AT TIME ZONE 'UTC' AS received_at FROM heartbeats - WHERE device_serial = :serial + WHERE device_serial = :serial{range_clause} ORDER BY received_at DESC LIMIT :limit OFFSET :offset """), - {"serial": device_serial, "limit": limit, "offset": offset}, + {**params, "limit": limit, "offset": offset}, ) rows = rows_result.mappings().all() @@ -207,7 +264,7 @@ async def get_latest_heartbeats() -> list: text(""" SELECT DISTINCT ON (device_serial) id, device_serial, device_id, firmware_version, ip_address, - gateway, uptime_ms, uptime_display, + gateway, uptime_ms, uptime_display, rssi, free_heap, state, ok, received_at AT TIME ZONE 'UTC' AS received_at FROM heartbeats ORDER BY device_serial, received_at DESC @@ -286,6 +343,548 @@ async def get_alerts(device_serial: str) -> list: return [_row_to_dict(r) for r in rows] +# --------------------------------------------------------------------------- +# Device alert events (insert-only history — WARNING/CRITICAL/FAILED only, +# CLEARED is not logged here). Used to answer "when was the most recent +# issue on this device", surviving past the alert being resolved. +# --------------------------------------------------------------------------- + +async def insert_alert_event(device_serial: str, subsystem: str, state: str, + message: str | None = None): + async with AsyncSessionLocal() as session: + await session.execute( + text(""" + INSERT INTO device_alert_events (device_serial, subsystem, state, message, occurred_at) + VALUES (:serial, :subsystem, :state, :message, now()) + """), + {"serial": device_serial, "subsystem": subsystem, "state": state, "message": message}, + ) + await session.commit() + + +async def get_latest_alert_events() -> list: + """Most recent alert event per device, across all devices — one row each.""" + async with AsyncSessionLocal() as session: + result = await session.execute( + text(""" + SELECT DISTINCT ON (device_serial) + id, device_serial, subsystem, state, message, + occurred_at AT TIME ZONE 'UTC' AS occurred_at + FROM device_alert_events + ORDER BY device_serial, occurred_at DESC + """) + ) + rows = result.mappings().all() + + return [_row_to_dict(r) for r in rows] + + +async def get_alert_events(device_serial: str, limit: int = 100, + offset: int = 0) -> tuple[list, int]: + """Full alert transition history for one device (WARNING/CRITICAL/FAILED — no CLEARED).""" + async with AsyncSessionLocal() as session: + count_result = await session.execute( + text("SELECT COUNT(*) FROM device_alert_events WHERE device_serial = :serial"), + {"serial": device_serial}, + ) + total = count_result.scalar() + + rows_result = await session.execute( + text(""" + SELECT id, device_serial, subsystem, state, message, + occurred_at AT TIME ZONE 'UTC' AS occurred_at + FROM device_alert_events + WHERE device_serial = :serial + ORDER BY occurred_at DESC + LIMIT :limit OFFSET :offset + """), + {"serial": device_serial, "limit": limit, "offset": offset}, + ) + rows = rows_result.mappings().all() + + return [_row_to_dict(r) for r in rows], total + + +async def get_latest_alert_event(device_serial: str) -> dict | None: + async with AsyncSessionLocal() as session: + result = await session.execute( + text(""" + SELECT id, device_serial, subsystem, state, message, + occurred_at AT TIME ZONE 'UTC' AS occurred_at + FROM device_alert_events + WHERE device_serial = :serial + ORDER BY occurred_at DESC + LIMIT 1 + """), + {"serial": device_serial}, + ) + row = result.mappings().first() + + return _row_to_dict(row) if row else None + + +# --------------------------------------------------------------------------- +# Device boot events (insert-only history from firmware's boot_report — see +# mqtt/logger.py::_handle_info). One row per boot, fault or not. Powers the +# Health tab's boot timeline + restart chart. +# --------------------------------------------------------------------------- + +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) -> int: + async with AsyncSessionLocal() as session: + result = 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, occurred_at) + VALUES + (:serial, :boot_count, :reset_reason, :is_fault, :free_heap, + :crash_task, :crash_pc, :crash_exc_cause, :crash_exc_vaddr, now()) + RETURNING id + """), + { + "serial": device_serial, + "boot_count": boot_count, + "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, + }, + ) + row = result.fetchone() + await session.commit() + return row[0] + + +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]: + async with AsyncSessionLocal() as session: + params: dict = {"serial": device_serial} + range_clause = _time_range_clause("occurred_at", since, until, params) + + count_result = await session.execute( + text(f"SELECT COUNT(*) FROM device_boot_events WHERE device_serial = :serial{range_clause}"), + params, + ) + total = count_result.scalar() + + 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 + FROM device_boot_events + WHERE device_serial = :serial{range_clause} + ORDER BY occurred_at DESC + LIMIT :limit OFFSET :offset + """), + {**params, "limit": limit, "offset": offset}, + ) + rows = rows_result.mappings().all() + + return [_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 + FROM device_boot_events + WHERE device_serial = :serial + ORDER BY occurred_at DESC + LIMIT 1 + """), + {"serial": device_serial}, + ) + row = result.mappings().first() + + return _row_to_dict(row) if row else None + + +# --------------------------------------------------------------------------- +# Device ping samples — backend-computed RTT, one row per ping/pong round +# trip. Sampler loop lives in mqtt/client.py; see MqttManager.ping_loop(). +# --------------------------------------------------------------------------- + +async def insert_ping_sample(device_serial: str, rtt_ms: int) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text(""" + INSERT INTO device_ping_samples (device_serial, rtt_ms, sampled_at) + VALUES (:serial, :rtt_ms, now()) + RETURNING id + """), + {"serial": device_serial, "rtt_ms": rtt_ms}, + ) + row = result.fetchone() + await session.commit() + return row[0] + + +async def get_ping_samples(device_serial: str, limit: int = 200, offset: int = 0, + since: datetime | None = None, + until: datetime | None = None) -> tuple[list, int]: + async with AsyncSessionLocal() as session: + params: dict = {"serial": device_serial} + range_clause = _time_range_clause("sampled_at", since, until, params) + + count_result = await session.execute( + text(f"SELECT COUNT(*) FROM device_ping_samples WHERE device_serial = :serial{range_clause}"), + params, + ) + total = count_result.scalar() + + rows_result = await session.execute( + text(f""" + SELECT id, device_serial, rtt_ms, + sampled_at AT TIME ZONE 'UTC' AS sampled_at + FROM device_ping_samples + WHERE device_serial = :serial{range_clause} + ORDER BY sampled_at DESC + LIMIT :limit OFFSET :offset + """), + {**params, "limit": limit, "offset": offset}, + ) + rows = rows_result.mappings().all() + + return [_row_to_dict(r) for r in rows], total + + +async def get_latest_ping_samples() -> list: + # Fleet-wide "last known" ping RTT for DeviceList — same DISTINCT ON + # pattern as get_latest_heartbeats(); devices with no RTC never appear. + async with AsyncSessionLocal() as session: + rows_result = await session.execute( + text(""" + SELECT DISTINCT ON (device_serial) + id, device_serial, rtt_ms, + sampled_at AT TIME ZONE 'UTC' AS sampled_at + FROM device_ping_samples + ORDER BY device_serial, sampled_at DESC + """) + ) + rows = rows_result.mappings().all() + + return [_row_to_dict(r) for r in rows] + + +# --------------------------------------------------------------------------- +# Device diagnostics reports — firmware's periodic (5-min) diagnostics_report +# event: CPU temperature, WiFi reconnects, OTA state, stack high-water marks. +# See project-vesper's feature catalog F-056 and CommunicationRouter README +# for the full design rationale (why this is separate from the 30s heartbeat). +# --------------------------------------------------------------------------- + +async def insert_diagnostics_report( + device_serial: str, + cpu_temp_avg: float | None = None, + cpu_temp_min: float | None = None, + cpu_temp_max: float | None = None, + cpu_temp_samples: int | None = None, + wifi_reconnect_count: int | None = None, + wifi_last_disconnect_reason: str | None = None, + wifi_last_disconnect_uptime_ms: int | None = None, + ota_current_version: str | None = None, + ota_update_available: bool | None = None, + ota_available_version: str | None = None, + ota_last_check_uptime_ms: int | None = None, + ota_last_error: str | None = None, + stack_high_water: dict | None = None, + bell_strikes: dict | None = None, + bell_loads: dict | None = None, + cooling_active: bool | None = None, +) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text(""" + INSERT INTO device_diagnostics_reports + (device_serial, + cpu_temp_avg, cpu_temp_min, cpu_temp_max, cpu_temp_samples, + wifi_reconnect_count, wifi_last_disconnect_reason, wifi_last_disconnect_uptime_ms, + ota_current_version, ota_update_available, ota_available_version, + ota_last_check_uptime_ms, ota_last_error, + stack_high_water, bell_strikes, bell_loads, cooling_active, received_at) + VALUES + (:serial, + :cpu_temp_avg, :cpu_temp_min, :cpu_temp_max, :cpu_temp_samples, + :wifi_reconnect_count, :wifi_last_disconnect_reason, :wifi_last_disconnect_uptime_ms, + :ota_current_version, :ota_update_available, :ota_available_version, + :ota_last_check_uptime_ms, :ota_last_error, + :stack_high_water, :bell_strikes, :bell_loads, :cooling_active, now()) + RETURNING id + """), + { + "serial": device_serial, + "cpu_temp_avg": cpu_temp_avg, + "cpu_temp_min": cpu_temp_min, + "cpu_temp_max": cpu_temp_max, + "cpu_temp_samples": cpu_temp_samples, + "wifi_reconnect_count": wifi_reconnect_count, + "wifi_last_disconnect_reason": wifi_last_disconnect_reason, + "wifi_last_disconnect_uptime_ms": wifi_last_disconnect_uptime_ms, + "ota_current_version": ota_current_version, + "ota_update_available": ota_update_available, + "ota_available_version": ota_available_version, + "ota_last_check_uptime_ms": ota_last_check_uptime_ms, + "ota_last_error": ota_last_error, + "stack_high_water": json.dumps(stack_high_water) if stack_high_water else None, + "bell_strikes": json.dumps(bell_strikes) if bell_strikes else None, + "bell_loads": json.dumps(bell_loads) if bell_loads else None, + "cooling_active": cooling_active, + }, + ) + row = result.fetchone() + await session.commit() + return row[0] + + +async def get_diagnostics_reports(device_serial: str, limit: int = 200, offset: int = 0, + since: datetime | None = None, + until: datetime | None = None) -> tuple[list, int]: + async with AsyncSessionLocal() as session: + params: dict = {"serial": device_serial} + range_clause = _time_range_clause("received_at", since, until, params) + + count_result = await session.execute( + text(f"SELECT COUNT(*) FROM device_diagnostics_reports WHERE device_serial = :serial{range_clause}"), + params, + ) + total = count_result.scalar() + + rows_result = await session.execute( + text(f""" + SELECT id, device_serial, + cpu_temp_avg, cpu_temp_min, cpu_temp_max, cpu_temp_samples, + wifi_reconnect_count, wifi_last_disconnect_reason, wifi_last_disconnect_uptime_ms, + ota_current_version, ota_update_available, ota_available_version, + ota_last_check_uptime_ms, ota_last_error, + stack_high_water, bell_strikes, bell_loads, cooling_active, + received_at AT TIME ZONE 'UTC' AS received_at + FROM device_diagnostics_reports + WHERE device_serial = :serial{range_clause} + ORDER BY received_at DESC + LIMIT :limit OFFSET :offset + """), + {**params, "limit": limit, "offset": offset}, + ) + rows = rows_result.mappings().all() + + return [_row_to_dict(r) for r in rows], total + + +async def get_latest_diagnostics_report(device_serial: str) -> dict | None: + async with AsyncSessionLocal() as session: + result = await session.execute( + text(""" + SELECT id, device_serial, + cpu_temp_avg, cpu_temp_min, cpu_temp_max, cpu_temp_samples, + wifi_reconnect_count, wifi_last_disconnect_reason, wifi_last_disconnect_uptime_ms, + ota_current_version, ota_update_available, ota_available_version, + ota_last_check_uptime_ms, ota_last_error, + stack_high_water, bell_strikes, bell_loads, cooling_active, + received_at AT TIME ZONE 'UTC' AS received_at + FROM device_diagnostics_reports + WHERE device_serial = :serial + ORDER BY received_at DESC + LIMIT 1 + """), + {"serial": device_serial}, + ) + row = result.mappings().first() + + return _row_to_dict(row) if row else None + + +async def get_latest_diagnostics_reports() -> list: + # Fleet-wide "last known" CPU temp for DeviceList — one DISTINCT ON query + # instead of N per-device round trips, same pattern as get_latest_heartbeats(). + async with AsyncSessionLocal() as session: + rows_result = await session.execute( + text(""" + SELECT DISTINCT ON (device_serial) + id, device_serial, + cpu_temp_avg, cpu_temp_min, cpu_temp_max, cpu_temp_samples, + received_at AT TIME ZONE 'UTC' AS received_at + FROM device_diagnostics_reports + ORDER BY device_serial, received_at DESC + """) + ) + rows = rows_result.mappings().all() + + return [_row_to_dict(r) for r in rows] + + +# --------------------------------------------------------------------------- +# Per-device stat reset — clears QA/bench test data for a single device +# before shipping to a customer. Deliberately narrow: only tables that are +# pure accumulated history/telemetry are here. device_alerts (live current +# state) is included but callers should treat it as opt-in/advanced, not a +# default — clearing it could hide a real active fault. Nothing here touches +# device identity, customer assignment, notes, warranty, or audit logs — see +# backend/devices/service.py for the separate device_stats (Firestore) reset. +# --------------------------------------------------------------------------- + +async def delete_device_logs(device_serial: str) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text("DELETE FROM device_logs WHERE device_serial = :serial"), + {"serial": device_serial}, + ) + await session.commit() + return result.rowcount + + +async def delete_device_heartbeats(device_serial: str) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text("DELETE FROM heartbeats WHERE device_serial = :serial"), + {"serial": device_serial}, + ) + await session.commit() + return result.rowcount + + +async def delete_device_commands(device_serial: str) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text("DELETE FROM commands WHERE device_serial = :serial"), + {"serial": device_serial}, + ) + await session.commit() + return result.rowcount + + +async def delete_device_alert_events(device_serial: str) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text("DELETE FROM device_alert_events WHERE device_serial = :serial"), + {"serial": device_serial}, + ) + await session.commit() + return result.rowcount + + +async def delete_device_current_alerts(device_serial: str) -> int: + """Clears LIVE alert state (device_alerts), not history. Opt-in only — + see module note above. If the device currently has a real active fault, + clearing this hides it from the console until the next state change.""" + async with AsyncSessionLocal() as session: + result = await session.execute( + text("DELETE FROM device_alerts WHERE device_serial = :serial"), + {"serial": device_serial}, + ) + await session.commit() + return result.rowcount + + +async def delete_device_boot_events(device_serial: str) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text("DELETE FROM device_boot_events WHERE device_serial = :serial"), + {"serial": device_serial}, + ) + await session.commit() + return result.rowcount + + +async def delete_device_ping_samples(device_serial: str) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text("DELETE FROM device_ping_samples WHERE device_serial = :serial"), + {"serial": device_serial}, + ) + await session.commit() + return result.rowcount + + +async def delete_device_diagnostics_reports(device_serial: str) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text("DELETE FROM device_diagnostics_reports WHERE device_serial = :serial"), + {"serial": device_serial}, + ) + await session.commit() + return result.rowcount + + +# --------------------------------------------------------------------------- +# Device reports — critical, unsolicited, board-initiated events from the new +# control/reports MQTT topic (currently only bell_overload). Console-side +# storage is deliberately light: insert + list for historical/audit purposes. +# The physical tablets, not this console, are the real-time consumer of +# control/reports. See project-vesper's vesper_mqtt_topic_spec_v2.md. +# --------------------------------------------------------------------------- + +async def insert_report(device_serial: str, report_type: str, payload: dict | None) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text(""" + INSERT INTO device_reports (device_serial, report_type, payload, occurred_at) + VALUES (:serial, :report_type, :payload, now()) + RETURNING id + """), + { + "serial": device_serial, + "report_type": report_type, + "payload": json.dumps(payload) if payload else None, + }, + ) + row = result.fetchone() + await session.commit() + return row[0] + + +async def get_reports(device_serial: str, limit: int = 200, offset: int = 0, + since: datetime | None = None, + until: datetime | None = None) -> tuple[list, int]: + async with AsyncSessionLocal() as session: + params: dict = {"serial": device_serial} + range_clause = _time_range_clause("occurred_at", since, until, params) + + count_result = await session.execute( + text(f"SELECT COUNT(*) FROM device_reports WHERE device_serial = :serial{range_clause}"), + params, + ) + total = count_result.scalar() + + rows_result = await session.execute( + text(f""" + SELECT id, device_serial, report_type, payload, + occurred_at AT TIME ZONE 'UTC' AS occurred_at + FROM device_reports + WHERE device_serial = :serial{range_clause} + ORDER BY occurred_at DESC + LIMIT :limit OFFSET :offset + """), + {**params, "limit": limit, "offset": offset}, + ) + rows = rows_result.mappings().all() + + return [_row_to_dict(r) for r in rows], total + + +async def delete_device_reports(device_serial: str) -> int: + async with AsyncSessionLocal() as session: + result = await session.execute( + text("DELETE FROM device_reports WHERE device_serial = :serial"), + {"serial": device_serial}, + ) + await session.commit() + return result.rowcount + + # --------------------------------------------------------------------------- # Partition management # --------------------------------------------------------------------------- @@ -314,8 +913,20 @@ async def ensure_current_partitions(): logger.info("Partition check complete") -async def drop_old_partitions(keep_months: int = 6): - """Drop device_logs partitions older than keep_months.""" +def _global_retention_days() -> int: + """Read the admin-configured global log retention (days), falling back to + the env-configured default if Firestore is unreachable or unset.""" + try: + from settings.log_retention_service import get_log_retention + return get_log_retention().days + except Exception: + return settings.mqtt_data_retention_days + + +async def drop_old_partitions(keep_months: int | None = None): + """Drop device_logs partitions older than keep_months (default: global retention setting, in months).""" + if keep_months is None: + keep_months = max(1, round(_global_retention_days() / 30)) cutoff = _add_months(date.today().replace(day=1), -keep_months) async with AsyncSessionLocal() as session: result = await session.execute(text(""" @@ -361,7 +972,7 @@ async def partition_manager_loop(): # --------------------------------------------------------------------------- async def purge_old_data(retention_days: int | None = None): - days = retention_days or settings.mqtt_data_retention_days + days = retention_days or _global_retention_days() cutoff = datetime.now(timezone.utc) - timedelta(days=days) async with AsyncSessionLocal() as session: await session.execute( @@ -372,8 +983,20 @@ async def purge_old_data(retention_days: int | None = None): text("DELETE FROM commands WHERE sent_at < :cutoff"), {"cutoff": cutoff}, ) + await session.execute( + text("DELETE FROM device_ping_samples WHERE sampled_at < :cutoff"), + {"cutoff": cutoff}, + ) + await session.execute( + text("DELETE FROM device_diagnostics_reports WHERE received_at < :cutoff"), + {"cutoff": cutoff}, + ) + await session.execute( + text("DELETE FROM device_reports WHERE occurred_at < :cutoff"), + {"cutoff": cutoff}, + ) await session.commit() - logger.info(f"Purged heartbeats and commands older than {days} days") + logger.info(f"Purged heartbeats, commands, ping samples, diagnostics reports, and device reports older than {days} days") async def purge_loop(): @@ -403,9 +1026,22 @@ async def close_db(): # --------------------------------------------------------------------------- def _row_to_dict(row) -> dict: - """Convert a SQLAlchemy RowMapping to a plain dict with ISO string timestamps.""" + """Convert a SQLAlchemy RowMapping to a plain dict with ISO string timestamps. + + Every timestamp column here is selected via `AT TIME ZONE 'UTC'`, which hands + back a naive (tzinfo-less) datetime that IS UTC but doesn't say so. Without + attaching tzinfo before .isoformat(), the resulting string has no offset + suffix (e.g. "2026-07-14T16:03:21") — browsers then parse that as LOCAL + time, not UTC, silently shifting every displayed time by the browser's + UTC offset (3h for Greek summer/EEST, 2h for winter/EET). Attaching + timezone.utc makes isoformat() emit the explicit "+00:00" suffix so + `new Date(...)` on the frontend converts to the viewer's local time + correctly year-round, including across DST transitions. + """ d = dict(row) for key, val in d.items(): if isinstance(val, datetime): + if val.tzinfo is None: + val = val.replace(tzinfo=timezone.utc) d[key] = val.isoformat() return d diff --git a/backend/main.py b/backend/main.py index be206ca..a719749 100644 --- a/backend/main.py +++ b/backend/main.py @@ -116,6 +116,7 @@ async def startup(): mqtt_manager.start(asyncio.get_event_loop()) asyncio.create_task(db.partition_manager_loop()) asyncio.create_task(db.purge_loop()) + asyncio.create_task(mqtt_manager.ping_loop()) asyncio.create_task(nextcloud_keepalive_loop()) asyncio.create_task(crm_poll_loop()) sync_accounts = [a for a in get_mail_accounts() if a.get("sync_inbound") and a.get("imap_host")] diff --git a/backend/mqtt/client.py b/backend/mqtt/client.py index 103f806..f88a8d9 100644 --- a/backend/mqtt/client.py +++ b/backend/mqtt/client.py @@ -1,12 +1,19 @@ import json import logging import asyncio +import time from typing import Set import paho.mqtt.client as paho_mqtt from config import settings logger = logging.getLogger("mqtt.client") +PING_INTERVAL_SECONDS = 60 +# Only devices heard from within this window get pinged — no point spending +# broker traffic/RTT samples on a device that's already known offline; its +# heartbeat-derived "online" state will already reflect that on the console. +PING_ONLINE_WINDOW_SECONDS = 90 + class MqttManager: """Singleton MQTT client manager.""" @@ -61,12 +68,17 @@ class MqttManager: if reason_code == 0: self._connected = True logger.info("MQTT connected, subscribing to topics") + # v2 topic set — see vesper_mqtt_topic_spec_v2.md in the firmware repo. + # control/command is inbound-to-device only; the console never subscribes to it. client.subscribe([ - ("vesper/+/data", 1), + ("vesper/+/control/ack", 1), + ("vesper/+/control/reports", 1), ("vesper/+/status/heartbeat", 1), - ("vesper/+/status/alerts", 1), - ("vesper/+/status/info", 0), - ("vesper/+/logs", 1), + ("vesper/+/status/playback", 1), + ("vesper/+/system/alerts", 1), + ("vesper/+/system/info", 1), + ("vesper/+/system/logs", 0), + ("vesper/+/system/metrics", 0), ]) else: logger.error(f"MQTT connection failed: {reason_code}") @@ -132,10 +144,48 @@ class MqttManager: if not self._client or not self._connected: return False - topic = f"vesper/{device_serial}/control" - payload = json.dumps({"cmd": cmd, "contents": contents}) + topic = f"vesper/{device_serial}/control/command" + payload = json.dumps({"v": 2, "cmd": cmd, "contents": contents}) result = self._client.publish(topic, payload, qos=1) return result.rc == paho_mqtt.MQTT_ERR_SUCCESS + async def ping_loop(self): + """Periodically pings every recently-online device with a client + timestamp so mqtt/logger.py::_handle_data_response can compute RTT + from the echoed pong. Deliberately bypasses db.insert_command — this + is a background health check, not a user-initiated command, and + shouldn't clutter the Control tab's command history. + + Only available on RTC-equipped firmware builds (see API Reference — + ping is compiled out on agnus/agnus-mini). A pong simply never + arrives for those devices, so no ping-latency samples accumulate for + them; the Health tab handles an empty series as "unsupported". + """ + while True: + await asyncio.sleep(PING_INTERVAL_SECONDS) + try: + await self._ping_online_devices() + except Exception as e: + logger.error(f"Ping loop error: {e}") + + async def _ping_online_devices(self): + import database as db + heartbeats = await db.get_latest_heartbeats() + now = time.time() + for hb in heartbeats: + try: + from datetime import datetime + received = datetime.fromisoformat(hb["received_at"]) + age = now - received.timestamp() + except (ValueError, TypeError, KeyError): + continue + if age > PING_ONLINE_WINDOW_SECONDS: + continue + self.publish_command( + device_serial=hb["device_serial"], + cmd="ping", + contents={"ts": int(now * 1000)}, + ) + mqtt_manager = MqttManager() diff --git a/backend/mqtt/logger.py b/backend/mqtt/logger.py index 302e922..deddcea 100644 --- a/backend/mqtt/logger.py +++ b/backend/mqtt/logger.py @@ -1,4 +1,5 @@ import logging +import time import database as db logger = logging.getLogger("mqtt.logger") @@ -16,16 +17,23 @@ LEVEL_MAP = { async def handle_message(serial: str, topic_type: str, payload: dict): try: + # v2 topic set — see project-vesper's vesper_mqtt_topic_spec_v2.md. if topic_type == "status/heartbeat": await _handle_heartbeat(serial, payload) - elif topic_type == "status/alerts": + elif topic_type == "system/alerts": await _handle_alerts(serial, payload) - elif topic_type == "status/info": + elif topic_type == "system/info": await _handle_info(serial, payload) - elif topic_type == "logs": + elif topic_type == "system/logs": await _handle_log(serial, payload) - elif topic_type == "data": - await _handle_data_response(serial, payload) + elif topic_type == "system/metrics": + await _handle_metrics(serial, payload) + elif topic_type == "control/ack": + await _handle_ack(serial, payload) + elif topic_type == "control/reports": + await _handle_report(serial, payload) + elif topic_type == "status/playback": + pass # no console-side storage needed — WS broadcast already carries it live else: logger.debug(f"Unhandled topic type: {topic_type} for {serial}") except Exception as e: @@ -35,15 +43,21 @@ async def handle_message(serial: str, topic_type: str, payload: dict): async def _handle_heartbeat(serial: str, payload: dict): # Store silently — do not log as a visible event. # The console surfaces an alert only when the device goes silent (no heartbeat for 90s). - inner = payload.get("payload", {}) + # v2 heartbeat payload is FLAT — no {"status","type","payload"} wrapper, + # and field names changed: firmware_version -> fw_version, timestamp -> uptime_human. + # See vesper_mqtt_topic_spec_v2.md and mqtt-events.md for the full shape. await db.insert_heartbeat( device_serial=serial, - device_id=inner.get("device_id", ""), - firmware_version=inner.get("firmware_version", ""), - ip_address=inner.get("ip_address", ""), - gateway=inner.get("gateway", ""), - uptime_ms=inner.get("uptime_ms", 0), - uptime_display=inner.get("timestamp", ""), + device_id=payload.get("device_id", ""), + firmware_version=payload.get("fw_version", ""), + ip_address=payload.get("ip_address", ""), + gateway=payload.get("gateway", ""), + uptime_ms=payload.get("uptime_ms", 0), + uptime_display=payload.get("uptime_human", ""), + rssi=payload.get("rssi"), + free_heap=payload.get("free_heap"), + state=payload.get("state"), + ok=payload.get("ok"), ) @@ -58,6 +72,7 @@ async def _handle_log(serial: str, payload: dict): level=level, message=message, device_timestamp=device_timestamp, + source="log", ) @@ -72,23 +87,136 @@ async def _handle_alerts(serial: str, payload: dict): await db.delete_alert(serial, subsystem) else: 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")) async def _handle_info(serial: str, payload: dict): event_type = payload.get("type", "") + + # boot_report carries structured fields (boot_count, crash detail, etc.) — + # parsed into device_boot_events instead of flattened into a text log line, + # so the Health tab can query/chart it. Not routed through insert_log at + # all: a text summary here would just duplicate the structured row. + # + # 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) + return + data = payload.get("payload", {}) if event_type == "playback_started": - logger.debug(f"{serial}: playback started — melody_uid={data.get('melody_uid')}") + message = f"Playback started — melody_uid={data.get('melody_uid')}" elif event_type == "playback_stopped": - logger.debug(f"{serial}: playback stopped") + message = "Playback stopped" else: - logger.debug(f"{serial}: info event '{event_type}'") + message = f"Info event '{event_type}'" + (f" — {data}" if data else "") + + await db.insert_log( + device_serial=serial, + level="INFO", + message=message, + source="info", + ) -async def _handle_data_response(serial: str, payload: dict): +async def _handle_boot_report(serial: str, payload: dict): + """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 + nested under "payload" — CommunicationRouter::publishInfo() wraps every + info event as { "type": ..., "payload": {...} }. + """ + data = payload.get("payload", {}) + crash = data.get("crash") or {} + 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"), + ) + + +async def _handle_metrics(serial: str, payload: dict): + """Parses the firmware's periodic metrics payload (see + CommunicationRouter::reportDiagnosticsIfDue() in project-vesper) into + device_diagnostics_reports. Unlike status/info events, system/metrics has + NO {"type","payload"} wrapper — the payload IS the metrics object, + flat at the top level. See vesper_mqtt_topic_spec_v2.md / mqtt-events.md. + + cpu_temp is entirely absent from the firmware payload (not present as a + key at all, not present-with-nulls) when no samples were taken yet this + window — .get() on a missing "cpu_temp" key correctly falls back to {} + below, leaving every cpu_temp_* column null for that row. + """ + cpu_temp = payload.get("cpu_temp") or {} + wifi = payload.get("wifi_reconnects") or {} + ota = payload.get("ota") or {} + stack = payload.get("stack_high_water") or {} + bell_strikes = payload.get("bell_strikes") or {} + bell_loads = payload.get("bell_loads") or {} + + await db.insert_diagnostics_report( + device_serial=serial, + cpu_temp_avg=cpu_temp.get("avg"), + cpu_temp_min=cpu_temp.get("min"), + cpu_temp_max=cpu_temp.get("max"), + cpu_temp_samples=cpu_temp.get("samples"), + wifi_reconnect_count=wifi.get("lifetime_count"), + wifi_last_disconnect_reason=wifi.get("last_reason"), + wifi_last_disconnect_uptime_ms=wifi.get("last_at_uptime_ms"), + ota_current_version=ota.get("current_version"), + ota_update_available=ota.get("update_available"), + ota_available_version=ota.get("available_version"), + ota_last_check_uptime_ms=ota.get("last_check_uptime_ms"), + ota_last_error=ota.get("last_error"), + stack_high_water=stack if stack else None, + bell_strikes=bell_strikes if bell_strikes else None, + bell_loads=bell_loads if bell_loads else None, + cooling_active=payload.get("cooling_active"), + ) + + +async def _handle_report(serial: str, payload: dict): + """Parses control/reports — critical, unsolicited board-initiated events + (currently only bell_overload). Console-side storage is intentionally + light: history/audit only. The tablets are the real-time consumer.""" + report_type = payload.get("type", "") + if not report_type: + logger.warning(f"Malformed report payload from {serial}: {payload}") + return + + await db.insert_report( + device_serial=serial, + report_type=report_type, + payload=payload.get("payload"), + ) + logger.warning(f"Report '{report_type}' received from {serial}: {payload.get('payload')}") + + +async def _handle_ack(serial: str, payload: dict): status = payload.get("status", "") + # Health-check pings (see MqttManager.ping_loop) are published directly, + # bypassing db.insert_command, so scheduled liveness checks don't clutter + # the user-facing command history. The device echoes the caller's own + # send-time ("ts", epoch-ms) back unchanged, so RTT is computed here + # without needing to correlate against any pending-command bookkeeping. + if payload.get("type") == "pong": + ts = (payload.get("data") or {}).get("ts") + if isinstance(ts, (int, float)): + rtt_ms = max(0, int(time.time() * 1000) - int(ts)) + await db.insert_ping_sample(device_serial=serial, rtt_ms=rtt_ms) + return + pending = await db.get_pending_command(serial) if pending: cmd_status = "success" if status == "SUCCESS" else "error" @@ -98,4 +226,4 @@ async def _handle_data_response(serial: str, payload: dict): response_payload=payload, ) else: - logger.debug(f"Received data response for {serial} with no pending command") + logger.debug(f"Received control/ack for {serial} with no pending command") diff --git a/backend/mqtt/models.py b/backend/mqtt/models.py index 4a11451..6922aa2 100644 --- a/backend/mqtt/models.py +++ b/backend/mqtt/models.py @@ -27,6 +27,7 @@ class DeviceLogEntry(BaseModel): level: str message: str device_timestamp: Optional[int] = None + source: str = "log" received_at: str @@ -39,6 +40,10 @@ class HeartbeatEntry(BaseModel): gateway: Optional[str] = None uptime_ms: Optional[int] = None uptime_display: Optional[str] = None + rssi: Optional[int] = None + free_heap: Optional[int] = None + state: Optional[str] = None # "idle" | "playing" | "paused" | "error" | "booting" (v2 firmware only) + ok: Optional[bool] = None # overall device health, independent of playback state (v2 firmware only) received_at: str @@ -53,11 +58,21 @@ class CommandEntry(BaseModel): responded_at: Optional[str] = None +class AlertEventEntry(BaseModel): + id: int + device_serial: str + subsystem: str + state: str + message: Optional[str] = None + occurred_at: str + + class DeviceMqttStatus(BaseModel): device_serial: str online: bool last_heartbeat: Optional[HeartbeatEntry] = None seconds_since_heartbeat: Optional[int] = None + last_alert_event: Optional[AlertEventEntry] = None class MqttStatusResponse(BaseModel): @@ -96,3 +111,99 @@ class DeviceAlertEntry(BaseModel): class DeviceAlertsResponse(BaseModel): alerts: List[DeviceAlertEntry] + + +class AlertEventListResponse(BaseModel): + events: List[AlertEventEntry] + total: int + + +class BootEventEntry(BaseModel): + id: int + device_serial: str + boot_count: Optional[int] = None + reset_reason: Optional[str] = None + is_fault: bool = False + free_heap: Optional[int] = None + crash_task: Optional[str] = None + crash_pc: Optional[int] = None + crash_exc_cause: Optional[int] = None + crash_exc_vaddr: Optional[int] = None + occurred_at: str + + +class BootEventListResponse(BaseModel): + events: List[BootEventEntry] + total: int + + +class PingSampleEntry(BaseModel): + id: int + device_serial: str + rtt_ms: int + sampled_at: str + + +class PingSampleListResponse(BaseModel): + samples: List[PingSampleEntry] + total: int + + +class DiagnosticsReportEntry(BaseModel): + id: int + device_serial: str + cpu_temp_avg: Optional[float] = None + cpu_temp_min: Optional[float] = None + cpu_temp_max: Optional[float] = None + cpu_temp_samples: Optional[int] = None + wifi_reconnect_count: Optional[int] = None + wifi_last_disconnect_reason: Optional[str] = None + wifi_last_disconnect_uptime_ms: Optional[int] = None + ota_current_version: Optional[str] = None + ota_update_available: Optional[bool] = None + ota_available_version: Optional[str] = None + ota_last_check_uptime_ms: Optional[int] = None + ota_last_error: Optional[str] = None + stack_high_water: Optional[str] = None # JSON-encoded string — see insert_diagnostics_report + bell_strikes: Optional[str] = None # JSON-encoded string, {"0": count, ...} — v2 firmware only + bell_loads: Optional[str] = None # JSON-encoded string, {"0": load, ...} — v2 firmware only + cooling_active: Optional[bool] = None # v2 firmware only + received_at: str + + +class DiagnosticsReportListResponse(BaseModel): + reports: List[DiagnosticsReportEntry] + total: int + + +class DeviceReportEntry(BaseModel): + """A control/reports event (currently only bell_overload) — see + project-vesper's vesper_mqtt_topic_spec_v2.md. Console-side storage is + intentionally light: history/audit only, not a real-time UI surface.""" + id: int + device_serial: str + report_type: str + payload: Optional[str] = None # JSON-encoded string + occurred_at: str + + +class DeviceReportListResponse(BaseModel): + reports: List[DeviceReportEntry] + total: int + + +class LatestDiagnosticsEntry(BaseModel): + device_serial: str + cpu_temp_avg: Optional[float] = None + received_at: str + + +class LatestPingEntry(BaseModel): + device_serial: str + rtt_ms: int + sampled_at: str + + +class LatestMetricsResponse(BaseModel): + diagnostics: List[LatestDiagnosticsEntry] + pings: List[LatestPingEntry] diff --git a/backend/mqtt/router.py b/backend/mqtt/router.py index 3ddb2d4..79a083e 100644 --- a/backend/mqtt/router.py +++ b/backend/mqtt/router.py @@ -1,11 +1,14 @@ from fastapi import APIRouter, Depends, Query, WebSocket, WebSocketDisconnect -from typing import Optional +from typing import Optional, List from auth.models import TokenPayload from auth.dependencies import require_permission from mqtt.models import ( MqttCommandRequest, CommandSendResponse, MqttStatusResponse, DeviceMqttStatus, LogListResponse, HeartbeatListResponse, - CommandListResponse, HeartbeatEntry, + CommandListResponse, HeartbeatEntry, AlertEventEntry, + AlertEventListResponse, BootEventListResponse, PingSampleListResponse, + DiagnosticsReportListResponse, LatestMetricsResponse, + LatestDiagnosticsEntry, LatestPingEntry, DeviceReportListResponse, ) from mqtt.client import mqtt_manager import database as db @@ -19,6 +22,8 @@ async def get_all_device_status( _user: TokenPayload = Depends(require_permission("mqtt", "view")), ): heartbeats = await db.get_latest_heartbeats() + alert_events = await db.get_latest_alert_events() + alert_by_serial = {a["device_serial"]: a for a in alert_events} now = datetime.now(timezone.utc) devices = [] for hb in heartbeats: @@ -31,11 +36,14 @@ async def get_all_device_status( except (ValueError, TypeError): seconds_ago = 9999 + alert_event = alert_by_serial.get(hb["device_serial"]) + devices.append(DeviceMqttStatus( device_serial=hb["device_serial"], online=seconds_ago < 90, last_heartbeat=HeartbeatEntry(**hb), seconds_since_heartbeat=seconds_ago, + last_alert_event=AlertEventEntry(**alert_event) if alert_event else None, )) return MqttStatusResponse( devices=devices, @@ -43,6 +51,21 @@ async def get_all_device_status( ) +@router.get("/latest-metrics", response_model=LatestMetricsResponse) +async def get_latest_metrics( + _user: TokenPayload = Depends(require_permission("mqtt", "view")), +): + # Fleet-wide "last known" CPU temp + ping RTT, one query each — used by + # DeviceList to render optional columns without polling any device. + # Uptime/firmware/RSSI don't need this: they're already in /mqtt/status. + diag_reports = await db.get_latest_diagnostics_reports() + ping_samples = await db.get_latest_ping_samples() + return LatestMetricsResponse( + diagnostics=[LatestDiagnosticsEntry(**d) for d in diag_reports], + pings=[LatestPingEntry(**p) for p in ping_samples], + ) + + @router.post("/command/{device_serial}", response_model=CommandSendResponse) async def send_command( device_serial: str, @@ -81,14 +104,18 @@ async def send_command( async def get_device_logs( device_serial: str, level: Optional[str] = Query(None, description="Filter: INFO, WARN, ERROR"), + min_level: bool = Query(False, description="If true, level is a floor — also includes higher-severity levels"), search: Optional[str] = Query(None), + source: Optional[List[str]] = Query(None, description="Filter by source: log, info. Repeat param to include several."), limit: int = Query(100, ge=1, le=1000), offset: int = Query(0, ge=0), + since: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/after this time"), + until: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/before this time"), _user: TokenPayload = Depends(require_permission("mqtt", "view")), ): logs, total = await db.get_logs( - device_serial, level=level, search=search, - limit=limit, offset=offset, + device_serial, level=level, search=search, source=source, + min_level=min_level, limit=limit, offset=offset, since=since, until=until, ) return LogListResponse(logs=logs, total=total) @@ -96,12 +123,14 @@ async def get_device_logs( @router.get("/heartbeats/{device_serial}", response_model=HeartbeatListResponse) async def get_device_heartbeats( device_serial: str, - limit: int = Query(100, ge=1, le=1000), + limit: int = Query(100, ge=1, le=5000), offset: int = Query(0, ge=0), + since: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/after this time"), + until: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/before this time"), _user: TokenPayload = Depends(require_permission("mqtt", "view")), ): heartbeats, total = await db.get_heartbeats( - device_serial, limit=limit, offset=offset, + device_serial, limit=limit, offset=offset, since=since, until=until, ) return HeartbeatListResponse(heartbeats=heartbeats, total=total) @@ -119,6 +148,82 @@ async def get_device_commands( return CommandListResponse(commands=commands, total=total) +@router.get("/alert-events/{device_serial}", response_model=AlertEventListResponse) +async def get_device_alert_events( + device_serial: str, + limit: int = Query(100, ge=1, le=1000), + offset: int = Query(0, ge=0), + _user: TokenPayload = Depends(require_permission("mqtt", "view")), +): + events, total = await db.get_alert_events( + device_serial, limit=limit, offset=offset, + ) + return AlertEventListResponse(events=events, total=total) + + +@router.get("/boot-events/{device_serial}", response_model=BootEventListResponse) +async def get_device_boot_events( + device_serial: str, + limit: int = Query(100, ge=1, le=2000), + offset: int = Query(0, ge=0), + since: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/after this time"), + until: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/before this time"), + _user: TokenPayload = Depends(require_permission("mqtt", "view")), +): + events, total = await db.get_boot_events( + device_serial, limit=limit, offset=offset, since=since, until=until, + ) + return BootEventListResponse(events=events, total=total) + + +@router.get("/ping-samples/{device_serial}", response_model=PingSampleListResponse) +async def get_device_ping_samples( + device_serial: str, + limit: int = Query(200, ge=1, le=5000), + offset: int = Query(0, ge=0), + since: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/after this time"), + until: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/before this time"), + _user: TokenPayload = Depends(require_permission("mqtt", "view")), +): + samples, total = await db.get_ping_samples( + device_serial, limit=limit, offset=offset, since=since, until=until, + ) + return PingSampleListResponse(samples=samples, total=total) + + +@router.get("/diagnostics-reports/{device_serial}", response_model=DiagnosticsReportListResponse) +async def get_device_diagnostics_reports( + device_serial: str, + limit: int = Query(200, ge=1, le=5000), + offset: int = Query(0, ge=0), + since: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/after this time"), + until: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/before this time"), + _user: TokenPayload = Depends(require_permission("mqtt", "view")), +): + reports, total = await db.get_diagnostics_reports( + device_serial, limit=limit, offset=offset, since=since, until=until, + ) + return DiagnosticsReportListResponse(reports=reports, total=total) + + +@router.get("/reports/{device_serial}", response_model=DeviceReportListResponse) +async def get_device_reports( + device_serial: str, + limit: int = Query(200, ge=1, le=2000), + offset: int = Query(0, ge=0), + since: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/after this time"), + until: Optional[datetime] = Query(None, description="ISO timestamp — only rows at/before this time"), + _user: TokenPayload = Depends(require_permission("mqtt", "view")), +): + """Critical, unsolicited board-initiated events from control/reports + (currently only bell_overload). History/audit only — the tablets are the + real-time consumer of this data, not the console.""" + reports, total = await db.get_reports( + device_serial, limit=limit, offset=offset, since=since, until=until, + ) + return DeviceReportListResponse(reports=reports, total=total) + + @router.websocket("/ws") async def mqtt_websocket(websocket: WebSocket): """Live MQTT data stream. Auth via query param: ?token=JWT"""