""" Database Configuration for Workers - ServiceManagerWeb Async database session management para Celery workers """ from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession, async_sessionmaker from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column from sqlalchemy import DateTime, func from sqlalchemy import types as sa_types from sqlalchemy.types import TypeDecorator, CHAR from sqlalchemy.dialects.postgresql import UUID as PG_UUID from contextlib import asynccontextmanager from typing import AsyncGenerator import uuid from datetime import datetime class GUID(TypeDecorator): """UUID portable: UUID nativo en Postgres, CHAR(36) en otros dialectos (SQLite para tests).""" impl = CHAR cache_ok = True def load_dialect_impl(self, dialect): if dialect.name == "postgresql": return dialect.type_descriptor(PG_UUID(as_uuid=True)) return dialect.type_descriptor(CHAR(36)) def process_bind_param(self, value, dialect): if value is None: return None if dialect.name == "postgresql": return value if isinstance(value, uuid.UUID): return str(value) return str(uuid.UUID(str(value))) def process_result_value(self, value, dialect): if value is None: return None if not isinstance(value, uuid.UUID): return uuid.UUID(str(value)) return value from app.core.config import get_settings settings = get_settings() # Create async engine para workers engine = create_async_engine( settings.DATABASE_URL, echo=False, # Menos verbose en workers pool_size=5, max_overflow=10, pool_pre_ping=True, pool_recycle=3600, ) # Create session factory AsyncSessionLocal = async_sessionmaker( engine, class_=AsyncSession, expire_on_commit=False, autoflush=True, autocommit=False ) class Base(DeclarativeBase): """Base class para todos los modelos SQLAlchemy - compartida con backend.""" # Columnas comunes para auditoría id: Mapped[uuid.UUID] = mapped_column(primary_key=True, default=uuid.uuid4) created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) updated_at: Mapped[datetime] = mapped_column( DateTime(timezone=True), server_default=func.now(), onupdate=func.now() ) @asynccontextmanager async def get_async_session_context() -> AsyncGenerator[AsyncSession, None]: """ Context manager para obtener sesión de base de datos en workers. Usage: async with get_async_session_context() as db: # Usar db aquí pass Yields: AsyncSession: Sesión de base de datos """ async with AsyncSessionLocal() as session: try: yield session await session.commit() except Exception: await session.rollback() raise finally: await session.close()