/
systemsstrategyy
/
Treker
Обзор
Документация
Войти
/
systemsstrategyy
/
Treker
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
master
api/shared/db/session.py
174 строки
8 KB
HaGaSRus
Волна 4 (аудит «как работает»): 24 фикса корректности/надёжности/эксплуатации
07 июл 2026, 14:15
07 июл 2026, 14:15
0dcebaa
Код
Авторство
О чём код?
"""Async SQLAlchemy session factory и FastAPI-зависимость. Coding_Principles §3 (ADR-0001): Unit of Work — один HTTP-запрос = одна сессия = одна транзакция. `expire_on_commit=False` обязательно для async (§13) — иначе после commit() lazy-load бросит MissingGreenlet. Usage: from api.shared.db.session import get_db @router.get("/items") async def list_items(db: AsyncSession = Depends(get_db)): ... """ from __future__ import annotations import asyncio from collections.abc import AsyncGenerator from typing import Any import structlog from sqlalchemy import event from sqlalchemy.ext.asyncio import ( AsyncEngine, AsyncSession, async_sessionmaker, create_async_engine, ) from sqlalchemy.orm import Session as _SyncSession from api.core.settings import settings log = structlog.get_logger(__name__) def _create_engine() -> AsyncEngine: """Создаёт async engine из настроек. SQL-эхо включено только в dev.""" return create_async_engine( settings.database.url, echo=settings.is_dev and settings.debug, pool_size=settings.database.pool_size, max_overflow=settings.database.max_overflow, pool_pre_ping=True, # detect dead connections ) engine: AsyncEngine = _create_engine() SessionLocal: async_sessionmaker[AsyncSession] = async_sessionmaker( engine, class_=AsyncSession, expire_on_commit=False, # см. ADR-0001 §13 autoflush=False, ) async def get_db() -> AsyncGenerator[AsyncSession]: """FastAPI dependency — одна сессия на запрос. Coding_Principles §3 Unit of Work: коммит на успехе, rollback на исключении. Эквивалентно `async with session.begin()`, но обеспечивает явный close().""" async with SessionLocal() as session: # Request-сессии работают в режиме «defer после commit»: заявки на # Procrastinate-задачи уходят в очередь только если бизнес-транзакция # успешно закоммичена (аудит w4, ранг 2 — фантомные письма при rollback). enable_after_commit_defers(session) try: yield session await session.commit() except Exception: await session.rollback() raise finally: await session.close() # ── After-commit hook для Procrastinate-defer'ов ──────────────────────────── # # Проблема (аудит w4, ранги 2/16): `task.defer_async()` пишет job в очередь # через ОТДЕЛЬНЫЙ коннект Procrastinate и коммитится немедленно — независимо # от исхода бизнес-транзакции вызывающего. При rollback транзакции письмо / # telegram всё равно уходят (фантом), а при сбое commit батча check_due_dates # метка due_notified_at теряется и следующий tick дублирует рассылку. # # Минимальный after-commit паттерн вместо полного transactional outbox: # заявки на defer копятся в `session.info` и отправляются в Procrastinate # из sync-события SQLAlchemy `after_commit`. Событие синхронное — await # внутри него невозможен, поэтому отправка идёт через `loop.create_task` # (event loop гарантированно запущен: commit async-сессии исполняется в # greenlet на потоке loop'а). При rollback заявки отбрасываются. # # Tradeoff: crash процесса в окне «commit прошёл, create_task ещё не успел # отправить job» теряет уведомление. Осознанный выбор: потерять письмо # безопаснее, чем отправить фантомное или дублирующее. Полный outbox # (persistent-таблица заявок) — отдельная задача, если потеря станет заметной. _DEFER_MODE_KEY = "defer_tasks_after_commit" _DEFER_QUEUE_KEY = "after_commit_defer_queue" # Strong refs на фоновые dispatch-задачи: asyncio держит только weak ref, # без этого set задача может быть собрана GC до завершения. _dispatch_tasks: set[asyncio.Task[None]] = set() def enable_after_commit_defers(session: AsyncSession) -> None: """Включает для сессии режим «defer после commit». Вне режима `defer_task` дефёрит немедленно (легаси-поведение, на котором построены существующие юнит-тесты notify).""" # getattr-guard: тест-дублёры сессий (FakeSession в tests/core/test_db.py) # не имеют .info — для них режим просто не включается. info = getattr(session, "info", None) if info is None: return info[_DEFER_MODE_KEY] = True async def defer_task(session: AsyncSession, task: Any, **kwargs: Any) -> None: """Дефёрит procrastinate-task с учётом транзакции сессии. В режиме after-commit заявка копится в `session.info` и уходит в очередь только после успешного commit; при rollback — отбрасывается. Вне режима — немедленный `defer_async` (как раньше).""" # getattr-guard как в enable_after_commit_defers: тест-дублёры сессий # без .info работают в легаси-режиме (немедленный defer). info = getattr(session, "info", None) if info is not None and info.get(_DEFER_MODE_KEY): info.setdefault(_DEFER_QUEUE_KEY, []).append((task, kwargs)) return await task.defer_async(**kwargs) async def wait_after_commit_defers() -> None: """Дожидается завершения фоновых dispatch-задач (для тестов и shutdown).""" if _dispatch_tasks: await asyncio.gather(*list(_dispatch_tasks), return_exceptions=True) async def _dispatch_one(task: Any, kwargs: dict[str, Any]) -> None: """Отправляет одну заявку; сбой логируем громко, но не роняем loop.""" try: await task.defer_async(**kwargs) except Exception: log.exception( "after_commit.defer_failed", task_name=getattr(task, "name", repr(task)), ) @event.listens_for(_SyncSession, "after_commit") def _dispatch_after_commit(session: _SyncSession) -> None: """Sync-событие SQLAlchemy: после успешного commit отправляем заявки.""" pending = session.info.pop(_DEFER_QUEUE_KEY, None) if not pending: return try: loop = asyncio.get_running_loop() except RuntimeError: # Commit вне event loop'а (sync-скрипты) — доставить нечем; # логируем громко, чтобы потеря не была молчаливой. log.error("after_commit.defer_no_event_loop", lost=len(pending)) return for task, kwargs in pending: t = loop.create_task(_dispatch_one(task, kwargs)) _dispatch_tasks.add(t) t.add_done_callback(_dispatch_tasks.discard) @event.listens_for(_SyncSession, "after_rollback") def _discard_defers_on_rollback(session: _SyncSession) -> None: """Rollback бизнес-транзакции отменяет и накопленные defer-заявки.""" pending = session.info.pop(_DEFER_QUEUE_KEY, None) if pending: log.info("after_commit.defers_discarded_on_rollback", count=len(pending))