/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/db/session.py
391 строка
12 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""SQLAlchemy async session management с production-ready конфигурацией.""" from __future__ import annotations import asyncio import time from collections.abc import AsyncGenerator, AsyncIterator, Awaitable, Callable from contextlib import asynccontextmanager from typing import Any, cast import structlog from sqlalchemy import text from sqlalchemy.exc import ( DisconnectionError, OperationalError, SQLAlchemyError, ) from sqlalchemy.ext.asyncio import ( AsyncSession, async_sessionmaker, create_async_engine, ) from sqlalchemy.orm import DeclarativeBase from sqlalchemy.pool import AsyncAdaptedQueuePool from src.config import get_settings logger = structlog.get_logger() class Base(DeclarativeBase): """Базовый класс для всех моделей.""" pass # ============================================================================ # Engine Configuration # ============================================================================ _settings = get_settings() # Получаем настройки пула с fallback на дефолтные значения _pool_size: int = cast(int, getattr(_settings, "db_pool_size", 10)) _max_overflow: int = cast(int, getattr(_settings, "db_max_overflow", 20)) _pool_timeout: int = cast(int, getattr(_settings, "db_pool_timeout", 30)) # Production-ready engine configuration engine = create_async_engine( _settings.database_url, echo=_settings.debug, pool_pre_ping=True, # Проверка соединения перед использованием pool_size=_pool_size, max_overflow=_max_overflow, pool_timeout=_pool_timeout, pool_recycle=1800, # Пересоздавать соединения каждые 30 минут connect_args={ "server_settings": { "jit": "off", # Отключить JIT для стабильности } } if "postgresql" in _settings.database_url else {}, ) async_session_maker = async_sessionmaker( engine, class_=AsyncSession, expire_on_commit=False, autocommit=False, autoflush=False, ) # ============================================================================ # Session Management # ============================================================================ @asynccontextmanager async def get_session() -> AsyncIterator[AsyncSession]: """ Context manager для получения async session. Автоматически коммитит изменения при успешном завершении, или откатывает при ошибке. Usage: async with get_session() as session: result = await session.execute(query) # автоматический commit при выходе из context """ async with async_session_maker() as session: try: yield session await session.commit() logger.debug("session.committed") except Exception as e: await session.rollback() logger.error( "session.rollback", error=str(e), error_type=type(e).__name__, ) raise async def get_db() -> AsyncGenerator[AsyncSession, None]: """ FastAPI dependency для получения async session. Usage в FastAPI endpoint: @app.get("/items") async def get_items(session: AsyncSession = Depends(get_db)): result = await session.execute(query) """ async with get_session() as session: yield session # ============================================================================ # Transaction Helpers # ============================================================================ @asynccontextmanager async def transaction() -> AsyncIterator[AsyncSession]: """ Explicit transaction context manager. Используется когда нужен finer control над транзакциями. Usage: async with transaction() as session: await session.execute(insert_stmt1) await session.execute(insert_stmt2) # commit происходит автоматически """ async with async_session_maker() as session: try: yield session await session.commit() except Exception: await session.rollback() raise async def execute_in_transaction( operation: Callable[[AsyncSession], Awaitable[Any]], ) -> Any: """ Выполнить операцию в транзакции с retry logic. Operation принимает session явно: async def my_operation(session: AsyncSession) -> Result: repo = MyRepository(session, user_id, workspace_id) return await repo.create(...) result = await execute_in_transaction(my_operation) Автоматически повторяет при transient ошибках (deadlocks, connection issues) с экспоненциальной задержкой. Args: operation: async callable, принимающий AsyncSession Returns: Результат операции. """ max_retries = 3 retry_delay = 0.1 for attempt in range(max_retries): try: async with transaction() as session: return await operation(session) except (OperationalError, DisconnectionError) as e: if attempt == max_retries - 1: logger.error( "transaction.retry_exhausted", attempt=attempt + 1, error=str(e), ) raise logger.warning( "transaction.retry", attempt=attempt + 1, max_retries=max_retries, error=str(e), ) await asyncio.sleep(retry_delay * (2**attempt)) # Exponential backoff return None # недостижимо, но для type checker # ============================================================================ # Database Lifecycle # ============================================================================ async def init_db() -> None: """ Создать таблицы в базе данных. Используется только для разработки! В production используйте Alembic миграции. """ try: async with engine.begin() as conn: await conn.run_sync(Base.metadata.create_all) logger.info("database.tables_created") except SQLAlchemyError as e: logger.error( "database.init_failed", error=str(e), error_type=type(e).__name__, ) raise async def close_db() -> None: """ Закрыть соединение с базой данных. Вызывается при остановке приложения. Gracefully закрывает все соединения в пуле. """ try: await engine.dispose() logger.info("database.connections_closed") except Exception as e: logger.warning( "database.close_error", error=str(e), error_type=type(e).__name__, ) # ============================================================================ # Helper: получить AsyncAdaptedQueuePool с правильной типизацией # ============================================================================ def _get_queue_pool() -> AsyncAdaptedQueuePool: """ Получить engine pool с правильной типизацией. AsyncAdaptedQueuePool — async-compatible версия QueuePool, которая автоматически используется create_async_engine(). """ return cast(AsyncAdaptedQueuePool, engine.pool) # ============================================================================ # Health Check # ============================================================================ async def check_database_health() -> dict[str, Any]: """ Проверить здоровье базы данных. Returns: Dict с информацией о состоянии: - status: "healthy" | "degraded" | "unhealthy" - latency_ms: время отклика - pool_size: текущий размер пула - checkedout: количество активных соединений - error: описание ошибки (если есть) """ start_time = time.time() try: # Используем direct session для health check (без auto-commit) async with async_session_maker() as session: await session.execute(text("SELECT 1")) latency_ms = (time.time() - start_time) * 1000 pool = _get_queue_pool() return { "status": "healthy", "latency_ms": round(latency_ms, 2), "pool_size": pool.size(), "checkedout": pool.checkedout(), "overflow": pool.overflow(), } except OperationalError as e: latency_ms = (time.time() - start_time) * 1000 return { "status": "unhealthy", "latency_ms": round(latency_ms, 2), "error": str(e), "error_type": "OperationalError", } except Exception as e: latency_ms = (time.time() - start_time) * 1000 return { "status": "degraded", "latency_ms": round(latency_ms, 2), "error": str(e), "error_type": type(e).__name__, } # ============================================================================ # Utility Functions # ============================================================================ async def get_engine_stats() -> dict[str, Any]: """ Получить статистику engine и connection pool. Полезно для мониторинга и debugging. Returns: Dict с детальной информацией о пуле соединений. """ pool = _get_queue_pool() return { "pool_size": pool.size(), "checkedout": pool.checkedout(), "overflow": pool.overflow(), "checkedin": pool.checkedin(), "status": pool.status(), "echo": engine.echo, "url": str(engine.url), } async def execute_raw_sql( sql: str, params: dict[str, Any] | None = None, ) -> list[dict[str, Any]]: """ Выполнить raw SQL запрос и вернуть результаты. ⚠️ Используйте с осторожностью! Предпочитайте SQLAlchemy ORM/Core. Note: Для SELECT запросов. Для UPDATE/DELETE используйте transaction(). Args: sql: SQL query string params: Параметры для prepared statement Returns: List of dicts с результатами. """ async with get_session() as session: result = await session.execute(text(sql), params or {}) # Преобразуем Row objects в dicts return [dict(row._mapping) for row in result.fetchall()] # ============================================================================ # Exports # ============================================================================ __all__ = [ # noqa: RUF022 # Base "Base", # Engine & Session "engine", "async_session_maker", # Session Management "get_session", "get_db", "transaction", "execute_in_transaction", # Lifecycle "init_db", "close_db", # Health & Monitoring "check_database_health", "get_engine_stats", # Utilities "execute_raw_sql", ]