From 57c3ff5a5cde76c96e29145a145cfa1226ee1cbb Mon Sep 17 00:00:00 2001 From: kaifarikman Date: Sun, 27 Sep 2026 18:14:51 +0300 Subject: [PATCH] =?UTF-8?q?fix:=20=D1=83=D1=87=D1=91=D1=82=20=D1=81=D0=B2?= =?UTF-8?q?=D1=8F=D0=B7=D0=B8=20=D0=BA=D1=83=D1=80=D1=81=D0=B0=D0=BD=D1=82?= =?UTF-8?q?=D0=B0=20=D0=BF=D0=BE=20=D1=81=D0=BE=D0=BA=D0=B5=D1=82=D0=B0?= =?UTF-8?q?=D0=BC=20=D0=B8=20=D1=84=D0=B0=D0=B7=D0=B5=20=D0=94=D0=94=D0=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- backend/app/api/http/sessions.py | 26 +- backend/app/api/ws/call.py | 11 +- backend/app/api/ws/control.py | 1 + backend/app/api/ws/station.py | 11 +- backend/app/session/hub.py | 65 +++- backend/app/session/state.py | 12 +- backend/tests/test_checkpoint_model.py | 1 + backend/tests/test_session_signals.py | 339 ++++++++++++++++++- frontend/src/pages/instructor/Instructor.tsx | 5 +- frontend/src/shared/api/http.ts | 1 + 10 files changed, 425 insertions(+), 47 deletions(-) diff --git a/backend/app/api/http/sessions.py b/backend/app/api/http/sessions.py index a414536..be795a1 100644 --- a/backend/app/api/http/sessions.py +++ b/backend/app/api/http/sessions.py @@ -113,6 +113,7 @@ class ActiveSessionOut(BaseModel): dds_statuses: dict[str, str] dds_snapshot: StationSnapshot | None = None signals: list[Signal] = [] + presence_known: bool = True def _out(session) -> SessionOut: @@ -206,11 +207,8 @@ def _signals( подписчика в хабе этого узла; окно нужно, чтобы короткий разрыв соединения не сразу считался потерей связи курсанта. - `live=False` — снимок собран из checkpoint (сессию не держит хаб этого - процесса: рестарт до переподключения сокета или чужой узел в кластере). - Такой объект живёт один запрос и выбрасывается: писать в него - `socket_last_seen_at` бессмысленно — на следующем опросе присутствие - неизвестно снова. Честнее не утверждать «на связи», чем соврать. + `live=False` — снимок чужого узла: его сокеты здесь не видны, поэтому + присутствие неизвестно. Чтение реестра не меняет состояние занятия. """ settings = get_settings() signals: list[Signal] = [] @@ -226,16 +224,14 @@ def _signals( )) if live: connected = ( - hub.station_connected(state.session_id) if state.exercise is Exercise.DDS + hub.station_connected(state.session_id) if state.dds_phase else hub.trainee_connected(state.session_id) ) - if connected: - state.socket_last_seen_at = now offline_seconds = ( - (now - state.socket_last_seen_at).total_seconds() - if state.socket_last_seen_at is not None else 0 + (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 offline_seconds >= settings.signal_offline_window_seconds: + if not connected and offline_seconds >= settings.signal_offline_window_seconds: signals.append(Signal( kind="offline", severity="violated", text=f"Курсант не на связи {int(offline_seconds)} с", @@ -256,9 +252,8 @@ async def active( state.session_id: state for state in hub.active_sessions(who.login) } - #: Снимки без живой записи в хабе этого узла — рестарт до переподключения - #: сокета либо чужой узел в кластере. Присутствие сокета для них здесь - #: не проверяется: объект живёт один запрос, а не хаб (см. `_signals`). + #: Снимки без живой записи в хабе этого узла — чужой узел или сессия, + #: ожидающая восстановления. Её сокеты нельзя проверить этим процессом. checkpoint_only: set[UUID] = set() if db is not None: rows = ( @@ -293,7 +288,7 @@ async def active( for state in states.values(): elapsed = (max(0, int((now - state.started_at).total_seconds())) 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 [] @@ -327,6 +322,7 @@ 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, )) return result diff --git a/backend/app/api/ws/call.py b/backend/app/api/ws/call.py index 2d0f0d4..34e4b8d 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, @@ -390,9 +390,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 if state.exercise is Exercise.CARD: from app.api.ws.control import card_briefing diff --git a/backend/app/api/ws/control.py b/backend/app/api/ws/control.py index eb1ab26..179302f 100644 --- a/backend/app/api/ws/control.py +++ b/backend/app/api/ws/control.py @@ -104,6 +104,7 @@ def _build_state(session_id: UUID, event, who, scenario, scenarios, dds_service=identity.service or event.dds_service, attempt=identity.attempt, criteria=event.criteria, + socket_last_seen_at=now_utc(), ) state.timers.limits[TimerCode.DDS_ACK] = event.criteria.decision_time_limit_seconds * 1000 state.timers.limits[TimerCode.CARD_FILL] = event.criteria.card_fill_time_limit_seconds * 1000 diff --git a/backend/app/api/ws/station.py b/backend/app/api/ws/station.py index 8481c09..596341c 100644 --- a/backend/app/api/ws/station.py +++ b/backend/app/api/ws/station.py @@ -17,7 +17,7 @@ from uuid import UUID 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 ( CallEndReason, ErrorEvent, @@ -59,9 +59,14 @@ async def station(ws: WebSocket, session_id: UUID, role: str = "dds") -> None: entered = await session_socket(ws, session_id, (Role.INSTRUCTOR, Role.TRAINEE)) if entered is None: return - _who, state = entered + who, state = entered - with hub.station(session_id) as queue: + async with hub.trainee_socket( + session_id, station=True, trainee=who.role is Role.TRAINEE, + ) as queue: + if hub.is_lease_fenced(session_id): + await close_fenced(ws) + return sender = asyncio.create_task(pump(ws, queue)) try: # Карточка, переданная до подключения станции, не теряется: diff --git a/backend/app/session/hub.py b/backend/app/session/hub.py index 3a2bc29..bfe8aea 100644 --- a/backend/app/session/hub.py +++ b/backend/app/session/hub.py @@ -24,6 +24,7 @@ from app.domain.events import ( ) from app.session.state import SessionState from app.session.store import MemorySessionStore, Record, SessionLeaseLost, SessionStore +from app.session.timers import now_utc #: Очередь одного подписчика. Медленный наблюдатель не тормозит занятие: #: очередь ограничена, переполнение роняет соединение, а не сессию. @@ -53,6 +54,8 @@ class SessionHub: self._observers: dict[UUID, set[asyncio.Queue]] = {} self._trainees: dict[UUID, set[asyncio.Queue]] = {} self._stations: dict[UUID, set[asyncio.Queue]] = {} + self._trainee_calls: dict[UUID, set[asyncio.Queue]] = {} + self._trainee_stations: dict[UUID, set[asyncio.Queue]] = {} self._tickers: dict[UUID, asyncio.Task] = {} # Задачи, порождённые внутри операции, наследуют контекст; после # закрытия операции их события идут напрямую (`_Operation.open`). @@ -135,9 +138,10 @@ class SessionHub: if not self.store.persistent: return 0 restored = await self.store.restore_active() + adopted = 0 for state in restored: - self._adopt(state) - return len(restored) + adopted += await self._adopt(state) + return adopted async def maintain_lease(self) -> None: """Один оборот супервизора: продлить свои lease, подхватить просроченные чужие. @@ -163,17 +167,28 @@ class SessionHub: current = self._sessions.get(state.session_id) if current is not None and not current.lease_fenced: continue - self._adopt(state) + await self._adopt(state) async def supervise_lease(self, interval: float) -> None: while True: await asyncio.sleep(interval) await self.maintain_lease() - def _adopt(self, state: SessionState) -> None: + async def _adopt(self, state: SessionState) -> bool: self.stop_ticker(state.session_id) + if state.socket_connected_at_checkpoint: + # Разрыв произошёл при потере узла; старое время подключения не + # доказывает, что курсант отсутствовал всё это время. + state.socket_last_seen_at = now_utc() + state.socket_connected_at_checkpoint = False + try: + await self.store.commit(state) + except Exception: + log.exception("не удалось сохранить присутствие после takeover %s", state.session_id) + return False self.register(state) self.start_ticker(state.session_id) + return True async def save_all(self) -> None: """Снимок живых занятий при остановке узла: следующий владелец продолжит с него.""" @@ -285,6 +300,44 @@ class SessionHub: def station(self, session_id: UUID): return self._subscribe(self._stations, session_id) + async def _record_presence(self, session_id: UUID, *, station: bool) -> None: + state = self.get(session_id) + if state is None or state.ended or state.dds_phase != station: + return + async with self.operation(session_id): + state.socket_last_seen_at = now_utc() + present = self._trainee_stations if station else self._trainee_calls + state.socket_connected_at_checkpoint = bool(present.get(session_id)) + + @contextlib.asynccontextmanager + async def trainee_socket( + self, session_id: UUID, *, station: bool, trainee: bool, + ) -> AsyncIterator[asyncio.Queue]: + """Учесть только сокет курсанта; вещание преподавателю остаётся общим.""" + registry = self._stations if station else self._trainees + presence = self._trainee_stations if station else self._trainee_calls + with self._subscribe(registry, session_id) as queue: + if trainee: + presence.setdefault(session_id, set()).add(queue) + try: + if trainee: + try: + await self._record_presence(session_id, station=station) + except Exception: + if not self.is_lease_fenced(session_id): + raise + yield queue + finally: + if trainee: + presence[session_id].discard(queue) + if not presence[session_id]: + presence.pop(session_id) + try: + await self._record_presence(session_id, station=station) + except Exception: + if not self.is_lease_fenced(session_id): + raise + # ── вещание ── @staticmethod @@ -315,10 +368,10 @@ class SessionHub: return len(self._observers.get(session_id, set())) def station_connected(self, session_id: UUID) -> bool: - return bool(self._stations.get(session_id)) + return bool(self._trainee_stations.get(session_id)) def trainee_connected(self, session_id: UUID) -> bool: - return bool(self._trainees.get(session_id)) + return bool(self._trainee_calls.get(session_id)) # ── такт таймеров ── diff --git a/backend/app/session/state.py b/backend/app/session/state.py index 0e51d4d..28d0237 100644 --- a/backend/app/session/state.py +++ b/backend/app/session/state.py @@ -127,11 +127,12 @@ class PersistedSession(BaseModel): #: Подряд идущих первичных отказов без принятой карточки между ними — #: сигнал реестра преподавателя, не влияет на балл. consecutive_refusals: int = 0 - #: Когда реестр преподавателя последний раз видел подключённый сокет - #: курсанта (станция ДДС или звонок 112) — для сигнала "не на связи". - #: Персистится: без этого рестарт/failover обнуляет отсчёт окна - #: присутствия и сигнал никогда не срабатывает для восстановленных сессий. + #: Последнее подключение/отключение сокета курсанта или начало фазы ДДС. + #: Сохраняется вместе с занятием, чтобы чтение реестра не меняло состояние. socket_last_seen_at: datetime | None = None + #: На момент последнего checkpoint курсант был подключён на текущем АРМ. + #: После аварии узла момент потери связи неизвестен: окно начинается при takeover. + socket_connected_at_checkpoint: bool = False @field_validator("processed_station_commands") @classmethod @@ -214,6 +215,9 @@ class SessionState(PersistedSession): self.dispatched_at = now_utc() if self.exercise is Exercise.CALL: self._receive_call_card() + if self.handoff_to_dds: + self.socket_last_seen_at = self.dispatched_at + self.socket_connected_at_checkpoint = False return self.dispatched_card def _receive_call_card(self) -> None: diff --git a/backend/tests/test_checkpoint_model.py b/backend/tests/test_checkpoint_model.py index cd79f6a..144e06c 100644 --- a/backend/tests/test_checkpoint_model.py +++ b/backend/tests/test_checkpoint_model.py @@ -163,6 +163,7 @@ def full_state() -> SessionState: text_revealed_facts={"f_address": "улица Ленина, 14"}, consecutive_refusals=2, socket_last_seen_at=AT, + socket_connected_at_checkpoint=True, ) diff --git a/backend/tests/test_session_signals.py b/backend/tests/test_session_signals.py index 8b1d533..c867771 100644 --- a/backend/tests/test_session_signals.py +++ b/backend/tests/test_session_signals.py @@ -1,16 +1,26 @@ """Колонка сигналов реестра: очередь, повторные отказы, курсант не на связи.""" +import asyncio import time from datetime import UTC, datetime, timedelta +from types import SimpleNamespace from uuid import uuid4 import pytest from fastapi.testclient import TestClient +from starlette.websockets import WebSocketDisconnect +from app.api.auth import Principal from app.api.http import sessions as sessions_http +from app.api.ws import session as session_ws from app.config import get_settings +from app.domain.events import Exercise, SessionMode +from app.domain.roles import Role from app.main import app -from app.session.hub import hub +from app.session.checkpoint import dump_state, load_state +from app.session import hub as hub_module +from app.session.hub import SessionHub, hub +from app.session.state import SessionState from app.session.store import MemorySessionStore POOL = ["fire-apartment-l2", "t01-1-fire-container"] @@ -35,7 +45,12 @@ def client(monkeypatch): with TestClient(app) as test_client: test_client.post("/api/auth/dev-token") hub.store = MemorySessionStore() - yield test_client + before = set(hub._sessions) + try: + yield test_client + finally: + for session_id in set(hub._sessions) - before: + hub.drop(session_id) finally: get_settings.cache_clear() @@ -60,7 +75,7 @@ def read_until(socket, wanted): raise AssertionError(f"событие {wanted} не пришло; получены: {received}") -def start_two_card_dds(client): +def start_two_card_dds(client, scenario_ids=POOL): session_id = uuid4() context = client.websocket_connect(f"/ws/control/{session_id}") control = context.__enter__() @@ -70,7 +85,32 @@ def start_two_card_dds(client): "trainee": "Иванов", "mode": "training", "exercise": "dds", - "random_scenario_ids": POOL, + "random_scenario_ids": scenario_ids, + }) + wait_for(lambda: hub.get(session_id)) + return session_id, control + + +def start_card_handoff(client): + session_id = uuid4() + context = client.websocket_connect(f"/ws/control/{session_id}") + control = context.__enter__() + control.send_json({ + "type": "scenario.start", "scenario_id": POOL[0], "trainee": "Иванов", + "mode": "training", "exercise": "card", "handoff_to_dds": True, + "scenario_ids": POOL, + }) + wait_for(lambda: hub.get(session_id)) + return session_id, control + + +def start_call(client): + session_id = uuid4() + context = client.websocket_connect(f"/ws/control/{session_id}") + control = context.__enter__() + control.send_json({ + "type": "scenario.start", "scenario_id": POOL[0], "trainee": "Иванов", + "mode": "training", "exercise": "call", }) wait_for(lambda: hub.get(session_id)) return session_id, control @@ -85,6 +125,22 @@ def signal_kinds(row): return {signal["kind"] for signal in row["signals"]} +def assigned_trainee(state, monkeypatch): + state.trainee_id = uuid4() + who = Principal(login="курсант", full_name="Курсант", role=Role.TRAINEE, + trainee_id=state.trainee_id) + monkeypatch.setattr(session_ws, "principal_of", lambda _ws: who) + + +def registry_time(monkeypatch, since, seconds): + class Clock(datetime): + @classmethod + def now(cls, tz=None): + return since + timedelta(seconds=seconds) + + monkeypatch.setattr(sessions_http, "datetime", Clock) + + def test_signal_backlog_when_queue_reaches_threshold(client, monkeypatch): monkeypatch.setenv("SIGNAL_BACKLOG_THRESHOLD", "2") get_settings.cache_clear() @@ -134,19 +190,274 @@ def test_signal_refusals_after_two_consecutive_declines(client): control.__exit__(None, None, None) -def test_signal_offline_when_station_socket_is_closed_past_the_window(client): - session_id, control = start_two_card_dds(client) +def test_acceptance_clears_refusal_streak(client): + session_id, control = start_two_card_dds( + client, [*POOL, "t20-2-stroke"], + ) try: state = hub.get(session_id) - with client.websocket_connect(f"/ws/station/{session_id}?role=dds") as station: + with client.websocket_connect(f"/ws/station/{session_id}") as station: read_until(station, "station.state") - row = row_for(client, session_id) - assert "offline" not in signal_kinds(row), "сокет подключён — сигнала быть не должно" - - window = get_settings().signal_offline_window_seconds - state.socket_last_seen_at = datetime.now(UTC) - timedelta(seconds=window + 1) - row = row_for(client, session_id) - assert "offline" in signal_kinds(row) + for card in state.desk.ordered()[:2]: + station.send_json({"type": "card.open", "card_id": str(card.card_id)}) + read_until(station, "station.state") + station.send_json({ + "type": "card.status", "service": state.card_services(card)[0], + "status": "declined", "comment": "Не наша территория, передано дежурному", + }) + read_until(station, "station.state") + assert state.consecutive_refusals == 2 + third = state.desk.ordered()[2] + station.send_json({"type": "card.open", "card_id": str(third.card_id)}) + read_until(station, "station.state") + station.send_json({ + "type": "card.status", "service": state.card_services(third)[0], + "status": "accepted", "comment": "Карточка принята диспетчером", + }) + read_until(station, "station.state") + assert state.consecutive_refusals == 0 + assert "refusals" not in signal_kinds(row_for(client, session_id)) finally: hub.stop_ticker(session_id) control.__exit__(None, None, None) + + +def test_signal_offline_when_station_socket_is_closed_past_the_window(client, monkeypatch): + session_id, control = start_two_card_dds(client) + try: + state = hub.get(session_id) + assigned_trainee(state, monkeypatch) + with client.websocket_connect(f"/ws/station/{session_id}?role=dds") as station: + read_until(station, "station.state") + assert hub.station_connected(session_id) + assert not hub.station_connected(session_id) + disconnected_at = state.socket_last_seen_at + assert disconnected_at is not None + window = get_settings().signal_offline_window_seconds + registry_time(monkeypatch, disconnected_at, window - 1) + assert "offline" not in signal_kinds(row_for(client, session_id)) + registry_time(monkeypatch, disconnected_at, window + 1) + row = row_for(client, session_id) + assert "offline" in signal_kinds(row) + assert state.socket_last_seen_at == disconnected_at, "GET реестра не должен менять занятие" + finally: + hub.stop_ticker(session_id) + control.__exit__(None, None, None) + + +def test_signal_offline_without_first_connection(client, monkeypatch): + session_id, control = start_two_card_dds(client) + try: + state = hub.get(session_id) + started = state.socket_last_seen_at + assert started is not None + registry_time(monkeypatch, started, get_settings().signal_offline_window_seconds + 1) + assert "offline" in signal_kinds(row_for(client, session_id)) + finally: + hub.stop_ticker(session_id) + control.__exit__(None, None, None) + + +def test_call_socket_disconnection_triggers_offline(client, monkeypatch): + session_id, control = start_call(client) + try: + state = hub.get(session_id) + assigned_trainee(state, monkeypatch) + with client.websocket_connect(f"/ws/call/{session_id}") as call: + read_until(call, "call.incoming") + assert hub.trainee_connected(session_id) + assert not hub.trainee_connected(session_id) + disconnected_at = state.socket_last_seen_at + registry_time(monkeypatch, disconnected_at, + get_settings().signal_offline_window_seconds + 1) + assert "offline" in signal_kinds(row_for(client, session_id)) + finally: + hub.stop_ticker(session_id) + control.__exit__(None, None, None) + + +def test_teacher_socket_does_not_hide_offline(client, monkeypatch): + session_id, control = start_two_card_dds(client) + try: + state = hub.get(session_id) + with client.websocket_connect(f"/ws/station/{session_id}") as station: + read_until(station, "station.state") + assert not hub.station_connected(session_id) + registry_time(monkeypatch, state.socket_last_seen_at, + get_settings().signal_offline_window_seconds + 1) + assert "offline" in signal_kinds(row_for(client, session_id)) + finally: + hub.stop_ticker(session_id) + control.__exit__(None, None, None) + + +def test_last_of_two_trainee_sockets_starts_offline_window(client, monkeypatch): + session_id, control = start_two_card_dds(client) + try: + state = hub.get(session_id) + assigned_trainee(state, monkeypatch) + with client.websocket_connect(f"/ws/station/{session_id}") as first: + read_until(first, "station.state") + with client.websocket_connect(f"/ws/station/{session_id}") as second: + read_until(second, "station.state") + assert hub.station_connected(session_id) + registry_time(monkeypatch, state.socket_last_seen_at, + get_settings().signal_offline_window_seconds + 1) + assert "offline" not in signal_kinds(row_for(client, session_id)) + disconnected_at = state.socket_last_seen_at + registry_time(monkeypatch, disconnected_at, + get_settings().signal_offline_window_seconds + 1) + assert "offline" in signal_kinds(row_for(client, session_id)) + finally: + hub.stop_ticker(session_id) + control.__exit__(None, None, None) + + +def test_offline_window_survives_checkpoint_restore(client, monkeypatch): + session_id, control = start_two_card_dds(client) + try: + state = hub.get(session_id) + assigned_trainee(state, monkeypatch) + with client.websocket_connect(f"/ws/station/{session_id}") as station: + read_until(station, "station.state") + disconnected_at = state.socket_last_seen_at + restored = load_state(dump_state(state), datetime.now(UTC)) + hub.drop(session_id) + asyncio.run(hub._adopt(restored)) + assert restored.socket_last_seen_at == disconnected_at + registry_time(monkeypatch, disconnected_at, + get_settings().signal_offline_window_seconds + 1) + assert "offline" in signal_kinds(row_for(client, session_id)) + finally: + hub.stop_ticker(session_id) + control.__exit__(None, None, None) + + +def test_takeover_starts_window_when_socket_was_open(client, monkeypatch): + session_id, control = start_two_card_dds(client) + try: + state = hub.get(session_id) + assigned_trainee(state, monkeypatch) + with client.websocket_connect(f"/ws/station/{session_id}") as station: + read_until(station, "station.state") + assert state.socket_connected_at_checkpoint + snapshot = dump_state(state) + old_time = state.socket_last_seen_at + takeover_at = old_time + timedelta(minutes=10) + monkeypatch.setattr(hub_module, "now_utc", lambda: takeover_at) + restored = load_state(snapshot, datetime.now(UTC)) + + class CapturingStore(MemorySessionStore): + def __init__(self): + self.saved = None + + async def commit(self, state, _records=()): + self.saved = dump_state(state) + + store = CapturingStore() + local_hub = SessionHub(store) + try: + assert asyncio.run(local_hub._adopt(restored)) + assert restored.socket_last_seen_at == takeover_at + assert not restored.socket_connected_at_checkpoint + assert datetime.fromisoformat( + store.saved["socket_last_seen_at"].replace("Z", "+00:00") + ) == takeover_at + assert store.saved["socket_connected_at_checkpoint"] is False + finally: + local_hub.stop_ticker(session_id) + finally: + hub.stop_ticker(session_id) + control.__exit__(None, None, None) + + +def test_card_handoff_uses_trainee_station_presence(client, monkeypatch): + session_id, control = start_card_handoff(client) + try: + state = hub.get(session_id) + assigned_trainee(state, monkeypatch) + with monkeypatch.context() as clock: + with client.websocket_connect(f"/ws/call/{session_id}") as call: + read_until(call, "card.briefing") + call.send_json({"type": "card.submit"}) + read_until(call, "call.ended") + handoff_at = state.socket_last_seen_at + clock.setattr(hub_module, "now_utc", lambda: handoff_at + timedelta(minutes=1)) + assert state.socket_last_seen_at == handoff_at, "старый call-сокет не сбрасывает окно ДДС" + assert state.exercise.value == "card" and state.dds_phase + assert state.desk.scenarios + with client.websocket_connect(f"/ws/station/{session_id}") as station: + read_until(station, "station.state") + registry_time(monkeypatch, state.socket_last_seen_at, + get_settings().signal_offline_window_seconds + 1) + row = row_for(client, session_id) + assert row["dds_open_cards"] > 0 + assert "offline" not in signal_kinds(row) + finally: + hub.stop_ticker(session_id) + control.__exit__(None, None, None) + + +def test_checkpoint_on_another_node_marks_presence_unknown(client, monkeypatch): + session_id, control = start_two_card_dds(client) + try: + state = hub.get(session_id) + snapshot = dump_state(state) + checkpoint_at = datetime.now(UTC) + hub.drop(session_id) + + class CheckpointDB: + async def scalars(self, _query): + row = SimpleNamespace(id=session_id, owner_login="dev", + live_state=snapshot, checkpoint_at=checkpoint_at) + return SimpleNamespace(all=lambda: [row]) + + monkeypatch.setattr(sessions_http, "require", lambda _request, _role: + Principal(login="dev", full_name="Преподаватель", + role=Role.INSTRUCTOR)) + rows = asyncio.run(sessions_http.active(None, db=CheckpointDB())) + row = next(item for item in rows if item.session_id == session_id) + assert row.presence_known is False + assert "offline" not in {item.kind for item in row.signals} + finally: + hub.stop_ticker(session_id) + control.__exit__(None, None, None) + + +def test_failed_presence_commit_does_not_leave_connected_socket(): + class BrokenStore(MemorySessionStore): + async def commit(self, _state, _records=()): + raise OSError("test commit failure") + + local_hub = SessionHub(BrokenStore()) + state = local_hub.register(SessionState( + session_id=uuid4(), scenario_id="test", scenario_title="Тест", level="L1", + mode=SessionMode.TRAINING, exercise=Exercise.DDS, + )) + + async def connect(): + async with local_hub.trainee_socket(state.session_id, station=True, trainee=True): + assert local_hub.is_lease_fenced(state.session_id) + + asyncio.run(connect()) + assert not local_hub.station_connected(state.session_id) + + +def test_failed_presence_commit_closes_station_with_fencing_code(client, monkeypatch): + class BrokenStore(MemorySessionStore): + async def commit(self, _state, _records=()): + raise OSError("test commit failure") + + session_id, control = start_two_card_dds(client) + try: + assigned_trainee(hub.get(session_id), monkeypatch) + hub.store = BrokenStore() + with client.websocket_connect(f"/ws/station/{session_id}") as station: + event = station.receive_json() + assert event["type"] == "error" and event["code"] == "internal" + with pytest.raises(WebSocketDisconnect) as closed: + station.receive_json() + assert closed.value.code == 1012 + assert not hub.station_connected(session_id) + finally: + control.__exit__(None, None, None) diff --git a/frontend/src/pages/instructor/Instructor.tsx b/frontend/src/pages/instructor/Instructor.tsx index 8a9f9f8..b618a04 100644 --- a/frontend/src/pages/instructor/Instructor.tsx +++ b/frontend/src/pages/instructor/Instructor.tsx @@ -390,10 +390,11 @@ export function Instructor() { ? `Просрочено: первичная реакция ${item.dds_overdue_cards}; отработка карточки ${item.dds_work_overdue_cards}` : "Нормативы не нарушены"} - + {item.signals.length ? item.signals.map((signal) => {signal.text}
) - : "Сигналов нет"} + : item.presence_known ? "Сигналов нет" : null} + {!item.presence_known && Связь курсанта не проверена}