856ff88440
后端(backend/,FastAPI :10085): - 全部数据 REST 接口:dashboard 快照 / 企业分区 / 播放列表 / 媒体 / 设置 - MQTT 控制通道:opc/display/command(切页/播放/卡片/通知), 管理端 POST /api/display/command → MQTT 广播 → 所有大屏同步响应 - 阿里云 DashScope:通义千问 LLM + 函数调用(工具经 MQTT 广播), paraformer-realtime-v2 语音识别 /api/ai/asr - 媒体资源统一由后端存储返回(上传/列表/静态服务) - SSE /api/events 保留作 MQTT 不可用时的兼容回退 前端: - config.js + .env.local 配置后端地址与 MQTT 账号(前端 dpm / 服务端 dpmserver) - mqtt.js 客户端 + useMqttControl(MQTT 驱动切页/媒体/卡片/通知) - DpmOverlays 全局覆盖层(通知 toast + 企业/分区/总览卡片) - useParkSim 改为后端 API 数据源(离线回退本地模拟) - AiChatPanel:对话走后端 LLM(工具调用),语音走本地录音 + 后端 ASR - MediaScreen:媒体控制走 MQTT(保留 SSE 回退),修复 useEffect TDZ Rust:lib.rs 移除内嵌 HTTP 服务器,只保留薄壳(窗口/权限/自启/npc 隧道) 安全:backend/.env、.env.local、media、data.json 已 gitignore
106 lines
3.5 KiB
Python
106 lines
3.5 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""MQTT 发布中心 —— 后端 → 前端控制通道
|
||
页面控制 / 媒体控制 / 卡片展示指令统一走 opc/display/command
|
||
数据快照按 DPM_MQTT_TICK_INTERVAL 推送 opc/dashboard/tick
|
||
"""
|
||
|
||
import json
|
||
import logging
|
||
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]
|
||
|
||
|
||
class MqttHub:
|
||
def __init__(self):
|
||
self.client = None
|
||
self.connected = False
|
||
self._lock = threading.RLock()
|
||
|
||
# ---------- 生命周期 ----------
|
||
def start(self):
|
||
if not settings.MQTT_ENABLED:
|
||
log.info("MQTT 已禁用(DPM_MQTT_ENABLED=0)")
|
||
return
|
||
try:
|
||
self.client = mqtt.Client(
|
||
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.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)
|
||
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 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 publish_command(self, action, params=None):
|
||
"""统一指令信封:{cmd_id, ts, action, params, published}
|
||
MQTT 广播 + SSE 兼容推送;published 标记是否真正发布成功"""
|
||
payload = {
|
||
"cmd_id": new_cmd_id(),
|
||
"ts": int(time.time() * 1000),
|
||
"action": action,
|
||
"params": params or {},
|
||
}
|
||
ok = self.publish(settings.TOPIC_COMMAND, payload)
|
||
payload["published"] = ok
|
||
if ok:
|
||
bus.emit(payload) # SSE 兼容通道(仅发布成功时推送)
|
||
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()
|