import asyncio import json import time from datetime import datetime from typing import Awaitable, Callable, Dict, Optional import aiomqtt from banban.service.binding import BindingService from banban.service.device_alarm import device_alarm_service from banban.service.device_setting import device_setting_service from banban.service.im import im_service from banban.service.location import location_service from config import settings from services.card_service import card_service from services.device_target_cache import device_target_cache from services.device_identity_initializer import ( DeviceIdentityInitializationError, device_identity_initializer, ) from services.device_update_manager import device_firmware_update_manager from services.offline_audio_cache import offline_audio_cache from services.task_manager import task_manager from utils.logger import session_logger as logger class TalkingQMQTTService: _instance = None _lock = asyncio.Lock() def __init__(self, config: dict): self.broker = config.get("broker", "broker.emqx.io") self.port = config.get("port", 1883) self.username = config.get("username") self.password = config.get("password") self.device_prefix = config.get("device_prefix", "TalkingQ") self.qos = config.get("qos", 0) self.keepalive = config.get("keepalive", 60) self.nfc_notice_interval = config.get("nfc_notice_interval", 600) self._client: Optional[aiomqtt.Client] = None self._connected = False self._connect_lock = asyncio.Lock() self._message_task: Optional[asyncio.Task] = None self._msg_handlers: Dict[str, Callable[[str, dict], Awaitable[None]]] = { "000": self._handle_device_info, "001": self._handle_gps_response, "002": self._handle_volume_response, "003": self._handle_ota_response, "004": self._handle_nfc_notice_response, "005": self._handle_nfc_listen_report, "006": self._handle_bind_response, "007": self._handle_open_response, "008": self._handle_sleep_schedule_response, "009": self._handle_remote_sleep_wake_response, "010": self._handle_alarm_report, "011": self._handle_short_press_message, "012": self._handle_device_identity_init, } @classmethod async def get_instance(cls, config: dict = None) -> Optional["TalkingQMQTTService"]: async with cls._lock: if cls._instance is None and config is not None: cls._instance = cls(config) return cls._instance @classmethod async def reset_instance(cls): async with cls._lock: if cls._instance is not None: await cls._instance.disconnect() cls._instance = None async def _message_loop(self): try: async for message in self._client.messages: try: topic = str(message.topic) parts = topic.split("/") if not parts or parts[0] != "device": continue payload = json.loads(message.payload.decode("utf-8")) msg_id = payload.get("msg_id") device_id = parts[1] if len(parts) >= 2 else "unknown" topic_kind = parts[2] if len(parts) >= 3 else "" if msg_id == "012" and topic_kind != "event": continue if msg_id != "012" and not device_id.startswith(f"{self.device_prefix}_"): continue handler = self._msg_handlers.get(msg_id) if handler: await handler(device_id, payload) else: logger.warning(device_id, "", f"unknown msg_id={msg_id}") except json.JSONDecodeError as exc: logger.error("", "", f"invalid mqtt payload: {exc}") except Exception as exc: logger.error("", "", f"mqtt message handling failed: {exc}") except asyncio.CancelledError: pass except Exception as exc: logger.error("", "", f"mqtt loop failed: {exc}") self._connected = False async def _schedule_persistence(self, device_id: str, label: str, coro) -> None: async def _runner(): try: await coro except Exception as exc: logger.error(device_id, "mqtt_persistence", f"{label} failed: {exc}", exc_info=True) await task_manager.create_task( _runner(), device_id=device_id, task_type="persistence", ) async def _handle_device_info(self, device_id: str, payload: dict): data = payload.get("data", {}) await self._schedule_persistence( device_id, "device_info", device_setting_service.insert_or_update( device_id=device_id, power=data.get("power"), signal_strength=data.get("signal"), version_str=data.get("version"), volume=data.get("voice"), ), ) await self._publish(f"device/{device_id}/event_resp", {"msg_id": "000", "status": "success"}) async def _handle_gps_response(self, device_id: str, payload: dict): if payload.get("status") != "success": logger.warning(device_id, "", f"[GPS] query failed: {payload}") return data = payload.get("data", {}) raw_device_time = data.get("device_time") parsed_device_time = None if isinstance(raw_device_time, str) and raw_device_time.strip(): try: parsed_device_time = datetime.fromisoformat(raw_device_time.strip()) except ValueError: parsed_device_time = None await self._schedule_persistence( device_id, "gps", location_service.report_mqtt_device_location( device_id=device_id, latitude=data.get("latitude"), longitude=data.get("longitude"), coord_type=data.get("coord_type"), accuracy_m=data.get("accuracy_m"), altitude_m=data.get("altitude_m"), speed_mps=data.get("speed_mps"), heading_deg=data.get("heading_deg"), source=data.get("source"), battery_pct=data.get("battery_pct"), device_time=parsed_device_time, ), ) async def _handle_volume_response(self, device_id: str, payload: dict): if payload.get("status") != "success": logger.warning(device_id, "", f"[volume] command failed: {payload}") return data = payload.get("data", {}) current_level = data.get("current_level") if current_level is None: return await self._schedule_persistence( device_id, "volume", device_setting_service.insert_or_update( device_id=device_id, power=None, volume=current_level, signal_strength=None, version_str=None, ), ) async def _handle_ota_response(self, device_id: str, payload: dict): status = payload.get("status") data = payload.get("data", {}) progress = data.get("progress") try: progress_value = float(progress) if progress is not None else None except (TypeError, ValueError): progress_value = None if status == "accepted": await self._schedule_persistence( device_id, "ota", device_firmware_update_manager.update_firmware_update( device_id=device_id, firmware_version=data.get("version") or data.get("new_version") or "unknown", update_status="accepted", progress=progress_value if progress_value is not None else 0.0, ), ) elif status == "updating": await self._schedule_persistence( device_id, "ota", device_firmware_update_manager.update_firmware_update( device_id=device_id, firmware_version=data.get("version") or data.get("new_version") or "updating", update_status="updating", progress=progress_value, ), ) elif status == "success": new_version = data.get("new_version") or data.get("version") await self._schedule_persistence( device_id, "ota", device_setting_service.insert_or_update( device_id=device_id, power=None, volume=None, signal_strength=None, version_str=new_version, ), ) await self._schedule_persistence( device_id, "ota", device_firmware_update_manager.update_firmware_update( device_id=device_id, firmware_version=new_version or "unknown", update_status="success", progress=100.0, ), ) elif status == "failed": await self._schedule_persistence( device_id, "ota", device_firmware_update_manager.update_firmware_update( device_id=device_id, firmware_version=data.get("version") or data.get("new_version") or "unknown", update_status="failed", progress=progress_value, ), ) logger.warning(device_id, "", f"[OTA] command failed: {payload}") else: logger.warning(device_id, "", f"[OTA] unknown status: {payload}") async def _handle_nfc_notice_response(self, device_id: str, payload: dict): status = payload.get("status") if status != "success": logger.warning(device_id, "", f"[NFC notice] command failed: {payload}") async def _handle_nfc_listen_report(self, device_id: str, payload: dict): params = payload.get("params", {}) nfc_uuid = params.get("uuid") logger.info(device_id, "", f"[NFC收听] 设备 {device_id} 请求收听留言, UUID={nfc_uuid}") await self._send_nfc_listen_response(device_id, nfc_uuid) async def _handle_bind_response(self, device_id: str, payload: dict): params = payload status = params.get("status") logger.info(device_id, "", f"[NFC绑定] 设备 {device_id} 请求绑定卡片返回状态: {status}") data = params.get("data", {}) nfc_uuid = data.get("uuid") # 卡片与设备绑定 # await card_service.activate_card(device_id=device_id, card_uuid=nfc_uuid) service = BindingService() result = await service.finalize_nfc_bind(device_id=device_id, card_uuid=nfc_uuid) if result is None: logger.warning(device_id, "", "[NFC bind] no pending bind session found") return logger.info(device_id, "", f"[NFC bind] bind completed, uuid={nfc_uuid}, result={result}") async def _handle_open_response(self, device_id: str, payload: dict): params = payload status = params.get("status") logger.info(device_id, "", f"[设备状态] 设备 {device_id} 请求返回设备状态: {status}") data = params.get("data", {}) device_status = data.get("status") logger.info(device_id, "", f"[设备状态] 设备 {device_id} 请求设备状态: {device_status}") async def _handle_sleep_schedule_response(self, device_id: str, payload: dict): status = payload.get("status") if status == "success": logger.info(device_id, "", f"[定时休眠] 设备 {device_id} 休眠时间设置成功") data = payload.get("data", {}) or {} start = data.get("start") end = data.get("end") timezone = data.get("timezone") or "Asia/Shanghai" if start and end: try: await self._schedule_persistence( device_id, "sleep_schedule", device_setting_service.upsert_sleep_schedule( device_id=device_id, sleep_mode=1, disable_time_start=datetime.strptime(start, "%H:%M").time(), disable_time_end=datetime.strptime(end, "%H:%M").time(), timezone=timezone, ), ) except ValueError: logger.warning(device_id, "", f"[定时休眠] 无法解析设备返回的时间: {payload}") else: logger.warning(device_id, "", f"[定时休眠] 设备 {device_id} 休眠时间设置失败: {payload}") async def _handle_remote_sleep_wake_response(self, device_id: str, payload: dict): status = payload.get("status") if status == "success": logger.info(device_id, "", f"[远程休眠唤醒] 设备 {device_id} 操作成功") else: logger.warning(device_id, "", f"[远程休眠唤醒] 设备 {device_id} 操作失败: {payload}") async def _handle_alarm_report(self, device_id: str, payload: dict): logger.info(device_id, "", f"[告警] 设备 {device_id} 发送紧急报警") await self._schedule_persistence( device_id, "alarm_event", device_alarm_service.record_alarm_event(device_id=device_id, source_msg_id="010"), ) await self._publish(f"device/{device_id}/event_resp", {"msg_id": "010", "status": "success"}) async def _handle_short_press_message(self, device_id: str, payload: dict): params = payload.get("params", {}) nfc_uuid = str(params.get("uuid") or "").strip() logger.info(device_id, "", f"[短按留言] 设备 {device_id} 短按发送留言, UUID={nfc_uuid}") if nfc_uuid == settings.talkingq_parent_message_uuid: await device_target_cache.set_parent_target(device_id, nfc_uuid) response_payload = { "msg_id": "011", "status": "success", "type": 0, "params": { "url": "http://175.24.73.253:8080/assets/audio/message_ok.mp3" }, } await self._publish(f"device/{device_id}/event_resp", response_payload) return media_file_key = str(params.get("media_file_key") or params.get("audio_url") or "").strip() if media_file_key: try: await im_service.create_device_parent_leave_message( device_id=device_id, media_file_key=media_file_key, media_duration_ms=params.get("media_duration_ms"), media_mime_type=params.get("media_mime_type"), media_size_bytes=params.get("media_size_bytes"), media_transcript_text=params.get("media_transcript_text"), client_msg_id=params.get("client_msg_id"), ext_json=params.get("ext_json") if isinstance(params.get("ext_json"), dict) else None, ) logger.info(device_id, "", f"[短按留言] 设备 {device_id} 留言已写入家长会话") await self._publish( f"device/{device_id}/event_resp", { "msg_id": "011", "status": "success", "type": 0, "params": { "url": "http://175.24.73.253:8080/assets/audio/message_ok.mp3" }, }, ) return except Exception as exc: logger.warning(device_id, "", f"[短按留言] 设备 {device_id} 留言写入失败: {exc}") await self._publish( f"device/{device_id}/event_resp", { "msg_id": "011", "status": "failed", "message": str(exc), }, ) return payload = {"msg_id": "011", "status": "failed", "message": "unsupported uuid"} await self._publish(f"device/{device_id}/event_resp", payload) async def _handle_device_identity_init(self, imei: str, payload: dict): try: result = await device_identity_initializer.initialize_by_imei(imei) except DeviceIdentityInitializationError as exc: logger.warning(imei, "", f"[设备初始化] IMEI初始化失败: {exc}") await self._publish( f"device/{imei}/event_resp", { "msg_id": "012", "status": "failed", "message": str(exc), }, ) return response_payload = { "msg_id": "012", "params": { "device_id": result.device_id, "device_sn": result.serial_number, }, } await self._publish(f"device/{imei}/event_resp", response_payload) async def _send_nfc_listen_response(self, device_id: str, nfc_uuid: str): topic = f"device/{device_id}/event_resp" card = await card_service.get_card_by_uuid(nfc_uuid) if not card: payload = { "msg_id": "005", "type": 0, "params": { "url_1": f"http://{settings.server_host}:{settings.server_port}/assets/audio/error_card.mp3" }, } await self._publish(topic, payload) return is_owner = await card_service.check_card_ownership(nfc_uuid, device_id) if is_owner: has_pending = await offline_audio_cache.has_pending_audio(device_id) if nfc_uuid == "53C22B6DA20001": payload = { "msg_id": "005", "type": 0, "params": { "url_1": f"http://{settings.server_host}:{settings.server_port}/assets/audio/test_zh.mp3" }, } await self._publish(topic, payload) return if has_pending: audio_urls = await offline_audio_cache.get_audio_urls(device_id) # 53D92B6DA20001 测试卡片 if nfc_uuid == "53C22B6DA20001": payload = { "msg_id": "005", "type": 0, "params": { "url_1": f"http://{settings.server_host}:{settings.server_port}/assets/audio/TalkingQ_XQSN00001005_2f654480.mp3" } } await self._publish(topic, payload) return if len(audio_urls) == 0: payload = { "msg_id": "005", "type": 0, "params": { "url_1": f"http://{settings.server_host}:{settings.server_port}/assets/audio/no_message.mp3" }, } else: params = {f"url_{k}": audio_url for k, audio_url in enumerate(audio_urls, start=1)} payload = {"msg_id": "005", "type": 0, "params": params} await offline_audio_cache.clear_audio_urls(device_id) await self._publish(topic, payload) return payload = { "msg_id": "005", "type": 0, "params": { "url_1": f"http://{settings.server_host}:{settings.server_port}/assets/audio/no_message.mp3" }, } await self._publish(topic, payload) return existing_card = await card_service.get_card_by_uuid(nfc_uuid) if existing_card: target_device_id = existing_card.device_id else: await card_service.activate_card(nfc_uuid, device_id) return await device_target_cache.set_target(device_id, target_device_id) payload = { "msg_id": "005", "type": 1, "params": { "url_1": f"http://{settings.server_host}:{settings.server_port}/assets/audio/welcome.mp3" }, } await self._publish(topic, payload) async def connect(self): async with self._connect_lock: if self._connected: return try: self._client = aiomqtt.Client( hostname=self.broker, port=self.port, username=self.username, password=self.password, identifier=f"talkingq_server_{int(time.time())}", keepalive=self.keepalive, ) await self._client.__aenter__() self._connected = True await self._client.subscribe("device/+/response", qos=self.qos) await self._client.subscribe("device/+/event", qos=self.qos) self._message_task = asyncio.create_task(self._message_loop()) logger.info("", "", f"mqtt connected: {self.broker}:{self.port}") except Exception as exc: self._connected = False logger.error("", "", f"mqtt connect failed: {exc}") async def disconnect(self): if self._message_task and not self._message_task.done(): self._message_task.cancel() try: await self._message_task except asyncio.CancelledError: pass self._message_task = None if self._client is not None: try: await self._client.__aexit__(None, None, None) except Exception: pass self._client = None self._connected = False async def _ensure_connected(self): if not self._connected: await self.connect() async def _publish(self, topic: str, payload: dict): await self._ensure_connected() try: await self._client.publish(topic, json.dumps(payload, ensure_ascii=False), qos=self.qos) logger.info("", "", f"mqtt publish topic={topic} payload={payload}") except Exception as exc: logger.error("", "", f"mqtt publish failed: {exc}") async def send_gps_query(self, device_id: str) -> str: await self._publish(f"device/{device_id}/command", {"msg_id": "001"}) return "001" async def send_volume_command(self, device_id: str, level: int) -> str: await self._publish( f"device/{device_id}/command", {"msg_id": "002", "params": {"level": level}}, ) return "002" async def send_ota_command(self, device_id: str, url: str, version: str) -> str: await self._publish( f"device/{device_id}/command", {"msg_id": "003", "params": {"url": url, "version": version}}, ) return "003" async def send_nfc_notice(self, device_id: str, url: str) -> str: topic = f"device/{device_id}/command" payload = { "msg_id": "004", "params": {"url": url} } await self._publish(topic, payload) return "004" async def send_bind_nfc_command(self, device_id: str) -> str: topic = f"device/{device_id}/command" payload = { "msg_id": "006", "params": { "url": f"http://{settings.server_host}:{settings.server_port}/assets/audio/bind_nfc_ready_zh.mp3" } } await self._publish(topic, payload) return "006" async def send_open_command(self, device_id: str, open_type: int, is_open: int) -> str: topic = f"device/{device_id}/command" if open_type not in [0, 1, 2]: raise ValueError("open_type must be 0, 1 or 2") if open_type == 0: payload = { "msg_id": "007", "type": open_type } elif open_type == 1: payload = { "msg_id": "007", "type": open_type } elif open_type == 2: payload = { "msg_id": "007", "type": open_type } await self._publish(topic, payload) return "007" async def send_sleep_schedule_command(self, device_id: str, start: str, end: str) -> str: topic = f"device/{device_id}/command" payload = { "msg_id": "008", "params": { "start": start, "end": end } } await self._publish(topic, payload) return "008" async def send_remote_sleep_wake_command(self, device_id: str, switch: str) -> str: if switch not in ["on", "off"]: raise ValueError("switch must be 'on' or 'off'") topic = f"device/{device_id}/command" payload = { "msg_id": "009", "params": { "switch": switch } } await self._publish(topic, payload) return "009"