Merge remote-tracking branch 'origin/main' into feat/dds-timers

# Conflicts:
#	backend/app/scoring/dispatcher.py
#	backend/app/scoring/export.py
#	backend/app/session/finish.py
#	docs/arch/CONTRACT.md
#	docs/spec/GAP.md
#	frontend/src/features/instructor/ScoreWeightsEditor.tsx
This commit is contained in:
gglamer 2026-09-27 20:56:09 +00:00
commit 4abdf0a901
61 changed files with 4084 additions and 167 deletions

View file

@ -28,7 +28,7 @@ from argon2 import PasswordHasher
from argon2.exceptions import VerifyMismatchError
from fastapi import APIRouter, HTTPException, Request, WebSocket
from pydantic import BaseModel, Field
from sqlalchemy import select
from sqlalchemy import select, update
from sqlalchemy.exc import IntegrityError
from starlette.websockets import WebSocketDisconnect
@ -56,6 +56,26 @@ _vanished: set[str] = set()
AUTH_GENERATION_SYNC_SECONDS = 1.0
AUTH_GENERATION_MAX_AGE_SECONDS = 2.0
_generations_synced_at: float | None = None
# Один lock на ещё неизвестный узлу логин: параллельные HTTP и WS входы
# разделяют один запрос к БД до очередной синхронизации (lct-42). Запись
# удаляется, когда lock больше никто не ждёт: иначе словарь растёт с каждым
# логином, который узел когда-либо проверял.
_lookup_locks: dict[str, "_LoginLookup"] = {}
# Логины, которых разовая проверка не нашла в users. Держатся до следующей
# сверки: повтор старой cookie не должен давать SELECT на каждый запрос.
_missing_logins: set[str] = set()
class _LoginLookup:
__slots__ = ("lock", "users")
def __init__(self) -> None:
self.lock = asyncio.Lock()
self.users = 0
class _AuthStateUnavailable(Exception):
"""Разовая проверка неизвестного узлу логина не смогла обратиться к БД."""
async def _close_revoked(ws: WebSocket) -> None:
@ -98,6 +118,7 @@ def prime_generations(values: dict[str, int]) -> None:
global _generations_synced_at
_generations.clear()
_generations.update(values)
_missing_logins.clear()
_generations_synced_at = time.monotonic()
@ -108,15 +129,38 @@ async def load_generations() -> None:
async def sync_generations() -> None:
"""Refresh shared account epochs and close sockets revoked on peer nodes."""
"""Сверить версии полномочий и закрыть отозванные на другом узле сокеты."""
known_before_query = set(_generations)
async with get_sessionmaker()() as db:
rows = (await db.execute(select(User.login, User.auth_version))).all()
current = {login: version for login, version in rows}
current = {login: version for login, version in rows}
# Версия узла выше БД, если отзыв не удалось записать (неудачный
# logout) или учётку пересоздали после удаления. Без записи в БД
# узлы расходятся навсегда: соседний узел выдаёт cookie со старой
# версией, а этот её отвергает. Поднимаем БД до версии узла — отзыв
# сохраняется, остальные узлы догоняют за одну сверку. При гонке со
# снимком БД уже выше, и UPDATE ничего не меняет.
ahead = {
login: _generations[login]
for login, version in current.items()
if login in _generations and _generations[login] > version
}
if ahead:
try:
for login, local in ahead.items():
await db.execute(_raise_auth_version(login, local))
await db.commit()
except Exception as exc: # noqa: BLE001 — повторим на следующей сверке
log.error("не удалось записать версию полномочий узла (%s)",
type(exc).__name__)
await db.rollback()
for login, version in current.items():
previous = _generations.get(login)
# Снимок БД мог устареть за время запроса; меньшая версия не должна
# отменять локальный отзыв, который уже поднял поколение.
if previous is None:
_generations[login] = version
elif previous != version:
elif version > previous:
invalidate_login(login, version)
# Account deletion is not exposed by the application. Still close active
# sockets if an operator removes one directly from the shared directory DB.
@ -127,13 +171,88 @@ async def sync_generations() -> None:
# lookup fall back to 0 and re-accept cookies issued before the revocation.
synthetic = {"dev"} if get_settings().dev_auth_bypass else set()
_vanished.intersection_update(_generations.keys() - current.keys())
for login in _generations.keys() - current.keys() - synthetic - _vanished:
# Разовая проверка могла найти учётку после снимка этого запроса.
# Её отсутствие в старом снимке не означает отзыв.
for login in known_before_query - current.keys() - synthetic - _vanished:
invalidate_login(login)
_vanished.add(login)
# Снимок только что прочитан: учётка, созданная до него, уже в кэше.
_missing_logins.clear()
global _generations_synced_at
_generations_synced_at = time.monotonic()
def _raise_auth_version(login: str, version: int):
# Только вверх: параллельная запись в БД могла уже поднять версию выше.
return (
update(User)
.where(User.login == login, User.auth_version < version)
.values(auth_version=version)
)
async def _resolve_unknown_login(login: str) -> int | None:
"""Проверить неизвестный узлу логин сразу, не дожидаясь опроса БД.
Отозванный логин остаётся в `_generations` (см. `_vanished`), поэтому
отсутствие в кэше означает, что этот узел ещё не видел учётку.
Возвращает версию или None, если строки в `users` нет. При недоступной
или зависшей БД вызывает `_AuthStateUnavailable`, чтобы вход был закрыт
с 503/1013.
"""
if login in _missing_logins:
return None
entry = _lookup_locks.get(login)
if entry is None:
entry = _lookup_locks[login] = _LoginLookup()
entry.users += 1
try:
# Тот же предел, что у сверки: при partition handshake не ждёт
# таймаута TCP, а очередь за lock не растягивает ожидание сверх него.
async with asyncio.timeout(AUTH_GENERATION_MAX_AGE_SECONDS):
async with entry.lock:
cached = _generations.get(login)
if cached is not None:
return cached # параллельный запрос уже получил версию
if login in _missing_logins:
return None
async with get_sessionmaker()() as db:
version = await db.scalar(
select(User.auth_version).where(User.login == login)
)
# Пока шёл SELECT, локальный отзыв или синхронизация могли
# записать новую версию. Старый ответ не должен вернуть
# отозванную cookie.
cached = _generations.get(login)
if cached is not None:
return cached
if version is None:
_missing_logins.add(login)
return None
_generations[login] = version
log.warning(
"учётка %s найдена разовой проверкой до сверки узла",
login_log_marker(login),
)
return version
except Exception as exc: # noqa: BLE001 — ошибка или таймаут БД закрывают вход
log.error("разовая проверка полномочий не удалась (%s)", type(exc).__name__)
raise _AuthStateUnavailable from exc
finally:
entry.users -= 1
if entry.users == 0 and _lookup_locks.get(login) is entry:
del _lookup_locks[login]
def login_log_marker(login: str) -> str:
"""Метка логина для журнала без самого логина.
По ней кластерный смоук доказывает, что узел прошёл через разовую
проверку, а не увидел учётку обычной сверкой.
"""
return hashlib.sha256(f"lct-login:{login}".encode()).hexdigest()[:12]
async def watch_generations() -> None:
"""Poll PostgreSQL once per node so remote logout/role changes close WS."""
while True:
@ -151,6 +270,21 @@ async def watch_generations() -> None:
await asyncio.sleep(AUTH_GENERATION_SYNC_SECONDS)
async def _send_auth_state_unavailable(scope, send) -> None:
if scope["type"] == "websocket":
await send({"type": "websocket.close", "code": 1013})
else:
await send({
"type": "http.response.start",
"status": 503,
"headers": [(b"content-type", b"application/json")],
})
await send({
"type": "http.response.body",
"body": b'{"detail":"auth_state_unavailable"}',
})
class AuthVersionMiddleware:
"""Check signed-cookie epochs against the fresh, DB-synchronized node cache."""
@ -186,22 +320,17 @@ class AuthVersionMiddleware:
synced_at is None
or time.monotonic() - synced_at > AUTH_GENERATION_MAX_AGE_SECONDS
):
if scope["type"] == "websocket":
await send({"type": "websocket.close", "code": 1013})
else:
await send({
"type": "http.response.start",
"status": 503,
"headers": [(b"content-type", b"application/json")],
})
await send({
"type": "http.response.body",
"body": b'{"detail":"auth_state_unavailable"}',
})
await _send_auth_state_unavailable(scope, send)
return
cookie_version = session.get("auth_generation")
version = _generations.get(login)
if version is None:
try:
version = await _resolve_unknown_login(login)
except _AuthStateUnavailable:
await _send_auth_state_unavailable(scope, send)
return
if version is None or cookie_version != version:
if session is not None:
session.clear()
@ -536,7 +665,20 @@ async def login(payload: LoginIn, request: Request) -> dict:
await audit_required(user.login, user.role, "login.blocked")
raise HTTPException(status_code=403, detail="blocked")
if _generations.get(user.login) != user.auth_version:
local_version = _generations.get(user.login)
if local_version is not None and local_version > user.auth_version:
# Отзыв на этом узле не дошёл до БД (неудачный logout). Сброс к версии
# БД вернул бы силу cookie, выданной до выхода; выдача cookie с
# версией узла без записи в БД развела бы узлы. Поднимаем БД.
try:
async with get_sessionmaker()() as db:
await db.execute(_raise_auth_version(user.login, local_version))
await db.commit()
except Exception as exc: # noqa: BLE001 — cookie без записанной версии не выдаём
log.error("вход: версию полномочий не удалось записать (%s)",
type(exc).__name__)
raise HTTPException(status_code=503, detail="auth_state_unavailable") from exc
elif local_version != user.auth_version:
invalidate_login(user.login, user.auth_version)
who = Principal(

View file

@ -23,6 +23,14 @@ 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,
@ -34,9 +42,13 @@ from app.domain.statuses import (
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
@ -117,6 +129,7 @@ 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
@ -298,7 +311,10 @@ async def active(
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.dds_phase else None
queue = station.queue_cards if station else []
@ -326,12 +342,15 @@ 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,
@ -657,6 +676,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

@ -31,6 +31,7 @@ from app.domain.events import (
PatchSource,
ScoreReady,
SessionMode,
SessionPaused,
StationState,
TimerTick,
TextTurnAccepted,
@ -192,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"):
@ -230,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,
@ -420,6 +432,8 @@ async def call(ws: WebSocket, session_id: UUID) -> None:
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
@ -462,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
@ -345,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({
@ -391,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

@ -3,9 +3,9 @@
Учебный исходящий звонок бригаде имитируется событиями JSON. Для целевого
локального прототипа не нужна внешняя телефония или аудиомодель.
Самая ценная механика цепочки — `card.bounce`: диспетчер видит, что не указан
этаж, и отбивает карточку обратно. Неполнота КИО перестаёт быть процентом
в отчёте и становится сорванным выездом с конкретной причиной.
`card.bounce` сервер отклоняет: ДДС не возвращает карточку в 112 при принятии
и не правит её поля. Об ошибке, которую вскрыл доклад бригады с места,
диспетчер сообщает в 112 по телефону (ответ заказчика, П.3).
Правила пульта — в `DdsDesk.apply`; здесь разбор команды и рассылка итога.
Отсев повторов и операция хранилища — в `run_command`.
@ -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)