feat: 实现设备注册和精准控制功能,支持双屏管理;更新相关逻辑以优化设备状态跟踪
This commit is contained in:
+106
-14
@@ -33,7 +33,7 @@ class MqttHub:
|
||||
self.client = None
|
||||
self.connected = False
|
||||
self._lock = threading.RLock()
|
||||
self._screens = {} # client_id -> last_seen_ts
|
||||
self.devices = {} # device_id(client_id) -> {role, parent_id, last_seen}
|
||||
self.last_command = None # {action, params, ts, published}
|
||||
|
||||
# ---------- 生命周期 ----------
|
||||
@@ -81,30 +81,45 @@ class MqttHub:
|
||||
log.warning("MQTT 断开(rc=%s),自动重连中...", rc)
|
||||
|
||||
def _on_message(self, client, userdata, msg):
|
||||
"""接收大屏心跳:记录 client_id 与时间"""
|
||||
"""接收大屏心跳:更新设备表在线时间(设备由注册表登记 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:
|
||||
self._screens[cid] = time.time()
|
||||
# 清理过期
|
||||
cutoff = time.time() - settings.SCREEN_TTL
|
||||
self._screens = {k: v for k, v in self._screens.items() if v > cutoff}
|
||||
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:
|
||||
cutoff = time.time() - settings.SCREEN_TTL
|
||||
return sum(1 for v in self._screens.values() if v > cutoff)
|
||||
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": self.screens_online(),
|
||||
"screens_online": len(devices),
|
||||
"devices": devices,
|
||||
"last_command": self.last_command,
|
||||
}
|
||||
|
||||
@@ -120,26 +135,103 @@ class MqttHub:
|
||||
log.warning("MQTT 发布失败: %s", e)
|
||||
return False
|
||||
|
||||
def publish_command(self, action, params=None):
|
||||
"""统一指令信封:{cmd_id, ts, action, params, published}
|
||||
MQTT 广播 + SSE 兼容推送;published 标记是否真正发布成功"""
|
||||
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 {},
|
||||
}
|
||||
ok = self.publish(settings.TOPIC_COMMAND, payload)
|
||||
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) # SSE 兼容通道(仅发布成功时推送)
|
||||
bus.emit(payload)
|
||||
return payload
|
||||
|
||||
def publish_tick(self, snapshot):
|
||||
|
||||
+18
-1
@@ -73,6 +73,8 @@ class StateBody(BaseModel):
|
||||
class DisplayCommandBody(BaseModel):
|
||||
action: str
|
||||
params: dict = Field(default_factory=dict)
|
||||
screen_id: str = "" # 目标屏幕唯一 client_id(精确到某台)
|
||||
screen_role: str = "" # main | secondary(按角色定位)
|
||||
|
||||
|
||||
class ChatMessage(BaseModel):
|
||||
@@ -487,13 +489,28 @@ _VALID_DISPLAY_ACTIONS = {
|
||||
}
|
||||
|
||||
|
||||
class RegisterBody(BaseModel):
|
||||
device_id: str = "" # 唯一设备标识(MQTT client_id)
|
||||
role: str = "" # main | secondary
|
||||
parent_id: str = "" # 副屏从属的主屏 device_id;主屏为空
|
||||
|
||||
|
||||
@router.post("/api/display/register")
|
||||
async def display_register(body: RegisterBody):
|
||||
"""设备注册:前端上报唯一 device_id + 角色 + 从属主屏,后端维护设备表供精准控制。"""
|
||||
hub.register_device(body.device_id, body.role, body.parent_id)
|
||||
return {"ok": True, "device_id": body.device_id, "role": body.role, "parent_id": body.parent_id}
|
||||
|
||||
|
||||
@router.post("/api/display/command")
|
||||
async def display_command_publish(body: DisplayCommandBody):
|
||||
"""管理端/任意客户端通过 REST 发控制指令 → 后端转 MQTT 广播给所有大屏"""
|
||||
if body.action not in _VALID_DISPLAY_ACTIONS:
|
||||
return JSONResponse({"ok": False, "error": f"不支持的指令: {body.action}"}, status_code=400)
|
||||
cmd = hub.publish_command(body.action, body.params)
|
||||
# 按屏幕终端路由:screen_id 精确到某客户端;screen_role 按主/副定位;皆空 → 全局广播
|
||||
cmd = hub.publish_command(body.action, body.params, body.screen_id or None, body.screen_role or None)
|
||||
return {"ok": cmd.get("published", False), "cmd_id": cmd["cmd_id"], "action": body.action,
|
||||
"screen_id": body.screen_id or "", "screen_role": body.screen_role or "both",
|
||||
"mqtt_connected": hub.connected}
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user