- 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 置顶条目
463 lines
16 KiB
Python
463 lines
16 KiB
Python
"""
|
||
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": "空闲",
|
||
"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}
|