modbus_rtu.py 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575
  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. # 注: 硬件通信固定 6 字节卡号;对 MQTT 的 8 字节(16 位)适配在 app.py 边界完成
  255. # (前 4 位 DTU 号 + 后 12 位卡号)。
  256. CARD_REG_ADDR = 0x0002 # 批量卡号读起始寄存器
  257. CARD_REG_QUANTITY = 0x48 # 72 寄存器 = 144 字节
  258. CARD_DATA_BYTES = 144 # 24 天线 × 6 字节
  259. BYTES_PER_TAG = 6
  260. ANTENNA_COUNT = 24
  261. NO_TAG_BYTES = b'\xff' * 6 # V4 无卡标记
  262. NO_TAG_HEX = 'ffffffffffff'
  263. BULK_BYTE_COUNT = 0x90 # 144,Modbus 响应字节计数字段
  264. BULK_RESPONSE_LEN = 149 # addr(1)+fc(1)+bytecount(1)+data(144)+crc(2)
  265. # 天线地址映射表(天线编号 -> Modbus寄存器地址)
  266. # V4: 所有天线共享寄存器 0x0002,由 read_all_antenna_cards 一次批量读出
  267. ANTENNA_ADDRESSES = {i: CARD_REG_ADDR for i in range(1, 25)}
  268. import struct
  269. import logging
  270. logger = logging.getLogger('serial_mqtt_gateway')
  271. def calculate_crc16(data: bytes) -> bytes:
  272. """计算Modbus CRC16校验"""
  273. crc = 0xFFFF
  274. for byte in data:
  275. crc ^= byte
  276. for _ in range(8):
  277. if crc & 0x0001:
  278. crc = (crc >> 1) ^ 0xA001
  279. else:
  280. crc >>= 1
  281. return struct.pack('<H', crc)
  282. def verify_crc16(data: bytes) -> bool:
  283. """验证Modbus CRC16校验"""
  284. if len(data) < 2:
  285. return False
  286. received_crc = struct.unpack('<H', data[-2:])[0]
  287. calculated_crc = struct.unpack('<H', calculate_crc16(data[:-2]))[0]
  288. return received_crc == calculated_crc
  289. class ModbusRTUClient:
  290. """Modbus RTU 客户端"""
  291. def __init__(self, serial_port):
  292. self.serial = serial_port
  293. self.default_timeout = 2.0
  294. self.test_cards = None # 测试注入:非 None 时 read_all_antenna_cards 返回构造数据
  295. def _build_request(self, device_address: int, function_code: int, data: bytes = b'') -> bytes:
  296. """构建Modbus请求帧"""
  297. pdu = bytes([device_address, function_code]) + data
  298. return pdu + calculate_crc16(pdu)
  299. def _send_receive(self, request: bytes, expected_min_len: int = 5, timeout: float = None) -> dict:
  300. """send request and receive response using sync method"""
  301. timeout = timeout or self.default_timeout
  302. if not self.serial.ser:
  303. return {"error": "serial not connected"}
  304. response = self.serial.send_and_wait(request, timeout=timeout, min_response_bytes=expected_min_len)
  305. if not response or len(response) < expected_min_len:
  306. return {"error": "response timeout or too short", "raw_data": response.hex() if response else ""}
  307. if not verify_crc16(response):
  308. logger.warning(f"CRC check failed: {response.hex()}")
  309. return {"error": "CRC check failed", "raw_data": response.hex()}
  310. return {
  311. "success": True,
  312. "device_address": response[0],
  313. "function_code": response[1],
  314. "data": response[2:-2].hex(),
  315. "raw_data": response.hex()
  316. }
  317. def read_all_antenna_cards(self, device_address: int, timeout: float = None) -> dict:
  318. """V4: 一次批量读取全部 24 路天线卡号。
  319. 寄存器 0x0002,数量 0x48(72 寄存器 = 144 字节)。
  320. 响应: [addr][0x03][0x90][144 数据字节][CRC_L][CRC_H] = 149 字节。
  321. 每张卡 6 字节(硬件卡号,12 位十六进制);无卡标记 = 6×0xFF。
  322. MQTT 侧的 8 字节(16 位) jumper_uid 适配在 app.py 边界完成。
  323. 返回:
  324. {'success': True, 'device_address', 'function_code': 0x03,
  325. 'cards': [{'antenna', 'card_str', 'present', 'uid'}, ...24],
  326. 'raw_data'}
  327. """
  328. # 测试注入: test_cards = {addr: {port_id: 'hex12'|None, ...}, ...}
  329. # 命中时绕过串口,返回构造数据(触发告警/同步流程测试)
  330. if self.test_cards is not None:
  331. panel_map = self.test_cards.get(device_address)
  332. if panel_map is not None:
  333. cards = []
  334. for i in range(ANTENNA_COUNT):
  335. ant = i + 1
  336. hex_str = panel_map.get(ant)
  337. if hex_str:
  338. tag = bytes.fromhex(hex_str)
  339. cards.append({
  340. 'antenna': ant,
  341. 'card_str': tag.hex(),
  342. 'present': True,
  343. 'uid': ':'.join(f'{b:02x}' for b in tag)
  344. })
  345. else:
  346. cards.append({
  347. 'antenna': ant, 'card_str': '', 'present': False, 'uid': ''
  348. })
  349. return {
  350. "success": True,
  351. "device_address": device_address,
  352. "function_code": 0x03,
  353. "cards": cards,
  354. "raw_data": "TEST_INJECTED"
  355. }
  356. timeout = timeout or self.default_timeout
  357. if not self.serial.ser:
  358. return {"error": "serial not connected"}
  359. request_data = struct.pack('>HH', CARD_REG_ADDR, CARD_REG_QUANTITY)
  360. request = self._build_request(device_address, 0x03, request_data)
  361. # 直调 send_and_wait 并要求完整 149 字节;不走 _send_receive
  362. # (其 data 字段会含 0x90 字节计数字节,且默认 min=5 会截断 149 字节帧)
  363. response = self.serial.send_and_wait(
  364. request, timeout=timeout, min_response_bytes=BULK_RESPONSE_LEN
  365. )
  366. if not response or len(response) < 5:
  367. return {"error": "response timeout or too short",
  368. "raw_data": response.hex() if response else ""}
  369. # Modbus 异常响应: addr + 0x83 + exception_code + CRC
  370. if response[1] & 0x80:
  371. return {"error": f"Modbus exception: {response[2]:#x}",
  372. "raw_data": response.hex()}
  373. if len(response) < BULK_RESPONSE_LEN:
  374. return {"error": f"response too short ({len(response)}/{BULK_RESPONSE_LEN})",
  375. "raw_data": response.hex()}
  376. if not verify_crc16(response):
  377. logger.warning(f"CRC check failed: {response.hex()}")
  378. return {"error": "CRC check failed", "raw_data": response.hex()}
  379. if response[1] != 0x03 or response[2] != BULK_BYTE_COUNT:
  380. return {"error": f"unexpected response: fc={response[1]:#x} count={response[2]:#x}",
  381. "raw_data": response.hex()}
  382. tag_data = response[3:3 + CARD_DATA_BYTES] # 144 字节
  383. cards = []
  384. for i in range(ANTENNA_COUNT):
  385. tag = tag_data[i * BYTES_PER_TAG:(i + 1) * BYTES_PER_TAG] # 6 字节
  386. present = tag != NO_TAG_BYTES
  387. cards.append({
  388. 'antenna': i + 1,
  389. 'card_str': tag.hex() if present else '',
  390. 'present': present,
  391. 'uid': ':'.join(f'{b:02x}' for b in tag) if present else ''
  392. })
  393. return {
  394. "success": True,
  395. "device_address": response[0],
  396. "function_code": 0x03,
  397. "cards": cards,
  398. "raw_data": response.hex()
  399. }
  400. def read_antenna_card(self, device_address: int, antenna_num: int, timeout: float = None) -> dict:
  401. """读取单根天线卡号(V4: 由批量读提取,不再发 V2 逐天线帧)。
  402. 调用 read_all_antenna_cards 后提取指定天线,返回兼容 schema。
  403. 无卡时 card_number_str 为空串(V2 是 '0000000000000000')。
  404. """
  405. if antenna_num < 1 or antenna_num > ANTENNA_COUNT:
  406. return {"error": f"invalid antenna number: {antenna_num}, must be 1-{ANTENNA_COUNT}"}
  407. result = self.read_all_antenna_cards(device_address, timeout=timeout)
  408. if 'error' in result:
  409. return result
  410. card = next((c for c in result.get('cards', []) if c['antenna'] == antenna_num), None)
  411. if card is None:
  412. return {"error": f"antenna {antenna_num} not found in bulk response",
  413. "raw_data": result.get('raw_data', '')}
  414. base = {k: v for k, v in result.items() if k != 'cards'}
  415. if card['present']:
  416. tag_bytes = bytes.fromhex(card['card_str']) # 6 字节
  417. card_number = int.from_bytes(tag_bytes, 'big')
  418. return {
  419. 'card_number': card_number,
  420. 'card_number_hex': f'0x{card_number:012x}',
  421. 'card_number_str': card['card_str'],
  422. 'antenna': antenna_num,
  423. 'uid': card['uid'],
  424. **base
  425. }
  426. return {
  427. 'card_number': 0,
  428. 'card_number_hex': '0x000000000000',
  429. 'card_number_str': '',
  430. 'antenna': antenna_num,
  431. 'uid': '',
  432. **base
  433. }
  434. def read_holding_registers(self, device_address: int, start_address: int, quantity: int = 1, timeout: float = None) -> dict:
  435. """读取保持寄存器"""
  436. request_data = struct.pack('>HH', start_address, quantity)
  437. request = self._build_request(device_address, 0x03, request_data)
  438. result = self._send_receive(request, timeout=timeout)
  439. if 'error' in result:
  440. return result
  441. data_str = result.get('data', '')
  442. data_bytes = bytes.fromhex(data_str) if isinstance(data_str, str) and data_str else b''
  443. registers = []
  444. for i in range(0, len(data_bytes), 2):
  445. if i + 1 < len(data_bytes):
  446. registers.append(struct.unpack('>H', data_bytes[i:i+2])[0])
  447. return {'registers': registers, **result}
  448. def write_single_register(self, device_address: int, register_address: int, value: int, timeout: float = None) -> dict:
  449. """写单个寄存器"""
  450. request_data = struct.pack('>HH', register_address, value)
  451. request = self._build_request(device_address, 0x06, request_data)
  452. return self._send_receive(request, timeout=timeout)
  453. def set_rgb_led(self, device_address: int, led_number: int, color: int, timeout: float = None) -> dict:
  454. """设置RGB灯
  455. Register 0x0001, value = (led_number << 8) | color
  456. color: 0=off, 1=flashing red, 2=flashing green, 3=flashing blue
  457. """
  458. value = (led_number << 8) | (color & 0xFF)
  459. return self.write_single_register(device_address, 0x0001, value, timeout)
  460. def scan_devices(self, max_address: int = 247) -> list:
  461. """scan online devices using Modbus 03"""
  462. import time
  463. devices = []
  464. for addr in range(1, min(max_address + 1, 248)):
  465. request = self._build_request(addr, 0x03, struct.pack('>HH', 0x0000, 0x0001))
  466. response = self.serial.send_and_wait(request, timeout=0.3, min_response_bytes=5)
  467. if response and len(response) >= 5 and verify_crc16(response):
  468. devices.append(addr)
  469. return devices