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

399 lines
18 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
"""Замер переподключения браузерного WSS после перезапуска frontend-proxy."""
import argparse
import asyncio
import json
import secrets
import sys
import time
from pathlib import Path
from uuid import uuid4
from playwright.async_api import async_playwright
from sqlalchemy import delete
from sqlalchemy import select
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
from load_browser import control, local_url, login_cookie
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT / "backend"))
from app.api.auth import hash_password # noqa: E402
from app.db.models import AuditLog, Session, Trainee, User # noqa: E402
async def create_temp_accounts(
database_url: str,
instructor_login: str,
instructor_password: str,
trainee_login: str,
trainee_password: str,
trainee_name: str,
):
"""Создать изолированную пару ролей и вернуть ID привязанного курсанта."""
engine = create_async_engine(database_url)
try:
factory = async_sessionmaker(engine, expire_on_commit=False)
async with factory() as db:
trainee = Trainee(name=trainee_name)
db.add(trainee)
await db.flush()
db.add_all([
User(
login=instructor_login,
full_name="Преподаватель проверки восстановления",
password_hash=hash_password(instructor_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()
return trainee.id
finally:
await engine.dispose()
async def cleanup_temp_accounts(
database_url: str,
logins: list[str],
trainee_name: str,
session_id,
*,
remove_accounts: bool,
) -> None:
engine = create_async_engine(database_url)
try:
factory = async_sessionmaker(engine, expire_on_commit=False)
async with factory() as db:
await db.execute(delete(Session).where(Session.id == session_id))
await db.execute(delete(AuditLog).where(AuditLog.object_id == str(session_id)))
if remove_accounts:
await db.execute(delete(AuditLog).where(AuditLog.actor.in_(logins)))
await db.execute(delete(User).where(User.login.in_(logins)))
await db.execute(delete(Trainee).where(Trainee.name == trainee_name))
await db.commit()
finally:
await engine.dispose()
async def compose_service(project: str, action: str, service: str) -> None:
process = await asyncio.create_subprocess_exec(
"docker", "compose", "--project-name", project, action, service,
cwd=ROOT,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
)
stdout, stderr = await process.communicate()
if process.returncode:
detail = (stderr or stdout).decode(errors="replace").strip()
raise RuntimeError(f"docker compose {action} {service}: {detail}")
async def run(args: argparse.Namespace) -> int:
frontend = local_url(args.frontend_url, {"https"})
backend = local_url(args.backend_url, {"http"})
if not args.compose_project.strip():
raise ValueError("--compose-project обязателен: recovery test останавливает Nginx")
if not 0 <= args.outage_seconds <= 30:
raise ValueError("--outage-seconds должен быть от 0 до 30")
ephemeral = not args.login and not args.password
if bool(args.login) != bool(args.password):
raise ValueError("login и password задаются вместе либо не задаются")
instructor_login = args.login
instructor_password = args.password
trainee_login = args.trainee_login
trainee_password = args.trainee_password
trainee_name = ""
trainee_id = None
session_uuid = uuid4()
if ephemeral:
suffix = session_uuid.hex[:10]
instructor_login = f"recovery-instructor-{suffix}"
instructor_password = secrets.token_urlsafe(24)
trainee_login = f"recovery-trainee-{suffix}"
trainee_password = secrets.token_urlsafe(24)
trainee_name = f"recovery-trainee-{suffix}"
trainee_id = await create_temp_accounts(
args.database_url,
instructor_login,
instructor_password,
trainee_login,
trainee_password,
trainee_name,
)
else:
if not trainee_login or not trainee_password:
raise ValueError("для существующих аккаунтов нужны trainee-login и trainee-password")
engine = create_async_engine(args.database_url)
try:
async with async_sessionmaker(engine, expire_on_commit=False)() as db:
row = (await db.execute(
select(Trainee.id, Trainee.name)
.join(User, User.trainee_id == Trainee.id)
.where(
User.login == trainee_login,
User.role == "trainee",
User.blocked.is_(False),
)
)).one_or_none()
if row is None:
raise ValueError("учётка курсанта не найдена или не привязана к профилю")
trainee_id, trainee_name = row
finally:
await engine.dispose()
assert instructor_login and instructor_password and trainee_login and trainee_password
instructor_cookie = login_cookie(backend, instructor_login, instructor_password)
trainee_cookie = login_cookie(backend, trainee_login, trainee_password)
cookie_name, cookie_value = trainee_cookie.split("=", 1)
ws_base = "ws" + backend.removeprefix("http")
session_id = str(session_uuid)
started = False
frontend_stopped = False
try:
await control(ws_base, session_id, instructor_cookie, {
"type": "scenario.start",
"scenario_id": args.scenario,
"trainee": trainee_name,
"trainee_id": str(trainee_id) if trainee_id else None,
"mode": "training",
"exercise": "card",
})
started = True
async with async_playwright() as playwright:
browser = await playwright.chromium.launch(headless=True)
context = await browser.new_context(
viewport={"width": 1280, "height": 720},
ignore_https_errors=args.ignore_https_errors,
)
await context.add_cookies([{
"name": cookie_name,
"value": cookie_value,
"url": frontend,
"httpOnly": True,
"secure": True,
"sameSite": "Lax",
}])
page = await context.new_page()
errors: list[str] = []
page.on("pageerror", lambda error: errors.append(str(error)))
await page.add_init_script("""(() => {
window.__lctDropOperatorKioAck = false;
window.__lctOperatorAckDropped = false;
let owner = WebSocket.prototype;
let descriptor;
while (owner && !descriptor) {
descriptor = Object.getOwnPropertyDescriptor(owner, 'onmessage');
if (!descriptor) owner = Object.getPrototypeOf(owner);
}
if (!owner || !descriptor?.set || !descriptor?.get) return;
Object.defineProperty(owner, 'onmessage', {
configurable: descriptor.configurable,
enumerable: descriptor.enumerable,
get() { return descriptor.get.call(this); },
set(handler) {
const socket = this;
const wrapped = (event) => {
if (window.__lctDropOperatorKioAck && typeof event.data === 'string') {
try {
const message = JSON.parse(event.data);
if (message.type === 'kio.patch' && message.source === 'operator') {
window.__lctDropOperatorKioAck = false;
window.__lctOperatorAckDropped = true;
socket.close(4000, 'test: discard checkpoint acknowledgement');
return;
}
} catch { /* non-JSON frames are passed to the application */ }
}
if (handler) handler.call(socket, event);
};
descriptor.set.call(socket, handler ? wrapped : null);
},
});
window.__lctAckInterceptorInstalled = Boolean(descriptor?.set && descriptor?.get);
})()""")
await page.goto(
f"{frontend}/trainee?session={session_id}",
wait_until="domcontentloaded",
timeout=args.timeout * 1000,
)
address = page.locator(".kio-arm-row-address input")
await address.wait_for(
timeout=args.timeout * 1000,
)
if not await page.evaluate("window.__lctAckInterceptorInstalled"):
raise RuntimeError("browser could not install the WebSocket acknowledgement test hook")
await page.locator(".arm-topbar-state .state-ok").wait_for(
timeout=args.timeout * 1000,
)
await page.evaluate("""() => {
const node = document.querySelector('.arm-topbar-state');
window.__lctRecoveryStates = [node?.textContent ?? 'missing'];
new MutationObserver(() => window.__lctRecoveryStates.push(
node?.textContent ?? 'missing'
)).observe(node, {attributes: true, childList: true, subtree: true});
}""")
buffered_value = "улица Буферная, дом 30"
started_at = time.perf_counter()
# Останавливаем proxy и вводим значение только после того, как
# browser отметил канал отключённым: это проверяет offline outbox.
await compose_service(args.compose_project, "stop", "frontend")
frontend_stopped = True
outage_started_at = time.perf_counter()
await page.wait_for_function(
"window.__lctRecoveryStates.some(value => !value.includes('АРМ подключён'))",
timeout=args.timeout * 1000,
)
# The operator edits only after the channel is known to be down.
# This exercises the disconnected outbox path instead of relying
# on a frame which might already have reached the server.
await address.fill(buffered_value)
await page.locator(".kio-arm-row-address.kio-arm-pending").wait_for(
timeout=args.timeout * 1000,
)
outage_remaining = args.outage_seconds - (time.perf_counter() - outage_started_at)
if outage_remaining > 0:
await asyncio.sleep(outage_remaining)
restore_started_at = time.perf_counter()
proxy_unavailable_seconds = restore_started_at - outage_started_at
await compose_service(args.compose_project, "start", "frontend")
frontend_stopped = False
await page.wait_for_function(
"window.__lctRecoveryStates.at(-1).includes('АРМ подключён')",
timeout=args.timeout * 1000,
)
recovery_seconds = time.perf_counter() - restore_started_at
await page.wait_for_function(
"!document.querySelector('.kio-arm-row-address')?.classList.contains('kio-arm-pending')",
timeout=args.timeout * 1000,
)
buffered_value_after = await page.locator(".kio-arm-row-address input").input_value()
# Now send while the WSS is open, allow PostgreSQL checkpoint/echo
# to complete, then deliberately discard that one echo and close
# the browser socket. This targets the ambiguous commit/ack window:
# reconnect must recover the stored value without duplicating a
# non-replace-style action or leaving the field pending.
committed_value = "улица Подтверждённая, дом 31"
await page.evaluate("window.__lctDropOperatorKioAck = true")
await address.fill(committed_value)
await page.wait_for_function(
"window.__lctOperatorAckDropped === true",
timeout=args.timeout * 1000,
)
await page.wait_for_function(
"window.__lctRecoveryStates.some(value => !value.includes('АРМ подключён'))",
timeout=args.timeout * 1000,
)
await page.wait_for_function(
"window.__lctRecoveryStates.at(-1).includes('АРМ подключён')",
timeout=args.timeout * 1000,
)
await page.wait_for_function(
"!document.querySelector('.kio-arm-row-address')?.classList.contains('kio-arm-pending')",
timeout=args.timeout * 1000,
)
committed_value_after = await page.locator(".kio-arm-row-address input").input_value()
ack_deliberately_dropped = await page.evaluate("window.__lctOperatorAckDropped")
pending_cleared = not await page.locator(".kio-arm-row-address").evaluate(
"element => element.classList.contains('kio-arm-pending')"
)
total_seconds = time.perf_counter() - started_at
states = await page.evaluate("window.__lctRecoveryStates")
await context.close()
await browser.close()
result = {
"frontend": frontend,
"session_id": session_id,
"component_restarted": "frontend (Nginx TLS/WSS proxy) container stop/start",
"observed_states": states,
"requested_outage_seconds": args.outage_seconds,
"proxy_unavailable_seconds": proxy_unavailable_seconds,
"recovery_after_restore_seconds": recovery_seconds,
"total_seconds": total_seconds,
"browser_errors": errors,
"buffered_patch": {
"field": "address",
"value": buffered_value_after,
"confirmed_by_server": buffered_value_after == buffered_value,
},
"lost_ack_recovery": {
"ack_deliberately_dropped": ack_deliberately_dropped,
"value_after_reconnect": committed_value_after,
"confirmed_by_recovered_server_state": committed_value_after == committed_value,
"pending_cleared": pending_cleared,
},
"pass_recovery_30s": (
args.outage_seconds <= 30
and proxy_unavailable_seconds <= 30.5
and recovery_seconds <= 30
and not errors
and buffered_value_after == buffered_value
and committed_value_after == committed_value
and ack_deliberately_dropped
and pending_cleared
),
"scope": (
"one live card exercise; Nginx is unavailable while the trainee enters an address; "
"the disconnected KIO patch is queued and confirmed after WSS reconnect; a second "
"KIO patch is committed while connected, its checkpoint echo is deliberately "
"discarded, and reconnect recovers that state; backend remains running"
),
}
rendered = json.dumps(result, ensure_ascii=False, indent=2)
print(rendered)
if args.output:
with open(args.output, "w", encoding="utf-8") as stream:
stream.write(rendered + "\n")
return 0 if result["pass_recovery_30s"] else 1
finally:
if frontend_stopped:
await compose_service(args.compose_project, "start", "frontend")
if started:
await control(ws_base, session_id, instructor_cookie, {"type": "session.stop"})
if ephemeral:
await cleanup_temp_accounts(
args.database_url,
[instructor_login, trainee_login],
trainee_name,
session_uuid,
remove_accounts=ephemeral,
)
if __name__ == "__main__":
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--frontend-url", default="https://127.0.0.1:5443")
parser.add_argument("--backend-url", default="http://127.0.0.1:8000")
parser.add_argument("--compose-project", required=True,
help="точное имя изолированного Compose project; тест останавливает его frontend")
parser.add_argument("--login")
parser.add_argument("--password")
parser.add_argument("--trainee-login")
parser.add_argument("--trainee-password")
parser.add_argument(
"--database-url",
default="postgresql+asyncpg://lct:lct@127.0.0.1:5432/lct",
)
parser.add_argument("--scenario", default="fire-apartment-l2")
parser.add_argument("--timeout", type=float, default=30)
parser.add_argument("--outage-seconds", type=float, default=0,
help="держать frontend недоступным столько секунд (0–30), вводить карточку после разрыва")
parser.add_argument("--ignore-https-errors", action="store_true")
parser.add_argument("--output")
try:
raise SystemExit(asyncio.run(run(parser.parse_args())))
except Exception as exc:
print(f"recovery check failed: {type(exc).__name__}: {exc}")
raise SystemExit(2)