# -*- coding: utf-8 -*- """MQTT 发布/订阅中心 —— 后端 ↔ 前端控制通道 发布:opc/display/command(页面/媒体/卡片控制)、opc/dashboard/tick(数据快照) 订阅:opc/display/heartbeat(大屏在线心跳,统计在线大屏数) """ import json import logging import os import threading import time import uuid import paho.mqtt.client as mqtt from .config import settings from .event_bus import bus log = logging.getLogger("dpm.mqtt") def new_cmd_id(): return uuid.uuid4().hex[:12] def _unique_client_id(base): """每个客户端进程使用唯一 client_id,避免与其它实例竞争被 broker 踢下线""" return f"{base}-{os.getpid()}-{uuid.uuid4().hex[:6]}" class MqttHub: def __init__(self): self.client = None self.connected = False self._lock = threading.RLock() self.devices = {} # device_id(client_id) -> {role, parent_id, last_seen} self.last_command = None # {action, params, ts, published} # ---------- 生命周期 ---------- def start(self): if not settings.MQTT_ENABLED: log.info("MQTT 已禁用(DPM_MQTT_ENABLED=0)") return try: self.client = mqtt.Client( client_id=_unique_client_id(settings.MQTT_CLIENT_ID), clean_session=True, protocol=mqtt.MQTTv311, ) if settings.MQTT_USERNAME: self.client.username_pw_set(settings.MQTT_USERNAME, settings.MQTT_PASSWORD) self.client.on_connect = self._on_connect self.client.on_disconnect = self._on_disconnect self.client.on_message = self._on_message self.client.connect_async(settings.MQTT_HOST, settings.MQTT_PORT, keepalive=30) self.client.loop_start() log.info("MQTT 连接中 %s:%s ...", settings.MQTT_HOST, settings.MQTT_PORT) except Exception as e: # noqa: BLE001 log.warning("MQTT 初始化失败(数据 API 不受影响): %s", e) def stop(self): if self.client: try: self.client.loop_stop() self.client.disconnect() except Exception: # noqa: BLE001 pass self.client = None def _on_connect(self, client, userdata, flags, rc): if rc == 0: self.connected = True log.info("MQTT 已连接 %s:%s", settings.MQTT_HOST, settings.MQTT_PORT) client.subscribe(settings.TOPIC_HEARTBEAT, qos=0) else: log.warning("MQTT 连接失败 rc=%s", rc) def _on_disconnect(self, client, userdata, rc): self.connected = False if rc != 0: log.warning("MQTT 断开(rc=%s),自动重连中...", rc) def _on_message(self, client, userdata, msg): """接收大屏心跳:更新设备表在线时间(设备由注册表登记 role/parent)""" if msg.topic == settings.TOPIC_HEARTBEAT: try: payload = json.loads(msg.payload.decode("utf-8")) cid = payload.get("client_id") or msg.topic with self._lock: cur = self.devices.get(cid, {"role": "", "parent_id": "", "tenant_id": ""}) cur["last_seen"] = time.time() if payload.get("tenant_id"): cur["tenant_id"] = payload["tenant_id"] self.devices[cid] = cur self._prune() except Exception: # noqa: BLE001 pass def _online_ids(self, tenant_id: str = "") -> list[str]: with self._lock: cutoff = time.time() - settings.SCREEN_TTL return [ cid for cid, v in self.devices.items() if v.get("last_seen", 0) > cutoff and (not tenant_id or v.get("tenant_id") == tenant_id) ] def _prune(self): """剔除超时未上报的设备""" cutoff = time.time() - settings.SCREEN_TTL self.devices = {k: v for k, v in self.devices.items() if v.get("last_seen", 0) > cutoff} # ---------- 状态查询 ---------- def screens_online(self): with self._lock: return sum(1 for v in self.devices.values() if v.get("last_seen", 0) > time.time() - settings.SCREEN_TTL) def status(self): with self._lock: self._prune() devices = [ { "device_id": cid, "role": v.get("role", ""), "parent_id": v.get("parent_id", ""), } for cid, v in self.devices.items() ] return { "mqtt_connected": self.connected, "mqtt_host": f"{settings.MQTT_HOST}:{settings.MQTT_PORT}", "screens_online": len(devices), "devices": devices, "last_command": self.last_command, } # ---------- 发布 ---------- def publish(self, topic, payload, qos=1, retain=False): if not self.client or not self.connected: log.debug("MQTT 未连接,丢弃发布 %s", topic) return False try: info = self.client.publish(topic, json.dumps(payload, ensure_ascii=False), qos=qos, retain=retain) return info.rc == mqtt.MQTT_ERR_SUCCESS except Exception as e: # noqa: BLE001 log.warning("MQTT 发布失败: %s", e) return False def register_device(self, device_id, role="", parent_id="", tenant_id=""): """前端上报设备:唯一 device_id + 角色(main/secondary) + 从属主屏 parent_id + 所属园区 tenant。""" device_id = (device_id or "").strip() if not device_id: return with self._lock: cur = self.devices.get(device_id, {}) if role: cur["role"] = role if parent_id: cur["parent_id"] = parent_id if tenant_id: cur["tenant_id"] = tenant_id cur["last_seen"] = time.time() self.devices[device_id] = cur self._prune() log.info("mqtt: 设备注册 id=%s role=%s parent=%s tenant=%s", device_id, role or cur.get("role"), parent_id or cur.get("parent_id"), tenant_id) def publish_command(self, action, params=None, screen_id=None, screen_role=None, tenant_id=""): """按客户端(设备)路由下发指令(多租户:只遍历该租户在线设备): · screen_id 精确到某台设备 → 只发到该设备独有 topic · screen_role(main/secondary) → 解析到该角色的设备,发到其独有 topic · 两者皆空(全部)→ 遍历该租户在线设备,依次下发到各自独有 topic 客户端只订阅自己的独有频道,因此定向指令绝不会被其它屏幕收到。""" payload = { "cmd_id": new_cmd_id(), "ts": int(time.time() * 1000), "action": action, "params": params or {}, } target = None requested = bool(screen_id or screen_role in ("main", "secondary")) if screen_id: target = (screen_id or "").strip() elif screen_role in ("main", "secondary"): with self._lock: for cid, v in self.devices.items(): if v.get("role") == screen_role and (not tenant_id or v.get("tenant_id") == tenant_id): target = cid break if target: # 定向:只发到目标设备独有频道 return self._publish_to(payload, target, screen_role, tenant_id) if requested: # 指定了目标但未找到对应设备:不广播,避免误发到全部 log.warning("mqtt: 目标设备未找到(id=%s role=%s tenant=%s),未下发", screen_id, screen_role, tenant_id) payload["published"] = False with self._lock: self.last_command = {"action": action, "params": params or {}, "ts": payload["ts"], "published": False, "topic": "none"} return payload # 全部:遍历该租户在线设备,依次下发到各自独有频道 online = self._online_ids(tenant_id) ok = True published_topics = [] for cid in online: p = dict(payload) p["screen_id"] = cid with self._lock: p["screen_role"] = self.devices.get(cid, {}).get("role", "") topic = f"{settings.TOPIC_COMMAND}/{tenant_id}/{cid}" if tenant_id else f"{settings.TOPIC_COMMAND}/{cid}" if self.publish(topic, p): bus.emit(p) published_topics.append(cid) else: ok = False payload["published"] = ok payload["screen_id"] = "" payload["screen_role"] = "both" with self._lock: self.last_command = { "action": action, "params": params or {}, "ts": payload["ts"], "published": ok, "topic": "all", } log.info("mqtt: 全部下发 %s 到 %d 台设备", action, len(published_topics)) return payload def _publish_to(self, payload, target, screen_role="", tenant_id=""): """发到指定设备独有频道,并记录 last_command。""" topic = f"{settings.TOPIC_COMMAND}/{tenant_id}/{target}" if tenant_id else f"{settings.TOPIC_COMMAND}/{target}" payload["screen_id"] = target with self._lock: payload["screen_role"] = self.devices.get(target, {}).get("role", screen_role or "") ok = self.publish(topic, payload) payload["published"] = ok with self._lock: self.last_command = { "action": payload["action"], "params": payload.get("params") or {}, "ts": payload["ts"], "published": ok, "topic": topic, } if ok: bus.emit(payload) return payload def publish_tick(self, snapshot, tenant_id=""): payload = {"ts": int(time.time() * 1000), "snapshot": snapshot} topic = f"{settings.TOPIC_TICK}/{tenant_id}" if tenant_id else settings.TOPIC_TICK self.publish(topic, payload, qos=0) hub = MqttHub()