493 lines
25 KiB
Python
493 lines
25 KiB
Python
"""Сокет диспетчера ДДС.
|
||
|
||
Учебный исходящий звонок бригаде имитируется событиями JSON. Для целевого
|
||
локального прототипа не нужна внешняя телефония или аудиомодель.
|
||
|
||
Самая ценная механика цепочки — `card.bounce`: диспетчер видит, что не указан
|
||
этаж, и отбивает карточку обратно. Неполнота КИО перестаёт быть процентом
|
||
в отчёте и становится сорванным выездом с конкретной причиной.
|
||
"""
|
||
|
||
import asyncio
|
||
import logging
|
||
import re
|
||
from uuid import UUID
|
||
|
||
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
|
||
from pydantic import TypeAdapter, ValidationError
|
||
|
||
from app.api.auth import principal_of, websocket_origin_allowed
|
||
from app.domain.events import (
|
||
CallEndReason,
|
||
CommandAck,
|
||
ErrorEvent,
|
||
ErrorKind,
|
||
Exercise,
|
||
PhoneLine,
|
||
PhoneReport,
|
||
ScoreReady,
|
||
SessionEnded,
|
||
StationState,
|
||
StationToServer,
|
||
)
|
||
from app.domain.roles import Role
|
||
from app.domain.statuses import (
|
||
PRIMARY,
|
||
SERVICE_STATUS_LABELS,
|
||
PhoneCallPending,
|
||
PhoneLineRecord,
|
||
PhoneReportRecord,
|
||
ServiceStatus,
|
||
StatusError,
|
||
current,
|
||
)
|
||
from app.domain.timers import TimerCode
|
||
from app.scoring.address import address_matches
|
||
from app.session.dds import deliver_due_cards
|
||
from app.session.finish import finish, score_current_dds
|
||
from app.session.hub import LEASE_FENCED_MESSAGE, hub
|
||
from app.session.state import now_utc
|
||
|
||
log = logging.getLogger(__name__)
|
||
router = APIRouter()
|
||
|
||
_adapter = TypeAdapter(StationToServer)
|
||
|
||
|
||
def _start_dds_work_timer(state) -> None:
|
||
"""Start the three-minute work clock once, when the card is opened."""
|
||
timer = state.timers.timers.get(TimerCode.DDS_WORK)
|
||
if timer is None or timer.started_at is None:
|
||
state.on_event("dds.open")
|
||
|
||
REPORT_PHASES = ("dispatched", "arrived", "working", "completed")
|
||
REQUIRED_STATUS = {
|
||
"dispatched": ServiceStatus.ACCEPTED,
|
||
"arrived": ServiceStatus.RESPONDING,
|
||
"working": ServiceStatus.ARRIVED,
|
||
"completed": ServiceStatus.WORKING,
|
||
}
|
||
STATUS_AT_OR_AFTER = {
|
||
"dispatched": {ServiceStatus.ACCEPTED, ServiceStatus.RESPONDING,
|
||
ServiceStatus.ARRIVED, ServiceStatus.WORKING},
|
||
"arrived": {ServiceStatus.RESPONDING, ServiceStatus.ARRIVED, ServiceStatus.WORKING},
|
||
"working": {ServiceStatus.ARRIVED, ServiceStatus.WORKING},
|
||
"completed": {ServiceStatus.WORKING, ServiceStatus.COMPLETED},
|
||
}
|
||
def _error(session_id: UUID, message: str) -> None:
|
||
hub.to_station(session_id, ErrorEvent(code=ErrorKind.UNSUPPORTED_EVENT, message=message))
|
||
|
||
|
||
def _line(session_id: UUID, state, speaker: str, text: str) -> None:
|
||
call = state.phone_pending
|
||
if call is None:
|
||
return
|
||
line = PhoneLineRecord(service=call.service, crew=call.crew,
|
||
speaker=speaker, text=text, at=now_utc())
|
||
state.phone_lines.append(line)
|
||
state.dds_log.append((f"phone.line.{speaker}", line.at, f"{call.crew}: {text}"))
|
||
hub.to_station(session_id, PhoneLine(**line.model_dump()))
|
||
|
||
|
||
def _address_matches(expected: str | None, supplied: str) -> bool:
|
||
"""Не даём сообщить бригаде другой номер дома/другую улицу."""
|
||
return address_matches(expected, supplied)
|
||
|
||
|
||
def _incident_matches(state, supplied: str) -> bool:
|
||
text = supplied.casefold()
|
||
# Описание часто начинается с адреса: его нельзя считать совпадением
|
||
# характера происшествия. Берём только название сценария и признаки ЕКП.
|
||
source = " ".join((state.scenario_title, " ".join(state.kio.signs)))
|
||
anchors = {word[:4] for word in re.findall(r"[а-яё]{5,}", source.casefold())}
|
||
return len(supplied.strip()) >= 8 and any(anchor in text for anchor in anchors)
|
||
|
||
|
||
def _has_purpose(text: str, stems: tuple[str, ...]) -> bool:
|
||
normalized = text.casefold()
|
||
return len(text.strip()) >= 8 and any(stem in normalized for stem in stems)
|
||
|
||
|
||
def _report_text(phase: str, crew: str, address: str) -> str:
|
||
match phase:
|
||
case "dispatched":
|
||
return f"{crew}: вызов по адресу {address} принят, выезжаем. О прибытии доложу."
|
||
case "arrived":
|
||
return f"{crew}: прибыли по адресу {address}. Уточняем обстановку на месте."
|
||
case "working":
|
||
return f"{crew}: обстановка уточнена, приступили к работам. Сообщим о завершении."
|
||
case _:
|
||
return f"{crew}: работы завершены. Дальнейшая помощь от нашей бригады не требуется."
|
||
|
||
|
||
def _finish_phone_call(session_id: UUID, state) -> None:
|
||
call = state.phone_pending
|
||
assert call is not None
|
||
text = _report_text(call.phase, call.crew, state.dispatched_card.address or "из карточки")
|
||
_line(session_id, state, "crew", text)
|
||
report = PhoneReportRecord(service=call.service, crew=call.crew,
|
||
phase=call.phase, text=text, at=now_utc())
|
||
state.phone_reports.append(report)
|
||
state.dds_log.append(("phone.report", now_utc(), f"{call.crew}: {call.phase}"))
|
||
state.phone_pending = None
|
||
hub.to_station(session_id, PhoneReport(**report.model_dump()))
|
||
|
||
|
||
async def _finish_dds(session_id: UUID, state) -> None:
|
||
state.ended_at = now_utc()
|
||
state.end_reason = CallEndReason.COMPLETE
|
||
state.capture_active_dds()
|
||
for card in state.dds_live_cards:
|
||
timer = card.timers.timers.get(TimerCode.DDS_WORK)
|
||
if timer is not None and timer.started_at is not None:
|
||
card.timers.on_event("dds.finish")
|
||
hub.stop_ticker(session_id)
|
||
hub.to_station(session_id, SessionEnded(reason=CallEndReason.COMPLETE))
|
||
hub.to_observers(session_id, SessionEnded(reason=CallEndReason.COMPLETE))
|
||
if hub.journal:
|
||
await hub.journal.session_ended(session_id, state.ended_at, CallEndReason.COMPLETE.value)
|
||
await finish(session_id, state)
|
||
hub.to_station(session_id, ScoreReady(session_id=session_id))
|
||
|
||
|
||
async def _handle(session_id: UUID, state, event) -> None:
|
||
if state.ended:
|
||
_error(session_id, "Занятие уже завершено")
|
||
return
|
||
match event.type:
|
||
case "card.ack":
|
||
# Подтверждение приёма — это статус «Принята» у главной службы.
|
||
# Кнопка осталась ради живой цепочки 112 → ДДС (lct-20), где
|
||
# диспетчер один и выбирать службу не из чего.
|
||
if not event.comment.strip():
|
||
_error(session_id, "Для подтверждения приёма добавьте комментарий с основанием")
|
||
return
|
||
if any(action == "card.ack" for action, _at, _detail in state.dds_log):
|
||
return
|
||
state.on_event("card.ack")
|
||
state.dds_log.append(("card.ack", now_utc(), None))
|
||
services = state.managed_services()
|
||
if services:
|
||
_start_dds_work_timer(state)
|
||
try:
|
||
state.set_service_status(
|
||
services[0], ServiceStatus.ACCEPTED, event.comment, author="диспетчер"
|
||
)
|
||
except StatusError:
|
||
pass # статус уже стоит: повторное нажатие ничего не меняет
|
||
case "card.status":
|
||
if event.service not in state.managed_services():
|
||
_error(session_id, "Можно менять статусы только своей ДДС")
|
||
return
|
||
_start_dds_work_timer(state)
|
||
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")
|
||
if event.status in {
|
||
ServiceStatus.COMPLETED, ServiceStatus.DECLINED, ServiceStatus.REFUSED,
|
||
}:
|
||
state.on_event("dds.complete")
|
||
case "crew.select":
|
||
if event.crew not in state.crew_options():
|
||
_error(session_id, "Выберите бригаду из списка доступных")
|
||
return
|
||
service = state.crew_service(event.crew)
|
||
assigned = state.crew_assignments.get(service)
|
||
if assigned and assigned != event.crew and any(
|
||
report.service == service for report in state.phone_reports
|
||
):
|
||
_error(session_id, "После первого доклада бригаду этой службы менять нельзя")
|
||
return
|
||
if state.phone_pending is not None:
|
||
_error(session_id, "Завершите текущий разговор перед сменой бригады")
|
||
return
|
||
if state.crew_selected == event.crew and assigned == event.crew:
|
||
return
|
||
state.crew_selected = event.crew
|
||
state.crew_assignments[service] = event.crew
|
||
state.dds_log.append(("crew.select", now_utc(), event.crew))
|
||
case "phone.dial":
|
||
if state.phone_pending is not None:
|
||
_error(session_id, "Разговор уже идёт: передайте сведения или завершите звонок")
|
||
return
|
||
crew = state.crew_selected
|
||
service = state.crew_service(crew) if crew else None
|
||
if service is None:
|
||
_error(session_id, "Сначала выберите бригаду")
|
||
return
|
||
previous = [report for report in state.phone_reports if report.service == service]
|
||
if len(previous) >= len(REPORT_PHASES):
|
||
_error(session_id, "Все доклады этой бригады уже получены")
|
||
return
|
||
phase = REPORT_PHASES[len(previous)]
|
||
required = REQUIRED_STATUS[phase]
|
||
if current(state.status_log, service) not in STATUS_AT_OR_AFTER[phase]:
|
||
_error(session_id, f"Перед звонком нужен статус «{SERVICE_STATUS_LABELS[required]}» этой службы")
|
||
return
|
||
state.phone_pending = PhoneCallPending(service=service, crew=crew, phase=phase)
|
||
state.dds_log.append(("phone.dial", now_utc(), crew))
|
||
greeting = (f"{crew}, старший группы на связи. Назовите адрес, характер происшествия "
|
||
"и что требуется от бригады."
|
||
if phase == "dispatched" else
|
||
f"{crew}, старший группы на связи. Слушаю ваш запрос по карточке.")
|
||
_line(session_id, state, "crew", greeting)
|
||
case "phone.brief":
|
||
call = state.phone_pending
|
||
if call is None or call.phase != "dispatched":
|
||
_error(session_id, "Сначала соединитесь со старшим группы для передачи вызова")
|
||
return
|
||
if not _address_matches(state.dispatched_card.address, event.address):
|
||
_error(session_id, "Проверьте адрес: улица и номер дома должны совпадать с карточкой")
|
||
return
|
||
if not _incident_matches(state, event.incident):
|
||
_error(session_id, "Уточните характер происшествия по данным карточки")
|
||
return
|
||
if not _has_purpose(event.request, ("выезд", "выех", "направ", "прибыт",
|
||
"реагир", "помощ", "подтверд", "долож")):
|
||
_error(session_id, "Сформулируйте задачу: выезд, помощь или доклад бригады")
|
||
return
|
||
_line(session_id, state, "dispatcher", f"Адрес: {event.address.strip()}. "
|
||
f"Происшествие: {event.incident.strip()}. {event.request.strip()}")
|
||
_finish_phone_call(session_id, state)
|
||
case "phone.check":
|
||
call = state.phone_pending
|
||
if call is None or call.phase == "dispatched":
|
||
_error(session_id, "Сначала передайте вызов, затем запросите обстановку")
|
||
return
|
||
if not _has_purpose(event.text, ("обстанов", "статус", "прибыл", "доех",
|
||
"выех", "работ", "заверш", "мест",
|
||
"ход", "сообщ", "долож", "уточн")):
|
||
_error(session_id, "Спросите обстановку, прибытие или ход работ по карточке")
|
||
return
|
||
_line(session_id, state, "dispatcher", event.text.strip())
|
||
_finish_phone_call(session_id, state)
|
||
case "phone.hangup":
|
||
if state.phone_pending is None:
|
||
_error(session_id, "Нет активного разговора")
|
||
return
|
||
state.dds_log.append(("phone.hangup", now_utc(), state.phone_pending.crew))
|
||
state.phone_pending = None
|
||
case "card.reply":
|
||
if (state.exercise is not Exercise.DDS and not state.handoff_to_dds) or not state.dispatched_card or (
|
||
event.card_id != state.dispatched_card.card_id
|
||
):
|
||
_error(session_id, "Ответ относится не к текущей карточке")
|
||
return
|
||
# A browser may lose the acknowledgement after the server has
|
||
# committed this replace-style value. Reconnect retries are safe:
|
||
# don't create another journal row
|
||
# when the current card already contains exactly this text.
|
||
if state.reply_text == event.text:
|
||
return
|
||
state.reply_text = event.text
|
||
state.reply_log.append((now_utc(), event.text))
|
||
case "card.open":
|
||
if (state.exercise is not Exercise.DDS and not state.handoff_to_dds) or not state.activate_dds_card(event.card_id):
|
||
_error(session_id, "Карточка отсутствует в текущей очереди")
|
||
return
|
||
_start_dds_work_timer(state)
|
||
hub.to_station(session_id, state.card_received_event())
|
||
# CardReceived carries the contents, while StationState carries
|
||
# the status journal and current queue. Send both on every switch
|
||
# so the newly opened card cannot briefly inherit the previous
|
||
# card's status snapshot until the next periodic tick.
|
||
hub.to_station(session_id, StationState(snapshot=state.station_snapshot()))
|
||
case "card.next":
|
||
if (state.exercise is not Exercise.DDS and not state.handoff_to_dds) or not state.dispatched_card or (
|
||
event.card_id != state.dispatched_card.card_id
|
||
):
|
||
_error(session_id, "Следующая карточка недоступна: ID текущей не совпадает")
|
||
return
|
||
if any(item.card_id == event.card_id for item in state.dds_completed):
|
||
_error(session_id, "Эта карточка уже завершена")
|
||
return
|
||
state.dds_completed.append(score_current_dds(state))
|
||
completed_index = state.dds_card_index
|
||
state.dds_live_cards = [
|
||
item for item in state.dds_live_cards if item.card_id != event.card_id
|
||
]
|
||
remaining = sorted(state.dds_live_cards, key=lambda item: item.original_index)
|
||
if remaining:
|
||
following = next(
|
||
(item for item in remaining if item.original_index > completed_index),
|
||
remaining[0],
|
||
)
|
||
state.activate_dds_card(following.card_id, capture=False)
|
||
hub.to_station(session_id, state.card_received_event())
|
||
else:
|
||
# Keep the lesson alive if selected cards have not arrived yet.
|
||
# The next delivery may become the active card immediately or
|
||
# after its configured interval; no completed card is reused.
|
||
state.dds_active_card_id = None
|
||
state.dispatched_card = None
|
||
state.dispatched_at = None
|
||
state.dds_card_index = completed_index
|
||
active_before_delivery = state.dds_active_card_id
|
||
deliver_due_cards(state)
|
||
if state.dds_active_card_id and state.dds_active_card_id != active_before_delivery:
|
||
hub.to_station(session_id, state.card_received_event())
|
||
hub.to_station(session_id, StationState(snapshot=state.station_snapshot()))
|
||
if not state.dds_live_cards and len(state.dds_completed) >= len(state.dds_scenarios):
|
||
await _finish_dds(session_id, state)
|
||
case "station.finish":
|
||
if state.exercise is not Exercise.DDS and not state.handoff_to_dds:
|
||
_error(session_id, "Операторское занятие завершается после звонка 112")
|
||
return
|
||
await _finish_dds(session_id, state)
|
||
case "card.bounce":
|
||
_error(
|
||
session_id,
|
||
"ДДС не проверяет заполнение карточки: замечания передаёт служба контроля 112",
|
||
)
|
||
return
|
||
case "zone.decision":
|
||
previous = next(
|
||
(detail for action, _at, detail in reversed(state.dds_log)
|
||
if action == "zone.decision"),
|
||
None,
|
||
)
|
||
decision = "в зоне" if event.in_zone else "не в зоне"
|
||
if previous is not None:
|
||
if previous != decision:
|
||
_error(session_id, "Решение по зоне уже записано для этой карточки")
|
||
return
|
||
state.on_event("zone.decision")
|
||
state.dds_log.append(("zone.decision", now_utc(), decision))
|
||
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))
|
||
state.capture_active_dds()
|
||
await hub.checkpoint(session_id)
|
||
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())
|
||
if (isinstance(event, ErrorEvent) and event.code is ErrorKind.INTERNAL
|
||
and event.message == LEASE_FENCED_MESSAGE):
|
||
await ws.close(code=1012)
|
||
return
|
||
|
||
|
||
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:
|
||
if not websocket_origin_allowed(ws):
|
||
await ws.close(code=1008)
|
||
return
|
||
if hub.is_lease_fenced(session_id):
|
||
await ws.close(code=1012)
|
||
return
|
||
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
|
||
if who.role is Role.INSTRUCTOR and state.owner_login != who.login:
|
||
await ws.send_text(
|
||
ErrorEvent(code=ErrorKind.SESSION_NOT_FOUND, message="Занятие не запущено").model_dump_json()
|
||
)
|
||
await ws.close()
|
||
return
|
||
if who.role is Role.TRAINEE and (
|
||
state.trainee_id is None or state.trainee_id != who.trainee_id
|
||
):
|
||
await _reject(ws, "Занятие не назначено этому обучающемуся")
|
||
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
|
||
raw_command_id = payload.get("_command_id") if isinstance(payload, dict) else None
|
||
try:
|
||
command_id = UUID(raw_command_id) if raw_command_id is not None else None
|
||
except (ValueError, TypeError, AttributeError):
|
||
hub.to_station(
|
||
session_id,
|
||
ErrorEvent(code=ErrorKind.UNSUPPORTED_EVENT,
|
||
message="Некорректный идентификатор команды"),
|
||
)
|
||
continue
|
||
if command_id is not None and str(command_id) in state.processed_station_commands:
|
||
# The checkpoint already proves this exact command committed.
|
||
# Re-ack it without rerunning its business transition.
|
||
hub.to_station(session_id, CommandAck(command_id=command_id))
|
||
continue
|
||
async with hub.durable_transition(session_id):
|
||
await _handle(session_id, state, event)
|
||
if command_id is not None:
|
||
state.processed_station_commands.append(str(command_id))
|
||
del state.processed_station_commands[:-512]
|
||
# Commit the state+dedupe ID before acknowledging. The
|
||
# transition context can have already flushed other
|
||
# events; an explicit checkpoint here makes the
|
||
# command/ACK boundary independent of that batch state.
|
||
await hub.checkpoint(session_id)
|
||
# The hub stages non-error events until the checkpoint
|
||
# transaction has committed, including this ack.
|
||
hub.to_station(session_id, CommandAck(command_id=command_id))
|
||
except WebSocketDisconnect:
|
||
return
|
||
except Exception: # noqa: BLE001 — failed durable transition may fence the owner
|
||
if not state.lease_fenced:
|
||
raise
|
||
log.info(
|
||
"закрытие станционного WebSocket после fencing занятия %s",
|
||
session_id,
|
||
)
|
||
# `hub.checkpoint` broadcasts a structured fence event before
|
||
# propagating the failed database write. Let the sender deliver
|
||
# that event and close with 1012 instead of an opaque 1006.
|
||
await asyncio.gather(sender, return_exceptions=True)
|
||
return
|
||
finally:
|
||
sender.cancel()
|