diff --git a/backend/app/api/ws/control.py b/backend/app/api/ws/control.py index 7b73095..5c073a8 100644 --- a/backend/app/api/ws/control.py +++ b/backend/app/api/ws/control.py @@ -58,6 +58,9 @@ router = APIRouter() _adapter = TypeAdapter(InstructorToServer) +# Пульт может молчать; без входящих кадров потеря lease замечается этим опросом. +_FENCE_POLL_SECONDS = 1.0 + def _dds_ineligible_scenarios(scenarios): """Консультация и передача региона не являются готовыми карточками ДДС.""" @@ -378,11 +381,16 @@ async def control(ws: WebSocket, session_id: UUID) -> None: await _reject(ws, "Недостаточно прав для этого экрана") return event_stream = hub.begin_event_stream(session_id) + # Чтение живёт дольше тика опроса: wait_for отменял бы его, а отмена после + # того, как receive уже забрал кадр, теряет команду или websocket.disconnect. + # Команда обрабатывается здесь, в задаче, открывшей event stream. + read: asyncio.Future | None = None try: while True: - try: - payload = await asyncio.wait_for(ws.receive_json(), timeout=1) - except TimeoutError: + if read is None: + read = asyncio.ensure_future(ws.receive_json()) + done, _ = await asyncio.wait({read}, timeout=_FENCE_POLL_SECONDS) + if not done: if hub.is_lease_fenced(session_id): await ws.send_text(ErrorEvent( code=ErrorKind.INTERNAL, message=LEASE_FENCED_MESSAGE @@ -390,6 +398,7 @@ async def control(ws: WebSocket, session_id: UUID) -> None: await ws.close(code=1012) return continue + payload, read = read.result(), None try: event = _adapter.validate_python(payload) except ValidationError: @@ -541,4 +550,6 @@ async def control(ws: WebSocket, session_id: UUID) -> None: except WebSocketDisconnect: return finally: + if read is not None: + read.cancel() await hub.end_event_stream(event_stream) diff --git a/backend/tests/test_ws_ownership.py b/backend/tests/test_ws_ownership.py index 469470a..abf955b 100644 --- a/backend/tests/test_ws_ownership.py +++ b/backend/tests/test_ws_ownership.py @@ -195,3 +195,63 @@ def test_control_command_checkpoint_failure_returns_fencing_error(monkeypatch): finally: hub.journal = old_journal hub.drop(session_id) + + +def test_control_does_not_lose_command_read_during_fencing_poll(monkeypatch): + """Опрос fencing по таймеру не отменяет чтение, уже забравшее кадр из сокета. + + Сокет отдаёт команду, но возвращается из receive позже тика опроса — так + выглядит задержка планировщика под нагрузкой. Отмена чтения на тике + выбрасывает уже прочитанную команду. + """ + control_module = importlib.import_module("app.api.ws.control") + monkeypatch.setattr(control_module, "_FENCE_POLL_SECONDS", 0.02) + session_id = uuid4() + state = SessionState( + session_id=session_id, scenario_id="test", scenario_title="Тест", level="L1", + mode="training", owner_login="lease-owner", + ) + hub.register(state) + who = Principal(login="lease-owner", full_name="Преподаватель", role=Role.INSTRUCTOR) + monkeypatch.setattr(control_module, "websocket_origin_allowed", lambda _ws: True) + monkeypatch.setattr(control_module, "principal_of", lambda _ws: who) + + class SlowReturnSocket: + def __init__(self): + self.wire = [{"type": "no.such.command"}] + self.lost = [] + self.closed = None + self.sent = [] + + async def accept(self): + pass + + async def receive_json(self): + if not self.wire: + await hub.fence(state) + await asyncio.sleep(5) + message = self.wire.pop(0) + try: + await asyncio.sleep(0.1) + except asyncio.CancelledError: + self.lost.append(message) + raise + return message + + async def send_text(self, message): + self.sent.append(json.loads(message)) + + async def close(self, code=None): + self.closed = code + + socket = SlowReturnSocket() + try: + with hub.observer(session_id) as observer_queue: + asyncio.run(control_module.control(socket, session_id)) + assert socket.lost == [] + event = observer_queue.get_nowait() + assert event.code == "unsupported_event" + assert socket.closed == 1012 + assert socket.sent[-1]["message"].endswith("переподключитесь.") + finally: + hub.drop(session_id)