diff --git a/backend/app/api/ws/call.py b/backend/app/api/ws/call.py index cf375f2..90d2956 100644 --- a/backend/app/api/ws/call.py +++ b/backend/app/api/ws/call.py @@ -30,7 +30,6 @@ from app.domain.events import ( KioState, PatchSource, ScoreReady, - SessionEnded, SessionMode, StationState, TimerTick, @@ -39,17 +38,15 @@ from app.domain.events import ( TranscriptAppend, TraineeToServer, ) -from app.domain.kio import ResponseStatus from app.dialog.slots import TurnResult from app.domain.roles import Role from app.scenarios import store from app.session.dds import prepare_handoff_queue -from app.session.finish import finish, refresh_archived_report, release_score +from app.session.finish import end_session, refresh_archived_report, release_score from app.session.hub import LEASE_FENCED_MESSAGE, hub from app.session.state import now_utc from app.session.store import ( HintRecorded, - LessonEnded, LessonStarted, SelfAssessed, UtteranceAppended, @@ -242,9 +239,7 @@ async def _handle(session_id: UUID, state, event) -> None: case "card.submit": state.on_event("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.dispatch() if state.handoff_to_dds: prepare_handoff_queue( state, @@ -253,20 +248,14 @@ async def _handle(session_id: UUID, state, event) -> None: max_waiting=state.desk.max_waiting, ) state.pending_dds_scenarios = [] - else: - 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)) if state.handoff_to_dds: + hub.to_trainee(session_id, CallEnded(reason=CallEndReason.COMPLETE)) hub.to_observers(session_id, state.snapshot()) else: - hub.to_observers(session_id, SessionEnded(reason=CallEndReason.COMPLETE)) - hub.record(session_id, LessonEnded(state.ended_at, CallEndReason.COMPLETE.value)) - await finish(session_id, state) + await end_session(session_id, state, CallEndReason.COMPLETE) case "call.answer": first_answer = state.started_at is None if first_answer: @@ -362,15 +351,7 @@ async def _handle(session_id: UUID, state, event) -> None: await release_score(session_id, state) case "call.hangup": - if state.voice is not None: - await state.voice.close() - state.ended_at = now_utc() - state.end_reason = CallEndReason.HANGUP - hub.stop_ticker(session_id) - hub.to_trainee(session_id, CallEnded(reason=CallEndReason.HANGUP)) - hub.to_observers(session_id, SessionEnded(reason=CallEndReason.HANGUP)) - hub.record(session_id, LessonEnded(state.ended_at, CallEndReason.HANGUP.value)) - await finish(session_id, state) + await end_session(session_id, state, CallEndReason.HANGUP) def _start_voice(session_id: UUID, state, *, initial_statement: bool = True) -> None: diff --git a/backend/app/api/ws/control.py b/backend/app/api/ws/control.py index bf07bb4..801f81f 100644 --- a/backend/app/api/ws/control.py +++ b/backend/app/api/ws/control.py @@ -40,7 +40,6 @@ from app.domain.events import ( InstructorToServer, ModeSet, ReferenceStarted, - ScoreReady, SessionEnded, StationState, ) @@ -49,9 +48,9 @@ from app.domain.timers import TimerCode from app.scenarios import store from app.session.dds import prepare_queue from app.session.hub import LEASE_FENCED_MESSAGE, hub -from app.session.finish import finish, override_score +from app.session.finish import end_session, override_score from app.session.state import SessionState, now_utc -from app.session.store import LessonEnded, LessonIdentity, LessonRequest, NoteAdded, ScoreOverridden +from app.session.store import LessonIdentity, LessonRequest, NoteAdded, ScoreOverridden from app.voice.models import get_voice_models from app.voice.pipeline import FILLERS, prefetch @@ -331,28 +330,9 @@ async def _start(session_id: UUID, event, who=None) -> None: async def _stop(session_id: UUID) -> None: state = hub.get(session_id) - if state is None or state.ended: + if state is None: return - state.ended_at = now_utc() - if state.exercise is Exercise.CARD and state.dispatched_card is None: - state.on_event("card.end") - state.end_reason = CallEndReason.INSTRUCTOR - if state.voice is not None: - await state.voice.close() - hub.stop_ticker(session_id) - hub.to_observers(session_id, SessionEnded(reason=CallEndReason.INSTRUCTOR)) - - if state.exercise in {Exercise.DDS, Exercise.CARD}: - if state.exercise is Exercise.DDS or state.handoff_to_dds and state.dispatched_card is not None: - hub.to_station(session_id, SessionEnded(reason=CallEndReason.INSTRUCTOR)) - else: - hub.to_trainee(session_id, CallEnded(reason=CallEndReason.INSTRUCTOR)) - if state.exercise is Exercise.DDS or state.handoff_to_dds and state.dispatched_card is not None: - hub.to_station(session_id, ScoreReady(session_id=session_id)) - else: - hub.to_trainee(session_id, CallEnded(reason=CallEndReason.INSTRUCTOR)) - hub.record(session_id, LessonEnded(state.ended_at, CallEndReason.INSTRUCTOR.value)) - await finish(session_id, state) + await end_session(session_id, state, CallEndReason.INSTRUCTOR) async def _command(session_id: UUID, event, who) -> None: diff --git a/backend/app/api/ws/station.py b/backend/app/api/ws/station.py index 737c7ec..b6e9aa3 100644 --- a/backend/app/api/ws/station.py +++ b/backend/app/api/ws/station.py @@ -24,15 +24,12 @@ from app.domain.events import ( CommandAck, ErrorEvent, ErrorKind, - ScoreReady, - SessionEnded, StationState, StationToServer, ) from app.domain.roles import Role -from app.session.finish import finish +from app.session.finish import end_session from app.session.hub import LEASE_FENCED_MESSAGE, hub -from app.session.store import LessonEnded log = logging.getLogger(__name__) router = APIRouter() @@ -44,16 +41,6 @@ def _error(session_id: UUID, message: str) -> None: hub.to_station(session_id, ErrorEvent(code=ErrorKind.UNSUPPORTED_EVENT, message=message)) -async def _finish_dds(session_id: UUID, state) -> None: - ended_at = state.end(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)) - hub.record(session_id, LessonEnded(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: outcome = state.desk.apply(event, state) if outcome.error is not None: @@ -62,7 +49,7 @@ async def _handle(session_id: UUID, state, event) -> None: for item in outcome.events: hub.to_station(session_id, item) if outcome.finished: - await _finish_dds(session_id, state) + await end_session(session_id, state, CallEndReason.COMPLETE) if not outcome.changed: return hub.to_station(session_id, StationState(snapshot=state.station_snapshot())) diff --git a/backend/app/session/dds.py b/backend/app/session/dds.py index 1e4110a..67832de 100644 --- a/backend/app/session/dds.py +++ b/backend/app/session/dds.py @@ -575,12 +575,12 @@ class DdsDesk(BaseModel): outcome.events.append(session.card_received_event()) outcome.events.append(StationState(snapshot=session.station_snapshot())) if not self.cards and len(self.completed) >= len(self.scenarios): - self._finish_work() + self.finish_work() outcome.finished = True case "station.finish": if not dds_lesson: return _refused("Операторское занятие завершается после звонка 112") - self._finish_work() + self.finish_work() outcome.finished = True case "card.bounce": return _refused( @@ -608,7 +608,7 @@ class DdsDesk(BaseModel): card.dds_log.append(("crew.arrived", now_utc(), None)) return outcome - def _finish_work(self) -> None: + def finish_work(self) -> None: """Конец занятия останавливает трёхминутный таймер у каждой начатой карточки.""" for card in self.cards.values(): timer = card.timers.timers.get(TimerCode.DDS_WORK) diff --git a/backend/app/session/finish.py b/backend/app/session/finish.py index 8e7e8d9..405183d 100644 --- a/backend/app/session/finish.py +++ b/backend/app/session/finish.py @@ -12,7 +12,7 @@ import logging import time from uuid import UUID -from app.domain.events import Exercise, Metric, ScoreReady +from app.domain.events import CallEnded, CallEndReason, Exercise, Metric, ScoreReady, SessionEnded from app.domain.statuses import ServiceStatus, current from app.domain.taxonomy import Competency, ErrorCode, Finding, FindingSource from app.domain.timers import TimerCode @@ -29,7 +29,13 @@ from app.scoring.timing import time_credit from app.scoring.weights import apply_weights from app.session.hub import hub from app.session.state import DdsCardRecord, DdsLiveCard, now_utc -from app.session.store import ScoreArchived, ScoreCalculated, ScoreOverridden, apply_score_override +from app.session.store import ( + LessonEnded, + ScoreArchived, + ScoreCalculated, + ScoreOverridden, + apply_score_override, +) log = logging.getLogger(__name__) @@ -353,10 +359,53 @@ async def finish(session_id: UUID, state) -> None: log.info("сессия %s: оценка %.1f, отметок %d", session_id, result.score, len(result.findings)) # Оценка и аудит уходят в commit операции, завершившей занятие; ScoreReady - # ждёт того же commit — без записи итог не выдаётся. + # рассылает `end_session` и ждёт того же commit — без записи итог не выдаётся. hub.record(session_id, ScoreCalculated(result.score, state.score)) + + +def _ends_on_station(state) -> bool: + """Итог занятия видит пульт ДДС: упражнение ДДС или связка, где карточка + уже передана диспетчеру. Иначе конец звонка узнаёт курсант 112.""" + return state.exercise is Exercise.DDS or ( + state.handoff_to_dds and state.dispatched_card is not None + ) + + +async def end_session(session_id: UUID, state, reason: CallEndReason) -> None: + """Единственный переход занятия в «завершено» — для станции, пульта и звонка. + + Вызывается внутри уже открытой `hub.operation`: строки журнала и события + уходят её commit, сверка `(node_id, epoch)` остаётся за вызывающим. + Порядок один для всех причин: таймеры фиксируются событием до оценки, + журнал пишется до оценки, рассылка — после, `ScoreReady` — последним. + """ + if state.ended: + return + ended_at = state.end(reason) + if state.exercise is Exercise.CARD and state.dispatched_card is None: + state.on_event("card.end") + # Норматив отработки фиксируется событием, а не текущим значением часов. + state.desk.finish_work() + + if state.voice is not None: + await state.voice.close() + hub.stop_ticker(session_id) + + hub.record(session_id, LessonEnded(ended_at, reason.value)) + await finish(session_id, state) + + hub.to_observers(session_id, SessionEnded(reason=reason)) + on_station = _ends_on_station(state) + if on_station: + hub.to_station(session_id, SessionEnded(reason=reason)) + else: + hub.to_trainee(session_id, CallEnded(reason=reason)) + if state.score is None: + return hub.to_observers(session_id, ScoreReady(session_id=session_id)) await release_score(session_id, state) + if on_station: + hub.to_station(session_id, ScoreReady(session_id=session_id)) async def refresh_archived_report(session_id: UUID, state) -> None: diff --git a/backend/tests/test_end_session.py b/backend/tests/test_end_session.py new file mode 100644 index 0000000..423b3a1 --- /dev/null +++ b/backend/tests/test_end_session.py @@ -0,0 +1,233 @@ +"""Завершение занятия — одна операция `end_session`. + +Порядок шагов проверяется без websocket: фейковый хаб пишет журнал вызовов, +фейковый голос и фейковая оценка отмечают себя в том же журнале. +""" + +import asyncio +import time +from pathlib import Path +from uuid import uuid4 + +import pytest + +from app.domain.events import CallEndReason, Exercise, SessionMode +from app.domain.kio import KIO +from app.domain.timers import TimerCode +from app.scenarios.loader import load_file +from app.session import finish as finish_module +from app.session.dds import prepare_handoff_queue, prepare_queue +from app.session.hub import hub +from app.session.state import SessionState +from app.session.store import LessonEnded, MemorySessionStore + +LIBRARY = Path(__file__).resolve().parents[2] / "scenarios" + + +def fire(): + return load_file(LIBRARY / "fire-apartment-l2.yaml", LIBRARY) + + +def dds_state() -> SessionState: + """Пульт ДДС на две карточки, вторая открыта: `DDS_WORK` идёт.""" + first = fire() + scenarios = [first.model_copy(deep=True, update={"id": f"desk-{index}"}) + for index in range(2)] + state = SessionState( + session_id=uuid4(), scenario_id=first.id, scenario_title=first.title, + level=first.level.value, mode=SessionMode.TRAINING, exercise=Exercise.DDS, + ) + prepare_queue(state, scenarios) + card = state.desk.ordered()[1] + state.desk.open(card.card_id) + # Карточка принята: три минуты отработки идут с этого момента. + card.on_event("card.ack") + card.on_event("dds.open") + return state + + +def card_state(*, handoff: bool = False) -> SessionState: + scenario = fire() + state = SessionState( + session_id=uuid4(), scenario_id=scenario.id, scenario_title=scenario.title, + level=scenario.level.value, mode=SessionMode.TRAINING, + exercise=Exercise.CARD, handoff_to_dds=handoff, scenario=scenario, + ) + state.kio = KIO(address="улица Ленина, 14", notify=["Служба 101"]) + if handoff: + state.dispatch() + prepare_handoff_queue(state, []) + return state + + +def work_stopped(state) -> bool: + return all( + card.timers.measured_ms(TimerCode.DDS_WORK) is not None + for card in state.desk.cards.values() + if TimerCode.DDS_WORK in card.timers.timers + ) + + +class FakeHub: + def __init__(self, log: list, state) -> None: + self.log = log + self.state = state + + def stop_ticker(self, _session_id) -> None: + self.log.append(("ticker", work_stopped(self.state))) + + def record(self, _session_id, record) -> None: + self.log.append(("record", type(record).__name__)) + + def to_station(self, _session_id, event) -> None: + self.log.append(("station", event.type)) + + def to_trainee(self, _session_id, event) -> None: + self.log.append(("trainee", event.type)) + + def to_observers(self, _session_id, event) -> None: + self.log.append(("observers", event.type)) + + +class FakeVoice: + def __init__(self, log: list, state) -> None: + self.log = log + self.state = state + + async def close(self) -> None: + self.log.append(("voice", work_stopped(self.state))) + + +@pytest.fixture +def fake(monkeypatch): + def install(state): + log: list = [] + monkeypatch.setattr(finish_module, "hub", FakeHub(log, state)) + + async def fake_finish(_session_id, finished): + log.append(("finish", finished.ended)) + finished.score = {"score": 1.0} + + monkeypatch.setattr(finish_module, "finish", fake_finish) + return log + return install + + +def test_end_session_runs_steps_in_fixed_order(fake): + state = dds_state() + log = fake(state) + state.voice = FakeVoice(log, state) + + asyncio.run(finish_module.end_session(state.session_id, state, CallEndReason.INSTRUCTOR)) + + assert state.ended and state.end_reason is CallEndReason.INSTRUCTOR + # таймеры → голос → тикер → журнал → оценка → рассылка + assert log == [ + ("voice", True), + ("ticker", True), + ("record", "LessonEnded"), + ("finish", True), + ("observers", "session.ended"), + ("station", "session.ended"), + ("observers", "score.ready"), + ("station", "score.ready"), + ] + + +def test_card_without_handoff_tells_trainee_and_releases_score(fake): + state = card_state() + log = fake(state) + + asyncio.run(finish_module.end_session(state.session_id, state, CallEndReason.COMPLETE)) + + events = [entry for entry in log if entry[0] in {"observers", "station", "trainee"}] + assert events == [ + ("observers", "session.ended"), + ("trainee", "call.ended"), + ("observers", "score.ready"), + ("trainee", "score.ready"), + ] + + +def test_card_handed_to_dds_ends_on_station(fake): + state = card_state(handoff=True) + log = fake(state) + + asyncio.run(finish_module.end_session(state.session_id, state, CallEndReason.INSTRUCTOR)) + + events = [entry for entry in log if entry[0] in {"observers", "station", "trainee"}] + assert events == [ + ("observers", "session.ended"), + ("station", "session.ended"), + ("observers", "score.ready"), + ("trainee", "score.ready"), + ("station", "score.ready"), + ] + + +def test_card_stopped_before_submit_stops_fill_timer(fake): + state = card_state() + state.on_event("card.start") + fake(state) + + asyncio.run(finish_module.end_session(state.session_id, state, CallEndReason.INSTRUCTOR)) + + assert state.timers.measured_ms(TimerCode.CARD_FILL) is not None + + +def test_second_end_session_writes_and_sends_nothing(fake): + state = dds_state() + log = fake(state) + asyncio.run(finish_module.end_session(state.session_id, state, CallEndReason.COMPLETE)) + ended_at = state.ended_at + log.clear() + + asyncio.run(finish_module.end_session(state.session_id, state, CallEndReason.INSTRUCTOR)) + + assert log == [] + assert state.ended_at == ended_at + assert state.end_reason is CallEndReason.COMPLETE + + +@pytest.fixture +def frozen_clock(monkeypatch): + clock = {"now": time.monotonic()} + monkeypatch.setattr(time, "monotonic", lambda: clock["now"]) + + class NoCoaching: + def model_dump(self, **_kwargs): + return {} + + async def no_coach(_metrics): + return NoCoaching() + + monkeypatch.setattr(finish_module, "coach", no_coach) + monkeypatch.setattr(hub, "store", MemorySessionStore()) + return clock + + +def test_instructor_stop_in_dds_stops_work_timer_and_keeps_score(frozen_clock): + before, after = dds_state(), dds_state() + frozen_clock["now"] += 95 + + async def old_stop(state): + # Прежний `_stop`: оценка брала время отработки по `current_ms()`. + hub.register(state) + async with hub.operation(state.session_id): + ended_at = state.end(CallEndReason.INSTRUCTOR) + hub.record(state.session_id, LessonEnded(ended_at, CallEndReason.INSTRUCTOR.value)) + await finish_module.finish(state.session_id, state) + + async def new_stop(state): + hub.register(state) + async with hub.operation(state.session_id): + await finish_module.end_session(state.session_id, state, CallEndReason.INSTRUCTOR) + + asyncio.run(old_stop(before)) + asyncio.run(new_stop(after)) + + active = after.desk.cards[after.desk.active_id] + assert active.timers.measured_ms(TimerCode.DDS_WORK) == 95_000 + assert after.score["score_auto"] == before.score["score_auto"] + assert after.score["metrics"] == before.score["metrics"] + assert after.score["findings"] == before.score["findings"]