""" Phase 5 — MQTT live data functions backed by Postgres. device_logs is a partitioned table; heartbeats and commands are plain tables. All three are accessed via raw SQL (not ORM) because device_logs partitioning does not play well with SQLAlchemy's declarative ORM. device_alerts is an ORM model (devices/orm.py) and is handled here via raw SQL to keep a single consistent interface for callers that used to import from database.core. """ import asyncio import json import logging from datetime import date, datetime, timedelta, timezone from sqlalchemy import text from config import settings from database.postgres import AsyncSessionLocal logger = logging.getLogger("database.pg_mqtt") # --------------------------------------------------------------------------- # Insert operations # --------------------------------------------------------------------------- async def insert_log(device_serial: str, level: str, message: str, 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, source) VALUES (:serial, :level, :message, :ts, now(), :source) RETURNING id """), {"serial": device_serial, "level": level, "message": message, "ts": device_timestamp, "source": source}, ) row = result.fetchone() await session.commit() return row[0] async def insert_heartbeat(device_serial: str, device_id: str, firmware_version: str, ip_address: str, 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, rssi, free_heap, state, ok, received_at) VALUES (:serial, :device_id, :fw, :ip, :gw, :uptime_ms, :uptime_display, :rssi, :free_heap, :state, :ok, now()) RETURNING id """), { "serial": device_serial, "device_id": device_id, "fw": firmware_version, "ip": ip_address, "gw": gateway, "uptime_ms": uptime_ms, "uptime_display": uptime_display, "rssi": rssi, "free_heap": free_heap, "state": state, "ok": ok, }, ) row = result.fetchone() await session.commit() return row[0] async def insert_command(device_serial: str, command_name: str, command_payload: dict) -> int: async with AsyncSessionLocal() as session: result = await session.execute( text(""" INSERT INTO commands (device_serial, command_name, command_payload, sent_at) VALUES (:serial, :name, :payload, now()) RETURNING id """), { "serial": device_serial, "name": command_name, "payload": json.dumps(command_payload), }, ) row = result.fetchone() await session.commit() return row[0] async def update_command_response(command_id: int, status: str, response_payload: dict | None = None): async with AsyncSessionLocal() as session: await session.execute( text(""" UPDATE commands SET status = :status, response_payload = :payload, responded_at = now() WHERE id = :id """), { "id": command_id, "status": status, "payload": json.dumps(response_payload) if response_payload else None, }, ) await session.commit() # --------------------------------------------------------------------------- # Query operations # --------------------------------------------------------------------------- LEVEL_ORDER = ["INFO", "WARN", "ERROR"] async def get_logs(device_serial: str, level: str | None = None, 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: 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( text(f"SELECT COUNT(*) FROM device_logs WHERE {where}"), params ) total = count_result.scalar() rows_result = await session.execute( text(f""" SELECT id, device_serial, level, message, device_timestamp, source, received_at AT TIME ZONE 'UTC' AS received_at FROM device_logs WHERE {where} ORDER BY received_at DESC LIMIT :limit OFFSET :offset """), params, ) rows = rows_result.mappings().all() return [_row_to_dict(r) for r in rows], total 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(f"SELECT COUNT(*) FROM heartbeats WHERE device_serial = :serial{range_clause}"), params, ) total = count_result.scalar() rows_result = await session.execute( text(f""" SELECT id, device_serial, device_id, firmware_version, ip_address, 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{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_commands(device_serial: str, limit: int = 100, offset: int = 0) -> tuple[list, int]: async with AsyncSessionLocal() as session: count_result = await session.execute( text("SELECT COUNT(*) FROM commands WHERE device_serial = :serial"), {"serial": device_serial}, ) total = count_result.scalar() rows_result = await session.execute( text(""" SELECT id, device_serial, command_name, command_payload, status, response_payload, sent_at AT TIME ZONE 'UTC' AS sent_at, responded_at AT TIME ZONE 'UTC' AS responded_at FROM commands WHERE device_serial = :serial ORDER BY sent_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_heartbeats() -> list: async with AsyncSessionLocal() as session: rows_result = await session.execute( text(""" SELECT DISTINCT ON (device_serial) id, device_serial, device_id, firmware_version, ip_address, 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 """) ) rows = rows_result.mappings().all() return [_row_to_dict(r) for r in rows] async def get_pending_command(device_serial: str) -> dict | None: async with AsyncSessionLocal() as session: result = await session.execute( text(""" SELECT id, device_serial, command_name, command_payload, status, response_payload, sent_at AT TIME ZONE 'UTC' AS sent_at, responded_at AT TIME ZONE 'UTC' AS responded_at FROM commands WHERE device_serial = :serial AND status = 'pending' ORDER BY sent_at DESC LIMIT 1 """), {"serial": device_serial}, ) row = result.mappings().fetchone() return _row_to_dict(row) if row else None # --------------------------------------------------------------------------- # Device alerts # --------------------------------------------------------------------------- async def upsert_alert(device_serial: str, subsystem: str, state: str, 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: result = await session.execute( text(""" INSERT INTO device_alerts (device_serial, subsystem, state, message, updated_at) VALUES (:serial, :subsystem, :state, :message, now()) ON CONFLICT (device_serial, subsystem) DO UPDATE SET state = EXCLUDED.state, message = EXCLUDED.message, updated_at = EXCLUDED.updated_at WHERE device_alerts.state IS DISTINCT FROM EXCLUDED.state OR device_alerts.message IS DISTINCT FROM EXCLUDED.message RETURNING id """), {"serial": device_serial, "subsystem": subsystem, "state": state, "message": message}, ) changed = result.fetchone() is not None await session.commit() return changed async def delete_alert(device_serial: str, subsystem: str): async with AsyncSessionLocal() as session: await session.execute( text("DELETE FROM device_alerts WHERE device_serial = :serial AND subsystem = :subsystem"), {"serial": device_serial, "subsystem": subsystem}, ) await session.commit() async def get_alerts(device_serial: str) -> list: async with AsyncSessionLocal() as session: result = await session.execute( text(""" SELECT id, device_serial, subsystem, state, message, updated_at AT TIME ZONE 'UTC' AS updated_at FROM device_alerts WHERE device_serial = :serial ORDER BY updated_at DESC """), {"serial": device_serial}, ) rows = result.mappings().all() 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. # --------------------------------------------------------------------------- _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: 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 device's most recent row: that boot is already recorded. 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. `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( text(""" SELECT boot_count, reset_reason FROM device_boot_events WHERE device_serial = :serial ORDER BY occurred_at DESC LIMIT 1 """), {"serial": device_serial}, ) row = latest.first() if row is not None and row.boot_count == boot_count and row.reset_reason == reset_reason: return None 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, 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, CAST(CAST(:crash AS TEXT) AS JSONB), CAST(CAST(:pre_crash AS TEXT) AS JSONB), now()) RETURNING id """), { "serial": device_serial, "boot_count": boot_count, "reset_reason": reset_reason, "is_fault": is_fault, "free_heap": free_heap, **_crash_columns(crash), "crash": _json_or_none(crash), "pre_crash": _json_or_none(pre_crash), }, ) row = result.fetchone() await session.commit() 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]: 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 {_BOOT_COLUMNS} 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 [_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(f""" SELECT {_BOOT_COLUMNS} FROM device_boot_events WHERE device_serial = :serial ORDER BY occurred_at DESC LIMIT 1 """), {"serial": device_serial}, ) row = result.mappings().first() 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] # --------------------------------------------------------------------------- # 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 # --------------------------------------------------------------------------- def _add_months(d: date, months: int) -> date: month = d.month - 1 + months year = d.year + month // 12 month = month % 12 + 1 return d.replace(year=year, month=month, day=1) async def ensure_current_partitions(): """Create device_logs partitions for the current and next month if missing.""" async with AsyncSessionLocal() as session: for month_offset in (0, 1): d = _add_months(date.today().replace(day=1), month_offset) partition_name = f"device_logs_{d.strftime('%Y_%m')}" start = d.isoformat() end = _add_months(d, 1).isoformat() await session.execute(text(f""" CREATE TABLE IF NOT EXISTS {partition_name} PARTITION OF device_logs FOR VALUES FROM ('{start}') TO ('{end}') """)) await session.commit() logger.info("Partition check complete") 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(""" SELECT tablename FROM pg_tables WHERE schemaname = 'public' AND tablename LIKE 'device_logs_%' """)) partitions = [r[0] for r in result.fetchall()] for name in partitions: # name format: device_logs_YYYY_MM parts = name.split("_") if len(parts) != 4: continue try: partition_date = date(int(parts[2]), int(parts[3]), 1) except ValueError: continue if partition_date < cutoff: async with AsyncSessionLocal() as session: await session.execute(text(f"DROP TABLE IF EXISTS {name}")) await session.commit() logger.info(f"Dropped old partition: {name}") async def partition_manager_loop(): """Runs once on startup, then monthly thereafter.""" await ensure_current_partitions() while True: # Sleep ~30 days, wake up and ensure next month's partition exists await asyncio.sleep(30 * 24 * 3600) try: await ensure_current_partitions() await drop_old_partitions() except Exception as e: logger.error(f"Partition manager error: {e}") # --------------------------------------------------------------------------- # Cleanup (replaces SQLite purge_loop — now a no-op since Postgres uses # partition drops instead of row-by-row deletes for device_logs; heartbeats # and commands are still purged by row deletion) # --------------------------------------------------------------------------- async def purge_old_data(retention_days: int | None = None): days = retention_days or _global_retention_days() cutoff = datetime.now(timezone.utc) - timedelta(days=days) async with AsyncSessionLocal() as session: await session.execute( text("DELETE FROM heartbeats WHERE received_at < :cutoff"), {"cutoff": cutoff}, ) await session.execute( 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, commands, ping samples, diagnostics reports, and device reports older than {days} days") async def purge_loop(): while True: await asyncio.sleep(86400) try: await purge_old_data() except Exception as e: logger.error(f"Purge failed: {e}") # --------------------------------------------------------------------------- # Stub — no longer needed but kept so nothing that imports init_db/close_db breaks # --------------------------------------------------------------------------- async def init_db(): """No-op: Postgres schema is managed by Alembic, not runtime init.""" logger.info("Postgres MQTT backend active — no SQLite init needed") async def close_db(): """No-op: SQLAlchemy engine lifecycle is managed by the process.""" # --------------------------------------------------------------------------- # Helpers # --------------------------------------------------------------------------- def _row_to_dict(row) -> dict: """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