modbus_rtu.py 23 KB

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