Files
CRM_AGENTES_CARGA/backend/core/celery_app.py
marcos be45f950b8 fix(crm): registra las tareas del carril en Celery, sin lo cual nada se drenaba
Delta de core/celery_app.py que se me quedo fuera al portar el carril. El sintoma era
enganoso: el encolado se veia perfecto -- fila en crm.efc_sync_outbox, status pending,
sin error -- pero el worker rechazaba la entrega con "Received unregistered task of
type 'expediente_gateway.deliver_outbox_row'" y la fila se quedaba en pending con 0
intentos PARA SIEMPRE. Ni el despacho inmediato ni el barrido existian.

  - include: api.v1.modules.crm.expediente_gateway.tasks
  - beat: sweep_outbox y sweep_file_outbox cada 120 s, sweep_expediente_gaps cada
    300 s. Los intervalos son los del carril de referencia de Anexo22. El reintento
    NO es exponencial a proposito: el backoff corto vive en el cliente HTTP y el
    largo es este barrido de intervalo fijo.

Verificado de punta a punta con los dos sistemas cableados. Desde EFC, no desde el
CRM:

  pedimento_app : CRM-2-EXP2026-08-002        (el storage_token del CRM)
  patente/aduana/clave_pedimento/regimen: None  <- provisional de verdad
  pedimento_expediente: estado=provisional, crm_expediente_id=3, folio=EXP2026-08-002
  organizacion  : Aduanasoft (hub_tenant_slug=aduanasoft, is_verified=True)
  licencia      : 5 GB, asi que la subida no falla por cuota

El resolver mapeo tenant 11 -> organizacion por slug, que es el puente 1:1 acordado.
Las cinco tareas quedan registradas en el worker y los tres barridos en el beat.

Ref: T2026-08-046

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-10 12:27:01 -06:00

150 lines
4.9 KiB
Python

import os
import logging
from celery import Celery
from celery.signals import task_postrun, task_prerun
from core.database import reset_rls_context_tokens, rls_company_var, rls_tenant_var
from core.config import settings
logger = logging.getLogger(__name__)
valkey_url = settings.VALKEY_URL
print(f"DEBUG: Celery Broker URL: {valkey_url}")
logger.info(
"Initializing Celery app app_version=%s environment=%s broker=%s",
settings.APP_VERSION,
settings.ENVIRONMENT,
valkey_url,
)
# Configurar broker y backend explícitamente en el constructor
celery_app = Celery(
"app_tasks",
broker=valkey_url,
backend=valkey_url,
)
celery_app.set_default()
_RLS_TOKENS_ATTR = "_rls_context_tokens"
def _coerce_int(value) -> int | None:
if value is None or value == "":
return None
try:
return int(value)
except (TypeError, ValueError):
return None
@task_prerun.connect
def _set_rls_context_from_task(task_id=None, task=None, args=None, kwargs=None, **_):
"""Fija las ContextVars de RLS para la ejecución de la tarea.
Las rutas propagan ``tenant_id`` / ``company_id`` vía Celery headers en
:func:`track_and_dispatch`. Aquí los materializamos en ContextVars para
que cualquier sesión que se abra durante la tarea (incluidos los helpers
``scoped_core_db`` y llamadas directas a ``CoreSessionLocal()``) aplique
``SET LOCAL`` automáticamente.
"""
headers = {}
request = getattr(task, "request", None) if task is not None else None
if request is not None:
headers = getattr(request, "headers", None) or {}
tenant_id = _coerce_int(headers.get("rls_tenant_id"))
company_id = _coerce_int(headers.get("rls_company_id"))
if tenant_id is None:
logger.warning(
"Celery task_prerun missing rls_tenant_id task=%s task_id=%s headers=%s",
getattr(task, "name", "<unknown>"),
task_id,
headers,
)
logger.info(
"Celery task_prerun RLS context task=%s task_id=%s tenant_id=%s company_id=%s",
getattr(task, "name", "<unknown>"),
task_id,
tenant_id,
company_id,
)
token_t = rls_tenant_var.set(tenant_id)
token_c = rls_company_var.set(company_id)
setattr(task, _RLS_TOKENS_ATTR, (token_t, token_c))
@task_postrun.connect
def _reset_rls_context_from_task(task_id=None, task=None, **_):
"""Restaura las ContextVars al terminar la tarea (evita fuga entre tareas
cuando un worker reutiliza el mismo hilo)."""
tokens = getattr(task, _RLS_TOKENS_ATTR, None) if task is not None else None
if tokens is None:
return
token_t, token_c = tokens
logger.info(
"Celery task_postrun clearing RLS context task=%s task_id=%s tenant_id=%s company_id=%s",
getattr(task, "name", "<unknown>"),
task_id,
rls_tenant_var.get(),
rls_company_var.get(),
)
reset_rls_context_tokens(token_t, token_c)
delattr(task, _RLS_TOKENS_ATTR)
celery_app.conf.update(
include=[
"api.v1.modules.core.help_center.tasks",
# Sin esta línea el worker rechaza las tareas del carril con "Received unregistered
# task of type 'expediente_gateway.deliver_outbox_row'": la fila queda en `pending`
# con 0 intentos y NUNCA se drena, aunque el encolado se vea perfecto.
"api.v1.modules.crm.expediente_gateway.tasks",
# Agrega aquí las tareas de tu proyecto:
# "api.v1.modules.example.tasks",
]
)
# Configuraciones adicionales
celery_app.conf.update(
task_track_started=True,
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="America/Mexico_City",
enable_utc=True,
)
celery_app.conf.beat_schedule = {
"sync-from-hub-every-minute": {
"task": "sync_from_hub_task",
"schedule": 60.0, # Run every 60 seconds
},
"cleanup-orphan-layout-imports-hourly": {
"task": "cleanup_orphan_layout_imports",
"schedule": 3600.0,
},
# Carril CRM -> EFC. Los tres intervalos vienen del carril de referencia de Anexo22: 120 s para
# las dos colas y 300 s para la reconciliación. El reintento NO es exponencial a propósito —el
# backoff corto vive en el cliente HTTP y el largo es este barrido de intervalo fijo.
#
# Estos barridos son la red que atrapa la ventana entre el encolado y el commit: el despacho
# inmediato puede llegar al worker antes de que la transacción confirme, no encontrar la fila
# y darse por vencido. Aquí se recoge.
"efc-sweep-outbox-every-2-min": {
"task": "expediente_gateway.sweep_outbox",
"schedule": 120.0,
},
"efc-sweep-file-outbox-every-2-min": {
"task": "expediente_gateway.sweep_file_outbox",
"schedule": 120.0,
},
"efc-sweep-expediente-gaps-every-5-min": {
"task": "expediente_gateway.sweep_expediente_gaps",
"schedule": 300.0,
},
}
if __name__ == "__main__":
celery_app.start()