Implement local task 09 training workflows

This commit is contained in:
andreysk0304 2026-09-21 17:40:54 +03:00
commit cec84ffcd0
45 changed files with 2679 additions and 366 deletions

View file

@ -10,15 +10,116 @@
`hint.shown` из живой сессии, эталонные вопросы — только в разборе.
"""
from fastapi import APIRouter, HTTPException
from typing import Any
from fastapi import APIRouter, Depends, HTTPException, Request
from pydantic import BaseModel, Field
from sqlalchemy.ext.asyncio import AsyncSession
from app.api.auth import audit, require
from app.db.base import get_session
from app.domain.roles import Role
from app.scenarios import store
from app.scenarios.editor import validate
from app.scenarios.loader import ScenarioError
router = APIRouter(prefix="/api/scenarios", tags=["scenarios"])
HIDDEN_FROM_TRAINEE = {"facts", "ground_truth", "tree", "checklist"}
class TemplateDraftIn(BaseModel):
source_id: str = Field(min_length=1)
title: str | None = Field(default=None, min_length=1, max_length=200)
def _draft_out(row) -> dict:
return {
"id": row.id,
"status": row.status,
"generation": "template_copy",
"body": row.body,
}
@router.post("/drafts/from-template", status_code=201)
async def create_template_draft(
body: TemplateDraftIn, request: Request, db: AsyncSession = Depends(get_session)
) -> dict:
who = require(request, Role.INSTRUCTOR)
source = store.get(body.source_id)
if source is None:
raise HTTPException(status_code=404, detail="published_source_not_found")
row = await store.create_draft(db, source=source, title=body.title)
await audit(who.login, who.role.value, "scenario.draft.create", row.id, f"template:{source.id}")
return _draft_out(row)
@router.get("/drafts/{scenario_id}")
async def read_draft(
scenario_id: str, request: Request, db: AsyncSession = Depends(get_session)
) -> dict:
require(request, Role.INSTRUCTOR)
row = await store.draft(db, scenario_id)
if row is None:
raise HTTPException(status_code=404, detail="draft_not_found")
return _draft_out(row)
@router.patch("/drafts/{scenario_id}")
async def patch_draft(
scenario_id: str,
body: dict[str, Any],
request: Request,
db: AsyncSession = Depends(get_session),
) -> dict:
who = require(request, Role.INSTRUCTOR)
row = await store.draft(db, scenario_id)
if row is None:
raise HTTPException(status_code=404, detail="draft_not_found")
try:
row = await store.update_draft(db, row, body)
except ScenarioError as exc:
raise HTTPException(status_code=422, detail=str(exc)) from exc
await audit(who.login, who.role.value, "scenario.draft.update", row.id)
return _draft_out(row)
@router.post("/drafts/{scenario_id}/validate")
async def validate_draft(
scenario_id: str, request: Request, db: AsyncSession = Depends(get_session)
) -> dict:
require(request, Role.INSTRUCTOR)
row = await store.draft(db, scenario_id)
if row is None:
raise HTTPException(status_code=404, detail="draft_not_found")
try:
scenario = validate(row.body)
except ScenarioError as exc:
return {"valid": False, "errors": [str(exc)]}
return {
"valid": True,
"errors": [],
"ground_truth": scenario.ground_truth.model_dump(mode="json"),
}
@router.post("/drafts/{scenario_id}/approve")
async def approve_draft(
scenario_id: str, request: Request, db: AsyncSession = Depends(get_session)
) -> dict:
who = require(request, Role.INSTRUCTOR)
row = await store.draft(db, scenario_id)
if row is None:
raise HTTPException(status_code=404, detail="draft_not_found")
try:
scenario = await store.approve_draft(db, row)
except ScenarioError as exc:
raise HTTPException(status_code=422, detail=str(exc)) from exc
await audit(who.login, who.role.value, "scenario.approve", scenario.id)
return {"id": scenario.id, "status": "published", "title": scenario.title}
@router.get("")
async def listing() -> list[dict]:
return [

View file

@ -7,16 +7,17 @@
from datetime import datetime
from uuid import UUID
from fastapi import APIRouter, Depends, HTTPException, Query, Request
from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response
from pydantic import BaseModel
from sqlalchemy.ext.asyncio import AsyncSession
from app.api.auth import audit, require
from app.db import repo
from app.db.base import get_session
from app.domain.events import SessionMode, SessionReport
from app.domain.events import Exercise, SessionMode, SessionReport
from app.scenarios import store
from app.scoring.report import build as build_report
from app.scoring.export import to_csv, to_pdf
from app.domain.roles import Role
from app.session.hub import hub
@ -57,7 +58,8 @@ def _out(session) -> SessionOut:
@router.post("", response_model=SessionOut, status_code=201)
async def create(body: SessionCreate, db: AsyncSession = Depends(get_session)) -> SessionOut:
async def create(body: SessionCreate, request: Request, db: AsyncSession = Depends(get_session)) -> SessionOut:
require(request, Role.INSTRUCTOR)
group = await repo.ensure_group(db, body.group) if body.group else None
trainee = await repo.ensure_trainee(db, body.trainee, group) if body.trainee else None
session = await repo.create_session(
@ -71,10 +73,13 @@ async def create(body: SessionCreate, db: AsyncSession = Depends(get_session)) -
@router.get("/{session_id}", response_model=SessionOut)
async def read(session_id: UUID, db: AsyncSession = Depends(get_session)) -> SessionOut:
async def read(session_id: UUID, request: Request, db: AsyncSession = Depends(get_session)) -> SessionOut:
who = require(request)
session = await repo.get_session(db, session_id)
if session is None:
raise HTTPException(status_code=404, detail="session_not_found")
if who.role is Role.TRAINEE and session.trainee_id != who.trainee_id:
raise HTTPException(status_code=403, detail="not_your_session")
return _out(session)
@ -84,16 +89,19 @@ class ChecklistItemOut(BaseModel):
@router.get("/{session_id}/checklist", response_model=list[ChecklistItemOut])
async def checklist(session_id: UUID) -> list[ChecklistItemOut]:
async def checklist(session_id: UUID, request: Request) -> list[ChecklistItemOut]:
"""Чек-лист для самооценки — **только после конца звонка**.
Во время звонка это содержимое подсказок: отдать его значит выдать
в контрольном режиме то, чего там быть не должно. После звонка курсант
по нему отмечает, что, по его мнению, пропустил.
"""
who = require(request)
state = hub.get(session_id)
if state is None:
raise HTTPException(status_code=404, detail="session_not_found")
if who.role is Role.TRAINEE and state.trainee_id != who.trainee_id:
raise HTTPException(status_code=403, detail="not_your_session")
if not state.ended:
raise HTTPException(status_code=409, detail="call_not_ended")
scenario = state.scenario or store.get(state.scenario_id)
@ -136,11 +144,37 @@ async def report(session_id: UUID, request: Request) -> SessionReport:
state, scenario = _live(session_id)
if who.role is Role.TRAINEE and state.trainee_id != who.trainee_id:
raise HTTPException(status_code=403, detail="not_your_session")
if who.role is Role.TRAINEE and state.exercise is Exercise.CALL and not state.self_assessed:
raise HTTPException(status_code=409, detail="self_assessment_required")
if state.score is None:
raise HTTPException(status_code=409, detail="score_not_ready")
return build_report(session_id, state, scenario)
@router.get("/{session_id}/report.csv")
async def report_csv(session_id: UUID, request: Request) -> Response:
"""Те же права и готовность оценки, что у JSON-разбора."""
data = await report(session_id, request)
return Response(
content=to_csv(data), media_type="text/csv; charset=utf-8",
headers={"Content-Disposition": f'attachment; filename="session-{session_id}-report.csv"'},
)
@router.get("/{session_id}/report.pdf")
async def report_pdf(session_id: UUID, request: Request) -> Response:
"""Печатный разбор; генерация полностью локальна."""
data = await report(session_id, request)
try:
content = to_pdf(data)
except RuntimeError as exc:
raise HTTPException(status_code=503, detail=str(exc)) from exc
return Response(
content=content, media_type="application/pdf",
headers={"Content-Disposition": f'attachment; filename="session-{session_id}-report.pdf"'},
)
@router.patch("/{session_id}/report", response_model=SessionReport)
async def override(session_id: UUID, body: ScoreOverride, request: Request) -> SessionReport:
"""Тренажёр готовит материал, преподаватель имеет последнее слово.
@ -177,6 +211,8 @@ async def listing(
who = require(request)
# Обучающийся видит только свою историю, что бы он ни передал в фильтре.
if who.role is Role.TRAINEE:
if who.trainee_id is None:
raise HTTPException(status_code=403, detail="trainee_profile_required")
trainee = who.trainee_id
rows = await repo.history(
db,

View file

@ -14,13 +14,18 @@ from pydantic import TypeAdapter, ValidationError
from app.domain.events import (
StationState,
CallIncoming,
CallEnded,
CallEndReason,
CallStarted,
ErrorEvent,
ErrorKind,
Exercise,
HintShown,
KioState,
KioPatchOut,
PatchSource,
ScoreReady,
SessionEnded,
SessionMode,
TimerTick,
@ -33,6 +38,7 @@ 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
from app.domain.kio import ResponseStatus
from app.voice.models import get_voice_models
from app.voice.pipeline import VoiceSession
@ -82,7 +88,37 @@ def _next_hint(state) -> tuple[str, str] | None:
async def _handle(session_id: UUID, state, event) -> None:
if state.ended and event.type != "self_assessment.submit":
hub.to_trainee(session_id, ErrorEvent(
code=ErrorKind.UNSUPPORTED_EVENT, message="Занятие уже завершено",
))
return
if state.exercise is Exercise.DDS or (
state.exercise is Exercise.CARD and event.type not in {"kio.patch", "card.submit"}
) or (state.exercise is Exercise.CALL and event.type == "card.submit"):
hub.to_trainee(session_id, ErrorEvent(
code=ErrorKind.UNSUPPORTED_EVENT, message="Действие недоступно в этом упражнении",
))
return
match event.type:
case "card.submit":
state.kio.registered_at = state.started_at or now_utc()
state.kio.response_status = ResponseStatus.TRANSFERRED
state.dispatched_card = state.kio.model_copy(deep=True)
state.dispatched_at = now_utc()
state.ended_at = state.dispatched_at
state.end_reason = CallEndReason.COMPLETE
hub.stop_ticker(session_id)
hub.to_station(session_id, state.card_received_event())
hub.to_station(session_id, StationState(snapshot=state.station_snapshot()))
hub.to_observers(session_id, KioState(kio=state.kio))
hub.to_trainee(session_id, CallEnded(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)
case "call.answer":
state.on_event("call.answer")
state.started_at = now_utc()
@ -93,7 +129,17 @@ async def _handle(session_id: UUID, state, event) -> None:
_start_voice(session_id, state)
case "kio.patch":
old_code, old_notify = state.kio.incident_code, list(state.kio.notify)
state.patch_kio(event.fields)
hub.to_trainee(session_id, KioPatchOut(
fields=event.fields, source=PatchSource.OPERATOR,
))
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,
))
# Наблюдателю уходит карточка целиком: рассинхрон на внешнем мониторе
# посреди занятия дороже лишних килобайт.
hub.to_observers(session_id, KioState(kio=state.kio))
@ -121,8 +167,19 @@ async def _handle(session_id: UUID, state, event) -> None:
await hub.journal.hint(session_id, checklist_id, question, now_utc())
case "dds.dispatch":
if event.service is None and not (state.kio.incident_code or state.kio.notify):
hub.to_trainee(session_id, ErrorEvent(
code=ErrorKind.UNSUPPORTED_EVENT,
message="Укажите ДДС или проставьте признаки для маршрутизации по ЕКП",
))
return
state.on_event("dds.dispatch")
state.dispatch(event.service.value)
state.dispatch(event.service.value if event.service else None)
hub.to_trainee(session_id, KioPatchOut(
fields={"response_status": state.kio.response_status.value,
"dds": state.kio.dds.value if state.kio.dds else None},
source=PatchSource.AUTO,
))
# Карточка замораживается снимком и уходит диспетчеру: оператор
# не должен иметь возможности дописать поле задним числом.
hub.to_station(session_id, state.card_received_event())
@ -227,8 +284,40 @@ async def call(ws: WebSocket, session_id: UUID) -> None:
)
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.trainee(session_id) as queue:
if state.exercise is Exercise.CARD:
from app.api.ws.control import card_briefing
hub.to_trainee(session_id, card_briefing(state))
if state.ended:
hub.to_trainee(session_id, CallEnded(reason=state.end_reason or CallEndReason.COMPLETE))
if state.score is not None:
hub.to_trainee(session_id, ScoreReady(session_id=session_id))
elif state.exercise is Exercise.CALL:
scenario = state.scenario or store.get(state.scenario_id)
required = ([field for field in state.required_fields if field != "dds"]
if scenario and scenario.ground_truth.incident_code
else list(state.required_fields))
hub.to_trainee(session_id, CallIncoming(
scenario_id=state.scenario_id, caller_number="+7 (495) 000-00-00",
level=state.level, mode=state.mode, required_fields=required,
))
hub.to_trainee(session_id, KioPatchOut(
fields=state.kio.model_dump(mode="json"), source=PatchSource.OPERATOR,
))
if state.started_at is not None:
hub.to_trainee(session_id, CallStarted(started_at=state.started_at))
if state.ended:
hub.to_trainee(session_id, CallEnded(reason=state.end_reason or CallEndReason.HANGUP))
if state.score is not None and state.self_assessed:
hub.to_trainee(session_id, ScoreReady(session_id=session_id))
hub.to_trainee(session_id, TimerTick(timers=state.timers.snapshot()))
writer = asyncio.create_task(_pump(ws, queue))
try:
while True:
@ -238,7 +327,13 @@ async def call(ws: WebSocket, session_id: UUID) -> None:
# Бинарные кадры — аудио, текстовые — события. Направление определяется
# каналом, обёртки JSON вокруг звука нет (docs/arch/CONTRACT.md).
if message.get("bytes") is not None:
_on_audio(session_id, state, message["bytes"])
if state.exercise is Exercise.CALL:
_on_audio(session_id, state, message["bytes"])
else:
hub.to_trainee(session_id, ErrorEvent(
code=ErrorKind.UNSUPPORTED_EVENT,
message="Аудио не используется в этом упражнении",
))
continue
try:
payload = json.loads(message.get("text") or "")

View file

@ -9,6 +9,7 @@
"""
import logging
import re
from uuid import UUID, uuid4
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
@ -21,8 +22,9 @@ from app.domain.events import (
ErrorKind,
ScoreReady,
CallIncoming,
ErrorEvent,
ErrorKind,
CardBriefing,
Exercise,
StationState,
InstructorToServer,
InstructorNoteShown,
ModeSet,
@ -42,6 +44,7 @@ from app.api.auth import audit, principal_of
from app.domain.roles import Role
from app.session.hub import hub
from app.session.state import SessionState, now_utc
from app.domain.kio import KIO, ResponseStatus
from app.voice.models import get_voice_models
from app.voice.pipeline import FILLERS, prefetch
@ -51,6 +54,27 @@ router = APIRouter()
_adapter = TypeAdapter(InstructorToServer)
def card_briefing(state: SessionState) -> CardBriefing:
"""Учебная текстовая вводная — исходные реплики, а не эталон карточки.
В отсутствие диалога факты раскрываются сразу. Если факт уточняется,
показываем и уточнение: иначе правильно заполнить карточку невозможно.
"""
scenario = state.scenario
lines = ["Учебная текстовая вводная: сведения заявителя приведены ниже.",
scenario.first_line]
for fact in scenario.facts:
lines.append(f"• {fact.value}")
if fact.refined:
lines.append(f" Уточнено: {fact.refined}")
return CardBriefing(
scenario_id=scenario.id, mode=state.mode, text="\n".join(lines),
required_fields=([field for field in state.required_fields if field != "dds"]
if scenario.ground_truth.incident_code else list(state.required_fields)),
card=state.kio,
)
async def _start(session_id: UUID, event, who=None) -> None:
scenario = store.get(event.scenario_id)
if scenario is None:
@ -61,9 +85,10 @@ async def _start(session_id: UUID, event, who=None) -> None:
return
attempt = 1
recorded_trainee_id = event.trainee_id
if hub.journal:
attempt = await hub.journal.start_lesson(
session_id, scenario.id, event.mode.value, event.trainee
attempt, recorded_trainee_id = await hub.journal.start_lesson(
session_id, scenario.id, event.mode.value, event.trainee, event.trainee_id
)
# Занятие собирается целиком и только потом регистрируется: иначе
@ -75,25 +100,59 @@ async def _start(session_id: UUID, event, who=None) -> None:
scenario_title=scenario.title,
level=scenario.level.value,
mode=event.mode,
exercise=event.exercise,
scenario=scenario.model_copy(deep=True),
required_fields=scenario.required_fields,
trainee_name=event.trainee,
trainee_id=getattr(event, "trainee_id", None),
trainee_id=recorded_trainee_id,
attempt=attempt,
)
embedder = get_embedder()
if embedder is not None:
state.slots = SlotMachine(state.scenario, embedder)
state.persona = PersonaState(state.scenario.persona)
state.caller = build_caller(scenario.id)
if event.exercise is Exercise.CALL:
embedder = get_embedder()
if embedder is not None:
state.slots = SlotMachine(state.scenario, embedder)
state.persona = PersonaState(state.scenario.persona)
state.caller = build_caller(scenario.id)
elif event.exercise is Exercise.DDS:
# В ДДС поступает уже оформленная учебная карточка. Содержимое берётся
# из утверждённого сценария, а не из действий несуществующего оператора.
truth = scenario.ground_truth
address_fact = next((fact.value for fact in scenario.facts if "address" in fact.id), "")
floor = re.search(r"(\d+)[-‑–]?й?\s*этаж", address_fact, re.IGNORECASE)
service = truth.dds.value if truth.dds else None
fallback = {"01": "Служба 101", "02": "МВД", "03": "Скорая помощь", "04": "Аварийная служба"}
state.kio = KIO(
registered_at=now_utc(), caller_number="+7 (495) 000-00-00",
address=truth.address or address_fact or None,
floor=floor.group(1) if floor else None,
incident_type=truth.incident_type, incident_code=truth.incident_code,
dds=truth.dds, signs=list(scenario.signs),
notify=list(truth.notify) or ([fallback[service]] if service in fallback else []),
victims_count=truth.victims,
description="; ".join(fact.value for fact in scenario.facts[:3]) or scenario.first_line,
)
if service:
state.dispatch(service)
else:
state.kio.response_status = ResponseStatus.TRANSFERRED
state.dispatched_card = state.kio.model_copy(deep=True)
state.dispatched_at = now_utc()
state.started_at = state.dispatched_at
else:
state.started_at = now_utc()
hub.register(state)
if event.exercise is not Exercise.CALL and hub.journal and state.started_at is not None:
await hub.journal.session_started(session_id, state.started_at)
# Первая реплика и филлеры синтезируются, пока курсант не снял трубку:
# «Алло! Помогите!» должно прозвучать мгновенно (docs/arch/BACKEND.md).
models = get_voice_models()
if models is not None:
asyncio.create_task(prefetch(models, [scenario.first_line, *FILLERS.values()]))
state.on_event("call.incoming")
if event.exercise is Exercise.CALL:
models = get_voice_models()
if models is not None:
asyncio.create_task(prefetch(models, [scenario.first_line, *FILLERS.values()]))
state.on_event("call.incoming")
elif event.exercise is Exercise.DDS:
state.on_event("dds.dispatch")
hub.start_ticker(session_id)
if who is not None:
# Запуск занятия меняет чужой результат — значит попадает в журнал
@ -103,29 +162,45 @@ async def _start(session_id: UUID, event, who=None) -> None:
f"{scenario.id}, режим {event.mode.value}, курсант {event.trainee}")
)
hub.to_trainee(
session_id,
CallIncoming(
scenario_id=scenario.id,
caller_number="+7 (495) 000-00-00",
level=scenario.level,
mode=event.mode,
scenario=scenario.model_copy(deep=True),
required_fields=scenario.required_fields,
),
)
if event.exercise is Exercise.CALL:
hub.to_trainee(
session_id,
CallIncoming(
scenario_id=scenario.id,
caller_number="+7 (495) 000-00-00",
level=scenario.level,
mode=event.mode,
required_fields=[field for field in scenario.required_fields
if field != "dds" or not scenario.ground_truth.incident_code],
),
)
elif event.exercise is Exercise.DDS:
hub.to_station(session_id, state.card_received_event())
hub.to_station(session_id, StationState(snapshot=state.station_snapshot()))
else:
hub.to_trainee(session_id, card_briefing(state))
hub.to_observers(session_id, ModeSet(mode=event.mode))
hub.to_observers(session_id, state.snapshot())
async def _stop(session_id: UUID) -> None:
state = hub.get(session_id)
if state is None:
if state is None or state.ended:
return
state.ended_at = now_utc()
state.end_reason = CallEndReason.INSTRUCTOR
hub.stop_ticker(session_id)
hub.to_observers(session_id, SessionEnded(reason=CallEndReason.INSTRUCTOR))
if state.exercise in {Exercise.DDS, Exercise.CARD}:
from app.session.finish import finish
if state.exercise is Exercise.DDS:
hub.to_station(session_id, SessionEnded(reason=CallEndReason.INSTRUCTOR))
else:
hub.to_trainee(session_id, CallEnded(reason=CallEndReason.INSTRUCTOR))
await finish(session_id, state)
if state.exercise is Exercise.DDS:
hub.to_station(session_id, ScoreReady(session_id=session_id))
if hub.journal:
await hub.journal.session_ended(session_id, state.ended_at, CallEndReason.INSTRUCTOR.value)
@ -204,6 +279,12 @@ async def control(ws: WebSocket, session_id: UUID) -> None:
state = hub.get(session_id)
if state is None:
continue
if state.exercise is not Exercise.CALL:
hub.to_observers(session_id, ErrorEvent(
code=ErrorKind.UNSUPPORTED_EVENT,
message="Директивы звонящему доступны только в голосовом упражнении",
))
continue
result = apply_directive(state, event.directive)
if result.needs_network:
hub.to_observers(session_id, ErrorEvent(

View file

@ -1,7 +1,7 @@
"""Сокет диспетчера ДДС.
**Только JSON: аудио здесь нет вообще** — значит, нет ни VAD, ни распознавания,
ни синтеза. Голосовой контур станция не трогает (docs/arch/CONTRACT.md).
Учебный исходящий звонок бригаде имитируется событиями JSON. Для целевого
локального прототипа не нужна внешняя телефония или аудиомодель.
Самая ценная механика цепочки — `card.bounce`: диспетчер видит, что не указан
этаж, и отбивает карточку обратно. Неполнота КИО перестаёт быть процентом
@ -15,20 +15,45 @@ 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.domain.events import (
CallEndReason, ErrorEvent, ErrorKind, Exercise, PhoneReport, ScoreReady,
SessionEnded, StationState, StationToServer,
)
from app.domain.statuses import PRIMARY, PhoneReportRecord, ServiceStatus, StatusError, current
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
from app.session.finish import finish
log = logging.getLogger(__name__)
router = APIRouter()
_adapter = TypeAdapter(StationToServer)
REPORT_PHASES = ("dispatched", "arrived", "working", "completed")
REPORT_TEXT = {
"dispatched": "Бригада выехала к месту происшествия.",
"arrived": "Бригада прибыла на место происшествия.",
"working": "Бригада приступила к проведению работ.",
"completed": "Работы завершены, бригада освобождена.",
}
REPORT_FOR_STATUS = {
ServiceStatus.RESPONDING: "dispatched",
ServiceStatus.ARRIVED: "arrived",
ServiceStatus.WORKING: "working",
ServiceStatus.COMPLETED: "completed",
}
def _error(session_id: UUID, message: str) -> None:
hub.to_station(session_id, ErrorEvent(code=ErrorKind.UNSUPPORTED_EVENT, message=message))
async def _handle(session_id: UUID, state, event) -> None:
if state.ended:
_error(session_id, "Занятие уже завершено")
return
match event.type:
case "card.ack":
# Подтверждение приёма — это статус «Принята» у главной службы.
@ -43,6 +68,16 @@ async def _handle(session_id: UUID, state, event) -> None:
except StatusError:
pass # статус уже стоит: повторное нажатие ничего не меняет
case "card.status":
if event.service not in state.notified_services():
_error(session_id, "Служба отсутствует в списке оповещения карточки")
return
phase = REPORT_FOR_STATUS.get(event.status)
if state.exercise is Exercise.DDS and phase is not None and not any(
report.service == event.service and report.phase == phase
for report in state.phone_reports
):
_error(session_id, f"Статус «{event.status.value}» требует доклада бригады")
return
try:
state.set_service_status(
event.service, event.status, event.comment, author="диспетчер"
@ -56,6 +91,55 @@ async def _handle(session_id: UUID, state, event) -> None:
# Первичный статус останавливает норматив 30 секунд.
if event.status in PRIMARY:
state.on_event("card.ack")
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
state.crew_selected = event.crew
state.crew_assignments[service] = event.crew
case "phone.dial":
crew = state.crew_selected
service = state.crew_service(crew) if crew else None
if service is None:
_error(session_id, "Сначала выберите бригаду")
return
if current(state.status_log, service) not in {
ServiceStatus.ACCEPTED, ServiceStatus.RESPONDING,
ServiceStatus.ARRIVED, ServiceStatus.WORKING,
}:
_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)]
report = PhoneReportRecord(
service=service, crew=crew, phase=phase,
text=REPORT_TEXT[phase], at=now_utc(),
)
state.phone_reports.append(report)
hub.to_station(session_id, PhoneReport(**report.model_dump()))
case "station.finish":
if state.exercise is not Exercise.DDS:
_error(session_id, "Операторское занятие завершается после звонка 112")
return
state.ended_at = now_utc()
state.end_reason = CallEndReason.COMPLETE
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))
case "card.bounce":
# Карточка вернулась: в разборе это E6 с конкретной причиной.
state.bounced_fields = list(event.missing_fields)
@ -105,6 +189,11 @@ async def station(ws: WebSocket, session_id: UUID, role: str = "dds") -> None:
)
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))