chore: обновить паузу занятия от main

# Conflicts:
#	backend/app/scoring/dispatcher.py
#	backend/app/session/store.py
This commit is contained in:
kaifarikman 2026-09-27 23:30:18 +03:00
commit 0d58878cbe
51 changed files with 4545 additions and 87 deletions

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: