90 lines
3.1 KiB
Python
90 lines
3.1 KiB
Python
#!/usr/bin/env python3
|
||
"""验证 curd:add 立即推送修复"""
|
||
import socket, base64, os, struct, json, time
|
||
|
||
IP = '198.120.0.100'; PORT = 8000
|
||
|
||
def ws_handshake(sock):
|
||
key = base64.b64encode(os.urandom(16)).decode()
|
||
req = (f"GET /ws HTTP/1.1\r\nHost: {IP}:{PORT}\r\nUpgrade: websocket\r\n"
|
||
f"Connection: Upgrade\r\nSec-WebSocket-Key: {key}\r\n"
|
||
f"Sec-WebSocket-Version: 13\r\n\r\n")
|
||
sock.sendall(req.encode())
|
||
resp = b""
|
||
while b"\r\n\r\n" not in resp:
|
||
c = sock.recv(4096)
|
||
if not c: break
|
||
resp += c
|
||
return resp
|
||
|
||
def ws_send(sock, msg):
|
||
if isinstance(msg, dict): msg = json.dumps(msg, ensure_ascii=False)
|
||
p = msg.encode(); L = len(p); mask = os.urandom(4)
|
||
m = bytes(p[i] ^ mask[i%4] for i in range(L))
|
||
f = bytearray([0x81])
|
||
if L<126: f.append(0x80|L)
|
||
elif L<65536: f.append(0x80|126); f.extend(struct.pack('>H',L))
|
||
else: f.append(0x80|127); f.extend(struct.pack('>Q',L))
|
||
f.extend(mask); f.extend(m); sock.sendall(bytes(f))
|
||
|
||
def ws_recv_all(sock, t=2.0):
|
||
sock.settimeout(t); data=b""
|
||
try:
|
||
while True:
|
||
c=sock.recv(65536)
|
||
if not c: break
|
||
data+=c
|
||
except socket.timeout: pass
|
||
msgs=[]; offset=0
|
||
while offset+2<=len(data):
|
||
b0,b1=data[offset],data[offset+1]; op=b0&0xF; L=b1&0x7F; offset+=2
|
||
if L==126: L=struct.unpack('>H',data[offset:offset+2])[0]; offset+=2
|
||
elif L==127: L=struct.unpack('>Q',data[offset:offset+8])[0]; offset+=8
|
||
if op in (1,2) and offset+L<=len(data):
|
||
p=data[offset:offset+L]
|
||
try: msgs.append(json.loads(p.decode()))
|
||
except: msgs.append(p.decode(errors='replace'))
|
||
offset+=L
|
||
return msgs
|
||
|
||
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||
sock.settimeout(10); sock.connect((IP,PORT))
|
||
ws_handshake(sock); print("[✓] 握手成功")
|
||
|
||
# 注册 dc out
|
||
ws_send(sock, {"type":"cmd","data":"datacenter out"})
|
||
time.sleep(1); ws_recv_all(sock,1)
|
||
|
||
# 测试立即推送
|
||
target = "sys.ch.tcp_s0_if"
|
||
print(f"\n[→] curd:add saddr={target}")
|
||
ws_send(sock, {"saddr":target,"signal_type":"out","curd":"add","setting_zone":"0","signal_data":""})
|
||
|
||
# 立即检查推送(100ms 内应有立即推送)
|
||
time.sleep(0.3)
|
||
msgs = ws_recv_all(sock, 2)
|
||
found = False
|
||
for msg in msgs:
|
||
if isinstance(msg, dict) and 'out' in msg and isinstance(msg['out'], list):
|
||
elapsed = "立即" # curd:add后直接推送
|
||
print(f"[✓] add后即收到 out 推送 ({elapsed}): {len(msg['out'])} 个信号")
|
||
for o in msg['out']:
|
||
print(f" saddr={o.get('saddr')} val={o.get('val')}")
|
||
found = True
|
||
|
||
if not found:
|
||
print("[✗] add后未立即收到推送,可能需要等待 ws_task 轮询")
|
||
time.sleep(2)
|
||
msgs = ws_recv_all(sock, 3)
|
||
for msg in msgs:
|
||
if isinstance(msg, dict) and 'out' in msg:
|
||
print(f"[i] 延迟收到推送: {len(msg['out'])} 个信号")
|
||
|
||
# 清理
|
||
ws_send(sock, {"saddr":target,"signal_type":"out","curd":"del","setting_zone":"0","signal_data":""})
|
||
time.sleep(0.5); ws_recv_all(sock, 1)
|
||
try: sock.sendall(b'\x88\x80'+os.urandom(4))
|
||
except: pass
|
||
sock.close()
|
||
print("\n[✓] 测试完成")
|