vd960DBN:
- dev_initialize_pub() 改为 V1.03 格式 (msg_id/cmd/ts/data/extra_info)
- iot_mqtt_srv: 订阅改双主题 dld960/{sn}/srv, 响应 topic 改 dev
- SUBACK 后自动发 initialize 告知服务器上线
- net_srv.h: 加 dev_initialize_pub 声明
DBNMQTTool:
- protocol.py: 加 CMD_INITIALIZE
- device_manager: DeviceInfo 加 extra_info, 加 mark_online()
- main.py: 接收 initialize 消息, 自动标记设备上线
- docs: 协议 V1.03 + devlog V3.1
124 lines
4.1 KiB
Python
124 lines
4.1 KiB
Python
"""
|
|
设备管理器 — 管理已发现的设备及其状态
|
|
"""
|
|
|
|
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)
|
|
extra_info: dict = field(default_factory=dict) # from initialize 消息
|
|
|
|
|
|
class DeviceManager:
|
|
"""管理设备列表和状态"""
|
|
|
|
def __init__(self):
|
|
self._devices: dict[str, DeviceInfo] = {}
|
|
self._lock = threading.RLock() # 可重入锁: update_xxx 内部调用 get_or_create
|
|
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 mark_online(self, dev_serial: str, extra_info: Optional[dict] = None):
|
|
"""V1.03: 设备上线初始化标记"""
|
|
with self._lock:
|
|
dev = self.get_or_create(dev_serial)
|
|
dev.online = True
|
|
dev.last_seen = datetime.now()
|
|
if extra_info:
|
|
dev.extra_info = extra_info
|
|
dev.soft_ver = extra_info.get("version", dev.soft_ver)
|
|
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]
|