Skip to main content

HexCore PyPI Downloads

HexCore es un módulo base reutilizable para proyectos Python que implementan arquitectura hexagonal y event handling.


Skills del Proyecto

Este repositorio cuenta con un conjunto de skills adicionales para extender y personalizar funcionalidades en VS Code y otros entornos compatibles. Puedes encontrarlas en:


¿Qué provee HexCore?

  • Clases base y abstracciones para entidades, repositorios, servicios y unidad de trabajo (UoW), siguiendo los principios de DDD y arquitectura hexagonal.
  • Interfaces y contratos para caché, eventos y manejo de dependencias, desacoplando la lógica de negocio de la infraestructura.
  • Utilidades para event sourcing y event dispatching listas para usar en cualquier proyecto.
  • Estructura flexible para que puedas construir microservicios o aplicaciones monolíticas desacopladas y testeables.

Instalación

pip install hexcore

Templates de Proyecto (CLI)

HexCore incluye templates base para bootstrap de proyectos:

hexcore init mi_proyecto --template hexagonal
hexcore init mi_proyecto --template vertical-slice
  • hexagonal: crea src/domain, src/application, src/infrastructure.
  • vertical-slice: crea src/features, src/shared/domain, src/shared/application, src/shared/infrastructure.

En ambos templates se generan:

  • config.py en raíz con repository_discovery_paths de ejemplo.
  • estructura de migraciones con Alembic.

Configuración v2 (Folder-Agnostic)

Desde v2, HexCore usa configuración explícita y no aplica fallback implícito para descubrir repositorios.

1. Configuración visible en raíz

Define un archivo config.py en la raíz del proyecto:

from hexcore.config import ServerConfig

config = ServerConfig(
    repository_discovery_paths={
        "myapp.features.users.infrastructure.repositories",
        "myapp.features.billing.infrastructure.repositories",
    }
)

2. Prioridad para cargar configuración

LazyConfig resuelve módulos en este orden:

  1. HEXCORE_CONFIG_MODULE
  2. HEXCORE_CONFIG_MODULES (lista separada por comas)
  3. módulos configurados por LazyConfig.set_config_modules(...)
  4. config por defecto (raíz del proyecto)

3. Regla de discovery en v2

  • Si repository_discovery_paths está vacío, no se cargan módulos de repositorios.
  • UoW falla con error explícito para evitar comportamiento ambiguo.

Pautas de Colaboración

¡Gracias por tu interés en contribuir a HexCore! Para mantener una colaboración organizada y eficiente, sigue estas pautas:

1. Código de Conducta

Mantén siempre una comunicación respetuosa y profesional. Revisa el Código de Conducta antes de interactuar.

2. Cómo Contribuir

  • Forkea el repositorio y crea una rama para tu contribución (feature/nombre, fix/nombre, etc.).
  • Realiza tus cambios en la rama y asegúrate de que el código funcione correctamente.
  • Escribe una descripción clara y detallada en tu pull request (PR).
  • Relaciona los issues relevantes en tu PR si aplica.

3. Estilo y Formato de Código

  • Sigue la guía de estilos de Python (PEP8).
  • Usa comentarios cuando sea necesario para clarificar el propósito del código.
  • Idealmente, incluye pruebas unitarias para nuevas funciones y arreglos.

4. Revisión de Pull Requests

  • Todos los PR serán revisados antes de ser aceptados. Se pueden solicitar cambios o aclaraciones.
  • Responde a los comentarios de los revisores para facilitar el proceso.

5. Issues

  • Describe claramente los problemas que encuentres.
  • Proporciona información relevante (logs, versiones, pasos para reproducir, etc.).

6. Comunicación

  • Usa los issues y las discusiones para preguntas, sugerencias o propuestas.
  • Si tienes dudas sobre cómo empezar, puedes abrir un issue para orientación.

7. Licencia

Al contribuir, aceptas que tu código será distribuido bajo la licencia del repositorio.


Documentación Básica

Estructura principal

HexCore se organiza con los siguientes submódulos y carpetas:

  • src/domain/: Módulos de dominio, entidades, repositorios, servicios, objetos de valor, eventos, enums y excepciones.
    src/domain/{modulo}/
      ├─ __init__.py
      ├─ entities.py
      ├─ repositories.py
      ├─ services.py
      ├─ value_objects.py
      ├─ events.py
      ├─ enums.py
      └─ exceptions.py
    
  • src/application/: Casos de uso (UseCase) y DTOs para orquestar la lógica de negocio.
  • src/infrastructure/: Implementaciones técnicas (ORM/ODM, CLI, caché, base de datos, repositorios, unit of work).
  • src/infrastructure/database/models/: Modelos SQLAlchemy para base de datos relacional.
  • src/infrastructure/database/documents/: Documentos Beanie para MongoDB.
  • tests/: Pruebas para módulos de dominio e infraestructura.

Abstracciones de Entidades y Eventos

BaseEntity

Clase base para entidades de dominio. Provee atributos comunes y gestión de eventos.

from hexcore.domain.base import BaseEntity

class User(BaseEntity):
    id: UUID
    name: str

DomainEvent y eventos de entidad

Abstracciones para eventos de dominio y para ciclo de vida de entidades.

from hexcore.domain.events import DomainEvent, EntityCreatedEvent

class UserCreatedEvent(EntityCreatedEvent[User]):
    pass

user = User(...)
event = UserCreatedEvent(entity_id=user.id, payload={"name": user.name})

Implementaciones de Repositorios

SQLAlchemyCommonImplementationsRepo

Repositorio genérico para modelos SQLAlchemy con métodos CRUD reutilizables.

class SQLAlchemyCommonImplementationsRepo(BaseSQLAlchemyRepository[T], HasBasicArgs[T, M], t.Generic[T, M]):
    # Métodos principales: get_by_id, list_all, save, delete
    ...

Ejemplo:

class UserRepository(SQLAlchemyCommonImplementationsRepo[UserEntity, UserModel]):
    def __init__(self, uow):
        super().__init__(
            entity_cls=UserEntity,
            model_cls=UserModel,
            not_found_exception=UserNotFoundException,
            fields_resolvers=None,
            fields_serializers=None,
            uow=uow
        )

BeanieODMCommonImplementationsRepo

Repositorio genérico para documentos Beanie ODM (MongoDB) con métodos CRUD reutilizables.

class BeanieODMCommonImplementationsRepo(IBaseRepository[T], HasBasicArgs[T, D], t.Generic[T, D]):
    # Métodos principales: get_by_id, list_all, save, delete
    ...

Ejemplo:

class UserRepository(BeanieODMCommonImplementationsRepo[UserEntity, UserDocument]):
    def __init__(self, uow):
        super().__init__(
            entity_cls=UserEntity,
            document_cls=UserDocument,
            not_found_exception=UserNotFoundException,
            fields_resolvers=None,
            fields_serializers=None,
            uow=uow
        )

Inicialización y Descubrimiento de Documentos Beanie

Para inicializar y registrar automáticamente todos los documentos Beanie:

from hexcore.infrastructure.repositories.orms.beanie.utils import init_beanie_documents

await init_beanie_documents()

Conversión entre modelos/documentos y entidades

Ambos repositorios utilizan to_entity_from_model_or_document para convertir modelos ORM/ODM en entidades del dominio, aplicando resolvers para atributos complejos.



Arquitectura CQRS en HexCore

HexCore v2 integra de forma nativa soporte para el patrón CQRS (Command Query Responsibility Segregation), permitiendo separar conceptual y técnicamente las operaciones de escritura (Commands) de las de lectura (Queries).

¿Cómo funciona el CQRS en HexCore?

El sistema se basa en 3 buses principales, configurables e independientes:

  1. AbstractCommandBus: Despacha inteniones de mutación (Command) a un único AbstractCommandHandler. Los commands modifican el estado del sistema. La transacción la gestiona el handler (el patrón que enseñan los ejemplos de use case); si preferís que la gestione el bus, añadí TransactionMiddleware explícitamente con su uow_factory.
  2. AbstractQueryBus: Despacha intenciones de lectura (Query) a un único AbstractQueryHandler. Retornan un resultado sin mutar el estado.
  3. EventBus: Distribuye eventos de dominio (DomainEvent) a múltiples suscriptores asíncronamente (vía subscribe/publish).

La configuración de CQRS se activa mediante el CQRSConfig en tu ServerConfig:

from hexcore.config import ServerConfig
from hexcore.application.cqrs.config import CQRSConfig, BusConfig

config = ServerConfig(
    cqrs=CQRSConfig(
        command_bus=BusConfig(
            # Sin middlewares por defecto. Los que no necesitan configuración se
            # pueden declarar por dotted path:
            middlewares=["hexcore.infrastructure.cqrs.middlewares.LoggingMiddleware"]
        ),
        # Puedes sustituir el backend en memoria por uno distribuido (Ej: Celery, Procrastinate)
        # backend="mi_app.infrastructure.ProcrastinateCommandBus" 
    )
)

TransactionMiddleware no es el default. Comitea después del handler, así que con un handler que ya gestiona su propia transacción comitearías dos veces. Y necesita un uow_factory construido con tu engine, cosa que no se puede expresar como dotted path: instancialo a mano y pasá el pipeline al bus.

TransactionMiddleware(uow_factory=lambda: SqlAlchemyUnitOfWork(session=session_factory()))

Guía de Migración: De Casos de Uso Clásicos a CQRS

Si ya tienes una aplicación escrita con la abstracción UseCase de HexCore, puedes migrar progresivamente a CQRS sin reescribir todo tu código, utilizando los adaptadores incluidos.

Paso 1: Usar el adaptador UseCaseCommandHandler

En lugar de instanciar un UseCase directamente en tu endpoint, envuélvelo en un comando:

import hexcore.cqrs as cqrs

# 1. Tienes tu UseCase legado, con sus dependencias
class CreateUserUseCase(UseCase[CreateUserCommand, UserDTO]):
    def __init__(self, uow: IUnitOfWork) -> None:
        self.uow = uow

    async def execute(self, request: CreateUserCommand) -> UserDTO:
        # logica legacy
        ...

# 2. Lo registras en el registry de CQRS utilizando el adaptador
registry = cqrs.HandlerRegistry()
registry.register_command_handler(
    CreateUserCommand,
    cqrs.UseCaseCommandHandler(CreateUserUseCase(uow)),
)

# O con un factory, si quieres resolver las dependencias en el momento del dispatch:
registry.register_command_handler(
    CreateUserCommand,
    cqrs.HandlerRegistry.factory(
        lambda: cqrs.UseCaseCommandHandler(CreateUserUseCase(build_uow()))
    ),
)

El método es register_command_handler (y register_query_handler), no register_command.

Paso 2: Consumirlo desde el endpoint usando el CommandBus

@router.post("/users")
async def create_user(
    cmd: CreateUserCommand, 
    # factory inyectado por dependencias
    bus: AbstractCommandBus = Depends(get_command_bus)
):
    # El bus despacha el comando al UseCase legacy de forma transparente
    result = await bus.dispatch(cmd)
    return result

Paso 3 (Final): Refactor a Handler Puro

Cuando estés listo, convierte tu UseCase directamente en un AbstractCommandHandler:

from hexcore.domain.cqrs import AbstractCommandHandler

class CreateUserHandler(AbstractCommandHandler[CreateUserCommand, UserDTO]):
    def __init__(self, uow: IUnitOfWork):
        self.uow = uow

    async def handle(self, command: CreateUserCommand) -> UserDTO:
        # Lógica refactorizada
        return dto

Guía: Almacenamiento Híbrido para Queries (Mongo, Redis, SQL)

La mayor ventaja de CQRS es optimizar las lecturas. HexCore permite que tus Commands escriban en una base de datos relacional (SQLAlchemy) fuertemente normalizada, mientras que los Queries leen de vistas desnormalizadas súper rápidas en MongoDB o Redis.

1. Sincronización a través del EventBus (La Proyección)

Cuando un Command modifica SQL, dispara un Evento de Dominio. Un handler de eventos intercepta este evento y actualiza el "Read Model" en MongoDB o Redis.

from hexcore.domain.events import EventBus, DomainEvent

class UserCreatedEvent(DomainEvent):
    user_id: str
    full_name: str
    email: str

async def project_user_to_mongodb(event: UserCreatedEvent):
    """Proyecta el evento en la BD de lectura (MongoDB)"""
    doc = UserReadDocument(
        id=event.user_id, 
        name=event.full_name, 
        email=event.email
    )
    await doc.insert() # usando Beanie (Mongo)

# Registrar la proyección
event_bus.subscribe(UserCreatedEvent, project_user_to_mongodb)

2. Query Handler leyendo del Read Model

Tu QueryHandler nunca toca SQL, simplemente ataca directamente a Mongo o Redis para máxima velocidad.

from hexcore.domain.cqrs import AbstractQueryHandler, Query

class GetUserQuery(Query[UserReadDTO]):
    user_id: str

class GetUserQueryHandler(AbstractQueryHandler[GetUserQuery, UserReadDTO]):
    async def handle(self, query: GetUserQuery) -> UserReadDTO:
        # Consulta ultra rápida a la colección de lectura en MongoDB
        doc = await UserReadDocument.get(query.user_id)
        
        # O desde Redis:
        # data = await redis_client.get(f"user:{query.user_id}")
        
        if not doc:
            raise UserNotFoundException()
        return UserReadDTO(**doc.dict())

Con este esquema, alcanzas una alta escalabilidad: tus endpoints GET son despachados por el QueryBus respondiendo en milisegundos desde Mongo/Redis, y tus operaciones POST/PUT/DELETE van por el CommandBus transaccionando con ACID en SQL.

3. Definición de Modelos de Lectura (Proyecciones)

Una pregunta frecuente es: ¿HexCore genera automáticamente estos modelos de lectura? La respuesta es No. El patrón CQRS sugiere que tus modelos de lectura estén diseñados específicamente para lo que tus interfaces visuales (UI) o APIs van a consultar. Por lo tanto, debes definir estos modelos manualmente.

Si usas MongoDB (Beanie) para lecturas: Debes crear un documento manual optimizado. Por ejemplo, en lugar de tener joins, puedes embeber datos:

from beanie import Document

# Modelo desnormalizado optimizado para la lectura
class UserReadDocument(Document):
    id: str  # ID referenciado de la tabla SQL
    name: str
    email: str
    total_purchases_cache: int = 0  # Dato pre-calculado por eventos

    class Settings:
        name = "users_read_projections"

Si usas PostgreSQL/MySQL (SQLAlchemy) para lecturas: Si prefieres mantenerte 100% en SQL pero aislando lecturas, puedes crear tablas específicas para proyecciones (Materialized Views o tablas planas):

from sqlalchemy.orm import declarative_base
from sqlalchemy import Column, String, Integer

Base = declarative_base()

class UserReadProjection(Base):
    __tablename__ = 'users_read_projection'
    
    # Modelo totalmente plano sin relaciones ForeignKey complejas
    id = Column(String, primary_key=True)
    full_name = Column(String)
    email = Column(String)
    total_purchases_cache = Column(Integer, default=0)

En ambos casos, es tu EventBus (o un consumidor como Procrastinate) el encargado de instanciar estos modelos manuales y persistirlos cada vez que se detecte un cambio en los modelos de escritura.


Integración con Task Queues (Smart Routing)

HexCore v2 hace que la delegación de tareas a Celery, Procrastinate o ARQ sea increíblemente sencilla y mágica a través del patrón de Smart Routing.

Ya no necesitas instanciar buses separados para código síncrono y asíncrono. HexCore enruta automáticamente tus comandos y eventos hacia las colas de background usando simples decoradores.

1. Decoradores de Background

HexCore ofrece 3 decoradores esenciales en hexcore.domain.cqrs.decorators para cubrir todos los casos de uso:

  1. @background_command(queue="..."): Aplícalo sobre una clase Command. Todo el comando y su handler se ejecutarán asíncronamente en el Worker. Ideal para operaciones pesadas iniciadas por el usuario (ej. Generar un reporte PDF masivo).
  2. @background_handler(queue="..."): Aplícalo sobre una función que maneje un evento (DomainEvent). Permite que un solo evento dispare algunas acciones síncronas rápidas y otras asíncronas lentas (ej. Enviar emails).
  3. @background_task(queue="..."): Aplícalo sobre funciones o utilidades genéricas que no pertenecen al modelo estricto de CQRS (ej. Limpiar base de datos, tareas tipo CRON).

Ejemplos de uso:

from hexcore.domain.cqrs.decorators import background_command, background_handler, background_task
from hexcore.domain.cqrs.commands import Command

# 1. Comando de ejecución asíncrona obligatoria
@background_command(queue="high_priority")
class SendEmailCommand(Command):
    user_id: str
    template: str

# 2. Handler de evento asíncrono
@background_handler(queue="analytics")
async def send_analytics_on_user_created(event: UserCreatedEvent):
    # Lógica costosa...
    pass

# 3. Tarea genérica (Non-CQRS)
@background_task(queue="maintenance")
async def clean_old_records_task(days_retention: int):
    # Limpieza de base de datos...
    pass

2. Usar un Enqueuer (Adaptador)

No hace falta escribirlo: HexCore trae ProcrastinateEnqueuer y CeleryEnqueuer listos, y registran las tareas del consumidor con los nombres hexcore.process_command, hexcore.process_handler y hexcore.process_task.

from hexcore.infrastructure.task_queues.procrastinate_adapter import ProcrastinateEnqueuer

enqueuer = ProcrastinateEnqueuer(procrastinate_app)

Si necesitas otro broker, implementa ITaskEnqueuer (4 métodos). Dos advertencias:

  • enqueue_event no es un pass. Una cola de tareas no puede hacer fan-out a "todos los suscriptores", así que los adaptadores oficiales lanzan NotImplementedError en vez de perder el evento en silencio. Para ejecutar un suscriptor concreto en background usa @background_handler (el EventBus llamará a enqueue_handler); para fan-out real usa RedisEventBus o PostgresEventBus.
  • Si envuelves corutinas en un worker síncrono (Celery), no uses asyncio.run() por tarea: cierra el event loop y deja el pool del AsyncEngine atado a un loop muerto. HexCore usa un loop persistente por proceso, expuesto como run_in_worker_loop(coro).

3. Configurar tus Buses con un Adaptador Oficial

HexCore provee adaptadores plug & play para Celery y Procrastinate. Simplemente importa el enqueuer, pásale tu app y configúralo en los buses de memoria.

Si además deseas persistencia o distribución de Eventos entre múltiples workers/servidores (Pub/Sub), puedes cambiar el InMemoryEventBus por RedisEventBus, PostgresEventBus o RabbitMQEventBus:

from hexcore.application.cqrs.in_memory_buses import InMemoryCommandBus
from hexcore.infrastructure.task_queues.celery_adapter import CeleryEnqueuer
from hexcore.infrastructure.cqrs.redis_bus import RedisEventBus
from celery import Celery
import redis.asyncio as redis

# 1. Adaptador de Task Queue (Para Comandos asíncronos y Event Handlers asíncronos)
app = Celery("my_app", broker="redis://localhost:6379/0")
enqueuer = CeleryEnqueuer(app)
serializer = PydanticSerializer()

command_bus = InMemoryCommandBus(registry=registry, enqueuer=enqueuer, serializer=serializer)

# 2. Event Bus (Para enviar los Eventos por la red)
redis_client = redis.from_url("redis://localhost:6379/0")
event_bus = RedisEventBus(
    redis_client=redis_client,
    serializer=serializer,
    stream_name="hexcore:events",
    group_name="api_workers",
    enqueuer=enqueuer  # <-- Importante para inyectarle la habilidad de Smart Routing
)

Tip: También dispones de PostgresEventBus(pool, serializer, channel_name) que usa LISTEN/NOTIFY nativo si quieres 0 dependencias externas aparte de tu BD de siempre.

4. Ejecutar tareas genéricas

Para encolar la tarea genérica (@background_task), la llamas indirectamente pasándola por el enqueuer:

# Así se encola una tarea genérica sin CQRS:
await enqueuer.enqueue_task(
    task_name=clean_old_records_task.__cqrs_task_name__, 
    payload={"days_retention": 30},
    queue=clean_old_records_task.__cqrs_queue__
)

5. Levantar el Worker (Consumidor Universal)

En el entrypoint de tu worker, register_hexcore_*_tasks autoconfigura las rutas hexcore.process_command, hexcore.process_handler y hexcore.process_task:

import hexcore.cqrs as cqrs
from hexcore.infrastructure.task_queues.celery_adapter import register_hexcore_celery_tasks

# Le pasas el MISMO bus que usa el proceso web: el consumer marca el mensaje como
# "viene del worker", así que el bus lo ejecuta en vez de reencolarlo.
consumer = cqrs.CQRSConsumer(command_bus, event_bus)

register_hexcore_celery_tasks(app, consumer)  # idempotente: llamarla dos veces no revienta

Un worker que sólo procesa comandos puede omitir el event bus: cqrs.CQRSConsumer(command_bus).

Y el entrypoint completo del worker (con scheduler, muerte mutua y SIGTERM) es una llamada:

await cqrs.run_procrastinate_worker(
    procrastinate_app,
    queues=["default", "reactive"],
    scheduler=cqrs.DynamicScheduler(repo, enqueuer, lock_provider=lock),
    on_startup=[lambda: cqrs.seed_cron_jobs(CRON_JOBS)],
)

6. Ejecutar un comando "aquí y ahora"

No hay una API separada para esto, y es a propósito: el contrato es que el bus decide por contexto.

  • Fuera de un worker, un @background_command se encola.
  • Dentro de un worker (es decir, cuando el mensaje viene del CQRSConsumer) el mismo bus lo ejecuta localmente. Por eso puedes —y debes— compartir un único bus entre la app web y el worker.
  • Si un handler despacha a propósito otro @background_command, ese sí se encola: el contexto de worker se consume en el primer dispatch.

Si necesitas comprobarlo desde tu código, cqrs.is_worker_execution() responde si el mensaje en curso viene de una cola.


Tareas Periódicas Dinámicas (Cronjobs en Caliente)

HexCore incluye un DynamicScheduler que te permite programar tareas (@background_task) para que se ejecuten periódicamente. La ventaja clave es que lee la configuración desde un repositorio (como tu Base de Datos), permitiendo activar, desactivar o cambiar los horarios sin necesidad de reiniciar tus servidores.

1. Usa el repositorio SQL de serie (o implementa el tuyo)

Si tu cron vive en SQL —el caso normal— no escribas nada: HexCore trae la tabla, el repositorio y el seed (extra [sql]).

import hexcore.cqrs as cqrs

CRON_JOBS = [
    cqrs.cron_job(clean_old_records_task, "*/5 * * * *", payload={"days_retention": 30}),
    cqrs.cron_job(cerrar_caja, "0 3 * * *"),
]

await cqrs.create_cron_tables()        # o una migración de Alembic
await cqrs.seed_cron_jobs(CRON_JOBS)   # idempotente, y NO pisa lo editado en BD

repo = cqrs.SqlAlchemyCronJobRepository()

cron_job() deriva el task_name de __cqrs_task_name__: escribirlo a mano es cómo se acaba con un cron que encola una tarea ya renombrada, y el fallo aparece en el worker.

Si tu configuración vive en otro sitio (Mongo, Redis, un YAML), implementa ICronJobRepository:

from datetime import datetime

import hexcore.cqrs as cqrs


class MiCronRepository(cqrs.ICronJobRepository):
    async def get_active_jobs(self) -> list[cqrs.CronJobDefinition]:
        return [
            cqrs.CronJobDefinition(
                job_id="clean-db",
                task_name="mi_app.tasks.clean_old_records_task",  # un @background_task
                cron_expression="*/5 * * * *",
                payload={"days_retention": 30},
                queue="maintenance",
            )
        ]

    async def update_last_run(self, job_id: str, run_time: datetime) -> None:
        # Importa implementarlo: es lo que deduplica el encolado entre ticks.
        ...

2. Levanta el Scheduler

En un proceso en background de tu API o en un microservicio separado, arranca el Scheduler inyectándole tu Enqueuer favorito (Celery, Procrastinate):

import hexcore.cqrs as cqrs

scheduler = cqrs.DynamicScheduler(
    repository=repo,
    enqueuer=enqueuer,
    tick_interval_seconds=60,
)

# Lo normal es dejar que el runner lo supervise junto al worker: si uno de los dos
# muere, se cancela el otro y el proceso sale para que el orquestador lo reinicie.
await cqrs.run_procrastinate_worker(procrastinate_app, scheduler=scheduler)

El Scheduler evalúa las expresiones con croniter y delega la carga pesada al enqueuer. Tu Worker no necesita saber de horarios, sólo ejecuta las tareas cuando le llegan.

Cómo decide si toca ejecutar. No compara contra el minuto actual, sino que busca si hubo alguna ocurrencia entre la última ejecución (last_run_at) y ahora. Dos consecuencias que importan:

  • Un minuto saltado por drift del tick no pierde la ejecución.
  • update_last_run deduplica de verdad, así que un tick_interval_seconds < 60 no duplica el encolado dentro del mismo proceso. Entre réplicas sí hace falta lock, y el scheduler emite un RuntimeWarning si detecta tick sub-minuto sin lock_provider.

catch_up_window_seconds (1 hora por defecto) acota el catch-up: un scheduler que estuvo caído una semana no dispara ocurrencias antiguas.

3. Distributed Locks (Evitar ejecuciones dobles)

Si corres tu aplicación en múltiples contenedores o réplicas (ej. Kubernetes), podrías tener múltiples instancias del DynamicScheduler ejecutándose al mismo tiempo. Para evitar que el mismo cronjob se encole dos veces en el mismo minuto, HexCore soporta Locks Distribuidos.

Puedes inyectar un proveedor de locks (ILockProvider) usando Redis o PostgreSQL (si lo usas como tu DB). Al inyectarlo, el Scheduler bloqueará atómicamente la tarea a través de toda tu red.

Usando Redis

from hexcore.infrastructure.cqrs.redis_lock import RedisLockProvider
import redis.asyncio as redis

redis_client = redis.from_url("redis://localhost:6379/0")
lock_provider = RedisLockProvider(redis_client)

scheduler = DynamicScheduler(
    repository=repo, 
    enqueuer=enqueuer, 
    lock_provider=lock_provider
)

Usando PostgreSQL (asyncpg)

Si usas Procrastinate o bases de datos SQL y no quieres levantar Redis:

from hexcore.infrastructure.cqrs.postgres_lock import PostgresLockProvider

lock_provider = PostgresLockProvider(my_asyncpg_pool)
await lock_provider.setup() # Crea la tabla y el índice, y purga lo expirado

scheduler = DynamicScheduler(
    repository=repo, 
    enqueuer=enqueuer, 
    lock_provider=lock_provider
)

El provider purga las filas expiradas solo (en setup() y cada 100 adquisiciones), así que la tabla de locks no crece sin límite. purge_expired() es pública si prefieres purgar desde un job propio, y purge_every=0 desactiva la purga automática.

Qué pasa si el lock no responde

Si Redis (o Postgres) se cae, acquire_lock no puede decidir, y las dos respuestas posibles son malas de formas distintas. La decisión es tuya y explícita:

RedisLockProvider(redis_client, on_error="skip")   # default: no correr. El cron se detiene.
RedisLockProvider(redis_client, on_error="raise")  # propagar, para que el supervisor lo vea.

En los logs, "no pude decidir" es critical y "el lock estaba tomado por otra réplica" —el caso normal— es debug.


Referencias

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

hexcore-5.0.0.tar.gz (192.0 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

hexcore-5.0.0-py3-none-any.whl (241.9 kB view details)

Uploaded Python 3

File details

Details for the file hexcore-5.0.0.tar.gz.

File metadata

  • Download URL: hexcore-5.0.0.tar.gz
  • Upload date:
  • Size: 192.0 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/7.0.0 CPython/3.12.13

File hashes

Hashes for hexcore-5.0.0.tar.gz
Algorithm Hash digest
SHA256 f15acfc1d6cfea8381b81f3d9e319d587bcca34d2f102f33debdd44f7069d72a
MD5 33e7df056ce0bae6ccf44ffd8d1a7bd8
BLAKE2b-256 186e07f4629aa62444a3cc86d02840570945fe67d6c8d1675de72cd2d6ca0bfd

See more details on using hashes here.

File details

Details for the file hexcore-5.0.0-py3-none-any.whl.

File metadata

  • Download URL: hexcore-5.0.0-py3-none-any.whl
  • Upload date:
  • Size: 241.9 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/7.0.0 CPython/3.12.13

File hashes

Hashes for hexcore-5.0.0-py3-none-any.whl
Algorithm Hash digest
SHA256 bd988628e896c7d7cfea5ae30e4bbf416d506f7bdad58ab5ede07ec9f7f809fa
MD5 69870f090d1d76c8b6975c63c146f355
BLAKE2b-256 c785fc09aff5be69c3bb2f113eaed283bd0ab2d131772d22c0600705dbe71696

See more details on using hashes here.

Release history Release notifications | RSS feed

9.0.1

2 files

9.0.0

2 files

8.0.0

2 files

7.0.0

2 files

6.2.1

2 files

6.2.0

2 files

6.1.0

2 files

6.0.2

2 files

6.0.1

2 files

6.0.0

2 files

This release

5.0.0 This release

2 files

4.0.0

2 files

3.0.0

2 files

2.5.0

2 files

2.4.0

2 files

2.3.0

2 files

2.2.0

2 files

2.1.0

2 files

2.0.6

2 files

2.0.5

2 files

2.0.4

2 files

2.0.3

2 files

2.0.2

2 files

2.0.1

2 files

2.0.0

2 files

1.8.0

2 files

1.7.0

2 files

1.6.8

2 files

1.6.7

2 files

1.6.6

2 files

1.6.5

2 files

1.6.4

2 files

1.6.3

2 files

1.6.2

2 files

1.6.1

2 files

1.6.0

2 files

1.5.1

2 files

1.5.0

2 files

1.4.2

2 files

1.4.1

2 files

1.4.0

2 files

1.3.2

2 files

1.3.1

2 files

1.3.0

2 files

1.2.0

2 files

1.1.0

2 files

1.0.2

2 files

1.0.1

2 files

1.0.0

2 files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page