fix: учёт связи курсанта по сокетам и фазе ДДС
This commit is contained in:
parent
2757a34ed5
commit
57c3ff5a5c
10 changed files with 420 additions and 42 deletions
|
|
@ -113,6 +113,7 @@ class ActiveSessionOut(BaseModel):
|
|||
dds_statuses: dict[str, str]
|
||||
dds_snapshot: StationSnapshot | None = None
|
||||
signals: list[Signal] = []
|
||||
presence_known: bool = True
|
||||
|
||||
|
||||
def _out(session) -> SessionOut:
|
||||
|
|
@ -206,11 +207,8 @@ def _signals(
|
|||
подписчика в хабе этого узла; окно нужно, чтобы короткий разрыв
|
||||
соединения не сразу считался потерей связи курсанта.
|
||||
|
||||
`live=False` — снимок собран из checkpoint (сессию не держит хаб этого
|
||||
процесса: рестарт до переподключения сокета или чужой узел в кластере).
|
||||
Такой объект живёт один запрос и выбрасывается: писать в него
|
||||
`socket_last_seen_at` бессмысленно — на следующем опросе присутствие
|
||||
неизвестно снова. Честнее не утверждать «на связи», чем соврать.
|
||||
`live=False` — снимок чужого узла: его сокеты здесь не видны, поэтому
|
||||
присутствие неизвестно. Чтение реестра не меняет состояние занятия.
|
||||
"""
|
||||
settings = get_settings()
|
||||
signals: list[Signal] = []
|
||||
|
|
@ -226,16 +224,14 @@ def _signals(
|
|||
))
|
||||
if live:
|
||||
connected = (
|
||||
hub.station_connected(state.session_id) if state.exercise is Exercise.DDS
|
||||
hub.station_connected(state.session_id) if state.dds_phase
|
||||
else hub.trainee_connected(state.session_id)
|
||||
)
|
||||
if connected:
|
||||
state.socket_last_seen_at = now
|
||||
offline_seconds = (
|
||||
(now - state.socket_last_seen_at).total_seconds()
|
||||
if state.socket_last_seen_at is not None else 0
|
||||
(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 offline_seconds >= settings.signal_offline_window_seconds:
|
||||
if not connected and offline_seconds >= settings.signal_offline_window_seconds:
|
||||
signals.append(Signal(
|
||||
kind="offline", severity="violated",
|
||||
text=f"Курсант не на связи {int(offline_seconds)} с",
|
||||
|
|
@ -256,9 +252,8 @@ async def active(
|
|||
state.session_id: state
|
||||
for state in hub.active_sessions(who.login)
|
||||
}
|
||||
#: Снимки без живой записи в хабе этого узла — рестарт до переподключения
|
||||
#: сокета либо чужой узел в кластере. Присутствие сокета для них здесь
|
||||
#: не проверяется: объект живёт один запрос, а не хаб (см. `_signals`).
|
||||
#: Снимки без живой записи в хабе этого узла — чужой узел или сессия,
|
||||
#: ожидающая восстановления. Её сокеты нельзя проверить этим процессом.
|
||||
checkpoint_only: set[UUID] = set()
|
||||
if db is not None:
|
||||
rows = (
|
||||
|
|
@ -293,7 +288,7 @@ async def active(
|
|||
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 []
|
||||
|
|
@ -327,6 +322,7 @@ 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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
# Карточка, переданная до подключения станции, не теряется:
|
||||
|
|
|
|||
|
|
@ -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,8 @@ 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._tickers: dict[UUID, asyncio.Task] = {}
|
||||
# Задачи, порождённые внутри операции, наследуют контекст; после
|
||||
# закрытия операции их события идут напрямую (`_Operation.open`).
|
||||
|
|
@ -135,9 +138,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, подхватить просроченные чужие.
|
||||
|
|
@ -163,17 +167,28 @@ 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 Exception:
|
||||
log.exception("не удалось сохранить присутствие после takeover %s", state.session_id)
|
||||
return False
|
||||
self.register(state)
|
||||
self.start_ticker(state.session_id)
|
||||
return True
|
||||
|
||||
async def save_all(self) -> None:
|
||||
"""Снимок живых занятий при остановке узла: следующий владелец продолжит с него."""
|
||||
|
|
@ -285,6 +300,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
|
||||
|
|
@ -315,10 +368,10 @@ class SessionHub:
|
|||
return len(self._observers.get(session_id, set()))
|
||||
|
||||
def station_connected(self, session_id: UUID) -> bool:
|
||||
return bool(self._stations.get(session_id))
|
||||
return bool(self._trainee_stations.get(session_id))
|
||||
|
||||
def trainee_connected(self, session_id: UUID) -> bool:
|
||||
return bool(self._trainees.get(session_id))
|
||||
return bool(self._trainee_calls.get(session_id))
|
||||
|
||||
# ── такт таймеров ──
|
||||
|
||||
|
|
|
|||
|
|
@ -127,11 +127,12 @@ class PersistedSession(BaseModel):
|
|||
#: Подряд идущих первичных отказов без принятой карточки между ними —
|
||||
#: сигнал реестра преподавателя, не влияет на балл.
|
||||
consecutive_refusals: int = 0
|
||||
#: Когда реестр преподавателя последний раз видел подключённый сокет
|
||||
#: курсанта (станция ДДС или звонок 112) — для сигнала "не на связи".
|
||||
#: Персистится: без этого рестарт/failover обнуляет отсчёт окна
|
||||
#: присутствия и сигнал никогда не срабатывает для восстановленных сессий.
|
||||
#: Последнее подключение/отключение сокета курсанта или начало фазы ДДС.
|
||||
#: Сохраняется вместе с занятием, чтобы чтение реестра не меняло состояние.
|
||||
socket_last_seen_at: datetime | None = None
|
||||
#: На момент последнего checkpoint курсант был подключён на текущем АРМ.
|
||||
#: После аварии узла момент потери связи неизвестен: окно начинается при takeover.
|
||||
socket_connected_at_checkpoint: bool = False
|
||||
|
||||
@field_validator("processed_station_commands")
|
||||
@classmethod
|
||||
|
|
@ -214,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:
|
||||
|
|
|
|||
Loading…
Reference in a new issue