提交 03272305 authored 作者: 陈泽健's avatar 陈泽健

feat(device-sim): P0~P3 全链路完成 + 部署三台 + HANDOFF 记录 + 同步 P0 修复回源码

- 门口屏模拟器 P0~P3 完成(MQTT 修复 + 动态 X-SIGN + result 解密 + MQTT-HTTP 联动)
- 部署到 5.60/5.44/5.202,容器内真实 Broker 联动验证 PASS
- 同步 P0 修复(mqtt_manager.py)回源码树
Co-Authored-By: 's avatarClaude <noreply@anthropic.com>
上级 a761d9ce
# HANDOFF — 设备模拟模块 # HANDOFF — 设备模拟模块
> **生成时间**: 2026-08-24 > ⚠️ **多窗口并行开发注意**:如需开多个 Claude Code 窗口同时开发,必须先阅读 `Docs/多窗口并行开发指南.md`(端口隔离 / 数据库错开 / Playwright 互斥 / Git 协作)。本仓库当前确有并行窗口在工作(性能测试、安全测试等模块的未提交改动与本模块无关,勿混入本模块提交)。
> **生成时间**: 2026-09-10
> **当前分支**: `platform-auto-test` > **当前分支**: `platform-auto-test`
> **最近提交**: `21d65661` docs(device-sim): 更新 HANDOFF 记录设备模拟模块部署到 5.60 > **最近提交**: `beb0f4af` fix(security): 修复5.60安全用例JSON双重编码500错误 + 部署脚本防复发(并行窗口提交)
> **未提交改动**: 仅 `backend/data/test_platform.db`(本地 SQLite,已 gitignore)+ 新增部署/验证脚本(`check_mounts_44_202.py` / `deploy_44_202.py` / `verify_44_202.py`,均按 .gitignore 规则不提交) > **本模块未提交改动**(设备模拟 P0~P3,均已部署到三台服务器,代码待提交):
> - `backend/app/services/mqtt_manager.py` — P0:connect() 复用存活连接 + `_teardown`(暂存副本 `scripts/mqtt_manager_p0_fix.py`)
> - `backend/app/simulators/door_token_client.py` — P2:动态 X-SIGN + result 解密(暂存副本 `scripts/door_token_client_p3.py`)
> - `backend/app/simulators/door_simulator.py` — P3:MQTT-HTTP 联动(暂存副本 `scripts/door_simulator_p3.py`)
> - 辅助脚本:`deploy_mqtt_fix.py` / `deploy_p3_linkage.py` / `p3_linkage_verify.py` / `diag_544_mqtt_broker.py` / `diag_544_mqtt_deep.py` / `_remote_54x_*.py`(远端回传校验副本)
> - 其余未提交文件(性能/安全/项目模块等)属并行窗口,与本模块无关
--- ---
## 📊 会话进度记录 ## 📊 会话进度记录
### 2026-09-10 会话 O:门口屏 P0~P3 全链路完成(MQTT 修复 + 动态 X-SIGN + result 解密 + MQTT-HTTP 联动,已部署三台并验证通过)
**会话目标**:按优先级完成门口屏模拟器 P0~P3 四项工作,并同步部署到 5.60 / 5.44 / 5.202。
**状态**:✅ 全部完成 + 三台部署成功 + 容器内真实 Broker 联动验证 PASS
---
#### P0:5.44 MQTT Broker 闪断(rc=7)修复
- 根因:模拟器批量断开/重建连接时 `mqtt_manager.connect()` 无条件先 disconnect 已有连接,导致同环境所有设备断流 + Broker 侧连接风暴。
- 修复:`mqtt_manager.py` connect() 已有存活连接时直接复用;增加 `_teardown` 清理路径。已按 `deploy_mqtt_fix.py` 部署三台。
#### P1:5.44 / 5.202 补传文件
- 两台缺失 `door_token_client.py` / `paperless_simulator.py`(历史部署遗漏),已补传并容器内 import 验证。
#### P2:动态 X-SIGN + 消息列表 result 解密(`door_token_client.py`)
- **动态 X-SIGN**(逆向前端 JS,与 http_client.py 算法一致):`x_random(8-16位)` + `x_timestamp(秒级)``sign_str = timestamp + JSON(params) + random`(无参数时 `timestamp + random`)→ `SHA256` → 以 `SHA256(bearer_token)` 为密钥源派生 `key=hex[16:32]``iv=hex[0:8]+hex[-8:]` 做 AES-CBC 加密;无 token 场景走 RandomCode 登录式签名(60 位随机串 + 随头回传)。
- **result 解密**(响应拦截器逆向):服务端返回 Base64 密文时,`AES-CFB` 解密,密钥派生同上;**关键:pycryptodome 必须显式 `segment_size=128`**(CryptoJS CFB 为 CFB-128,默认 8 解出乱码),ZeroPadding 去尾(`rstrip(b"\x00")`);失败时与前端 catch 同语义返回原数据不抛异常。
- 已在 5.44 `/api/company/getSystemConfigData` 真实服务端实测闭环(RandomCode 签名 → 服务端加密 → 解出明文 JSON)。解密数据缓存于 `last_message_list_data` / `last_face_page_data` / `last_global_config_data`
- ⚠️ 5.48(门口屏 /exapi 后端)不可达,消息列表真实明文样本仍缺失;解密能力已在 5.44 同源接口验证。
#### P3:消息列表结果与 MQTT 主题联动(`door_simulator.py`)
- **MQTT → HTTP**`_on_command` 收到 `meeting_message` 主题(`/meeting/message/{room_id}/`)推送时,立即起线程调用一次消息列表(模拟真实门口屏"推送触发拉取")。
- **HTTP → MQTT**:拉取成功且解密出消息时,向 `android_broadcast``/uams/android/broadcast`)发布 `message_sync` 摘要(`messageCount` + `latest`,防御性解析列表型字段)。
- `start()` 顺序调整:先取 token 回填 topic_params → 重解析主题 → MQTT 注册订阅 → 启动消息列表定时轮询。
#### 🔥 本会话新踩坑(已修复)
- **`_on_command` 内 `threading.Thread` 报 `name 'threading' is not defined`**:P3 新代码使用了 `threading` 但模块未 import,路径 A(推送触发拉取)首次容器内验证即崩。**教训:给既有大文件加并发代码时先核对 import 区**。修复后重部署三台。
#### 验证结果(5.60 容器内,真实 Broker 192.168.5.44:1883)
- 路径 A:向 meeting_message 推送 → `_on_command` 触发即时拉取 → `message_sync` 回发 ✅
- 路径 B:直接 `_call_message_list_once()``message_sync` 回发,`messageCount=2``latest` 内容正确 ✅
- 本地 CFB-128 解密回环(RandomCode 派生密钥 + ZeroPadding)✅
- 三台 health 检查均 200(5.60:80,5.44/5.202:8081)✅
- 验证脚本:`backend/scripts/p3_linkage_verify.py`;部署脚本:`backend/scripts/deploy_p3_linkage.py`(含标记校验 + 备份 + 健康检查,单台重跑:`python scripts/deploy_p3_linkage.py 560|544|5202`
#### 待办(下会话)
- 5.48 可达后:用真实消息列表密文回归 `decrypt_result` + 真实 `getMsgPageList` 全链路。
- 门口屏设备当前全部为 stopped 状态(容器重启后内存态重置,历史上报时间 08-21),如需联调在设备模拟页批量启动即可。
### 2026-08-24 会话 N:测试管理平台前后端部署到 5.44 / 5.202(已部署并验证通过) ### 2026-08-24 会话 N:测试管理平台前后端部署到 5.44 / 5.202(已部署并验证通过)
**会话目标**:将当前测试管理平台的前后端代码更新部署到 192.168.5.44 和 192.168.5.202 两台测试管理平台。 **会话目标**:将当前测试管理平台的前后端代码更新部署到 192.168.5.44 和 192.168.5.202 两台测试管理平台。
......
...@@ -13,6 +13,7 @@ from __future__ import annotations ...@@ -13,6 +13,7 @@ from __future__ import annotations
import json import json
import logging import logging
import random
import threading import threading
import time import time
from datetime import datetime from datetime import datetime
...@@ -138,14 +139,30 @@ class MqttManager: ...@@ -138,14 +139,30 @@ class MqttManager:
logger.info(f"MQTT 复用已有连接: env_config_id={env_config_id}") logger.info(f"MQTT 复用已有连接: env_config_id={env_config_id}")
return True return True
# 存在但已断开的连接,先清理 # 存在但已断开的连接:优先等待 paho 内置自动重连成功(最多 2 秒)。
if existing: # 不要立即 disconnect+重建 —— 那会与 paho 自动重连竞态,产生两个
self.disconnect(env_config_id) # 并发连接互踢,触发 EMQX session takeover 死循环(rc=7 闪断根因)。
if existing and existing.client:
for _ in range(20):
if existing.connected:
logger.info(
f"MQTT 自动重连成功,继续复用: env_config_id={env_config_id}")
return True
time.sleep(0.1)
logger.warning(
f"MQTT 自动重连未恢复,重建连接: env_config_id={env_config_id}")
self._teardown(existing)
try: try:
conn = MqttConnection(env_config_id) conn = MqttConnection(env_config_id)
# 唯一化 client_id:同一环境每次物理连接/重连使用不同 id,
# 杜绝相同 client_id 触发 EMQX session takeover 互踢死循环
# (P0 闪断根因修复,详见 backend/scripts/diag_544_takeover.py)。
base_cid = client_id or f"sim_{env_config_id[:8]}"
unique_cid = (f"{base_cid[:48]}_{int(time.time() * 1000) % 100000000}"
f"_{random.randint(100, 999)}")
client = mqtt.Client( client = mqtt.Client(
client_id=client_id or f"sim_{env_config_id[:8]}", client_id=unique_cid,
protocol=mqtt.MQTTv311 protocol=mqtt.MQTTv311
) )
...@@ -161,8 +178,9 @@ class MqttManager: ...@@ -161,8 +178,9 @@ class MqttManager:
if use_tls: if use_tls:
client.tls_set() client.tls_set()
# 设置自动重连(paho-mqtt 内置机制) # 设置自动重连(paho-mqtt 内置机制)。上限 5 秒保证断开后快速恢复,
client.reconnect_delay_set(min_delay=1, max_delay=60) # 让 connect() 的"等待自动重连"窗口(2 秒内)更容易命中复用路径
client.reconnect_delay_set(min_delay=1, max_delay=5)
# 连接 # 连接
client.connect(host, port, keepalive=60) client.connect(host, port, keepalive=60)
...@@ -177,7 +195,7 @@ class MqttManager: ...@@ -177,7 +195,7 @@ class MqttManager:
self._connections[env_config_id] = conn self._connections[env_config_id] = conn
if conn.connected: if conn.connected:
logger.info(f"MQTT 连接成功: {host}:{port}") logger.info(f"MQTT 连接成功: {host}:{port}, cid={unique_cid}")
return True return True
else: else:
logger.warning(f"MQTT 连接超时: {host}:{port}") logger.warning(f"MQTT 连接超时: {host}:{port}")
...@@ -187,6 +205,27 @@ class MqttManager: ...@@ -187,6 +205,27 @@ class MqttManager:
logger.error(f"MQTT 连接异常: {host}:{port}, error={e}") logger.error(f"MQTT 连接异常: {host}:{port}, error={e}")
return False return False
@staticmethod
def _teardown(conn: MqttConnection) -> None:
"""彻底关闭一个 MQTT 连接的底层网络线程。
loop_stop 会等待网络线程退出,期间把 paho 的自动重连关掉,
避免旧连接仍在自动重连、与应用层新建连接竞态互踢。
"""
client = conn.client
if not client:
return
try:
client.loop_stop()
except Exception as e:
logger.warning(f"MQTT teardown loop_stop 异常: {e}")
try:
client.disconnect()
except Exception as e:
logger.warning(f"MQTT teardown disconnect 异常: {e}")
conn.client = None
conn.connected = False
def disconnect(self, env_config_id: str) -> None: def disconnect(self, env_config_id: str) -> None:
""" """
断开 MQTT 连接 断开 MQTT 连接
...@@ -196,13 +235,10 @@ class MqttManager: ...@@ -196,13 +235,10 @@ class MqttManager:
""" """
with self._manager_lock: with self._manager_lock:
conn = self._connections.pop(env_config_id, None) conn = self._connections.pop(env_config_id, None)
if conn and conn.client: if conn:
try: if conn.client:
conn.client.loop_stop()
conn.client.disconnect()
logger.info(f"MQTT 断开连接: env_config_id={env_config_id}") logger.info(f"MQTT 断开连接: env_config_id={env_config_id}")
except Exception as e: self._teardown(conn)
logger.warning(f"MQTT 断开异常: {e}")
def publish(self, env_config_id: str, topic: str, payload: dict) -> bool: def publish(self, env_config_id: str, topic: str, payload: dict) -> bool:
""" """
......
...@@ -12,6 +12,7 @@ ...@@ -12,6 +12,7 @@
import json import json
import logging import logging
import random import random
import threading
import time import time
from typing import Optional from typing import Optional
...@@ -199,6 +200,21 @@ class DoorSimulator(BaseSimulator): ...@@ -199,6 +200,21 @@ class DoorSimulator(BaseSimulator):
command = payload.get("command", "") command = payload.get("command", "")
logger.info(f"门口屏收到指令: device_id={self.device_id}, command={command}") logger.info(f"门口屏收到指令: device_id={self.device_id}, command={command}")
# P3 联动(MQTT → HTTP):平台向会议消息主题推送新消息时,真实门口屏
# 会立即拉取一次消息列表获取完整内容,此处模拟该"推送触发拉取"行为。
meeting_msg_topic = self.get_resolved_topic("meeting_message")
if (topic and meeting_msg_topic and topic == meeting_msg_topic
and self._token_client and self._authorization):
logger.info(
f"[联动] 收到会议消息主题推送,触发即时消息拉取: "
f"device_id={self.device_id}, topic={topic}"
)
threading.Thread(
target=self._call_message_list_once,
daemon=True,
name=f"msgpush_{self.device_id}",
).start()
# 记录接收到的平台下发指令(direction=subscribe) # 记录接收到的平台下发指令(direction=subscribe)
command_topic = f"{self.get_topic_prefix()}/{self.device_id}/command" command_topic = f"{self.get_topic_prefix()}/{self.device_id}/command"
self._notify_report( self._notify_report(
...@@ -413,16 +429,24 @@ class DoorSimulator(BaseSimulator): ...@@ -413,16 +429,24 @@ class DoorSimulator(BaseSimulator):
authorization=self._authorization, authorization=self._authorization,
conference_id=conference_id, conference_id=conference_id,
) )
summary = self._summarize_decrypted_messages(
getattr(self._token_client, "last_message_list_data", None)
)
self._notify_report( self._notify_report(
topic=f"message_list_{'success' if success else 'failed'}", topic=f"message_list_{'success' if success else 'failed'}",
payload={ payload={
"device_id": self.device_id, "device_id": self.device_id,
"conference_id": conference_id, "conference_id": conference_id,
"success": success, "success": success,
"summary": summary or None,
}, },
direction="publish", direction="publish",
status="success" if success else "failed", status="success" if success else "failed",
) )
# P3 联动(HTTP → MQTT):拉取成功且解密出消息时,向真实上报主题
# 发布 message_sync 摘要(模拟设备同步到新消息后的上报行为)
if success:
self._publish_message_sync()
return success return success
except Exception as e: except Exception as e:
logger.warning(f"[MsgList] 调用异常: device_id={self.device_id}, error={e}") logger.warning(f"[MsgList] 调用异常: device_id={self.device_id}, error={e}")
...@@ -435,6 +459,80 @@ class DoorSimulator(BaseSimulator): ...@@ -435,6 +459,80 @@ class DoorSimulator(BaseSimulator):
) )
return False return False
@staticmethod
def _summarize_decrypted_messages(data) -> dict:
"""
从解密后的消息列表数据中提取摘要(HTTP → MQTT 联动用)
真实消息列表解密后的结构未知(5.48 不可达未取得明文样本),
此处做防御性解析:识别 dict 中的列表型字段或顶层 list,
提取消息条数与最新一条的标题/时间类字段。
"""
if data is None:
return {}
items = None
if isinstance(data, list):
items = data
elif isinstance(data, dict):
# 常见分页结构:list/records/rows/data 等键存消息数组
for key in ("list", "records", "rows", "data", "msgs", "messages"):
if isinstance(data.get(key), list):
items = data[key]
break
if items is None:
# 兜底:取第一个列表型字段
for v in data.values():
if isinstance(v, list):
items = v
break
if not isinstance(items, list):
return {}
summary = {"message_count": len(items)}
if items and isinstance(items[0], dict):
first = items[0]
for key in ("title", "content", "createTime", "sendTime", "messageName"):
if first.get(key):
summary["latest"] = str(first[key])[:100]
break
return summary
def _publish_message_sync(self) -> None:
"""将消息列表解密结果摘要发布到真实 MQTT 上报主题(HTTP → MQTT 联动)"""
try:
data = getattr(self._token_client, "last_message_list_data", None)
summary = self._summarize_decrypted_messages(data)
if not summary.get("message_count"):
return
# 优先 android_broadcast(门口屏心跳同主题),其次第一个 publish 主题
topic = self.get_resolved_topic("android_broadcast")
if not topic:
publish_topics = self.get_publish_topics()
topic = publish_topics[0]["topic"] if publish_topics else None
if not topic:
return
payload = {
"type": "message_sync",
"clientId": self.device_id,
"messageCount": summary["message_count"],
}
if summary.get("latest"):
payload["latest"] = summary["latest"]
ok = self.mqtt.publish(self.env_config_id, topic, payload)
self._notify_report(
topic=topic,
payload=payload,
direction="publish",
status="success" if ok else "failed",
)
logger.info(
f"[联动] message_sync 已发布: device_id={self.device_id}, "
f"topic={topic}, count={summary['message_count']}"
)
except Exception as e:
logger.warning(
f"[联动] message_sync 发布异常: device_id={self.device_id}, error={e}"
)
def _message_list_loop(self) -> None: def _message_list_loop(self) -> None:
""" """
消息列表 + 人脸页面 + 全局配置定时拉取循环 消息列表 + 人脸页面 + 全局配置定时拉取循环
......
Markdown 格式
0%
您添加了 0 到此讨论。请谨慎行事。
请先完成此评论的编辑!
注册 或者 后发表评论