From c19e4652846b16df9ba7bbca6e50eac081b955e8 Mon Sep 17 00:00:00 2001 From: wangfq Date: Mon, 6 Jul 2026 21:40:26 +0800 Subject: [PATCH] =?UTF-8?q?feat(DBNMQTTool):=20=E6=A8=A1=E6=8B=9F=E4=B8=8A?= =?UTF-8?q?=E6=8A=A5=20+=20=E5=8D=8F=E8=AE=AETopic=20+=20=E8=87=AA?= =?UTF-8?q?=E5=AE=9A=E4=B9=89Topic?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增三个 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 通配符匹配) - 自定义订阅消息自动路由到接收区 --- DBNMQTTool/main.py | 436 ++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 434 insertions(+), 2 deletions(-) diff --git a/DBNMQTTool/main.py b/DBNMQTTool/main.py index 8f0120b..a9713ea 100644 --- a/DBNMQTTool/main.py +++ b/DBNMQTTool/main.py @@ -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) + # ================================================================ # 设备操作 # ================================================================