Files
CRM_AGENTES_CARGA/backend/core/celery_app.py
marcos 02feb973c9 feat(crm): carril con reintentos que entrega los documentos del CRM a EFC
Clon del gateway Anexo22 -> EFC que ya esta en produccion, con los nombres
cambiados. No es una reinterpretacion: la maquina de reintentos de tres capas,
el corte en 4xx y las cuatro guardas de idempotencia se conservan tal cual.

Sin esto, "EFC es la fuente unica" obligaria a llamar a EFC dentro del request
del usuario, y un EFC caido le haria perder su trabajo. Con outbox + Celery, la
subida responde 201 siempre y el sistema entrega cuando EFC vuelve, sin
duplicar: el crm_document_ref viaja con la subida y EFC devuelve 200 con el
documento que ya existia en vez de crear otro. Es lo que cubre el timeout
ambiguo -EFC commiteo y contesto tarde-, donde el CRM no puede saber si entro.

TRES DESVIACIONES DELIBERADAS DEL ORIGINAL, las tres con su razon en el codigo:

1. SAVEPOINT en vez de db.rollback() en el except del encolado. En Anexo22 el
   outbox vive en OTRA base que el pedimento, asi que su rollback solo revertia
   la sesion del outbox. El CRM es mono-base y el encolado corre DENTRO de la
   transaccion del usuario: heredar ese rollback tumbaba la solicitud y el
   expediente recien creados -exactamente lo contrario de best-effort, y en
   silencio-. Lo destapo un test y asi se manifestaba:
   "InvalidRequestError: Instance '<ServiceRequest>' is not persistent within
   this Session". Con el savepoint el fallo deshace solo la fila del outbox.

2. source_table junto a source_id en la guarda _ya_entregado. El CRM tiene DOS
   tablas de documentos con secuencias independientes: crm.documents.id = 5 y
   ops.shipment_documents.id = 5 son documentos distintos. Con el id solo, haber
   entregado el primero haria que el segundo se saltara para siempre sin un solo
   error visible. Hay test que lo fija.

3. EFC_UPLOAD_TIMEOUT_MS aparte de EFC_API_TIMEOUT_MS. Los 8 s de los metadatos
   no alcanzan para un archivo de 25 MB, y el timeout debe quedar POR DEBAJO del
   proxy_read_timeout del nginx de EFC: si el CRM esperara mas, veria un 504
   opaco sin saber si el documento entro.

El cliente HTTP llega con las pruebas que el carril de referencia NO tiene
-verificado: en Anexo22 no hay ni un test de EfcClient._request-, asi que alli
el bucle de reintentos, el backoff y el corte en 4xx nunca se ejercitan. Ese
hueco no se clona: 16 casos contra httpx.MockTransport, sin tocar la red.

Las migraciones se GENERAN y se dejan SIN aplicar.

BLOQUEADO: docker-compose.prod.yml no se toco. El ticket pide las 8 variables en
api, worker y beat, pero Orquestacion.md 13.14 y 4.5 lo prohiben expresamente
("ni tocarlo"), y el orquestador manda sobre el ticket. Queda como paso manual
en el reporte; sin el, worker y beat no ven EFC_API_URL y el carril queda
apagado en produccion, que es degradar limpio y no romper.

Verificacion: pytest tests/ EXIT=0, 151 passed 1 skipped (baseline 70 passed).
Cuatro roturas deliberadas y restauradas: quitando source_table de la guarda el
test de la ambiguedad se puso rojo; reintentando los 4xx los tres tests del
corte dieron "assert 3 == 1"; borrando el objeto local antes de subir cayeron
los tres del corte directo; y disparando el ensure por cualquier 404 se rompio
el test del code.

Ticket: T2026-08-046 (fase 6 de 7)
2026-08-07 18:05:48 -06:00

143 lines
4.4 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",
"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 reconciliacion. El reintento NO es exponencial a proposito —el
# backoff corto vive en el cliente HTTP y el largo es este barrido de intervalo fijo.
"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()