diff --git a/backend/app/session/hub.py b/backend/app/session/hub.py index bfe8aea..2925dab 100644 --- a/backend/app/session/hub.py +++ b/backend/app/session/hub.py @@ -56,6 +56,7 @@ class SessionHub: 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`). @@ -119,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: @@ -158,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: # следующий оборот попробует снова @@ -183,9 +197,13 @@ class SessionHub: 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 diff --git a/backend/tests/test_session_signals.py b/backend/tests/test_session_signals.py index c867771..e75127e 100644 --- a/backend/tests/test_session_signals.py +++ b/backend/tests/test_session_signals.py @@ -461,3 +461,33 @@ def test_failed_presence_commit_closes_station_with_fencing_code(client, monkeyp 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())