Files
server-core/app/park/mqtt.py
T
Pine 4d4d38fe0f feat(park): 园区子应用 app/park 骨架迁入 — dispatcher /park 前缀分流
迁自 park-desktop/backend/app:config(数据目录→serverdata/park)/storage/sim_engine/
mqtt/event_bus + 精简 routers(health/settings/config/dashboard/park/companies/display/
events) + app.py 子应用(lifespan 起 MQTT hub+tick 线程, include_router prefix=/park)。
dispatcher.py 增 ROUTE_PARK_PREFIX + _prefix_app,与 app.training(/api)、app.main 分流。
pyproject 增 paho-mqtt(MQTT 必需)。AI/知识库/语音/视觉重模块延后全量迁入。
2026-08-24 16:57:27 +08:00

243 lines
9.4 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- 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": ""})
cur["last_seen"] = time.time()
self.devices[cid] = cur
self._prune()
except Exception: # noqa: BLE001
pass
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=""):
"""前端上报设备:唯一 device_id + 角色(main/secondary) + 从属主屏 parent_id。"""
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
cur["last_seen"] = time.time()
self.devices[device_id] = cur
self._prune()
log.info("mqtt: 设备注册 id=%s role=%s parent=%s", device_id, role or cur.get("role"), parent_id or cur.get("parent_id"))
def publish_command(self, action, params=None, screen_id=None, screen_role=None):
"""按客户端(设备)路由下发指令:
· 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:
target = cid
break
if target:
# 定向:只发到目标设备独有频道
return self._publish_to(payload, target, screen_role)
if requested:
# 指定了目标但未找到对应设备:不广播,避免误发到全部
log.warning("mqtt: 目标设备未找到(id=%s role=%s),未下发", screen_id, screen_role)
payload["published"] = False
with self._lock:
self.last_command = {"action": action, "params": params or {}, "ts": payload["ts"], "published": False, "topic": "none"}
return payload
# 全部:遍历在线设备列表,依次下发到各自独有频道
with self._lock:
cutoff = time.time() - settings.SCREEN_TTL
online = [cid for cid, v in self.devices.items() if v.get("last_seen", 0) > cutoff]
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", "")
if self.publish(f"{settings.TOPIC_COMMAND}/{cid}", 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=""):
"""发到指定设备独有频道,并记录 last_command。"""
topic = 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):
payload = {"ts": int(time.time() * 1000), "snapshot": snapshot}
self.publish(settings.TOPIC_TICK, payload, qos=0)
hub = MqttHub()