Files
service_manager/workers/app/core/database.py
2026-03-03 09:29:53 -07:00

105 lines
2.9 KiB
Python

"""
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()