from typing import Optional import threading # ========== 地址配置协议 ========== def build_broadcast_query() -> bytes: """构建广播查询指令 DTU定时发送广播指令查询未配置地址的从机 格式: 00 40 00 40 00 (无CRC, 固定5字节) """ return bytes([0x00, 0x40, 0x00, 0x40, 0x00]) def parse_device_response(data: bytes) -> Optional[dict]: """解析从机应答 从机应答格式: 00 + 41 + 12个UID + CRC_L + CRC_H 例如: 00 41 18 00 40 00 14 00 00 59 59 54 30 56 AE 34 """ if len(data) < 16: return {'error': '数据长度不足', 'raw_data': data.hex()} if data[1] != 0x41: return {'error': f"非预期功能码: {data[1]:#x}", 'raw_data': data.hex()} uid_bytes = data[2:14] uid_hex = uid_bytes.hex() # CRC 校验失败则拒绝(传入完整 16 字节帧,最后两字节为 CRC) if not verify_crc16(data): return {'error': 'CRC校验失败', 'raw_data': data.hex()} # UID 有效性检查:拒绝明显是噪声或碰撞产生的数据 zero_count = uid_bytes.count(0) if zero_count > 6: return {'error': f"UID中零字节过多({zero_count}/12)", 'raw_data': data.hex()} # 检查碰撞:UID 中不应有大量重复字节 from collections import Counter counts = Counter(uid_bytes) if counts.most_common(1)[0][1] > 8: return {'error': f"UID中单字节重复过多({counts.most_common(1)})", 'raw_data': data.hex()} return { 'function_code': 0x41, 'uid': uid_hex, 'uid_readable': ':'.join(f'{b:02x}' for b in uid_bytes), 'source_address': data[0], 'raw_data': data.hex() } def build_confirm_address(device_address: int, uid: bytes) -> bytes: """构建确认地址指令 DTU遍历已存储的UID和地址,如果存在匹配的UID则发送确认指令 格式: add + 44 + 00 + 12个UID + CRC_L + CRC_H """ if len(uid) != 12: raise ValueError(f"UID长度必须是12字节,当前: {len(uid)}") data = bytes([device_address, 0x44, 0x00]) + uid return data + calculate_crc16(data) def build_assign_address(uid: bytes, address: int) -> bytes: """构建分配地址指令 如果DTU中没有该UID的记录,则分配新地址 格式: 00 + 42 + 12个UID + add + CRC_L + CRC_H """ if len(uid) != 12: raise ValueError(f"UID长度必须是12字节,当前: {len(uid)}") if address < 1 or address > 247: raise ValueError(f"设备地址必须在1-247之间,当前: {address}") data = bytes([0x00, 0x42]) + uid + bytes([address]) return data + calculate_crc16(data) def parse_address_assignment_response(data: bytes) -> dict: """解析地址分配响应 模块应答格式: add + 43 + 00 + CRC_L + CRC_H 例如: 01 43 00 11 30 """ if len(data) < 5: return {'error': '数据长度不足', 'raw_data': data.hex()} device_address = data[0] if data[1] != 0x43: return {'error': f"非预期功能码: {data[1]:#x}", 'raw_data': data.hex()} if not verify_crc16(data): logger.warning(f"CRC校验失败: {data.hex()}") return { 'device_address': device_address, 'function_code': 0x43, 'raw_data': data.hex() } class AddressConfigProtocol: """地址配置协议处理器""" def __init__(self, serial_port): self.serial = serial_port self.default_timeout = 5.0 # 设备广播响应可能较慢,给足时间 self.stored_devices = {} self.config_lock = threading.RLock() # 地址配置期间占用,避免与轮询冲突 def add_stored_device(self, uid: str, address: int): """添加已存储的设备""" self.stored_devices[uid] = address def get_stored_devices(self) -> dict: """获取已存储的设备列表""" return self.stored_devices.copy() def broadcast_query(self, timeout: float = None) -> list: """send broadcast query and collect all valid 00 41 responses""" import time with self.config_lock: request = build_broadcast_query() logger.info(f"send broadcast: {request.hex()}") timeout = timeout or self.default_timeout # 先清空输入缓冲区,避免残留数据干扰 self.serial.flush_input() time.sleep(0.05) # 发送广播并等待响应;所有串口操作统一走发送队列, # 这里不再额外直接读串口,避免在收集窗口内把后续队列命令的响应抢走。 response = self.serial.send_and_wait(request, timeout=timeout, min_response_bytes=16) responses = [] if response and len(response) >= 16: logger.info(f"got raw response ({len(response)}B): {response.hex()}") # RS485 半双工会回显已发送的数据,导致 response 开头可能有 5 字节回显 # 在 response 中搜索所有 00 41 标记,尝试从中提取 16 字节帧解析 for i in range(len(response) - 15 + 1): if response[i] == 0x00 and response[i + 1] == 0x41: chunk = response[i:i + 16] if len(chunk) >= 16: parsed = parse_device_response(chunk) if "error" not in parsed: if not any(p.get('uid') == parsed['uid'] for p in responses): responses.append(parsed) else: logger.info(f"no response or too short: {len(response) if response else 0}B") logger.info(f"broadcast done, got {len(responses)} valid unique responses") return responses def _next_available_address(self) -> int: """从 1~247 中找一个未被 stored_devices 占用的地址""" used = set(self.stored_devices.values()) for addr in range(1, 248): if addr not in used: return addr raise RuntimeError("Modbus 地址 1~247 已全部占用") def process_responses(self, responses: list) -> list: """处理广播查询响应:给每个应答设备分配地址(优先使用已记录地址)""" import time with self.config_lock: results = [] for resp in responses: if 'error' in resp: results.append({'status': 'error', 'message': resp['error']}) continue uid = resp['uid'] src_addr = resp.get('source_address', 0) device_address = self.stored_devices.get(uid) # 优先使用地址表中已保存的地址;没有则分配下一个可用地址 if device_address: new_address = device_address logger.info(f"UID={uid} 源地址={src_addr},地址表存在记录,分配已有地址={new_address}") else: new_address = self._next_available_address() logger.info(f"UID={uid} 源地址={src_addr},地址表无记录,分配新地址={new_address}") request = build_assign_address(bytes.fromhex(uid), new_address) logger.info(f"发送分配地址帧: {request.hex()}") self.serial.flush_input() time.sleep(0.05) success, msg = self.serial.send_raw(request) if not success: logger.error(f"UID={uid} 分配地址帧发送失败: {msg}") results.append({ 'status': 'error', 'uid': uid, 'message': f"send failed: {msg}" }) continue # 等待并解析从机的 43 确认响应 time.sleep(0.1) ack = self.serial.send_and_wait(b'', timeout=0.5, min_response_bytes=5) ack_parsed = None if ack and len(ack) >= 5: # 过滤回显/噪声,找到 43 响应帧 for i in range(len(ack) - 4): if ack[i + 1] == 0x43: chunk = ack[i:i + 5] ack_parsed = parse_address_assignment_response(chunk) break if ack_parsed and 'error' not in ack_parsed: logger.info(f"UID={uid} 地址={new_address} 已确认 (响应: {ack_parsed.get('raw_data')})") self.stored_devices[uid] = new_address results.append({ 'status': 'assigned', 'uid': uid, 'address': new_address, 'request': request.hex(), 'ack': ack_parsed.get('raw_data') }) else: # 即使未收到确认也保存地址:部分从机不返回 43,但实际已生效 logger.warning(f"UID={uid} 地址={new_address} 未收到 43 确认,仍视为已分配") self.stored_devices[uid] = new_address results.append({ 'status': 'assigned', 'uid': uid, 'address': new_address, 'request': request.hex(), 'ack': ack.hex() if ack else None }) return results def auto_configure(self, timeout: float = None) -> dict: """自动配置所有未配置的设备""" with self.config_lock: logger.info("开始自动配置设备...") responses = self.broadcast_query(timeout) if not responses: return { 'success': True, 'message': '未发现任何从机设备', 'discovered': 0, 'configured': 0, 'results': [] } results = self.process_responses(responses) assigned = sum(1 for r in results if r.get('status') == 'assigned') errors = sum(1 for r in results if r.get('status') == 'error') return { 'success': errors == 0, 'discovered': len(responses), 'confirmed': 0, 'assigned': assigned, 'errors': errors, 'results': results, 'stored_devices': self.get_stored_devices() } def save_config(self, filepath: str) -> bool: """保存配置到文件""" import json try: with open(filepath, 'w') as f: json.dump(self.stored_devices, f, indent=2) return True except Exception as e: logger.error(f"保存配置失败: {e}") return False def load_config(self, filepath: str) -> bool: """从文件加载配置""" import json try: with open(filepath, 'r') as f: self.stored_devices = json.load(f) return True except Exception as e: logger.error(f"加载配置失败: {e}") return False # ========== Modbus RTU Client ========== # ========== V4 批量读卡协议常量 ========== # V4: 一次性批量读 24 路卡号;每张卡 6 字节;无卡标记 6×0xFF。 CARD_REG_ADDR = 0x0002 # 批量卡号读起始寄存器 CARD_REG_QUANTITY = 0x48 # 72 寄存器 = 144 字节 CARD_DATA_BYTES = 144 # 24 天线 × 6 字节 BYTES_PER_TAG = 6 ANTENNA_COUNT = 24 NO_TAG_BYTES = b'\xff' * 6 # V4 无卡标记 NO_TAG_HEX = 'ffffffffffff' BULK_BYTE_COUNT = 0x90 # 144,Modbus 响应字节计数字段 BULK_RESPONSE_LEN = 149 # addr(1)+fc(1)+bytecount(1)+data(144)+crc(2) # 天线地址映射表(天线编号 -> Modbus寄存器地址) # V4: 所有天线共享寄存器 0x0002,由 read_all_antenna_cards 一次批量读出 ANTENNA_ADDRESSES = {i: CARD_REG_ADDR for i in range(1, 25)} import struct import logging logger = logging.getLogger('serial_mqtt_gateway') def calculate_crc16(data: bytes) -> bytes: """计算Modbus CRC16校验""" crc = 0xFFFF for byte in data: crc ^= byte for _ in range(8): if crc & 0x0001: crc = (crc >> 1) ^ 0xA001 else: crc >>= 1 return struct.pack(' bool: """验证Modbus CRC16校验""" if len(data) < 2: return False received_crc = struct.unpack(' bytes: """构建Modbus请求帧""" pdu = bytes([device_address, function_code]) + data return pdu + calculate_crc16(pdu) def _send_receive(self, request: bytes, expected_min_len: int = 5, timeout: float = None) -> dict: """send request and receive response using sync method""" timeout = timeout or self.default_timeout if not self.serial.ser: return {"error": "serial not connected"} response = self.serial.send_and_wait(request, timeout=timeout, min_response_bytes=expected_min_len) if not response or len(response) < expected_min_len: return {"error": "response timeout or too short", "raw_data": response.hex() if response else ""} if not verify_crc16(response): logger.warning(f"CRC check failed: {response.hex()}") return {"error": "CRC check failed", "raw_data": response.hex()} return { "success": True, "device_address": response[0], "function_code": response[1], "data": response[2:-2].hex(), "raw_data": response.hex() } def read_all_antenna_cards(self, device_address: int, timeout: float = None) -> dict: """V4: 一次批量读取全部 24 路天线卡号。 寄存器 0x0002,数量 0x48(72 寄存器 = 144 字节)。 响应: [addr][0x03][0x90][144 数据字节][CRC_L][CRC_H] = 149 字节。 每张卡 6 字节;无卡标记 = FF FF FF FF FF FF。 返回: {'success': True, 'device_address', 'function_code': 0x03, 'cards': [{'antenna', 'card_str', 'present', 'uid'}, ...24], 'raw_data'} """ # 测试注入: test_cards = {addr: {port_id: 'hex12'|None, ...}, ...} # 命中时绕过串口,返回构造数据(触发告警/同步流程测试) if self.test_cards is not None: panel_map = self.test_cards.get(device_address) if panel_map is not None: cards = [] for i in range(ANTENNA_COUNT): ant = i + 1 hex_str = panel_map.get(ant) if hex_str: tag = bytes.fromhex(hex_str) cards.append({ 'antenna': ant, 'card_str': tag.hex(), 'present': True, 'uid': ':'.join(f'{b:02x}' for b in tag) }) else: cards.append({ 'antenna': ant, 'card_str': '', 'present': False, 'uid': '' }) return { "success": True, "device_address": device_address, "function_code": 0x03, "cards": cards, "raw_data": "TEST_INJECTED" } timeout = timeout or self.default_timeout if not self.serial.ser: return {"error": "serial not connected"} request_data = struct.pack('>HH', CARD_REG_ADDR, CARD_REG_QUANTITY) request = self._build_request(device_address, 0x03, request_data) # 直调 send_and_wait 并要求完整 149 字节;不走 _send_receive # (其 data 字段会含 0x90 字节计数字节,且默认 min=5 会截断 149 字节帧) response = self.serial.send_and_wait( request, timeout=timeout, min_response_bytes=BULK_RESPONSE_LEN ) if not response or len(response) < 5: return {"error": "response timeout or too short", "raw_data": response.hex() if response else ""} # Modbus 异常响应: addr + 0x83 + exception_code + CRC if response[1] & 0x80: return {"error": f"Modbus exception: {response[2]:#x}", "raw_data": response.hex()} if len(response) < BULK_RESPONSE_LEN: return {"error": f"response too short ({len(response)}/{BULK_RESPONSE_LEN})", "raw_data": response.hex()} if not verify_crc16(response): logger.warning(f"CRC check failed: {response.hex()}") return {"error": "CRC check failed", "raw_data": response.hex()} if response[1] != 0x03 or response[2] != BULK_BYTE_COUNT: return {"error": f"unexpected response: fc={response[1]:#x} count={response[2]:#x}", "raw_data": response.hex()} tag_data = response[3:3 + CARD_DATA_BYTES] # 144 字节 cards = [] for i in range(ANTENNA_COUNT): tag = tag_data[i * BYTES_PER_TAG:(i + 1) * BYTES_PER_TAG] # 6 字节 present = tag != NO_TAG_BYTES cards.append({ 'antenna': i + 1, 'card_str': tag.hex() if present else '', 'present': present, 'uid': ':'.join(f'{b:02x}' for b in tag) if present else '' }) return { "success": True, "device_address": response[0], "function_code": 0x03, "cards": cards, "raw_data": response.hex() } def read_antenna_card(self, device_address: int, antenna_num: int, timeout: float = None) -> dict: """读取单根天线卡号(V4: 由批量读提取,不再发 V2 逐天线帧)。 调用 read_all_antenna_cards 后提取指定天线,返回兼容 schema。 无卡时 card_number_str 为空串(V2 是 '0000000000000000')。 """ if antenna_num < 1 or antenna_num > ANTENNA_COUNT: return {"error": f"invalid antenna number: {antenna_num}, must be 1-{ANTENNA_COUNT}"} result = self.read_all_antenna_cards(device_address, timeout=timeout) if 'error' in result: return result card = next((c for c in result.get('cards', []) if c['antenna'] == antenna_num), None) if card is None: return {"error": f"antenna {antenna_num} not found in bulk response", "raw_data": result.get('raw_data', '')} base = {k: v for k, v in result.items() if k != 'cards'} if card['present']: tag_bytes = bytes.fromhex(card['card_str']) # 6 字节 card_number = int.from_bytes(tag_bytes, 'big') return { 'card_number': card_number, 'card_number_hex': f'0x{card_number:012x}', 'card_number_str': card['card_str'], 'antenna': antenna_num, 'uid': card['uid'], **base } return { 'card_number': 0, 'card_number_hex': '0x000000000000', 'card_number_str': '', 'antenna': antenna_num, 'uid': '', **base } def read_holding_registers(self, device_address: int, start_address: int, quantity: int = 1, timeout: float = None) -> dict: """读取保持寄存器""" request_data = struct.pack('>HH', start_address, quantity) request = self._build_request(device_address, 0x03, request_data) result = self._send_receive(request, timeout=timeout) if 'error' in result: return result data_str = result.get('data', '') data_bytes = bytes.fromhex(data_str) if isinstance(data_str, str) and data_str else b'' registers = [] for i in range(0, len(data_bytes), 2): if i + 1 < len(data_bytes): registers.append(struct.unpack('>H', data_bytes[i:i+2])[0]) return {'registers': registers, **result} def write_single_register(self, device_address: int, register_address: int, value: int, timeout: float = None) -> dict: """写单个寄存器""" request_data = struct.pack('>HH', register_address, value) request = self._build_request(device_address, 0x06, request_data) return self._send_receive(request, timeout=timeout) def set_rgb_led(self, device_address: int, led_number: int, color: int, timeout: float = None) -> dict: """设置RGB灯 Register 0x0001, value = (led_number << 8) | color color: 0=off, 1=flashing red, 2=flashing green, 3=flashing blue """ value = (led_number << 8) | (color & 0xFF) return self.write_single_register(device_address, 0x0001, value, timeout) def scan_devices(self, max_address: int = 247) -> list: """scan online devices using Modbus 03""" import time devices = [] for addr in range(1, min(max_address + 1, 248)): request = self._build_request(addr, 0x03, struct.pack('>HH', 0x0000, 0x0001)) response = self.serial.send_and_wait(request, timeout=0.3, min_response_bytes=5) if response and len(response) >= 5 and verify_crc16(response): devices.append(addr) return devices