refactor: запись хода занятия — один commit SessionStore на операцию вместо журнала и checkpoint

Правка балла с пульта и из отчёта идёт одной доменной операцией: раньше WS-путь не обновлял full_report живой сессии.
This commit is contained in:
gglamer 2026-09-26 22:48:24 +00:00
commit 47cb85ee02
30 changed files with 1406 additions and 1161 deletions

View file

@ -11,7 +11,6 @@ import logging
from contextvars import ContextVar, Token
from collections.abc import AsyncIterator, Iterator
from datetime import UTC, datetime
from typing import Protocol
from uuid import UUID
from pydantic import BaseModel
@ -24,6 +23,7 @@ from app.domain.events import (
TimerTick,
)
from app.session.state import SessionState
from app.session.store import MemorySessionStore, Record, SessionStore
#: Очередь одного подписчика. Медленный наблюдатель не тормозит занятие:
#: очередь ограничена, переполнение роняет соединение, а не сессию.
@ -34,51 +34,30 @@ LEASE_FENCED_MESSAGE = "Занятие передано другому backend-
log = logging.getLogger(__name__)
def _current_task():
try:
return asyncio.current_task()
except RuntimeError: # synchronous tests and tooling have no running loop
return None
class _Operation:
"""Отложенная публикация одной операции: события ждут её commit."""
class Journal(Protocol):
"""Запись в БД. Вынесена за хаб: без базы занятие должно идти,
но молчать об ошибке записи нельзя."""
async def start_lesson(
self, session_id: UUID, scenario_id: str, mode: str, trainee_name: str | None,
trainee_id: UUID | None = None, owner_login: str | None = None,
backend_node_id: str | None = None,
) -> tuple[int, UUID | None, str | None, int] | None: ...
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
) -> bool: ...
async def score(self, session_id: UUID, score_auto: float, report: dict) -> bool: ...
async def score_snapshot(self, session_id: UUID, report: dict) -> None: ...
async def score_override(
self, session_id: UUID, score_final: float, author: str, comment: str,
) -> bool: ...
async def session_started(self, session_id: UUID, at) -> None: ...
async def session_ended(self, session_id: UUID, at, reason: str) -> None: ...
async def checkpoint(self, state: SessionState) -> None: ...
async def restore_active(self) -> list[SessionState]: ...
async def renew(self, session_id: UUID) -> None: ...
async def claim_expired(self, session_id: UUID | None = None) -> list[SessionState]: ...
def __init__(self, session_id: UUID) -> None:
self.session_id = session_id
self.records: list[Record] = []
self.events: list[tuple[dict[UUID, set[asyncio.Queue]], BaseModel]] = []
self.open = True
#: Такт, которому нечего сохранять, не пишет снимок каждую секунду.
self.persist = True
class SessionHub:
def __init__(self, journal: Journal | None = None) -> None:
self.journal = journal
def __init__(self, store: SessionStore | None = None) -> None:
self.store: SessionStore = store if store is not None else MemorySessionStore()
self._sessions: dict[UUID, SessionState] = {}
self._observers: dict[UUID, set[asyncio.Queue]] = {}
self._trainees: dict[UUID, set[asyncio.Queue]] = {}
self._stations: dict[UUID, set[asyncio.Queue]] = {}
self._tickers: dict[UUID, asyncio.Task] = {}
self._event_batch: ContextVar[dict | None] = ContextVar(
f"session-event-batch-{id(self)}", default=None
# Задачи, порождённые внутри операции, наследуют контекст; после
# закрытия операции их события идут напрямую (`_Operation.open`).
self._operation: ContextVar[_Operation | None] = ContextVar(
f"session-operation-{id(self)}", default=None
)
# ── реестр ──
@ -139,84 +118,67 @@ class SessionHub:
self._sessions.pop(session_id, None)
self.stop_ticker(session_id)
async def checkpoint(self, session_id: UUID) -> None:
"""Зафиксировать подтверждённое состояние, если журнал доступен."""
state = self._sessions.get(session_id)
if state is not None and state.lease_fenced:
raise RuntimeError(LEASE_FENCED_MESSAGE)
if state is not None and self.journal is not None:
try:
await self.journal.checkpoint(state)
except Exception:
self._discard_event_batch(session_id)
await self.fence(state)
raise
self._flush_event_batch(session_id)
def begin_event_stream(self, session_id: UUID) -> Token:
"""Stage controller output until each explicit checkpoint in its loop."""
return self._event_batch.set({
"session_id": session_id, "events": [], "committed": False,
"persistent": True, "owner_task": _current_task(),
})
async def end_event_stream(self, token: Token) -> None:
batch = self._event_batch.get()
try:
if batch is not None and batch["events"]:
session_id = batch["session_id"]
batch["events"].clear()
state = self._sessions.get(session_id)
if self.journal is not None and state is not None and not state.lease_fenced:
await self.fence(state)
finally:
self._event_batch.reset(token)
# ── операция: один commit, потом события ──
@contextlib.asynccontextmanager
async def durable_transition(self, session_id: UUID):
"""Do not publish state-changing events until its checkpoint commits."""
batch = {
"session_id": session_id, "events": [], "committed": False,
"persistent": False, "owner_task": _current_task(),
}
token: Token = self._event_batch.set(batch)
async def operation(self, session_id: UUID) -> AsyncIterator[_Operation]:
"""Изменение занятия фиксируется одним `store.commit` на выходе.
События копятся до коммита. Исключение внутри операции или сбой
коммита отбрасывают их и закрывают занятие на узле (fail closed).
"""
op = _Operation(session_id)
token: Token = self._operation.set(op)
try:
yield
if not batch["committed"]:
await self.checkpoint(session_id)
else:
self._flush_event_batch(session_id)
except Exception:
self._discard_event_batch(session_id)
yield op
state = self._sessions.get(session_id)
if self.journal is not None and state is not None and not state.lease_fenced:
if state is not None and op.persist:
if state.lease_fenced:
raise RuntimeError(LEASE_FENCED_MESSAGE)
op.persist = False # второй commit из обработчика отмены не нужен
await asyncio.shield(self.store.commit(state, op.records))
except asyncio.CancelledError:
# Задачу сокета отменили посреди операции — это не сбой хранилища.
# Сделанное фиксируется, события слать уже некому.
op.events.clear()
state = self._sessions.get(session_id)
if state is not None and op.persist and not state.lease_fenced:
try:
await asyncio.shield(self.store.commit(state, op.records))
except Exception:
await self.fence(state)
raise
except Exception:
op.events.clear()
state = self._sessions.get(session_id)
if state is not None and not state.lease_fenced:
await self.fence(state)
raise
finally:
self._event_batch.reset(token)
op.open = False
self._operation.reset(token)
for registry, event in op.events:
self._put(registry.get(session_id, set()), event)
def _discard_event_batch(self, session_id: UUID) -> None:
batch = self._event_batch.get()
if batch is not None and batch["session_id"] == session_id:
batch["events"].clear()
def record(self, session_id: UUID, record: Record) -> None:
"""Строка журнала уходит в commit текущей операции, не отдельной транзакцией."""
op = self._operation.get()
if op is None or not op.open or op.session_id != session_id:
raise RuntimeError(f"запись занятия {session_id} вне операции")
op.records.append(record)
def _flush_event_batch(self, session_id: UUID) -> None:
batch = self._event_batch.get()
if batch is None or batch["session_id"] != session_id:
return
pending, batch["events"] = batch["events"], []
batch["committed"] = not batch.get("persistent", False)
for registry, target_session_id, event in pending:
self._put(registry.get(target_session_id, set()), event)
async def commit(self, session_id: UUID, *records: Record) -> None:
"""Операция из одних строк — реплика голосового контура."""
async with self.operation(session_id):
for record in records:
self.record(session_id, record)
def _send(self, registry: dict[UUID, set[asyncio.Queue]], session_id: UUID,
event: BaseModel) -> None:
batch = self._event_batch.get()
if isinstance(event, ErrorEvent):
self._put(registry.get(session_id, set()), event)
elif (batch is not None and batch["session_id"] == session_id
and batch["owner_task"] is _current_task()):
batch["events"].append((registry, session_id, event))
op = self._operation.get()
if (op is not None and op.open and op.session_id == session_id
and not isinstance(event, ErrorEvent)):
op.events.append((registry, event))
else:
self._put(registry.get(session_id, set()), event)
@ -316,12 +278,11 @@ class SessionHub:
if state is None or state.ended:
return
if state.dds_phase:
active_before = state.desk.active_id
delivered = state.desk.deliver_due()
if state.desk.active_id != active_before and state.desk.active_id:
self.to_station(session_id, state.card_received_event())
if delivered:
await self.checkpoint(session_id)
async with self.operation(session_id) as op:
active_before = state.desk.active_id
op.persist = bool(state.desk.deliver_due())
if state.desk.active_id != active_before and state.desk.active_id:
self.to_station(session_id, state.card_received_event())
# Keep the pending count and countdown live even while
# the active dispatcher card is being handled.
self.to_station(session_id, StationState(snapshot=state.station_snapshot()))