Merge branch 'feat/session-signals'

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

View file

@ -8,6 +8,7 @@ import logging
import time import time
from collections.abc import AsyncIterator from collections.abc import AsyncIterator
from datetime import UTC, datetime from datetime import UTC, datetime
from typing import Literal
from uuid import UUID from uuid import UUID
from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response 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.db.models import Group, Score, Session, Trainee
from app.domain.events import Exercise, SessionMode, SessionReport from app.domain.events import Exercise, SessionMode, SessionReport
from app.domain.roles import Role from app.domain.roles import Role
from app.domain.statuses import SERVICE_STATUS_LABELS, StationSnapshot, current from app.domain.statuses import SERVICE_STATUS_LABELS, DdsQueueCard, StationSnapshot, current
from app.domain.timers import TimerCode from app.domain.timers import TimerCode, TimerState
from app.scoring.export import to_csv, to_pdf from app.scoring.export import to_csv, to_pdf
from app.scoring.report import build as build_report from app.scoring.report import build as build_report
from app.session.access import can_access 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.finish import override_score
from app.session.score import scoring_scenario from app.session.score import scoring_scenario
from app.session.hub import hub from app.session.hub import hub
from app.session.state import SessionState
from app.session.store import ScoreOverridden, apply_score_override from app.session.store import ScoreOverridden, apply_score_override
from app.voice.recording import recording_path from app.voice.recording import recording_path
@ -86,6 +88,15 @@ class DdsHistoryOut(BaseModel):
recipient_services: list[str] = [] recipient_services: list[str] = []
class Signal(BaseModel):
"""Одна строка колонки сигналов реестра. Цвет несёт смысл только через
`severity`; текст обязателен и не заменяется цветом."""
kind: Literal["backlog", "refusals", "offline"]
severity: TimerState
text: str
class ActiveSessionOut(BaseModel): class ActiveSessionOut(BaseModel):
session_id: UUID session_id: UUID
trainee_name: str | None trainee_name: str | None
@ -101,6 +112,8 @@ class ActiveSessionOut(BaseModel):
dds_work_overdue_cards: int dds_work_overdue_cards: int
dds_statuses: dict[str, str] dds_statuses: dict[str, str]
dds_snapshot: StationSnapshot | None = None dds_snapshot: StationSnapshot | None = None
signals: list[Signal] = []
presence_known: bool = True
def _out(session) -> SessionOut: def _out(session) -> SessionOut:
@ -184,6 +197,52 @@ async def dds_history(
return result 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]) @router.get("/active", response_model=list[ActiveSessionOut])
async def active( async def active(
request: Request, request: Request,
@ -197,6 +256,9 @@ async def active(
state.session_id: state state.session_id: state
for state in hub.active_sessions(who.login) for state in hub.active_sessions(who.login)
} }
#: Снимки без живой записи в хабе этого узла — чужой узел или сессия,
#: ожидающая восстановления. Её сокеты нельзя проверить этим процессом.
checkpoint_only: set[UUID] = set()
if db is not None: if db is not None:
rows = ( rows = (
await db.scalars( await db.scalars(
@ -225,11 +287,12 @@ async def active(
state.owner_login = row.owner_login state.owner_login = row.owner_login
if not state.ended: if not state.ended:
states[state.session_id] = state states[state.session_id] = state
checkpoint_only.add(state.session_id)
for state in states.values(): for state in states.values():
elapsed = (max(0, int((now - state.started_at).total_seconds())) elapsed = (max(0, int((now - state.started_at).total_seconds()))
if state.started_at else 0) 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 [] queue = station.queue_cards if station else []
card = state.desk.active card = state.desk.active
status_log = card.status_log if card is not None else [] status_log = card.status_log if card is not None else []
@ -262,6 +325,8 @@ async def active(
), ),
dds_statuses=latest_statuses, dds_statuses=latest_statuses,
dds_snapshot=station, 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 return result

View file

@ -14,7 +14,7 @@ from uuid import UUID, uuid4
from fastapi import APIRouter, WebSocket, WebSocketDisconnect from fastapi import APIRouter, WebSocket, WebSocketDisconnect
from pydantic import TypeAdapter, ValidationError 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 ( from app.domain.events import (
BgStart, BgStart,
CallEnded, CallEnded,
@ -390,9 +390,14 @@ async def call(ws: WebSocket, session_id: UUID) -> None:
) )
if entered is None: if entered is None:
return 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: if state.exercise is Exercise.CARD:
from app.api.ws.control import card_briefing 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, dds_service=identity.service or event.dds_service,
attempt=identity.attempt, attempt=identity.attempt,
criteria=event.criteria, 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.DDS_ACK] = event.criteria.decision_time_limit_seconds * 1000
state.timers.limits[TimerCode.CARD_FILL] = event.criteria.card_fill_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 fastapi import APIRouter, WebSocket, WebSocketDisconnect
from pydantic import TypeAdapter, ValidationError 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 ( from app.domain.events import (
CallEndReason, CallEndReason,
ErrorEvent, 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)) entered = await session_socket(ws, session_id, (Role.INSTRUCTOR, Role.TRAINEE))
if entered is None: if entered is None:
return 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)) sender = asyncio.create_task(pump(ws, queue))
try: try:
# Карточка, переданная до подключения станции, не теряется: # Карточка, переданная до подключения станции, не теряется:

View file

@ -120,6 +120,11 @@ class Settings(BaseSettings):
def limit_ms(self, code: TimerCode) -> int: def limit_ms(self, code: TimerCode) -> int:
return self.timer_limits_ms.get(code, NORMATIVES[code].limit_ms) 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: def validate_deployment_security(self) -> None:
"""Reject known development authentication defaults on a production app.""" """Reject known development authentication defaults on a production app."""
if self.app_env != "production": if self.app_env != "production":

View file

@ -446,6 +446,7 @@ class DdsDesk(BaseModel):
services[0], ServiceStatus.ACCEPTED, command.comment, services[0], ServiceStatus.ACCEPTED, command.comment,
author="диспетчер", author="диспетчер",
) )
session.consecutive_refusals = 0
except StatusError: except StatusError:
pass # статус уже стоит: повторное нажатие ничего не меняет pass # статус уже стоит: повторное нажатие ничего не меняет
case "card.status": case "card.status":
@ -462,6 +463,13 @@ class DdsDesk(BaseModel):
# Первичный статус останавливает норматив 30 секунд. # Первичный статус останавливает норматив 30 секунд.
if command.status in PRIMARY: if command.status in PRIMARY:
card.on_event("card.ack") card.on_event("card.ack")
# Сигнал реестра "повторные отказы" считает подряд идущие
# отказы, а не сумму за занятие.
session.consecutive_refusals = (
session.consecutive_refusals + 1
if command.status is ServiceStatus.DECLINED
else 0
)
if command.status in { if command.status in {
ServiceStatus.COMPLETED, ServiceStatus.DECLINED, ServiceStatus.REFUSED, 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.state import SessionState
from app.session.store import MemorySessionStore, Record, SessionLeaseLost, SessionStore 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._observers: dict[UUID, set[asyncio.Queue]] = {}
self._trainees: dict[UUID, set[asyncio.Queue]] = {} self._trainees: dict[UUID, set[asyncio.Queue]] = {}
self._stations: 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] = {} self._tickers: dict[UUID, asyncio.Task] = {}
# Задачи, порождённые внутри операции, наследуют контекст; после # Задачи, порождённые внутри операции, наследуют контекст; после
# закрытия операции их события идут напрямую (`_Operation.open`). # закрытия операции их события идут напрямую (`_Operation.open`).
@ -116,6 +120,7 @@ class SessionHub:
def drop(self, session_id: UUID) -> None: def drop(self, session_id: UUID) -> None:
self._sessions.pop(session_id, None) self._sessions.pop(session_id, None)
self._pending_adoptions.pop(session_id, None)
self.stop_ticker(session_id) self.stop_ticker(session_id)
def live_count(self) -> int: def live_count(self) -> int:
@ -135,9 +140,10 @@ class SessionHub:
if not self.store.persistent: if not self.store.persistent:
return 0 return 0
restored = await self.store.restore_active() restored = await self.store.restore_active()
adopted = 0
for state in restored: for state in restored:
self._adopt(state) adopted += await self._adopt(state)
return len(restored) return adopted
async def maintain_lease(self) -> None: async def maintain_lease(self) -> None:
"""Один оборот супервизора: продлить свои lease, подхватить просроченные чужие. """Один оборот супервизора: продлить свои lease, подхватить просроченные чужие.
@ -154,6 +160,18 @@ class SessionHub:
log.warning("продление lease занятия %s не удалось", state.session_id, log.warning("продление lease занятия %s не удалось", state.session_id,
exc_info=True) exc_info=True)
await self.fence(state) 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: try:
claimed = await self.store.claim_expired() claimed = await self.store.claim_expired()
except Exception: # следующий оборот попробует снова except Exception: # следующий оборот попробует снова
@ -163,17 +181,32 @@ class SessionHub:
current = self._sessions.get(state.session_id) current = self._sessions.get(state.session_id)
if current is not None and not current.lease_fenced: if current is not None and not current.lease_fenced:
continue continue
self._adopt(state) await self._adopt(state)
async def supervise_lease(self, interval: float) -> None: async def supervise_lease(self, interval: float) -> None:
while True: while True:
await asyncio.sleep(interval) await asyncio.sleep(interval)
await self.maintain_lease() await self.maintain_lease()
def _adopt(self, state: SessionState) -> None: async def _adopt(self, state: SessionState) -> bool:
self.stop_ticker(state.session_id) 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.register(state)
self.start_ticker(state.session_id) self.start_ticker(state.session_id)
return True
async def save_all(self) -> None: async def save_all(self) -> None:
"""Снимок живых занятий при остановке узла: следующий владелец продолжит с него.""" """Снимок живых занятий при остановке узла: следующий владелец продолжит с него."""
@ -285,6 +318,44 @@ class SessionHub:
def station(self, session_id: UUID): def station(self, session_id: UUID):
return self._subscribe(self._stations, session_id) 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 @staticmethod
@ -314,6 +385,12 @@ class SessionHub:
def observer_count(self, session_id: UUID) -> int: def observer_count(self, session_id: UUID) -> int:
return len(self._observers.get(session_id, set())) 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: 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. #: a lost WebSocket acknowledgement cannot apply an operation twice.
processed_station_commands: list[str] = Field(default_factory=list) processed_station_commands: list[str] = Field(default_factory=list)
text_revealed_facts: dict[str, str] = Field(default_factory=dict) 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") @field_validator("processed_station_commands")
@classmethod @classmethod
@ -206,6 +215,9 @@ class SessionState(PersistedSession):
self.dispatched_at = now_utc() self.dispatched_at = now_utc()
if self.exercise is Exercise.CALL: if self.exercise is Exercise.CALL:
self._receive_call_card() 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 return self.dispatched_card
def _receive_call_card(self) -> None: def _receive_call_card(self) -> None:

View file

@ -161,6 +161,9 @@ def full_state() -> SessionState:
resolve_comment="передано в другой регион", resolve_comment="передано в другой регион",
processed_station_commands=[str(uuid4())], processed_station_commands=[str(uuid4())],
text_revealed_facts={"f_address": "улица Ленина, 14"}, 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, useActiveSessions, useIncidentGroups, useReviewScenarioSubmission, useScenarioSubmissions,
useScenarios, useSessions, useTrainees, useScenarios, useSessions, useTrainees,
} from "@/shared/api/http"; } from "@/shared/api/http";
import type { Signal } from "@/shared/api/http";
import { sessionIdFromUrl } from "@/shared/api/session"; import { sessionIdFromUrl } from "@/shared/api/session";
import type { Level, SessionMode, SessionReport } from "@/shared/types/generated"; 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; 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 { function sessionDateLabel(value: string | null): string | null {
if (!value) return null; if (!value) return null;
const date = new Date(value); const date = new Date(value);
@ -338,7 +345,7 @@ export function Instructor() {
Список ваших активных сессий обновляется автоматически каждые 2 секунды. Список ваших активных сессий обновляется автоматически каждые 2 секунды.
</p> </p>
<div className="dds-table-wrap"><table className="grid"> <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> <tbody>
{(activeSessions.data ?? []).map((item) => ( {(activeSessions.data ?? []).map((item) => (
<tr key={item.session_id}> <tr key={item.session_id}>
@ -383,12 +390,18 @@ export function Instructor() {
? `Просрочено: первичная реакция ${item.dds_overdue_cards}; отработка карточки ${item.dds_work_overdue_cards}` ? `Просрочено: первичная реакция ${item.dds_overdue_cards}; отработка карточки ${item.dds_work_overdue_cards}`
: "Нормативы не нарушены"} : "Нормативы не нарушены"}
</td> </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)}> <td><button type="button" onClick={() => openActiveSession(item)}>
{sessionId === item.session_id ? "Открыто" : "Наблюдать"} {sessionId === item.session_id ? "Открыто" : "Наблюдать"}
</button></td> </button></td>
</tr> </tr>
))} ))}
{!activeSessions.data?.length && <tr><td colSpan={6}> {!activeSessions.data?.length && <tr><td colSpan={7}>
{activeSessions.isLoading ? "Загружаем активные занятия…" : "Активных занятий пока нет."} {activeSessions.isLoading ? "Загружаем активные занятия…" : "Активных занятий пока нет."}
</td></tr>} </td></tr>}
</tbody> </tbody>

View file

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