#!/usr/bin/env python3 """ 数据中心 → 信号管理页面 (out) 添加信号功能测试 测试 RK3568 远端 RTU 的 WebSocket 前端功能 关键发现: - dc_data 通过二进制 WebSocket 帧发送(ws_send_binary) - curd:add 后不会立即推送,需要等 ws_task 定时轮询(约 1 秒周期) - curd:del 后会立即推送 """ import socket, base64, hashlib, struct, json, os, time, sys IP = '198.120.0.100' PORT = 8000 # =========================== WebSocket 基础工具 =========================== def ws_handshake(sock): """执行 WebSocket 握手""" key = base64.b64encode(os.urandom(16)).decode() req = ( f"GET /ws HTTP/1.1\r\n" f"Host: {IP}:{PORT}\r\n" f"Upgrade: websocket\r\n" f"Connection: Upgrade\r\n" f"Sec-WebSocket-Key: {key}\r\n" f"Sec-WebSocket-Version: 13\r\n" f"\r\n" ) sock.sendall(req.encode()) resp = b"" while b"\r\n\r\n" not in resp: chunk = sock.recv(4096) if not chunk: raise Exception("Handshake failed: no response") resp += chunk if b"101" not in resp.split(b"\r\n")[0]: raise Exception(f"Handshake failed: {resp.split(b'\r\n')[0].decode()}") print("[✓] WebSocket 握手成功") def ws_send(sock, msg): """发送带掩码的 WebSocket 文本帧""" if isinstance(msg, dict): msg = json.dumps(msg) b = msg.encode() L = len(b) mask = os.urandom(4) masked = bytes(b[i] ^ mask[i % 4] for i in range(L)) frame = bytearray([0x81]) # FIN + text opcode if L < 126: frame.append(0x80 | L) elif L < 65536: frame.append(0x80 | 126) frame.extend(struct.pack('>H', L)) else: frame.append(0x80 | 127) frame.extend(struct.pack('>Q', L)) frame.extend(mask) frame.extend(masked) sock.sendall(bytes(frame)) def ws_recv(sock, timeout=5.0): """接收一个 WebSocket 帧,返回 (opcode, payload_bytes)""" sock.settimeout(timeout) try: header = sock.recv(2) if len(header) < 2: return None, None opcode = header[0] & 0x0F masked = (header[1] & 0x80) != 0 length = header[1] & 0x7F if length == 126: length = struct.unpack('>H', sock.recv(2))[0] elif length == 127: length = struct.unpack('>Q', sock.recv(8))[0] if masked: mask_key = sock.recv(4) payload = b"" while len(payload) < length: chunk = sock.recv(min(length - len(payload), 65536)) if not chunk: break payload += chunk if masked: payload = bytes(payload[i] ^ mask_key[i % 4] for i in range(len(payload))) return opcode, payload except socket.timeout: return None, None except Exception as e: print(f" [!] recv error: {e}") return None, None def try_parse_json(payload): """尝试将 payload 解析为 JSON。dc_data 是二进制帧中嵌入的 JSON""" if payload is None: return None try: text = payload.decode('utf-8') return json.loads(text) except: return None def drain_sock(sock, timeout=1.0): """排空 socket 缓冲区中的所有待处理帧,返回所有解析后的消息""" msgs = [] while True: opcode, payload = ws_recv(sock, timeout=timeout) if opcode is None: break # 文本帧 if opcode == 0x01 and payload: obj = try_parse_json(payload) if obj: msgs.append(('text', obj)) else: msgs.append(('text', payload.decode('utf-8', errors='replace'))) # 二进制帧(dc_data 在这里) elif opcode == 0x02 and payload: obj = try_parse_json(payload) if obj: msgs.append(('binary', obj)) else: msgs.append(('binary', payload)) elif opcode == 0x08: msgs.append(('close', None)) break elif opcode == 0x09: # ping pass return msgs # =========================== 测试主流程 =========================== def main(): print("=" * 70) print(" RTU 数据中心 → 信号管理页(out) 添加信号功能测试") print(f" 目标: {IP}:{PORT}") print("=" * 70) # ---- 连接 ---- sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) sock.settimeout(10) try: sock.connect((IP, PORT)) print(f"[✓] TCP 连接成功 → {IP}:{PORT}") except Exception as e: print(f"[✗] TCP 连接失败: {e}") sys.exit(1) ws_handshake(sock) # 原始 socket 级别检查是否有数据 print("[i] 检查原始 socket 是否有数据到达...") sock.settimeout(3.0) try: raw_data = sock.recv(65536) if raw_data: print(f"[i] 收到原始数据: {len(raw_data)} 字节") print(f" 前 50 字节(hex): {raw_data[:50].hex()}") print(f" 前 200 字节(str): {raw_data[:200]}") # 尝试解析 WebSocket 帧 if len(raw_data) >= 2: b0, b1 = raw_data[0], raw_data[1] opcode = b0 & 0x0F fin = (b0 >> 7) & 1 print(f" FIN={fin} opcode={opcode} ({'text' if opcode==1 else 'binary' if opcode==2 else 'close' if opcode==8 else 'other'})") else: print("[i] 原始 socket 无数据") except socket.timeout: print("[i] 原始 socket 超时,无数据") sock.settimeout(10) # 等待初始推送 (cmd_list 等) time.sleep(1.0) msgs = drain_sock(sock, timeout=2.0) print(f"[i] 初始消息数: {len(msgs)}") for kind, msg in msgs: if isinstance(msg, dict): print(f" [{kind}] type={msg.get('type','?')}, keys={list(msg.keys())}") if msg.get('type') == 'cmd_list': print(f" [i] 收到 cmd_list,命令数: {len(msg.get('cmds', []))}") elif isinstance(msg, str): print(f" [{kind}] str: {msg[:200]}") elif isinstance(msg, bytes): print(f" [{kind}] bytes: {len(msg)}B, 前200字节: {msg[:200]}") # ---- 步骤 1: 通过终端命令注册 dc_data 信号 ---- print("\n" + "=" * 70) print(" 步骤 1: 通过 datacenter 命令注册所有类型的信号") print("=" * 70) types_cn = {'out': '注册', 'in': '链接注册', 'yk': '遥控', 'ao': '参数', 'param': '定值'} all_dc_data = {} for st in ['out', 'in', 'yk', 'ao', 'param']: print(f"\n [→] 执行: datacenter {st}") ws_send(sock, {"type": "cmd", "cmd": f"datacenter {st}"}) time.sleep(2.0) # 等待更长时间 # dc_data 通过二进制帧返回 msgs = drain_sock(sock, timeout=3.0) print(f" [i] 收到 {len(msgs)} 条消息") for kind, msg in msgs: if isinstance(msg, dict): print(f" [{kind}] type={msg.get('type','?')}") if msg.get('type') == 'dc_data': all_dc_data[st] = msg for sig_type in ['out', 'in', 'yk', 'ao', 'param']: count = len(msg.get(sig_type, [])) if count > 0: print(f" {sig_type}: {count} 个信号") for k in msg: if k != 'type' and isinstance(msg[k], list): print(f" {k}: {len(msg[k])} 个信号") for item in msg[k][:3]: print(f" - saddr={item.get('saddr','?')} desc={item.get('desc','?')}") elif isinstance(msg, str): print(f" [{kind}] str({len(msg)}): {msg[:300]}") elif isinstance(msg, bytes): print(f" [{kind}] bytes({len(msg)}): {msg[:200]}") # ---- 步骤 2: 检查 dc_data 并选择 out 信号 ---- print("\n" + "=" * 70) print(" 步骤 2: 从 dc_data 选择 out 信号添加到配置页") print("=" * 70) dc_out = all_dc_data.get('out', {}) out_signals = dc_out.get('out', []) if not out_signals: print(" [!] 没有可用的 out 信号!") print(" 可能原因: 服务端 datacenter out 未注册任何信号") sock.close() return # 选择前 3 个 out 信号(只选择直控或选控类型用于测试) test_signals = [s for s in out_signals if s.get('ctrl_type', 0) in [1, 2]] if not test_signals: test_signals = out_signals[:3] # 降级:选前 3 个 test_signals = test_signals[:3] print(f" 从 {len(out_signals)} 个 out 信号中选择 {len(test_signals)} 个测试:") for i, sig in enumerate(test_signals): saddr = sig.get('saddr', '') print(f" [{i+1}] saddr={saddr} desc={sig.get('desc','?')} ctrl_type={sig.get('ctrl_type','?')}") # 先排空已有消息 drain_sock(sock, timeout=0.5) # ---- 步骤 3: 执行 curd:add 添加信号 ---- print("\n" + "=" * 70) print(" 步骤 3: 发送 curd:add 添加信号到配置页(out)") print("=" * 70) added_saddrs = set() for i, sig in enumerate(test_signals): saddr = sig.get('saddr', '') add_msg = { "saddr": saddr, "signal_type": "out", "curd": "add", "setting_zone": "0", "signal_data": "" } print(f"\n [{i+1}] 添加: saddr={saddr}") ws_send(sock, add_msg) added_saddrs.add(saddr) time.sleep(2.0) # 2 秒间隔 # ---- 步骤 4: 等待 ws_task 推送 out 数据 ---- print("\n" + "=" * 70) print(" 步骤 4: 等待 ws_task 定时推送 out 数据(最多 10 秒)") print("=" * 70) out_push_received = False out_signals_list = [] start = time.time() while time.time() - start < 10.0: msgs = drain_sock(sock, timeout=2.0) for kind, msg in msgs: if isinstance(msg, dict): out_data = msg.get('out') if out_data is not None and isinstance(out_data, list) and len(out_data) > 0: out_push_received = True out_signals_list = out_data elapsed = time.time() - start print(f"\n [✓] 收到 out 数据推送! (耗时 {elapsed:.1f}s, {kind}帧)") print(f" out 当前注册信号数: {len(out_data)}") for j, o in enumerate(out_data): marker = " ← 新添加" if o.get('saddr') in added_saddrs else "" print(f" [{j+1}] saddr={o.get('saddr','?')} val={o.get('val','?')} desc={o.get('desc','?')}{marker}") break elif msg.get('type') == 'dc_data': pass # 忽略 dc_data 推送 elif msg.get('type'): pass elif isinstance(msg, str) and len(msg) > 5: pass # 可能是 cmd 输出 if out_push_received: break # ---- 步骤 5: 验证结果 ---- print("\n" + "=" * 70) print(" 步骤 5: 验证结果") print("=" * 70) if not out_push_received: print(" [✗] 未收到 out 数据推送") print(" 可能原因:") print(" 1. ws_task 定时器未触发(等待时间不足)") print(" 2. add_signal 失败(saddr 无效或信号已存在)") print(" 3. 服务端没有 has_change(新信号初始值与实际值相同)") else: registered_saddrs = {o.get('saddr') for o in out_signals_list} all_found = added_saddrs.issubset(registered_saddrs) missing = added_saddrs - registered_saddrs if all_found: print(f" [✓] 所有 {len(added_saddrs)} 个测试信号已成功注册到 out!") else: print(f" [✗] 以下 {len(missing)} 个信号未能在 out 中找到: {missing}") if missing: print(f"\n 当前 out 中注册的信号 saddr: {registered_saddrs}") # ---- 步骤 6: 清理 — 删除测试添加的信号 ---- print("\n" + "=" * 70) print(" 步骤 6: 清理 — 删除测试添加的信号") print("=" * 70) for sig in test_signals: saddr = sig.get('saddr', '') del_msg = { "saddr": saddr, "signal_type": "out", "curd": "del", "setting_zone": "0", "signal_data": "" } print(f" [→] 删除: saddr={saddr}") ws_send(sock, del_msg) time.sleep(0.5) # del 后立即推送,检查删除结果 msgs = drain_sock(sock, timeout=2.0) for kind, msg in msgs: if isinstance(msg, dict) and msg.get('out') is not None: print(f" [i] 删除后 out 信号数: {len(msg.get('out', []))}") # ---- 总结 ---- print("\n" + "=" * 70) print(" 测试总结") print("=" * 70) print(f" 数据中心 out 信号数: {len(out_signals)}") print(f" 测试添加信号数: {len(test_signals)}") if out_push_received: print(f" out 推送接收: ✅") print(f" 信号注册成功: {'✅ 全部通过' if all_found else f'❌ 缺少 {len(missing)} 个'}") else: print(f" out 推送接收: ❌ 未收到") print("=" * 70) # 发送关闭帧 close_frame = struct.pack('>H', 1000) frame = bytearray([0x88, 0x02]) frame.extend(close_frame) try: sock.sendall(bytes(frame)) except: pass sock.close() if __name__ == '__main__': main()