Files
vd_960/DBNMQTTool/dbn_mqtt_tool/protocol.py
T
wangfq 51298693da fix(vd960DBN): OTA 刷写成功误判未启动 — 状态回 idle + 补 ota_report 主动上报
现场: 远程平台发起 OTA, 设备物理刷写成功 (70块全ACK, flash DONE),
但平台 ota_status 轮询只见 state=ready → 误判'刷写未启动'。

根因:
1. ota_flash_done 刷写成功后台侧状态置 READY (与协议状态机
   'FLASHING──成功──▶IDLE(清槽/保留)' 不符), 平台无法区分
   '待刷 ready' 与 '已刷完' → 持续 ready 判未启动
2. 固件漏实现协议 §5.5 ota_report 主动上报 (done/failed 重发 3次×5s),
   平台收不到成功信号, 只能靠 ota_status 轮询兜底

修复 (协议 V1.09 + 固件):
- 协议: 刷写成功状态明确回 idle (镜像保留, last_result=0, 可重刷);
  ota_report 升级为刷写结果主依据; 平台判定指引 (idle+size>0+last_result=0
  =成功; 不得以轮询未见 flashing 或持续 ready 判未启动)
- ota_srv.c: ota_flash_done → OTA_STATE_IDLE (镜像保留) + 上报 done;
  ota_flash_fail → 补 ota_report failed (保留 event_report 告警);
  ota_cmd_begin 兼容 idle+size/crc32 一致 → 直接回 ready (免下载重刷)
- iot_mqtt_srv.c: 新增 ota_report 上报状态机 (立即首发 + 5s×3 重发,
  同 msg_id/ts), 共享 _iot_pub_payload (event_report 复用, RAM 零新增)
- 单测: test_flash_flow/test_flash_ready_timeout 断言 ota_report 上报,
  state 断言 READY→IDLE; 新增 test_begin_reflash; 10/10 全过 + 22 py 断言
2026-08-20 18:56:21 +08:00

463 lines
16 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
DLD960 IoT MQTT 协议定义
基于《DLD960_IoT_MQTT协议.md》V1.03
双主题结构:{sn}/srv(服务器下发)、{sn}/dev(设备上报)
消息类型由 JSON `cmd` 字段区分。
"""
import json
import time
import struct
import zlib
from typing import Optional, Any
from dataclasses import dataclass, field, asdict
# ============================================================
# Topic 结构(V1.01 — 双主题)
# ============================================================
def topic_down(dev_serial: str) -> str:
"""服务器下发 — 设备订阅此主题接收所有命令"""
return f"dld960/{dev_serial}/srv"
def topic_up(dev_serial: str) -> str:
"""设备上报 — 设备发布到此主题(数据、事件、心跳、响应)"""
return f"dld960/{dev_serial}/dev"
# 保持向后兼容的别名(4.7 iot_topic_set 仍使用 topic_pub / topic_sub
def topic_dev_pub(dev_serial: str) -> str:
"""设备发布主题(= topic_up"""
return topic_up(dev_serial)
def topic_dev_sub(dev_serial: str) -> str:
"""设备订阅主题(= topic_down"""
return topic_down(dev_serial)
# 服务器端通配符订阅(监听所有设备)
TOPIC_ALL_DEVICE_UP = "dld960/+/dev"
# ============================================================
# 命令枚举
# ============================================================
CMD_DEV_SERIAL_SET = "dev_serial_set"
CMD_DEV_INFO_QUERY = "dev_info_query"
CMD_SSC_NET_SET = "ssc_net_set"
CMD_SSC_NET_QUERY = "ssc_net_query"
CMD_IOT_NET_SET = "iot_net_set"
CMD_IOT_NET_QUERY = "iot_net_query"
CMD_IOT_TOPIC_SET = "iot_topic_set"
CMD_IOT_TOPIC_QUERY = "iot_topic_query"
CMD_PWD_VERIFY = "pwd_verify"
CMD_PWD_SET = "pwd_set"
CMD_FACTORY_RESET = "factory_reset"
CMD_DEVICE_RESET = "device_reset"
CMD_LOOP_PARAM_SET = "loop_param_set"
CMD_LOOP_PARAM_QUERY = "loop_param_query"
CMD_REPORT_CONFIG = "report_config"
# V1.06: 脱机事件日志
CMD_LOG_STAT = "log_stat"
CMD_LOG_QUERY = "log_query"
CMD_LOG_CLEAR = "log_clear"
# 设备上报
CMD_LOOP_DATA = "loop_data"
CMD_EVENT_REPORT = "event_report"
CMD_HEARTBEAT = "heartbeat"
CMD_INITIALIZE = "initialize" # V1.03: 设备上电初始化
# 配置类命令 → topic
CONFIG_SET_COMMANDS = {CMD_DEV_SERIAL_SET, CMD_SSC_NET_SET, CMD_IOT_NET_SET, CMD_IOT_TOPIC_SET, CMD_LOOP_PARAM_SET}
CONFIG_QUERY_COMMANDS = {CMD_DEV_INFO_QUERY, CMD_SSC_NET_QUERY, CMD_IOT_NET_QUERY, CMD_IOT_TOPIC_QUERY, CMD_LOOP_PARAM_QUERY}
CTRL_COMMANDS = {CMD_PWD_VERIFY, CMD_PWD_SET, CMD_FACTORY_RESET, CMD_DEVICE_RESET, CMD_REPORT_CONFIG}
LOG_COMMANDS = {CMD_LOG_STAT, CMD_LOG_QUERY, CMD_LOG_CLEAR} # V1.06 脱机日志
# V1.07: log_* 命令日志流 (与 BLE 0x28/0x2A 同语义)
STREAM_EVENT = "event"
STREAM_SNAPSHOT = "snapshot"
# ============================================================
# 错误码
# ============================================================
ERR_SUCCESS = 0
ERR_PARAM_ERROR = 1
ERR_PWD_FAILED = 2
ERR_DEVICE_BUSY = 3
ERR_UNSUPPORTED = 4
ERR_INTERNAL = 5
ERR_DATA_TOO_LONG = 6
ERROR_MSGS = {
0: "成功",
1: "参数错误",
2: "密码验证失败",
3: "设备忙",
4: "不支持的命令",
5: "内部错误",
6: "数据超长",
}
# ============================================================
# 频率档位 / 输出模式 / 事件类型
# ============================================================
FREQ_LEVELS = ["high", "mid_high", "mid_low", "low"]
OUTPUT_MODES = ["exist", "enter_pulse", "leave_pulse", "direction"]
EVENT_TYPES = ["car_enter", "car_leave", "loop_cut", "loop_restore"]
# 脱机事件日志类型描述 (log_query 响应 records[].type, MQTT V1.06)
LOG_EVENT_TYPE_DESC = {
"boot": "上电/复位",
"iot_connect": "MQTT TCP 连接成功",
"iot_ready": "MQTT 订阅完成→发 initialize",
"iot_disconnect": "MQTT 断连",
"iot_reconn": "重连退避",
"evt_retry": "event ACK 超时重发",
"evt_giveup": "event 重试耗尽挂起",
"coil": "线圈事件",
"time_anchor": "时钟同步锚点",
"log_clear": "日志清除(审计)",
"unknown": "未知类型",
}
# 传感快照 channels[].misc_type 描述 (log_query stream=snapshot, MQTT V1.07)
SNAP_MISC_TYPE_DESC = {
"time": "时间量(50ms)",
"cut_count": "线圈断开次数",
"flow_count": "车流量",
"relay_count": "继电器输出次数",
}
# ============================================================
# 消息构建
# ============================================================
_msg_id_counter = 0
def next_msg_id() -> int:
global _msg_id_counter
_msg_id_counter += 1
return _msg_id_counter
def build_request(cmd: str, data: Optional[dict] = None) -> dict:
"""构建通用请求消息"""
msg: dict[str, Any] = {
"msg_id": next_msg_id(),
"cmd": cmd,
"ts": int(time.time()),
}
if data is not None:
msg["data"] = data
return msg
def build_response(cmd: str, msg_id: int, code: int = 0, msg_text: str = "success", data: Optional[dict] = None) -> dict:
"""构建响应消息"""
resp: dict[str, Any] = {
"msg_id": msg_id,
"cmd": cmd,
"ts": int(time.time()),
"code": code,
"msg": msg_text,
}
if data is not None:
resp["data"] = data
return resp
def get_topic_for_cmd(cmd: str, dev_serial: str) -> str:
"""获取命令对应的下发 topic。所有命令统一使用 topic_down"""
return topic_down(dev_serial)
# ============================================================
# 特定命令的 data 构建器
# ============================================================
def data_dev_serial_set(dev_serial: str) -> dict:
return {"dev_serial": dev_serial}
def data_ssc_net_set(dev_ip: str, subnet_mask: str, route_ip: str,
lssc_ip: str, dns: str, port: int) -> dict:
return {
"dev_ip": dev_ip,
"subnet_mask": subnet_mask,
"route_ip": route_ip,
"lssc_ip": lssc_ip,
"dns": dns,
"port": port,
}
def data_iot_net_set(host: str, port: int = 1883,
client_id: str = "", username: str = "", password: str = "") -> dict:
return {
"host": host,
"port": port,
"client_id": client_id,
"username": username,
"password": password,
}
def data_iot_topic_set(client_id_enable: bool, topic_pub: str, topic_sub: str) -> dict:
return {
"client_id_enable": client_id_enable,
"topic_pub": topic_pub,
"topic_sub": topic_sub,
}
def data_pwd_verify(password: str) -> dict:
return {"password": password}
def data_pwd_set(old_password: str, new_password: str) -> dict:
return {"old_password": old_password, "new_password": new_password}
def data_loop_param_set(channels: list[dict], auto_mode: bool = False) -> dict:
return {"auto_mode": auto_mode, "channels": channels}
def data_report_config(sensor_type: int = 12, enable: bool = True,
once: bool = False, env_eval: bool = False,
interval: int = 5, ack_required: bool = False,
timeout: int = 0) -> dict:
return {
"sensor_type": sensor_type,
"enable": enable,
"once": once,
"env_eval": env_eval,
"interval": interval,
"ack_required": ack_required,
"timeout": timeout,
}
def data_log_query(start_seq: int = 1, count: int = 4, stream: str = STREAM_EVENT) -> dict:
"""V1.07 log_query 分页拉取脱机日志
stream=event: 事件流, count 上限 4 (设备侧超限按 4 处理)
stream=snapshot: 快照流, count 上限 2 (hex 原始字节上报); 显式传 stream 字段
"""
d: dict[str, Any] = {
"start_seq": max(1, int(start_seq)),
"count": min(4, max(1, int(count))),
}
if stream == STREAM_SNAPSHOT:
d["stream"] = STREAM_SNAPSHOT
d["count"] = min(2, max(1, int(count)))
return d
def data_log_clear(stream: str = STREAM_EVENT) -> dict:
"""V1.07 log_clear 清除脱机日志 (审计留痕)
stream=snapshot: 快照流清除 (~45ms, 审计写事件流)
event 流省略 data (设备缺省 event, 兼容老固件)
"""
if stream == STREAM_SNAPSHOT:
return {"stream": STREAM_SNAPSHOT}
return {}
# ============================================================
# OTA 远程升级 (V1.08: Loop MCU 先存后刷, ROADMAP P1.4 ①)
# 协议依据《DLD960_IoT_MQTT协议.md》§4.19~4.24 / §5.5
# ============================================================
CMD_OTA_BEGIN = "ota_begin"
CMD_OTA_DATA = "ota_data"
CMD_OTA_END = "ota_end"
CMD_OTA_ABORT = "ota_abort"
CMD_OTA_FLASH = "ota_flash"
CMD_OTA_STATUS = "ota_status"
CMD_OTA_REPORT = "ota_report" # dev→srv 主动上报
OTA_TARGET_LOOP = "loop"
OTA_CHUNK_SIZE = 256 # 单片 256B(协议常量, 不接受协商)
OTA_MAX_SIZE = 96 * 1024 # 镜像上限 96KB (Slot 100KB - 4KB 边界余量)
OTA_SLOT_A = "a"
OTA_SLOT_B = "b"
OTA_STATE_DESC = {
"idle": "空闲(镜像保留,已刷写完成则 idle+size>0)",
"downloading": "下载中",
"ready": "校验通过(可刷写)",
"flashing": "刷写中",
"flash_failed": "刷写失败",
"aborted": "已中止",
}
OTA_STAGE_DESC = {
"begin": "会话开启",
"downloading": "分片落盘",
"ready": "校验通过",
"flashing": "刷写中",
"done": "刷写成功(Loop 已重启)",
"failed": "刷写失败",
}
OTA_ERROR_DESC = {
0x1001: "启动帧无响应",
0x1002: "地址帧错误",
0x1003: "数据块 ACK 超限",
0x1004: "全镜像校验失败",
0x1005: "安全窗口拒绝后强制失败",
}
def ota_crc32(data: bytes) -> int:
"""CRC-32/ISO-HDLC (协议 §4.19.1; zlib.crc32 即此算法, 设备查表法一致)"""
return zlib.crc32(data) & 0xFFFFFFFF
def ota_split_bin(data: bytes, chunk_size: int = OTA_CHUNK_SIZE) -> list[tuple[int, int, str]]:
"""bin → 分片列表 [(offset, crc32, hex_str), ...] (单片 256B, hex 512 字符)"""
if len(data) > OTA_MAX_SIZE:
raise ValueError(f"镜像 {len(data)}B 超上限 {OTA_MAX_SIZE}B (96KB)")
return [(off, ota_crc32(data[off:off + chunk_size]), data[off:off + chunk_size].hex())
for off in range(0, len(data), chunk_size)]
def data_ota_begin(size: int, crc32: int, version: str = "",
target: str = OTA_TARGET_LOOP, force: bool = False) -> dict:
"""ota_begin 请求 data (开启会话 / 断点续传定位)"""
return {"target": target, "size": size, "crc32": crc32,
"version": version, "force": force}
def data_ota_data(offset: int, crc32: int, data_hex: str,
target: str = OTA_TARGET_LOOP) -> dict:
"""ota_data 请求 data (分片下发, offset 256 对齐)"""
return {"target": target, "offset": offset, "crc32": crc32, "data": data_hex}
def data_ota_end(crc32: int, target: str = OTA_TARGET_LOOP) -> dict:
"""ota_end 请求 data (结束下载, 全镜像 CRC32 复核)"""
return {"target": target, "crc32": crc32}
def data_ota_abort(target: str = OTA_TARGET_LOOP) -> dict:
"""ota_abort 请求 data (中止会话, 释放暂存)"""
return {"target": target}
def data_ota_flash(slot: str = OTA_SLOT_A, force: bool = False,
target: str = OTA_TARGET_LOOP) -> dict:
"""ota_flash 请求 data (触发本地 ISP 刷写, 仅 ready 态)"""
return {"target": target, "slot": slot, "force": force}
# ============================================================
# 解析设备上报
# ============================================================
def parse_topic_dev_serial(topic: str) -> Optional[str]:
"""从 topic 中提取设备序列码"""
parts = topic.split("/")
if len(parts) >= 2 and parts[0] == "dld960":
return parts[1]
return None
# ============================================================
# 脱机日志 hex 解析 (V1.07: log_query 响应 records[].hex 原始字节)
# 字段表与《DLD960 BLE 协议》§6.4 (SnapRec) / §7 (OfflogEvt) 一致
# ============================================================
# OfflogEvt 事件类型 (BLE 协议 §7 事件类型表)
OFFLOG_TYPE_NAMES = {
0x01: "boot", 0x10: "iot_connect", 0x11: "iot_ready", 0x12: "iot_disconnect",
0x13: "iot_reconn", 0x30: "evt_retry", 0x31: "evt_giveup", 0x40: "coil",
0x50: "time_anchor", 0x70: "log_clear",
}
_OFFLOG_SUB_NAMES = {1: "car_enter", 2: "car_leave", 3: "loop_cut", 4: "loop_restore"}
_OFFLOG_DISC_REASONS = {1: "断开", 2: "超时", 3: "CONNACK拒绝", 4: "连接超时"}
def parse_offlog_hex(h: str) -> dict:
"""OfflogEvt 32B hex → dict (magic/type/seq/ts_ms/unix_ts/boot_seq/payload)"""
b = bytes.fromhex(h)
if len(b) != 32:
return {"error": f"bad len {len(b)}"}
magic, etype, elen, flags, seq, ts_ms, unix_ts, boot_seq, rsvd, payload = \
struct.unpack('<BBBBIIIHH12s', b)
return {"magic": magic, "type": etype, "len": elen, "flags": flags, "seq": seq,
"ts_ms": ts_ms, "unix_ts": unix_ts, "boot_seq": boot_seq, "payload": payload}
def offlog_payload_desc(etype: int, payload: bytes) -> str:
"""事件 payload 人类可读描述 (BLE 协议 §7 payload 定义, 大端)"""
if etype == 0x01: # boot: 复位原因寄存器
rst = int.from_bytes(payload[0:4], 'big') if len(payload) >= 4 else 0
bits = [n for n, b in [("LPWR", 31), ("WWDG", 30), ("IWDG", 29),
("SFT", 28), ("POR", 27), ("NRST", 26)] if rst & (1 << b)]
return f"rst=0x{rst:08X} ({'+'.join(bits) if bits else 'none'})"
if etype == 0x12: # iot_disconnect
r = payload[0] if payload else 0
return f"reason={r} ({_OFFLOG_DISC_REASONS.get(r, '?')})"
if etype == 0x13: # iot_reconn
return f"backoff_ms={int.from_bytes(payload[0:4], 'big')}" if len(payload) >= 4 else ""
if etype == 0x30: # evt_retry
mid = int.from_bytes(payload[0:4], 'big') if len(payload) >= 4 else 0
retry = payload[4] if len(payload) >= 5 else 0
return f"msg_id={mid} retry={retry}"
if etype == 0x31: # evt_giveup
mid = int.from_bytes(payload[0:4], 'big') if len(payload) >= 4 else 0
return f"msg_id={mid}"
if etype == 0x40: # coil
sub = _OFFLOG_SUB_NAMES.get(payload[0], payload[0]) if payload else "?"
ch = payload[1] if len(payload) >= 2 else 0
val = int.from_bytes(payload[2:6], 'big') if len(payload) >= 6 else 0
return f"sub={sub} ch={ch} value={val}"
if etype in (0x10, 0x11, 0x50, 0x70): # 无 payload
return ""
return payload.hex() if payload else ""
def parse_snap_hex(h: str) -> dict:
"""SnapRec 64B hex → dict (seq/boot_seq/ts_ms/coil_count/channels[])"""
b = bytes.fromhex(h)
if len(b) != 64:
return {"error": f"bad len {len(b)}"}
magic, slen, flags, rsvd, seq, ts_ms, boot_seq, rsvd2, coils = \
struct.unpack('<BBBBIIHH48s', b)
coil_count = min(slen // 12, 4)
channels = []
for c in range(coil_count):
p = coils[c * 12:(c + 1) * 12]
cfg, cond = p[0], p[1]
freq = p[2] | (p[3] << 8) | (p[4] << 16)
var = p[5] | (p[6] << 8) | (p[7] << 16)
if var & 0x800000:
var -= 0x1000000
misc = int.from_bytes(p[8:12], 'little')
channels.append({
"ch": c + 1,
"freq_level": ("high", "mid_high", "mid_low", "low")[(cfg >> 6) & 3],
"direction": (cfg >> 5) & 1,
"freq_type": (cfg >> 4) & 1,
"sensitivity": cfg & 0x0F,
"condition": (cond >> 4) & 0x0F,
"loop_ok": not ((cond >> 3) & 1),
"has_car": bool((cond >> 2) & 1),
"misc_type": ("time", "cut_count", "flow_count", "relay_count")[cond & 3],
"freq": freq, "variation": var, "misc": misc,
})
return {"seq": seq, "boot_seq": boot_seq, "ts_ms": ts_ms,
"coil_count": coil_count, "channels": channels}