Files
vd_960/DBNMQTTool/dbn_mqtt_tool/protocol.py
T
wangfq d39ab5ae89 feat(DBNMQTTool): 脱机日志支持快照流 stream=snapshot (2026-08-18)
- protocol.py: STREAM_EVENT/STREAM_SNAPSHOT 常量, data_log_query 加 stream 参数(快照流 count≤2), 新增 data_log_clear(stream), SNAP_MISC_TYPE_DESC
- main.py: 日志流选择器 + 三命令按流发请求 + log_stat/log_query 响应区分流(快照 channels 逐字段展示)
- 验证: protocol 断言 6 例 + offscreen MainWindow 快照/事件流展示实测全过
- devlog 置顶 + .gitignore (venv/)
2026-08-18 09:38:33 +08:00

285 lines
8.6 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
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 (64B×2 记录); 显式传 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 {}
# ============================================================
# 解析设备上报
# ============================================================
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