"""Сокет диспетчера ДДС. **Только JSON: аудио здесь нет вообще** — значит, нет ни VAD, ни распознавания, ни синтеза. Голосовой контур станция не трогает (docs/arch/CONTRACT.md). Самая ценная механика цепочки — `card.bounce`: диспетчер видит, что не указан этаж, и отбивает карточку обратно. Неполнота КИО перестаёт быть процентом в отчёте и становится сорванным выездом с конкретной причиной. """ import asyncio import logging from uuid import UUID from fastapi import APIRouter, WebSocket, WebSocketDisconnect from pydantic import TypeAdapter, ValidationError from app.domain.events import ErrorEvent, ErrorKind, StationState, StationToServer from app.domain.statuses import PRIMARY, ServiceStatus, StatusError from app.api.auth import principal_of from app.domain.roles import Role from app.session.hub import hub from app.session.state import now_utc log = logging.getLogger(__name__) router = APIRouter() _adapter = TypeAdapter(StationToServer) async def _handle(session_id: UUID, state, event) -> None: match event.type: case "card.ack": # Подтверждение приёма — это статус «Принята» у главной службы. # Кнопка осталась ради живой цепочки 112 → ДДС (lct-20), где # диспетчер один и выбирать службу не из чего. state.on_event("card.ack") state.dds_log.append(("card.ack", now_utc(), None)) services = state.notified_services() if services: try: state.set_service_status(services[0], ServiceStatus.ACCEPTED) except StatusError: pass # статус уже стоит: повторное нажатие ничего не меняет case "card.status": try: state.set_service_status( event.service, event.status, event.comment, author="диспетчер" ) except StatusError as exc: hub.to_station( session_id, ErrorEvent(code=ErrorKind.UNSUPPORTED_EVENT, message=str(exc)), ) return # Первичный статус останавливает норматив 30 секунд. if event.status in PRIMARY: state.on_event("card.ack") case "card.bounce": # Карточка вернулась: в разборе это E6 с конкретной причиной. state.bounced_fields = list(event.missing_fields) state.dds_log.append(("card.bounce", now_utc(), event.comment)) case "zone.decision": state.on_event("zone.decision") state.dds_log.append(("zone.decision", now_utc(), "в зоне" if event.in_zone else "не в зоне")) case "crew.dispatched": state.kio = state.kio.model_copy(update={"dispatch_order_at": event.at}) state.dds_log.append(("crew.dispatched", now_utc(), None)) case "crew.arrived": state.on_event("crew.arrived") state.kio = state.kio.model_copy(update={"arrival_at": event.at}) state.dds_log.append(("crew.arrived", now_utc(), None)) hub.to_station(session_id, StationState(snapshot=state.station_snapshot())) hub.to_observers(session_id, state.snapshot()) 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 _reject(ws: WebSocket, message: str) -> None: """Отказ до входа в цикл: сокет закрывается с объяснением, а не молча.""" await ws.send_text( ErrorEvent(code=ErrorKind.FORBIDDEN, message=message).model_dump_json() ) await ws.close() @router.websocket("/ws/station/{session_id}") async def station(ws: WebSocket, session_id: UUID, role: str = "dds") -> None: await ws.accept() # За АРМ ДДС садится обучающийся, преподаватель смотрит и подменяет. who = principal_of(ws) if who is None or who.role not in (Role.INSTRUCTOR, Role.TRAINEE): await _reject(ws, "Недостаточно прав для этого экрана") return 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.station(session_id) as queue: sender = asyncio.create_task(_pump(ws, queue)) try: # Карточка, переданная до подключения станции, не теряется: # диспетчер садится за АРМ, когда вызов уже идёт. if state.dispatched_card is not None: hub.to_station(session_id, state.card_received_event()) hub.to_station(session_id, StationState(snapshot=state.station_snapshot())) while True: payload = await ws.receive_json() try: event = _adapter.validate_python(payload) except ValidationError: hub.to_station( session_id, ErrorEvent(code=ErrorKind.UNSUPPORTED_EVENT, message=str(payload)[:200]), ) continue await _handle(session_id, state, event) except WebSocketDisconnect: return finally: sender.cancel()