You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

268 lines
9.8 KiB

#!/usr/bin/env python3
"""
chat_integration_test.py — integration test: 2 chat nodes over UDP, sync messages.
Запускает два процесса utun в фоне, создаёт канал, пишет 4 сообщения,
подключает второй узел по invite-ссылке, проверяет members и messages.
Usage:
python3 tools/chat_integration_test.py
"""
import asyncio
import json
import os
import socket
import sys
import tempfile
import time
sys.path.insert(0, os.path.dirname(__file__))
from chat_client import ChatClient, ChatClientError
UTUN_BIN = os.path.join(os.path.dirname(__file__), "..", "src", "utun")
READY_TIMEOUT = 3.0 # max wait for utun ready
SYNC_TIMEOUT = 8.0 # max wait for join + sync
REQUEST_TIMEOUT = 5.0 # per-request timeout
TOTAL_TIMEOUT = 25.0 # overall test timeout
def find_free_port():
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
s.bind(("127.0.0.1", 0))
port = s.getsockname()[1]
s.close()
return port
def write_config(path, etcp_port, ctrl_port, db_subdir):
content = f"""[global]
db_path={db_subdir}
[server: srv]
addr=127.0.0.1:{etcp_port}
type=public
[chatserver]
db_path={db_subdir}
headless_control_bind=127.0.0.1:{ctrl_port}
[allowed_keys]
allow_all=1
"""
with open(path, "w") as f:
f.write(content)
async def kill_proc(proc, label):
if proc is None or proc.returncode is not None:
return
try:
proc.terminate()
try:
await asyncio.wait_for(proc.wait(), timeout=3.0)
except asyncio.TimeoutError:
proc.kill()
await proc.wait()
except ProcessLookupError:
pass
async def wait_ready(cli, timeout=READY_TIMEOUT):
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
try:
await cli.ping()
return True
except ChatClientError:
await asyncio.sleep(0.1)
return False
def check(name, expr, detail=""):
if not expr:
detail = f" ({detail})" if detail else ""
raise AssertionError(f"FAIL: {name}{detail}")
print(f" OK: {name}")
async def main():
etcp_a = find_free_port()
etcp_b = find_free_port()
ctrl_a = find_free_port()
ctrl_b = find_free_port()
print(f"ports: etcp={etcp_a},{etcp_b} ctrl={ctrl_a},{ctrl_b}")
tmpdir = tempfile.mkdtemp(prefix="utun_chat_test_")
db_a = os.path.join(tmpdir, "db_a")
db_b = os.path.join(tmpdir, "db_b")
os.makedirs(db_a, exist_ok=True)
os.makedirs(db_b, exist_ok=True)
config_a = os.path.join(tmpdir, "a.conf")
config_b = os.path.join(tmpdir, "b.conf")
write_config(config_a, etcp_a, ctrl_a, db_a)
write_config(config_b, etcp_b, ctrl_b, db_b)
proc_a = None
proc_b = None
try:
# ── Start utun processes ──
print("\n--- Starting utun ---")
log_a = os.path.join(tmpdir, "utun_a.log")
log_b = os.path.join(tmpdir, "utun_b.log")
pid_a = os.path.join(tmpdir, "utun_a.pid")
pid_b = os.path.join(tmpdir, "utun_b.pid")
proc_a = await asyncio.create_subprocess_exec(
UTUN_BIN, "-f", "-p", pid_a, "-l", log_a, "-c", config_a,
stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL,
)
proc_b = await asyncio.create_subprocess_exec(
UTUN_BIN, "-f", "-p", pid_b, "-l", log_b, "-c", config_b,
stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL,
)
print(f" proc_a pid={proc_a.pid} proc_b pid={proc_b.pid}")
await asyncio.sleep(0.5)
# ── Open connections ──
print("\n--- Connecting to headless control ---")
async with ChatClient(port=ctrl_a, timeout=REQUEST_TIMEOUT) as cli_a:
if not await wait_ready(cli_a):
raise RuntimeError(f"node A not ready after {READY_TIMEOUT}s")
print(f" node A ready")
async with ChatClient(port=ctrl_b, timeout=REQUEST_TIMEOUT) as cli_b:
if not await wait_ready(cli_b):
raise RuntimeError(f"node B not ready after {READY_TIMEOUT}s")
print(f" node B ready")
# ── Create channel ──
print("\n--- Create channel ---")
await cli_a.create_channel("TestGroup")
channels = await cli_a.channels()
check("channel created", isinstance(channels, list) and len(channels) == 1,
f"channels={json.dumps(channels)}")
ch_id = str(channels[0]["id"])
print(f" channel_id={ch_id} name={channels[0]['name']}")
# ── Write 4 messages ──
print("\n--- Write messages ---")
for i in range(4):
await cli_a.send(ch_id, f"Message {i+1}")
msgs_a = await cli_a.messages(ch_id, count=10)
check("4 messages on A", isinstance(msgs_a, list) and len(msgs_a) == 4,
f"got {len(msgs_a) if isinstance(msgs_a, list) else '?'}")
for i, m in enumerate(msgs_a):
print(f" [#{i+1}] ts={m.get('ts','?')} author={m.get('author_id','?')}")
# ── Invite ──
print("\n--- Invite link ---")
invite = await cli_a.invite(ch_id)
link = invite.get("link", "")
check("invite link", isinstance(link, str) and link.startswith("utun://"),
f"link={'...' + link[-20:] if link else 'MISSING'}")
print(f" link={link[:80]}...")
# ── Connect node B ──
print("\n--- Connect node B ---")
conn = await cli_b.connect_channel(link)
check("connect accepted", isinstance(conn, dict) and conn.get("connecting"),
f"resp={json.dumps(conn)}")
print(f" connecting={conn.get('connecting')}")
# Wait for join + sync
print(f"\n--- Wait sync (max {SYNC_TIMEOUT}s) ---")
deadline = time.monotonic() + SYNC_TIMEOUT
member_count = 0
while time.monotonic() < deadline:
try:
members = await cli_b.members(ch_id)
member_count = len(members) if isinstance(members, list) else 0
if member_count >= 2:
break
except ChatClientError:
pass
await asyncio.sleep(0.15)
check("2 members on B", member_count >= 2,
f"got {member_count} members after {SYNC_TIMEOUT}s")
# Small pause to let db_sync fully settle
await asyncio.sleep(0.5)
# ── Verify members on B ──
members = await cli_b.members(ch_id)
print(f"\n--- Members on B ({len(members)}) ---")
for m in members:
connected = "✓" if m.get("connected") else "✗"
online = "✓" if m.get("online") else "✗"
print(f" {m['node_id']} name={m.get('name','?')} online={online} connected={connected}")
check("owner online", any(m.get("online") for m in members),
f"members={json.dumps(members)}")
check("owner connected", any(m.get("connected") for m in members),
f"members={json.dumps(members)}")
# ── Verify messages on B ──
if proc_b.returncode is not None:
raise RuntimeError(f"node B died with code {proc_b.returncode}")
print(f"\n--- Messages on B ---")
try:
msgs_b = await cli_b.messages(ch_id, count=10)
except ChatClientError as e:
print(f" messages error: {e}")
print(f" proc_b.returncode={proc_b.returncode}")
raise
check("4 messages on B", isinstance(msgs_b, list) and len(msgs_b) == 4,
f"got {len(msgs_b) if isinstance(msgs_b, list) else '?'}")
for i, m in enumerate(msgs_b):
print(f" [#{i+1}] ts={m.get('ts','?')} author={m.get('author_id','?')}")
# ── Verify members on A ──
members_a = await cli_a.members(ch_id)
print(f"\n--- Members on A ({len(members_a)}) ---")
for m in members_a:
print(f" {m['node_id']} name={m.get('name','?')} online={m.get('online')} connected={m.get('connected')}")
check("2 members on A", isinstance(members_a, list) and len(members_a) >= 2,
f"got {len(members_a) if isinstance(members_a, list) else '?'}")
print("\n=== TEST PASSED ===")
return 0
except Exception as e:
print(f"\n=== TEST FAILED: {e} ===", file=sys.stderr)
import traceback
traceback.print_exc()
return 1
finally:
print("\n--- Cleanup ---")
await kill_proc(proc_a, "proc_a")
await kill_proc(proc_b, "proc_b")
for f in [config_a, config_b]:
try: os.unlink(f)
except OSError: pass
try: os.rmdir(db_a)
except OSError: pass
try: os.rmdir(db_b)
except OSError: pass
try: os.rmdir(tmpdir)
except OSError: pass
print(f" temp dir cleaned: {tmpdir}")
if __name__ == "__main__":
async def _run():
try:
return await asyncio.wait_for(main(), timeout=TOTAL_TIMEOUT)
except asyncio.TimeoutError:
print(f"\n=== TEST FAILED: total timeout {TOTAL_TIMEOUT}s ===", file=sys.stderr)
return 1
sys.exit(asyncio.run(_run()))