375 lines
13 KiB
Python
375 lines
13 KiB
Python
#!/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()
|