Merge branch 'feat/session-signals'

Сигналы реестра занятий: очередь, повторные отказы, курсант не на связи
This commit is contained in:
GGlamer 2026-09-27 23:04:21 +03:00 • committed by GitHub
commit 3bebdbb803
12 changed files with 737 additions and 16 deletions

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
@ -23,8 +24,8 @@ 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.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.session.access import can_access
@ -32,6 +33,7 @@ from app.session.checkpoint import load_state
from app.session.finish import override_score
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 +88,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
@ -101,6 +112,8 @@ class ActiveSessionOut(BaseModel):
dds_work_overdue_cards: int
dds_statuses: dict[str, str]
dds_snapshot: StationSnapshot | None = None
signals: list[Signal] = []
presence_known: bool = True
def _out(session) -> SessionOut:
@ -184,6 +197,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 +256,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 +287,12 @@ 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()))
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 []
@ -262,6 +325,8 @@ async def active(
),
dds_statuses=latest_statuses,
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

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,
@ -390,9 +390,14 @@ 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
if state.exercise is Exercise.CARD:
from app.api.ws.control import card_briefing

View file

@ -104,6 +104,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

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,
@ -59,9 +59,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:
# Карточка, переданная до подключения станции, не теряется:

View file

@ -120,6 +120,11 @@ class Settings(BaseSettings):
def limit_ms(self, code: TimerCode) -> int:
return self.timer_limits_ms.get(code, NORMATIVES[code].limit_ms)
# Пороги сигналов реестра живых занятий (пульт преподавателя).
signal_backlog_threshold: int = Field(default=4, ge=1)
signal_refusals_threshold: int = Field(default=2, ge=1)
signal_offline_window_seconds: int = Field(default=45, ge=1)
def validate_deployment_security(self) -> None:
"""Reject known development authentication defaults on a production app."""
if self.app_env != "production":

View file

@ -446,6 +446,7 @@ class DdsDesk(BaseModel):
services[0], ServiceStatus.ACCEPTED, command.comment,
author="диспетчер",
)
session.consecutive_refusals = 0
except StatusError:
pass # статус уже стоит: повторное нажатие ничего не меняет
case "card.status":
@ -462,6 +463,13 @@ class DdsDesk(BaseModel):
# Первичный статус останавливает норматив 30 секунд.
if command.status in PRIMARY:
card.on_event("card.ack")
# Сигнал реестра "повторные отказы" считает подряд идущие
# отказы, а не сумму за занятие.
session.consecutive_refusals = (
session.consecutive_refusals + 1
if command.status is ServiceStatus.DECLINED
else 0
)
if command.status in {
ServiceStatus.COMPLETED, ServiceStatus.DECLINED, ServiceStatus.REFUSED,
}:

View file

@ -24,6 +24,7 @@ from app.domain.events import (
)
from app.session.state import SessionState
from app.session.store import MemorySessionStore, Record, SessionLeaseLost, SessionStore
from app.session.timers import now_utc
#: Очередь одного подписчика. Медленный наблюдатель не тормозит занятие:
#: очередь ограничена, переполнение роняет соединение, а не сессию.
@ -53,6 +54,9 @@ class SessionHub:
self._observers: dict[UUID, set[asyncio.Queue]] = {}
self._trainees: dict[UUID, set[asyncio.Queue]] = {}
self._stations: dict[UUID, set[asyncio.Queue]] = {}
self._trainee_calls: dict[UUID, set[asyncio.Queue]] = {}
self._trainee_stations: dict[UUID, set[asyncio.Queue]] = {}
self._pending_adoptions: dict[UUID, SessionState] = {}
self._tickers: dict[UUID, asyncio.Task] = {}
# Задачи, порождённые внутри операции, наследуют контекст; после
# закрытия операции их события идут напрямую (`_Operation.open`).
@ -116,6 +120,7 @@ class SessionHub:
def drop(self, session_id: UUID) -> None:
self._sessions.pop(session_id, None)
self._pending_adoptions.pop(session_id, None)
self.stop_ticker(session_id)
def live_count(self) -> int:
@ -135,9 +140,10 @@ class SessionHub:
if not self.store.persistent:
return 0
restored = await self.store.restore_active()
adopted = 0
for state in restored:
self._adopt(state)
return len(restored)
adopted += await self._adopt(state)
return adopted
async def maintain_lease(self) -> None:
"""Один оборот супервизора: продлить свои lease, подхватить просроченные чужие.
@ -154,6 +160,18 @@ class SessionHub:
log.warning("продление lease занятия %s не удалось", state.session_id,
exc_info=True)
await self.fence(state)
for state in list(self._pending_adoptions.values()):
try:
await self.store.commit(state)
except SessionLeaseLost:
self._pending_adoptions.pop(state.session_id, None)
except Exception:
log.warning("повторное сохранение занятия %s после takeover не удалось",
state.session_id, exc_info=True)
else:
self._pending_adoptions.pop(state.session_id, None)
self.register(state)
self.start_ticker(state.session_id)
try:
claimed = await self.store.claim_expired()
except Exception: # следующий оборот попробует снова
@ -163,17 +181,32 @@ class SessionHub:
current = self._sessions.get(state.session_id)
if current is not None and not current.lease_fenced:
continue
self._adopt(state)
await self._adopt(state)
async def supervise_lease(self, interval: float) -> None:
while True:
await asyncio.sleep(interval)
await self.maintain_lease()
def _adopt(self, state: SessionState) -> None:
async def _adopt(self, state: SessionState) -> bool:
self.stop_ticker(state.session_id)
if state.socket_connected_at_checkpoint:
# Разрыв произошёл при потере узла; старое время подключения не
# доказывает, что курсант отсутствовал всё это время.
state.socket_last_seen_at = now_utc()
state.socket_connected_at_checkpoint = False
try:
await self.store.commit(state)
except SessionLeaseLost:
return False
except Exception:
log.exception("не удалось сохранить присутствие после takeover %s", state.session_id)
self._pending_adoptions[state.session_id] = state
return False
self._pending_adoptions.pop(state.session_id, None)
self.register(state)
self.start_ticker(state.session_id)
return True
async def save_all(self) -> None:
"""Снимок живых занятий при остановке узла: следующий владелец продолжит с него."""
@ -285,6 +318,44 @@ class SessionHub:
def station(self, session_id: UUID):
return self._subscribe(self._stations, session_id)
async def _record_presence(self, session_id: UUID, *, station: bool) -> None:
state = self.get(session_id)
if state is None or state.ended or state.dds_phase != station:
return
async with self.operation(session_id):
state.socket_last_seen_at = now_utc()
present = self._trainee_stations if station else self._trainee_calls
state.socket_connected_at_checkpoint = bool(present.get(session_id))
@contextlib.asynccontextmanager
async def trainee_socket(
self, session_id: UUID, *, station: bool, trainee: bool,
) -> AsyncIterator[asyncio.Queue]:
"""Учесть только сокет курсанта; вещание преподавателю остаётся общим."""
registry = self._stations if station else self._trainees
presence = self._trainee_stations if station else self._trainee_calls
with self._subscribe(registry, session_id) as queue:
if trainee:
presence.setdefault(session_id, set()).add(queue)
try:
if trainee:
try:
await self._record_presence(session_id, station=station)
except Exception:
if not self.is_lease_fenced(session_id):
raise
yield queue
finally:
if trainee:
presence[session_id].discard(queue)
if not presence[session_id]:
presence.pop(session_id)
try:
await self._record_presence(session_id, station=station)
except Exception:
if not self.is_lease_fenced(session_id):
raise
# ── вещание ──
@staticmethod
@ -314,6 +385,12 @@ class SessionHub:
def observer_count(self, session_id: UUID) -> int:
return len(self._observers.get(session_id, set()))
def station_connected(self, session_id: UUID) -> bool:
return bool(self._trainee_stations.get(session_id))
def trainee_connected(self, session_id: UUID) -> bool:
return bool(self._trainee_calls.get(session_id))
# ── такт таймеров ──
def start_ticker(self, session_id: UUID) -> None:

View file

@ -124,6 +124,15 @@ class PersistedSession(BaseModel):
#: 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)
#: Подряд идущих первичных отказов без принятой карточки между ними —
#: сигнал реестра преподавателя, не влияет на балл.
consecutive_refusals: int = 0
#: Последнее подключение/отключение сокета курсанта или начало фазы ДДС.
#: Сохраняется вместе с занятием, чтобы чтение реестра не меняло состояние.
socket_last_seen_at: datetime | None = None
#: На момент последнего checkpoint курсант был подключён на текущем АРМ.
#: После аварии узла момент потери связи неизвестен: окно начинается при takeover.
socket_connected_at_checkpoint: bool = False
@field_validator("processed_station_commands")
@classmethod
@ -206,6 +215,9 @@ class SessionState(PersistedSession):
self.dispatched_at = now_utc()
if self.exercise is Exercise.CALL:
self._receive_call_card()
if self.handoff_to_dds:
self.socket_last_seen_at = self.dispatched_at
self.socket_connected_at_checkpoint = False
return self.dispatched_card
def _receive_call_card(self) -> None:

View file

@ -161,6 +161,9 @@ def full_state() -> SessionState:
resolve_comment="передано в другой регион",
processed_station_commands=[str(uuid4())],
text_revealed_facts={"f_address": "улица Ленина, 14"},
consecutive_refusals=2,
socket_last_seen_at=AT,
socket_connected_at_checkpoint=True,
)

View file

@ -0,0 +1,517 @@
"""Колонка сигналов реестра: очередь, повторные отказы, курсант не на связи."""
import asyncio
import time
from datetime import UTC, datetime, timedelta
from uuid import uuid4
import pytest
from fastapi.testclient import TestClient
from sqlalchemy.engine.result import IteratorResult, SimpleResultMetaData
from starlette.websockets import WebSocketDisconnect
from app.api.auth import Principal
from app.api.http import sessions as sessions_http
from app.api.ws import session as session_ws
from app.config import get_settings
from app.db.models import Session
from app.domain.events import Exercise, SessionMode
from app.domain.roles import Role
from app.main import app
from app.session.checkpoint import dump_state, load_state
from app.session import hub as hub_module
from app.session.hub import SessionHub, hub
from app.session.state import SessionState
from app.session.store import MemorySessionStore
POOL = ["fire-apartment-l2", "t01-1-fire-container"]
@pytest.fixture
def client(monkeypatch):
monkeypatch.setenv("DEV_AUTH_BYPASS", "true")
get_settings.cache_clear()
async def audit_override(*_args, **_kwargs):
return None
async def optional_session_override():
yield None
monkeypatch.setattr(sessions_http, "audit_required", audit_override)
monkeypatch.setitem(
app.dependency_overrides, sessions_http.optional_session, optional_session_override
)
try:
with TestClient(app) as test_client:
test_client.post("/api/auth/dev-token")
hub.store = MemorySessionStore()
before = set(hub._sessions)
try:
yield test_client
finally:
for session_id in set(hub._sessions) - before:
hub.drop(session_id)
finally:
get_settings.cache_clear()
def wait_for(predicate, timeout=3):
end = time.monotonic() + timeout
while time.monotonic() < end:
value = predicate()
if value:
return value
time.sleep(0.02)
raise AssertionError("состояние не обновилось")
def read_until(socket, wanted):
received = []
for _ in range(20):
event = socket.receive_json()
received.append(event["type"])
if event["type"] == wanted:
return event
raise AssertionError(f"событие {wanted} не пришло; получены: {received}")
def start_two_card_dds(client, scenario_ids=POOL):
session_id = uuid4()
context = client.websocket_connect(f"/ws/control/{session_id}")
control = context.__enter__()
control.send_json({
"type": "scenario.start",
"scenario_id": POOL[0],
"trainee": "Иванов",
"mode": "training",
"exercise": "dds",
"random_scenario_ids": scenario_ids,
})
wait_for(lambda: hub.get(session_id))
return session_id, control
def start_card_handoff(client):
session_id = uuid4()
context = client.websocket_connect(f"/ws/control/{session_id}")
control = context.__enter__()
control.send_json({
"type": "scenario.start", "scenario_id": POOL[0], "trainee": "Иванов",
"mode": "training", "exercise": "card", "handoff_to_dds": True,
"scenario_ids": POOL,
})
wait_for(lambda: hub.get(session_id))
return session_id, control
def start_call(client):
session_id = uuid4()
context = client.websocket_connect(f"/ws/control/{session_id}")
control = context.__enter__()
control.send_json({
"type": "scenario.start", "scenario_id": POOL[0], "trainee": "Иванов",
"mode": "training", "exercise": "call",
})
wait_for(lambda: hub.get(session_id))
return session_id, control
def row_for(client, session_id):
rows = client.get("/api/sessions/active").json()
return next(item for item in rows if item["session_id"] == str(session_id))
def signal_kinds(row):
return {signal["kind"] for signal in row["signals"]}
def assigned_trainee(state, monkeypatch):
state.trainee_id = uuid4()
who = Principal(login="курсант", full_name="Курсант", role=Role.TRAINEE,
trainee_id=state.trainee_id)
monkeypatch.setattr(session_ws, "principal_of", lambda _ws: who)
def registry_time(monkeypatch, since, seconds):
class Clock(datetime):
@classmethod
def now(cls, tz=None):
return since + timedelta(seconds=seconds)
monkeypatch.setattr(sessions_http, "datetime", Clock)
def test_signal_backlog_when_queue_reaches_threshold(client, monkeypatch):
monkeypatch.setenv("SIGNAL_BACKLOG_THRESHOLD", "2")
get_settings.cache_clear()
session_id, control = start_two_card_dds(client)
try:
state = hub.get(session_id)
assert len(state.desk.cards) == 2, "оба билета должны прийти сразу без интервала"
row = row_for(client, session_id)
assert row["dds_open_cards"] == 2
assert "backlog" in signal_kinds(row)
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_declined_card_leaves_backlog(client, monkeypatch):
monkeypatch.setenv("SIGNAL_BACKLOG_THRESHOLD", "2")
get_settings.cache_clear()
session_id, control = start_two_card_dds(client)
try:
state = hub.get(session_id)
first_card = state.desk.ordered()[0]
with client.websocket_connect(f"/ws/station/{session_id}?role=dds") as station:
read_until(station, "station.state")
station.send_json({
"type": "card.status", "service": state.card_services(first_card)[0],
"status": "declined", "comment": "не наш адрес, передано в УК",
})
read_until(station, "station.state")
row = row_for(client, session_id)
assert row["dds_open_cards"] == 2, "отклонённая карточка ждёт «Следующей» на пульте"
assert "backlog" not in signal_kinds(row)
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_signal_refusals_after_two_consecutive_declines(client):
session_id, control = start_two_card_dds(client)
try:
state = hub.get(session_id)
first_card, second_card = state.desk.ordered()
with client.websocket_connect(f"/ws/station/{session_id}?role=dds") as station:
read_until(station, "station.state")
first_service = state.card_services(first_card)[0]
station.send_json({
"type": "card.status", "service": first_service, "status": "declined",
"comment": "не наш адрес, передано в УК",
})
read_until(station, "station.state")
row = row_for(client, session_id)
assert "refusals" not in signal_kinds(row), "одного отказа недостаточно для сигнала"
station.send_json({"type": "card.open", "card_id": str(second_card.card_id)})
read_until(station, "station.state")
second_service = state.card_services(second_card)[0]
station.send_json({
"type": "card.status", "service": second_service, "status": "declined",
"comment": "не наша территория, передано в ОМВД",
})
read_until(station, "station.state")
row = row_for(client, session_id)
assert state.consecutive_refusals == 2
assert "refusals" in signal_kinds(row)
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_acceptance_clears_refusal_streak(client):
session_id, control = start_two_card_dds(
client, [*POOL, "t20-2-stroke"],
)
try:
state = hub.get(session_id)
with client.websocket_connect(f"/ws/station/{session_id}") as station:
read_until(station, "station.state")
for card in state.desk.ordered()[:2]:
station.send_json({"type": "card.open", "card_id": str(card.card_id)})
read_until(station, "station.state")
station.send_json({
"type": "card.status", "service": state.card_services(card)[0],
"status": "declined", "comment": "Не наша территория, передано дежурному",
})
read_until(station, "station.state")
assert state.consecutive_refusals == 2
third = state.desk.ordered()[2]
station.send_json({"type": "card.open", "card_id": str(third.card_id)})
read_until(station, "station.state")
station.send_json({
"type": "card.status", "service": state.card_services(third)[0],
"status": "accepted", "comment": "Карточка принята диспетчером",
})
read_until(station, "station.state")
assert state.consecutive_refusals == 0
assert "refusals" not in signal_kinds(row_for(client, session_id))
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_signal_offline_when_station_socket_is_closed_past_the_window(client, monkeypatch):
session_id, control = start_two_card_dds(client)
try:
state = hub.get(session_id)
assigned_trainee(state, monkeypatch)
with client.websocket_connect(f"/ws/station/{session_id}?role=dds") as station:
read_until(station, "station.state")
assert hub.station_connected(session_id)
assert not hub.station_connected(session_id)
disconnected_at = state.socket_last_seen_at
assert disconnected_at is not None
window = get_settings().signal_offline_window_seconds
registry_time(monkeypatch, disconnected_at, window - 1)
assert "offline" not in signal_kinds(row_for(client, session_id))
registry_time(monkeypatch, disconnected_at, window + 1)
row = row_for(client, session_id)
assert "offline" in signal_kinds(row)
assert state.socket_last_seen_at == disconnected_at, "GET реестра не должен менять занятие"
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_signal_offline_without_first_connection(client, monkeypatch):
session_id, control = start_two_card_dds(client)
try:
state = hub.get(session_id)
started = state.socket_last_seen_at
assert started is not None
registry_time(monkeypatch, started, get_settings().signal_offline_window_seconds + 1)
assert "offline" in signal_kinds(row_for(client, session_id))
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_call_socket_disconnection_triggers_offline(client, monkeypatch):
session_id, control = start_call(client)
try:
state = hub.get(session_id)
assigned_trainee(state, monkeypatch)
with client.websocket_connect(f"/ws/call/{session_id}") as call:
read_until(call, "call.incoming")
assert hub.trainee_connected(session_id)
assert not hub.trainee_connected(session_id)
disconnected_at = state.socket_last_seen_at
registry_time(monkeypatch, disconnected_at,
get_settings().signal_offline_window_seconds + 1)
assert "offline" in signal_kinds(row_for(client, session_id))
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_teacher_socket_does_not_hide_offline(client, monkeypatch):
session_id, control = start_two_card_dds(client)
try:
state = hub.get(session_id)
with client.websocket_connect(f"/ws/station/{session_id}") as station:
read_until(station, "station.state")
assert not hub.station_connected(session_id)
registry_time(monkeypatch, state.socket_last_seen_at,
get_settings().signal_offline_window_seconds + 1)
assert "offline" in signal_kinds(row_for(client, session_id))
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_last_of_two_trainee_sockets_starts_offline_window(client, monkeypatch):
session_id, control = start_two_card_dds(client)
try:
state = hub.get(session_id)
assigned_trainee(state, monkeypatch)
with client.websocket_connect(f"/ws/station/{session_id}") as first:
read_until(first, "station.state")
with client.websocket_connect(f"/ws/station/{session_id}") as second:
read_until(second, "station.state")
assert hub.station_connected(session_id)
registry_time(monkeypatch, state.socket_last_seen_at,
get_settings().signal_offline_window_seconds + 1)
assert "offline" not in signal_kinds(row_for(client, session_id))
disconnected_at = state.socket_last_seen_at
registry_time(monkeypatch, disconnected_at,
get_settings().signal_offline_window_seconds + 1)
assert "offline" in signal_kinds(row_for(client, session_id))
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_offline_window_survives_checkpoint_restore(client, monkeypatch):
session_id, control = start_two_card_dds(client)
try:
state = hub.get(session_id)
assigned_trainee(state, monkeypatch)
with client.websocket_connect(f"/ws/station/{session_id}") as station:
read_until(station, "station.state")
disconnected_at = state.socket_last_seen_at
restored = load_state(dump_state(state), datetime.now(UTC))
hub.drop(session_id)
asyncio.run(hub._adopt(restored))
assert restored.socket_last_seen_at == disconnected_at
registry_time(monkeypatch, disconnected_at,
get_settings().signal_offline_window_seconds + 1)
assert "offline" in signal_kinds(row_for(client, session_id))
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_takeover_starts_window_when_socket_was_open(client, monkeypatch):
session_id, control = start_two_card_dds(client)
try:
state = hub.get(session_id)
assigned_trainee(state, monkeypatch)
with client.websocket_connect(f"/ws/station/{session_id}") as station:
read_until(station, "station.state")
assert state.socket_connected_at_checkpoint
snapshot = dump_state(state)
old_time = state.socket_last_seen_at
takeover_at = old_time + timedelta(minutes=10)
monkeypatch.setattr(hub_module, "now_utc", lambda: takeover_at)
restored = load_state(snapshot, datetime.now(UTC))
class CapturingStore(MemorySessionStore):
def __init__(self):
self.saved = None
async def commit(self, state, _records=()):
self.saved = dump_state(state)
store = CapturingStore()
local_hub = SessionHub(store)
try:
assert asyncio.run(local_hub._adopt(restored))
assert restored.socket_last_seen_at == takeover_at
assert not restored.socket_connected_at_checkpoint
assert datetime.fromisoformat(
store.saved["socket_last_seen_at"].replace("Z", "+00:00")
) == takeover_at
assert store.saved["socket_connected_at_checkpoint"] is False
finally:
local_hub.stop_ticker(session_id)
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_card_handoff_uses_trainee_station_presence(client, monkeypatch):
session_id, control = start_card_handoff(client)
try:
state = hub.get(session_id)
assigned_trainee(state, monkeypatch)
with monkeypatch.context() as clock:
with client.websocket_connect(f"/ws/call/{session_id}") as call:
read_until(call, "card.briefing")
call.send_json({"type": "card.submit"})
read_until(call, "call.ended")
handoff_at = state.socket_last_seen_at
clock.setattr(hub_module, "now_utc", lambda: handoff_at + timedelta(minutes=1))
assert state.socket_last_seen_at == handoff_at, "старый call-сокет не сбрасывает окно ДДС"
assert state.exercise.value == "card" and state.dds_phase
assert state.desk.scenarios
with client.websocket_connect(f"/ws/station/{session_id}") as station:
read_until(station, "station.state")
registry_time(monkeypatch, state.socket_last_seen_at,
get_settings().signal_offline_window_seconds + 1)
row = row_for(client, session_id)
assert row["dds_open_cards"] > 0
assert "offline" not in signal_kinds(row)
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_checkpoint_on_another_node_marks_presence_unknown(client, monkeypatch):
session_id, control = start_two_card_dds(client)
try:
state = hub.get(session_id)
snapshot = dump_state(state)
checkpoint_at = datetime.now(UTC)
hub.drop(session_id)
row = Session(id=session_id, owner_login="dev",
live_state=snapshot, checkpoint_at=checkpoint_at)
class CheckpointDB:
async def scalars(self, _query):
return IteratorResult(SimpleResultMetaData(["Session"]), iter([(row,)])).scalars()
monkeypatch.setattr(sessions_http, "require", lambda _request, _role:
Principal(login="dev", full_name="Преподаватель",
role=Role.INSTRUCTOR))
rows = asyncio.run(sessions_http.active(None, db=CheckpointDB()))
row = next(item for item in rows if item.session_id == session_id)
assert row.presence_known is False
assert "offline" not in {item.kind for item in row.signals}
finally:
hub.stop_ticker(session_id)
control.__exit__(None, None, None)
def test_failed_presence_commit_does_not_leave_connected_socket():
class BrokenStore(MemorySessionStore):
async def commit(self, _state, _records=()):
raise OSError("test commit failure")
local_hub = SessionHub(BrokenStore())
state = local_hub.register(SessionState(
session_id=uuid4(), scenario_id="test", scenario_title="Тест", level="L1",
mode=SessionMode.TRAINING, exercise=Exercise.DDS,
))
async def connect():
async with local_hub.trainee_socket(state.session_id, station=True, trainee=True):
assert local_hub.is_lease_fenced(state.session_id)
asyncio.run(connect())
assert not local_hub.station_connected(state.session_id)
def test_failed_presence_commit_closes_station_with_fencing_code(client, monkeypatch):
class BrokenStore(MemorySessionStore):
async def commit(self, _state, _records=()):
raise OSError("test commit failure")
session_id, control = start_two_card_dds(client)
try:
assigned_trainee(hub.get(session_id), monkeypatch)
hub.store = BrokenStore()
with client.websocket_connect(f"/ws/station/{session_id}") as station:
event = station.receive_json()
assert event["type"] == "error" and event["code"] == "internal"
with pytest.raises(WebSocketDisconnect) as closed:
station.receive_json()
assert closed.value.code == 1012
assert not hub.station_connected(session_id)
finally:
control.__exit__(None, None, None)
def test_takeover_retries_failed_presence_checkpoint():
class FlakyStore(MemorySessionStore):
def __init__(self):
self.commits = 0
async def commit(self, _state, _records=()):
self.commits += 1
if self.commits == 1:
raise OSError("temporary database outage")
store = FlakyStore()
local_hub = SessionHub(store)
state = SessionState(
session_id=uuid4(), scenario_id="test", scenario_title="Тест", level="L1",
mode=SessionMode.TRAINING, exercise=Exercise.DDS,
socket_last_seen_at=datetime.now(UTC) - timedelta(minutes=10),
socket_connected_at_checkpoint=True,
)
async def recover():
assert not await local_hub._adopt(state)
assert local_hub.get(state.session_id) is None
await local_hub.maintain_lease()
assert local_hub.get(state.session_id) is state
assert store.commits == 2
await local_hub.shutdown()
asyncio.run(recover())

View file

@ -22,6 +22,7 @@ import {
useActiveSessions, useIncidentGroups, useReviewScenarioSubmission, useScenarioSubmissions,
useScenarios, useSessions, useTrainees,
} from "@/shared/api/http";
import type { Signal } from "@/shared/api/http";
import { sessionIdFromUrl } from "@/shared/api/session";
import type { Level, SessionMode, SessionReport } from "@/shared/types/generated";
@ -44,6 +45,12 @@ function modeLabel(mode: SessionMode): string {
return MODES.find((item) => item.value === mode)?.label ?? mode;
}
function worstSignalSeverity(signals: Signal[]): "warn" | "violated" | null {
return signals.some((signal) => signal.severity === "violated")
? "violated"
: signals.length ? "warn" : null;
}
function sessionDateLabel(value: string | null): string | null {
if (!value) return null;
const date = new Date(value);
@ -338,7 +345,7 @@ export function Instructor() {
Список ваших активных сессий обновляется автоматически каждые 2 секунды.
</p>
<div className="dds-table-wrap"><table className="grid">
<thead><tr><th>Курсант</th><th>Упражнение · карточка</th><th>Ход</th><th>ДДС</th><th>Задержки</th><th></th></tr></thead>
<thead><tr><th>Курсант</th><th>Упражнение · карточка</th><th>Ход</th><th>ДДС</th><th>Задержки</th><th>Сигналы</th><th></th></tr></thead>
<tbody>
{(activeSessions.data ?? []).map((item) => (
<tr key={item.session_id}>
@ -383,12 +390,18 @@ export function Instructor() {
? `Просрочено: первичная реакция ${item.dds_overdue_cards}; отработка карточки ${item.dds_work_overdue_cards}`
: "Нормативы не нарушены"}
</td>
<td className={worstSignalSeverity(item.signals) ? `state-${worstSignalSeverity(item.signals)}` : item.presence_known ? "state-ok" : undefined}>
{item.signals.length
? item.signals.map((signal) => <span key={signal.kind}>{signal.text}<br /></span>)
: item.presence_known ? "Сигналов нет" : null}
{!item.presence_known && <span>Связь курсанта не проверена</span>}
</td>
<td><button type="button" onClick={() => openActiveSession(item)}>
{sessionId === item.session_id ? "Открыто" : "Наблюдать"}
</button></td>
</tr>
))}
{!activeSessions.data?.length && <tr><td colSpan={6}>
{!activeSessions.data?.length && <tr><td colSpan={7}>
{activeSessions.isLoading ? "Загружаем активные занятия…" : "Активных занятий пока нет."}
</td></tr>}
</tbody>

View file

@ -3,7 +3,9 @@
import { useMutation, useQuery, useQueryClient } from "@tanstack/react-query";
import type { Level, SessionMode, SessionReport, StationSnapshot } from "@/shared/types/generated";
import type {
Level, SessionMode, SessionReport, StationSnapshot, TimerState,
} from "@/shared/types/generated";
import { reportErrorMessage, reportRetryDelay, shouldRetryReport } from "./report.mjs";
@ -94,6 +96,12 @@ export interface SessionInfo {
end_reason: string | null;
}
export interface Signal {
kind: "backlog" | "refusals" | "offline";
severity: TimerState;
text: string;
}
export interface ActiveSessionInfo {
session_id: string;
trainee_name: string | null;
@ -109,6 +117,8 @@ export interface ActiveSessionInfo {
dds_work_overdue_cards: number;
dds_statuses: Record<string, string>;
dds_snapshot: StationSnapshot | null;
signals: Signal[];
presence_known: boolean;
}
export const useActiveSessions = () =>