#!/usr/bin/env python3 """ DBN MQTT Tool — DLD960 IoT MQTT 设备管理工具 跨平台桌面软件 (PySide6) """ import sys import os import json import queue import time from datetime import datetime from typing import Optional sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) from PySide6.QtWidgets import ( QApplication, QMainWindow, QWidget, QVBoxLayout, QHBoxLayout, QLabel, QLineEdit, QPushButton, QTreeWidget, QTreeWidgetItem, QTabWidget, QTextEdit, QGroupBox, QGridLayout, QCheckBox, QSplitter, QMessageBox, QHeaderView, QFrame, QComboBox, QSpinBox, QPlainTextEdit, QFileDialog, QProgressBar, ) from PySide6.QtCore import Qt, QTimer, Signal, Slot, QThread 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 from dbn_mqtt_tool.protocol import ( CMD_DEV_INFO_QUERY, CMD_SSC_NET_QUERY, CMD_IOT_NET_QUERY, CMD_IOT_TOPIC_QUERY, CMD_LOOP_PARAM_QUERY, CMD_REPORT_CONFIG, CMD_SSC_NET_SET, CMD_IOT_NET_SET, CMD_IOT_TOPIC_SET, CMD_PWD_VERIFY, CMD_PWD_SET, CMD_FACTORY_RESET, CMD_DEVICE_RESET, CMD_LOOP_DATA, CMD_EVENT_REPORT, CMD_HEARTBEAT, CMD_INITIALIZE, CMD_LOG_STAT, CMD_LOG_QUERY, CMD_LOG_CLEAR, STREAM_EVENT, STREAM_SNAPSHOT, ERROR_MSGS, FREQ_LEVELS, OUTPUT_MODES, EVENT_TYPES, LOG_EVENT_TYPE_DESC, SNAP_MISC_TYPE_DESC, OFFLOG_TYPE_NAMES, parse_offlog_hex, offlog_payload_desc, parse_snap_hex, data_ssc_net_set, data_iot_net_set, data_iot_topic_set, data_pwd_verify, data_pwd_set, data_report_config, data_log_query, data_log_clear, build_request, next_msg_id, topic_down, topic_up, # V1.08 OTA CMD_OTA_BEGIN, CMD_OTA_DATA, CMD_OTA_END, CMD_OTA_ABORT, CMD_OTA_FLASH, CMD_OTA_STATUS, CMD_OTA_REPORT, OTA_TARGET_LOOP, OTA_CHUNK_SIZE, OTA_SLOT_A, OTA_SLOT_B, OTA_STATE_DESC, OTA_STAGE_DESC, OTA_ERROR_DESC, data_ota_begin, data_ota_data, data_ota_end, data_ota_abort, data_ota_flash, ota_crc32, ota_split_bin, ) # OTA 相关命令集合 (用于响应分发) OTA_CMD_SET = {CMD_OTA_BEGIN, CMD_OTA_DATA, CMD_OTA_END, CMD_OTA_ABORT, CMD_OTA_FLASH, CMD_OTA_STATUS} class OtaDownloadThread(QThread): """OTA 分片下载线程 (协议 V1.08): ota_begin → ota_data×N → ota_end 断点续传: ota_begin 响应 offset 为设备已接收字节数, 从此处续传 重定位: ota_data code=1 时按 data.offset 调整 (缺片/乱序/CRC) """ progress = Signal(int, int) # sent, total log = Signal(str) finished_ok = Signal(str) finished_err = Signal(str) def __init__(self, mqtt_client: MqttClient, dev_serial: str, bin_data: bytes, version: str, force: bool, resp_queue: "queue.Queue"): super().__init__() self._mqtt = mqtt_client self._sn = dev_serial self._bin = bin_data self._version = version self._force = force self._resp_queue = resp_queue self._stop = False def stop(self): self._stop = True def _send(self, cmd: str, data: dict) -> int: from dbn_mqtt_tool.protocol import get_topic_for_cmd msg = build_request(cmd, data) self._mqtt.publish(get_topic_for_cmd(cmd, self._sn), msg, qos=1) return msg["msg_id"] def _wait_resp(self, mid: int, cmd: str, timeout: float = 5.0) -> dict: """等待与 msg_id 匹配的响应 (resp_queue 由 _on_message 填充)""" deadline = time.time() + timeout while time.time() < deadline and not self._stop: try: c, payload = self._resp_queue.get(timeout=0.2) except queue.Empty: continue if c == cmd and payload.get("msg_id") == mid: return payload raise TimeoutError(f"{cmd} 响应超时 (msg_id={mid})") def run(self): try: size = len(self._bin) total_crc = ota_crc32(self._bin) chunks = ota_split_bin(self._bin) self.log.emit(f"[OTA] 镜像 {size}B CRC32=0x{total_crc:08X} → {len(chunks)} 片 @{OTA_CHUNK_SIZE}B") # 1. ota_begin (断点续传定位) mid = self._send(CMD_OTA_BEGIN, data_ota_begin( size, total_crc, version=self._version, force=self._force)) resp = self._wait_resp(mid, CMD_OTA_BEGIN) if resp.get("code") != 0: raise RuntimeError(f"ota_begin 失败: code={resp.get('code')} {resp.get('msg')} " f"(err_code={resp.get('data', {}).get('err_code', '-')})") d = resp.get("data", {}) offset = int(d.get("offset", 0)) self.log.emit(f"[OTA] begin OK slot={d.get('slot', '?')} state={d.get('state', '?')} " f"offset={offset} (断点续传点)") # 2. 逐片下发 (停等: 每片等响应) retry = 0 while offset < size and not self._stop: idx = offset // OTA_CHUNK_SIZE if idx >= len(chunks): # 设备 received 溢出, 直接截断到 size break off, crc, hexs = chunks[idx] mid = self._send(CMD_OTA_DATA, data_ota_data(off, crc, hexs)) resp = self._wait_resp(mid, CMD_OTA_DATA) code = resp.get("code") if code == 0: offset = int(resp.get("data", {}).get("received", offset + OTA_CHUNK_SIZE)) retry = 0 self.progress.emit(min(offset, size), size) elif code == 1: # 缺片/乱序/CRC 失败 → 按设备指示重定位 (data.offset=received) new_off = int(resp.get("data", {}).get("offset", offset)) if new_off == offset: retry += 1 if retry > 3: raise RuntimeError("分片连续失败 (>3 次), 建议 ota_abort 后重来") else: retry = 0 offset = new_off else: raise RuntimeError(f"ota_data@{off} 失败: code={code} {resp.get('msg')} " f"(err_code={resp.get('data', {}).get('err_code', '-')})") self.log.emit(f"[OTA] data@{off} → received={offset}") if self._stop: self.log.emit("[OTA] 已停止") return self.progress.emit(size, size) # 3. ota_end (全镜像 CRC32 复核) mid = self._send(CMD_OTA_END, data_ota_end(total_crc)) resp = self._wait_resp(mid, CMD_OTA_END, timeout=10) if resp.get("code") != 0: raise RuntimeError(f"ota_end 失败: code={resp.get('code')} {resp.get('msg')} " f"crc_ok={resp.get('data', {}).get('crc_ok', '-')}") self.log.emit("[OTA] end OK crc_ok=true → state=ready, 可触发 ota_flash") self.finished_ok.emit("下载+校验完成, 可触发刷写 (ota_flash)") except Exception as e: self.log.emit(f"[OTA] ✗ {e}") self.finished_err.emit(str(e)) class MainWindow(QMainWindow): _mqtt_status = Signal(bool, str) _mqtt_msg = Signal(str, str, object) _devices_changed = Signal() def __init__(self): super().__init__() self.setWindowTitle("DBN MQTT Tool — DLD960 IoT 设备管理") self.resize(1280, 800) self.setMinimumSize(960, 640) self._mqtt = MqttClient() self._devmgr = DeviceManager() self._mqtt.on_status_change(lambda c, m: self._mqtt_status.emit(c, m)) self._mqtt.on_message(lambda t, s, p: self._mqtt_msg.emit(t, s, p)) self._devmgr.on_change(lambda: self._devices_changed.emit()) self._mqtt_status.connect(self._on_status) self._mqtt_msg.connect(self._on_message) self._devices_changed.connect(self._refresh_devices) self._build_ui() self._log("DBN MQTT Tool 已启动") # ================================================================ # UI 构建 # ================================================================ def _build_ui(self): central = QWidget() self.setCentralWidget(central) root = QVBoxLayout(central) root.setContentsMargins(8, 8, 8, 8) root.setSpacing(6) # -- 顶栏 -- root.addWidget(self._build_top_bar()) # -- 主体 -- splitter = QSplitter(Qt.Horizontal) left = self._build_device_panel() splitter.addWidget(left) self._notebook = QTabWidget() self._notebook.addTab(self._build_info_tab(), "设备信息") self._notebook.addTab(self._build_config_tab(), "参数配置") self._notebook.addTab(self._build_data_tab(), "实时数据") self._notebook.addTab(self._build_log_tab(), "日志") self._notebook.addTab(self._build_ota_tab(), "OTA") self._notebook.addTab(self._build_simulate_tab(), "模拟上报") self._notebook.addTab(self._build_proto_topic_tab(), "协议Topic") self._notebook.addTab(self._build_custom_topic_tab(), "自定义Topic") splitter.addWidget(self._notebook) splitter.setSizes([280, 960]) root.addWidget(splitter, 1) def _build_top_bar(self) -> QWidget: bar = QWidget() layout = QHBoxLayout(bar) layout.setContentsMargins(0, 0, 0, 0) layout.setSpacing(6) layout.addWidget(QLabel("Broker:")) self._host_edit = QLineEdit("121.37.20.199") self._host_edit.setMaximumWidth(140) layout.addWidget(self._host_edit) layout.addWidget(QLabel("Port:")) self._port_edit = QLineEdit("1883") self._port_edit.setMaximumWidth(60) layout.addWidget(self._port_edit) layout.addWidget(QLabel("用户:")) self._user_edit = QLineEdit() self._user_edit.setMaximumWidth(90) layout.addWidget(self._user_edit) layout.addWidget(QLabel("密码:")) self._pass_edit = QLineEdit() self._pass_edit.setMaximumWidth(90) self._pass_edit.setEchoMode(QLineEdit.Password) layout.addWidget(self._pass_edit) self._btn_conn = QPushButton("连接") self._btn_conn.setFixedWidth(70) self._btn_conn.clicked.connect(self._toggle_connect) layout.addWidget(self._btn_conn) self._status_label = QLabel("● 未连接") self._status_label.setStyleSheet("color: red; font-weight: bold;") layout.addWidget(self._status_label) layout.addStretch() return bar def _build_device_panel(self) -> QWidget: w = QWidget() layout = QVBoxLayout(w) layout.setContentsMargins(0, 0, 0, 0) layout.addWidget(QLabel("设备列表")) self._dev_tree = QTreeWidget() self._dev_tree.setHeaderLabels(["设备序列码", "型号"]) self._dev_tree.setRootIsDecorated(False) self._dev_tree.header().setStretchLastSection(True) self._dev_tree.itemSelectionChanged.connect(self._on_dev_select) layout.addWidget(self._dev_tree, 1) btn_row = QHBoxLayout() btn_query = QPushButton("查询设备信息") btn_query.clicked.connect(self._query_dev_info) btn_row.addWidget(btn_query) btn_all = QPushButton("查询全部") btn_all.clicked.connect(self._query_all) btn_row.addWidget(btn_all) layout.addLayout(btn_row) return w def _build_info_tab(self) -> QWidget: w = QWidget() layout = QVBoxLayout(w) self._info_text = QTextEdit() self._info_text.setReadOnly(True) self._info_text.setFont(QFont("Consolas", 10)) layout.addWidget(self._info_text) return w def _build_config_tab(self) -> QWidget: w = QWidget() layout = QVBoxLayout(w) # -- SSC 网络 -- g1 = QGroupBox("SSC 网络配置") g1_layout = QGridLayout(g1) fields_ssc = [ ("设备 IP", "ssc_ip"), ("子网掩码", "ssc_mask"), ("网关", "ssc_gw"), ("LSSC IP", "ssc_lssc"), ("DNS", "ssc_dns"), ("端口", "ssc_port"), ] for i, (label, name) in enumerate(fields_ssc): g1_layout.addWidget(QLabel(label + ":"), i // 2, (i % 2) * 2) edit = QLineEdit() edit.setMaximumWidth(140) g1_layout.addWidget(edit, i // 2, (i % 2) * 2 + 1) setattr(self, f"_edit_{name}", edit) btn_row = QHBoxLayout() b1 = QPushButton("查询") b1.clicked.connect(self._query_ssc_net) b2 = QPushButton("设置") b2.clicked.connect(self._set_ssc_net) b2.setEnabled(False) b2.setToolTip("MQTT 固件暂未实现 ssc_net_set,请用 TCP JSON 或蓝牙通道配置") btn_row.addWidget(b1) btn_row.addWidget(b2) btn_row.addStretch() g1_layout.addLayout(btn_row, 3, 0, 1, 4) layout.addWidget(g1) # -- IoT 网络 -- g2 = QGroupBox("IoT 网络配置") g2_layout = QGridLayout(g2) fields_iot = [ ("MQTT Host", "iot_host"), ("MQTT Port", "iot_port"), ("Client ID", "iot_cid"), ("用户名", "iot_user"), ("密码", "iot_pass"), ] for i, (label, name) in enumerate(fields_iot): g2_layout.addWidget(QLabel(label + ":"), i, 0) edit = QLineEdit() edit.setMinimumWidth(200) if "pass" in name: edit.setEchoMode(QLineEdit.Password) g2_layout.addWidget(edit, i, 1) setattr(self, f"_edit_{name}", edit) btn_row2 = QHBoxLayout() b3 = QPushButton("查询") b3.clicked.connect(self._query_iot_net) b4 = QPushButton("设置") b4.clicked.connect(self._set_iot_net) b4.setEnabled(False) b4.setToolTip("MQTT 固件暂未实现 iot_net_set,请用 TCP JSON 或蓝牙通道配置") btn_row2.addWidget(b3) btn_row2.addWidget(b4) btn_row2.addStretch() g2_layout.addLayout(btn_row2, len(fields_iot), 0, 1, 2) layout.addWidget(g2) # -- Topic -- g3 = QGroupBox("IoT Topic 配置") g3_layout = QGridLayout(g3) self._chk_cid = QCheckBox("Client ID 启用") self._chk_cid.setChecked(True) g3_layout.addWidget(self._chk_cid, 0, 0, 1, 2) g3_layout.addWidget(QLabel("Topic 发布:"), 1, 0) self._edit_topic_pub = QLineEdit() self._edit_topic_pub.setMinimumWidth(250) g3_layout.addWidget(self._edit_topic_pub, 1, 1) g3_layout.addWidget(QLabel("Topic 订阅:"), 2, 0) self._edit_topic_sub = QLineEdit() g3_layout.addWidget(self._edit_topic_sub, 2, 1) btn_row3 = QHBoxLayout() b5 = QPushButton("查询") b5.clicked.connect(self._query_iot_topic) b6 = QPushButton("设置") b6.clicked.connect(self._set_iot_topic) b6.setEnabled(False) b6.setToolTip("MQTT 固件暂未实现 iot_topic_set,请用 TCP JSON 或蓝牙通道配置") btn_row3.addWidget(b5) btn_row3.addWidget(b6) btn_row3.addStretch() g3_layout.addLayout(btn_row3, 3, 0, 1, 2) layout.addWidget(g3) # -- 控制 -- g4 = QGroupBox("控制命令") g4_layout = QGridLayout(g4) g4_layout.addWidget(QLabel("密码:"), 0, 0) self._edit_pwd = QLineEdit() self._edit_pwd.setMaximumWidth(100) self._edit_pwd.setEchoMode(QLineEdit.Password) g4_layout.addWidget(self._edit_pwd, 0, 1) g4_layout.addWidget(QLabel("新密码:"), 0, 2) self._edit_new_pwd = QLineEdit() self._edit_new_pwd.setMaximumWidth(100) self._edit_new_pwd.setEchoMode(QLineEdit.Password) g4_layout.addWidget(self._edit_new_pwd, 0, 3) btn_row4 = QHBoxLayout() for label, slot in [ ("验证密码", self._do_pwd_verify), ("设置密码", self._do_pwd_set), ("出厂初始化", self._do_factory_reset), ("设备复位", self._do_device_reset), ]: btn = QPushButton(label) btn.clicked.connect(slot) if label != "验证密码": btn.setEnabled(False) btn.setToolTip(f"MQTT 固件暂未实现 {label},请用 TCP JSON 或蓝牙通道") btn_row4.addWidget(btn) btn_row4.addStretch() g4_layout.addLayout(btn_row4, 1, 0, 1, 4) layout.addWidget(g4) # -- 主动上报 -- g5 = QGroupBox("主动上报配置 (report_config)") g5_layout = QGridLayout(g5) g5_layout.addWidget(QLabel("传感器类型:"), 0, 0) self._edit_report_sensor_type = QSpinBox() self._edit_report_sensor_type.setRange(0, 255) self._edit_report_sensor_type.setValue(12) g5_layout.addWidget(self._edit_report_sensor_type, 0, 1) g5_layout.addWidget(QLabel("上报间隔(s):"), 0, 2) self._edit_report_interval = QSpinBox() self._edit_report_interval.setRange(0, 86400) self._edit_report_interval.setValue(5) g5_layout.addWidget(self._edit_report_interval, 0, 3) g5_layout.addWidget(QLabel("超时(ms):"), 0, 4) self._edit_report_timeout = QSpinBox() self._edit_report_timeout.setRange(0, 60000) g5_layout.addWidget(self._edit_report_timeout, 0, 5) self._chk_report_enable = QCheckBox("启用上报") self._chk_report_enable.setChecked(True) g5_layout.addWidget(self._chk_report_enable, 1, 0) self._chk_report_once = QCheckBox("单次上报") g5_layout.addWidget(self._chk_report_once, 1, 1) self._chk_report_env = QCheckBox("环境评估") g5_layout.addWidget(self._chk_report_env, 1, 2) self._chk_report_ack = QCheckBox("需应答") g5_layout.addWidget(self._chk_report_ack, 1, 3) btn_row5 = QHBoxLayout() b7 = QPushButton("查询") b7.clicked.connect(self._query_report_config) b8 = QPushButton("设置") b8.clicked.connect(self._set_report_config) btn_row5.addWidget(b7) btn_row5.addWidget(b8) btn_row5.addStretch() g5_layout.addLayout(btn_row5, 2, 0, 1, 6) layout.addWidget(g5) # -- 脱机日志 (V1.06 事件流 / V1.07 快照流) -- g6 = QGroupBox("脱机日志 (log_stat / log_query / log_clear)") g6_layout = QGridLayout(g6) g6_layout.addWidget(QLabel("日志流:"), 0, 0) self._combo_log_stream = QComboBox() self._combo_log_stream.addItem("事件日志 (event)", STREAM_EVENT) self._combo_log_stream.addItem("传感快照 (snapshot)", STREAM_SNAPSHOT) self._combo_log_stream.setCurrentIndex(0) g6_layout.addWidget(self._combo_log_stream, 0, 1) g6_layout.addWidget(QLabel("起始序号:"), 0, 2) self._edit_log_start_seq = QSpinBox() self._edit_log_start_seq.setRange(1, 2**31 - 1) self._edit_log_start_seq.setValue(1) g6_layout.addWidget(self._edit_log_start_seq, 0, 3) g6_layout.addWidget(QLabel("条数:"), 0, 4) self._edit_log_count = QSpinBox() self._edit_log_count.setRange(1, 4) self._edit_log_count.setValue(4) g6_layout.addWidget(self._edit_log_count, 0, 5) btn_row6 = QHBoxLayout() b9 = QPushButton("统计") b9.clicked.connect(self._query_log_stat) self._btn_log_prev = QPushButton("◀ 上一页") self._btn_log_prev.clicked.connect(self._page_log_prev) self._btn_log_prev.setEnabled(False) self._btn_log_prev.setToolTip("起始序号 = 当前页第一条 - 条数(下限 1)") b10 = QPushButton("拉取日志") b10.clicked.connect(self._query_log) self._btn_log_next = QPushButton("下一页 ▶") self._btn_log_next.clicked.connect(self._page_log_next) self._btn_log_next.setEnabled(False) self._btn_log_next.setToolTip("起始序号 = 当前页最后一条 seq + 1") b11 = QPushButton("清除日志") b11.clicked.connect(self._clear_log) btn_row6.addWidget(b9) btn_row6.addWidget(self._btn_log_prev) btn_row6.addWidget(b10) btn_row6.addWidget(self._btn_log_next) btn_row6.addWidget(b11) btn_row6.addStretch() g6_layout.addLayout(btn_row6, 1, 0, 1, 4) self._log_rec_text = QPlainTextEdit() self._log_rec_text.setReadOnly(True) # 脱机日志翻页状态 (翻页锚点) self._log_page_last_seq = 0 # 当前页最后一条 seq; 0 = 无有效页 self._log_rec_text.setMaximumHeight(150) self._log_rec_text.setFont(QFont("Consolas", 9)) g6_layout.addWidget(self._log_rec_text, 2, 0, 1, 4) layout.addWidget(g6) layout.addStretch() return w def _build_data_tab(self) -> QWidget: w = QWidget() layout = QVBoxLayout(w) layout.addWidget(QLabel("线圈传感数据 (loop_data)")) self._loop_text = QTextEdit() self._loop_text.setReadOnly(True) self._loop_text.setFont(QFont("Consolas", 10)) layout.addWidget(self._loop_text, 2) layout.addWidget(QLabel("事件上报 (event_report)")) self._event_text = QTextEdit() self._event_text.setReadOnly(True) self._event_text.setFont(QFont("Consolas", 10)) layout.addWidget(self._event_text, 1) return w def _build_log_tab(self) -> QWidget: w = QWidget() layout = QVBoxLayout(w) self._log_text = QTextEdit() self._log_text.setReadOnly(True) self._log_text.setFont(QFont("Consolas", 9)) layout.addWidget(self._log_text, 1) btn_clear = QPushButton("清空") btn_clear.clicked.connect(self._log_text.clear) layout2 = QHBoxLayout() layout2.addStretch() layout2.addWidget(btn_clear) layout.addLayout(layout2) return w def _build_ota_tab(self) -> QWidget: """OTA 升级页签 (协议 V1.08: Loop 远程 OTA, 先存后刷)""" w = QWidget() layout = QVBoxLayout(w) # -- 固件选择 -- row1 = QHBoxLayout() row1.addWidget(QLabel("固件 bin:")) self._ota_file = QLineEdit() self._ota_file.setReadOnly(True) self._ota_file.setPlaceholderText("选择 Loop 固件 bin (≤96KB, 0x08003400 起 APP 映像)") row1.addWidget(self._ota_file, 1) btn_file = QPushButton("选择…") btn_file.clicked.connect(self._ota_choose_file) row1.addWidget(btn_file) layout.addLayout(row1) # -- 参数 -- row2 = QHBoxLayout() row2.addWidget(QLabel("目标版本:")) self._ota_version = QLineEdit("1.1.0") self._ota_version.setMaximumWidth(110) row2.addWidget(self._ota_version) self._ota_force = QCheckBox("force") self._ota_force.setToolTip("强制: 覆盖现有镜像 / 跳过安全检查 (高风险)") row2.addWidget(self._ota_force) row2.addWidget(QLabel("slot:")) self._ota_slot = QComboBox() self._ota_slot.addItems([OTA_SLOT_A, OTA_SLOT_B]) self._ota_slot.setMaximumWidth(60) row2.addWidget(self._ota_slot) self._ota_info = QLabel("未选择固件") row2.addWidget(self._ota_info, 1) layout.addLayout(row2) # -- 操作按钮 -- row3 = QHBoxLayout() b_begin = QPushButton("① 下载 (begin/data/end)") b_begin.setToolTip("开启会话→分片下发→全镜像校验; 已下载过则断点续传") b_begin.clicked.connect(self._ota_start_download) row3.addWidget(b_begin) b_flash = QPushButton("② 刷写 (ota_flash)") b_flash.setToolTip("仅 state=ready 可触发; 安全检查 (有车拒绝) 后异步执行") b_flash.clicked.connect(self._ota_send_flash) row3.addWidget(b_flash) b_abort = QPushButton("中止 (ota_abort)") b_abort.clicked.connect(self._ota_send_abort) row3.addWidget(b_abort) b_status = QPushButton("查询状态 (ota_status)") b_status.clicked.connect(self._ota_send_status) row3.addWidget(b_status) b_clear = QPushButton("清空日志") b_clear.clicked.connect(lambda: self._ota_log_text.clear()) row3.addWidget(b_clear) layout.addLayout(row3) # -- 进度 -- self._ota_progress = QProgressBar() self._ota_progress.setRange(0, 100) self._ota_progress.setValue(0) layout.addWidget(self._ota_progress) # -- 状态日志 -- self._ota_log_text = QPlainTextEdit() self._ota_log_text.setReadOnly(True) self._ota_log_text.setFont(QFont("Consolas", 9)) self._ota_log_text.setMaximumBlockCount(2000) layout.addWidget(self._ota_log_text, 1) # -- 响应队列 (下载线程 ← _on_message) -- self._ota_resp_queue: "queue.Queue" = queue.Queue() self._ota_thread: Optional[OtaDownloadThread] = None return w # ================================================================ # OTA 升级 (V1.08) # ================================================================ def _ota_choose_file(self): path, _ = QFileDialog.getOpenFileName(self, "选择固件 bin", "", "固件文件 (*.bin);;所有文件 (*)") if path: self._ota_file.setText(path) try: with open(path, "rb") as f: data = f.read() if len(data) > 96 * 1024: QMessageBox.warning(self, "固件过大", f"{len(data)}B 超上限 96KB (Slot 100KB-4KB)") crc = ota_crc32(data) self._ota_info.setText(f"{len(data)}B · CRC32=0x{crc:08X}") except OSError as e: QMessageBox.critical(self, "读取失败", str(e)) def _ota_load_bin(self) -> Optional[bytes]: path = self._ota_file.text() if not path: QMessageBox.warning(self, "提示", "请先选择固件 bin") return None try: data = open(path, "rb").read() except OSError as e: QMessageBox.critical(self, "读取失败", str(e)) return None if len(data) > 96 * 1024: QMessageBox.warning(self, "固件过大", f"{len(data)}B 超上限 96KB") return None return data def _ota_log(self, text: str): self._ota_log_text.appendPlainText(f"[{datetime.now().strftime('%H:%M:%S')}] {text}") def _ota_start_download(self): sn = self._require_dev() if not sn: return if self._ota_thread and self._ota_thread.isRunning(): QMessageBox.warning(self, "提示", "OTA 下载已在运行, 先中止") return data = self._ota_load_bin() if data is None: return self._ota_log(f"开始 OTA 下载: sn={sn} size={len(data)}B") self._ota_progress.setValue(0) self._ota_thread = OtaDownloadThread( self._mqtt, sn, data, version=self._ota_version.text().strip() or "0.0.0", force=self._ota_force.isChecked(), resp_queue=self._ota_resp_queue, ) self._ota_thread.progress.connect(self._ota_on_progress) self._ota_thread.log.connect(self._ota_log) self._ota_thread.finished_ok.connect( lambda m: (self._ota_log(f"✓ {m}"), QMessageBox.information(self, "OTA", m))) self._ota_thread.finished_err.connect( lambda m: QMessageBox.critical(self, "OTA 失败", m)) self._ota_thread.start() def _ota_on_progress(self, sent: int, total: int): pct = int(sent * 100 / total) if total else 0 self._ota_progress.setValue(pct) self._ota_progress.setFormat(f"{sent}/{total} B ({pct}%)") def _ota_send_flash(self): sn = self._require_dev() if not sn: return confirm = QMessageBox.question( self, "确认刷写", "触发 ota_flash 将把暂存镜像刷入 Loop MCU (约 6~8s 窗口, 期间检测中断)。\n" "请确认现场无车压线圈。继续?", QMessageBox.Yes | QMessageBox.No) if confirm != QMessageBox.Yes: return self._send_cmd(CMD_OTA_FLASH, data_ota_flash( slot=self._ota_slot.currentText(), force=self._ota_force.isChecked())) def _ota_send_abort(self): if self._ota_thread and self._ota_thread.isRunning(): self._ota_thread.stop() self._send_cmd(CMD_OTA_ABORT, data_ota_abort()) def _ota_send_status(self): self._send_cmd(CMD_OTA_STATUS) def _ota_show_status(self, data: dict): """显示 ota_status / ota_report 状态""" line = [] state = data.get("state", "") stage = data.get("stage", "") prog = data.get("progress", {}) if state: line.append(f"state={state}({OTA_STATE_DESC.get(state, '?')})") if stage: line.append(f"stage={stage}({OTA_STAGE_DESC.get(stage, '?')})") if prog: line.append(f"progress={prog.get('sent', 0)}/{prog.get('total', 0)}B") if data.get("received") is not None: line.append(f"received={data['received']}B") if data.get("last_result"): line.append(f"last_result={data['last_result']}") if data.get("last_error"): err = OTA_ERROR_DESC.get(data["last_error"], data["last_error"]) line.append(f"last_error={err}") self._ota_log("状态: " + (" | ".join(line) if line else json.dumps(data, ensure_ascii=False))) # ================================================================ # 模拟设备上报 # ================================================================ 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) # -- initialize -- g_init = QGroupBox("上电初始化 (initialize) — V1.03") g_init_l = QVBoxLayout(g_init) self._sim_init = QPlainTextEdit() self._sim_init.setFont(QFont("Consolas", 10)) self._sim_init.setPlainText(self._sample_initialize_data()) g_init_l.addWidget(self._sim_init) g_init_l2 = QHBoxLayout() b_init = QPushButton("发送") b_init.clicked.connect(self._sim_send_initialize) g_init_l2.addWidget(b_init) g_init_l2.addStretch() g_init_l.addLayout(g_init_l2) layout.addWidget(g_init) # -- 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, "level": "high", "iscar": False, "loop_ok": True, "freq": 105280, "diff": 20, "sens": 7, "cndtn": 0, "misc": {"type": "time", "value": 0}}, {"ch": 2, "level": "mid_high", "iscar": True, "loop_ok": True, "freq": 98700, "diff": 1500, "sens": 7, "cndtn": 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_initialize_data(self) -> str: return json.dumps({ "msg_id": 1, "cmd": "initialize", "ts": 1719000000, "data": { "dev_serial": "DC045A49718F", "model": "DLD960", "hard_ver": "1.0", "soft_ver": "1.0", "loop_ver": "1.2.3", "loop_hw_ver": "1.0.0", "extra_info": { "code": "869756049404948", "csq": "21", "location": "113.9237976,022.6400375" } } }, 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_send("[模拟]", topic, payload) except json.JSONDecodeError as e: QMessageBox.warning(self, "JSON 错误", str(e)) except Exception as e: QMessageBox.critical(self, "发布失败", f"{type(e).__name__}: {e}") def _sim_send_loop(self): sn = self._sim_sn.text() self._sim_publish(topic_up(sn), self._sim_loop.toPlainText()) def _sim_send_event(self): sn = self._sim_sn.text() self._sim_publish(topic_up(sn), self._sim_event.toPlainText()) def _sim_send_initialize(self): sn = self._sim_sn.text() self._sim_publish(topic_up(sn), self._sim_init.toPlainText()) def _sim_send_hb(self): sn = self._sim_sn.text() self._sim_publish(topic_up(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(False) 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(False) 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() QTreeWidgetItem(srv, [topic_down(sn), "设备订阅此主题接收所有命令"]) # 设备上报 — 单主题 dev = self._proto_dev_tree dev.clear() item = QTreeWidgetItem(dev, [topic_up(sn), "设备发布:loop_data / event_report / heartbeat / 命令响应"]) # 子条目列出 cmd 类型 for cmd_name, desc in [ ("initialize", "设备上电初始化 (V1.03)"), ("loop_data", "线圈传感数据上报"), ("event_report", "事件上报(有车/无车/故障)"), ("heartbeat", "设备心跳"), ("响应", "各命令查询/设置响应(msg_id 匹配)"), ]: QTreeWidgetItem(item, [f" cmd={cmd_name}", desc]) def _proto_publish_item(self, item: QTreeWidgetItem, col: int): # 如果点中的是子条目(无有效 topic),取父条目的 topic topic = item.text(0) if item.parent() and not topic.startswith("dld960/"): topic = item.parent().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 raw = self._proto_payload.toPlainText().strip() try: if raw: payload = json.loads(raw) else: payload = build_request(CMD_DEV_INFO_QUERY) self._mqtt.publish(topic, payload) self._log_send("[协议]", topic, payload) except json.JSONDecodeError as e: QMessageBox.warning(self, "JSON 错误", str(e)) except Exception as e: QMessageBox.critical(self, "发布失败", f"{type(e).__name__}: {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 raw = self._custom_payload.toPlainText().strip() try: payload = json.loads(raw) if raw else {} qos = int(self._custom_qos.currentText()) self._mqtt.publish(topic, payload, qos=qos) self._log_send("[自定义]", topic, payload) except json.JSONDecodeError as e: QMessageBox.warning(self, "JSON 错误", str(e)) except Exception as e: QMessageBox.critical(self, "发布失败", f"{type(e).__name__}: {e}") self._log(f"[自定义] 发布失败: {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 try: qos = int(self._custom_sub_qos.currentText()) self._mqtt.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})") except Exception as e: QMessageBox.critical(self, "订阅失败", f"{type(e).__name__}: {e}") self._log(f"[订阅] 失败: {e}") def _custom_unsubscribe(self): items = self._custom_sub_list.selectedItems() if not items: return for item in items: topic = item.text(0) try: if topic in self._custom_subs: self._mqtt.unsubscribe(topic) del self._custom_subs[topic] idx = self._custom_sub_list.indexOfTopLevelItem(item) self._custom_sub_list.takeTopLevelItem(idx) self._log(f"[取消订阅] {topic}") except Exception as e: self._log(f"[取消订阅] 失败: {e}") # ================================================================ # MQTT # ================================================================ def _toggle_connect(self): if self._mqtt.connected: self._mqtt.disconnect() else: cfg = BrokerConfig( host=self._host_edit.text(), port=int(self._port_edit.text() or "1883"), username=self._user_edit.text(), password=self._pass_edit.text(), ) self._mqtt.connect(cfg) self._log(f"正在连接 {cfg.host}:{cfg.port} ...") @Slot(bool, str) def _on_status(self, connected: bool, msg: str): if connected: self._btn_conn.setText("断开") self._status_label.setText("● 已连接") self._status_label.setStyleSheet("color: green; font-weight: bold;") else: self._btn_conn.setText("连接") self._status_label.setText("● 未连接") self._status_label.setStyleSheet("color: red; font-weight: bold;") self._log(msg) @Slot(str, str, object) def _on_message(self, topic: str, dev_serial: str, payload: dict): cmd = payload.get("cmd", "") ts = datetime.fromtimestamp(payload.get("ts", 0)).strftime("%H:%M:%S") # 响应类消息(含 code 字段) if "code" in payload: code = payload.get("code", -1) pmsg = payload.get("msg", "") data = payload.get("data", {}) self._log_recv(topic, payload, f"{cmd} code={code} {pmsg}") if code == 0: if cmd == CMD_DEV_INFO_QUERY: self._devmgr.update_from_dev_info(dev_serial, data) self._show_json(data) elif cmd == CMD_SSC_NET_QUERY: self._apply_ssc_net_data(data) elif cmd == CMD_IOT_NET_QUERY: self._apply_iot_net_data(data) elif cmd == CMD_IOT_TOPIC_QUERY: self._apply_iot_topic_data(data) elif cmd == CMD_LOOP_PARAM_QUERY: self._show_json({cmd: data}) elif cmd == CMD_REPORT_CONFIG: self._apply_report_config_data(data) elif cmd == CMD_LOG_STAT: self._apply_log_stat_data(data) elif cmd == CMD_LOG_QUERY: self._apply_log_query_data(data) elif cmd == CMD_LOG_CLEAR: self._log(f"脱机日志已清除 (审计留痕, 需重新统计确认)") elif cmd in OTA_CMD_SET: # OTA 响应 → 喂下载线程响应队列 + 状态显示 self._ota_resp_queue.put((cmd, payload)) self._ota_show_status(payload.get("data", {})) else: if cmd in OTA_CMD_SET: self._ota_resp_queue.put((cmd, payload)) self._ota_log(f"{cmd} 失败 code={code} {pmsg} " f"(err_code={payload.get('data', {}).get('err_code', '-')})") self._show_json({"error": f"code={code} {pmsg} ({ERROR_MSGS.get(code, '?')})"}) elif cmd == CMD_LOOP_DATA: data = payload.get("data", {}) self._devmgr.update_loop_data(dev_serial, data) chs = len(data.get("channels", [])) self._log_recv(topic, payload, f"loop_data ({chs}ch)") self._append_text(self._loop_text, json.dumps(payload, indent=2, ensure_ascii=False)) elif cmd == CMD_EVENT_REPORT: data = payload.get("data", {}) self._devmgr.update_event(dev_serial, data) evt_count = len(data.get("events", [])) self._log_recv(topic, payload, f"event ({evt_count} events)") self._append_text(self._event_text, json.dumps(payload, indent=2, ensure_ascii=False)) elif cmd == CMD_HEARTBEAT: data = payload.get("data", {}) self._devmgr.update_heartbeat(dev_serial, data) self._log_recv(topic, payload, f"heartbeat uptime={data.get('uptime', 0)}s") elif cmd == CMD_INITIALIZE: # V1.03: 设备上线初始化 data = payload.get("data", {}) dev_sn = data.get("dev_serial", dev_serial) model = data.get("model", "") hard_ver = data.get("hard_ver", "") soft_ver = data.get("soft_ver", "") extra = data.get("extra_info", {}) self._devmgr.mark_online(dev_sn, model=model, hard_ver=hard_ver, soft_ver=soft_ver, extra_info=extra) self._log_recv(topic, payload, f"initialize sn={dev_sn} model={model}") self._show_json({"online": dev_sn, "model": model, "extra": extra}) elif cmd == CMD_OTA_REPORT: # V1.08: OTA 进度/结果主动上报 (无 code 字段) data = payload.get("data", {}) self._ota_show_status(data) self._log_recv(topic, payload, f"ota_report stage={data.get('stage', '?')}") # 旧协议 / 无 cmd 字段 的消息(如 Initialize) elif "Method" in payload: self._log_recv(topic, payload, f"Method={payload.get('Method', '?')}") else: self._log_recv(topic, payload) # 自定义订阅消息 → 显示到自定义 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) # ================================================================ # 设备操作 # ================================================================ def _selected_dev(self) -> Optional[str]: items = self._dev_tree.selectedItems() return items[0].text(0).lstrip("●○ ") if items else None def _require_dev(self) -> Optional[str]: sn = self._selected_dev() if not sn: QMessageBox.warning(self, "提示", "请先选择一个设备") return None if not self._mqtt.connected: QMessageBox.warning(self, "提示", "请先连接 MQTT Broker") return None return sn def _send_cmd(self, cmd: str, data: Optional[dict] = None): sn = self._require_dev() if not sn: return try: mid = self._mqtt.send_command(sn, cmd, data) from dbn_mqtt_tool.protocol import build_request, get_topic_for_cmd msg = build_request(cmd, data) topic = get_topic_for_cmd(cmd, sn) self._log_send(f"[命令] {sn}", topic, msg) except Exception as e: QMessageBox.critical(self, "发送失败", str(e)) def _query_dev_info(self): self._send_cmd(CMD_DEV_INFO_QUERY) def _query_all(self): sn = self._require_dev() if not sn: return for cmd in [CMD_DEV_INFO_QUERY, CMD_SSC_NET_QUERY, CMD_IOT_NET_QUERY, CMD_IOT_TOPIC_QUERY]: try: self._mqtt.send_command(sn, cmd) except Exception: pass def _query_ssc_net(self): self._send_cmd(CMD_SSC_NET_QUERY) def _set_ssc_net(self): self._send_cmd(CMD_SSC_NET_SET, data_ssc_net_set( dev_ip=self._edit_ssc_ip.text(), subnet_mask=self._edit_ssc_mask.text(), route_ip=self._edit_ssc_gw.text(), lssc_ip=self._edit_ssc_lssc.text(), dns=self._edit_ssc_dns.text(), port=int(self._edit_ssc_port.text() or "0"), )) def _query_iot_net(self): self._send_cmd(CMD_IOT_NET_QUERY) def _set_iot_net(self): self._send_cmd(CMD_IOT_NET_SET, data_iot_net_set( host=self._edit_iot_host.text(), port=int(self._edit_iot_port.text() or "1883"), client_id=self._edit_iot_cid.text(), username=self._edit_iot_user.text(), password=self._edit_iot_pass.text(), )) def _query_iot_topic(self): self._send_cmd(CMD_IOT_TOPIC_QUERY) def _set_iot_topic(self): self._send_cmd(CMD_IOT_TOPIC_SET, data_iot_topic_set( client_id_enable=self._chk_cid.isChecked(), topic_pub=self._edit_topic_pub.text(), topic_sub=self._edit_topic_sub.text(), )) def _apply_ssc_net_data(self, data: dict): self._edit_ssc_ip.setText(str(data.get("dev_ip", ""))) self._edit_ssc_mask.setText(str(data.get("subnet_mask", ""))) self._edit_ssc_gw.setText(str(data.get("route_ip", ""))) self._edit_ssc_lssc.setText(str(data.get("lssc_ip", ""))) self._edit_ssc_dns.setText(str(data.get("dns", ""))) self._edit_ssc_port.setText(str(data.get("port", ""))) def _apply_iot_net_data(self, data: dict): self._edit_iot_host.setText(str(data.get("host", ""))) self._edit_iot_port.setText(str(data.get("port", ""))) self._edit_iot_cid.setText(str(data.get("client_id", ""))) self._edit_iot_user.setText(str(data.get("username", ""))) self._edit_iot_pass.setText(str(data.get("password", ""))) def _apply_iot_topic_data(self, data: dict): self._chk_cid.setChecked(bool(data.get("client_id_enable", True))) self._edit_topic_pub.setText(str(data.get("topic_pub", ""))) self._edit_topic_sub.setText(str(data.get("topic_sub", ""))) def _apply_report_config_data(self, data: dict): self._edit_report_sensor_type.setValue(int(data.get("sensor_type", 12))) self._edit_report_interval.setValue(int(data.get("interval", 5))) self._edit_report_timeout.setValue(int(data.get("timeout", 0))) self._chk_report_enable.setChecked(bool(data.get("enable", True))) self._chk_report_once.setChecked(bool(data.get("once", False))) self._chk_report_env.setChecked(bool(data.get("env_eval", False))) self._chk_report_ack.setChecked(bool(data.get("ack_required", False))) def _do_pwd_verify(self): pwd = self._edit_pwd.text() if not pwd: QMessageBox.warning(self, "提示", "请输入密码") return self._send_cmd(CMD_PWD_VERIFY, data_pwd_verify(pwd)) def _do_pwd_set(self): old = self._edit_pwd.text() new = self._edit_new_pwd.text() if not old or not new: QMessageBox.warning(self, "提示", "请输入旧密码和新密码") return self._send_cmd(CMD_PWD_SET, data_pwd_set(old, new)) def _do_factory_reset(self): r = QMessageBox.question(self, "确认", "确定要恢复出厂设置吗?\n此操作不可撤销!", QMessageBox.Yes | QMessageBox.No) if r == QMessageBox.Yes: self._send_cmd(CMD_FACTORY_RESET) def _do_device_reset(self): r = QMessageBox.question(self, "确认", "确定要复位设备吗?设备将重启。", QMessageBox.Yes | QMessageBox.No) if r == QMessageBox.Yes: self._send_cmd(CMD_DEVICE_RESET) def _query_report_config(self): self._send_cmd(CMD_REPORT_CONFIG) def _set_report_config(self): self._send_cmd(CMD_REPORT_CONFIG, data_report_config( sensor_type=self._edit_report_sensor_type.value(), enable=self._chk_report_enable.isChecked(), once=self._chk_report_once.isChecked(), env_eval=self._chk_report_env.isChecked(), interval=self._edit_report_interval.value(), ack_required=self._chk_report_ack.isChecked(), timeout=self._edit_report_timeout.value(), )) # ---- 脱机日志 (V1.06 事件流 / V1.07 快照流) ---- def _selected_log_stream(self) -> str: """当前 UI 选择的日志流 (event/snapshot)""" return self._combo_log_stream.currentData() or STREAM_EVENT def _query_log_stat(self): """log_stat: 日志统计/分页定位 (按当前流)""" stream = self._selected_log_stream() data = {"stream": STREAM_SNAPSHOT} if stream == STREAM_SNAPSHOT else None self._send_cmd(CMD_LOG_STAT, data) def _query_log(self): """log_query: 按全局序号分页拉取 (事件流 count≤4 / 快照流 count≤2)""" self._send_cmd(CMD_LOG_QUERY, data_log_query( start_seq=self._edit_log_start_seq.value(), count=self._edit_log_count.value(), stream=self._selected_log_stream(), )) def _page_log_prev(self): """上一页: 起始序号 = max(1, 当前起始 - 条数), 自动拉取""" start = self._edit_log_start_seq.value() count = self._edit_log_count.value() new_start = max(1, start - count) self._edit_log_start_seq.setValue(new_start) self._log(f"翻页: 上一页 → 起始序号 {new_start}") self._query_log() def _page_log_next(self): """下一页: 起始序号 = 当前页最后一条 seq + 1, 自动拉取""" if self._log_page_last_seq > 0: new_start = self._log_page_last_seq + 1 else: # 无有效页锚点 (拉取返回空), 按页大小步进 new_start = self._edit_log_start_seq.value() + self._edit_log_count.value() self._edit_log_start_seq.setValue(new_start) self._log(f"翻页: 下一页 → 起始序号 {new_start}") self._query_log() def _clear_log(self): """log_clear: 清除脱机日志 (按当前流, 审计留痕, 设备侧阻塞)""" stream = self._selected_log_stream() target = "传感快照" if stream == STREAM_SNAPSHOT else "脱机事件日志" block = "~45ms (逻辑清除+当前扇区)" if stream == STREAM_SNAPSHOT else "~2.8s (63 扇区擦除)" r = QMessageBox.question(self, "确认", f"确定要清除设备{target}吗?\n" f"清除动作本身会写入审计记录,不可撤销!\n" f"设备侧阻塞约 {block}。", QMessageBox.Yes | QMessageBox.No) if r == QMessageBox.Yes: self._send_cmd(CMD_LOG_CLEAR, data_log_clear(stream)) def _apply_log_stat_data(self, data: dict): stream = data.get("stream", "event") stream_name = "传感快照" if stream == "snapshot" else "事件日志" lines = [ f"日志统计[{stream_name}]: enabled={data.get('enabled', False)} boot_seq={data.get('boot_seq', 0)}", f"记录数: {data.get('count', 0)} / {data.get('capacity', 0)}", f"全局序号范围: {data.get('seq_first', 0)} ~ {data.get('seq_last', 0)}", f"提示: 拉取日志用 log_query, 起始序号填 seq_first" + (" (快照流 count 上限 2)" if stream == "snapshot" else ""), ] self._log_rec_text.setPlainText("\n".join(lines)) self._log(f"log_stat[{stream_name}]: count={data.get('count', 0)} seq={data.get('seq_first', 0)}~{data.get('seq_last', 0)}") def _apply_log_query_data(self, data: dict): records = data.get("records", []) start_seq = data.get("start_seq", 0) # V1.07: 记录为 {"seq":N,"hex":"..."} 原始字节; hex 长度 128=快照(64B) / 64=事件(32B) is_snap = bool(records) and len(records[0].get("hex", "")) == 128 stream_name = "传感快照" if is_snap else "事件日志" lines = [f"log_query[{stream_name}] start_seq={start_seq} → 返回 {len(records)} 条:"] for rec in records: hex_str = rec.get("hex", "") if is_snap: s = parse_snap_hex(hex_str) if "error" in s: lines.append(f" seq={rec.get('seq', 0)} 解析失败: {s['error']}") continue lines.append( f" seq={s['seq']} boot={s['boot_seq']} " f"boot+{s['ts_ms']}ms ({s['coil_count']}ch)" ) for ch in s["channels"]: mt = ch["misc_type"] misc_desc = SNAP_MISC_TYPE_DESC.get(mt, mt) lines.append( f" ch{ch['ch']}: {ch['freq_level']} " f"dir={ch['direction']} ftype={ch['freq_type']} " f"sens={ch['sensitivity']} cond={ch['condition']} " f"loop={'OK' if ch['loop_ok'] else '断'} " f"car={'有' if ch['has_car'] else '无'} " f"freq={ch['freq']}Hz Δ={ch['variation']} " f"[{misc_desc}] misc={ch['misc']}" ) else: e = parse_offlog_hex(hex_str) if "error" in e: lines.append(f" seq={rec.get('seq', 0)} 解析失败: {e['error']}") continue tname = OFFLOG_TYPE_NAMES.get(e["type"], f"0x{e['type']:02x}") desc = LOG_EVENT_TYPE_DESC.get(tname, tname) if e["unix_ts"]: ts_str = datetime.fromtimestamp(e["unix_ts"]).strftime("%Y-%m-%d %H:%M:%S") else: ts_str = f"boot+{e['ts_ms']}ms (未同步)" pdesc = offlog_payload_desc(e["type"], e["payload"]) lines.append( f" seq={e['seq']} boot={e['boot_seq']} " f"{ts_str} [{desc}] {pdesc}" ) if not records: lines.append(" (无记录 / 起始序号越界, 可先用 log_stat 确认范围)") self._log_rec_text.setPlainText("\n".join(lines)) self._log(f"log_query[{stream_name}]: start_seq={start_seq} got {len(records)}") # 更新翻页状态: 锚点 = 当前页最后一条 seq; 拉取成功即启用翻页按钮 if records: self._log_page_last_seq = records[-1].get("seq", 0) else: self._log_page_last_seq = 0 self._btn_log_prev.setEnabled(True) self._btn_log_next.setEnabled(True) # ================================================================ # UI 更新 # ================================================================ @Slot() def _refresh_devices(self): tree = self._dev_tree existing = {} for i in range(tree.topLevelItemCount()): item = tree.topLevelItem(i) sn = item.text(0).lstrip("●○ ") existing[sn] = item for sn, dev in self._devmgr.devices.items(): status = "●" if dev.online else "○" model = dev.model or "?" if sn in existing: existing[sn].setText(0, f"{status} {sn}") existing[sn].setText(1, model) del existing[sn] else: item = QTreeWidgetItem([f"{status} {sn}", model]) tree.addTopLevelItem(item) for sn, item in existing.items(): idx = tree.indexOfTopLevelItem(item) tree.takeTopLevelItem(idx) def _on_dev_select(self): sn = self._selected_dev() if sn and sn in self._devmgr.devices: dev = self._devmgr.devices[sn] info = { "dev_serial": dev.dev_serial, "model": dev.model, "hard_ver": dev.hard_ver, "soft_ver": dev.soft_ver, "product_code": dev.product_code, "net_enabled": dev.net_enabled, "iot_enabled": dev.iot_enabled, "bus": dev.bus, "online": dev.online, "last_seen": dev.last_seen.strftime("%Y-%m-%d %H:%M:%S") if dev.last_seen else "—", } self._show_json(info) def _show_json(self, data: dict): self._info_text.setPlainText(json.dumps(data, indent=2, ensure_ascii=False)) def _append_text(self, widget: QTextEdit, text: str): ts = datetime.now().strftime("%H:%M:%S") widget.append(f"--- {ts} ---\n{text}\n") widget.moveCursor(QTextCursor.End) def _log(self, msg: str): ts = datetime.now().strftime("%H:%M:%S") self._log_text.append(f"[{ts}] {msg}") self._log_text.moveCursor(QTextCursor.End) # ---- 详细收发日志 ---- def _log_send(self, label: str, topic: str, payload: dict): """记录发送消息: 标签 + topic + 精简 payload""" payload_str = json.dumps(payload, ensure_ascii=False) if len(payload_str) > 200: payload_str = payload_str[:200] + "..." self._log(f"{label} → {topic}") self._log(f" PASTE: {payload_str}") def _log_recv(self, topic: str, payload: dict, summary: str = ""): """记录接收消息: topic + summary + 精简 payload""" payload_str = json.dumps(payload, ensure_ascii=False) if len(payload_str) > 300: payload_str = payload_str[:300] + "..." self._log(f"← {topic}{' ' + summary if summary else ''}") self._log(f" RECV: {payload_str}") def main(): app = QApplication(sys.argv) app.setStyle("Fusion") win = MainWindow() win.show() sys.exit(app.exec()) if __name__ == "__main__": main()