diff --git a/backend/app/api/http/sessions.py b/backend/app/api/http/sessions.py index e6d3569..1996e95 100644 --- a/backend/app/api/http/sessions.py +++ b/backend/app/api/http/sessions.py @@ -8,6 +8,7 @@ import logging import time from collections.abc import AsyncIterator from datetime import UTC, datetime +from typing import Literal from uuid import UUID from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response @@ -23,8 +24,8 @@ from app.db.base import get_session from app.db.models import Group, Score, Session, Trainee from app.domain.events import Exercise, SessionMode, SessionReport from app.domain.roles import Role -from app.domain.statuses import SERVICE_STATUS_LABELS, StationSnapshot, current -from app.domain.timers import TimerCode +from app.domain.statuses import SERVICE_STATUS_LABELS, DdsQueueCard, StationSnapshot, current +from app.domain.timers import TimerCode, TimerState from app.scoring.export import to_csv, to_pdf from app.scoring.report import build as build_report from app.session.access import can_access @@ -32,6 +33,7 @@ from app.session.checkpoint import load_state from app.session.finish import override_score from app.session.score import scoring_scenario from app.session.hub import hub +from app.session.state import SessionState from app.session.store import ScoreOverridden, apply_score_override from app.voice.recording import recording_path @@ -86,6 +88,15 @@ class DdsHistoryOut(BaseModel): recipient_services: list[str] = [] +class Signal(BaseModel): + """Одна строка колонки сигналов реестра. Цвет несёт смысл только через + `severity`; текст обязателен и не заменяется цветом.""" + + kind: Literal["backlog", "refusals", "offline"] + severity: TimerState + text: str + + class ActiveSessionOut(BaseModel): session_id: UUID trainee_name: str | None @@ -101,6 +112,8 @@ class ActiveSessionOut(BaseModel): dds_work_overdue_cards: int dds_statuses: dict[str, str] dds_snapshot: StationSnapshot | None = None + signals: list[Signal] = [] + presence_known: bool = True paused: bool = False @@ -185,6 +198,53 @@ async def dds_history( return result +def _signals( + state: SessionState, queue: list[DdsQueueCard], now: datetime, *, live: bool +) -> list[Signal]: + """Сигналы реестра: очередь, повторные отказы, курсант не на связи. + + Пауза занятия (`session.pause`) их не трогает и ложных срабатываний не + даёт: очередь на паузе не выдаёт карточки, команды станции отклоняются, + поэтому счётчики стоят; присутствие сокета — факт и на паузе. + Присутствие — по наличию подписчика в хабе этого узла; окно нужно, чтобы + короткий разрыв соединения не сразу считался потерей связи курсанта. + + `live=False` — снимок чужого узла: его сокеты здесь не видны, поэтому + присутствие неизвестно. Чтение реестра не меняет состояние занятия. + """ + settings = get_settings() + signals: list[Signal] = [] + # Необработанная — без первичного статуса: он останавливает норматив + # DDS_ACK. Принятая или отклонённая карточка висит в очереди до + # «Следующей», но реакции курсанта уже не ждёт. + pending = sum(not card.timer_stopped for card in queue) + if pending >= settings.signal_backlog_threshold: + signals.append(Signal( + kind="backlog", severity=TimerState.WARN, + text=f"Очередь: {pending} необработанных карточек", + )) + if state.consecutive_refusals >= settings.signal_refusals_threshold: + signals.append(Signal( + kind="refusals", severity=TimerState.WARN, + text=f"Подряд отказов: {state.consecutive_refusals}", + )) + if live: + connected = ( + hub.station_connected(state.session_id) if state.dds_phase + else hub.trainee_connected(state.session_id) + ) + offline_seconds = ( + (now - (state.socket_last_seen_at or state.started_at)).total_seconds() + if (state.socket_last_seen_at or state.started_at) is not None else 0 + ) + if not connected and offline_seconds >= settings.signal_offline_window_seconds: + signals.append(Signal( + kind="offline", severity=TimerState.VIOLATED, + text=f"Курсант не на связи {int(offline_seconds)} с", + )) + return sorted(signals, key=lambda signal: signal.severity is not TimerState.VIOLATED) + + @router.get("/active", response_model=list[ActiveSessionOut]) async def active( request: Request, @@ -198,6 +258,9 @@ async def active( state.session_id: state for state in hub.active_sessions(who.login) } + #: Снимки без живой записи в хабе этого узла — чужой узел или сессия, + #: ожидающая восстановления. Её сокеты нельзя проверить этим процессом. + checkpoint_only: set[UUID] = set() if db is not None: rows = ( await db.scalars( @@ -226,6 +289,7 @@ async def active( state.owner_login = row.owner_login if not state.ended: states[state.session_id] = state + checkpoint_only.add(state.session_id) for state in states.values(): paused_ms = state.total_paused_ms @@ -233,7 +297,7 @@ async def active( paused_ms += max(0, int((now - state.paused_at).total_seconds() * 1000)) elapsed = (max(0, int((now - state.started_at).total_seconds() - paused_ms / 1000)) if state.started_at else 0) - station = state.station_snapshot() if state.exercise is Exercise.DDS else None + station = state.station_snapshot() if state.dds_phase else None queue = station.queue_cards if station else [] card = state.desk.active status_log = card.status_log if card is not None else [] @@ -268,6 +332,8 @@ async def active( ), dds_statuses=latest_statuses, dds_snapshot=station, + signals=_signals(state, queue, now, live=state.session_id not in checkpoint_only), + presence_known=state.session_id not in checkpoint_only, paused=state.paused, )) return result diff --git a/backend/app/api/ws/call.py b/backend/app/api/ws/call.py index 0ef1cee..37ac20b 100644 --- a/backend/app/api/ws/call.py +++ b/backend/app/api/ws/call.py @@ -14,7 +14,7 @@ from uuid import UUID, uuid4 from fastapi import APIRouter, WebSocket, WebSocketDisconnect from pydantic import TypeAdapter, ValidationError -from app.api.ws.session import pump, run_command, session_socket +from app.api.ws.session import close_fenced, pump, run_command, session_socket from app.domain.events import ( BgStart, CallEnded, @@ -40,6 +40,7 @@ from app.domain.events import ( TraineeToServer, ) from app.dialog.slots import TurnResult +from app.domain.kio import PatchRejected, get_field from app.domain.roles import Role from app.scenarios import store from app.session.dds import prepare_handoff_queue @@ -193,6 +194,9 @@ async def _handle(session_id: UUID, state, event) -> None: )) return if state.paused: + # Отказ временный, в отличие от отклонённого значения: `kio.patch` + # без `auto`-отката. useCall держит правку в pending и повторяет её + # после `session.paused: false` — откат стёр бы ввод курсанта. hub.to_trainee(session_id, ErrorEvent( code=ErrorKind.UNSUPPORTED_EVENT, message="Пауза, ждите преподавателя", )) @@ -279,16 +283,37 @@ async def _handle(session_id: UUID, state, event) -> None: case "kio.patch": old_code, old_notify = state.kio.incident_code, list(state.kio.notify) - state.patch_kio(event.fields) + old_extra = list(state.kio.notify_extra) + try: + state.patch_kio(event.fields) + except (ValidationError, PatchRejected) as exc: + reason = exc.errors()[0].get("msg", "") if isinstance(exc, ValidationError) else str(exc) + # Недопустимое значение — отказ курсанту, а не сбой операции: + # сбой закрыл бы занятие на этом узле. + hub.to_trainee(session_id, ErrorEvent( + code=ErrorKind.UNSUPPORTED_EVENT, + message=f"Поле карточки не принято: {reason}"[:200], + )) + # Пакет отклонён целиком, а фронт держит его в pending до эха. + # Серверное значение с `auto` снимает pending — иначе на экране + # останутся правки, которых на сервере нет. + hub.to_trainee(session_id, KioPatchOut( + fields={path: get_field(state.kio, path) for path in event.fields}, + source=PatchSource.AUTO, + )) + return hub.to_trainee(session_id, KioPatchOut( fields=event.fields, source=PatchSource.OPERATOR, )) + auto: dict = {} if (state.kio.incident_code, state.kio.notify) != (old_code, old_notify): - hub.to_trainee(session_id, KioPatchOut( - fields={"incident_code": state.kio.incident_code, - "notify": list(state.kio.notify)}, - source=PatchSource.AUTO, - )) + auto |= {"incident_code": state.kio.incident_code, "notify": list(state.kio.notify)} + # Добавки сервер чистит сам (повторы, службы из notify) — фронту + # нужно итоговое значение, а не эхо присланного. + if state.kio.notify_extra != event.fields.get("notify_extra", old_extra): + auto["notify_extra"] = list(state.kio.notify_extra) + if auto: + hub.to_trainee(session_id, KioPatchOut(fields=auto, source=PatchSource.AUTO)) # Наблюдателю уходит карточка целиком: рассинхрон на внешнем мониторе # посреди занятия дороже лишних килобайт. hub.to_observers(session_id, KioState(kio=state.kio)) @@ -399,9 +424,14 @@ async def call(ws: WebSocket, session_id: UUID) -> None: ) if entered is None: return - _who, state = entered + who, state = entered - with hub.trainee(session_id) as queue: + async with hub.trainee_socket( + session_id, station=False, trainee=who.role is Role.TRAINEE, + ) as queue: + if hub.is_lease_fenced(session_id): + await close_fenced(ws) + return # Каждый новый канал получает серверное состояние, включая снятую паузу: # клиент мог потерять событие resume во время переподключения. queue.put_nowait(SessionPaused(paused=state.paused)) diff --git a/backend/tests/test_checkpoint_model.py b/backend/tests/test_checkpoint_model.py index 4dd821a..5cc764e 100644 --- a/backend/tests/test_checkpoint_model.py +++ b/backend/tests/test_checkpoint_model.py @@ -162,6 +162,9 @@ def full_state() -> SessionState: resolve_comment="передано в другой регион", processed_station_commands=[str(uuid4())], text_revealed_facts={"f_address": "улица Ленина, 14"}, + consecutive_refusals=2, + socket_last_seen_at=AT, + socket_connected_at_checkpoint=True, paused=True, paused_at=AT, total_paused_ms=15_000,