feat(DBNMQTTool): 模拟上报 + 协议Topic + 自定义Topic

新增三个 Tab:
1. 模拟上报 — 可编辑 loop_data/event_report/heartbeat JSON,
   支持单次发送和周期上报(间隔可调)
2. 协议Topic — 按协议文档列出所有 dld960/{sn}/srv/... 和 dev/... topic,
   双击填充到发布区,支持编辑 JSON 载荷
3. 自定义Topic — 自由输入 topic 发布,支持订阅/取消订阅,
   接收消息实时展示

新增:
- QComboBox, QSpinBox, QPlainTextEdit 控件
- paho.mqtt.client 导入(topic_matches_sub 通配符匹配)
- 自定义订阅消息自动路由到接收区
This commit is contained in:
wangfq
2026-07-06 21:40:26 +08:00
parent 5210051f19
commit c19e465284
+434 -2
View File
@@ -17,9 +17,11 @@ from PySide6.QtWidgets import (
QLabel, QLineEdit, QPushButton, QTreeWidget, QTreeWidgetItem,
QTabWidget, QTextEdit, QGroupBox, QGridLayout, QCheckBox,
QSplitter, QMessageBox, QHeaderView, QFrame,
QComboBox, QSpinBox, QPlainTextEdit,
)
from PySide6.QtCore import Qt, QTimer, Signal, Slot
from PySide6.QtGui import QFont, QColor, QTextCursor
import paho.mqtt.client as mqtt
from dbn_mqtt_tool.mqtt_client import MqttClient, BrokerConfig
from dbn_mqtt_tool.device_manager import DeviceManager
@@ -28,12 +30,15 @@ from dbn_mqtt_tool.protocol import (
CMD_IOT_TOPIC_QUERY, CMD_LOOP_PARAM_QUERY,
CMD_SSC_NET_SET, CMD_IOT_NET_SET, CMD_IOT_TOPIC_SET,
CMD_PWD_VERIFY, CMD_PWD_SET, CMD_FACTORY_RESET, CMD_DEVICE_RESET,
ERROR_MSGS,
CMD_LOOP_DATA, CMD_EVENT_REPORT, CMD_HEARTBEAT,
ERROR_MSGS, FREQ_LEVELS, OUTPUT_MODES, EVENT_TYPES,
data_ssc_net_set, data_iot_net_set, data_iot_topic_set,
data_pwd_verify, data_pwd_set,
build_request, next_msg_id,
topic_config_set, topic_config_query, topic_ctrl,
topic_config_resp, topic_data_loop, topic_data_event, topic_status,
)
class MainWindow(QMainWindow):
_mqtt_status = Signal(bool, str)
_mqtt_msg = Signal(str, str, object)
@@ -84,6 +89,9 @@ class MainWindow(QMainWindow):
right.addTab(self._build_config_tab(), "参数配置")
right.addTab(self._build_data_tab(), "实时数据")
right.addTab(self._build_log_tab(), "日志")
right.addTab(self._build_simulate_tab(), "模拟上报")
right.addTab(self._build_proto_topic_tab(), "协议Topic")
right.addTab(self._build_custom_topic_tab(), "自定义Topic")
splitter.addWidget(right)
splitter.setSizes([280, 960])
@@ -310,6 +318,418 @@ class MainWindow(QMainWindow):
layout.addLayout(layout2)
return w
# ================================================================
# 模拟设备上报
# ================================================================
def _build_simulate_tab(self) -> QWidget:
w = QWidget()
layout = QVBoxLayout(w)
# -- 设备序列码 --
row1 = QHBoxLayout()
row1.addWidget(QLabel("模拟设备 SN:"))
self._sim_sn = QLineEdit("DC045A49718F")
self._sim_sn.setMaximumWidth(160)
row1.addWidget(self._sim_sn)
row1.addStretch()
layout.addLayout(row1)
# -- loop_data --
g1 = QGroupBox("线圈数据上报 (loop_data)")
g1l = QVBoxLayout(g1)
self._sim_loop = QPlainTextEdit()
self._sim_loop.setFont(QFont("Consolas", 10))
self._sim_loop.setMaximumBlockCount(2000)
self._sim_loop.setPlainText(self._sample_loop_data())
g1l.addWidget(self._sim_loop)
btn_row1 = QHBoxLayout()
btn_row1.addWidget(QLabel("上报间隔(s):"))
self._sim_loop_interval = QSpinBox()
self._sim_loop_interval.setRange(1, 3600)
self._sim_loop_interval.setValue(5)
btn_row1.addWidget(self._sim_loop_interval)
b1 = QPushButton("发送一次")
b1.clicked.connect(self._sim_send_loop)
btn_row1.addWidget(b1)
self._sim_loop_btn = QPushButton("开始周期上报")
self._sim_loop_btn.setCheckable(True)
self._sim_loop_btn.toggled.connect(self._sim_toggle_loop)
btn_row1.addWidget(self._sim_loop_btn)
btn_row1.addStretch()
g1l.addLayout(btn_row1)
layout.addWidget(g1)
# -- event_report --
g2 = QGroupBox("事件上报 (event_report)")
g2l = QVBoxLayout(g2)
self._sim_event = QPlainTextEdit()
self._sim_event.setFont(QFont("Consolas", 10))
self._sim_event.setPlainText(self._sample_event_data())
g2l.addWidget(self._sim_event)
g2l2 = QHBoxLayout()
b2 = QPushButton("发送")
b2.clicked.connect(self._sim_send_event)
g2l2.addWidget(b2)
g2l2.addStretch()
g2l.addLayout(g2l2)
layout.addWidget(g2)
# -- heartbeat --
g3 = QGroupBox("心跳上报 (heartbeat)")
g3l = QVBoxLayout(g3)
self._sim_hb = QPlainTextEdit()
self._sim_hb.setFont(QFont("Consolas", 10))
self._sim_hb.setPlainText(self._sample_heartbeat_data())
g3l.addWidget(self._sim_hb)
g3l2 = QHBoxLayout()
g3l2.addWidget(QLabel("间隔(s):"))
self._sim_hb_interval = QSpinBox()
self._sim_hb_interval.setRange(5, 3600)
self._sim_hb_interval.setValue(60)
g3l2.addWidget(self._sim_hb_interval)
b3 = QPushButton("发送一次")
b3.clicked.connect(self._sim_send_hb)
g3l2.addWidget(b3)
self._sim_hb_btn = QPushButton("开始周期上报")
self._sim_hb_btn.setCheckable(True)
self._sim_hb_btn.toggled.connect(self._sim_toggle_hb)
g3l2.addWidget(self._sim_hb_btn)
g3l2.addStretch()
g3l.addLayout(g3l2)
layout.addWidget(g3)
# timers
self._sim_loop_timer = QTimer()
self._sim_loop_timer.timeout.connect(self._sim_send_loop)
self._sim_hb_timer = QTimer()
self._sim_hb_timer.timeout.connect(self._sim_send_hb)
return w
def _sample_loop_data(self) -> str:
return json.dumps({
"msg_id": 100,
"cmd": "loop_data",
"ts": 1719000100,
"data": {
"channels": [
{"ch": 1, "freq_level": "high", "has_car": False, "loop_ok": True,
"freq_current": 105280, "freq_diff": 20, "sensitivity": 7, "condition": 0,
"misc": {"type": "time", "value": 0}},
{"ch": 2, "freq_level": "mid_high", "has_car": True, "loop_ok": True,
"freq_current": 98700, "freq_diff": 1500, "sensitivity": 7, "condition": 2,
"misc": {"type": "time", "value": 350}},
]
}
}, indent=2)
def _sample_event_data(self) -> str:
return json.dumps({
"msg_id": 101,
"cmd": "event_report",
"ts": 1719000200,
"data": {"events": [
{"type": "car_enter", "ch": 2, "value": 0},
{"type": "car_leave", "ch": 2, "value": 350},
]}
}, indent=2)
def _sample_heartbeat_data(self) -> str:
return json.dumps({
"msg_id": 200,
"cmd": "heartbeat",
"ts": 1719000060,
"data": {"uptime": 3600, "loop_status": [True, True, False, True],
"net_status": True, "iot_status": True}
}, indent=2)
def _sim_publish(self, topic: str, text: str):
if not self._mqtt.connected:
QMessageBox.warning(self, "提示", "请先连接 MQTT Broker")
return
try:
payload = json.loads(text)
self._mqtt.publish(topic, payload, qos=1)
self._log(f"[模拟] → {topic}")
except json.JSONDecodeError as e:
QMessageBox.warning(self, "JSON 错误", str(e))
def _sim_send_loop(self):
sn = self._sim_sn.text()
self._sim_publish(topic_data_loop(sn), self._sim_loop.toPlainText())
def _sim_send_event(self):
sn = self._sim_sn.text()
self._sim_publish(topic_data_event(sn), self._sim_event.toPlainText())
def _sim_send_hb(self):
sn = self._sim_sn.text()
self._sim_publish(topic_status(sn), self._sim_hb.toPlainText())
def _sim_toggle_loop(self, checked: bool):
if checked:
self._sim_loop_timer.start(self._sim_loop_interval.value() * 1000)
self._sim_loop_btn.setText("停止")
else:
self._sim_loop_timer.stop()
self._sim_loop_btn.setText("开始周期上报")
def _sim_toggle_hb(self, checked: bool):
if checked:
self._sim_hb_timer.start(self._sim_hb_interval.value() * 1000)
self._sim_hb_btn.setText("停止")
else:
self._sim_hb_timer.stop()
self._sim_hb_btn.setText("开始周期上报")
# ================================================================
# 协议 Topic 订阅/发布
# ================================================================
def _build_proto_topic_tab(self) -> QWidget:
w = QWidget()
layout = QVBoxLayout(w)
row1 = QHBoxLayout()
row1.addWidget(QLabel("设备 SN:"))
self._proto_sn = QLineEdit()
self._proto_sn.setMaximumWidth(150)
self._proto_sn.textChanged.connect(self._refresh_proto_topics)
row1.addWidget(self._proto_sn)
row1.addStretch()
layout.addLayout(row1)
# 服务器下发 topics
g1 = QGroupBox("服务器下发 Topic(可发布)")
self._proto_srv_tree = QTreeWidget()
self._proto_srv_tree.setHeaderLabels(["Topic", "用途", "示例命令"])
self._proto_srv_tree.setRootIsDecorated(True)
self._proto_srv_tree.itemDoubleClicked.connect(self._proto_publish_item)
g1l = QVBoxLayout(g1)
g1l.addWidget(self._proto_srv_tree)
layout.addWidget(g1)
# 设备上报 topics
g2 = QGroupBox("设备上报 Topic(可订阅/发布)")
self._proto_dev_tree = QTreeWidget()
self._proto_dev_tree.setHeaderLabels(["Topic", "用途"])
self._proto_dev_tree.setRootIsDecorated(True)
self._proto_dev_tree.itemDoubleClicked.connect(self._proto_publish_item)
g2l = QVBoxLayout(g2)
g2l.addWidget(self._proto_dev_tree)
layout.addWidget(g2)
# 发布载荷编辑
g3 = QGroupBox("发布载荷 (JSON)")
g3l = QVBoxLayout(g3)
g3l2 = QHBoxLayout()
g3l2.addWidget(QLabel("Topic:"))
self._proto_pub_topic = QLineEdit()
g3l2.addWidget(self._proto_pub_topic, 1)
g3l.addLayout(g3l2)
self._proto_payload = QPlainTextEdit()
self._proto_payload.setFont(QFont("Consolas", 10))
self._proto_payload.setMaximumBlockCount(1000)
g3l.addWidget(self._proto_payload, 1)
btn_row = QHBoxLayout()
b1 = QPushButton("发布")
b1.clicked.connect(self._proto_publish)
btn_row.addWidget(b1)
btn_row.addStretch()
g3l.addLayout(btn_row)
layout.addWidget(g3)
return w
def _refresh_proto_topics(self):
sn = self._proto_sn.text().strip() or "{sn}"
# 服务器下发
srv = self._proto_srv_tree
srv.clear()
def _add_srv(parent, name, tpl, desc, sample_cmd=""):
topic = tpl.format(sn=sn)
item = QTreeWidgetItem(parent or srv, [topic, desc, sample_cmd])
return item
root_set = _add_srv(None, "config/set", "dld960/{sn}/srv/config/set",
"设置配置参数", "dev_serial_set / iot_net_set / loop_param_set")
root_query = _add_srv(None, "config/query", "dld960/{sn}/srv/config/query",
"查询配置参数", "dev_info_query / ssc_net_query")
root_ctrl = _add_srv(None, "ctrl", "dld960/{sn}/srv/ctrl",
"控制命令", "pwd_verify / factory_reset / device_reset")
# 设备上报
dev = self._proto_dev_tree
dev.clear()
for name, tpl, desc in [
("config/resp", "dld960/{sn}/dev/config/resp", "配置查询/设置响应"),
("data/loop", "dld960/{sn}/dev/data/loop", "线圈传感数据上报"),
("data/event", "dld960/{sn}/dev/data/event", "事件上报"),
("status", "dld960/{sn}/dev/status", "设备状态/心跳"),
]:
QTreeWidgetItem(dev, [tpl.format(sn=sn), desc])
srv.expandAll()
dev.expandAll()
def _proto_publish_item(self, item: QTreeWidgetItem, col: int):
topic = item.text(0)
self._proto_pub_topic.setText(topic)
self._proto_payload.clear()
self._notebook.setCurrentWidget(self._proto_payload.parent().parent())
def _proto_publish(self):
topic = self._proto_pub_topic.text()
if not topic or not self._mqtt.connected:
QMessageBox.warning(self, "提示", "请输入 Topic 并连接 Broker")
return
try:
if self._proto_payload.toPlainText().strip():
payload = json.loads(self._proto_payload.toPlainText())
else:
payload = build_request(CMD_DEV_INFO_QUERY)
self._mqtt.publish(topic, payload)
self._log(f"[协议] → {topic}")
except json.JSONDecodeError as e:
QMessageBox.warning(self, "JSON 错误", str(e))
# ================================================================
# 自定义 Topic 订阅/发布
# ================================================================
def _build_custom_topic_tab(self) -> QWidget:
w = QWidget()
layout = QVBoxLayout(w)
# 发布区
g1 = QGroupBox("发布 (Publish)")
g1l = QVBoxLayout(g1)
row1 = QHBoxLayout()
row1.addWidget(QLabel("Topic:"))
self._custom_pub_topic = QLineEdit()
row1.addWidget(self._custom_pub_topic, 1)
g1l.addLayout(row1)
row2 = QHBoxLayout()
row2.addWidget(QLabel("QoS:"))
self._custom_qos = QComboBox()
self._custom_qos.addItems(["0", "1", "2"])
self._custom_qos.setCurrentIndex(1)
row2.addWidget(self._custom_qos)
row2.addStretch()
g1l.addLayout(row2)
self._custom_payload = QPlainTextEdit()
self._custom_payload.setFont(QFont("Consolas", 10))
self._custom_payload.setMaximumBlockCount(2000)
g1l.addWidget(self._custom_payload, 1)
btn_pub = QPushButton("发布")
btn_pub.clicked.connect(self._custom_publish)
g1l.addWidget(btn_pub)
layout.addWidget(g1)
# 订阅区
g2 = QGroupBox("订阅 (Subscribe) / 接收")
g2l = QVBoxLayout(g2)
row3 = QHBoxLayout()
row3.addWidget(QLabel("Topic:"))
self._custom_sub_topic = QLineEdit()
row3.addWidget(self._custom_sub_topic, 1)
row3.addWidget(QLabel("QoS:"))
self._custom_sub_qos = QComboBox()
self._custom_sub_qos.addItems(["0", "1", "2"])
self._custom_sub_qos.setCurrentIndex(1)
row3.addWidget(self._custom_sub_qos)
btn_sub = QPushButton("订阅")
btn_sub.clicked.connect(self._custom_subscribe)
row3.addWidget(btn_sub)
g2l.addLayout(row3)
# 已订阅列表 + 接收消息
self._custom_sub_list = QTreeWidget()
self._custom_sub_list.setHeaderLabels(["已订阅 Topic", "QoS"])
self._custom_sub_list.setRootIsDecorated(False)
g2l.addWidget(self._custom_sub_list)
self._custom_recv = QPlainTextEdit()
self._custom_recv.setReadOnly(True)
self._custom_recv.setFont(QFont("Consolas", 10))
self._custom_recv.setMaximumBlockCount(2000)
g2l.addWidget(self._custom_recv, 1)
row4 = QHBoxLayout()
btn_unsub = QPushButton("取消订阅")
btn_unsub.clicked.connect(self._custom_unsubscribe)
row4.addWidget(btn_unsub)
row4.addStretch()
btn_clear = QPushButton("清空接收")
btn_clear.clicked.connect(self._custom_recv.clear)
row4.addWidget(btn_clear)
g2l.addLayout(row4)
layout.addWidget(g2)
# 保存自定义订阅列表,用于消息分发
self._custom_subs: dict[str, int] = {}
return w
def _custom_publish(self):
if not self._mqtt.connected:
QMessageBox.warning(self, "提示", "请先连接 Broker")
return
topic = self._custom_pub_topic.text()
if not topic:
return
try:
payload = json.loads(self._custom_payload.toPlainText()) if self._custom_payload.toPlainText().strip() else {}
qos = int(self._custom_qos.currentText())
self._mqtt.publish(topic, payload, qos=qos)
self._log(f"[自定义] → {topic}")
except json.JSONDecodeError as e:
QMessageBox.warning(self, "JSON 错误", str(e))
def _custom_subscribe(self):
if not self._mqtt.connected:
QMessageBox.warning(self, "提示", "请先连接 Broker")
return
topic = self._custom_sub_topic.text()
if not topic:
return
qos = int(self._custom_sub_qos.currentText())
self._mqtt._client.subscribe(topic, qos=qos)
self._custom_subs[topic] = qos
# 更新列表
found = False
for i in range(self._custom_sub_list.topLevelItemCount()):
item = self._custom_sub_list.topLevelItem(i)
if item.text(0) == topic:
item.setText(1, str(qos))
found = True
break
if not found:
self._custom_sub_list.addTopLevelItem(QTreeWidgetItem([topic, str(qos)]))
self._log(f"[订阅] {topic} (QoS={qos})")
def _custom_unsubscribe(self):
items = self._custom_sub_list.selectedItems()
if not items:
return
for item in items:
topic = item.text(0)
if topic in self._custom_subs:
self._mqtt._client.unsubscribe(topic)
del self._custom_subs[topic]
idx = self._custom_sub_list.indexOfTopLevelItem(item)
self._custom_sub_list.takeTopLevelItem(idx)
self._log(f"[取消订阅] {topic}")
# ================================================================
# MQTT
# ================================================================
@@ -377,6 +797,18 @@ class MainWindow(QMainWindow):
self._devmgr.update_heartbeat(dev_serial, data)
self._log(f"[{ts}] {dev_serial} ← heartbeat (uptime={data.get('uptime', 0)}s)")
# 自定义订阅消息 → 显示到自定义 Topic 接收区
for sub_topic in self._custom_subs:
if mqtt.topic_matches_sub(sub_topic, topic):
self._append_custom_recv(topic, payload)
break
def _append_custom_recv(self, topic: str, payload: dict):
ts = datetime.now().strftime("%H:%M:%S")
text = f"--- {ts} ---\nTopic: {topic}\n{json.dumps(payload, indent=2, ensure_ascii=False)}\n"
self._custom_recv.appendPlainText(text)
self._custom_recv.moveCursor(QTextCursor.End)
# ================================================================
# 设备操作
# ================================================================