feat: один прод-compose со всеми моделями из готовых образов

This commit is contained in:
gglamer 2026-09-28 21:14:38 +00:00
commit a694702d8f
44 changed files with 1126 additions and 2075 deletions

View file

@ -1,27 +0,0 @@
#!/usr/bin/env bash
set -euo pipefail
repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
compose=(docker compose -f "$repo_root/docker-compose.yml" \
-f "$repo_root/docker-compose.production.yml" -f "$repo_root/docker-compose.tls.yml")
if env -u POSTGRES_PASSWORD "${compose[@]}" config --quiet >/dev/null 2>&1; then
echo "ОШИБКА: production Compose запустился без POSTGRES_PASSWORD" >&2
exit 1
fi
secret="compose-check-0123456789abcdef0123456789abcdef"
case "$secret" in
*[!A-Za-z0-9._~-]*) echo "ОШИБКА: тестовый пароль не URL-safe" >&2; exit 1 ;;
esac
resolved="$(POSTGRES_PASSWORD="$secret" "${compose[@]}" config --format json 2>/dev/null)"
printf '%s' "$resolved" | jq -e --arg secret "$secret" '
.services.postgres.environment.POSTGRES_PASSWORD == $secret and
.services.backend.environment.DATABASE_URL == ("postgresql+asyncpg://lct:" + $secret + "@postgres:5432/lct") and
.services.backup.environment.DATABASE_URL == ("postgresql+asyncpg://lct:" + $secret + "@postgres:5432/lct") and
.services.backend.environment.APP_ENV == "production" and
.services.backend.environment.DEMO_NO_DB == "false" and
.services.backend.environment.DEV_AUTH_BYPASS == "false" and
.services.backend.environment.SECURE_COOKIES == "true"
' >/dev/null
echo "Production Compose требует уникальный пароль и передаёт его PostgreSQL/backend/backup."

View file

@ -21,8 +21,7 @@ def remove_smoke_recordings(
run = runner or subprocess.run
result = run(
[
"docker", "compose", "-f", "docker-compose.yml", "-f", "docker-compose.sip.yml",
"exec", "-T", "sip", "rm", "-f", "--",
"docker", "compose", "exec", "-T", "sip", "rm", "-f", "--",
*(f"/recordings/{name}" for name in safe_names),
],
cwd=root,
@ -34,8 +33,7 @@ def remove_smoke_recordings(
return False
verification = run(
[
"docker", "compose", "-f", "docker-compose.yml", "-f", "docker-compose.sip.yml",
"exec", "-T", "sip", "find", "/recordings", "-maxdepth", "1",
"docker", "compose", "exec", "-T", "sip", "find", "/recordings", "-maxdepth", "1",
"-type", "f", "-name", "*.wav",
],
cwd=root,

View file

@ -117,8 +117,7 @@ def receive_final(sock: socket.socket, timeout: float = 3.0) -> tuple[str, dict[
def compose_password(user: str) -> str:
result = subprocess.run(
[
"docker", "compose", "-f", "docker-compose.yml", "-f", "docker-compose.sip.yml",
"exec", "-T", "sip", "cat", "/var/lib/lct-sip/credentials.env",
"docker", "compose", "exec", "-T", "sip", "cat", "/var/lib/lct-sip/credentials.env",
],
check=True,
capture_output=True,

View file

@ -59,8 +59,7 @@ async def cleanup_account(database_url: str, login: str) -> None:
def recordings() -> dict[str, int]:
result = subprocess.run(
[
"docker", "compose", "-f", "docker-compose.yml", "-f", "docker-compose.sip.yml",
"exec", "-T", "sip", "find", "/recordings", "-maxdepth", "1", "-type", "f",
"docker", "compose", "exec", "-T", "sip", "find", "/recordings", "-maxdepth", "1", "-type", "f",
"-name", "*.wav", "-printf", "%f %s\\n",
],
cwd=ROOT,
@ -85,7 +84,7 @@ async def configure(page, extension: str, password: str, target: str, timeout_ms
async def run(args: argparse.Namespace) -> int:
frontend = local_url(args.frontend_url, {"https"})
frontend = local_url(args.frontend_url, {"http", "https"})
backend = local_url(args.backend_url, {"http"})
login = f"webrtc-smoke-{secrets.token_hex(5)}"
app_password = secrets.token_urlsafe(24)
@ -219,7 +218,7 @@ async def run(args: argparse.Namespace) -> int:
if __name__ == "__main__":
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--frontend-url", default="https://127.0.0.1:5443")
parser.add_argument("--frontend-url", default="http://127.0.0.1:5173")
parser.add_argument("--backend-url", default="http://127.0.0.1:8000")
parser.add_argument(
"--database-url", default="postgresql+asyncpg://lct:lct@127.0.0.1:5432/lct"

View file

@ -1,760 +0,0 @@
#!/usr/bin/env python3
"""Live smoke: verify a complete DDS exercise across fenced cluster takeover.
Only use with a disposable local Compose cluster and its own PostgreSQL database.
The script creates temporary users/session, stops and restarts one named backend,
and removes the temporary rows on completion.
"""
from __future__ import annotations
import argparse
import asyncio
from contextlib import AsyncExitStack
import json
import os
import secrets
import ssl
import subprocess
import sys
import time
import urllib.error
import urllib.request
from pathlib import Path
from uuid import UUID, uuid4
import websockets
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT / "backend"))
def compose(args: argparse.Namespace, *command: str, check: bool = True) -> str:
env = os.environ.copy()
env.update({
"POSTGRES_PORT": str(args.postgres_port),
"BACKEND_PORT": str(args.backend_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",
]
if args.failure_mode == "partition":
files.append("docker-compose.partition-test.yml")
result = subprocess.run(
["docker", "compose", "-p", args.project,
*[part for file in files for part in ("-f", str(ROOT / file))], *command],
cwd=ROOT, env=env, check=False, capture_output=True, text=True,
)
if check and result.returncode:
raise subprocess.CalledProcessError(
result.returncode, result.args, output=result.stdout, stderr=result.stderr
)
return result.stdout
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=10) as response:
cookie = response.headers.get("Set-Cookie", "").split(";", 1)[0]
if not cookie.startswith("lct_session="):
raise RuntimeError("login did not issue lct_session cookie")
return cookie
async def send_control(ws_base: str, session_id: UUID, cookie: str, payload: dict) -> float:
async with websockets.connect(
f"{ws_base}/ws/control/{session_id}",
additional_headers={"Cookie": cookie},
ssl=ssl._create_unverified_context() if ws_base.startswith("wss://") else None,
open_timeout=10,
ping_interval=None,
) as socket:
await socket.send(json.dumps(payload, ensure_ascii=False))
sent_at = time.monotonic()
await asyncio.sleep(0.15)
return sent_at
async def probe_ws(ws_base: str, path: str, cookie: str) -> dict:
try:
async with websockets.connect(
f"{ws_base}{path}", additional_headers={"Cookie": cookie},
ssl=ssl._create_unverified_context() if ws_base.startswith("wss://") else None,
open_timeout=6, ping_interval=None,
) as socket:
try:
message = await asyncio.wait_for(socket.recv(), timeout=4)
event = json.loads(message) if isinstance(message, str) else "binary"
except TimeoutError:
event = "connected_no_initial_event"
return {"connected": True, "initial_event": event}
except Exception as exc: # report exact endpoint failure in the result
return {"connected": False, "error": f"{type(exc).__name__}: {exc}"}
async def station_command(
ws_base: str, session_id: UUID, cookie: str, payload: dict | None = None,
) -> dict:
"""Read the authoritative station snapshot, optionally issue one command."""
async with websockets.connect(
f"{ws_base}/ws/station/{session_id}", additional_headers={"Cookie": cookie},
ssl=ssl._create_unverified_context() if ws_base.startswith("wss://") else None,
open_timeout=10, ping_interval=None,
) as socket:
async def receive_snapshot() -> dict:
for _ in range(60):
raw = await asyncio.wait_for(socket.recv(), timeout=10)
event = json.loads(raw) if isinstance(raw, str) else {}
if event.get("type") == "error":
raise RuntimeError(f"station rejected smoke: {event.get('message')}")
if event.get("type") == "station.state":
return event["snapshot"]
raise RuntimeError("station did not send its state snapshot")
snapshot = await receive_snapshot()
if payload is None:
return snapshot
await socket.send(json.dumps(payload, ensure_ascii=False))
return await receive_snapshot()
async def replay_station_command(
ws_base: str, session_id: UUID, cookie: str, payload: dict, command_id: str,
attempts: int = 5,
) -> bool:
"""Replay one already committed command and require its durable ACK."""
last_error = None
for attempt in range(attempts):
try:
async with websockets.connect(
f"{ws_base}/ws/station/{session_id}", additional_headers={"Cookie": cookie},
ssl=ssl._create_unverified_context() if ws_base.startswith("wss://") else None,
open_timeout=10, ping_interval=None,
) as socket:
for _ in range(60):
raw = await asyncio.wait_for(socket.recv(), timeout=10)
event = json.loads(raw) if isinstance(raw, str) else {}
if event.get("type") == "error":
raise RuntimeError(f"station reconnect rejected: {event.get('message')}")
if event.get("type") == "station.state":
break
else:
raise RuntimeError("replacement station did not send a state snapshot")
replay = {**payload, "_command_id": command_id}
await socket.send(json.dumps(replay, ensure_ascii=False))
for _ in range(60):
raw = await asyncio.wait_for(socket.recv(), timeout=10)
event = json.loads(raw) if isinstance(raw, str) else {}
if event.get("type") == "error":
raise RuntimeError(f"station replay rejected: {event.get('message')}")
if event.get("type") == "command.ack":
return event.get("command_id") == command_id
raise RuntimeError("replacement owner did not acknowledge replayed command")
except (TimeoutError, OSError, websockets.exceptions.WebSocketException) as exc:
last_error = exc
if attempt + 1 < attempts:
await asyncio.sleep(1)
raise RuntimeError(f"station replay did not reconnect: {last_error}")
def ws_connect(ws_base: str, path: str, cookie: str):
return websockets.connect(
f"{ws_base}{path}", additional_headers={"Cookie": cookie},
ssl=ssl._create_unverified_context() if ws_base.startswith("wss://") else None,
open_timeout=8, ping_interval=None,
)
# Текст `_close_auth_state_unavailable` в backend/app/api/auth.py: 1013 с ним —
# закрытие от устаревшего кэша поколений, а не произвольный 1013.
AUTH_UNAVAILABLE_REASON = "Состояние доступа временно недоступно"
async def wait_auth_fresh(base: str, cookie: str, timeout: float = 10) -> bool:
"""Узел принял cookie: кэш поколений свежий и логин ему известен."""
request = urllib.request.Request(f"{base}/api/auth/me", headers={"Cookie": cookie})
deadline = time.monotonic() + timeout
while True:
try:
with urllib.request.urlopen(request, timeout=2) as response:
if response.status == 200:
return True
except OSError:
pass
if time.monotonic() > deadline:
return False
await asyncio.sleep(0.2)
async def wait_for_fence(socket, timeout: float = 10, ack_id: str | None = None) -> dict:
acked = False
def outcome(**closure) -> dict:
if ack_id is not None:
closure["command_acked"] = acked
return closure
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
try:
message = await asyncio.wait_for(socket.recv(), timeout=deadline - time.monotonic())
except TimeoutError:
return outcome(closed=False, fence_notice=False, reason="timeout")
except Exception as exc:
close = getattr(exc, "rcvd", None) or getattr(exc, "sent", None)
close_code = getattr(close, "code", None) or getattr(exc, "code", None)
return outcome(
closed=True,
fence_notice=close_code == 1012,
close_code=close_code,
reason=getattr(close, "reason", "") if close else str(exc),
)
if isinstance(message, str):
try:
event = json.loads(message)
except json.JSONDecodeError:
continue
if ack_id is not None and event.get("type") == "command.ack" \
and event.get("command_id") == ack_id:
acked = True
if event.get("message") == "Занятие передано другому backend-узлу; переподключитесь.":
return outcome(closed=True, fence_notice=True, close_code=1012,
reason=event.get("message"))
return outcome(closed=False, fence_notice=False, reason="timeout")
async def main(args: argparse.Namespace) -> int:
from sqlalchemy import delete, select
from sqlalchemy.engine import make_url
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
from sqlalchemy.pool import NullPool
from app.api.auth import hash_password, login_log_marker
from app.db.models import AuditLog, Session, Trainee, User
database_url = (
f"postgresql+asyncpg://lct:{args.postgres_password}@127.0.0.1:"
f"{args.postgres_port}/lct"
)
if make_url(database_url).host not in {"127.0.0.1", "localhost"}:
raise ValueError("failover smoke accepts loopback PostgreSQL only")
engine = create_async_engine(database_url)
factory = async_sessionmaker(engine, expire_on_commit=False)
observer_engine = create_async_engine(database_url, poolclass=NullPool)
run_id = uuid4().hex[:12]
session_id = args.session_id or uuid4()
trainee = Trainee(name=f"Failover smoke {run_id}")
instructor_login = f"fo-i-{run_id}"
trainee_login = f"fo-t-{run_id}"
instructor_password = secrets.token_urlsafe(24)
trainee_password = secrets.token_urlsafe(24)
owner_service = "backend"
owner_stopped = False
fault_service = None
fault_stopped = False
result: dict = {"session_id": str(session_id), "checks": {}}
backend = f"http://127.0.0.1:{args.backend_port}"
result["login_node"] = "backend-a"
ws_base = args.frontend_url.replace("https://", "wss://").replace("http://", "ws://")
instructor_cookie = trainee_cookie = None
try:
async with factory() as db:
db.add(trainee)
await db.flush()
db.add_all([
User(login=instructor_login, full_name="Failover smoke instructor",
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()
users_committed_at = time.monotonic()
instructor_cookie = login(backend, instructor_login, instructor_password)
trainee_cookie = login(backend, trainee_login, trainee_password)
# Учётки созданы прямо в БД, вход — на узле A, а сокет занятия Nginx
# может отдать узлу B до его сверки поколений. Ждать сверку не нужно:
# B проверяет неизвестный логин разовым SELECT (lct-42). Ожидание здесь
# лишило бы доказательности проверку ниже: опрос B сам вызвал бы разовую
# проверку, а за секунду B узнал бы логин сверкой.
control_sent_at = await send_control(ws_base, session_id, instructor_cookie, {
"type": "scenario.start", "scenario_id": args.scenario,
"trainee": trainee.name, "trainee_id": str(trainee.id),
"mode": "training", "exercise": "dds",
})
result["users_commit_to_control_seconds"] = round(control_sent_at - users_committed_at, 4)
if args.expect_initial_owner == "backend-b":
result["checks"]["control_sent_within_first_second"] = (
result["users_commit_to_control_seconds"] < 1
)
# Сверка B идёт раз в секунду в случайной фазе, и время само по себе
# не доказывает, что B не знал логин. Доказательство — запись B о
# разовой проверке именно этого логина.
result["checks"]["backend_b_resolved_login_one_shot"] = (
login_log_marker(instructor_login)
in compose(args, "logs", "--no-color", "backend-b")
)
async def session_row():
async with factory() as db:
return await db.scalar(select(Session).where(Session.id == session_id))
async def wait_for_command_checkpoint(command_id: str, timeout: float = 10) -> bool:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
# Use a new physical connection, avoiding a pooled observer's
# stale read transaction while the websocket writer commits.
async with observer_engine.connect() as db:
payload = (await db.execute(
select(Session.live_state).where(Session.id == session_id)
)).scalar_one_or_none()
if payload and command_id in payload.get("processed_station_commands", []):
return True
await asyncio.sleep(0.1)
return False
deadline = time.monotonic() + 12
row = None
while time.monotonic() < deadline:
row = await session_row()
if row and row.backend_node_id and row.checkpoint_at and row.live_state:
break
await asyncio.sleep(0.2)
if not row or not row.live_state:
raise RuntimeError("session did not persist an owner and recovery checkpoint")
if row.backend_node_id not in {"backend-a", "backend-b"}:
raise RuntimeError(f"unexpected session owner: {row.backend_node_id}")
owner_service = "backend" if row.backend_node_id == "backend-a" else "backend-b"
old_epoch = row.backend_fencing_epoch
result["initial_owner"] = row.backend_node_id
result["initial_epoch"] = old_epoch
if args.expect_initial_owner and row.backend_node_id != args.expect_initial_owner:
raise RuntimeError(
f"session routed to {row.backend_node_id}, expected {args.expect_initial_owner}"
)
if args.expect_initial_owner == "backend-b":
result["checks"]["login_a_initial_control_b"] = True
# Commit real DDS work before the fault. These actions must survive the
# checkpoint handoff and remain part of the final scored report.
snapshot = await station_command(ws_base, session_id, trainee_cookie)
service = snapshot.get("managed_service") or next(iter(snapshot["services"]), None)
if not service:
raise RuntimeError("the DDS exercise has no managed service")
crew = next(
(item for item in snapshot["crew_options"] if item.startswith(service + " — ")),
None,
)
if not crew:
raise RuntimeError(f"no crew option is available for {service}")
steps_before_fault = []
for status, detail in (
("accepted", "card accepted for processing"),
("crew.select", crew),
("responding", "crew reported departure"),
):
if status == "crew.select":
command = {"type": status, "crew": crew}
else:
command = {
"type": "card.status", "service": service, "status": status,
"comment": f"Основание: доклад по карточке.\nСведения: {detail}; {service} notified.",
}
snapshot = await station_command(ws_base, session_id, trainee_cookie, command)
steps_before_fault.append(status)
result["dds_progress_before_fault"] = steps_before_fault
result["dds_service"] = service
result["dds_crew"] = crew
# Commit a non-replace station command, but close the client socket
# without reading command.ack. The live database read proves the ID
# and business state are committed before the owner fails.
lost_ack_command_id = str(uuid4())
lost_ack_command = {
"type": "card.status", "service": service, "status": "arrived",
"comment": (
"Основание: доклад по карточке.\n"
f"Сведения: прибытие до отказа; {service} notified."
),
}
async with ws_connect(
ws_base, f"/ws/station/{session_id}", trainee_cookie,
) as lost_ack_socket:
for _ in range(60):
raw = await asyncio.wait_for(lost_ack_socket.recv(), timeout=10)
event = json.loads(raw) if isinstance(raw, str) else {}
if event.get("type") == "station.state":
break
if event.get("type") == "error":
raise RuntimeError(f"station rejected lost-ACK setup: {event.get('message')}")
else:
raise RuntimeError("station did not become ready for lost-ACK command")
await lost_ack_socket.send(json.dumps({
**lost_ack_command, "_command_id": lost_ack_command_id,
}, ensure_ascii=False))
# Observe protocol output in the harness, but deliberately do not
# dispatch the ACK to the simulated browser/outbox. This is the
# lost-ack boundary: the command is committed, client state stays
# pending, and the socket is closed before the application consumes
# command.ack.
deadline = time.monotonic() + 12
ack_seen = False
while time.monotonic() < deadline:
raw = await asyncio.wait_for(
lost_ack_socket.recv(), timeout=max(0.1, deadline - time.monotonic())
)
event = json.loads(raw) if isinstance(raw, str) else {}
if event.get("type") == "error":
raise RuntimeError(
f"station command failed before lost-ACK simulation: {event.get('message')}"
)
if event.get("type") == "command.ack":
ack_seen = event.get("command_id") == lost_ack_command_id
break
result["checks"]["lost_ack_command_ack_emitted"] = ack_seen
if not ack_seen:
raise RuntimeError("server did not emit command.ack for the DDS status command")
result["checks"]["lost_ack_command_committed"] = await wait_for_command_checkpoint(
lost_ack_command_id,
)
if not result["checks"]["lost_ack_command_committed"]:
raise RuntimeError("station command ID did not reach the PostgreSQL checkpoint")
# Do not consume any frame after send: the browser's pending
# command remains unacknowledged when this transport is closed.
await lost_ack_socket.close()
result["lost_ack_command_id"] = lost_ack_command_id
result["lost_ack_ack_dispatched_to_browser"] = False
if args.failure_mode == "kill":
fault_service = owner_service
compose(args, "kill", fault_service)
owner_stopped = True
else:
proxy_name = "db-proxy-a" if row.backend_node_id == "backend-a" else "db-proxy-b"
fault_service = proxy_name
async with AsyncExitStack() as channels:
active_sockets = {}
for channel, path, cookie in (
("control", f"/ws/control/{session_id}", instructor_cookie),
("call", f"/ws/call/{session_id}", trainee_cookie),
("observe", f"/ws/observe/{session_id}", instructor_cookie),
("station", f"/ws/station/{session_id}", trainee_cookie),
):
active_sockets[channel] = await channels.enter_async_context(
ws_connect(ws_base, path, cookie)
)
compose(args, "stop", proxy_name)
fault_stopped = True
# Реальная мутация старому владельцу уже без БД: commit не пройдёт,
# ACK не придёт, а id не должен попасть в checkpoint нового
# владельца — прямое доказательство запрета записи после потери БД.
# Отправка сразу, до `compose ps`: через 2 с auth закроет сокет.
boundary_command_id = str(uuid4())
boundary = {"command_id": boundary_command_id, "status": "working"}
try:
await active_sockets["station"].send(json.dumps({
"type": "card.status", "service": service, "status": "working",
"comment": (
"Основание: доклад по карточке.\n"
f"Сведения: начало работ после обрыва БД; {service} notified."
),
"_command_id": boundary_command_id,
}, ensure_ascii=False))
boundary["sent"] = True
except websockets.exceptions.WebSocketException as exc:
boundary["sent"] = False
boundary["error"] = f"{type(exc).__name__}: {exc}"
result["boundary_command_sent_to_old_owner"] = boundary
result["checks"]["boundary_command_sent_after_partition"] = boundary["sent"]
running_services = compose(args, "ps", "--status", "running", "--services").splitlines()
result["checks"]["owner_process_running"] = owner_service in running_services
if owner_service not in running_services:
raise RuntimeError("owner process did not remain running during DB partition")
closures = await asyncio.gather(*(
wait_for_fence(
socket, timeout=args.takeover_timeout,
ack_id=boundary_command_id if channel == "station" else None,
)
for channel, socket in active_sockets.items()
))
result["existing_channel_closures"] = dict(zip(active_sockets, closures))
result["existing_channel_fence_notices"] = {
channel: closure["fence_notice"]
for channel, closure in result["existing_channel_closures"].items()
}
result["checks"].update({
f"existing_{channel}_closed": closure["closed"]
for channel, closure in result["existing_channel_closures"].items()
})
# Без БД старый владелец закрывает каналы fail-closed одним из двух
# путей: auth не может сверить поколения учёток дольше 2 с — 1013,
# supervisor lease замечает потерю владения (цикл 5 с) — 1012.
# Auth обычно успевает первым; оба кода корректны, фронт их не
# различает. Фактический код остаётся в existing_channel_closures.
result["existing_channel_close_codes"] = {
channel: closure.get("close_code")
for channel, closure in result["existing_channel_closures"].items()
}
result["checks"].update({
f"existing_{channel}_closed_fail_closed": (
closure["fence_notice"] or (
closure.get("close_code") == 1013
and closure.get("reason") == AUTH_UNAVAILABLE_REASON
)
)
for channel, closure in result["existing_channel_closures"].items()
})
result["checks"]["boundary_command_not_acked"] = (
result["existing_channel_closures"]["station"].get("command_acked") is False
)
deadline = time.monotonic() + args.takeover_timeout
takeover = None
while time.monotonic() < deadline:
candidate = await session_row()
if candidate and candidate.backend_node_id != row.backend_node_id:
takeover = candidate
break
await asyncio.sleep(0.5)
if takeover is None:
raise RuntimeError("other backend did not take over the expired lease")
result["takeover_owner"] = takeover.backend_node_id
result["takeover_epoch"] = takeover.backend_fencing_epoch
result["takeover_seconds"] = round(args.takeover_timeout - max(0, deadline - time.monotonic()), 2)
result["checks"]["owner_changed"] = takeover.backend_node_id != row.backend_node_id
result["checks"]["fencing_epoch_incremented"] = takeover.backend_fencing_epoch > old_epoch
result["checks"]["checkpoint_restored"] = bool(takeover.live_state and takeover.checkpoint_at)
result["checks"]["lost_ack_command_replayed"] = await replay_station_command(
ws_base, session_id, trainee_cookie, lost_ack_command,
lost_ack_command_id, attempts=args.reconnect_attempts,
)
if args.failure_mode == "partition":
compose(args, "start", fault_service)
fault_stopped = False
await asyncio.sleep(6) # let the surviving old process observe the newer epoch
old_backend_port = args.backend_port if owner_service == "backend" else args.backend_b_port
old_backend = f"http://127.0.0.1:{old_backend_port}"
try:
with urllib.request.urlopen(f"{old_backend}/api/health", timeout=5) as response:
result["checks"]["old_owner_process_healthy"] = response.status == 200
except Exception as exc:
result["checks"]["old_owner_process_healthy"] = False
result["old_owner_health_error"] = f"{type(exc).__name__}: {exc}"
# Auth на старом узле снова свежий: 1013 до accept исключён, и 403
# на рукопожатии ниже остаётся только за fencing занятия (1012).
result["checks"]["old_owner_auth_fresh"] = await wait_auth_fresh(
old_backend, instructor_cookie,
)
stale = await probe_ws(
old_backend.replace("http://", "ws://"),
f"/ws/control/{session_id}", instructor_cookie,
)
result["stale_owner_control"] = stale
result["checks"]["stale_owner_control_rejected"] = (
not stale["connected"] and "HTTP 403" in stale.get("error", "")
)
# БД вернулась к старому владельцу, но его команда после обрыва так
# и не должна появиться в checkpoint: fencing epoch отверг запись.
result["checks"]["boundary_command_absent_from_checkpoint"] = (
not await wait_for_command_checkpoint(boundary_command_id, timeout=1)
)
request = urllib.request.Request(
f"{args.frontend_url}/api/sessions/{session_id}",
headers={"Cookie": instructor_cookie},
)
rest_started = time.monotonic()
rest_deadline = rest_started + args.rest_timeout
rest_attempt = 0
while time.monotonic() < rest_deadline:
rest_attempt += 1
try:
with urllib.request.urlopen(
request,
timeout=min(3.0, max(0.2, rest_deadline - time.monotonic())),
context=ssl._create_unverified_context(),
) as response:
result["checks"]["session_rest"] = response.status == 200
if result["checks"]["session_rest"]:
result.pop("rest_error", None)
break
except urllib.error.HTTPError as exc:
result["rest_error"] = f"HTTP {exc.code}"
except Exception as exc:
result["rest_error"] = f"{type(exc).__name__}: {exc}"
await asyncio.sleep(min(1, max(0, rest_deadline - time.monotonic())))
result["checks"].setdefault("session_rest", False)
result["rest_attempts"] = rest_attempt
result["rest_elapsed_seconds"] = round(time.monotonic() - rest_started, 2)
async def reconnect_probe(path: str, cookie: str) -> dict:
last = {}
for attempt in range(1, args.reconnect_attempts + 1):
last = await probe_ws(ws_base, path, cookie)
if last["connected"]:
last["attempts"] = attempt
return last
await asyncio.sleep(1)
last["attempts"] = args.reconnect_attempts
return last
probes = {
"control": await reconnect_probe(f"/ws/control/{session_id}", instructor_cookie),
"call": await reconnect_probe(f"/ws/call/{session_id}", trainee_cookie),
"observe": await reconnect_probe(f"/ws/observe/{session_id}", instructor_cookie),
"station": await reconnect_probe(f"/ws/station/{session_id}", trainee_cookie),
}
result["websockets"] = probes
result["checks"].update({f"ws_{name}": probe["connected"] for name, probe in probes.items()})
# Reconcile the race status against the recovered checkpoint. If it
# committed before the fault, do not repeat it; otherwise issue it now.
recovered_snapshot = await station_command(ws_base, session_id, trainee_cookie)
recovered_status = recovered_snapshot["statuses"].get(service)
result["boundary_status_after_takeover"] = recovered_status
result["boundary_command_outcome"] = (
"committed before takeover" if recovered_status == "arrived"
else "not in recovered checkpoint; reconciled after takeover"
)
result["checks"]["boundary_command_reconciled"] = recovered_status in {
"responding", "arrived",
}
if args.failure_mode == "partition":
# `working` после обрыва не применён: статус остался на lost-ack `arrived`.
result["checks"]["boundary_command_not_applied"] = recovered_status == "arrived"
result["dds_progress_after_takeover"] = []
if recovered_status == "responding":
recovered_snapshot = await station_command(ws_base, session_id, trainee_cookie, {
"type": "card.status", "service": service, "status": "arrived",
"comment": "Основание: доклад по карточке.\n"
f"Сведения: бригада сообщила о прибытии; {service} notified.",
})
result["dds_progress_after_takeover"].append("arrived")
# Continue the same card after owner recovery and require a stored final
# score, not merely successful handshakes on the replacement owner.
for status, detail in (
("working", "crew started work"),
("completed", "crew completed work"),
):
snapshot = await station_command(ws_base, session_id, trainee_cookie, {
"type": "card.status", "service": service, "status": status,
"comment": f"Основание: доклад по карточке.\nСведения: {detail}; {service} notified.",
})
result["dds_progress_after_takeover"].append(status)
async with websockets.connect(
f"{ws_base}/ws/station/{session_id}",
additional_headers={"Cookie": trainee_cookie},
ssl=ssl._create_unverified_context() if ws_base.startswith("wss://") else None,
open_timeout=10, ping_interval=None,
) as socket:
await socket.send(json.dumps({"type": "station.finish"}))
deadline = time.monotonic() + 15
while time.monotonic() < deadline:
raw = await asyncio.wait_for(socket.recv(), timeout=10)
event = json.loads(raw) if isinstance(raw, str) else {}
if event.get("type") == "error":
raise RuntimeError(f"station finish rejected: {event.get('message')}")
if event.get("type") == "score.ready":
result["checks"]["dds_finished"] = True
break
else:
raise RuntimeError("DDS exercise did not finish after takeover")
report_request = urllib.request.Request(
f"{args.frontend_url}/api/sessions/{session_id}/report",
headers={"Cookie": instructor_cookie},
)
with urllib.request.urlopen(
report_request, timeout=10, context=ssl._create_unverified_context(),
) as response:
report = json.loads(response.read())
result["dds_final_score"] = report.get("score_auto")
card_results = report.get("card_results", [])
if len(card_results) != 1:
raise RuntimeError(f"expected one DDS card result, got {len(card_results)}")
final_card = card_results[0]
result["dds_card_metrics"] = {
metric["key"]: metric["passed"] for metric in final_card["metrics"]
}
required_metric_keys = {
"dds_ack", "dds_decision", "dds_crew", "dds_progress",
"dds_completion", "dds_reply",
}
result["checks"]["dds_business_metrics_passed"] = all(
result["dds_card_metrics"].get(key) is True for key in required_metric_keys
)
from collections import Counter
observed_status_counts = Counter(
action.get("status") for action in final_card["actions"]
if action.get("type") == "card.status" and action.get("service") == service
)
expected_statuses = {"accepted", "responding", "arrived", "working", "completed"}
result["checks"]["dds_statuses_persisted_exactly_once"] = all(
observed_status_counts[status] == 1 for status in expected_statuses
)
result["dds_status_counts"] = dict(observed_status_counts)
result["checks"]["lost_ack_status_recorded_once"] = (
observed_status_counts["arrived"] == 1
)
passed = all(result["checks"].values())
result["pass"] = passed
print(json.dumps(result, ensure_ascii=False, indent=2))
return 0 if passed else 1
finally:
if instructor_cookie:
try:
await send_control(ws_base, session_id, instructor_cookie, {"type": "session.stop"})
except Exception:
pass
if fault_stopped and fault_service:
compose(args, "start", fault_service, check=False)
if owner_stopped:
compose(args, "start", owner_service, check=False)
async with factory() as db:
row = await db.scalar(select(Session).where(Session.id == session_id))
if row:
await db.execute(delete(Session).where(Session.id == session_id))
await db.execute(delete(AuditLog).where(AuditLog.object_id == str(session_id)))
await db.execute(delete(AuditLog).where(AuditLog.actor.in_([instructor_login, trainee_login])))
await db.execute(delete(User).where(User.login.in_([instructor_login, trainee_login])))
await db.execute(delete(Trainee).where(Trainee.name == trainee.name))
await db.commit()
await engine.dispose()
await observer_engine.dispose()
if __name__ == "__main__":
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--project", required=True)
parser.add_argument("--postgres-port", type=int, required=True)
parser.add_argument("--backend-port", type=int, required=True)
parser.add_argument("--backend-b-port", type=int, default=18001)
parser.add_argument("--frontend-port", type=int, required=True)
parser.add_argument("--tls-port", type=int, required=True)
parser.add_argument("--postgres-password", required=True)
parser.add_argument("--session-secret", required=True)
parser.add_argument("--tls-cert-dir", required=True)
parser.add_argument("--frontend-url", required=True)
parser.add_argument("--scenario", default="fire-apartment-l2")
parser.add_argument("--takeover-timeout", type=float, default=30)
parser.add_argument("--reconnect-attempts", type=int, default=5)
parser.add_argument("--rest-timeout", type=float, default=60)
parser.add_argument("--failure-mode", choices=("kill", "partition"), default="kill")
parser.add_argument("--expect-initial-owner", choices=("backend-a", "backend-b"))
parser.add_argument("--session-id", type=UUID)
raise SystemExit(asyncio.run(main(parser.parse_args())))

View file

@ -1,402 +0,0 @@
#!/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)))

View file

@ -2,18 +2,29 @@
set -euo pipefail
repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
project="lct-test-$$-$RANDOM$RANDOM"
compose=(docker compose -p "$project" -f "$repo_root/docker-compose.test-db.yml")
container="lct-test-$$-$RANDOM$RANDOM"
cleanup() {
trap - EXIT INT TERM
"${compose[@]}" down --volumes --remove-orphans >/dev/null || true
docker rm --force --volumes "$container" >/dev/null 2>&1 || true
}
trap cleanup EXIT
printf 'Изолированный тестовый Compose project: %s\n' "$project"
"${compose[@]}" up --pull never --detach --wait
port_mapping=$("${compose[@]}" port postgres 5432)
# Временная БД без compose-файла и без тома: dev-база не затрагивается,
# порт на loopback выбирает Docker.
printf 'Изолированная тестовая PostgreSQL: %s\n' "$container"
docker run --detach --rm --name "$container" \
--env POSTGRES_USER=lct_test --env POSTGRES_PASSWORD=lct_test --env POSTGRES_DB=lct_test \
--publish 127.0.0.1::5432 \
--health-cmd 'pg_isready -U lct_test -d lct_test' --health-interval 2s --health-retries 30 \
postgres:16-alpine >/dev/null
for _ in $(seq 60); do
status=$(docker inspect --format '{{.State.Health.Status}}' "$container")
[[ $status == healthy ]] && break
sleep 1
done
[[ $status == healthy ]] || { echo "PostgreSQL не поднялась: $status" >&2; exit 1; }
port_mapping=$(docker port "$container" 5432/tcp | head -n1)
db_port="${port_mapping##*:}"
database_url="postgresql+asyncpg://lct_test:lct_test@127.0.0.1:${db_port}/lct_test"
evidence_dir="${LCT_TEST_EVIDENCE_DIR:-$repo_root/.local/evidence}"

View file

@ -1,35 +0,0 @@
#!/usr/bin/env bash
set -euo pipefail
repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
project="lct-production-check-$$"
port=$((20000 + $$ % 40000))
secret="prod-check-$(od -An -N16 -tx1 /dev/urandom | tr -d ' \n')"
compose=(docker compose -p "$project" -f "$repo_root/docker-compose.yml" \
-f "$repo_root/docker-compose.production.yml")
cleanup() {
POSTGRES_PASSWORD="$secret" POSTGRES_PORT="$port" \
"${compose[@]}" down --volumes --remove-orphans >/dev/null 2>&1 || true
}
trap cleanup EXIT INT TERM
if env -u POSTGRES_PASSWORD "${compose[@]}" config --quiet >/dev/null 2>&1; then
echo "ОШИБКА: production Compose запустился без POSTGRES_PASSWORD" >&2
exit 1
fi
POSTGRES_PASSWORD="$secret" POSTGRES_PORT="$port" \
"${compose[@]}" up --pull never --detach --wait postgres
network="${project}_default"
docker run --pull never --rm --network "$network" -e "PGPASSWORD=$secret" \
postgres:16-alpine psql -h postgres -U lct -d lct -Atc 'select 1' | rg -x '1' >/dev/null
if docker run --pull never --rm --network "$network" -e PGPASSWORD=wrong \
postgres:16-alpine psql -h postgres -U lct -d lct -Atc 'select 1' >/dev/null 2>&1; then
echo "ОШИБКА: production PostgreSQL принял неверный пароль" >&2
exit 1
fi
echo "Production PostgreSQL принял пароль из backend-сети и отклонил неверный."

View file

@ -100,7 +100,7 @@ async def compose_service(project: str, action: str, service: str) -> None:
async def run(args: argparse.Namespace) -> int:
frontend = local_url(args.frontend_url, {"https"})
frontend = local_url(args.frontend_url, {"http", "https"})
backend = local_url(args.backend_url, {"http"})
if not args.compose_project.strip():
raise ValueError("--compose-project обязателен: recovery test останавливает Nginx")
@ -374,7 +374,7 @@ async def run(args: argparse.Namespace) -> int:
if __name__ == "__main__":
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--frontend-url", default="https://127.0.0.1:5443")
parser.add_argument("--frontend-url", default="http://127.0.0.1:5173")
parser.add_argument("--backend-url", default="http://127.0.0.1:8000")
parser.add_argument("--compose-project", required=True,
help="точное имя изолированного Compose project; тест останавливает его frontend")

View file

@ -1,48 +0,0 @@
#!/usr/bin/env bash
set -euo pipefail
repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
project="lct-sip-smoke-$(date +%Y%m%d)-$$-${RANDOM:-0}"
compose=(docker compose -p "$project" -f "$repo_root/docker-compose.yml" -f "$repo_root/docker-compose.sip.yml")
sip_port="${SIP_SMOKE_PORT:-15170}"
sips_port="${SIP_SMOKE_SIPS_PORT:-15171}"
websocket_port="${SIP_SMOKE_WS_PORT:-18098}"
rtp_start="${SIP_SMOKE_RTP_START:-12000}"
rtp_end=$((rtp_start + 99))
rtp_offset=$((rtp_start - 10000))
tls_dir="$(mktemp -d /tmp/lct-sip-smoke-tls-XXXXXX)"
sip_smoke_password="$(openssl rand -hex 24)"
cleanup() {
trap - EXIT INT TERM
"${compose[@]}" down --volumes --remove-orphans >/dev/null || true
case "$tls_dir" in
/tmp/lct-sip-smoke-tls-*) rm -rf -- "$tls_dir" ;;
*) printf 'Refusing to remove unexpected TLS temp path: %s\n' "$tls_dir" >&2 ;;
esac
}
trap cleanup EXIT
trap 'exit 130' INT
trap 'exit 143' TERM
python3 -c 'import socket,sys; ports=[("tcp",int(sys.argv[1])),("udp",int(sys.argv[1])),("tcp",int(sys.argv[2])),("tcp",int(sys.argv[3]))]+[("udp",p) for p in range(int(sys.argv[4]),int(sys.argv[5])+1)]; sockets=[]
try:
for kind,port in ports:
sock=socket.socket(socket.AF_INET,socket.SOCK_STREAM if kind=="tcp" else socket.SOCK_DGRAM)
sock.bind(("127.0.0.1",port)); sockets.append(sock)
except OSError as exc:
raise SystemExit(f"SIP smoke port {port}/{kind} is unavailable: {exc}")
finally:
[sock.close() for sock in sockets]' "$sip_port" "$sips_port" "$websocket_port" "$rtp_start" "$rtp_end"
printf 'Starting isolated SIP smoke project: %s\n' "$project"
SIP_6001_PASSWORD="$sip_smoke_password" \
SIP_PORT="$sip_port" SIPS_PORT="$sips_port" SIP_WS_PORT="$websocket_port" \
RTP_PORT_START="$rtp_start" RTP_PORT_END="$rtp_end" \
SIP_TLS_DIR="$tls_dir" SIP_EXTERNAL_MEDIA_ADDRESS=127.0.0.1 \
"${compose[@]}" up --pull never --no-build --detach --wait --no-deps sip
SIP_PASSWORD="$sip_smoke_password" python3 "$repo_root/scripts/smoke_sip.py" \
--host 127.0.0.1 --port "$sip_port" --rtp-host-offset "$rtp_offset" \
--output "$repo_root/docs/evidence/sip-rtp-2026-09-25.json"

View file

@ -1,83 +0,0 @@
#!/usr/bin/env bash
set -euo pipefail
repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
project="lct-webrtc-smoke-$$-${RANDOM:-0}"
compose=(docker compose -p "$project" \
-f "$repo_root/docker-compose.yml" \
-f "$repo_root/docker-compose.tls.yml" \
-f "$repo_root/docker-compose.sip.yml" \
-f "$repo_root/docker-compose.webrtc-test.yml")
temp_root="$(mktemp -d /tmp/lct-webrtc-smoke-XXXXXX)"
# Keep all ports in one checked block; RTP gets its own 100-port block.
base_port=$((30000 + ($$ % 15000)))
export COMPOSE_PROJECT_NAME="$project"
export BIND_HOST=127.0.0.1
export POSTGRES_PORT=$base_port
export BACKEND_PORT=$((base_port + 1))
export FRONTEND_PORT=$((base_port + 2))
export TLS_PORT=$((base_port + 3))
export SIP_PORT=$((base_port + 4))
export SIPS_PORT=$((base_port + 5))
export SIP_WS_PORT=$((base_port + 6))
export RTP_PORT_START=$((base_port + 100))
export RTP_PORT_END=$((base_port + 199))
export SIP_RTP_START="$RTP_PORT_START"
export SIP_RTP_END="$RTP_PORT_END"
export POSTGRES_PASSWORD="$(openssl rand -hex 24)"
export SESSION_SECRET="$(openssl rand -hex 48)"
export TLS_CERT_DIR="$temp_root/frontend-tls"
export SIP_TLS_DIR="$temp_root/sip-tls"
export SIP_EXTERNAL_MEDIA_ADDRESS=127.0.0.1
mkdir -p "$TLS_CERT_DIR" "$SIP_TLS_DIR"
cleanup() {
status=$?
trap - EXIT INT TERM
if ((status != 0)); then
failure_log="/tmp/${project}-compose.log"
"${compose[@]}" logs --no-color --tail 250 >"$failure_log" 2>&1 || true
printf 'Isolated WebRTC diagnostics saved to %s\n' "$failure_log" >&2
fi
"${compose[@]}" down --volumes --remove-orphans >/dev/null || true
case "$temp_root" in
/tmp/lct-webrtc-smoke-*) rm -rf -- "$temp_root" ;;
*) printf 'Refusing to remove unexpected temporary path: %s\n' "$temp_root" >&2 ;;
esac
exit "$status"
}
trap cleanup EXIT
trap 'exit 130' INT
trap 'exit 143' TERM
python3 -c 'import socket,sys
ports=[("tcp",int(sys.argv[1])),("tcp",int(sys.argv[2])),("tcp",int(sys.argv[3])),
("tcp",int(sys.argv[4])),("tcp",int(sys.argv[5])),("udp",int(sys.argv[5])),
("tcp",int(sys.argv[6])),("tcp",int(sys.argv[7]))]
ports.extend(("udp",p) for p in range(int(sys.argv[8]),int(sys.argv[9])+1))
sockets=[]
try:
for kind,port in ports:
sock=socket.socket(socket.AF_INET,socket.SOCK_STREAM if kind=="tcp" else socket.SOCK_DGRAM)
sock.bind(("127.0.0.1",port)); sockets.append(sock)
except OSError as exc:
raise SystemExit(f"isolated WebRTC smoke port {port}/{kind} unavailable: {exc}")
finally:
[sock.close() for sock in sockets]' \
"$POSTGRES_PORT" "$BACKEND_PORT" "$FRONTEND_PORT" "$TLS_PORT" "$SIP_PORT" \
"$SIPS_PORT" "$SIP_WS_PORT" "$RTP_PORT_START" "$RTP_PORT_END"
printf 'Starting isolated WebRTC smoke project: %s\n' "$project"
cd "$repo_root"
npm --prefix frontend run build
"${compose[@]}" build backend sip
"${compose[@]}" up --no-build --pull never --detach --wait --wait-timeout 300 \
postgres backend frontend sip
python3 "$repo_root/scripts/smoke_webrtc_browser.py" \
--frontend-url "https://127.0.0.1:$TLS_PORT" \
--backend-url "http://127.0.0.1:$BACKEND_PORT" \
--database-url "postgresql+asyncpg://lct:${POSTGRES_PASSWORD}@127.0.0.1:${POSTGRES_PORT}/lct" \
--timeout "${WEBRTC_TEST_TIMEOUT:-60}" \
--output "$repo_root/docs/evidence/webrtc-browser-isolated-2026-09-26.json"

View file

@ -1,25 +0,0 @@
#!/usr/bin/env bash
set -euo pipefail
repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
if [[ -n "${POSTGRES_PASSWORD+x}" ]]; then
password="$POSTGRES_PASSWORD"
elif [[ -f "$repo_root/.env" ]]; then
matches="$(grep -c '^POSTGRES_PASSWORD=' "$repo_root/.env" || true)"
if [[ "$matches" != "1" ]]; then
echo "Укажите ровно один POSTGRES_PASSWORD в .env или окружении." >&2
exit 2
fi
password="$(sed -n 's/^POSTGRES_PASSWORD=//p' "$repo_root/.env")"
else
echo "Для production запуска скопируйте .env.example в .env и задайте POSTGRES_PASSWORD." >&2
exit 2
fi
if [[ ! "$password" =~ ^[A-Za-z0-9._~-]{32,}$ ]]; then
echo "POSTGRES_PASSWORD должен содержать не менее 32 символов: латиница, цифры, точка, дефис, подчёркивание или тильда. Значение не показано." >&2
exit 2
fi
echo "Production DB secret задан и удовлетворяет формату (само значение не отображается)."