feat(DBNMQTTool): OTA 页签 + 协议单测 (MQTT V1.08 工具先行)

- protocol.py: OTA 模块 — 命令常量/ota_crc32(ISO-HDLC)/ota_split_bin(256B×96KB)/
  data_ota_* 构建器/状态·阶段·错误码描述
- main.py: OTA 页签 — bin 选择(≤96KB+CRC32 显示)/版本·force·slot 参数/
  下载(OtaDownloadThread: begin→data×N→end, 断点续传·缺片重定位·CRC 重发)/
  刷写(二次确认)/中止/状态查询; 进度条 + ota_report/ota_status 实时显示
- tests/test_ota_protocol.py: 22 断言 — CRC32 标准向量/分片边界/构建器/
  OtaDeviceSim 设备状态机(顺序·幂等·乱序·单片CRC·全镜像复核·续传)/完整下载流
- devlog 置顶条目
This commit is contained in:
wangfq
2026-08-20 11:53:50 +08:00
parent 2ccf0804a8
commit 6c1f276002
4 changed files with 755 additions and 1 deletions
+325 -1
View File
@@ -7,6 +7,8 @@ DBN MQTT Tool — DLD960 IoT MQTT 设备管理工具
import sys
import os
import json
import queue
import time
from datetime import datetime
from typing import Optional
@@ -18,8 +20,9 @@ from PySide6.QtWidgets import (
QTabWidget, QTextEdit, QGroupBox, QGridLayout, QCheckBox,
QSplitter, QMessageBox, QHeaderView, QFrame,
QComboBox, QSpinBox, QPlainTextEdit,
QFileDialog, QProgressBar,
)
from PySide6.QtCore import Qt, QTimer, Signal, Slot
from PySide6.QtCore import Qt, QTimer, Signal, Slot, QThread
from PySide6.QtGui import QFont, QColor, QTextCursor
import paho.mqtt.client as mqtt
@@ -39,8 +42,128 @@ from dbn_mqtt_tool.protocol import (
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)
@@ -91,6 +214,7 @@ class MainWindow(QMainWindow):
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")
@@ -429,6 +553,192 @@ class MainWindow(QMainWindow):
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)))
# ================================================================
# 模拟设备上报
# ================================================================
@@ -946,7 +1256,15 @@ class MainWindow(QMainWindow):
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:
@@ -982,6 +1300,12 @@ class MainWindow(QMainWindow):
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', '?')}")