diff --git a/backend/app/api/ws/__init__.py b/backend/app/api/ws/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/backend/app/api/ws/call.py b/backend/app/api/ws/call.py new file mode 100644 index 0000000..94e4699 --- /dev/null +++ b/backend/app/api/ws/call.py @@ -0,0 +1,154 @@ +"""Сокет курсанта-оператора 112. + +Здесь только JSON-часть: приём вызова, правки карточки, подсказки, передача +в ДДС, завершение. Бинарные аудиокадры и голосовой контур — карточка lct-06. +""" + +import asyncio +import logging +from uuid import UUID + +from fastapi import APIRouter, WebSocket, WebSocketDisconnect +from pydantic import TypeAdapter, ValidationError + +from app.domain.events import ( + CallEnded, + CallEndReason, + CallStarted, + ErrorEvent, + ErrorKind, + HintShown, + KioState, + SessionEnded, + SessionMode, + TimerTick, + TraineeToServer, +) +from app.scenarios import store +from app.session.hub import hub +from app.session.state import now_utc + +log = logging.getLogger(__name__) +router = APIRouter() + +_adapter = TypeAdapter(TraineeToServer) + + +def _next_hint(state) -> tuple[str, str] | None: + """Следующий пункт чек-листа, который ещё не подсказывали. + + Порядок временный: пока нет слот-автомата (lct-07), подсказка идёт + по порядку чек-листа, а не по реально неотработанным пунктам. + Перевод на слот-автомат — карточка lct-14. + """ + scenario = store.get(state.scenario_id) + if scenario is None: + return None + for item in scenario.checklist: + if item.id not in state.hints_shown and item.question: + return item.id, item.question + return None + + +async def _handle(session_id: UUID, state, event) -> None: + match event.type: + case "call.answer": + state.on_event("call.answer") + state.started_at = now_utc() + hub.to_trainee(session_id, CallStarted(started_at=state.started_at)) + hub.to_observers(session_id, state.snapshot()) + if hub.journal: + await hub.journal.session_started(session_id, state.started_at) + + case "kio.patch": + state.patch_kio(event.fields) + # Наблюдателю уходит карточка целиком: рассинхрон на внешнем мониторе + # посреди занятия дороже лишних килобайт. + hub.to_observers(session_id, KioState(kio=state.kio)) + + case "hint.request": + if state.mode is SessionMode.EXAM: + # На экзамене опоры нет — это часть нормы контроля. + hub.to_trainee( + session_id, + ErrorEvent( + code=ErrorKind.HINT_DENIED_IN_EXAM, + message="В контрольном режиме подсказки недоступны", + ), + ) + return + nxt = _next_hint(state) + if nxt is None: + return + checklist_id, question = nxt + state.hints_shown.append(checklist_id) + shown = HintShown(checklist_id=checklist_id, question=question) + hub.broadcast(session_id, shown) + if hub.journal: + await hub.journal.hint(session_id, checklist_id, question, now_utc()) + + case "dds.dispatch": + state.on_event("dds.dispatch") + state.dispatch(event.service.value) + hub.to_observers(session_id, KioState(kio=state.kio)) + hub.to_observers(session_id, TimerTick(timers=state.timers.snapshot())) + + case "callback.dial": + state.on_event("callback.dial") + + case "self_assessment.submit": + if hub.journal: + await hub.journal.self_assessment( + session_id, event.missed, event.comment, now_utc() + ) + + case "call.hangup": + state.ended_at = now_utc() + state.end_reason = CallEndReason.HANGUP + hub.stop_ticker(session_id) + hub.to_trainee(session_id, CallEnded(reason=CallEndReason.HANGUP)) + hub.to_observers(session_id, SessionEnded(reason=CallEndReason.HANGUP)) + if hub.journal: + await hub.journal.session_ended( + session_id, state.ended_at, CallEndReason.HANGUP.value + ) + + +async def _pump(ws: WebSocket, queue: asyncio.Queue) -> None: + while True: + event = await queue.get() + await ws.send_text(event.model_dump_json()) + + +@router.websocket("/ws/call/{session_id}") +async def call(ws: WebSocket, session_id: UUID) -> None: + await ws.accept() + + state = hub.get(session_id) + if state is None: + await ws.send_text( + ErrorEvent( + code=ErrorKind.SESSION_NOT_FOUND, message="Занятие ещё не запущено преподавателем" + ).model_dump_json() + ) + await ws.close() + return + + with hub.trainee(session_id) as queue: + writer = asyncio.create_task(_pump(ws, queue)) + try: + while True: + payload = await ws.receive_json() + try: + event = _adapter.validate_python(payload) + except ValidationError: + hub.to_trainee( + session_id, + ErrorEvent(code=ErrorKind.UNSUPPORTED_EVENT, message=str(payload)[:200]), + ) + continue + await _handle(session_id, state, event) + except WebSocketDisconnect: + return + finally: + writer.cancel() diff --git a/backend/app/api/ws/control.py b/backend/app/api/ws/control.py new file mode 100644 index 0000000..4379e04 --- /dev/null +++ b/backend/app/api/ws/control.py @@ -0,0 +1,146 @@ +"""Канал преподавателя: только передача. + +**Ни одной команды, меняющей карточку курсанта.** Преподаватель управляет +ситуацией, а не работой обучаемого, иначе оценка перестаёт быть оценкой +курсанта (docs/arch/CONTRACT.md). + +Ответы сюда не идут — канал односторонний. Всё, что сервер хочет сказать +преподавателю, уходит на его же сокет `observe`. +""" + +import logging +from uuid import UUID, uuid4 + +from fastapi import APIRouter, WebSocket, WebSocketDisconnect +from pydantic import TypeAdapter, ValidationError + +from app.domain.events import ( + CallEndReason, + CallIncoming, + ErrorEvent, + ErrorKind, + InstructorToServer, + InstructorNoteShown, + ModeSet, + ReferenceStarted, + SessionEnded, +) +from app.scenarios import store +from app.session.hub import hub +from app.session.state import SessionState, now_utc + +log = logging.getLogger(__name__) +router = APIRouter() + +_adapter = TypeAdapter(InstructorToServer) + + +async def _start(session_id: UUID, event) -> None: + scenario = store.get(event.scenario_id) + if scenario is None: + hub.to_observers( + session_id, + ErrorEvent(code=ErrorKind.SCENARIO_INVALID, message=f"Нет сценария {event.scenario_id}"), + ) + return + + attempt = 1 + if hub.journal: + attempt = await hub.journal.start_lesson( + session_id, scenario.id, event.mode.value, event.trainee + ) + + state = hub.register( + SessionState( + session_id=session_id, + scenario_id=scenario.id, + scenario_title=scenario.title, + level=scenario.level.value, + mode=event.mode, + required_fields=scenario.required_fields, + trainee_name=event.trainee, + attempt=attempt, + ) + ) + state.on_event("call.incoming") + hub.start_ticker(session_id) + + hub.to_trainee( + session_id, + CallIncoming( + scenario_id=scenario.id, + caller_number="+7 (495) 000-00-00", + level=scenario.level, + mode=event.mode, + required_fields=scenario.required_fields, + ), + ) + hub.to_observers(session_id, ModeSet(mode=event.mode)) + hub.to_observers(session_id, state.snapshot()) + + +async def _stop(session_id: UUID) -> None: + state = hub.get(session_id) + if state is None: + return + state.ended_at = now_utc() + state.end_reason = CallEndReason.INSTRUCTOR + hub.stop_ticker(session_id) + hub.to_observers(session_id, SessionEnded(reason=CallEndReason.INSTRUCTOR)) + if hub.journal: + await hub.journal.session_ended(session_id, state.ended_at, CallEndReason.INSTRUCTOR.value) + + +@router.websocket("/ws/control/{session_id}") +async def control(ws: WebSocket, session_id: UUID) -> None: + await ws.accept() + try: + while True: + payload = await ws.receive_json() + try: + event = _adapter.validate_python(payload) + except ValidationError: + hub.to_observers( + session_id, + ErrorEvent(code=ErrorKind.UNSUPPORTED_EVENT, message=str(payload)[:200]), + ) + continue + + match event.type: + case "scenario.start": + await _start(session_id, event) + case "session.stop": + await _stop(session_id) + case "instructor_note.add": + hub.to_observers( + session_id, + InstructorNoteShown( + transcript_ref=event.transcript_ref, + text=event.text, + author="преподаватель", + ), + ) + if hub.journal: + await hub.journal.note( + session_id, event.transcript_ref, event.text, "преподаватель" + ) + case "reference.play": + state = hub.get(session_id) + if state is not None: + hub.to_observers(session_id, ReferenceStarted(scenario_id=state.scenario_id)) + case "director.inject": + # Поведение звонящего — карточка lct-07, пульт — lct-22. + # До них директива копится в состоянии и видна в разборе. + state = hub.get(session_id) + if state is not None: + state.directives.append(event.directive) + case _: + hub.to_observers( + session_id, + ErrorEvent( + code=ErrorKind.UNSUPPORTED_EVENT, + message=f"{event.type} ещё не реализовано", + ), + ) + except WebSocketDisconnect: + return diff --git a/backend/app/api/ws/observe.py b/backend/app/api/ws/observe.py new file mode 100644 index 0000000..5a85ffb --- /dev/null +++ b/backend/app/api/ws/observe.py @@ -0,0 +1,62 @@ +"""Канал наблюдателя: внешний монитор и пульт преподавателя. + +**На этом сокете нет ни одного обработчика входящих сообщений.** Внешний +монитор физически не может повлиять на занятие: у него нет ни аудио, ни канала +записи. Преподаватель смотрит здесь, а пишет через отдельный `control`. + +Входящие кадры читаются и выбрасываются, не разбираясь: чтение нужно ровно +затем, чтобы заметить разрыв соединения. Без него задача сокета висела бы +на очереди до первой неудачной отправки, а закрытая вкладка монитора оставляла +бы за собой подписку. Разбора, диспетчеризации и эффекта у входящих нет. +""" + +import asyncio +from uuid import UUID + +from fastapi import APIRouter, WebSocket, WebSocketDisconnect + +from app.domain.events import ErrorEvent, ErrorKind +from app.session.hub import hub + +router = APIRouter() + + +async def _pump(ws: WebSocket, queue: asyncio.Queue) -> None: + while True: + event = await queue.get() + await ws.send_text(event.model_dump_json()) + + +async def _wait_for_disconnect(ws: WebSocket) -> None: + """Единственное назначение — дождаться разрыва. Содержимое кадров + не читается и никуда не передаётся.""" + while True: + message = await ws.receive() + if message["type"] == "websocket.disconnect": + return + + +@router.websocket("/ws/observe/{session_id}") +async def observe(ws: WebSocket, session_id: UUID) -> None: + await ws.accept() + + state = hub.get(session_id) + if state is None: + await ws.send_text( + ErrorEvent(code=ErrorKind.SESSION_NOT_FOUND, message="Занятие не запущено").model_dump_json() + ) + await ws.close() + return + + # Снимок при подключении обязателен: монитор в классе включают посреди + # занятия, и он должен показать текущее состояние, а не ждать событий. + await ws.send_text(state.snapshot().model_dump_json()) + + with hub.observer(session_id) as queue: + sender = asyncio.create_task(_pump(ws, queue)) + try: + await _wait_for_disconnect(ws) + except WebSocketDisconnect: + pass + finally: + sender.cancel() diff --git a/backend/app/db/repo.py b/backend/app/db/repo.py index f2925e3..a297b1d 100644 --- a/backend/app/db/repo.py +++ b/backend/app/db/repo.py @@ -29,6 +29,7 @@ async def create_session( mode: str, trainee_id: UUID | None = None, group_id: UUID | None = None, + session_id: UUID | None = None, ) -> Session: session = Session( scenario_id=scenario_id, @@ -37,11 +38,43 @@ async def create_session( group_id=group_id, attempt=await next_attempt(db, trainee_id, scenario_id), ) + if session_id is not None: + session.id = session_id db.add(session) await db.commit() return session +async def ensure_session( + db: AsyncSession, + *, + session_id: UUID, + scenario_id: str, + mode: str, + trainee_name: str | None = None, + group_name: str | None = None, +) -> Session: + """Занятие, запущенное с пульта, должно иметь строку в журнале. + + Иначе реплики, подсказки и пометки не к чему привязать: они уходят + в нарушение внешнего ключа, а профиль курсанта остаётся пустым. + """ + existing = await db.get(Session, session_id) + if existing is not None: + return existing + + group = await ensure_group(db, group_name) if group_name else None + trainee = await ensure_trainee(db, trainee_name, group) if trainee_name else None + return await create_session( + db, + scenario_id=scenario_id, + mode=mode, + trainee_id=trainee.id if trainee else None, + group_id=group.id if group else None, + session_id=session_id, + ) + + async def get_session(db: AsyncSession, session_id: UUID) -> Session | None: return await db.get(Session, session_id) diff --git a/backend/app/main.py b/backend/app/main.py index 79b321c..8202256 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -8,8 +8,14 @@ from pathlib import Path from app.api.http import scenarios as scenarios_api from app.api.http import sessions +from app.api.ws import call as call_ws +from app.api.ws import control as control_ws +from app.api.ws import observe as observe_ws from app.config import get_settings +from app.db.base import get_sessionmaker from app.scenarios import store +from app.session.hub import hub +from app.session.journal import DbJournal from app.scenarios.loader import ScenarioError @@ -26,14 +32,22 @@ async def lifespan(app: FastAPI): raise RuntimeError(f"библиотека сценариев не прошла проверку: {exc}") from exc app.state.scenarios_loaded = len(loaded) + # Журнал: всё, что не записано, для оценки не существует. + hub.journal = DbJournal(get_sessionmaker()) + # Прогрев моделей — карточка lct-06. app.state.models_ready = False yield + await hub.shutdown() + app = FastAPI(title="Учебный симулятор занятия для системы 112", lifespan=lifespan) app.include_router(sessions.router) app.include_router(scenarios_api.router) +app.include_router(call_ws.router) +app.include_router(observe_ws.router) +app.include_router(control_ws.router) @app.get("/api/health") diff --git a/backend/app/session/__init__.py b/backend/app/session/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/backend/app/session/hub.py b/backend/app/session/hub.py new file mode 100644 index 0000000..6e99bf7 --- /dev/null +++ b/backend/app/session/hub.py @@ -0,0 +1,142 @@ +"""Реестр живых сессий и подписки наблюдателей. + +Дельты курсанту, полное состояние наблюдателям — сознательно: у курсанта одно +соединение и важна задержка, наблюдателей несколько и подключаются они +в произвольный момент (docs/arch/CONTRACT.md). +""" + +import asyncio +import contextlib +from collections.abc import AsyncIterator, Iterator +from typing import Protocol +from uuid import UUID + +from pydantic import BaseModel + +from app.domain.events import TimerTick +from app.session.state import SessionState + +#: Очередь одного подписчика. Медленный наблюдатель не тормозит занятие: +#: очередь ограничена, переполнение роняет соединение, а не сессию. +QUEUE_SIZE = 256 + +TICK_SECONDS = 1.0 + + +class Journal(Protocol): + """Запись в БД. Вынесена за хаб: без базы занятие должно идти, + но молчать об ошибке записи нельзя.""" + + async def start_lesson( + self, session_id: UUID, scenario_id: str, mode: str, trainee_name: str | None + ) -> int: ... + async def utterance(self, session_id: UUID, entry) -> None: ... + async def hint(self, session_id: UUID, checklist_id: str, question: str, at) -> None: ... + async def note(self, session_id: UUID, ref: str, text: str, author: str) -> None: ... + async def self_assessment(self, session_id: UUID, missed: list[str], comment: str, at) -> None: ... + async def session_started(self, session_id: UUID, at) -> None: ... + async def session_ended(self, session_id: UUID, at, reason: str) -> None: ... + + +class SessionHub: + def __init__(self, journal: Journal | None = None) -> None: + self.journal = journal + self._sessions: dict[UUID, SessionState] = {} + self._observers: dict[UUID, set[asyncio.Queue]] = {} + self._trainees: dict[UUID, set[asyncio.Queue]] = {} + self._tickers: dict[UUID, asyncio.Task] = {} + + # ── реестр ── + + def register(self, state: SessionState) -> SessionState: + self._sessions[state.session_id] = state + return state + + def get(self, session_id: UUID) -> SessionState | None: + return self._sessions.get(session_id) + + def drop(self, session_id: UUID) -> None: + self._sessions.pop(session_id, None) + self.stop_ticker(session_id) + + # ── подписки ── + + @contextlib.contextmanager + def _subscribe(self, registry: dict[UUID, set[asyncio.Queue]], session_id: UUID) -> Iterator[asyncio.Queue]: + queue: asyncio.Queue = asyncio.Queue(maxsize=QUEUE_SIZE) + registry.setdefault(session_id, set()).add(queue) + try: + yield queue + finally: + registry.get(session_id, set()).discard(queue) + + def observer(self, session_id: UUID): + return self._subscribe(self._observers, session_id) + + def trainee(self, session_id: UUID): + return self._subscribe(self._trainees, session_id) + + # ── вещание ── + + @staticmethod + def _put(queues: set[asyncio.Queue], event: BaseModel) -> None: + for queue in list(queues): + try: + queue.put_nowait(event) + except asyncio.QueueFull: + queues.discard(queue) + + def to_observers(self, session_id: UUID, event: BaseModel) -> None: + self._put(self._observers.get(session_id, set()), event) + + def to_trainee(self, session_id: UUID, event: BaseModel) -> None: + self._put(self._trainees.get(session_id, set()), event) + + def broadcast(self, session_id: UUID, event: BaseModel) -> None: + self.to_trainee(session_id, event) + self.to_observers(session_id, event) + + def observer_count(self, session_id: UUID) -> int: + return len(self._observers.get(session_id, set())) + + # ── такт таймеров ── + + def start_ticker(self, session_id: UUID) -> None: + """`timer.tick` раз в секунду, а не на каждое изменение: таймеров дюжина, + а UI всё равно рисует секунды.""" + if session_id in self._tickers: + return + self._tickers[session_id] = asyncio.create_task(self._tick(session_id)) + + def stop_ticker(self, session_id: UUID) -> None: + task = self._tickers.pop(session_id, None) + if task is not None: + task.cancel() + + async def shutdown(self) -> None: + """Погасить все такты при остановке приложения. + + Без этого задачи тикеров переживают выключение и держат событийный цикл: + первым это ловит не продакшен, а тест, который не может закрыть клиент. + """ + for session_id in list(self._tickers): + self.stop_ticker(session_id) + + async def _tick(self, session_id: UUID) -> None: + try: + while True: + await asyncio.sleep(TICK_SECONDS) + state = self.get(session_id) + if state is None or state.ended: + return + self.broadcast(session_id, TimerTick(timers=state.timers.snapshot())) + except asyncio.CancelledError: + raise + + +hub = SessionHub() + + +async def drain(queue: asyncio.Queue) -> AsyncIterator[BaseModel]: + while True: + yield await queue.get() diff --git a/backend/app/session/journal.py b/backend/app/session/journal.py new file mode 100644 index 0000000..e52432c --- /dev/null +++ b/backend/app/session/journal.py @@ -0,0 +1,103 @@ +"""Запись хода занятия в БД. + +Профиль курсанта, дельта попыток и аналитика группы строятся по журналу, +а не по памяти процесса: всё, что здесь не записано, для оценки не существует. +""" + +import logging +from datetime import datetime +from uuid import UUID + +from sqlalchemy import update +from sqlalchemy.ext.asyncio import async_sessionmaker + +from app.db import repo +from app.db.models import SelfAssessment, Session + +log = logging.getLogger(__name__) + + +class DbJournal: + def __init__(self, sessionmaker: async_sessionmaker) -> None: + self._sessionmaker = sessionmaker + + async def _write(self, action, *args, **kwargs) -> None: + """Ошибка записи не роняет занятие, но и не проглатывается молча: + занятие идёт дальше, в логе остаётся след.""" + try: + async with self._sessionmaker() as db: + await action(db, *args, **kwargs) + except Exception: # noqa: BLE001 — журнал не должен ронять живую сессию + log.exception("журнал: запись не удалась") + + async def start_lesson( + self, session_id: UUID, scenario_id: str, mode: str, trainee_name: str | None + ) -> int: + """Завести сессию в журнале и вернуть номер попытки. + + Если база недоступна, занятие всё равно идёт: номер попытки + деградирует до первого, и это видно в логе. + """ + try: + async with self._sessionmaker() as db: + row = await repo.ensure_session( + db, + session_id=session_id, + scenario_id=scenario_id, + mode=mode, + trainee_name=trainee_name, + ) + return row.attempt + except Exception: # noqa: BLE001 — журнал не должен ронять живую сессию + log.exception("журнал: сессию завести не удалось") + return 1 + + async def utterance(self, session_id: UUID, entry) -> None: + await self._write( + repo.append_utterance, + session_id=session_id, + ref=entry.ref, + speaker=entry.speaker.value, + text=entry.text, + at=entry.at, + mood=entry.mood.value if entry.mood else None, + ) + + async def hint(self, session_id: UUID, checklist_id: str, question: str, at: datetime) -> None: + await self._write( + repo.record_hint, session_id=session_id, checklist_id=checklist_id, question=question, at=at + ) + + async def note(self, session_id: UUID, ref: str, text: str, author: str) -> None: + await self._write(repo.add_note, session_id=session_id, transcript_ref=ref, text=text, author=author) + + async def self_assessment( + self, session_id: UUID, missed: list[str], comment: str, at: datetime + ) -> None: + async def action(db): + db.add( + SelfAssessment( + session_id=session_id, missed=missed, comment=comment, submitted_at=at + ) + ) + await db.commit() + + await self._write(lambda db: action(db)) + + async def session_started(self, session_id: UUID, at: datetime) -> None: + async def action(db): + await db.execute(update(Session).where(Session.id == session_id).values(started_at=at)) + await db.commit() + + await self._write(lambda db: action(db)) + + async def session_ended(self, session_id: UUID, at: datetime, reason: str) -> None: + async def action(db): + await db.execute( + update(Session) + .where(Session.id == session_id) + .values(ended_at=at, end_reason=reason) + ) + await db.commit() + + await self._write(lambda db: action(db)) diff --git a/backend/app/session/state.py b/backend/app/session/state.py new file mode 100644 index 0000000..bd9b298 --- /dev/null +++ b/backend/app/session/state.py @@ -0,0 +1,99 @@ +"""Состояние живой сессии: карточка, транскрипт, таймеры, режим, попытка. + +Живёт в памяти процесса — поэтому воркер uvicorn ровно один: с двумя +преподаватель подключился бы к другому процессу, чем курсант, и увидел +пустой экран (docs/arch/STACK.md). +""" + +from dataclasses import dataclass, field +from datetime import datetime, timezone +from typing import Any +from uuid import UUID + +from app.domain.events import ( + CallEndReason, + Mood, + SessionMode, + SessionSnapshot, + Speaker, + TranscriptEntry, +) +from app.domain.kio import KIO, apply_patch +from app.session.timers import SessionTimers + + +def now_utc() -> datetime: + """Часы серверные. Метрика, посчитанная по часам браузера, недоказуема.""" + return datetime.now(timezone.utc) + + +@dataclass +class SessionState: + session_id: UUID + scenario_id: str + scenario_title: str + level: str + mode: SessionMode + required_fields: list[str] = field(default_factory=list) + trainee_name: str | None = None + attempt: int = 1 + + kio: KIO = field(default_factory=KIO) + transcript: list[TranscriptEntry] = field(default_factory=list) + timers: SessionTimers = field(default_factory=SessionTimers) + hints_shown: list[str] = field(default_factory=list) + directives: list[str] = field(default_factory=list) + + started_at: datetime | None = None + ended_at: datetime | None = None + end_reason: CallEndReason | None = None + dispatched_card: KIO | None = None + + def on_event(self, event_type: str) -> None: + """Единственная точка, где событие двигает таймеры.""" + self.timers.on_event(event_type) + + def append(self, speaker: Speaker, text: str, mood: Mood | None = None) -> TranscriptEntry: + entry = TranscriptEntry( + ref=f"u{len(self.transcript) + 1}", + speaker=speaker, + text=text, + at=now_utc(), + mood=mood, + ) + self.transcript.append(entry) + return entry + + def patch_kio(self, fields: dict[str, Any]) -> KIO: + self.kio = apply_patch(self.kio, fields) + return self.kio + + def dispatch(self, service: str) -> KIO: + """Карточка замораживается снимком: оператор не должен иметь + возможности дописать задним числом поле, которое забыл.""" + self.kio = apply_patch(self.kio, {"dds": service, "response_status": "transferred"}) + self.dispatched_card = self.kio.model_copy(deep=True) + return self.dispatched_card + + @property + def ended(self) -> bool: + return self.ended_at is not None + + def snapshot(self) -> SessionSnapshot: + """Полное состояние. Монитор в классе включают посреди занятия — + он обязан показать текущее, а не ждать следующего события.""" + return SessionSnapshot( + session_id=self.session_id, + scenario_id=self.scenario_id, + scenario_title=self.scenario_title, + level=self.level, + mode=self.mode, + trainee_name=self.trainee_name, + started_at=self.started_at, + kio=self.kio, + required_fields=self.required_fields, + transcript=list(self.transcript), + timers=self.timers.snapshot(), + hints_used=len(self.hints_shown), + ended=self.ended, + ) diff --git a/backend/app/session/timers.py b/backend/app/session/timers.py new file mode 100644 index 0000000..7a35b3b --- /dev/null +++ b/backend/app/session/timers.py @@ -0,0 +1,106 @@ +"""Таймеры сессии. + +**Правило: таймеры останавливаются событиями, а не таймаутами.** Иначе метрика +времени опроса становится недоказуемой, а вся оценка держится на том, что её +можно предъявить и проверить (docs/arch/CONTRACT.md). + +Связь «событие → таймер» объявлена таблицей, а не разбросана по обработчикам: +так видно целиком, что чем запускается, и норматив нельзя потерять по дороге. +""" + +import time +from dataclasses import dataclass, field + +from app.domain.timers import NORMATIVES, TimerCode, TimerSnapshot, state_for + +#: Какое событие какой таймер запускает. +STARTS: dict[str, tuple[TimerCode, ...]] = { + "call.incoming": (TimerCode.ANSWER,), + "call.answer": (TimerCode.INTERVIEW, TimerCode.DDS_NOTIFY), + "dds.dispatch": (TimerCode.DDS_ACK, TimerCode.CLOSE), + "card.received": (TimerCode.ZONE_CHECK,), + "callback.dial": (TimerCode.CALLBACK,), +} + +#: Какое событие какой таймер останавливает. +STOPS: dict[str, tuple[TimerCode, ...]] = { + "call.answer": (TimerCode.ANSWER,), + "dds.dispatch": (TimerCode.INTERVIEW, TimerCode.DDS_NOTIFY), + "card.ack": (TimerCode.DDS_ACK,), + "zone.decision": (TimerCode.ZONE_CHECK,), + "crew.arrived": (TimerCode.CLOSE,), + "call.started": (TimerCode.CALLBACK,), +} + + +@dataclass +class Timer: + code: TimerCode + started_at: float | None = None + elapsed_ms: int = 0 + attempt: int = 1 + stopped: bool = False + + def start(self, now: float) -> None: + if self.stopped: + # Повторный запуск после остановки — это новая попытка (обратный дозвон). + self.attempt += 1 + self.stopped = False + self.elapsed_ms = 0 + if self.started_at is None: + self.started_at = now + + def stop(self, now: float) -> None: + if self.started_at is not None and not self.stopped: + self.elapsed_ms = int((now - self.started_at) * 1000) + self.stopped = True + + def current_ms(self, now: float) -> int: + if self.stopped or self.started_at is None: + return self.elapsed_ms + return int((now - self.started_at) * 1000) + + +@dataclass +class SessionTimers: + """Набор таймеров одной сессии. `limits` приходит из конфига — + норматив меняется значением, а не правкой кода.""" + + limits: dict[TimerCode, int] = field( + default_factory=lambda: {code: norm.limit_ms for code, norm in NORMATIVES.items()} + ) + timers: dict[TimerCode, Timer] = field(default_factory=dict) + + def on_event(self, event_type: str, now: float | None = None) -> None: + now = time.monotonic() if now is None else now + for code in STARTS.get(event_type, ()): + self.timers.setdefault(code, Timer(code)).start(now) + for code in STOPS.get(event_type, ()): + self.timers.setdefault(code, Timer(code)).stop(now) + + def snapshot(self, now: float | None = None) -> list[TimerSnapshot]: + """Только запущенные таймеры: показывать нули по нормативам, + до которых занятие ещё не дошло, значит пугать курсанта зря.""" + now = time.monotonic() if now is None else now + result: list[TimerSnapshot] = [] + for code, timer in self.timers.items(): + limit = self.limits[code] + elapsed = timer.current_ms(now) + result.append( + TimerSnapshot( + code=code, + elapsed_ms=elapsed, + limit_ms=limit, + state=state_for(elapsed, limit), + attempt=timer.attempt, + stopped=timer.stopped, + ) + ) + return result + + def measured_ms(self, code: TimerCode) -> int | None: + """Зафиксированное событием значение — то, что пойдёт в оценку.""" + timer = self.timers.get(code) + if timer is None or not timer.stopped: + return None + return timer.elapsed_ms diff --git a/backend/tests/test_ws.py b/backend/tests/test_ws.py new file mode 100644 index 0000000..9c50ba8 --- /dev/null +++ b/backend/tests/test_ws.py @@ -0,0 +1,170 @@ +"""Каналы занятия целиком: преподаватель запускает, курсант работает, +наблюдатель смотрит и не может вмешаться. + +Проверяются пункты приёмки карточки lct-05, а не отдельные функции. +""" + +import contextlib +import time +from uuid import uuid4 + +import pytest +from fastapi.testclient import TestClient + +from app.main import app +from app.session.hub import hub + + +@pytest.fixture +def client(): + with TestClient(app) as test_client: + hub.journal = None # тесты не пишут в БД: проверяется поведение каналов + yield test_client + + +def wait_for(predicate, timeout: float = 3.0): + """Ждёт истинного значения. Для чисел не годится: опрос в тесте длится + 0 мс, и ноль — законный результат. Для них есть wait_value.""" + return _wait(predicate, timeout, lambda value: bool(value)) + + +def wait_value(predicate, timeout: float = 3.0): + """Ждёт любого значения, кроме None: ноль миллисекунд — тоже измерение.""" + return _wait(predicate, timeout, lambda value: value is not None) + + +def _wait(predicate, timeout: float, ready): + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + value = predicate() + if ready(value): + return value + time.sleep(0.02) + raise AssertionError("не дождались") + + +def read_until(ws, event_type: str, limit: int = 20) -> dict: + """Пропускает такты таймера: они идут раз в секунду и мешают разбирать поток.""" + for _ in range(limit): + message = ws.receive_json() + if message["type"] == event_type: + return message + raise AssertionError(f"событие {event_type} не пришло") + + +@contextlib.contextmanager +def lesson(client, mode: str = "training"): + """Занятие, запущенное преподавателем. + + Сокет закрывается выходом из контекстного менеджера, а не `close()`: + `close()` шлёт кадр закрытия, но не дожидается завершения задачи + соединения, и тестовый клиент потом не может закрыться. + """ + session_id = uuid4() + with client.websocket_connect(f"/ws/control/{session_id}") as control: + control.send_json( + { + "type": "scenario.start", + "scenario_id": "fire-apartment-l2", + "trainee": "Иванов И.И.", + "mode": mode, + } + ) + wait_for(lambda: hub.get(session_id)) + yield session_id, control + + +def test_observer_joining_midway_sees_the_whole_state(client): + """Монитор в классе включают посреди занятия — он обязан показать + текущее состояние, а не ждать следующего события.""" + with lesson(client) as (session_id, _): + with client.websocket_connect(f"/ws/call/{session_id}") as trainee: + trainee.send_json({"type": "call.answer"}) + trainee.send_json({"type": "kio.patch", "fields": {"floor": "5"}}) + wait_for(lambda: hub.get(session_id).kio.floor == "5") + + with client.websocket_connect(f"/ws/observe/{session_id}") as observer: + snapshot = observer.receive_json() + + assert snapshot["type"] == "session.snapshot" + assert snapshot["scenario_title"] == "Пожар в квартире, паникующий заявитель" + assert snapshot["kio"]["floor"] == "5", "карточка в снимке отстаёт от состояния" + assert snapshot["mode"] == "training" + assert snapshot["required_fields"], "обязательные поля не доехали до наблюдателя" + + +def test_interview_timer_starts_on_answer_and_stops_on_dispatch(client): + """Центральная метрика продукта. Таймер останавливается событием, + а не таймаутом: иначе время опроса недоказуемо.""" + from app.domain.timers import TimerCode + + with lesson(client) as (session_id, _): + state = hub.get(session_id) + with client.websocket_connect(f"/ws/call/{session_id}") as trainee: + trainee.send_json({"type": "call.answer"}) + wait_for(lambda: TimerCode.INTERVIEW in state.timers.timers) + assert state.timers.measured_ms(TimerCode.INTERVIEW) is None, "таймер уже остановлен" + + trainee.send_json({"type": "dds.dispatch", "service": "01"}) + measured = wait_value(lambda: state.timers.measured_ms(TimerCode.INTERVIEW)) + + assert measured >= 0 + assert state.kio.dds.value == "01" + assert state.dispatched_card is not None, "карточка не заморожена снимком" + + +def test_card_edit_reaches_observer(client): + with lesson(client) as (session_id, _): + with client.websocket_connect(f"/ws/observe/{session_id}") as observer: + observer.receive_json() # снимок при подключении + with client.websocket_connect(f"/ws/call/{session_id}") as trainee: + trainee.send_json({"type": "kio.patch", "fields": {"address": "улица Ленина, 14"}}) + message = read_until(observer, "kio.state") + + assert message["kio"]["address"] == "улица Ленина, 14" + + +def test_observer_cannot_influence_the_lesson(client): + """Разделение прав архитектурное: на сокете наблюдателя нет ни одного + обработчика входящих, поэтому влиять нечем.""" + with lesson(client) as (session_id, _): + state = hub.get(session_id) + with client.websocket_connect(f"/ws/observe/{session_id}") as observer: + observer.receive_json() + observer.send_json({"type": "kio.patch", "fields": {"address": "подмена"}}) + observer.send_json({"type": "session.stop"}) + time.sleep(0.3) + + assert state.kio.address is None, "наблюдатель изменил карточку" + assert not state.ended, "наблюдатель завершил занятие" + + +def test_hint_is_denied_in_exam_and_given_in_training(client): + with lesson(client, mode="exam") as (session_id, _): + with client.websocket_connect(f"/ws/call/{session_id}") as trainee: + trainee.send_json({"type": "hint.request"}) + message = read_until(trainee, "error") + assert message["code"] == "hint_denied_in_exam" + + with lesson(client, mode="training") as (session_id, _): + with client.websocket_connect(f"/ws/call/{session_id}") as trainee: + trainee.send_json({"type": "hint.request"}) + message = read_until(trainee, "hint.shown") + assert message["question"], "подсказка без вопроса" + assert hub.get(session_id).hints_shown, "использование подсказки не зафиксировано" + + +def test_instructor_cannot_touch_the_card(client): + """В канале `control` нет ни одной команды, меняющей карточку курсанта.""" + from app.domain.events import InstructorToServer + from typing import get_args + + union = get_args(get_args(InstructorToServer)[0]) + fields = {name for model in union for name in model.model_fields} + assert "fields" not in fields and "kio" not in fields + + +def test_call_socket_refuses_session_that_was_not_started(client): + with client.websocket_connect(f"/ws/call/{uuid4()}") as trainee: + message = trainee.receive_json() + assert message["type"] == "error" and message["code"] == "session_not_found"