lct-05: состояние сессии, таймеры по событиям, четыре канала
Таймеры объявлены таблицей «событие → старт/стоп», а не разбросаны по обработчикам: так видно целиком, что чем запускается, и норматив нельзя потерять по дороге. Опрос стартует на ответе курсанта и останавливается передачей в ДДС — событием, а не таймаутом. Канал наблюдателя без единого обработчика входящих: кадры читаются и выбрасываются, не разбираясь, только чтобы заметить разрыв. Без чтения задача сокета висела бы на очереди до первой отправки, а закрытая вкладка монитора оставляла бы подписку. Тест проверяет, что наблюдатель не может ни изменить карточку, ни завершить занятие. Пульт теперь заводит сессию в журнале: запущенное с него занятие жило только в памяти, и реплики с подсказками уходили в нарушение внешнего ключа — журнал об этом честно сообщал в лог. Тикер таймеров гасится при остановке приложения, иначе задачи переживают выключение и держат событийный цикл. Проверено вживую через uvicorn: монитор, подключённый посреди занятия, получает снимок с заполненной карточкой; такт таймера идёт раз в секунду.
This commit is contained in:
parent
6b2c96bab6
commit
e7a0798ba7
12 changed files with 1029 additions and 0 deletions
0
backend/app/session/__init__.py
Normal file
0
backend/app/session/__init__.py
Normal file
142
backend/app/session/hub.py
Normal file
142
backend/app/session/hub.py
Normal file
|
|
@ -0,0 +1,142 @@
|
|||
"""Реестр живых сессий и подписки наблюдателей.
|
||||
|
||||
Дельты курсанту, полное состояние наблюдателям — сознательно: у курсанта одно
|
||||
соединение и важна задержка, наблюдателей несколько и подключаются они
|
||||
в произвольный момент (docs/arch/CONTRACT.md).
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import contextlib
|
||||
from collections.abc import AsyncIterator, Iterator
|
||||
from typing import Protocol
|
||||
from uuid import UUID
|
||||
|
||||
from pydantic import BaseModel
|
||||
|
||||
from app.domain.events import TimerTick
|
||||
from app.session.state import SessionState
|
||||
|
||||
#: Очередь одного подписчика. Медленный наблюдатель не тормозит занятие:
|
||||
#: очередь ограничена, переполнение роняет соединение, а не сессию.
|
||||
QUEUE_SIZE = 256
|
||||
|
||||
TICK_SECONDS = 1.0
|
||||
|
||||
|
||||
class Journal(Protocol):
|
||||
"""Запись в БД. Вынесена за хаб: без базы занятие должно идти,
|
||||
но молчать об ошибке записи нельзя."""
|
||||
|
||||
async def start_lesson(
|
||||
self, session_id: UUID, scenario_id: str, mode: str, trainee_name: str | None
|
||||
) -> int: ...
|
||||
async def utterance(self, session_id: UUID, entry) -> None: ...
|
||||
async def hint(self, session_id: UUID, checklist_id: str, question: str, at) -> None: ...
|
||||
async def note(self, session_id: UUID, ref: str, text: str, author: str) -> None: ...
|
||||
async def self_assessment(self, session_id: UUID, missed: list[str], comment: str, at) -> None: ...
|
||||
async def session_started(self, session_id: UUID, at) -> None: ...
|
||||
async def session_ended(self, session_id: UUID, at, reason: str) -> None: ...
|
||||
|
||||
|
||||
class SessionHub:
|
||||
def __init__(self, journal: Journal | None = None) -> None:
|
||||
self.journal = journal
|
||||
self._sessions: dict[UUID, SessionState] = {}
|
||||
self._observers: dict[UUID, set[asyncio.Queue]] = {}
|
||||
self._trainees: dict[UUID, set[asyncio.Queue]] = {}
|
||||
self._tickers: dict[UUID, asyncio.Task] = {}
|
||||
|
||||
# ── реестр ──
|
||||
|
||||
def register(self, state: SessionState) -> SessionState:
|
||||
self._sessions[state.session_id] = state
|
||||
return state
|
||||
|
||||
def get(self, session_id: UUID) -> SessionState | None:
|
||||
return self._sessions.get(session_id)
|
||||
|
||||
def drop(self, session_id: UUID) -> None:
|
||||
self._sessions.pop(session_id, None)
|
||||
self.stop_ticker(session_id)
|
||||
|
||||
# ── подписки ──
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _subscribe(self, registry: dict[UUID, set[asyncio.Queue]], session_id: UUID) -> Iterator[asyncio.Queue]:
|
||||
queue: asyncio.Queue = asyncio.Queue(maxsize=QUEUE_SIZE)
|
||||
registry.setdefault(session_id, set()).add(queue)
|
||||
try:
|
||||
yield queue
|
||||
finally:
|
||||
registry.get(session_id, set()).discard(queue)
|
||||
|
||||
def observer(self, session_id: UUID):
|
||||
return self._subscribe(self._observers, session_id)
|
||||
|
||||
def trainee(self, session_id: UUID):
|
||||
return self._subscribe(self._trainees, session_id)
|
||||
|
||||
# ── вещание ──
|
||||
|
||||
@staticmethod
|
||||
def _put(queues: set[asyncio.Queue], event: BaseModel) -> None:
|
||||
for queue in list(queues):
|
||||
try:
|
||||
queue.put_nowait(event)
|
||||
except asyncio.QueueFull:
|
||||
queues.discard(queue)
|
||||
|
||||
def to_observers(self, session_id: UUID, event: BaseModel) -> None:
|
||||
self._put(self._observers.get(session_id, set()), event)
|
||||
|
||||
def to_trainee(self, session_id: UUID, event: BaseModel) -> None:
|
||||
self._put(self._trainees.get(session_id, set()), event)
|
||||
|
||||
def broadcast(self, session_id: UUID, event: BaseModel) -> None:
|
||||
self.to_trainee(session_id, event)
|
||||
self.to_observers(session_id, event)
|
||||
|
||||
def observer_count(self, session_id: UUID) -> int:
|
||||
return len(self._observers.get(session_id, set()))
|
||||
|
||||
# ── такт таймеров ──
|
||||
|
||||
def start_ticker(self, session_id: UUID) -> None:
|
||||
"""`timer.tick` раз в секунду, а не на каждое изменение: таймеров дюжина,
|
||||
а UI всё равно рисует секунды."""
|
||||
if session_id in self._tickers:
|
||||
return
|
||||
self._tickers[session_id] = asyncio.create_task(self._tick(session_id))
|
||||
|
||||
def stop_ticker(self, session_id: UUID) -> None:
|
||||
task = self._tickers.pop(session_id, None)
|
||||
if task is not None:
|
||||
task.cancel()
|
||||
|
||||
async def shutdown(self) -> None:
|
||||
"""Погасить все такты при остановке приложения.
|
||||
|
||||
Без этого задачи тикеров переживают выключение и держат событийный цикл:
|
||||
первым это ловит не продакшен, а тест, который не может закрыть клиент.
|
||||
"""
|
||||
for session_id in list(self._tickers):
|
||||
self.stop_ticker(session_id)
|
||||
|
||||
async def _tick(self, session_id: UUID) -> None:
|
||||
try:
|
||||
while True:
|
||||
await asyncio.sleep(TICK_SECONDS)
|
||||
state = self.get(session_id)
|
||||
if state is None or state.ended:
|
||||
return
|
||||
self.broadcast(session_id, TimerTick(timers=state.timers.snapshot()))
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
|
||||
|
||||
hub = SessionHub()
|
||||
|
||||
|
||||
async def drain(queue: asyncio.Queue) -> AsyncIterator[BaseModel]:
|
||||
while True:
|
||||
yield await queue.get()
|
||||
103
backend/app/session/journal.py
Normal file
103
backend/app/session/journal.py
Normal file
|
|
@ -0,0 +1,103 @@
|
|||
"""Запись хода занятия в БД.
|
||||
|
||||
Профиль курсанта, дельта попыток и аналитика группы строятся по журналу,
|
||||
а не по памяти процесса: всё, что здесь не записано, для оценки не существует.
|
||||
"""
|
||||
|
||||
import logging
|
||||
from datetime import datetime
|
||||
from uuid import UUID
|
||||
|
||||
from sqlalchemy import update
|
||||
from sqlalchemy.ext.asyncio import async_sessionmaker
|
||||
|
||||
from app.db import repo
|
||||
from app.db.models import SelfAssessment, Session
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class DbJournal:
|
||||
def __init__(self, sessionmaker: async_sessionmaker) -> None:
|
||||
self._sessionmaker = sessionmaker
|
||||
|
||||
async def _write(self, action, *args, **kwargs) -> None:
|
||||
"""Ошибка записи не роняет занятие, но и не проглатывается молча:
|
||||
занятие идёт дальше, в логе остаётся след."""
|
||||
try:
|
||||
async with self._sessionmaker() as db:
|
||||
await action(db, *args, **kwargs)
|
||||
except Exception: # noqa: BLE001 — журнал не должен ронять живую сессию
|
||||
log.exception("журнал: запись не удалась")
|
||||
|
||||
async def start_lesson(
|
||||
self, session_id: UUID, scenario_id: str, mode: str, trainee_name: str | None
|
||||
) -> int:
|
||||
"""Завести сессию в журнале и вернуть номер попытки.
|
||||
|
||||
Если база недоступна, занятие всё равно идёт: номер попытки
|
||||
деградирует до первого, и это видно в логе.
|
||||
"""
|
||||
try:
|
||||
async with self._sessionmaker() as db:
|
||||
row = await repo.ensure_session(
|
||||
db,
|
||||
session_id=session_id,
|
||||
scenario_id=scenario_id,
|
||||
mode=mode,
|
||||
trainee_name=trainee_name,
|
||||
)
|
||||
return row.attempt
|
||||
except Exception: # noqa: BLE001 — журнал не должен ронять живую сессию
|
||||
log.exception("журнал: сессию завести не удалось")
|
||||
return 1
|
||||
|
||||
async def utterance(self, session_id: UUID, entry) -> None:
|
||||
await self._write(
|
||||
repo.append_utterance,
|
||||
session_id=session_id,
|
||||
ref=entry.ref,
|
||||
speaker=entry.speaker.value,
|
||||
text=entry.text,
|
||||
at=entry.at,
|
||||
mood=entry.mood.value if entry.mood else None,
|
||||
)
|
||||
|
||||
async def hint(self, session_id: UUID, checklist_id: str, question: str, at: datetime) -> None:
|
||||
await self._write(
|
||||
repo.record_hint, session_id=session_id, checklist_id=checklist_id, question=question, at=at
|
||||
)
|
||||
|
||||
async def note(self, session_id: UUID, ref: str, text: str, author: str) -> None:
|
||||
await self._write(repo.add_note, session_id=session_id, transcript_ref=ref, text=text, author=author)
|
||||
|
||||
async def self_assessment(
|
||||
self, session_id: UUID, missed: list[str], comment: str, at: datetime
|
||||
) -> None:
|
||||
async def action(db):
|
||||
db.add(
|
||||
SelfAssessment(
|
||||
session_id=session_id, missed=missed, comment=comment, submitted_at=at
|
||||
)
|
||||
)
|
||||
await db.commit()
|
||||
|
||||
await self._write(lambda db: action(db))
|
||||
|
||||
async def session_started(self, session_id: UUID, at: datetime) -> None:
|
||||
async def action(db):
|
||||
await db.execute(update(Session).where(Session.id == session_id).values(started_at=at))
|
||||
await db.commit()
|
||||
|
||||
await self._write(lambda db: action(db))
|
||||
|
||||
async def session_ended(self, session_id: UUID, at: datetime, reason: str) -> None:
|
||||
async def action(db):
|
||||
await db.execute(
|
||||
update(Session)
|
||||
.where(Session.id == session_id)
|
||||
.values(ended_at=at, end_reason=reason)
|
||||
)
|
||||
await db.commit()
|
||||
|
||||
await self._write(lambda db: action(db))
|
||||
99
backend/app/session/state.py
Normal file
99
backend/app/session/state.py
Normal file
|
|
@ -0,0 +1,99 @@
|
|||
"""Состояние живой сессии: карточка, транскрипт, таймеры, режим, попытка.
|
||||
|
||||
Живёт в памяти процесса — поэтому воркер uvicorn ровно один: с двумя
|
||||
преподаватель подключился бы к другому процессу, чем курсант, и увидел
|
||||
пустой экран (docs/arch/STACK.md).
|
||||
"""
|
||||
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
from uuid import UUID
|
||||
|
||||
from app.domain.events import (
|
||||
CallEndReason,
|
||||
Mood,
|
||||
SessionMode,
|
||||
SessionSnapshot,
|
||||
Speaker,
|
||||
TranscriptEntry,
|
||||
)
|
||||
from app.domain.kio import KIO, apply_patch
|
||||
from app.session.timers import SessionTimers
|
||||
|
||||
|
||||
def now_utc() -> datetime:
|
||||
"""Часы серверные. Метрика, посчитанная по часам браузера, недоказуема."""
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
@dataclass
|
||||
class SessionState:
|
||||
session_id: UUID
|
||||
scenario_id: str
|
||||
scenario_title: str
|
||||
level: str
|
||||
mode: SessionMode
|
||||
required_fields: list[str] = field(default_factory=list)
|
||||
trainee_name: str | None = None
|
||||
attempt: int = 1
|
||||
|
||||
kio: KIO = field(default_factory=KIO)
|
||||
transcript: list[TranscriptEntry] = field(default_factory=list)
|
||||
timers: SessionTimers = field(default_factory=SessionTimers)
|
||||
hints_shown: list[str] = field(default_factory=list)
|
||||
directives: list[str] = field(default_factory=list)
|
||||
|
||||
started_at: datetime | None = None
|
||||
ended_at: datetime | None = None
|
||||
end_reason: CallEndReason | None = None
|
||||
dispatched_card: KIO | None = None
|
||||
|
||||
def on_event(self, event_type: str) -> None:
|
||||
"""Единственная точка, где событие двигает таймеры."""
|
||||
self.timers.on_event(event_type)
|
||||
|
||||
def append(self, speaker: Speaker, text: str, mood: Mood | None = None) -> TranscriptEntry:
|
||||
entry = TranscriptEntry(
|
||||
ref=f"u{len(self.transcript) + 1}",
|
||||
speaker=speaker,
|
||||
text=text,
|
||||
at=now_utc(),
|
||||
mood=mood,
|
||||
)
|
||||
self.transcript.append(entry)
|
||||
return entry
|
||||
|
||||
def patch_kio(self, fields: dict[str, Any]) -> KIO:
|
||||
self.kio = apply_patch(self.kio, fields)
|
||||
return self.kio
|
||||
|
||||
def dispatch(self, service: str) -> KIO:
|
||||
"""Карточка замораживается снимком: оператор не должен иметь
|
||||
возможности дописать задним числом поле, которое забыл."""
|
||||
self.kio = apply_patch(self.kio, {"dds": service, "response_status": "transferred"})
|
||||
self.dispatched_card = self.kio.model_copy(deep=True)
|
||||
return self.dispatched_card
|
||||
|
||||
@property
|
||||
def ended(self) -> bool:
|
||||
return self.ended_at is not None
|
||||
|
||||
def snapshot(self) -> SessionSnapshot:
|
||||
"""Полное состояние. Монитор в классе включают посреди занятия —
|
||||
он обязан показать текущее, а не ждать следующего события."""
|
||||
return SessionSnapshot(
|
||||
session_id=self.session_id,
|
||||
scenario_id=self.scenario_id,
|
||||
scenario_title=self.scenario_title,
|
||||
level=self.level,
|
||||
mode=self.mode,
|
||||
trainee_name=self.trainee_name,
|
||||
started_at=self.started_at,
|
||||
kio=self.kio,
|
||||
required_fields=self.required_fields,
|
||||
transcript=list(self.transcript),
|
||||
timers=self.timers.snapshot(),
|
||||
hints_used=len(self.hints_shown),
|
||||
ended=self.ended,
|
||||
)
|
||||
106
backend/app/session/timers.py
Normal file
106
backend/app/session/timers.py
Normal file
|
|
@ -0,0 +1,106 @@
|
|||
"""Таймеры сессии.
|
||||
|
||||
**Правило: таймеры останавливаются событиями, а не таймаутами.** Иначе метрика
|
||||
времени опроса становится недоказуемой, а вся оценка держится на том, что её
|
||||
можно предъявить и проверить (docs/arch/CONTRACT.md).
|
||||
|
||||
Связь «событие → таймер» объявлена таблицей, а не разбросана по обработчикам:
|
||||
так видно целиком, что чем запускается, и норматив нельзя потерять по дороге.
|
||||
"""
|
||||
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
from app.domain.timers import NORMATIVES, TimerCode, TimerSnapshot, state_for
|
||||
|
||||
#: Какое событие какой таймер запускает.
|
||||
STARTS: dict[str, tuple[TimerCode, ...]] = {
|
||||
"call.incoming": (TimerCode.ANSWER,),
|
||||
"call.answer": (TimerCode.INTERVIEW, TimerCode.DDS_NOTIFY),
|
||||
"dds.dispatch": (TimerCode.DDS_ACK, TimerCode.CLOSE),
|
||||
"card.received": (TimerCode.ZONE_CHECK,),
|
||||
"callback.dial": (TimerCode.CALLBACK,),
|
||||
}
|
||||
|
||||
#: Какое событие какой таймер останавливает.
|
||||
STOPS: dict[str, tuple[TimerCode, ...]] = {
|
||||
"call.answer": (TimerCode.ANSWER,),
|
||||
"dds.dispatch": (TimerCode.INTERVIEW, TimerCode.DDS_NOTIFY),
|
||||
"card.ack": (TimerCode.DDS_ACK,),
|
||||
"zone.decision": (TimerCode.ZONE_CHECK,),
|
||||
"crew.arrived": (TimerCode.CLOSE,),
|
||||
"call.started": (TimerCode.CALLBACK,),
|
||||
}
|
||||
|
||||
|
||||
@dataclass
|
||||
class Timer:
|
||||
code: TimerCode
|
||||
started_at: float | None = None
|
||||
elapsed_ms: int = 0
|
||||
attempt: int = 1
|
||||
stopped: bool = False
|
||||
|
||||
def start(self, now: float) -> None:
|
||||
if self.stopped:
|
||||
# Повторный запуск после остановки — это новая попытка (обратный дозвон).
|
||||
self.attempt += 1
|
||||
self.stopped = False
|
||||
self.elapsed_ms = 0
|
||||
if self.started_at is None:
|
||||
self.started_at = now
|
||||
|
||||
def stop(self, now: float) -> None:
|
||||
if self.started_at is not None and not self.stopped:
|
||||
self.elapsed_ms = int((now - self.started_at) * 1000)
|
||||
self.stopped = True
|
||||
|
||||
def current_ms(self, now: float) -> int:
|
||||
if self.stopped or self.started_at is None:
|
||||
return self.elapsed_ms
|
||||
return int((now - self.started_at) * 1000)
|
||||
|
||||
|
||||
@dataclass
|
||||
class SessionTimers:
|
||||
"""Набор таймеров одной сессии. `limits` приходит из конфига —
|
||||
норматив меняется значением, а не правкой кода."""
|
||||
|
||||
limits: dict[TimerCode, int] = field(
|
||||
default_factory=lambda: {code: norm.limit_ms for code, norm in NORMATIVES.items()}
|
||||
)
|
||||
timers: dict[TimerCode, Timer] = field(default_factory=dict)
|
||||
|
||||
def on_event(self, event_type: str, now: float | None = None) -> None:
|
||||
now = time.monotonic() if now is None else now
|
||||
for code in STARTS.get(event_type, ()):
|
||||
self.timers.setdefault(code, Timer(code)).start(now)
|
||||
for code in STOPS.get(event_type, ()):
|
||||
self.timers.setdefault(code, Timer(code)).stop(now)
|
||||
|
||||
def snapshot(self, now: float | None = None) -> list[TimerSnapshot]:
|
||||
"""Только запущенные таймеры: показывать нули по нормативам,
|
||||
до которых занятие ещё не дошло, значит пугать курсанта зря."""
|
||||
now = time.monotonic() if now is None else now
|
||||
result: list[TimerSnapshot] = []
|
||||
for code, timer in self.timers.items():
|
||||
limit = self.limits[code]
|
||||
elapsed = timer.current_ms(now)
|
||||
result.append(
|
||||
TimerSnapshot(
|
||||
code=code,
|
||||
elapsed_ms=elapsed,
|
||||
limit_ms=limit,
|
||||
state=state_for(elapsed, limit),
|
||||
attempt=timer.attempt,
|
||||
stopped=timer.stopped,
|
||||
)
|
||||
)
|
||||
return result
|
||||
|
||||
def measured_ms(self, code: TimerCode) -> int | None:
|
||||
"""Зафиксированное событием значение — то, что пойдёт в оценку."""
|
||||
timer = self.timers.get(code)
|
||||
if timer is None or not timer.stopped:
|
||||
return None
|
||||
return timer.elapsed_ms
|
||||
Loading…
Reference in a new issue