modbus_rtu.py 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555
  1. from typing import Optional
  2. import threading
  3. # ========== 地址配置协议 ==========
  4. def build_broadcast_query() -> bytes:
  5. """构建广播查询指令
  6. DTU定时发送广播指令查询未配置地址的从机
  7. 格式: 00 40 00 40 00 (无CRC, 固定5字节)
  8. """
  9. return bytes([0x00, 0x40, 0x00, 0x40, 0x00])
  10. def parse_device_response(data: bytes) -> Optional[dict]:
  11. """解析从机应答
  12. 从机应答格式: 00 + 41 + 12个UID + CRC_L + CRC_H
  13. 例如: 00 41 18 00 40 00 14 00 00 59 59 54 30 56 AE 34
  14. """
  15. if len(data) < 16:
  16. return {'error': '数据长度不足', 'raw_data': data.hex()}
  17. if data[1] != 0x41:
  18. return {'error': f"非预期功能码: {data[1]:#x}", 'raw_data': data.hex()}
  19. uid_bytes = data[2:14]
  20. uid_hex = uid_bytes.hex()
  21. # CRC 校验失败则拒绝(传入完整 16 字节帧,最后两字节为 CRC)
  22. if not verify_crc16(data):
  23. return {'error': 'CRC校验失败', 'raw_data': data.hex()}
  24. # UID 有效性检查:拒绝明显是噪声或碰撞产生的数据
  25. zero_count = uid_bytes.count(0)
  26. if zero_count > 6:
  27. return {'error': f"UID中零字节过多({zero_count}/12)", 'raw_data': data.hex()}
  28. # 检查碰撞:UID 中不应有大量重复字节
  29. from collections import Counter
  30. counts = Counter(uid_bytes)
  31. if counts.most_common(1)[0][1] > 8:
  32. return {'error': f"UID中单字节重复过多({counts.most_common(1)})", 'raw_data': data.hex()}
  33. return {
  34. 'function_code': 0x41,
  35. 'uid': uid_hex,
  36. 'uid_readable': ':'.join(f'{b:02x}' for b in uid_bytes),
  37. 'source_address': data[0],
  38. 'raw_data': data.hex()
  39. }
  40. def build_confirm_address(device_address: int, uid: bytes) -> bytes:
  41. """构建确认地址指令
  42. DTU遍历已存储的UID和地址,如果存在匹配的UID则发送确认指令
  43. 格式: add + 44 + 00 + 12个UID + CRC_L + CRC_H
  44. """
  45. if len(uid) != 12:
  46. raise ValueError(f"UID长度必须是12字节,当前: {len(uid)}")
  47. data = bytes([device_address, 0x44, 0x00]) + uid
  48. return data + calculate_crc16(data)
  49. def build_assign_address(uid: bytes, address: int) -> bytes:
  50. """构建分配地址指令
  51. 如果DTU中没有该UID的记录,则分配新地址
  52. 格式: 00 + 42 + 12个UID + add + CRC_L + CRC_H
  53. """
  54. if len(uid) != 12:
  55. raise ValueError(f"UID长度必须是12字节,当前: {len(uid)}")
  56. if address < 1 or address > 247:
  57. raise ValueError(f"设备地址必须在1-247之间,当前: {address}")
  58. data = bytes([0x00, 0x42]) + uid + bytes([address])
  59. return data + calculate_crc16(data)
  60. def parse_address_assignment_response(data: bytes) -> dict:
  61. """解析地址分配响应
  62. 模块应答格式: add + 43 + 00 + CRC_L + CRC_H
  63. 例如: 01 43 00 11 30
  64. """
  65. if len(data) < 5:
  66. return {'error': '数据长度不足', 'raw_data': data.hex()}
  67. device_address = data[0]
  68. if data[1] != 0x43:
  69. return {'error': f"非预期功能码: {data[1]:#x}", 'raw_data': data.hex()}
  70. if not verify_crc16(data):
  71. logger.warning(f"CRC校验失败: {data.hex()}")
  72. return {
  73. 'device_address': device_address,
  74. 'function_code': 0x43,
  75. 'raw_data': data.hex()
  76. }
  77. class AddressConfigProtocol:
  78. """地址配置协议处理器"""
  79. def __init__(self, serial_port):
  80. self.serial = serial_port
  81. self.default_timeout = 5.0 # 设备广播响应可能较慢,给足时间
  82. self.stored_devices = {}
  83. self.config_lock = threading.RLock() # 地址配置期间占用,避免与轮询冲突
  84. def add_stored_device(self, uid: str, address: int):
  85. """添加已存储的设备"""
  86. self.stored_devices[uid] = address
  87. def get_stored_devices(self) -> dict:
  88. """获取已存储的设备列表"""
  89. return self.stored_devices.copy()
  90. def broadcast_query(self, timeout: float = None) -> list:
  91. """send broadcast query and collect all valid 00 41 responses"""
  92. import time
  93. with self.config_lock:
  94. request = build_broadcast_query()
  95. logger.info(f"send broadcast: {request.hex()}")
  96. timeout = timeout or self.default_timeout
  97. # 先清空输入缓冲区,避免残留数据干扰
  98. self.serial.flush_input()
  99. time.sleep(0.05)
  100. # 发送广播并等待响应;所有串口操作统一走发送队列,
  101. # 这里不再额外直接读串口,避免在收集窗口内把后续队列命令的响应抢走。
  102. response = self.serial.send_and_wait(request, timeout=timeout, min_response_bytes=16)
  103. responses = []
  104. if response and len(response) >= 16:
  105. logger.info(f"got raw response ({len(response)}B): {response.hex()}")
  106. # RS485 半双工会回显已发送的数据,导致 response 开头可能有 5 字节回显
  107. # 在 response 中搜索所有 00 41 标记,尝试从中提取 16 字节帧解析
  108. for i in range(len(response) - 15 + 1):
  109. if response[i] == 0x00 and response[i + 1] == 0x41:
  110. chunk = response[i:i + 16]
  111. if len(chunk) >= 16:
  112. parsed = parse_device_response(chunk)
  113. if "error" not in parsed:
  114. if not any(p.get('uid') == parsed['uid'] for p in responses):
  115. responses.append(parsed)
  116. else:
  117. logger.info(f"no response or too short: {len(response) if response else 0}B")
  118. logger.info(f"broadcast done, got {len(responses)} valid unique responses")
  119. return responses
  120. def _next_available_address(self) -> int:
  121. """从 1~247 中找一个未被 stored_devices 占用的地址"""
  122. used = set(self.stored_devices.values())
  123. for addr in range(1, 248):
  124. if addr not in used:
  125. return addr
  126. raise RuntimeError("Modbus 地址 1~247 已全部占用")
  127. def process_responses(self, responses: list) -> list:
  128. """处理广播查询响应:给每个应答设备分配地址(优先使用已记录地址)"""
  129. import time
  130. with self.config_lock:
  131. results = []
  132. for resp in responses:
  133. if 'error' in resp:
  134. results.append({'status': 'error', 'message': resp['error']})
  135. continue
  136. uid = resp['uid']
  137. src_addr = resp.get('source_address', 0)
  138. device_address = self.stored_devices.get(uid)
  139. # 优先使用地址表中已保存的地址;没有则分配下一个可用地址
  140. if device_address:
  141. new_address = device_address
  142. logger.info(f"UID={uid} 源地址={src_addr},地址表存在记录,分配已有地址={new_address}")
  143. else:
  144. new_address = self._next_available_address()
  145. logger.info(f"UID={uid} 源地址={src_addr},地址表无记录,分配新地址={new_address}")
  146. request = build_assign_address(bytes.fromhex(uid), new_address)
  147. logger.info(f"发送分配地址帧: {request.hex()}")
  148. self.serial.flush_input()
  149. time.sleep(0.05)
  150. success, msg = self.serial.send_raw(request)
  151. if not success:
  152. logger.error(f"UID={uid} 分配地址帧发送失败: {msg}")
  153. results.append({
  154. 'status': 'error',
  155. 'uid': uid,
  156. 'message': f"send failed: {msg}"
  157. })
  158. continue
  159. # 等待并解析从机的 43 确认响应
  160. time.sleep(0.1)
  161. ack = self.serial.send_and_wait(b'', timeout=0.5, min_response_bytes=5)
  162. ack_parsed = None
  163. if ack and len(ack) >= 5:
  164. # 过滤回显/噪声,找到 43 响应帧
  165. for i in range(len(ack) - 4):
  166. if ack[i + 1] == 0x43:
  167. chunk = ack[i:i + 5]
  168. ack_parsed = parse_address_assignment_response(chunk)
  169. break
  170. if ack_parsed and 'error' not in ack_parsed:
  171. logger.info(f"UID={uid} 地址={new_address} 已确认 (响应: {ack_parsed.get('raw_data')})")
  172. self.stored_devices[uid] = new_address
  173. results.append({
  174. 'status': 'assigned',
  175. 'uid': uid,
  176. 'address': new_address,
  177. 'request': request.hex(),
  178. 'ack': ack_parsed.get('raw_data')
  179. })
  180. else:
  181. # 即使未收到确认也保存地址:部分从机不返回 43,但实际已生效
  182. logger.warning(f"UID={uid} 地址={new_address} 未收到 43 确认,仍视为已分配")
  183. self.stored_devices[uid] = new_address
  184. results.append({
  185. 'status': 'assigned',
  186. 'uid': uid,
  187. 'address': new_address,
  188. 'request': request.hex(),
  189. 'ack': ack.hex() if ack else None
  190. })
  191. return results
  192. def auto_configure(self, timeout: float = None) -> dict:
  193. """自动配置所有未配置的设备"""
  194. with self.config_lock:
  195. logger.info("开始自动配置设备...")
  196. responses = self.broadcast_query(timeout)
  197. if not responses:
  198. return {
  199. 'success': True,
  200. 'message': '未发现任何从机设备',
  201. 'discovered': 0,
  202. 'configured': 0,
  203. 'results': []
  204. }
  205. results = self.process_responses(responses)
  206. assigned = sum(1 for r in results if r.get('status') == 'assigned')
  207. errors = sum(1 for r in results if r.get('status') == 'error')
  208. return {
  209. 'success': errors == 0,
  210. 'discovered': len(responses),
  211. 'confirmed': 0,
  212. 'assigned': assigned,
  213. 'errors': errors,
  214. 'results': results,
  215. 'stored_devices': self.get_stored_devices()
  216. }
  217. def save_config(self, filepath: str) -> bool:
  218. """保存配置到文件"""
  219. import json
  220. try:
  221. with open(filepath, 'w') as f:
  222. json.dump(self.stored_devices, f, indent=2)
  223. return True
  224. except Exception as e:
  225. logger.error(f"保存配置失败: {e}")
  226. return False
  227. def load_config(self, filepath: str) -> bool:
  228. """从文件加载配置"""
  229. import json
  230. try:
  231. with open(filepath, 'r') as f:
  232. self.stored_devices = json.load(f)
  233. return True
  234. except Exception as e:
  235. logger.error(f"加载配置失败: {e}")
  236. return False
  237. # ========== Modbus RTU Client ==========
  238. # ========== V4 批量读卡协议常量 ==========
  239. # V4: 一次性批量读 24 路卡号;每张卡 6 字节;无卡标记 6×0xFF。
  240. CARD_REG_ADDR = 0x0002 # 批量卡号读起始寄存器
  241. CARD_REG_QUANTITY = 0x48 # 72 寄存器 = 144 字节
  242. CARD_DATA_BYTES = 144 # 24 天线 × 6 字节
  243. BYTES_PER_TAG = 6
  244. ANTENNA_COUNT = 24
  245. NO_TAG_BYTES = b'\xff' * 6 # V4 无卡标记
  246. NO_TAG_HEX = 'ffffffffffff'
  247. BULK_BYTE_COUNT = 0x90 # 144,Modbus 响应字节计数字段
  248. BULK_RESPONSE_LEN = 149 # addr(1)+fc(1)+bytecount(1)+data(144)+crc(2)
  249. # 天线地址映射表(天线编号 -> Modbus寄存器地址)
  250. # V4: 所有天线共享寄存器 0x0002,由 read_all_antenna_cards 一次批量读出
  251. ANTENNA_ADDRESSES = {i: CARD_REG_ADDR for i in range(1, 25)}
  252. import struct
  253. import logging
  254. logger = logging.getLogger('serial_mqtt_gateway')
  255. def calculate_crc16(data: bytes) -> bytes:
  256. """计算Modbus CRC16校验"""
  257. crc = 0xFFFF
  258. for byte in data:
  259. crc ^= byte
  260. for _ in range(8):
  261. if crc & 0x0001:
  262. crc = (crc >> 1) ^ 0xA001
  263. else:
  264. crc >>= 1
  265. return struct.pack('<H', crc)
  266. def verify_crc16(data: bytes) -> bool:
  267. """验证Modbus CRC16校验"""
  268. if len(data) < 2:
  269. return False
  270. received_crc = struct.unpack('<H', data[-2:])[0]
  271. calculated_crc = struct.unpack('<H', calculate_crc16(data[:-2]))[0]
  272. return received_crc == calculated_crc
  273. class ModbusRTUClient:
  274. """Modbus RTU 客户端"""
  275. def __init__(self, serial_port):
  276. self.serial = serial_port
  277. self.default_timeout = 2.0
  278. self.test_cards = None # 测试注入:非 None 时 read_all_antenna_cards 返回构造数据
  279. def _build_request(self, device_address: int, function_code: int, data: bytes = b'') -> bytes:
  280. """构建Modbus请求帧"""
  281. pdu = bytes([device_address, function_code]) + data
  282. return pdu + calculate_crc16(pdu)
  283. def _send_receive(self, request: bytes, expected_min_len: int = 5, timeout: float = None) -> dict:
  284. """send request and receive response using sync method"""
  285. timeout = timeout or self.default_timeout
  286. if not self.serial.ser:
  287. return {"error": "serial not connected"}
  288. response = self.serial.send_and_wait(request, timeout=timeout, min_response_bytes=expected_min_len)
  289. if not response or len(response) < expected_min_len:
  290. return {"error": "response timeout or too short", "raw_data": response.hex() if response else ""}
  291. if not verify_crc16(response):
  292. logger.warning(f"CRC check failed: {response.hex()}")
  293. return {"error": "CRC check failed", "raw_data": response.hex()}
  294. return {
  295. "success": True,
  296. "device_address": response[0],
  297. "function_code": response[1],
  298. "data": response[2:-2].hex(),
  299. "raw_data": response.hex()
  300. }
  301. def read_all_antenna_cards(self, device_address: int, timeout: float = None) -> dict:
  302. """V4: 一次批量读取全部 24 路天线卡号。
  303. 寄存器 0x0002,数量 0x48(72 寄存器 = 144 字节)。
  304. 响应: [addr][0x03][0x90][144 数据字节][CRC_L][CRC_H] = 149 字节。
  305. 每张卡 6 字节;无卡标记 = FF FF FF FF FF FF。
  306. 返回:
  307. {'success': True, 'device_address', 'function_code': 0x03,
  308. 'cards': [{'antenna', 'card_str', 'present', 'uid'}, ...24],
  309. 'raw_data'}
  310. """
  311. # 测试注入: test_cards = {addr: {port_id: 'hex12'|None, ...}, ...}
  312. # 命中时绕过串口,返回构造数据(触发告警/同步流程测试)
  313. if self.test_cards is not None:
  314. panel_map = self.test_cards.get(device_address)
  315. if panel_map is not None:
  316. cards = []
  317. for i in range(ANTENNA_COUNT):
  318. ant = i + 1
  319. hex_str = panel_map.get(ant)
  320. if hex_str:
  321. tag = bytes.fromhex(hex_str)
  322. cards.append({
  323. 'antenna': ant,
  324. 'card_str': tag.hex(),
  325. 'present': True,
  326. 'uid': ':'.join(f'{b:02x}' for b in tag)
  327. })
  328. else:
  329. cards.append({
  330. 'antenna': ant, 'card_str': '', 'present': False, 'uid': ''
  331. })
  332. return {
  333. "success": True,
  334. "device_address": device_address,
  335. "function_code": 0x03,
  336. "cards": cards,
  337. "raw_data": "TEST_INJECTED"
  338. }
  339. timeout = timeout or self.default_timeout
  340. if not self.serial.ser:
  341. return {"error": "serial not connected"}
  342. request_data = struct.pack('>HH', CARD_REG_ADDR, CARD_REG_QUANTITY)
  343. request = self._build_request(device_address, 0x03, request_data)
  344. # 直调 send_and_wait 并要求完整 149 字节;不走 _send_receive
  345. # (其 data 字段会含 0x90 字节计数字节,且默认 min=5 会截断 149 字节帧)
  346. response = self.serial.send_and_wait(
  347. request, timeout=timeout, min_response_bytes=BULK_RESPONSE_LEN
  348. )
  349. if not response or len(response) < 5:
  350. return {"error": "response timeout or too short",
  351. "raw_data": response.hex() if response else ""}
  352. # Modbus 异常响应: addr + 0x83 + exception_code + CRC
  353. if response[1] & 0x80:
  354. return {"error": f"Modbus exception: {response[2]:#x}",
  355. "raw_data": response.hex()}
  356. if len(response) < BULK_RESPONSE_LEN:
  357. return {"error": f"response too short ({len(response)}/{BULK_RESPONSE_LEN})",
  358. "raw_data": response.hex()}
  359. if not verify_crc16(response):
  360. logger.warning(f"CRC check failed: {response.hex()}")
  361. return {"error": "CRC check failed", "raw_data": response.hex()}
  362. if response[1] != 0x03 or response[2] != BULK_BYTE_COUNT:
  363. return {"error": f"unexpected response: fc={response[1]:#x} count={response[2]:#x}",
  364. "raw_data": response.hex()}
  365. tag_data = response[3:3 + CARD_DATA_BYTES] # 144 字节
  366. cards = []
  367. for i in range(ANTENNA_COUNT):
  368. tag = tag_data[i * BYTES_PER_TAG:(i + 1) * BYTES_PER_TAG] # 6 字节
  369. present = tag != NO_TAG_BYTES
  370. cards.append({
  371. 'antenna': i + 1,
  372. 'card_str': tag.hex() if present else '',
  373. 'present': present,
  374. 'uid': ':'.join(f'{b:02x}' for b in tag) if present else ''
  375. })
  376. return {
  377. "success": True,
  378. "device_address": response[0],
  379. "function_code": 0x03,
  380. "cards": cards,
  381. "raw_data": response.hex()
  382. }
  383. def read_antenna_card(self, device_address: int, antenna_num: int, timeout: float = None) -> dict:
  384. """读取单根天线卡号(V4: 由批量读提取,不再发 V2 逐天线帧)。
  385. 调用 read_all_antenna_cards 后提取指定天线,返回兼容 schema。
  386. 无卡时 card_number_str 为空串(V2 是 '0000000000000000')。
  387. """
  388. if antenna_num < 1 or antenna_num > ANTENNA_COUNT:
  389. return {"error": f"invalid antenna number: {antenna_num}, must be 1-{ANTENNA_COUNT}"}
  390. result = self.read_all_antenna_cards(device_address, timeout=timeout)
  391. if 'error' in result:
  392. return result
  393. card = next((c for c in result.get('cards', []) if c['antenna'] == antenna_num), None)
  394. if card is None:
  395. return {"error": f"antenna {antenna_num} not found in bulk response",
  396. "raw_data": result.get('raw_data', '')}
  397. base = {k: v for k, v in result.items() if k != 'cards'}
  398. if card['present']:
  399. tag_bytes = bytes.fromhex(card['card_str']) # 6 字节
  400. card_number = int.from_bytes(tag_bytes, 'big')
  401. return {
  402. 'card_number': card_number,
  403. 'card_number_hex': f'0x{card_number:012x}',
  404. 'card_number_str': card['card_str'],
  405. 'antenna': antenna_num,
  406. 'uid': card['uid'],
  407. **base
  408. }
  409. return {
  410. 'card_number': 0,
  411. 'card_number_hex': '0x000000000000',
  412. 'card_number_str': '',
  413. 'antenna': antenna_num,
  414. 'uid': '',
  415. **base
  416. }
  417. def read_holding_registers(self, device_address: int, start_address: int, quantity: int = 1, timeout: float = None) -> dict:
  418. """读取保持寄存器"""
  419. request_data = struct.pack('>HH', start_address, quantity)
  420. request = self._build_request(device_address, 0x03, request_data)
  421. result = self._send_receive(request, timeout=timeout)
  422. if 'error' in result:
  423. return result
  424. data_str = result.get('data', '')
  425. data_bytes = bytes.fromhex(data_str) if isinstance(data_str, str) and data_str else b''
  426. registers = []
  427. for i in range(0, len(data_bytes), 2):
  428. if i + 1 < len(data_bytes):
  429. registers.append(struct.unpack('>H', data_bytes[i:i+2])[0])
  430. return {'registers': registers, **result}
  431. def write_single_register(self, device_address: int, register_address: int, value: int, timeout: float = None) -> dict:
  432. """写单个寄存器"""
  433. request_data = struct.pack('>HH', register_address, value)
  434. request = self._build_request(device_address, 0x06, request_data)
  435. return self._send_receive(request, timeout=timeout)
  436. def set_rgb_led(self, device_address: int, led_number: int, color: int, timeout: float = None) -> dict:
  437. """设置RGB灯
  438. Register 0x0001, value = (led_number << 8) | color
  439. color: 0=off, 1=flashing red, 2=flashing green, 3=flashing blue
  440. """
  441. value = (led_number << 8) | (color & 0xFF)
  442. return self.write_single_register(device_address, 0x0001, value, timeout)
  443. def scan_devices(self, max_address: int = 247) -> list:
  444. """scan online devices using Modbus 03"""
  445. import time
  446. devices = []
  447. for addr in range(1, min(max_address + 1, 248)):
  448. request = self._build_request(addr, 0x03, struct.pack('>HH', 0x0000, 0x0001))
  449. response = self.serial.send_and_wait(request, timeout=0.3, min_response_bytes=5)
  450. if response and len(response) >= 5 and verify_crc16(response):
  451. devices.append(addr)
  452. return devices