refactor: lease занятий ведёт хаб поверх SessionStore, main и admin не читают реестр хаба напрямую
This commit is contained in:
parent
47cb85ee02
commit
0d7bbf39b3
5 changed files with 248 additions and 74 deletions
|
|
@ -374,7 +374,7 @@ async def _runtime_metrics(db: AsyncSession) -> RuntimeMetrics:
|
||||||
)
|
)
|
||||||
return RuntimeMetrics(
|
return RuntimeMetrics(
|
||||||
**raw,
|
**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),
|
restored_sessions=getattr(app.state, "sessions_restored", 0),
|
||||||
completed_sessions_24h=completed or 0,
|
completed_sessions_24h=completed or 0,
|
||||||
)
|
)
|
||||||
|
|
@ -501,9 +501,9 @@ async def status(request: Request, db: AsyncSession = Depends(get_session)) -> l
|
||||||
name="Живых занятий",
|
name="Живых занятий",
|
||||||
ok=True,
|
ok=True,
|
||||||
detail=(
|
detail=(
|
||||||
f"активно {sum(not item.ended for item in hub._sessions.values())}; "
|
f"активно {hub.live_count()}; "
|
||||||
f"восстановлено после запуска {getattr(app.state, 'sessions_restored', 0)}"
|
f"восстановлено после запуска {getattr(app.state, 'sessions_restored', 0)}"
|
||||||
), # noqa: SLF001 — реестр в памяти процесса
|
),
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
|
|
|
||||||
|
|
@ -29,10 +29,12 @@ from app.scenarios import store
|
||||||
from app.scenarios.loader import ScenarioError
|
from app.scenarios.loader import ScenarioError
|
||||||
from app.session.hub import hub
|
from app.session.hub import hub
|
||||||
from app.session.pg_store import PostgresSessionStore
|
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
|
from app.voice.models import get_voice_models
|
||||||
|
|
||||||
LIBRARY = Path(__file__).resolve().parents[2] / "scenarios"
|
LIBRARY = Path(__file__).resolve().parents[2] / "scenarios"
|
||||||
|
#: Треть срока lease (`pg_store.LEASE_SECONDS`): два пропущенных оборота ещё не теряют занятие.
|
||||||
|
LEASE_SUPERVISOR_SECONDS = 5
|
||||||
|
|
||||||
# uvicorn настраивает только собственные логгеры: без этого INFO из модулей
|
# uvicorn настраивает только собственные логгеры: без этого INFO из модулей
|
||||||
# приложения («получено N кадров», замеры задержки голоса) молча теряется,
|
# приложения («получено N кадров», замеры задержки голоса) молча теряется,
|
||||||
|
|
@ -97,66 +99,27 @@ async def lifespan(app: FastAPI):
|
||||||
else PostgresSessionStore(get_sessionmaker(), node_id=settings.backend_node_id)
|
else PostgresSessionStore(get_sessionmaker(), node_id=settings.backend_node_id)
|
||||||
)
|
)
|
||||||
app.state.sessions_restored = 0
|
app.state.sessions_restored = 0
|
||||||
if hub.store.persistent:
|
try:
|
||||||
try:
|
async with asyncio.timeout(3):
|
||||||
async with asyncio.timeout(3):
|
app.state.sessions_restored = await hub.restore()
|
||||||
restored = await hub.store.restore_active()
|
if app.state.sessions_restored:
|
||||||
for state in restored:
|
logging.getLogger(__name__).info(
|
||||||
hub.register(state)
|
"восстановлено активных занятий: %d", app.state.sessions_restored
|
||||||
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
|
|
||||||
)
|
)
|
||||||
|
except (SQLAlchemyError, OSError, TimeoutError) as exc:
|
||||||
lease_task = None
|
logging.getLogger(__name__).warning(
|
||||||
if hub.store.persistent and settings.backend_node_id:
|
"не удалось восстановить активные занятия: %s", exc
|
||||||
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"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# 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
|
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):
|
with suppress(asyncio.CancelledError):
|
||||||
await lease_task
|
await lease_task
|
||||||
|
|
||||||
if hub.store.persistent:
|
await hub.save_all()
|
||||||
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.shutdown()
|
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)
|
app = FastAPI(title="Учебный симулятор занятия для системы 112", lifespan=lifespan)
|
||||||
|
|
|
||||||
|
|
@ -348,7 +348,7 @@ async def finish(session_id: UUID, state) -> None:
|
||||||
# запущенной модели разбор выходит со статусом «недоступно», а не ждёт её.
|
# запущенной модели разбор выходит со статусом «недоступно», а не ждёт её.
|
||||||
state.score["ai_coaching"] = (await coach(result.metrics)).model_dump(mode="json")
|
state.score["ai_coaching"] = (await coach(result.metrics)).model_dump(mode="json")
|
||||||
# Полный разбор хранится вместе с оценкой: PDF/CSV и история должны
|
# Полный разбор хранится вместе с оценкой: PDF/CSV и история должны
|
||||||
# переживать перезапуск backend, а не зависеть от объекта в hub._sessions.
|
# переживать перезапуск backend, а не зависеть от живого объекта в реестре хаба.
|
||||||
state.score["full_report"] = build_report(session_id, state, scenario).model_dump(mode="json")
|
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))
|
log.info("сессия %s: оценка %.1f, отметок %d", session_id, result.score, len(result.findings))
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -23,7 +23,7 @@ from app.domain.events import (
|
||||||
TimerTick,
|
TimerTick,
|
||||||
)
|
)
|
||||||
from app.session.state import SessionState
|
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._sessions.pop(session_id, None)
|
||||||
self.stop_ticker(session_id)
|
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, потом события ──
|
# ── операция: один commit, потом события ──
|
||||||
|
|
||||||
@contextlib.asynccontextmanager
|
@contextlib.asynccontextmanager
|
||||||
|
|
@ -269,6 +336,9 @@ class SessionHub:
|
||||||
"""
|
"""
|
||||||
for session_id in list(self._tickers):
|
for session_id in list(self._tickers):
|
||||||
self.stop_ticker(session_id)
|
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:
|
async def _tick(self, session_id: UUID) -> None:
|
||||||
try:
|
try:
|
||||||
|
|
|
||||||
152
backend/tests/test_session_lease.py
Normal file
152
backend/tests/test_session_lease.py
Normal file
|
|
@ -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 == []
|
||||||
Loading…
Reference in a new issue