Merge remote-tracking branch 'origin/main' into feat/cluster-auth-handshake

# Conflicts:
#	docs/arch/SCALE-OUT.md
#	scripts/test_cluster_failover.py
This commit is contained in:
gglamer 2026-09-27 20:48:17 +00:00
commit 240848a68d
58 changed files with 5789 additions and 152 deletions

View file

@ -48,3 +48,14 @@ async def signs(prefix: list[str] = Query(default=[]), group: int | None = None)
else None
),
}
@router.get("/services")
async def services() -> dict:
"""Каталог служб для ручного добавления в карточку — как окно «Добавьте службы».
Каталог — справочник боевого АРМ, а не содержимое вызова: отдавать его
курсанту так же безопасно, как группы происшествий.
"""
data = ekp.catalog()
return {"source": data.source, "services": [item.model_dump() for item in data.services]}

View file

@ -8,6 +8,7 @@ import logging
import time
from collections.abc import AsyncIterator
from datetime import UTC, datetime
from typing import Literal
from uuid import UUID
from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response
@ -22,16 +23,29 @@ from app.db import repo
from app.db.base import get_session
from app.db.models import Group, Score, Session, Trainee
from app.domain.events import Exercise, SessionMode, SessionReport
from app.domain.taxonomy import (
ERRORS,
ErrorCode,
Finding,
FindingDecision,
FindingReview,
FindingSource,
)
from app.domain.roles import Role
from app.domain.statuses import SERVICE_STATUS_LABELS, StationSnapshot, current
from app.domain.timers import TimerCode
from app.domain.statuses import SERVICE_STATUS_LABELS, DdsQueueCard, StationSnapshot, current
from app.domain.timers import TimerCode, TimerState
from app.scoring.export import to_csv, to_pdf
from app.scoring.report import build as build_report
from app.scoring.review import has_finding
from app.scoring.taxonomy import METRIC_MAP
from app.session.access import can_access
from app.session.checkpoint import load_state
from app.session.finish import override_score
from app.session.finish import change_findings, override_score
from app.session.store import FindingAdded, FindingReviewed
from app.session.timers import now_utc
from app.session.score import scoring_scenario
from app.session.hub import hub
from app.session.state import SessionState
from app.session.store import ScoreOverridden, apply_score_override
from app.voice.recording import recording_path
@ -86,6 +100,15 @@ class DdsHistoryOut(BaseModel):
recipient_services: list[str] = []
class Signal(BaseModel):
"""Одна строка колонки сигналов реестра. Цвет несёт смысл только через
`severity`; текст обязателен и не заменяется цветом."""
kind: Literal["backlog", "refusals", "offline"]
severity: TimerState
text: str
class ActiveSessionOut(BaseModel):
session_id: UUID
trainee_name: str | None
@ -100,7 +123,10 @@ class ActiveSessionOut(BaseModel):
dds_overdue_cards: int
dds_work_overdue_cards: int
dds_statuses: dict[str, str]
paused: bool = False
dds_snapshot: StationSnapshot | None = None
signals: list[Signal] = []
presence_known: bool = True
def _out(session) -> SessionOut:
@ -184,6 +210,52 @@ async def dds_history(
return result
def _signals(
state: SessionState, queue: list[DdsQueueCard], now: datetime, *, live: bool
) -> list[Signal]:
"""Сигналы реестра: очередь, повторные отказы, курсант не на связи.
Не завязаны на паузу занятия (`session.pause`, если появится) — очередь,
отказы и присутствие сокета не читают таймеры. Присутствие — по наличию
подписчика в хабе этого узла; окно нужно, чтобы короткий разрыв
соединения не сразу считался потерей связи курсанта.
`live=False` — снимок чужого узла: его сокеты здесь не видны, поэтому
присутствие неизвестно. Чтение реестра не меняет состояние занятия.
"""
settings = get_settings()
signals: list[Signal] = []
# Необработанная — без первичного статуса: он останавливает норматив
# DDS_ACK. Принятая или отклонённая карточка висит в очереди до
# «Следующей», но реакции курсанта уже не ждёт.
pending = sum(not card.timer_stopped for card in queue)
if pending >= settings.signal_backlog_threshold:
signals.append(Signal(
kind="backlog", severity=TimerState.WARN,
text=f"Очередь: {pending} необработанных карточек",
))
if state.consecutive_refusals >= settings.signal_refusals_threshold:
signals.append(Signal(
kind="refusals", severity=TimerState.WARN,
text=f"Подряд отказов: {state.consecutive_refusals}",
))
if live:
connected = (
hub.station_connected(state.session_id) if state.dds_phase
else hub.trainee_connected(state.session_id)
)
offline_seconds = (
(now - (state.socket_last_seen_at or state.started_at)).total_seconds()
if (state.socket_last_seen_at or state.started_at) is not None else 0
)
if not connected and offline_seconds >= settings.signal_offline_window_seconds:
signals.append(Signal(
kind="offline", severity=TimerState.VIOLATED,
text=f"Курсант не на связи {int(offline_seconds)} с",
))
return sorted(signals, key=lambda signal: signal.severity is not TimerState.VIOLATED)
@router.get("/active", response_model=list[ActiveSessionOut])
async def active(
request: Request,
@ -197,6 +269,9 @@ async def active(
state.session_id: state
for state in hub.active_sessions(who.login)
}
#: Снимки без живой записи в хабе этого узла — чужой узел или сессия,
#: ожидающая восстановления. Её сокеты нельзя проверить этим процессом.
checkpoint_only: set[UUID] = set()
if db is not None:
rows = (
await db.scalars(
@ -225,11 +300,15 @@ async def active(
state.owner_login = row.owner_login
if not state.ended:
states[state.session_id] = state
checkpoint_only.add(state.session_id)
for state in states.values():
elapsed = (max(0, int((now - state.started_at).total_seconds()))
paused_ms = state.total_paused_ms
if state.paused and state.paused_at is not None:
paused_ms += max(0, int((now - state.paused_at).total_seconds() * 1000))
elapsed = (max(0, int((now - state.started_at).total_seconds() - paused_ms / 1000))
if state.started_at else 0)
station = state.station_snapshot() if state.exercise is Exercise.DDS else None
station = state.station_snapshot() if state.dds_phase else None
queue = station.queue_cards if station else []
card = state.desk.active
status_log = card.status_log if card is not None else []
@ -255,13 +334,18 @@ async def active(
),
dds_work_overdue_cards=sum(
(timer := card.timers.timers.get(TimerCode.DDS_WORK)) is not None
and timer.started_at is not None
# Пауза держит `started_at=None`, хотя таймер уже шёл: судить
# по нему одному спрятало бы уже случившееся нарушение (lct-39).
and (timer.started_at is not None or timer.paused)
and not timer.stopped
and timer.current_ms(time.monotonic()) > card.timers.limits[TimerCode.DDS_WORK]
for card in state.desk.cards.values()
),
dds_statuses=latest_statuses,
paused=state.paused,
dds_snapshot=station,
signals=_signals(state, queue, now, live=state.session_id not in checkpoint_only),
presence_known=state.session_id not in checkpoint_only,
))
return result
@ -584,6 +668,157 @@ async def override(
return SessionReport.model_validate(report["full_report"])
def _required_text(value: str) -> str:
cleaned = value.strip()
if not cleaned:
raise ValueError("поле обязательно")
return cleaned
class FindingReviewIn(BaseModel):
"""Решение по отметке: причина обязательна и остаётся в разборе."""
decision: FindingDecision
reason: str = Field(min_length=1, max_length=1000)
_reason = field_validator("reason")(_required_text)
class FindingIn(BaseModel):
"""Отметка преподавателя: код, факт и норма — то же обоснование, что у автоматической."""
code: ErrorCode
fact: str = Field(min_length=1, max_length=1000)
norm: str = Field(min_length=1, max_length=1000)
#: Номер карточки очереди ДДС (с 1); без него отметка относится к занятию.
card: int | None = Field(default=None, ge=1)
#: Ключ идемпотентности с клиента: повтор запроса не добавит отметку второй раз.
client_id: UUID | None = None
_fact = field_validator("fact")(_required_text)
_norm = field_validator("norm")(_required_text)
async def _change_findings(
session_id: UUID, who, db: AsyncSession | None, make,
) -> SessionReport:
"""Общий путь решений по отметкам — тот же, что у правки итога.
`make(who, score)` строит запись изменения по сохранённой оценке или
отвечает HTTP-ошибкой. Роль проверяет маршрут: администратору сюда нельзя,
как и к правке балла.
"""
state = hub.get(session_id)
if state is not None:
_require_access(who, state)
if state.score is None:
raise HTTPException(status_code=409, detail="score_not_ready")
scenario = scoring_scenario(state)
if scenario is None:
raise HTTPException(status_code=409, detail="scenario_not_found")
change = make(who, state.score)
async with hub.operation(session_id):
if not _repeated(state.score, change):
change_findings(state, change)
return build_report(session_id, state, scenario)
if db is None:
raise HTTPException(status_code=404, detail="session_not_found")
session = await repo.get_session(db, session_id)
if session is None:
raise HTTPException(status_code=404, detail="session_not_found")
_require_access(who, session)
score = await db.scalar(select(Score).where(Score.session_id == session_id))
if score is None:
raise HTTPException(status_code=409, detail="score_not_ready")
if (score.report or {}).get("full_report") is None:
raise HTTPException(status_code=409, detail="report_not_archived")
current = dict(score.report)
change = make(who, current)
await db.rollback() # чтение закончено; запись — одной транзакцией хранилища
if not _repeated(current, change):
await hub.store.commit_archived(session_id, [change])
# Ответ — по строке после commit: пересчёт шёл под блокировкой и мог учесть
# решение из соседней вкладки, которого не было в прочитанном выше отчёте.
score = await db.scalar(
select(Score).where(Score.session_id == session_id)
.execution_options(populate_existing=True)
)
report = SessionReport.model_validate(score.report["full_report"])
await db.rollback()
return report
def _repeated(report: dict, change: FindingReviewed | FindingAdded) -> bool:
return isinstance(change, FindingAdded) and has_finding(report, change.finding.client_id)
@router.post("/{session_id}/findings/{index}/review", response_model=SessionReport)
async def review_finding(
session_id: UUID,
index: int,
body: FindingReviewIn,
request: Request,
db: AsyncSession | None = Depends(optional_session),
) -> SessionReport:
"""Подтвердить или снять отметку; снятая не штрафует свою метрику."""
who = require(request, Role.INSTRUCTOR)
def make(who, score: dict) -> FindingReviewed:
if not 0 <= index < len(score.get("findings", [])):
raise HTTPException(status_code=404, detail="finding_not_found")
return FindingReviewed(index=index, role=who.role.value, review=FindingReview(
decision=body.decision, reason=body.reason, author=who.login, at=now_utc(),
))
return await _change_findings(session_id, who, db, make)
@router.post("/{session_id}/findings", response_model=SessionReport)
async def add_finding(
session_id: UUID,
body: FindingIn,
request: Request,
db: AsyncSession | None = Depends(optional_session),
) -> SessionReport:
"""Отметка преподавателя штрафует связанную метрику по той же методике."""
who = require(request, Role.INSTRUCTOR)
def make(who, score: dict) -> FindingAdded:
service = None
prefix = ""
exercise = (score.get("full_report") or {}).get("exercise")
# E1–E6 в связке 112 → ДДС — ошибки приёма вызова: их метрики у занятия,
# а не у карточки очереди. С номером карточки отметка не нашла бы свою
# метрику 112 и штрафовала бы отдельной метрикой преподавателя.
if (body.card is not None and body.code.value.startswith("E")
and exercise != Exercise.DDS.value):
raise HTTPException(status_code=422, detail="call_code_without_card")
if body.card is not None:
cards = score.get("card_results") or []
if body.card > len(cards):
raise HTTPException(status_code=422, detail="card_not_found")
service = cards[body.card - 1].get("managed_service")
prefix = f"Карточка {body.card}: " + (f"{service}: " if service else "")
competency = next((item for code, item in METRIC_MAP.values() if code is body.code), None)
return FindingAdded(role=who.role.value, finding=Finding(
code=body.code,
source=FindingSource.INSTRUCTOR,
summary=f"{prefix}{ERRORS[body.code].title}",
fact=body.fact,
norm=body.norm,
ref="отметка преподавателя на разборе",
competency=competency,
at=now_utc(),
card=body.card,
service=service,
author=who.login,
client_id=body.client_id,
))
return await _change_findings(session_id, who, db, make)
@router.get("", response_model=list[SessionOut])
async def listing(
request: Request,

View file

@ -14,7 +14,7 @@ from uuid import UUID, uuid4
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
from pydantic import TypeAdapter, ValidationError
from app.api.ws.session import pump, run_command, session_socket
from app.api.ws.session import close_fenced, pump, run_command, session_socket
from app.domain.events import (
BgStart,
CallEnded,
@ -31,6 +31,7 @@ from app.domain.events import (
PatchSource,
ScoreReady,
SessionMode,
SessionPaused,
StationState,
TimerTick,
TextTurnAccepted,
@ -39,6 +40,7 @@ from app.domain.events import (
TraineeToServer,
)
from app.dialog.slots import TurnResult
from app.domain.kio import PatchRejected, get_field
from app.domain.roles import Role
from app.scenarios import store
from app.session.dds import prepare_handoff_queue
@ -191,6 +193,14 @@ async def _handle(session_id: UUID, state, event) -> None:
code=ErrorKind.UNSUPPORTED_EVENT, message="Занятие уже завершено",
))
return
if state.paused:
# Отказ временный, в отличие от отклонённого значения: `kio.patch`
# без `auto`-отката. useCall держит правку в pending и повторяет её
# после `session.paused: false` — откат стёр бы ввод курсанта.
hub.to_trainee(session_id, ErrorEvent(
code=ErrorKind.UNSUPPORTED_EVENT, message="Пауза, ждите преподавателя",
))
return
if (event.type == "text.turn" and state.exercise is not Exercise.CARD) or state.exercise is Exercise.DDS or (
state.exercise is Exercise.CARD and event.type not in {"kio.patch", "card.submit", "text.turn"}
) or (state.exercise is Exercise.CALL and event.type == "card.submit"):
@ -229,6 +239,9 @@ async def _handle(session_id: UUID, state, event) -> None:
code=ErrorKind.INTERNAL, message="Не удалось получить ответ заявителя. Попробуйте ещё раз.",
))
return
if state.paused:
# Ответ модели мог закончиться уже после команды преподавателя.
return
caller_entry = state.append(Speaker.CALLER, line.text, line.mood)
hub.to_trainee(session_id, CallerUtterance(
utterance_id=uuid4(), text=line.text,
@ -270,16 +283,37 @@ async def _handle(session_id: UUID, state, event) -> None:
case "kio.patch":
old_code, old_notify = state.kio.incident_code, list(state.kio.notify)
state.patch_kio(event.fields)
old_extra = list(state.kio.notify_extra)
try:
state.patch_kio(event.fields)
except (ValidationError, PatchRejected) as exc:
reason = exc.errors()[0].get("msg", "") if isinstance(exc, ValidationError) else str(exc)
# Недопустимое значение — отказ курсанту, а не сбой операции:
# сбой закрыл бы занятие на этом узле.
hub.to_trainee(session_id, ErrorEvent(
code=ErrorKind.UNSUPPORTED_EVENT,
message=f"Поле карточки не принято: {reason}"[:200],
))
# Пакет отклонён целиком, а фронт держит его в pending до эха.
# Серверное значение с `auto` снимает pending — иначе на экране
# останутся правки, которых на сервере нет.
hub.to_trainee(session_id, KioPatchOut(
fields={path: get_field(state.kio, path) for path in event.fields},
source=PatchSource.AUTO,
))
return
hub.to_trainee(session_id, KioPatchOut(
fields=event.fields, source=PatchSource.OPERATOR,
))
auto: dict = {}
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,
))
auto |= {"incident_code": state.kio.incident_code, "notify": list(state.kio.notify)}
# Добавки сервер чистит сам (повторы, службы из notify) — фронту
# нужно итоговое значение, а не эхо присланного.
if state.kio.notify_extra != event.fields.get("notify_extra", old_extra):
auto["notify_extra"] = list(state.kio.notify_extra)
if auto:
hub.to_trainee(session_id, KioPatchOut(fields=auto, source=PatchSource.AUTO))
# Наблюдателю уходит карточка целиком: рассинхрон на внешнем мониторе
# посреди занятия дороже лишних килобайт.
hub.to_observers(session_id, KioState(kio=state.kio))
@ -390,9 +424,16 @@ async def call(ws: WebSocket, session_id: UUID) -> None:
)
if entered is None:
return
_who, state = entered
who, state = entered
with hub.trainee(session_id) as queue:
async with hub.trainee_socket(
session_id, station=False, trainee=who.role is Role.TRAINEE,
) as queue:
if hub.is_lease_fenced(session_id):
await close_fenced(ws)
return
# Переподключившийся клиент сначала узнаёт, разрешён ли ввод.
queue.put_nowait(SessionPaused(paused=state.paused))
if state.exercise is Exercise.CARD:
from app.api.ws.control import card_briefing
@ -435,6 +476,10 @@ async def call(ws: WebSocket, session_id: UUID) -> None:
# Бинарные кадры — аудио, текстовые — события. Направление определяется
# каналом, обёртки JSON вокруг звука нет (docs/arch/CONTRACT.md).
if message.get("bytes") is not None:
if state.paused:
# JSON-команды отказывает `_handle`; бинарные кадры сюда
# не заходят, поэтому пауза глушит звук здесь же (lct-39).
continue
if state.exercise is Exercise.CALL:
_on_audio(session_id, state, message["bytes"])
else:

View file

@ -41,6 +41,7 @@ from app.domain.events import (
ModeSet,
ReferenceStarted,
SessionEnded,
SessionPaused,
StationState,
)
from app.domain.roles import Role
@ -51,7 +52,14 @@ from app.session.access import can_access
from app.session.hub import hub
from app.session.finish import end_session, override_score
from app.session.state import SessionState, now_utc
from app.session.store import LessonIdentity, LessonRequest, NoteAdded, ScoreOverridden
from app.session.store import (
LessonIdentity,
LessonPaused,
LessonRequest,
LessonResumed,
NoteAdded,
ScoreOverridden,
)
from app.voice.models import get_voice_models
from app.voice.pipeline import FILLERS, prefetch
@ -104,6 +112,7 @@ def _build_state(session_id: UUID, event, who, scenario, scenarios,
dds_service=identity.service or event.dds_service,
attempt=identity.attempt,
criteria=event.criteria,
socket_last_seen_at=now_utc(),
)
state.timers.limits[TimerCode.DDS_ACK] = event.criteria.decision_time_limit_seconds * 1000
state.timers.limits[TimerCode.CARD_FILL] = event.criteria.card_fill_time_limit_seconds * 1000
@ -344,6 +353,26 @@ async def _command(session_id: UUID, event, who) -> None:
await _start(session_id, event, who)
case "session.stop":
await _stop(session_id)
case "session.pause":
if state is None or state.ended or state.paused:
return
state.pause()
hub.record(session_id, LessonPaused(
at=state.paused_at, author=who.login, role=who.role.value,
))
hub.broadcast(session_id, SessionPaused(paused=True))
hub.to_station(session_id, SessionPaused(paused=True))
case "session.resume":
if state is None or state.ended or not state.paused:
return
paused_ms_before = state.total_paused_ms
state.resume()
hub.record(session_id, LessonResumed(
at=now_utc(), author=who.login, role=who.role.value,
paused_ms=state.total_paused_ms - paused_ms_before,
))
hub.broadcast(session_id, SessionPaused(paused=False))
hub.to_station(session_id, SessionPaused(paused=False))
case "instructor_note.add":
if state is not None:
state.notes.append({
@ -390,6 +419,14 @@ async def _command(session_id: UUID, event, who) -> None:
case "director.inject":
if state is None:
return
if state.paused:
# Обрыв на паузе запустил бы норматив обратного дозвона, и простой
# ушёл бы в него; курсант же за баннером паузы не может ответить.
hub.to_observers(session_id, ErrorEvent(
code=ErrorKind.UNSUPPORTED_EVENT,
message="Занятие на паузе",
))
return
if state.exercise is not Exercise.CALL:
hub.to_observers(session_id, ErrorEvent(
code=ErrorKind.UNSUPPORTED_EVENT,

View file

@ -17,7 +17,7 @@ from uuid import UUID
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
from pydantic import TypeAdapter, ValidationError
from app.api.ws.session import pump, run_command, session_socket
from app.api.ws.session import close_fenced, pump, run_command, session_socket
from app.domain.events import (
CallEndReason,
ErrorEvent,
@ -39,6 +39,9 @@ def _error(session_id: UUID, message: str) -> None:
async def _handle(session_id: UUID, state, event) -> None:
if state.paused:
_error(session_id, "Пауза, ждите преподавателя")
return
outcome = state.desk.apply(event, state)
if outcome.error is not None:
_error(session_id, outcome.error)
@ -59,9 +62,14 @@ async def station(ws: WebSocket, session_id: UUID, role: str = "dds") -> None:
entered = await session_socket(ws, session_id, (Role.INSTRUCTOR, Role.TRAINEE))
if entered is None:
return
_who, state = entered
who, state = entered
with hub.station(session_id) as queue:
async with hub.trainee_socket(
session_id, station=True, trainee=who.role is Role.TRAINEE,
) as queue:
if hub.is_lease_fenced(session_id):
await close_fenced(ws)
return
sender = asyncio.create_task(pump(ws, queue))
try:
# Карточка, переданная до подключения станции, не теряется: