lct-hack/scripts/test_cluster_scenario_registry.py
2026-09-26 17:13:45 +00:00

402 lines
20 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""Cross-process smoke: publish a trainee scenario on backend A, start it on B.
Use only with a disposable local Compose project and its own PostgreSQL. The
smoke creates temporary accounts, group, proposal, and DDS session and removes
their rows in ``finally``. It never writes evidence into the repository.
"""
from __future__ import annotations
import argparse
import asyncio
import json
import os
from pathlib import Path
import secrets
import ssl
import subprocess
import time
import urllib.error
import urllib.request
from uuid import UUID, uuid4
import websockets
def compose(args: argparse.Namespace, *command: str) -> str:
env = os.environ.copy()
env.update({
"POSTGRES_PORT": str(args.postgres_port),
"BACKEND_PORT": str(args.backend_a_port),
"BACKEND_B_PORT": str(args.backend_b_port),
"FRONTEND_PORT": str(args.frontend_port),
"TLS_PORT": str(args.tls_port),
"POSTGRES_PASSWORD": args.postgres_password,
"SESSION_SECRET": args.session_secret,
"TLS_CERT_DIR": args.tls_cert_dir,
})
files = [
"docker-compose.yml", "docker-compose.tls.yml",
"docker-compose.load-test.yml", "docker-compose.cluster.yml",
]
command_result = subprocess.run(
["docker", "compose", "-p", args.project,
*[part for file in files for part in ("-f", file)], *command],
check=False, capture_output=True, text=True, env=env,
)
if command_result.returncode:
raise RuntimeError(
f"docker compose {' '.join(command)} failed: {command_result.stderr.strip()}"
)
return command_result.stdout
def request_json(base: str, path: str, *, cookie: str | None = None,
payload: dict | None = None) -> dict | list:
headers = {"Content-Type": "application/json"}
if cookie:
headers["Cookie"] = cookie
data = json.dumps(payload, ensure_ascii=False).encode() if payload is not None else None
request = urllib.request.Request(f"{base}{path}", data=data, headers=headers,
method="POST" if data is not None else "GET")
with urllib.request.urlopen(request, timeout=15) as response:
return json.loads(response.read())
def login(base: str, login_name: str, password: str) -> str:
payload = json.dumps({"login": login_name, "password": password}).encode()
request = urllib.request.Request(
f"{base}/api/auth/login", data=payload, method="POST",
headers={"Content-Type": "application/json"},
)
with urllib.request.urlopen(request, timeout=15) as response:
cookie = response.headers.get("Set-Cookie", "").split(";", 1)[0]
if not cookie.startswith("lct_session="):
raise RuntimeError(f"{base} did not issue the session cookie")
return cookie
async def main(args: argparse.Namespace) -> int:
from sqlalchemy import delete, select
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
from app.api.auth import hash_password
from app.db.models import (
AuditLog, Group, Scenario as ScenarioRow, ScenarioSubmission,
Session, Trainee, User,
)
from app.domain import ekp
from app.scenarios import store
db_url = (f"postgresql+asyncpg://lct:{args.postgres_password}@127.0.0.1:"
f"{args.postgres_port}/lct")
engine = create_async_engine(db_url, pool_pre_ping=True)
factory = async_sessionmaker(engine, expire_on_commit=False)
suffix = uuid4().hex[:12]
teacher_login, trainee_login = f"registry-i-{suffix}", f"registry-t-{suffix}"
teacher_password, trainee_password = secrets.token_urlsafe(24), secrets.token_urlsafe(24)
teacher_name, trainee_name = "Registry smoke instructor", f"Registry smoke {suffix}"
group = Group(name=f"registry-smoke-{suffix}", owner_login=teacher_login)
trainee = Trainee(name=trainee_name)
scenario_id: str | None = None
session_id = uuid4()
session_ids: list[UUID] = []
submission_id: UUID | None = None
teacher_cookie = trainee_cookie = None
session_ws_base = None
ssl_context = None
result: dict = {"checks": {}}
a = f"http://127.0.0.1:{args.backend_a_port}"
b = f"http://127.0.0.1:{args.backend_b_port}"
try:
store.load_from_disk(Path(args.scenarios))
source = store.get(args.source_scenario)
if source is None or source.ground_truth.dds is None:
raise RuntimeError(f"invalid source scenario {args.source_scenario}")
trainee.group_id = None
async with factory() as db:
db.add(group)
await db.flush()
trainee.group_id = group.id
db.add(trainee)
await db.flush()
db.add_all([
User(login=teacher_login, full_name=teacher_name,
password_hash=hash_password(teacher_password), role="instructor",
blocked=False),
User(login=trainee_login, full_name=trainee_name,
password_hash=hash_password(trainee_password), role="trainee",
trainee_id=trainee.id, blocked=False),
])
await db.commit()
result["checks"]["backend_a_healthy"] = request_json(a, "/api/health").get("status") == "ok"
result["checks"]["backend_b_healthy"] = request_json(b, "/api/health").get("status") == "ok"
teacher_cookie = login(a, teacher_login, teacher_password)
trainee_cookie = login(a, trainee_login, trainee_password)
incident_group = ekp.incident(source.ground_truth.incident_code).group
created = request_json(a, "/api/scenario-submissions", cookie=trainee_cookie, payload={
"title": f"Межузловая проверка {suffix}", "level": source.level.value,
"kio": {
"caller_name": "Тестовый заявитель",
"caller_contact": "+7 900 000-00-00",
"address": f"Москва, учебная улица, дом {suffix[:3]}",
"description": "На учебном объекте обнаружено задымление и нужна бригада.",
"incident_group": incident_group,
"signs": source.signs,
"incident_type": source.type.value,
"dds": source.ground_truth.dds.value,
"victims_count": 0,
"fire": {"object_kind": "учебное помещение", "fire_nature": "задымление"},
},
})
submission_id = UUID(created["id"])
approved = request_json(a, f"/api/scenario-submissions/{submission_id}/review",
cookie=teacher_cookie,
payload={"decision": "approve", "comment": "Межузловой smoke."})
scenario_id = approved["scenario_id"]
result["scenario_id"] = scenario_id
result["published_on"] = "backend-a"
result["checks"]["submission_approved_on_a"] = approved["status"] == "approved"
# WS start is the first scenario operation on the selected target node.
# With --frontend-url, UUIDs are tried until Nginx consistent-hash routes
# one to B; every candidate is pinned for all session channels.
session_ws_base = args.frontend_url.replace("https://", "wss://").replace(
"http://", "ws://"
) if args.frontend_url else b.replace("http://", "ws://")
ssl_context = ssl._create_unverified_context() if session_ws_base.startswith("wss://") else None
target_backend = "backend-b"
max_route_attempts = 12 if args.frontend_url else 1
session_row = None
for route_attempt in range(max_route_attempts):
candidate_id = uuid4()
session_ids.append(candidate_id)
session_id = candidate_id
ws_url = f"{session_ws_base}/ws/control/{session_id}"
async with websockets.connect(
ws_url, additional_headers={"Cookie": teacher_cookie}, ssl=ssl_context,
open_timeout=15, ping_interval=None,
) as control:
await control.send(json.dumps({
"type": "scenario.start", "scenario_id": scenario_id,
"scenario_ids": [scenario_id], "trainee": trainee_name,
"trainee_id": str(trainee.id), "dds_service": created["kio"]["notify"][0],
"mode": "training", "exercise": "dds",
}, ensure_ascii=False))
async with factory() as db:
deadline = time.monotonic() + 8
session_row = None
while time.monotonic() < deadline:
session_row = await db.scalar(
select(Session).where(Session.id == candidate_id)
)
if session_row and session_row.backend_node_id:
break
await asyncio.sleep(0.1)
if session_row and session_row.backend_node_id == target_backend:
break
if session_row and session_row.backend_node_id:
result.setdefault("nginx_route_attempts", []).append({
"session_id": str(candidate_id),
"owner": session_row.backend_node_id,
})
await control.send(json.dumps({"type": "session.stop"}))
await asyncio.sleep(0.15)
else:
raise RuntimeError("session.start was not persisted through the selected route")
result["checks"]["session_started_on_backend_b"] = bool(
session_row and session_row.backend_node_id == target_backend
and session_row.scenario_id == scenario_id and session_row.live_state
)
if not result["checks"]["session_started_on_backend_b"]:
raise RuntimeError("could not route a fresh session to backend B")
result["session_id"] = str(session_id)
result["route_attempts"] = len(session_ids)
# Prove the user-visible catalog/detail routes while B is still alive.
catalog = request_json(b, "/api/scenarios", cookie=teacher_cookie)
entry = next((item for item in catalog if item["id"] == scenario_id), None)
result["checks"]["visible_in_backend_b_catalog"] = entry is not None
result["checks"]["owned_by_teacher_on_backend_b"] = bool(entry and entry["can_manage"])
detail = request_json(b, f"/api/scenarios/{scenario_id}", cookie=teacher_cookie)
result["checks"]["detail_loaded_from_backend_b"] = (
detail.get("student_card", {}).get("address")
== created["kio"]["address"]
)
if args.failure_mode == "kill-backend-b":
old_epoch = session_row.backend_fencing_epoch
compose(args, "kill", "backend-b")
result["checks"]["backend_b_killed_after_start"] = True
deadline = time.monotonic() + args.takeover_timeout
takeover_row = None
while time.monotonic() < deadline:
async with factory() as db:
takeover_row = await db.scalar(
select(Session).where(Session.id == session_id)
)
if takeover_row and takeover_row.backend_node_id == "backend-a":
break
await asyncio.sleep(0.25)
result["checks"]["takeover_to_backend_a"] = bool(
takeover_row and takeover_row.backend_node_id == "backend-a"
)
result["checks"]["fencing_epoch_incremented"] = bool(
takeover_row and takeover_row.backend_fencing_epoch > old_epoch
)
result["checks"]["checkpoint_restored_after_kill"] = bool(
takeover_row and takeover_row.live_state and takeover_row.checkpoint_at
)
if not all(result["checks"][key] for key in (
"takeover_to_backend_a", "fencing_epoch_incremented",
"checkpoint_restored_after_kill",
)):
raise RuntimeError("backend A did not take over the killed B session")
result["takeover_seconds"] = round(
args.takeover_timeout - max(0, deadline - time.monotonic()), 2
)
station_url = f"{session_ws_base}/ws/station/{session_id}"
station_owner_key = (
"dds_card_received_after_takeover" if args.failure_mode == "kill-backend-b"
else "dds_card_received_on_backend_b"
)
station_snapshot = None
received_card = None
last_station_error = None
for reconnect_attempt in range(8):
try:
async with websockets.connect(
station_url, additional_headers={"Cookie": trainee_cookie}, ssl=ssl_context,
open_timeout=10, ping_interval=None,
) as station:
for _ in range(60):
event = json.loads(await asyncio.wait_for(station.recv(), timeout=10))
if event.get("type") == "error":
raise RuntimeError(
f"station error after route/failover: {event.get('message')}"
)
if event.get("type") == "card.received":
received_card = event["card"]
elif event.get("type") == "station.state":
station_snapshot = event.get("snapshot")
if received_card is not None and station_snapshot is not None:
break
if received_card is None or station_snapshot is None:
raise RuntimeError("station reconnect did not restore card and snapshot")
break
except (TimeoutError, OSError, websockets.exceptions.WebSocketException,
RuntimeError) as exc:
last_station_error = exc
if reconnect_attempt == 7:
raise RuntimeError(f"station did not recover through proxy: {exc}") from exc
await asyncio.sleep(1)
result["checks"][station_owner_key] = received_card is not None
result["checks"]["approved_kio_preserved"] = bool(
received_card
and received_card.get("address") == created["kio"]["address"]
and received_card.get("caller_name") == created["kio"]["caller_name"]
)
if args.failure_mode == "kill-backend-b":
command_id = str(uuid4())
command = {
"type": "card.status", "service": station_snapshot["managed_service"],
"status": "accepted",
"comment": (
"Основание: доклад по карточке.\n"
f"Сведения: карточка повторно принята после takeover; {station_snapshot['managed_service']} notified."
),
"_command_id": command_id,
}
async with websockets.connect(
station_url, additional_headers={"Cookie": trainee_cookie}, ssl=ssl_context,
open_timeout=10, ping_interval=None,
) as station:
await station.send(json.dumps(command, ensure_ascii=False))
acked = False
for _ in range(40):
event = json.loads(await asyncio.wait_for(station.recv(), timeout=10))
if event.get("type") == "error":
raise RuntimeError(f"station command failed after takeover: {event.get('message')}")
if event.get("type") == "command.ack" and event.get("command_id") == command_id:
acked = True
break
result["checks"]["station_command_acked_after_takeover"] = acked
async with factory() as db:
restored = await db.scalar(select(Session).where(Session.id == session_id))
result["checks"]["post_takeover_command_checkpointed"] = bool(
restored and restored.backend_node_id == "backend-a"
and command_id in (restored.live_state or {}).get(
"processed_station_commands", []
)
)
result["pass"] = all(result["checks"].values())
print(json.dumps(result, ensure_ascii=False, indent=2))
return 0 if result["pass"] else 1
finally:
if args.failure_mode == "kill-backend-b" and args.project:
try:
compose(args, "start", "backend-b")
except Exception:
pass
if teacher_cookie and session_ws_base:
try:
stop_base = (
f"ws://127.0.0.1:{args.backend_a_port}"
if args.failure_mode == "kill-backend-b" else session_ws_base
)
stop_ssl = None if stop_base.startswith("ws://") else ssl_context
async with websockets.connect(
f"{stop_base}/ws/control/{session_id}",
additional_headers={"Cookie": teacher_cookie}, ssl=stop_ssl,
open_timeout=5, ping_interval=None,
) as control:
await control.send(json.dumps({"type": "session.stop"}))
except Exception:
pass
async with factory() as db:
await db.execute(delete(AuditLog).where(
AuditLog.object_id.in_([str(item) for item in session_ids])
))
if submission_id:
await db.execute(delete(ScenarioSubmission).where(
ScenarioSubmission.id == submission_id
))
if session_ids:
await db.execute(delete(Session).where(Session.id.in_(session_ids)))
if scenario_id:
await db.execute(delete(ScenarioRow).where(ScenarioRow.id == scenario_id))
await db.execute(delete(AuditLog).where(
AuditLog.actor.in_([teacher_login, trainee_login])
))
await db.execute(delete(User).where(User.login.in_([teacher_login, trainee_login])))
await db.execute(delete(Trainee).where(Trainee.name == trainee_name))
await db.execute(delete(Group).where(Group.id == group.id))
await db.commit()
await engine.dispose()
if __name__ == "__main__":
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--postgres-port", type=int, required=True)
parser.add_argument("--backend-a-port", type=int, required=True)
parser.add_argument("--backend-b-port", type=int, required=True)
parser.add_argument("--postgres-password", required=True)
parser.add_argument("--frontend-url", help="optional TLS Nginx URL; tests consistent-hash routing")
parser.add_argument("--failure-mode", choices=("none", "kill-backend-b"), default="none")
parser.add_argument("--project", help="Compose project, required for --failure-mode kill-backend-b")
parser.add_argument("--tls-cert-dir", default="")
parser.add_argument("--session-secret", default="")
parser.add_argument("--frontend-port", type=int, default=15180)
parser.add_argument("--tls-port", type=int, default=15443)
parser.add_argument("--takeover-timeout", type=float, default=45)
parser.add_argument("--source-scenario", default="t01-1-fire-container")
parser.add_argument("--scenarios", default="scenarios")
parsed = parser.parse_args()
if parsed.failure_mode == "kill-backend-b" and not (
parsed.frontend_url and parsed.project and parsed.session_secret and parsed.tls_cert_dir
):
parser.error("kill-backend-b requires --frontend-url, --project, --session-secret and --tls-cert-dir")
raise SystemExit(asyncio.run(main(parsed)))