diff --git a/backend/app/api/http/admin.py b/backend/app/api/http/admin.py index 4f3a6ad..e7e773d 100644 --- a/backend/app/api/http/admin.py +++ b/backend/app/api/http/admin.py @@ -374,7 +374,7 @@ async def _runtime_metrics(db: AsyncSession) -> RuntimeMetrics: ) return RuntimeMetrics( **raw, - active_sessions=sum(not item.ended for item in hub._sessions.values()), # noqa: SLF001 + active_sessions=hub.live_count(), restored_sessions=getattr(app.state, "sessions_restored", 0), completed_sessions_24h=completed or 0, ) @@ -501,9 +501,9 @@ async def status(request: Request, db: AsyncSession = Depends(get_session)) -> l name="Живых занятий", ok=True, detail=( - f"активно {sum(not item.ended for item in hub._sessions.values())}; " + f"активно {hub.live_count()}; " f"восстановлено после запуска {getattr(app.state, 'sessions_restored', 0)}" - ), # noqa: SLF001 — реестр в памяти процесса + ), ) ) try: diff --git a/backend/app/main.py b/backend/app/main.py index d63e4d6..1a89cd9 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -29,10 +29,12 @@ from app.scenarios import store from app.scenarios.loader import ScenarioError from app.session.hub import hub from app.session.pg_store import PostgresSessionStore -from app.session.store import MemorySessionStore, SessionLeaseLost +from app.session.store import MemorySessionStore from app.voice.models import get_voice_models LIBRARY = Path(__file__).resolve().parents[2] / "scenarios" +#: Треть срока lease (`pg_store.LEASE_SECONDS`): два пропущенных оборота ещё не теряют занятие. +LEASE_SUPERVISOR_SECONDS = 5 # uvicorn настраивает только собственные логгеры: без этого INFO из модулей # приложения («получено N кадров», замеры задержки голоса) молча теряется, @@ -97,66 +99,27 @@ async def lifespan(app: FastAPI): else PostgresSessionStore(get_sessionmaker(), node_id=settings.backend_node_id) ) app.state.sessions_restored = 0 - if hub.store.persistent: - try: - async with asyncio.timeout(3): - restored = await hub.store.restore_active() - for state in restored: - hub.register(state) - hub.start_ticker(state.session_id) - app.state.sessions_restored = len(restored) - if restored: - logging.getLogger(__name__).info( - "восстановлено активных занятий: %d", len(restored) - ) - except (SQLAlchemyError, OSError, TimeoutError) as exc: - logging.getLogger(__name__).warning( - "не удалось восстановить активные занятия: %s", exc + try: + async with asyncio.timeout(3): + app.state.sessions_restored = await hub.restore() + if app.state.sessions_restored: + logging.getLogger(__name__).info( + "восстановлено активных занятий: %d", app.state.sessions_restored ) - - lease_task = None - if hub.store.persistent and settings.backend_node_id: - async def supervise_session_ownership() -> None: - while True: - await asyncio.sleep(5) - store = hub.store - for state in list(hub._sessions.values()): - if state.ended or state.lease_fenced: - continue - try: - await store.renew(state.session_id) - except SessionLeaseLost: - await hub.fence(state) - except (SQLAlchemyError, OSError, TimeoutError): - logging.getLogger(__name__).warning( - "backend lease renewal failed for %s", state.session_id, - exc_info=True, - ) - await hub.fence(state) - except Exception: # noqa: BLE001 — unknown ownership state fails closed - logging.getLogger(__name__).exception( - "unexpected backend lease failure for %s", state.session_id - ) - await hub.fence(state) - try: - restored = await store.claim_expired() - for state in restored: - current = hub._sessions.get(state.session_id) - if current is not None and not current.lease_fenced: - continue - if current is not None: - hub.stop_ticker(state.session_id) - hub.register(state) - hub.start_ticker(state.session_id) - except (SQLAlchemyError, OSError, TimeoutError): - logging.getLogger(__name__).exception( - "не удалось проверить/восстановить занятия с истёкшей backend lease" - ) - - lease_task = asyncio.create_task( - supervise_session_ownership(), name="session-owner-lease-supervisor" + except (SQLAlchemyError, OSError, TimeoutError) as exc: + logging.getLogger(__name__).warning( + "не удалось восстановить активные занятия: %s", exc ) + # Lease держит только узел кластера: без node_id владелец один. + lease_task = ( + asyncio.create_task( + hub.supervise_lease(LEASE_SUPERVISOR_SECONDS), + name="session-owner-lease-supervisor", + ) + if hub.store.persistent and settings.backend_node_id else None + ) + # Эмбеддинги для слот-автомата — грузятся один раз, до первого занятия. app.state.embeddings_ready = not settings.demo_no_db and get_embedder() is not None @@ -178,19 +141,8 @@ async def lifespan(app: FastAPI): with suppress(asyncio.CancelledError): await lease_task - if hub.store.persistent: - for state in list(hub._sessions.values()): - if not state.ended and not state.lease_fenced: - try: - await hub.store.commit(state) - except Exception: # noqa: BLE001 — shutdown must release the process - logging.getLogger(__name__).exception( - "не удалось сохранить checkpoint %s при shutdown", state.session_id - ) + await hub.save_all() await hub.shutdown() - for state in list(hub._sessions.values()): - if state.voice is not None: - await state.voice.close() app = FastAPI(title="Учебный симулятор занятия для системы 112", lifespan=lifespan) diff --git a/backend/app/session/finish.py b/backend/app/session/finish.py index 952c630..8e7e8d9 100644 --- a/backend/app/session/finish.py +++ b/backend/app/session/finish.py @@ -348,7 +348,7 @@ async def finish(session_id: UUID, state) -> None: # запущенной модели разбор выходит со статусом «недоступно», а не ждёт её. state.score["ai_coaching"] = (await coach(result.metrics)).model_dump(mode="json") # Полный разбор хранится вместе с оценкой: PDF/CSV и история должны - # переживать перезапуск backend, а не зависеть от объекта в hub._sessions. + # переживать перезапуск backend, а не зависеть от живого объекта в реестре хаба. state.score["full_report"] = build_report(session_id, state, scenario).model_dump(mode="json") log.info("сессия %s: оценка %.1f, отметок %d", session_id, result.score, len(result.findings)) diff --git a/backend/app/session/hub.py b/backend/app/session/hub.py index 4b29112..056c616 100644 --- a/backend/app/session/hub.py +++ b/backend/app/session/hub.py @@ -23,7 +23,7 @@ from app.domain.events import ( TimerTick, ) from app.session.state import SessionState -from app.session.store import MemorySessionStore, Record, SessionStore +from app.session.store import MemorySessionStore, Record, SessionLeaseLost, SessionStore #: Очередь одного подписчика. Медленный наблюдатель не тормозит занятие: #: очередь ограничена, переполнение роняет соединение, а не сессию. @@ -118,6 +118,73 @@ class SessionHub: self._sessions.pop(session_id, None) self.stop_ticker(session_id) + def live_count(self) -> int: + """Сколько занятий этот узел сейчас ведёт — для диагностики администратора.""" + return sum(1 for state in self._live()) + + def _live(self) -> list[SessionState]: + return [ + state for state in self._sessions.values() + if not state.ended and not state.lease_fenced + ] + + # ── lease узла поверх хранилища ── + + async def restore(self) -> int: + """Поднять незавершённые занятия узла после перезапуска процесса.""" + if not self.store.persistent: + return 0 + restored = await self.store.restore_active() + for state in restored: + self._adopt(state) + return len(restored) + + async def maintain_lease(self) -> None: + """Один оборот супервизора: продлить свои lease, подхватить просроченные чужие. + + Непродлённая lease — не повод писать дальше «на авось»: любая ошибка + продления закрывает занятие на узле (docs/arch/SCALE-OUT.md). + """ + for state in self._live(): + try: + await self.store.renew(state.session_id) + except SessionLeaseLost: + await self.fence(state) + except Exception: # неизвестное владение закрывается + log.warning("продление lease занятия %s не удалось", state.session_id, + exc_info=True) + await self.fence(state) + try: + claimed = await self.store.claim_expired() + except Exception: # следующий оборот попробует снова + log.exception("не удалось подхватить занятия с истёкшей lease") + return + for state in claimed: + current = self._sessions.get(state.session_id) + if current is not None and not current.lease_fenced: + continue + 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: + self.stop_ticker(state.session_id) + self.register(state) + self.start_ticker(state.session_id) + + async def save_all(self) -> None: + """Снимок живых занятий при остановке узла: следующий владелец продолжит с него.""" + if not self.store.persistent: + return + for state in self._live(): + try: + await self.store.commit(state) + except Exception: # остановка должна освободить процесс + log.exception("не удалось сохранить снимок %s при остановке", state.session_id) + # ── операция: один commit, потом события ── @contextlib.asynccontextmanager @@ -269,6 +336,9 @@ class SessionHub: """ for session_id in list(self._tickers): self.stop_ticker(session_id) + for state in list(self._sessions.values()): + if state.voice is not None: + await state.voice.close() async def _tick(self, session_id: UUID) -> None: try: diff --git a/backend/tests/test_session_lease.py b/backend/tests/test_session_lease.py new file mode 100644 index 0000000..ed8cf5b --- /dev/null +++ b/backend/tests/test_session_lease.py @@ -0,0 +1,152 @@ +"""Lease занятий ведёт хаб поверх `SessionStore`, а не код запуска приложения.""" + +import asyncio +import re +from pathlib import Path + +from app.session.hub import LEASE_FENCED_MESSAGE, SessionHub +from app.session.state import now_utc +from app.session.store import MemorySessionStore, SessionLeaseLost +from tests.test_session_checkpoint import dds_state + +APP = Path(__file__).resolve().parents[1] / "app" + + +class LeaseStore(MemorySessionStore): + persistent = True + + def __init__(self) -> None: + super().__init__() + self.renewed = [] + self.renew_errors: dict = {} + self.expired = [] + self.active = [] + + async def renew(self, session_id): + self.renewed.append(session_id) + error = self.renew_errors.get(session_id) + if error is not None: + raise error + + async def claim_expired(self, session_id=None): + claimed, self.expired = self.expired, [] + return claimed + + async def restore_active(self): + return list(self.active) + + +_hubs: list[SessionHub] = [] + + +def make_hub(store) -> SessionHub: + item = SessionHub(store=store) + _hubs.append(item) + return item + + +def run(coro): + """Такты, запущенные подхватом, гасятся в том же цикле событий.""" + async def wrapper(): + try: + return await coro + finally: + while _hubs: + await _hubs.pop().shutdown() + return asyncio.run(wrapper()) + + +def test_lost_or_unconfirmed_lease_fences_only_that_session(): + store = LeaseStore() + local_hub = make_hub(store) + lost = local_hub.register(dds_state()) + broken = local_hub.register(dds_state()) + healthy = local_hub.register(dds_state()) + ended = local_hub.register(dds_state()) + ended.ended_at = now_utc() + store.renew_errors = { + lost.session_id: SessionLeaseLost("fenced"), + broken.session_id: OSError("partition"), + } + + with local_hub.trainee(lost.session_id) as trainee: + run(local_hub.maintain_lease()) + assert trainee.get_nowait().message == LEASE_FENCED_MESSAGE + + assert lost.lease_fenced and broken.lease_fenced + assert not healthy.lease_fenced + assert ended.session_id not in store.renewed, "завершённое занятие lease не держит" + + +def test_expired_sessions_are_adopted_but_live_owner_is_not_replaced(): + store = LeaseStore() + local_hub = make_hub(store) + live = local_hub.register(dds_state()) + stale = local_hub.register(dds_state()) + stale.lease_fenced = True + + live_copy = live.model_copy() + stale_copy = stale.model_copy() + stale_copy.lease_fenced = False + fresh = dds_state() + store.expired = [live_copy, stale_copy, fresh] + + async def scenario(): + await local_hub.maintain_lease() + return set(local_hub._tickers) + + tickers = run(scenario()) + assert local_hub.get(live.session_id) is live + assert local_hub.get(stale.session_id) is stale_copy + assert local_hub.get(fresh.session_id) is fresh + assert {stale.session_id, fresh.session_id} <= tickers + + +def test_restore_registers_active_sessions_with_tickers(): + store = LeaseStore() + store.active = [dds_state(), dds_state()] + local_hub = make_hub(store) + + async def scenario(): + restored = await local_hub.restore() + return restored, set(local_hub._tickers) + + restored, tickers = run(scenario()) + assert restored == 2 + assert {state.session_id for state in store.active} == tickers + assert local_hub.live_count() == 2 + + +def test_volatile_store_restores_nothing(): + local_hub = make_hub(MemorySessionStore()) + assert run(local_hub.restore()) == 0 + + +def test_save_all_commits_only_live_sessions(): + store = LeaseStore() + local_hub = make_hub(store) + live = local_hub.register(dds_state()) + local_hub.register(dds_state()).ended_at = now_utc() + local_hub.register(dds_state()).lease_fenced = True + + run(local_hub.save_all()) + assert [sid for sid, _records in store.commits] == [live.session_id] + + +def test_live_count_skips_ended_and_fenced(): + local_hub = make_hub(MemorySessionStore()) + local_hub.register(dds_state()) + local_hub.register(dds_state()).ended_at = now_utc() + local_hub.register(dds_state()).lease_fenced = True + assert local_hub.live_count() == 1 + + +def test_application_code_does_not_touch_hub_internals(): + offenders = [ + f"{path.relative_to(APP)}:{number}" + for path in APP.rglob("*.py") + if path.name != "hub.py" + for number, line in enumerate(path.read_text(encoding="utf-8").splitlines(), 1) + if re.search(r"\bhub\._", line) + ] + assert offenders == []