import logging import time import database as db logger = logging.getLogger("mqtt.logger") LEVEL_MAP = { "🟢 INFO": "INFO", "🟡 WARN": "WARN", "🔴 EROR": "ERROR", "INFO": "INFO", "WARN": "WARN", "ERROR": "ERROR", "EROR": "ERROR", } 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 == "system/alerts": await _handle_alerts(serial, payload) elif topic_type == "system/info": await _handle_info(serial, payload) elif topic_type == "system/logs": await _handle_log(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: logger.error(f"Error handling {topic_type} for {serial}: {e}") 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). # 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=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"), ) async def _handle_log(serial: str, payload: dict): raw_level = payload.get("level", "INFO") level = LEVEL_MAP.get(raw_level, "INFO") message = payload.get("message", "") device_timestamp = payload.get("timestamp") await db.insert_log( device_serial=serial, level=level, message=message, device_timestamp=device_timestamp, source="log", ) async def _handle_alerts(serial: str, payload: dict): subsystem = payload.get("subsystem", "") state = payload.get("state", "") if not subsystem or not state: logger.warning(f"Malformed alert payload from {serial}: {payload}") return if state == "CLEARED": await db.delete_alert(serial, subsystem) else: changed = await db.upsert_alert(serial, subsystem, state, payload.get("msg")) # Append-only history — survives past the alert being resolved, used to # answer "when was the most recent issue" even once it's cleared. # Skipped when nothing changed: alerts are retained, so the broker # re-sends every active one on each backend (re)subscribe. if changed: await db.insert_alert_event(serial, subsystem, state, payload.get("msg")) async def _handle_info(serial: str, payload: dict): 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": message = f"Playback started — melody_uid={data.get('melody_uid')}" elif event_type == "playback_stopped": message = "Playback stopped" else: 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_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" await db.update_command_response( command_id=pending["id"], status=cmd_status, response_payload=payload, ) else: logger.debug(f"Received control/ack for {serial} with no pending command")