refactor: снимок занятия — model_dump сохраняемой части сессии вместо ручного перечисления полей

This commit is contained in:
gglamer 2026-09-26 22:13:35 +00:00
commit 72c6a7cfa3
7 changed files with 546 additions and 471 deletions

View file

@ -1,374 +1,70 @@
"""Переносимый JSON-снимок незавершённого занятия.
Снимок хранится в PostgreSQL после каждого подтверждённого действия. Он не
содержит сокеты, аудиобуферы или объекты моделей: после перезапуска процесса
они создаются заново, а учебные данные, таймеры и состояние ДДС восстанавливаются.
Снимок хранится в PostgreSQL после каждого подтверждённого действия. Это
`model_dump` сохраняемой части сессии (`PersistedSession`): сокеты, аудиобуферы
и объекты моделей в неё не входят и после перезапуска процесса собираются
заново (`rebuild_live`).
Снимки прежних версий не читаются: снимок живёт только у незавершённого
занятия, а плохой снимок реестр и восстановление узла логируют и пропускают.
"""
import time
from datetime import UTC, datetime
from uuid import UUID
from datetime import datetime
from fastapi.encoders import jsonable_encoder
from pydantic import BaseModel, Field
from app.dialog.factory import build_caller
from app.dialog.persona import PersonaState
from app.dialog.persona import PersonaProgress, PersonaState
from app.dialog.runtime import get_embedder
from app.dialog.slots import SlotMachine
from app.domain.events import (
CallEndReason,
Exercise,
LessonCriteria,
Metric,
SessionMode,
TranscriptEntry,
)
from app.domain.kio import KIO
from app.domain.statuses import (
PhoneCallPending,
PhoneLineRecord,
PhoneReportRecord,
StatusEntry,
)
from app.domain.taxonomy import Finding
from app.domain.timers import TimerCode
from app.scenarios.schema import Scenario
from app.session.state import DdsCardRecord, DdsLiveCard, SessionState, now_utc
from app.session.timers import SessionTimers, Timer
from app.dialog.slots import SlotMachine, SlotProgress
from app.domain.events import Exercise
from app.session.state import SessionState
from app.session.timers import downtime_ms
CHECKPOINT_VERSION = 1
CHECKPOINT_VERSION = 2
def _dump_timers(timers: SessionTimers, now: float) -> dict:
return {
"limits": {code.value: limit for code, limit in timers.limits.items()},
"items": {
code.value: {
"elapsed_ms": timer.current_ms(now),
"attempt": timer.attempt,
"stopped": timer.stopped,
}
for code, timer in timers.timers.items()
},
}
class CallerProgress(BaseModel):
"""Прогресс голосового звонящего. В снимок пока не входит — см. `SessionState`."""
def _dump_live_card(item: DdsLiveCard, now: float) -> dict:
return {
"original_index": item.original_index,
"scenario": item.scenario.model_dump(mode="json"),
"kio": item.kio.model_dump(mode="json"),
"dispatched_card": item.dispatched_card.model_dump(mode="json"),
"dispatched_at": item.dispatched_at.isoformat(),
"timers": _dump_timers(item.timers, now),
"bounced_fields": item.bounced_fields,
"dds_log": [[action, at.isoformat(), detail] for action, at, detail in item.dds_log],
"status_log": [entry.model_dump(mode="json") for entry in item.status_log],
"crew_selected": item.crew_selected,
"crew_assignments": item.crew_assignments,
"phone_reports": [entry.model_dump(mode="json") for entry in item.phone_reports],
"phone_lines": [entry.model_dump(mode="json") for entry in item.phone_lines],
"phone_pending": item.phone_pending.model_dump(mode="json") if item.phone_pending else None,
"reply_text": item.reply_text,
"reply_log": [[at.isoformat(), text] for at, text in item.reply_log],
}
slots: SlotProgress = Field(default_factory=SlotProgress)
persona: PersonaProgress = Field(default_factory=PersonaProgress)
def dump_state(state: SessionState) -> dict:
"""Сериализовать только данные, необходимые для точного продолжения.
"""Сериализовать сохраняемую часть. Чистое чтение."""
return {"version": CHECKPOINT_VERSION, **state.model_dump(mode="json")}
Чистое чтение. Верхние ключи работы диспетчера повторяют активную
карточку пульта — так их писал формат до `DdsDesk`.
def rebuild_live(state: SessionState, caller: CallerProgress | None = None) -> None:
"""Собрать runtime-объекты звонящего; без прогресса — с чистого листа.
Карточка и оценка при этом остаются прежними.
"""
desk = state.desk
active = desk.active
now = time.monotonic()
payload = {
"version": CHECKPOINT_VERSION,
"session_id": str(state.session_id),
"scenario_id": state.scenario_id,
"scenario_title": state.scenario_title,
"level": state.level,
"mode": state.mode.value,
"owner_login": state.owner_login,
"backend_fencing_epoch": state.backend_fencing_epoch,
"exercise": state.exercise.value,
"handoff_to_dds": state.handoff_to_dds,
"required_fields": state.required_fields,
"trainee_name": state.trainee_name,
"trainee_id": str(state.trainee_id) if state.trainee_id else None,
"dds_service": state.dds_service,
"attempt": state.attempt,
"criteria": state.criteria.model_dump(mode="json"),
"kio": state.kio.model_dump(mode="json"),
"transcript": [entry.model_dump(mode="json") for entry in state.transcript],
"timers": _dump_timers(state.timers, now),
"hints_shown": state.hints_shown,
"hints_log": [[item, at.isoformat()] for item, at in state.hints_log],
"notes": state.notes,
"directives": state.directives,
"scenario": state.scenario.model_dump(mode="json") if state.scenario else None,
"audio_frames": state.audio_frames,
"bad_frames": state.bad_frames,
"self_assessed": state.self_assessed,
"self_assessment": state.self_assessment,
"score": state.score,
"started_at": state.started_at.isoformat() if state.started_at else None,
"ended_at": state.ended_at.isoformat() if state.ended_at else None,
"end_reason": state.end_reason.value if state.end_reason else None,
"dispatched_card": (
state.dispatched_card.model_dump(mode="json") if state.dispatched_card else None
),
"dispatched_at": state.dispatched_at.isoformat() if state.dispatched_at else None,
"bounced_fields": state.bounced_fields,
"dds_log": [
[action, at.isoformat(), detail] for action, at, detail in active.dds_log
] if active else [],
"status_log": [item.model_dump(mode="json") for item in active.status_log] if active else [],
"crew_selected": active.crew_selected if active else None,
"crew_assignments": active.crew_assignments if active else {},
"phone_reports": [item.model_dump(mode="json") for item in active.phone_reports] if active else [],
"phone_lines": [item.model_dump(mode="json") for item in active.phone_lines] if active else [],
"phone_pending": (
active.phone_pending.model_dump(mode="json") if active and active.phone_pending else None
),
"dds_scenarios": [item.model_dump(mode="json") for item in desk.scenarios],
"pending_dds_scenarios": [item.model_dump(mode="json") for item in state.pending_dds_scenarios],
"operator_kio": state.operator_kio.model_dump(mode="json") if state.operator_kio else None,
"operator_scenario": state.operator_scenario.model_dump(mode="json") if state.operator_scenario else None,
"dds_live_cards": [_dump_live_card(item, now) for item in desk.ordered()],
"dds_active_card_id": str(desk.active_id) if desk.active_id else None,
"dds_card_index": desk.card_index,
"dds_arrival_interval_seconds": desk.arrival_interval_seconds,
"dds_max_waiting": desk.max_waiting,
"dds_next_scenario_index": desk.next_index,
"dds_next_arrival_at": (
desk.next_arrival_at.isoformat() if desk.next_arrival_at else None
),
"dds_completed": [
{
"card_id": str(item.card_id),
"scenario_id": item.scenario_id,
"reply_text": item.reply_text,
"metrics": [metric.model_dump(mode="json") for metric in item.metrics],
"findings": [finding.model_dump(mode="json") for finding in item.findings],
"actions": item.actions,
"duration_ms": item.duration_ms,
"title": item.title,
"address": item.address,
"description": item.description,
"incident_type": item.incident_type,
"victims_count": item.victims_count,
"received_at": item.received_at.isoformat() if item.received_at else None,
"managed_service": item.managed_service,
"recipient_services": item.recipient_services,
}
for item in desk.completed
],
"reply_text": active.reply_text if active else "",
"reply_log": [[at.isoformat(), text] for at, text in active.reply_log] if active else [],
"text_revealed_facts": state.text_revealed_facts,
"resolved_outcome": state.resolved_outcome,
"resolve_comment": state.resolve_comment,
"processed_station_commands": state.processed_station_commands[-512:],
}
# В actions/score могут быть datetime/UUID из расчёта; JSONB должен
# получать только стандартные JSON-типы.
return jsonable_encoder(payload)
def _dt(value: str | None) -> datetime | None:
return datetime.fromisoformat(value) if value else None
def _restore_timers(payload: dict, saved_at: datetime) -> SessionTimers:
raw = payload.get("timers") or {}
limits = {
TimerCode(code): int(limit)
for code, limit in (raw.get("limits") or {}).items()
}
restored = SessionTimers(limits=limits or SessionTimers().limits)
now_mono = time.monotonic()
now_wall = datetime.now(UTC)
if saved_at.tzinfo is None:
saved_at = saved_at.replace(tzinfo=UTC)
downtime_ms = max(0, int((now_wall - saved_at).total_seconds() * 1000))
for raw_code, item in (raw.get("items") or {}).items():
code = TimerCode(raw_code)
stopped = bool(item.get("stopped"))
elapsed = max(0, int(item.get("elapsed_ms", 0)))
total = elapsed if stopped else elapsed + downtime_ms
restored.timers[code] = Timer(
code=code,
started_at=None if stopped else now_mono - total / 1000,
elapsed_ms=elapsed if stopped else 0,
attempt=max(1, int(item.get("attempt", 1))),
stopped=stopped,
)
return restored
def _restore_live_card(item: dict, saved_at: datetime) -> DdsLiveCard:
return DdsLiveCard(
original_index=int(item["original_index"]),
scenario=Scenario.model_validate(item["scenario"]),
kio=KIO.model_validate(item["kio"]),
dispatched_card=KIO.model_validate(item["dispatched_card"]),
dispatched_at=datetime.fromisoformat(item["dispatched_at"]),
timers=_restore_timers({"timers": item["timers"]}, saved_at),
bounced_fields=list(item.get("bounced_fields") or []),
dds_log=[(action, datetime.fromisoformat(at), detail)
for action, at, detail in item.get("dds_log", [])],
status_log=[StatusEntry.model_validate(entry)
for entry in item.get("status_log", [])],
crew_selected=item.get("crew_selected"),
crew_assignments=dict(item.get("crew_assignments") or {}),
phone_reports=[PhoneReportRecord.model_validate(entry)
for entry in item.get("phone_reports", [])],
phone_lines=[PhoneLineRecord.model_validate(entry)
for entry in item.get("phone_lines", [])],
phone_pending=(PhoneCallPending.model_validate(item["phone_pending"])
if item.get("phone_pending") else None),
reply_text=item.get("reply_text", ""),
reply_log=[(datetime.fromisoformat(at), text)
for at, text in item.get("reply_log", [])],
if state.exercise is not Exercise.CALL or state.scenario is None:
return
progress = caller or CallerProgress()
state.persona = PersonaState(state.scenario.persona, progress=progress.persona)
state.caller = build_caller(
state.scenario.id,
use_pregenerated=state.scenario.tree.pregenerated,
)
def _restore_completed(item: dict) -> DdsCardRecord:
return DdsCardRecord(
card_id=UUID(item["card_id"]),
scenario_id=item["scenario_id"],
reply_text=item.get("reply_text", ""),
metrics=[Metric.model_validate(metric) for metric in item.get("metrics", [])],
findings=[Finding.model_validate(finding) for finding in item.get("findings", [])],
actions=list(item.get("actions") or []),
duration_ms=int(item.get("duration_ms", 0)),
title=item.get("title"),
address=item.get("address"),
description=item.get("description"),
incident_type=item.get("incident_type"),
victims_count=item.get("victims_count"),
received_at=_dt(item.get("received_at")),
managed_service=item.get("managed_service"),
recipient_services=list(item.get("recipient_services") or []),
)
def _restore_desk(state: SessionState, payload: dict, saved_at: datetime) -> None:
desk = state.desk
live = payload.get("dds_live_cards", [])
desk.scenarios = [Scenario.model_validate(item) for item in payload.get("dds_scenarios", [])]
desk.completed = [_restore_completed(item) for item in payload.get("dds_completed", [])]
desk.limits = dict(state.timers.limits)
desk.card_index = int(payload.get("dds_card_index", 0))
desk.arrival_interval_seconds = int(payload.get("dds_arrival_interval_seconds", 0))
desk.max_waiting = int(payload.get("dds_max_waiting", 3))
desk.next_index = int(payload.get(
"dds_next_scenario_index",
max((item["original_index"] for item in live), default=-1) + 1,
))
desk.next_arrival_at = _dt(payload.get("dds_next_arrival_at"))
for item in live:
desk.add(_restore_live_card(item, saved_at))
if desk.cards:
# Legacy snapshots had no explicit active ID; newer snapshots may
# intentionally be between cards while waiting for the next arrival.
raw_id = payload.get("dds_active_card_id")
active_id = UUID(raw_id) if raw_id else None
if active_id is None and "dds_active_card_id" not in payload:
active_id = desk.ordered()[0].card_id
if active_id is not None:
desk.open(active_id)
elif (state.exercise is Exercise.CALL and state.dispatched_card is not None
and state.scenario is not None):
# До пульта живой диспетчер упражнения 112 писал прямо в сессию:
# его работа лежит в верхних ключах снимка.
desk.add(_restore_live_card({
**payload,
"original_index": 0,
"scenario": payload["scenario"],
"dispatched_card": payload["dispatched_card"],
"dispatched_at": payload.get("dispatched_at") or saved_at.isoformat(),
}, saved_at))
desk.open(state.dispatched_card.card_id)
if state.exercise is Exercise.CALL:
# Таймеры звонка и диспетчера — одна цепочка 112 → ДДС (см. dispatch).
for card in desk.cards.values():
card.timers = state.timers
if (desk.scenarios
and desk.next_index < len(desk.scenarios)
and desk.next_arrival_at is None):
# Old checkpoints had no delivery schedule; resume any remaining
# selected scenarios immediately rather than strand the session.
desk.next_arrival_at = now_utc()
embedder = get_embedder()
if embedder is not None:
state.slots = SlotMachine(state.scenario, embedder, progress=progress.slots)
def load_state(payload: dict, saved_at: datetime) -> SessionState:
"""Восстановить состояние; неизвестная версия отклоняется явно."""
if payload.get("version") != CHECKPOINT_VERSION:
raise ValueError("неподдерживаемая версия снимка занятия")
scenario = Scenario.model_validate(payload["scenario"]) if payload.get("scenario") else None
state = SessionState(
session_id=UUID(payload["session_id"]),
scenario_id=payload["scenario_id"],
scenario_title=payload["scenario_title"],
level=payload["level"],
mode=SessionMode(payload["mode"]),
owner_login=payload.get("owner_login"),
backend_fencing_epoch=int(payload.get("backend_fencing_epoch", 0)),
exercise=Exercise(payload["exercise"]),
handoff_to_dds=bool(payload.get("handoff_to_dds")),
required_fields=list(payload.get("required_fields") or []),
trainee_name=payload.get("trainee_name"),
trainee_id=UUID(payload["trainee_id"]) if payload.get("trainee_id") else None,
dds_service=payload.get("dds_service"),
attempt=int(payload.get("attempt", 1)),
criteria=LessonCriteria.model_validate(payload.get("criteria") or {}),
kio=KIO.model_validate(payload.get("kio") or {}),
transcript=[TranscriptEntry.model_validate(item) for item in payload.get("transcript", [])],
timers=_restore_timers(payload, saved_at),
hints_shown=list(payload.get("hints_shown") or []),
hints_log=[(item, datetime.fromisoformat(at))
for item, at in payload.get("hints_log", [])],
notes=list(payload.get("notes") or []),
directives=list(payload.get("directives") or []),
scenario=scenario,
audio_frames=int(payload.get("audio_frames", 0)),
bad_frames=int(payload.get("bad_frames", 0)),
self_assessed=bool(payload.get("self_assessed")),
self_assessment=payload.get("self_assessment"),
score=payload.get("score"),
started_at=_dt(payload.get("started_at")),
ended_at=_dt(payload.get("ended_at")),
end_reason=(CallEndReason(payload["end_reason"]) if payload.get("end_reason") else None),
dispatched_card=(
KIO.model_validate(payload["dispatched_card"])
if payload.get("dispatched_card") else None
),
dispatched_at=_dt(payload.get("dispatched_at")),
bounced_fields=list(payload.get("bounced_fields") or []),
pending_dds_scenarios=[Scenario.model_validate(item)
for item in payload.get("pending_dds_scenarios", [])],
operator_kio=(KIO.model_validate(payload["operator_kio"])
if payload.get("operator_kio") else None),
operator_scenario=(Scenario.model_validate(payload["operator_scenario"])
if payload.get("operator_scenario") else None),
text_revealed_facts=dict(payload.get("text_revealed_facts") or {}),
resolved_outcome=payload.get("resolved_outcome"),
resolve_comment=payload.get("resolve_comment", ""),
processed_station_commands=list(payload.get("processed_station_commands") or [])[-512:],
state = SessionState.model_validate(
payload, context={"downtime_ms": downtime_ms(saved_at)}
)
_restore_desk(state, payload, saved_at)
# Голосовые runtime-объекты не сериализуются. Их безопасно собрать заново;
# карточка и оценка при этом остаются прежними.
if state.exercise is Exercise.CALL and state.scenario is not None:
state.persona = PersonaState(state.scenario.persona)
state.caller = build_caller(
state.scenario.id,
use_pregenerated=state.scenario.tree.pregenerated,
)
embedder = get_embedder()
if embedder is not None:
state.slots = SlotMachine(state.scenario, embedder)
if state.exercise is Exercise.CALL:
# Таймеры звонка и диспетчера — одна цепочка 112 → ДДС (см. dispatch):
# в снимке это копии, в живой сессии — один объект.
for card in state.desk.cards.values():
card.timers = state.timers
rebuild_live(state)
return state

View file

@ -12,6 +12,8 @@ from datetime import datetime, timedelta
from typing import TYPE_CHECKING, Any
from uuid import UUID, uuid4
from pydantic import BaseModel, Field
from app.domain import ekp
from app.domain.events import Exercise, Metric, PhoneLine, PhoneReport, StationState
from app.domain.kio import KIO, ResponseStatus, apply_patch
@ -37,8 +39,7 @@ if TYPE_CHECKING:
from app.session.state import SessionState
@dataclass
class DdsCardRecord:
class DdsCardRecord(BaseModel):
card_id: UUID
scenario_id: str
reply_text: str
@ -53,7 +54,7 @@ class DdsCardRecord:
victims_count: int | None = None
received_at: datetime | None = None
managed_service: str | None = None
recipient_services: list[str] = field(default_factory=list)
recipient_services: list[str] = Field(default_factory=list)
@property
def score_auto(self) -> float:
@ -63,8 +64,7 @@ class DdsCardRecord:
return round(100 * passed / total, 1) if total else 0.0
@dataclass
class DdsLiveCard:
class DdsLiveCard(BaseModel):
"""Изолированное живое состояние одной одновременно выданной карточки."""
original_index: int
@ -73,16 +73,16 @@ class DdsLiveCard:
dispatched_card: KIO
dispatched_at: datetime
timers: SessionTimers
bounced_fields: list[str] = field(default_factory=list)
dds_log: list[tuple[str, datetime, str | None]] = field(default_factory=list)
status_log: list[StatusEntry] = field(default_factory=list)
bounced_fields: list[str] = Field(default_factory=list)
dds_log: list[tuple[str, datetime, str | None]] = Field(default_factory=list)
status_log: list[StatusEntry] = Field(default_factory=list)
crew_selected: str | None = None
crew_assignments: dict[str, str] = field(default_factory=dict)
phone_reports: list[PhoneReportRecord] = field(default_factory=list)
phone_lines: list[PhoneLineRecord] = field(default_factory=list)
crew_assignments: dict[str, str] = Field(default_factory=dict)
phone_reports: list[PhoneReportRecord] = Field(default_factory=list)
phone_lines: list[PhoneLineRecord] = Field(default_factory=list)
phone_pending: PhoneCallPending | None = None
reply_text: str = ""
reply_log: list[tuple[datetime, str]] = field(default_factory=list)
reply_log: list[tuple[datetime, str]] = Field(default_factory=list)
@property
def card_id(self) -> UUID:
@ -297,18 +297,17 @@ def _finish_phone_call(card: DdsLiveCard) -> list[Any]:
return [line, PhoneReport(**report.model_dump())]
@dataclass
class DdsDesk:
class DdsDesk(BaseModel):
"""Живые карточки занятия, активная из них и очередь поступления."""
scenarios: list[Scenario] = field(default_factory=list)
cards: dict[UUID, DdsLiveCard] = field(default_factory=dict)
scenarios: list[Scenario] = Field(default_factory=list)
cards: dict[UUID, DdsLiveCard] = Field(default_factory=dict)
active_id: UUID | None = None
#: Номер карточки на экране. После завершения последней из поступивших
#: остаётся прежним, пока не придёт следующая.
card_index: int = 0
completed: list[DdsCardRecord] = field(default_factory=list)
limits: dict[TimerCode, int] = field(default_factory=lambda: dict(SessionTimers().limits))
completed: list[DdsCardRecord] = Field(default_factory=list)
limits: dict[TimerCode, int] = Field(default_factory=lambda: dict(SessionTimers().limits))
arrival_interval_seconds: int = 0
max_waiting: int = 3
next_index: int = 0

View file

@ -6,11 +6,12 @@
"""
import time
from dataclasses import dataclass, field
from datetime import datetime
from typing import Any
from uuid import UUID
from pydantic import BaseModel, ConfigDict, Field, field_serializer, field_validator
from app.dialog.caller import TemplateCaller
from app.dialog.persona import PersonaState
from app.dialog.slots import SlotMachine
@ -40,11 +41,21 @@ from app.scenarios.schema import Scenario
from app.session.dds import DdsCardRecord, DdsDesk, DdsLiveCard
from app.session.timers import SessionTimers, now_utc
__all__ = ["DdsCardRecord", "DdsLiveCard", "SessionState", "now_utc"]
__all__ = ["DdsCardRecord", "DdsLiveCard", "PersistedSession", "SessionState", "now_utc"]
#: Сколько последних команд станции помнит защита от повтора.
MAX_STATION_COMMANDS = 512
@dataclass
class SessionState:
class PersistedSession(BaseModel):
"""Сохраняемая часть занятия — ровно то, что пишется в снимок.
Снимок — `model_dump` этой модели, загрузка — `model_validate`, поэтому
поле сессии нельзя завести, не решив, сохраняемое оно или живое: подкласс
обязан пометить каждое своё поле `Field(exclude=True)`, иначе класс не
создастся (см. `__pydantic_init_subclass__`).
"""
session_id: UUID
scenario_id: str
scenario_title: str
@ -54,12 +65,10 @@ class SessionState:
owner_login: str | None = None
#: Monotonic DB ownership generation; stale processes may not persist writes.
backend_fencing_epoch: int = 0
#: Runtime-only: set when this process loses or cannot confirm DB ownership.
lease_fenced: bool = False
exercise: Exercise = Exercise.CALL
#: После заполнения КИО занятие продолжится на АРМ ДДС, а не завершится.
handoff_to_dds: bool = False
required_fields: list[str] = field(default_factory=list)
required_fields: list[str] = Field(default_factory=list)
trainee_name: str | None = None
#: Чьё это занятие. Проставляется при запуске, когда курсант известен
#: по учётной записи: по нему разбор закрывается от чужих (lct-23).
@ -68,30 +77,20 @@ class SessionState:
#: остальные адресаты карточки показываются информационно.
dds_service: str | None = None
attempt: int = 1
criteria: LessonCriteria = field(default_factory=LessonCriteria)
criteria: LessonCriteria = Field(default_factory=LessonCriteria)
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)
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)
#: Когда именно подсказывали — в разборе видно, какой пункт и на какой минуте.
hints_log: list[tuple[str, datetime]] = field(default_factory=list)
notes: list[dict] = field(default_factory=list)
directives: list[str] = field(default_factory=list)
hints_log: list[tuple[str, datetime]] = Field(default_factory=list)
notes: list[dict] = Field(default_factory=list)
directives: list[str] = Field(default_factory=list)
# Звонящий. Автомата нет, если не скачана модель эмбеддингов:
# занятие идёт, подсказки откатываются на порядок чек-листа.
# Своя копия сценария на занятие: директивы преподавателя правят факты
# и эталон, и правка в одной группе не должна протекать в остальные.
scenario: Scenario | None = None
slots: SlotMachine | None = None
persona: PersonaState | None = None
caller: TemplateCaller | None = None
# Голосовой контур звонка. Нет — если голос выключен или моделей нет:
# тогда кадры микрофона только считаются.
voice: object | None = None
recorder: object | None = None
recording_path: str | None = None
# Аудио курсанта. До голосового контура (lct-06) кадры только считаются —
# этого достаточно, чтобы доказать, что звук доходит от микрофона до сервера.
@ -110,21 +109,73 @@ class SessionState:
dispatched_card: KIO | None = None
dispatched_at: datetime | None = None
#: Поля, из-за которых диспетчер вернул карточку, — основание E6.
bounced_fields: list[str] = field(default_factory=list)
bounced_fields: list[str] = Field(default_factory=list)
#: Готовые карточки связки 112→ДДС ждут, пока курсант не сдаст свою.
pending_dds_scenarios: list[Scenario] = field(default_factory=list)
pending_dds_scenarios: list[Scenario] = Field(default_factory=list)
#: Исходная часть упражнения 112→ДДС сохраняется отдельно от карточек пульта.
operator_kio: KIO | None = None
operator_scenario: Scenario | None = None
#: Работа диспетчера живёт только здесь — в карточках пульта.
desk: DdsDesk = field(default_factory=DdsDesk)
desk: DdsDesk = Field(default_factory=DdsDesk)
#: Чем курсант закрыл вызов, если не карточкой (lct-36).
resolved_outcome: str | None = None
resolve_comment: str = ""
#: Recently committed DDS command IDs; included in the durable checkpoint so
#: a lost WebSocket acknowledgement cannot apply an operation twice.
processed_station_commands: list[str] = field(default_factory=list)
text_revealed_facts: dict[str, str] = field(default_factory=dict)
processed_station_commands: list[str] = Field(default_factory=list)
text_revealed_facts: dict[str, str] = Field(default_factory=dict)
@field_validator("processed_station_commands")
@classmethod
def _recent_commands(cls, value: list[str]) -> list[str]:
return value[-MAX_STATION_COMMANDS:]
@field_serializer("processed_station_commands")
def _dump_recent_commands(self, value: list[str]) -> list[str]:
return value[-MAX_STATION_COMMANDS:]
@classmethod
def __pydantic_init_subclass__(cls, **kwargs: Any) -> None:
super().__pydantic_init_subclass__(**kwargs)
undecided = [name for name, info in cls.model_fields.items()
if name not in PersistedSession.model_fields and not info.exclude]
if undecided:
raise TypeError(
f"{cls.__name__}: поля {undecided} — сохраняемое поле объявляется в "
"PersistedSession, живое — с Field(exclude=True)"
)
class SessionState(PersistedSession):
"""Живое занятие: сохраняемая часть плюс runtime-объекты процесса.
Живые поля в снимок не попадают и пересобираются при загрузке
(`checkpoint.rebuild_live`). Прогресс звонящего (`slots`, `persona`) пока
тоже живой: после failover звонящий начинает с чистого листа — голос P2,
а потеря редкая. Чтобы сохранять его, достаточно добавить в
`PersistedSession` поле `CallerProgress` и передать его в `rebuild_live`.
"""
model_config = ConfigDict(arbitrary_types_allowed=True)
#: Runtime-only: set when this process loses or cannot confirm DB ownership.
lease_fenced: bool = Field(default=False, exclude=True)
# Звонящий. Автомата нет, если не скачана модель эмбеддингов:
# занятие идёт, подсказки откатываются на порядок чек-листа.
slots: SlotMachine | None = Field(default=None, exclude=True)
persona: PersonaState | None = Field(default=None, exclude=True)
caller: TemplateCaller | None = Field(default=None, exclude=True)
# Голосовой контур звонка. Нет — если голос выключен или моделей нет:
# тогда кадры микрофона только считаются.
voice: object | None = Field(default=None, exclude=True)
recorder: object | None = Field(default=None, exclude=True)
recording_path: str | None = Field(default=None, exclude=True)
def persisted(self) -> PersistedSession:
"""Сохраняемая часть без копирования — то, что уходит в снимок."""
return PersistedSession.model_construct(
**{name: getattr(self, name) for name in PersistedSession.model_fields}
)
def on_event(self, event_type: str) -> None:
"""Единственная точка, где событие двигает таймеры."""

View file

@ -9,8 +9,17 @@
"""
import time
from dataclasses import dataclass, field
from datetime import UTC, datetime
from typing import Any
from pydantic import (
BaseModel,
Field,
SerializerFunctionWrapHandler,
ValidationInfo,
model_serializer,
model_validator,
)
from app.domain.timers import NORMATIVES, TimerCode, TimerSnapshot, state_for
@ -57,14 +66,46 @@ STOPS: dict[str, tuple[TimerCode, ...]] = {
}
@dataclass
class Timer:
def downtime_ms(saved_at: datetime) -> int:
"""Сколько занятие пролежало в снимке: запущенный таймер считает и это время."""
if saved_at.tzinfo is None:
saved_at = saved_at.replace(tzinfo=UTC)
return max(0, int((now_utc() - saved_at).total_seconds() * 1000))
class Timer(BaseModel):
"""`started_at` — monotonic-отметка процесса, в другом процессе она ничего
не значит. Снимок хранит прошедшее время, а загрузка пересчитывает отметку
от своих часов (простой берётся из контекста `downtime_ms`)."""
code: TimerCode
started_at: float | None = None
elapsed_ms: int = 0
attempt: int = 1
stopped: bool = False
@model_serializer(mode="wrap")
def _dump(self, handler: SerializerFunctionWrapHandler) -> dict[str, Any]:
data = handler(self)
data.pop("started_at")
data["elapsed_ms"] = self.current_ms(time.monotonic())
data["started"] = self.started_at is not None
return data
@model_validator(mode="before")
@classmethod
def _restore(cls, data: Any, info: ValidationInfo) -> Any:
if not isinstance(data, dict) or "started" not in data:
return data
data = dict(data)
elapsed = max(0, int(data.get("elapsed_ms", 0)))
stopped = bool(data.get("stopped"))
downtime = 0 if stopped else (info.context or {}).get("downtime_ms", 0)
started = data.pop("started")
data["started_at"] = time.monotonic() - (elapsed + downtime) / 1000 if started else None
data["elapsed_ms"] = elapsed if stopped else 0
return data
def start(self, now: float) -> None:
if self.stopped:
# Повторный запуск после остановки — это новая попытка (обратный дозвон).
@ -85,22 +126,21 @@ class Timer:
return int((now - self.started_at) * 1000)
@dataclass
class SessionTimers:
class SessionTimers(BaseModel):
"""Набор таймеров одной сессии. `limits` приходит из конфига —
норматив меняется значением, а не правкой кода."""
limits: dict[TimerCode, int] = field(
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)
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)
self.timers.setdefault(code, Timer(code=code)).start(now)
for code in STOPS.get(event_type, ()):
self.timers.setdefault(code, Timer(code)).stop(now)
self.timers.setdefault(code, Timer(code=code)).stop(now)
def snapshot(self, now: float | None = None) -> list[TimerSnapshot]:
"""Только запущенные таймеры: показывать нули по нормативам,