feat: DBNMQTTool — DLD960 IoT MQTT 设备管理工具 (Python/tkinter 跨平台)

基于《DLD960_IoT_MQTT协议.md》V1.00 实现

功能:
- MQTT Broker 连接管理
- 设备自动发现(订阅 dld960/+/dev/# 通配符)
- 设备信息查询(dev_info_query)
- 网络配置(SSC/IoT TCP)+ Topic 配置
- 实时线圈数据监控(loop_data)+ 事件上报(event_report)
- 心跳监控(heartbeat)
- 控制命令:密码验证/设置、出厂初始化、设备复位

项目结构:
- main.py          主窗口 (tkinter GUI)
- dbn_mqtt_tool/
  - protocol.py    DLD960 IoT MQTT 协议定义
  - mqtt_client.py MQTT 客户端封装 (paho-mqtt)
  - device_manager.py 设备发现与状态管理
This commit is contained in:
wangfq
2026-07-06 18:14:38 +08:00
parent 220f778117
commit a2cfe4602c
5 changed files with 997 additions and 0 deletions
+111
View File
@@ -0,0 +1,111 @@
"""
设备管理器 — 管理已发现的设备及其状态
"""
import threading
from typing import Callable, Optional
from dataclasses import dataclass, field
from datetime import datetime
@dataclass
class DeviceInfo:
dev_serial: str
model: str = ""
hard_ver: str = ""
soft_ver: str = ""
product_code: str = ""
net_enabled: bool = False
iot_enabled: bool = False
bus: dict = field(default_factory=dict)
last_seen: Optional[datetime] = None
online: bool = False
loop_data: Optional[dict] = None
events: list = field(default_factory=list)
class DeviceManager:
"""管理设备列表和状态"""
def __init__(self):
self._devices: dict[str, DeviceInfo] = {}
self._lock = threading.Lock()
self._callbacks: list[Callable] = []
def on_change(self, callback: Callable):
self._callbacks.append(callback)
def _notify(self):
for cb in self._callbacks:
try:
cb()
except Exception:
pass
def get_or_create(self, dev_serial: str) -> DeviceInfo:
with self._lock:
if dev_serial not in self._devices:
self._devices[dev_serial] = DeviceInfo(dev_serial=dev_serial)
return self._devices[dev_serial]
def update_from_dev_info(self, dev_serial: str, data: dict):
with self._lock:
dev = self.get_or_create(dev_serial)
dev.dev_serial = data.get("dev_serial", dev_serial)
dev.model = data.get("model", dev.model)
dev.hard_ver = data.get("hard_ver", dev.hard_ver)
dev.soft_ver = data.get("soft_ver", dev.soft_ver)
dev.product_code = data.get("product_code", dev.product_code)
sc = data.get("sub_code", {})
dev.net_enabled = sc.get("net", dev.net_enabled)
dev.iot_enabled = sc.get("iot", dev.iot_enabled)
dev.bus = data.get("bus", dev.bus)
dev.last_seen = datetime.now()
dev.online = True
self._notify()
def update_loop_data(self, dev_serial: str, data: dict):
with self._lock:
dev = self.get_or_create(dev_serial)
dev.loop_data = data
dev.last_seen = datetime.now()
dev.online = True
self._notify()
def update_event(self, dev_serial: str, data: dict):
with self._lock:
dev = self.get_or_create(dev_serial)
events = data.get("events", [])
for evt in events:
evt["_ts"] = datetime.now()
dev.events.insert(0, evt)
dev.events = dev.events[:100] # 保留最近 100 条
dev.last_seen = datetime.now()
dev.online = True
self._notify()
def update_heartbeat(self, dev_serial: str, data: dict):
with self._lock:
dev = self.get_or_create(dev_serial)
dev.last_seen = datetime.now()
dev.online = True
self._notify()
def set_offline(self):
"""标记所有设备离线(断开连接时调用)"""
with self._lock:
now = datetime.now()
for dev in self._devices.values():
if dev.last_seen and (now - dev.last_seen).seconds > 120:
dev.online = False
self._notify()
@property
def devices(self) -> dict[str, DeviceInfo]:
with self._lock:
return dict(self._devices)
@property
def online_devices(self) -> list[DeviceInfo]:
with self._lock:
return [d for d in self._devices.values() if d.online]
+141
View File
@@ -0,0 +1,141 @@
"""
MQTT 客户端封装
"""
import json
import time
import threading
from typing import Callable, Optional
from dataclasses import dataclass
import paho.mqtt.client as mqtt
from .protocol import parse_topic_dev_serial
@dataclass
class BrokerConfig:
host: str = "121.37.20.199"
port: int = 1883
username: str = ""
password: str = ""
client_id: str = ""
keepalive: int = 60
class MqttClient:
"""DLD960 IoT MQTT 客户端"""
def __init__(self):
self._client: Optional[mqtt.Client] = None
self._config = BrokerConfig()
self._connected = False
self._lock = threading.Lock()
self._message_callbacks: list[Callable[[str, str, dict], None]] = []
self._status_callbacks: list[Callable[[bool, str], None]] = []
@property
def connected(self) -> bool:
return self._connected
def on_status_change(self, callback: Callable[[bool, str], None]):
self._status_callbacks.append(callback)
def on_message(self, callback: Callable[[str, str, dict], None]):
"""回调参数: (topic, dev_serial, payload_dict)"""
self._message_callbacks.append(callback)
def _notify_status(self, connected: bool, msg: str):
for cb in self._status_callbacks:
try:
cb(connected, msg)
except Exception:
pass
def _notify_message(self, topic: str, payload: dict):
dev_serial = parse_topic_dev_serial(topic) or ""
for cb in self._message_callbacks:
try:
cb(topic, dev_serial, payload)
except Exception:
pass
# ---- MQTT callbacks ----
def _on_connect(self, client, userdata, flags, reason_code, properties=None):
rc = reason_code if isinstance(reason_code, int) else reason_code.value
if rc == 0:
self._connected = True
self._notify_status(True, "已连接到 Broker")
# 订阅所有设备的上报 topic
from .protocol import (
TOPIC_ALL_CONFIG_RESP, TOPIC_ALL_DATA_LOOP,
TOPIC_ALL_DATA_EVENT, TOPIC_ALL_STATUS
)
client.subscribe(TOPIC_ALL_CONFIG_RESP, qos=1)
client.subscribe(TOPIC_ALL_DATA_LOOP, qos=1)
client.subscribe(TOPIC_ALL_DATA_EVENT, qos=1)
client.subscribe(TOPIC_ALL_STATUS, qos=0)
else:
self._connected = False
self._notify_status(False, f"连接失败 (rc={rc})")
def _on_disconnect(self, client, userdata, flags, reason_code, properties=None):
self._connected = False
self._notify_status(False, "已断开连接")
def _on_message(self, client, userdata, msg: mqtt.MQTTMessage):
try:
payload = json.loads(msg.payload.decode("utf-8"))
except (json.JSONDecodeError, UnicodeDecodeError):
payload = {"_raw": msg.payload.hex()}
self._notify_message(msg.topic, payload)
# ---- Public API ----
def connect(self, config: Optional[BrokerConfig] = None):
if config:
self._config = config
if self._client:
self.disconnect()
client_id = self._config.client_id or f"dbn_mqtt_tool_{int(time.time())}"
# MQTT 5.0 用 CallbackAPIVersion.VERSION2
self._client = mqtt.Client(
client_id=client_id,
protocol=mqtt.MQTTv311,
)
if self._config.username:
self._client.username_pw_set(self._config.username, self._config.password)
self._client.on_connect = self._on_connect
self._client.on_disconnect = self._on_disconnect
self._client.on_message = self._on_message
self._client.connect_async(self._config.host, self._config.port, self._config.keepalive)
self._client.loop_start()
def disconnect(self):
if self._client:
self._client.loop_stop()
self._client.disconnect()
self._client = None
self._connected = False
def publish(self, topic: str, payload: dict, qos: int = 1):
"""发送 JSON 消息到指定 topic"""
if not self._client or not self._connected:
raise ConnectionError("MQTT 未连接")
data = json.dumps(payload, ensure_ascii=False)
self._client.publish(topic, data, qos=qos)
def send_command(self, dev_serial: str, cmd: str, data: Optional[dict] = None, qos: int = 1) -> int:
"""向设备发送命令,返回 msg_id"""
from .protocol import build_request, get_topic_for_cmd
msg = build_request(cmd, data)
topic = get_topic_for_cmd(cmd, dev_serial)
self.publish(topic, msg, qos=qos)
return msg["msg_id"]
+247
View File
@@ -0,0 +1,247 @@
"""
DLD960 IoT MQTT 协议定义
基于《DLD960_IoT_MQTT协议.md》V1.00
"""
import json
import time
from typing import Optional, Any
from dataclasses import dataclass, field, asdict
# ============================================================
# Topic 结构
# ============================================================
def topic_config_set(dev_serial: str) -> str:
"""服务器下发 — 设置配置"""
return f"dld960/{dev_serial}/srv/config/set"
def topic_config_query(dev_serial: str) -> str:
"""服务器下发 — 查询配置"""
return f"dld960/{dev_serial}/srv/config/query"
def topic_ctrl(dev_serial: str) -> str:
"""服务器下发 — 控制命令"""
return f"dld960/{dev_serial}/srv/ctrl"
def topic_config_resp(dev_serial: str) -> str:
"""设备上报 — 配置响应"""
return f"dld960/{dev_serial}/dev/config/resp"
def topic_data_loop(dev_serial: str) -> str:
"""设备上报 — 线圈数据"""
return f"dld960/{dev_serial}/dev/data/loop"
def topic_data_event(dev_serial: str) -> str:
"""设备上报 — 事件"""
return f"dld960/{dev_serial}/dev/data/event"
def topic_status(dev_serial: str) -> str:
"""设备上报 — 状态/心跳"""
return f"dld960/{dev_serial}/dev/status"
# 订阅通配符(监听所有设备)
TOPIC_ALL_CONFIG_RESP = "dld960/+/dev/config/resp"
TOPIC_ALL_DATA_LOOP = "dld960/+/dev/data/loop"
TOPIC_ALL_DATA_EVENT = "dld960/+/dev/data/event"
TOPIC_ALL_STATUS = "dld960/+/dev/status"
# ============================================================
# 命令枚举
# ============================================================
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"
# 设备上报
CMD_LOOP_DATA = "loop_data"
CMD_EVENT_REPORT = "event_report"
CMD_HEARTBEAT = "heartbeat"
# 配置类命令 → 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}
# ============================================================
# 错误码
# ============================================================
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"]
# ============================================================
# 消息构建
# ============================================================
_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"""
if cmd in CONFIG_SET_COMMANDS:
return topic_config_set(dev_serial)
elif cmd in CONFIG_QUERY_COMMANDS:
return topic_config_query(dev_serial)
elif cmd in CTRL_COMMANDS:
return topic_ctrl(dev_serial)
else:
return topic_config_set(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 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