Files
server-core/app/park/mqtt.py
T
Pine 81f4dfd0ad feat(park): 园区多租户 — SQLite(tenants表+companies+kb) + 每租户引擎 + 认证 + 作用域隔离
- tenants.py:park_tenants/park_companies/park_kb_docs 三表(serverdata/park/park.db),
  每园区独立 name/intro(两段式)/auth(账密 salted hash)/data(大屏锚点)/agent/kb;
  ensure_default_tenant 自愈建默认园区+灌 39 家企业。
- sim_engine 重构为 SimEngine(data)——每租户独立实例,get_engine(tid)/refresh_engine(tid)。
- auth.py:/park/auth/login(账密→长效 tenant JWT)、/park/auth/me、require_tenant(401)。
- mqtt:register/publish_command/publish_tick 按 tenant_id 作用域(拓扑为 command/tid/clientId)。
- routers:/park/tenants CRUD + 全部数据/指令端点租户作用域(authorization token 或 ?tenant_id)。
- app.py lifespan init_db + 逐租户 tick/publish;tools._query_companies 读租户库。
- .gitignore 增 serverdata/park/。TestClient:建租户/登录/快照隔离/企业CRUD/指令鉴权(无token401) 全过。
2026-08-24 18:06:51 +08:00

255 lines
10 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": "", "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()