feat(crm): carril hacia EFC montado sobre el expediente existente (crm.cases)

Rebase del lado emisor de T2026-08-046 sobre esta rama. La entrega anterior partia
de feature/crm-cumplimiento-pdf (16-jul), 40 commits atras, y por eso construyo un
expediente PARALELO -- crm.expedientes con su propio generador de folio y su propia
migracion -- que duplicaba el que ya existe aqui. Dos expedientes y dos secuencias
peleando por el mismo namespace EXP no se fusionan; se tira el nuestro.

La estructura del expediente es de esta rama y no se toca: crm.cases es el
expediente, su folio vive en `reference` y el consecutivo lo reserva
crm/common/folios.py con bloqueo de fila. Nuestro aporte es SOLO la conexion:

  - crm.cases gana seis columnas efc_* (espejo de EFC, nunca el handle) y nada mas;
  - crm.efc_sync_outbox y crm.efc_file_outbox, el outbox transaccional, con
    expediente_ref -> crm.cases.id;
  - core/efc_client.py y crm/expediente_gateway/ (outbox, reintentos, barridos),
    clonados del gateway Anexo22 -> EFC que ya corre en produccion;
  - las ocho variables EFC_* en config. EFC_API_URL vacia = carril apagado.

Verificado contra la base real: next_folio(...,'EXP',None,with_direction=False)
devuelve EXP2026-08-001, identico al formato que el contrato con EFC exige, y
storage_token da CRM-{company}-{folio} de 22 caracteres sobre los 25 de
pedimento_app.

Se corrige un error del docstring de storage_token: decia que cabian companies de
7 digitos y son 6 (4+7+1+14 = 26 > 25). Ahora valida y falla ruidosamente en vez de
entregar un token recortado, que apuntaria a la carpeta de otro expediente y
mezclaria documentos en silencio.

El revision id de la migracion tirada (e6f7a8b9c0d1) chocaba con crm_catalog_items
de esta rama: dos migraciones distintas con el mismo id habrian roto alembic al
fusionar. La nueva es c5d6e7f8a9b0, aditiva sobre d4e5f6a7b8c9.

PENDIENTE: falta el pegamento que invocaba el carril desde los flujos de la app
(alta del provisional al mintear el folio, subida de documento -> outbox, rutas en
el router y UI). Por eso test_efc_outbox, test_gateway_rutas y tres casos de
test_contrato_efc todavia no colectan. El carril no esta cableado al router, asi
que la app funciona igual: backend y frontend responden 200.

Ref: T2026-08-046

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-08-10 10:35:44 -06:00
parent dfb3c8a08d
commit 5c4df590d4
19 changed files with 3411 additions and 2 deletions

View File

@@ -1,4 +1,4 @@
from sqlalchemy import ForeignKey, Integer, String, text
from sqlalchemy import ForeignKey, Integer, String, Text, text
from sqlalchemy.orm import Mapped, mapped_column
from api.v1.common.base_models import TenantScopedMixin, TimestampMixin
@@ -26,3 +26,27 @@ class Case(Base, TenantScopedMixin, TimestampMixin):
status: Mapped[str] = mapped_column(String(20), nullable=False, server_default=text("'abierto'"), index=True)
created_by: Mapped[str | None] = mapped_column(String(64), nullable=True)
updated_by: Mapped[str | None] = mapped_column(String(64), nullable=True)
# ── Espejo del carril hacia EFC (T2026-08-046) ──────────────────────────────────────────
# EFC es la fuente única de los documentos del expediente: cada expediente se refleja allá
# como un *pedimento provisional* y los archivos viven en su MinIO, no en el del CRM.
#
# Estas columnas son un ESPEJO, nunca el handle. El handle con el que el CRM habla de este
# expediente es su ``id`` y su ``reference``: ``efc_pedimento_id`` es un caché de la
# resolución, y del lado de EFC el ``pedimento_app`` es mutable —se reescribe al completar
# el provisional con la data aduanera real—, así que apoyarse en él rompería justo cuando
# llegue esa data. La liga vive en EFC, en la tabla desechable ``pedimento_expediente``.
efc_organizacion_id: Mapped[str | None] = mapped_column(String(36), nullable=True)
efc_pedimento_id: Mapped[str | None] = mapped_column(String(36), nullable=True)
# INMUTABLE una vez asignado: es la carpeta de MinIO donde EFC guarda los objetos de este
# expediente. Que no cambie nunca es lo que permite completar el pedimento sin mover ni un
# archivo.
efc_storage_token: Mapped[str | None] = mapped_column(String(25), nullable=True)
# PENDING | LINKED | FAILED
efc_link_state: Mapped[str] = mapped_column(
String(20), nullable=False, server_default=text("'PENDING'"), index=True
)
# El diagnóstico se guarda en la fila para que se vea en la ficha del expediente, sin
# obligar a nadie a ir a los logs del worker.
efc_error_code: Mapped[str | None] = mapped_column(String(60), nullable=True)
efc_error_detail: Mapped[str | None] = mapped_column(Text, nullable=True)

View File

@@ -0,0 +1,64 @@
"""Catálogo CERRADO de tipos de documento que EFC acepta del CRM.
Estas 22 claves son **exactamente** las de ``TIPOS_DOCUMENTO_CRM`` en
``api/record/views_integrations_crm.py`` de EFC. La lista está duplicada a mano en dos repos con
despliegue independiente, así que ``tests/test_doc_types_paridad.py`` la fija: si alguien agrega un
tipo de un solo lado, ese test se pone rojo antes de que un documento se rechace en producción.
Por qué es un conjunto cerrado y no texto libre, a diferencia del carril de Anexo22 —que manda el
tipo suelto y deja que EFC lo resuelva por nombre—: en el CRM ``doc_type`` es ``String(60)`` /
``String(30)`` **sin validación de backend**, los catálogos viven solo en TypeScript
(``frontend/src/lib/api/crm/format.ts``). Un typo crearía un ``DocumentType`` basura en el catálogo
**global** de EFC, que es compartido por todas las organizaciones y no se limpia solo.
Las tres fuentes del CRM y su origen:
- ``crm.documents`` → ``DOC_TYPES`` de ``format.ts``
- ``ops.shipment_documents`` → ``SHIPMENT_DOC_TYPES`` del mismo archivo
- ``fin.invoices`` → el PDF de factura (``factura_venta``)
``otro`` existe en las dos listas del CRM y significa lo mismo en ambas: es una sola entrada.
"""
# --- crm.documents ---------------------------------------------------------------------------
_TIPOS_DOCUMENTOS_CLIENTE = (
"constancia_fiscal",
"acta_constitutiva",
"identificacion",
"comprobante_domicilio",
"contrato",
"presentacion",
"certificacion",
"licencia",
"convenio",
"tarifario",
)
# --- ops.shipment_documents ------------------------------------------------------------------
_TIPOS_DOCUMENTOS_EMBARQUE = (
"MBL",
"HBL",
"MAWB",
"HAWB",
"CMR",
"factura_comercial",
"packing_list",
"carta_encomienda",
"carta_garantia",
"certificado_permiso",
)
# --- fin.invoices ----------------------------------------------------------------------------
_TIPOS_FACTURACION = ("factura_venta",)
# --- común a varias fuentes -------------------------------------------------------------------
_TIPOS_COMUNES = ("otro",)
EFC_DOC_TYPES: frozenset[str] = frozenset(
_TIPOS_DOCUMENTOS_CLIENTE + _TIPOS_DOCUMENTOS_EMBARQUE + _TIPOS_FACTURACION + _TIPOS_COMUNES
)
def is_valid_doc_type(doc_type: str | None) -> bool:
"""``True`` si EFC va a aceptar ese tipo. Se valida en el CRM para no gastar un viaje de red."""
return bool(doc_type) and doc_type in EFC_DOC_TYPES

View File

@@ -0,0 +1,131 @@
"""Outbox transaccional del carril CRM Agentes de Carga -> EFC.
DOS tablas separadas POR PROPÓSITO, igual que en el carril de referencia de Anexo22: una para los
expedientes (metadatos, JSON) y otra para los archivos (binarios que viven en MinIO y se referencian
por su ``s3_key``). Un worker de Celery las drena hacia EFC con reintentos.
**Diferencia con el original, y es necesaria:** aquí las filas se insertan en la MISMA transacción
que el expediente o el documento, porque el CRM es mono-base. En Anexo22 el outbox vivía en otra
base que el pedimento, y ese doble-commit es justamente lo que obligó a inventar el barrido de
huecos. Aquí el barrido se conserva —cubre lo creado antes de activar la integración y cualquier
crash— pero deja de ser el parche de una ventana estructural.
"""
from datetime import datetime
from typing import Optional
from sqlalchemy import JSON, Boolean, DateTime, ForeignKey, Index, Integer, String, Text, text
from sqlalchemy.orm import Mapped, mapped_column
from api.v1.common.base_models import TenantScopedMixin, TimestampMixin
from core.database import Base
# Tipo de trabajo (columna kind) del outbox de EXPEDIENTES.
KIND_EXPEDIENTE = "expediente"
KIND_COMPLETAR = "completar"
# Tipos del outbox de ARCHIVOS (efc_file_outbox).
FILE_KIND_DOCUMENTO = "documento"
# Tablas de origen posibles de un archivo. El CRM tiene DOS tablas de documentos con secuencias
# independientes, así que `source_id` por sí solo es ambiguo: crm.documents.id = 5 y
# ops.shipment_documents.id = 5 coexisten.
SOURCE_CRM_DOCUMENTS = "crm.documents"
SOURCE_OPS_SHIPMENT_DOCUMENTS = "ops.shipment_documents"
SOURCE_FIN_INVOICES = "fin.invoices"
# Estados (columna status).
STATUS_PENDING = "pending"
STATUS_SENT = "sent"
STATUS_FAILED = "failed"
# Tope de reintentos antes de marcar 'failed' (reconciliación / reintento manual).
# Heredado del carril de Anexo22. Con barridos de 120 s son ~17 minutos de insistencia antes de
# rendirse y dejar la fila visible para que una persona la reintente a mano.
MAX_ATTEMPTS = 8
class EfcSyncOutbox(Base, TenantScopedMixin, TimestampMixin):
"""Cola de metadatos hacia EFC: crear el expediente provisional y completarlo."""
__tablename__ = "efc_sync_outbox"
__table_args__ = (
Index("ix_crm_efc_sync_outbox_status", "status"),
Index("ix_crm_efc_sync_outbox_kind_status", "kind", "status"),
{"schema": "crm"},
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
kind: Mapped[str] = mapped_column(String(20), nullable=False)
# Datos para construir el request a EFC (folio, storage_token, tenant slug, company, y la data
# aduanera si el kind es 'completar').
payload: Mapped[dict] = mapped_column(JSON, nullable=False)
# id local del expediente (crm.cases.id) que originó la fila.
expediente_ref: Mapped[Optional[int]] = mapped_column(Integer, nullable=True, index=True)
# Ciclo de vida.
status: Mapped[str] = mapped_column(String(10), nullable=False, server_default=text(f"'{STATUS_PENDING}'"))
attempts: Mapped[int] = mapped_column(Integer, nullable=False, server_default=text("0"))
last_error: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
sent_at: Mapped[Optional[datetime]] = mapped_column(DateTime, nullable=True)
# Acuse de EFC al confirmar (trazabilidad).
efc_pedimento_id: Mapped[Optional[str]] = mapped_column(String(36), nullable=True)
class EfcFileOutbox(Base, TenantScopedMixin, TimestampMixin):
"""Cola de ARCHIVOS hacia EFC.
El binario vive en el MinIO del CRM (durable); esta fila referencia su ``s3_key`` y el expediente
destino. El worker lo sube a EFC y, con ``delete_local`` (corte directo), BORRA la copia local al
confirmar la entrega.
``delete_local`` **es el mecanismo de «EFC es la fuente única»**: "solo EFC" es el estado FINAL
(eventual), no el inmediato. Entre que el usuario sube el archivo y que EFC lo confirma, la copia
local es lo único que hay, y borrarla antes perdería el archivo si la entrega fallara.
``source_table`` es un añadido necesario sobre el original de Anexo22, que solo llevaba
``source_id``. El CRM tiene dos tablas de documentos con secuencias independientes, así que un
entero solo es ambiguo entre ellas. Es el mismo problema que Anexo22 resolvió con su mapa por
``kind``, y su comentario dice qué pasa si se ignora: un UPDATE con el id de otra tabla **vacía la
columna de un documento ajeno** que tuviera ese mismo entero — daño en el dato de otro, sin un
solo error visible. Un ``(kind, source_table)`` que no esté en el mapa **no toca nada**, en vez
de caer por omisión.
"""
__tablename__ = "efc_file_outbox"
__table_args__ = (
Index("ix_crm_efc_file_outbox_status", "status"),
Index("ix_crm_efc_file_outbox_kind_status", "kind", "status"),
Index("ix_crm_efc_file_outbox_source", "source_table", "source_id"),
{"schema": "crm"},
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
kind: Mapped[str] = mapped_column(String(30), nullable=False)
# Objeto en MinIO a subir + metadata para el upload a EFC.
s3_key: Mapped[str] = mapped_column(String(1024), nullable=False)
file_name: Mapped[str] = mapped_column(String(255), nullable=False)
content_type: Mapped[Optional[str]] = mapped_column(String(100), nullable=True)
efc_tipo: Mapped[str] = mapped_column(String(40), nullable=False) # tipo de documento en EFC
# Origen: la pareja (tabla, id) desambigua entre las dos secuencias de documentos del CRM.
source_table: Mapped[str] = mapped_column(String(30), nullable=False)
source_id: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
# El handle autoritativo que viaja a EFC y garantiza la idempotencia del lado de allá.
crm_document_ref: Mapped[Optional[str]] = mapped_column(String(64), nullable=True)
expediente_ref: Mapped[int] = mapped_column(
Integer, ForeignKey("crm.cases.id"), nullable=False, index=True
)
delete_local: Mapped[bool] = mapped_column(Boolean, nullable=False, server_default=text("true"))
# Ciclo de vida.
status: Mapped[str] = mapped_column(String(10), nullable=False, server_default=text(f"'{STATUS_PENDING}'"))
attempts: Mapped[int] = mapped_column(Integer, nullable=False, server_default=text("0"))
last_error: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
sent_at: Mapped[Optional[datetime]] = mapped_column(DateTime, nullable=True)
efc_document_id: Mapped[Optional[str]] = mapped_column(String(36), nullable=True)

View File

@@ -0,0 +1,58 @@
"""Endpoints de operación y observabilidad del carril CRM -> EFC.
Tablero mínimo para ver y reintentar la entrega de expedientes y documentos a EFC. Autenticado con
el auth normal del CRM y acotado por tenant/company, como el resto del módulo.
Montado bajo ``/v1/crm`` → ``/v1/crm/expediente-gateway/...``
"""
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy.orm import Session
from core.database import get_core_db
from core.security import get_current_user
from . import service
router = APIRouter(prefix="/expediente-gateway", tags=["EFC Gateway (ops)"])
@router.get("/outbox")
def list_outbox(
company_id: int = Query(..., description="Company ID"),
tipo: str | None = Query(None, description="Filtrar por tabla: sync|file"),
status: str | None = Query(None, description="Filtrar por status: pending|sent|failed"),
limit: int = Query(100, ge=1, le=500),
current_user: dict = Depends(get_current_user),
db: Session = Depends(get_core_db),
):
"""Filas de los dos outbox, para ver los fallos y su ``last_error``."""
return service.list_outbox(db, current_user["tenant_id"], company_id, tipo, status, limit)
@router.post("/outbox/{outbox_id}/retry")
def retry_outbox(
outbox_id: int,
company_id: int = Query(..., description="Company ID"),
tipo: str = Query("file", description="Tabla de la fila: sync|file"),
current_user: dict = Depends(get_current_user),
db: Session = Depends(get_core_db),
):
"""Reintento manual de una fila: la resetea a ``pending`` y la re-despacha.
Una fila inexistente devuelve **404 con mensaje específico**, no un 200 silencioso: el frontend
pinta el botón de reintento según lo que reciba, y un 200 le haría creer que la entrega volvió a
la cola cuando no hay nada que entregar.
"""
ok = service.retry_outbox_row(db, outbox_id, current_user["tenant_id"], company_id, tipo)
if not ok:
raise HTTPException(status_code=404, detail="Fila de outbox no encontrada")
return {"status": "requeued", "id": outbox_id}
@router.get("/metrics")
def metrics(
company_id: int = Query(..., description="Company ID"),
current_user: dict = Depends(get_current_user),
db: Session = Depends(get_core_db),
):
"""Conteo de los dos outbox por status (pending/sent/failed) para monitoreo."""
return service.outbox_metrics(db, current_user["tenant_id"], company_id)

View File

@@ -0,0 +1,731 @@
"""Carril CRM Agentes de Carga -> EFC: encolado, entrega y reconciliación.
Clon del gateway de Anexo22 (``anexo22/.../pedimentos/pedimento_gateway/service.py``), que es el
carril de referencia ya en producción. Quien conozca uno debe poder leer el otro, así que la tabla
de equivalencias va aquí:
====================================== ======================================
Anexo22 CRM
====================================== ======================================
``replicate_pedimento_best_effort`` ``replicate_expediente_best_effort``
``_enqueue_pedimento_outbox`` ``_enqueue_expediente_outbox``
``_dispatch_delivery`` igual
``deliver_row`` / ``_deliver_pedimento`` ``deliver_row`` / ``_deliver_expediente``
``_register_failure`` **idéntico**
``_ya_entregado(source_id, kind)`` ``_ya_entregado(source_table, source_id, kind)``
``deliver_file_row`` **idéntico**, con ensure-then-upload y ``delete_local``
``_register_file_failure`` **idéntico**
``_resolve_org_id`` + ``_org_id_cache`` igual — dict módulo-global, por worker, sin invalidación
``list_outbox`` / ``retry_outbox_row`` / ``outbox_metrics`` igual, para las dos tablas
``find_pedimento_gaps`` ``find_expediente_gaps``
====================================== ======================================
**La máquina de reintentos tiene tres capas y las tres se conservan:**
1. En el cliente HTTP: 3 intentos, backoff lineal ``0.15 * (attempt + 1)``, corte seco en 4xx.
2. En el worker: ``deliver_row`` **nunca lanza**; registra el fallo en la propia fila.
3. En el beat: barridos cada 120 s que re-despachan lo ``pending``.
No hay ``autoretry_for``, ``retry_backoff`` ni ``max_retries`` en las tareas: duplicarían el
mecanismo que ya está en el cliente y en el barrido.
**Cuatro guardas de idempotencia**, en este orden:
1. ``_ya_entregado(source_table, source_id, kind)`` antes de encolar.
2. ``if row.status == STATUS_SENT: return`` al entrar a entregar.
3. El ``crm_document_ref`` que viaja con la subida: EFC devuelve 200 con el que ya existía.
4. El UNIQUE parcial del lado de EFC — la única que garantiza la base.
**Por qué ``find_expediente_gaps`` sigue aquí aunque el CRM sea mono-base.** En Anexo22 el outbox se
commitea aparte del pedimento (dos bases distintas) y ese doble-commit es lo que obligó a inventar el
barrido de huecos. Aquí la fila del outbox va en la MISMA transacción que el expediente, así que esa
ventana no existe. El barrido se conserva porque cubre otras dos cosas: los expedientes creados
**antes** de activar la integración, y cualquier crash. Queda escrito para que el siguiente que lo
lea no lo borre creyendo que es redundante.
"""
import logging
from contextlib import contextmanager
from datetime import datetime, timezone
from typing import Optional
from sqlalchemy.orm import Session
from core.config import settings
from core.database import scoped_core_db
from core.efc_client import EfcClient, EfcClientError, efc_client
from ..cases.models import Case
from .models import (
KIND_COMPLETAR,
KIND_EXPEDIENTE,
MAX_ATTEMPTS,
STATUS_FAILED,
STATUS_PENDING,
STATUS_SENT,
EfcFileOutbox,
EfcSyncOutbox,
)
logger = logging.getLogger(__name__)
# Cache de organización EFC por slug de tenant. Dict módulo-global: vive por worker y NO se
# invalida, igual que el del carril de Anexo22. Es correcto porque la organización de un tenant no
# cambia de id: el resolver de EFC es idempotente y devuelve siempre la misma. Si algún día pudiera
# cambiar, reiniciar el worker la vuelve a resolver.
_org_id_cache: dict[str, str] = {}
@contextmanager
def _savepoint(db: Session):
"""Aísla un encolado dentro de la transacción del usuario con un SAVEPOINT.
**Esto es lo único del encolado que NO se clona del carril de Anexo22, y la razón es de fondo.**
Allá el outbox vive en otra base que el pedimento, así que su ``except`` podía hacer
``db.rollback()`` sin consecuencias: revertía la sesión del outbox y la del pedimento ni se
enteraba.
Aquí el CRM es mono-base y el encolado corre DENTRO de la transacción del usuario. Un
``db.rollback()`` en el ``except`` se llevaría por delante la solicitud y el expediente que el
usuario acaba de crear — exactamente lo contrario de best-effort, y sin un solo error visible
para él. Con el SAVEPOINT, un fallo del encolado deshace **solo** la fila del outbox y la
operación local sigue en pie para que el llamador la commitee.
"""
nested = db.begin_nested()
try:
yield nested
except Exception:
nested.rollback()
raise
# ══ Expediente: encolado y entrega ══════════════════════════════════════════
def replicate_expediente_best_effort(db: Session, expediente: Case) -> None:
"""Encola la réplica del expediente a EFC y dispara la entrega inmediata.
Best-effort en todo: si EFC no está configurado, o si el encolado o el despacho fallan, **no se
propaga el error**. El expediente local ya existe y la operación del usuario no se puede romper
porque un sistema de terceros no conteste. El barrido periódico recoge lo que quede pendiente.
"""
if not settings.EFC_API_URL:
return
row = _enqueue_expediente_outbox(db, expediente)
if row is None:
return
_dispatch_delivery(row.id, row.tenant_id, row.company_id)
def _enqueue_expediente_outbox(db: Session, expediente: Case) -> Optional[EfcSyncOutbox]:
"""Inserta la fila de outbox del expediente. Devuelve ``None`` si falla, sin romper nada.
A diferencia del original, **no commitea**: el CRM es mono-base, así que la fila viaja en la
misma transacción que el expediente. Eso cierra de raíz la ventana del doble-commit que en
Anexo22 obligó a inventar el barrido de huecos.
"""
try:
if _expediente_ya_encolado(db, expediente.id):
return None
# El slug del tenant NO se resuelve aquí: se rellena al ENTREGAR. Resolverlo ahora abriría
# una segunda sesión de base (``scoped_core_db``) dentro de la transacción del usuario, que
# es justo lo que el encolado debe evitar. Es además lo que hace el carril de referencia.
payload = {
"source": "crm",
"crm_company_id": expediente.company_id,
"crm_expediente_id": expediente.id,
"folio": expediente.reference,
"storage_token": expediente.efc_storage_token,
}
row = EfcSyncOutbox(
kind=KIND_EXPEDIENTE,
payload=payload,
expediente_ref=expediente.id,
status=STATUS_PENDING,
tenant_id=expediente.tenant_id,
company_id=expediente.company_id,
)
with _savepoint(db):
db.add(row)
db.flush()
return row
except Exception:
logger.warning(
"expediente_gateway: no se pudo encolar el expediente id=%s en el outbox",
getattr(expediente, "id", None), exc_info=True,
)
return None
def _expediente_ya_encolado(db: Session, expediente_id: int) -> bool:
"""¿Ya hay una fila viva de alta para este expediente? Evita encolar la misma réplica dos veces."""
return (
db.query(EfcSyncOutbox.id)
.filter(
EfcSyncOutbox.expediente_ref == expediente_id,
EfcSyncOutbox.kind == KIND_EXPEDIENTE,
EfcSyncOutbox.status.in_((STATUS_PENDING, STATUS_SENT)),
)
.first()
is not None
)
def enqueue_completar_best_effort(db: Session, expediente: Case, campos: dict) -> None:
"""Encola el completado del provisional en EFC con la data aduanera real."""
if not settings.EFC_API_URL:
return
try:
row = EfcSyncOutbox(
kind=KIND_COMPLETAR,
payload={
"source": "crm",
"crm_company_id": expediente.company_id,
"crm_expediente_id": expediente.id,
"folio": expediente.reference,
"pedimento": campos,
},
expediente_ref=expediente.id,
status=STATUS_PENDING,
tenant_id=expediente.tenant_id,
company_id=expediente.company_id,
)
with _savepoint(db):
db.add(row)
db.flush()
except Exception:
logger.warning(
"expediente_gateway: no se pudo encolar el completado del expediente id=%s",
getattr(expediente, "id", None), exc_info=True,
)
return
_dispatch_delivery(row.id, row.tenant_id, row.company_id)
def _dispatch_delivery(outbox_id: int, tenant_id: int, company_id: int) -> None:
"""Dispara la tarea de entrega propagando el contexto RLS por headers de Celery.
Best-effort: si el broker no responde, el barrido la recoge. Los headers son obligatorios —
``core/celery_app.py`` materializa el contexto de RLS a partir de ellos, y sin ellos la tarea
corre sin tenant y no ve nada.
"""
try:
from .tasks import deliver_outbox_row # import diferido: evita ciclo con celery_app
deliver_outbox_row.apply_async(
args=[outbox_id, tenant_id, company_id],
headers={"rls_tenant_id": str(tenant_id), "rls_company_id": str(company_id)},
)
except Exception:
logger.warning(
"expediente_gateway: no se pudo despachar la entrega outbox_id=%s (lo tomará el sweep)",
outbox_id, exc_info=True,
)
def deliver_row(db: Session, row: EfcSyncOutbox, client: Optional[EfcClient] = None) -> None:
"""Entrega una fila del outbox de expedientes a EFC. Actualiza estado y ``attempts``.
**No lanza nunca**: los fallos se registran en la propia fila para reconciliación. Un fallo no
puede matar al worker ni perder la intención de entregar.
"""
client = client or efc_client
if not client.is_configured:
logger.info("expediente_gateway: EFC no configurado; se deja pendiente row=%s", row.id)
return
if row.status == STATUS_SENT:
return
try:
if row.kind == KIND_EXPEDIENTE:
_deliver_expediente(db, row, client)
elif row.kind == KIND_COMPLETAR:
_deliver_completar(db, row, client)
else:
row.status = STATUS_FAILED
row.last_error = f"kind desconocido: {row.kind}"
db.commit()
except EfcClientError as exc:
_register_failure(db, row, exc, retryable=exc.retryable)
except Exception as exc: # noqa: BLE001 — cualquier fallo se registra, no rompe el worker
_register_failure(db, row, exc, retryable=True)
def _register_failure(db: Session, row: EfcSyncOutbox, exc: Exception, retryable: bool) -> None:
row.attempts = (row.attempts or 0) + 1
row.last_error = str(exc)[:2000]
if (not retryable) or row.attempts >= MAX_ATTEMPTS:
row.status = STATUS_FAILED
db.commit()
logger.warning(
"expediente_gateway: entrega falló row=%s attempts=%s retryable=%s status=%s: %s",
row.id, row.attempts, retryable, row.status, exc,
)
def _deliver_expediente(db: Session, row: EfcSyncOutbox, client: EfcClient) -> None:
payload = dict(row.payload or {})
org_id = _resolve_org_id(client, row.tenant_id)
payload["organizacion"] = {"efc_organizacion_id": org_id}
payload["crm_tenant_slug"] = _tenant_slug(row.tenant_id)[0] or ""
resp = client.ingest_expediente(payload)
efc = (resp or {}).get("efc") or {}
row.status = STATUS_SENT
row.sent_at = datetime.now(timezone.utc)
row.efc_pedimento_id = efc.get("pedimento_id")
_stamp_expediente_link(db, row.expediente_ref, org_id, efc.get("pedimento_id"))
db.commit()
logger.info(
"expediente_gateway: expediente replicado row=%s efc_pedimento_id=%s",
row.id, row.efc_pedimento_id,
)
def _deliver_completar(db: Session, row: EfcSyncOutbox, client: EfcClient) -> None:
payload = dict(row.payload or {})
org_id = _resolve_org_id(client, row.tenant_id)
payload["organizacion"] = {"efc_organizacion_id": org_id}
payload["crm_tenant_slug"] = _tenant_slug(row.tenant_id)[0] or ""
folio = payload.get("folio")
client.completar_expediente(folio, payload)
row.status = STATUS_SENT
row.sent_at = datetime.now(timezone.utc)
db.commit()
logger.info("expediente_gateway: expediente completado en EFC row=%s folio=%s", row.id, folio)
def _stamp_expediente_link(db: Session, expediente_id: Optional[int], org_id: str,
pedimento_id: Optional[str]) -> None:
"""Refleja en la fila del expediente que EFC ya lo tiene, para que la UI lo pinte.
Es un espejo, no un handle: el CRM sigue hablando de este expediente por su ``folio``. Se guarda
porque el proxy de descarga necesita el ``organizacion_id`` para preguntarle a EFC.
"""
if expediente_id is None:
return
expediente = db.query(Case).filter(Case.id == expediente_id).first()
if expediente is None:
return
expediente.efc_organizacion_id = org_id
if pedimento_id:
expediente.efc_pedimento_id = pedimento_id
expediente.efc_link_state = "LINKED"
expediente.efc_error_code = None
expediente.efc_error_detail = None
# ══ Archivos: encolado y entrega ════════════════════════════════════════════
def _ya_entregado(db: Session, source_table: str, source_id: int, kind: str) -> bool:
"""¿Este archivo ya se entregó al expediente? Evita re-encolar lo que ya está allá.
Sin esta guarda, un reintento encolaba otra entrega del mismo archivo — que además **falla al
leer el objeto local, porque la primera entrega ya lo borró** con ``delete_local``. Ruido en el
log y una fila del outbox condenada a ``failed``.
Lleva ``source_table`` además de ``source_id``, a diferencia del original: el CRM tiene dos
tablas de documentos con secuencias independientes, así que el id solo es ambiguo y esta guarda
se dispararía de más, saltándose la entrega de un documento distinto que casualmente comparte
entero.
"""
return (
db.query(EfcFileOutbox.id)
.filter(
EfcFileOutbox.source_table == source_table,
EfcFileOutbox.source_id == source_id,
EfcFileOutbox.kind == kind,
EfcFileOutbox.status == STATUS_SENT,
)
.first()
is not None
)
def enqueue_file_best_effort(
db: Session,
*,
kind: str,
s3_key: str,
file_name: str,
content_type: Optional[str],
efc_tipo: str,
source_table: str,
source_id: int,
crm_document_ref: str,
expediente_ref: int,
tenant_id: int,
company_id: int,
delete_local: bool = True,
) -> Optional[EfcFileOutbox]:
"""Encola un archivo hacia el expediente de EFC. Devuelve la fila, o ``None`` si no se encoló.
**No commitea**: la fila va en la misma transacción que el documento que la origina, de modo que
no puede existir un documento sin su intención de entrega ni al revés.
"""
if not settings.EFC_API_URL:
return None
if _ya_entregado(db, source_table, source_id, kind):
return None
try:
row = EfcFileOutbox(
kind=kind,
s3_key=s3_key,
file_name=file_name,
content_type=content_type,
efc_tipo=efc_tipo,
source_table=source_table,
source_id=source_id,
crm_document_ref=crm_document_ref,
expediente_ref=expediente_ref,
delete_local=delete_local,
status=STATUS_PENDING,
tenant_id=tenant_id,
company_id=company_id,
)
with _savepoint(db):
db.add(row)
db.flush()
return row
except Exception:
logger.warning(
"expediente_gateway: no se pudo encolar el archivo %s (%s:%s)",
s3_key, source_table, source_id, exc_info=True,
)
return None
def _dispatch_file_delivery(outbox_id: int, tenant_id: int, company_id: int) -> None:
try:
from .tasks import deliver_file_outbox_row # import diferido
deliver_file_outbox_row.apply_async(
args=[outbox_id, tenant_id, company_id],
headers={"rls_tenant_id": str(tenant_id), "rls_company_id": str(company_id)},
)
except Exception:
logger.warning(
"expediente_gateway: no se pudo despachar entrega de archivo outbox_id=%s (lo tomará el sweep)",
outbox_id, exc_info=True,
)
def deliver_file_row(db: Session, row: EfcFileOutbox, client: Optional[EfcClient] = None) -> None:
"""Sube el archivo de ``row.s3_key`` al expediente de EFC y, si ``delete_local``, borra la copia.
**Ensure-then-upload**: si EFC contesta 404 ``expediente_no_encontrado``, la creación del
provisional puede venir en camino (el outbox de expedientes y el de archivos son colas
distintas), así que se asegura el expediente y se reintenta el upload **una** vez.
**No lanza nunca**: como ``deliver_row``, registra el fallo en la propia fila.
"""
client = client or efc_client
if not client.is_configured or row.status == STATUS_SENT:
return
try:
org_id = _resolve_org_id(client, row.tenant_id)
expediente = db.query(Case).filter(Case.id == row.expediente_ref).first()
if expediente is None:
raise EfcClientError(
f"expediente {row.expediente_ref} no encontrado para el archivo '{row.kind}'",
retryable=True,
)
from core.storage_s3 import get_object_bytes
content = get_object_bytes(row.s3_key)
ct = row.content_type or "application/octet-stream"
try:
resp = client.upload_documento(
org_id, row.company_id, expediente.id, row.efc_tipo,
row.file_name, content, ct, crm_document_ref=row.crm_document_ref,
)
except EfcClientError as exc:
if exc.status_code == 404 and exc.code == "expediente_no_encontrado":
# La creación del provisional puede venir en camino: se asegura y se reintenta UNA vez.
client.ingest_expediente({
"source": "crm",
"crm_tenant_slug": (_tenant_slug(row.tenant_id)[0] or ""),
"crm_company_id": row.company_id,
"crm_expediente_id": expediente.id,
"folio": expediente.reference,
"storage_token": expediente.efc_storage_token,
"organizacion": {"efc_organizacion_id": org_id},
})
resp = client.upload_documento(
org_id, row.company_id, expediente.id, row.efc_tipo,
row.file_name, content, ct, crm_document_ref=row.crm_document_ref,
)
else:
raise
doc_id = resp.get("id") if isinstance(resp, dict) else None
if row.delete_local:
try:
from core.storage_s3 import delete_object_if_exists
delete_object_if_exists(row.s3_key)
except Exception:
# Ya está en EFC: no poder borrar la copia local no invalida la entrega.
logger.warning(
"expediente_gateway: no se pudo borrar el archivo local %s (ya en EFC)",
row.s3_key, exc_info=True,
)
row.status = STATUS_SENT
row.sent_at = datetime.now(timezone.utc)
row.efc_document_id = doc_id
db.commit()
_marcar_documento_entregado(db, row, doc_id)
logger.info(
"expediente_gateway: archivo entregado row=%s kind=%s efc_document_id=%s",
row.id, row.kind, doc_id,
)
except EfcClientError as exc:
_register_file_failure(db, row, exc, exc.retryable)
except Exception as exc: # noqa: BLE001
_register_file_failure(db, row, exc, True)
def _register_file_failure(db: Session, row: EfcFileOutbox, exc: Exception, retryable: bool) -> None:
row.attempts = (row.attempts or 0) + 1
row.last_error = str(exc)[:2000]
if (not retryable) or row.attempts >= MAX_ATTEMPTS:
row.status = STATUS_FAILED
db.commit()
_marcar_documento_fallido(db, row, exc)
logger.warning(
"expediente_gateway: entrega de archivo falló row=%s attempts=%s status=%s: %s",
row.id, row.attempts, row.status, exc,
)
# El mapa (kind, source_table) -> modelo del documento de origen. Un par que NO esté aquí **no toca
# nada**, en vez de caer por omisión sobre una tabla cualquiera: escribir con el id de otra tabla
# vaciaría las columnas de un documento ajeno que tuviera ese mismo entero — daño en el dato de otro,
# sin un solo error visible.
def _modelo_de_origen(source_table: str):
if source_table == "crm.documents":
from ..documents.models import Document
return Document
if source_table == "ops.shipment_documents":
from api.v1.modules.ops.shipments.models import ShipmentDocument
return ShipmentDocument
return None
def _fila_de_origen(db: Session, row: EfcFileOutbox):
modelo = _modelo_de_origen(row.source_table)
if modelo is None or row.source_id is None:
return None
return (
db.query(modelo)
.filter(
modelo.id == row.source_id,
modelo.tenant_id == row.tenant_id,
modelo.company_id == row.company_id,
)
.first()
)
def _marcar_documento_entregado(db: Session, row: EfcFileOutbox, doc_id) -> None:
"""Cierra la entrega en la fila del documento: el badge de la UI pasa a «En expediente»."""
documento = _fila_de_origen(db, row)
if documento is None:
return
documento.efc_document_id = str(doc_id) if doc_id else None
documento.efc_sync_state = "SYNCED"
documento.efc_synced_at = datetime.now(timezone.utc)
documento.efc_error_code = None
documento.efc_error_detail = None
if row.delete_local:
# El objeto local ya no está: dejar la key apuntaría a algo inexistente y la descarga se
# ramificaría por el camino equivocado.
documento.file_key = None
db.commit()
def _marcar_documento_fallido(db: Session, row: EfcFileOutbox, exc: Exception) -> None:
"""Refleja el fallo en la fila del documento para que la ficha lo muestre sin ir a los logs."""
documento = _fila_de_origen(db, row)
if documento is None:
return
documento.efc_attempts = row.attempts
documento.efc_error_detail = str(exc)[:2000]
documento.efc_error_code = getattr(exc, "code", None)
if row.status == STATUS_FAILED:
documento.efc_sync_state = "FAILED"
db.commit()
# ══ Organización ════════════════════════════════════════════════════════════
def _resolve_org_id(client: EfcClient, tenant_id: int) -> str:
slug, name = _tenant_slug(tenant_id)
if not slug:
raise EfcClientError(
f"tenant {tenant_id} sin slug; no se puede resolver la organización EFC.",
retryable=False,
)
if slug in _org_id_cache:
return _org_id_cache[slug]
resp = client.resolve_organizacion(slug, name)
org_id = resp.get("id") if isinstance(resp, dict) else None
if not org_id:
raise EfcClientError("El resolver de organización de EFC no devolvió id.", retryable=True)
_org_id_cache[slug] = org_id
return org_id
def _tenant_slug(tenant_id: int) -> tuple[Optional[str], Optional[str]]:
from api.v1.modules.core.tenants.models import Tenant
with scoped_core_db(tenant_id=tenant_id) as db:
t = db.query(Tenant).filter(Tenant.id == tenant_id).first()
if t is None:
return None, None
return t.slug, t.name
# ══ Tablero de ops ══════════════════════════════════════════════════════════
def _outbox_to_dict(r: EfcSyncOutbox) -> dict:
return {
"id": r.id,
"tabla": "sync",
"kind": r.kind,
"status": r.status,
"attempts": r.attempts,
"last_error": r.last_error,
"expediente_ref": r.expediente_ref,
"efc_pedimento_id": r.efc_pedimento_id,
"created_at": r.created_at.isoformat() if r.created_at else None,
"sent_at": r.sent_at.isoformat() if r.sent_at else None,
}
def _file_outbox_to_dict(r: EfcFileOutbox) -> dict:
return {
"id": r.id,
"tabla": "file",
"kind": r.kind,
"status": r.status,
"attempts": r.attempts,
"last_error": r.last_error,
"expediente_ref": r.expediente_ref,
"file_name": r.file_name,
"efc_tipo": r.efc_tipo,
"source_table": r.source_table,
"source_id": r.source_id,
"crm_document_ref": r.crm_document_ref,
"efc_document_id": r.efc_document_id,
"created_at": r.created_at.isoformat() if r.created_at else None,
"sent_at": r.sent_at.isoformat() if r.sent_at else None,
}
def list_outbox(db: Session, tenant_id: int, company_id: int, tipo: Optional[str] = None,
status: Optional[str] = None, limit: int = 100) -> list[dict]:
"""Lista filas de los DOS outbox para el tablero de ops. ``tipo`` ∈ ``sync`` | ``file``."""
salida: list[dict] = []
if tipo in (None, "", "sync"):
q = db.query(EfcSyncOutbox).filter(
EfcSyncOutbox.tenant_id == tenant_id, EfcSyncOutbox.company_id == company_id
)
if status:
q = q.filter(EfcSyncOutbox.status == status)
salida += [
_outbox_to_dict(r)
for r in q.order_by(EfcSyncOutbox.created_at.desc()).limit(limit).all()
]
if tipo in (None, "", "file"):
q = db.query(EfcFileOutbox).filter(
EfcFileOutbox.tenant_id == tenant_id, EfcFileOutbox.company_id == company_id
)
if status:
q = q.filter(EfcFileOutbox.status == status)
salida += [
_file_outbox_to_dict(r)
for r in q.order_by(EfcFileOutbox.created_at.desc()).limit(limit).all()
]
salida.sort(key=lambda d: (d.get("created_at") or ""), reverse=True)
return salida[:limit]
def retry_outbox_row(db: Session, outbox_id: int, tenant_id: int, company_id: int,
tipo: str = "file") -> bool:
"""Reintento manual: resetea la fila a ``pending`` (``attempts=0``) y la re-despacha.
Devuelve ``False`` si no existe para ese tenant/company — el llamador lo traduce a **404 con
mensaje específico**, no a un 200 silencioso: es contrato con el frontend, que pinta el botón
según lo que reciba.
"""
modelo = EfcSyncOutbox if tipo == "sync" else EfcFileOutbox
r = (
db.query(modelo)
.filter(modelo.id == outbox_id, modelo.tenant_id == tenant_id, modelo.company_id == company_id)
.first()
)
if r is None:
return False
r.status = STATUS_PENDING
r.attempts = 0
r.last_error = None
db.commit()
if tipo == "sync":
_dispatch_delivery(r.id, r.tenant_id, r.company_id)
else:
_reset_documento_pendiente(db, r)
_dispatch_file_delivery(r.id, r.tenant_id, r.company_id)
return True
def _reset_documento_pendiente(db: Session, row: EfcFileOutbox) -> None:
documento = _fila_de_origen(db, row)
if documento is None:
return
documento.efc_sync_state = "PENDING"
documento.efc_error_code = None
documento.efc_error_detail = None
db.commit()
def outbox_metrics(db: Session, tenant_id: int, company_id: int) -> dict:
"""Conteo de los dos outbox por status (monitoreo). Los conteos suman las dos tablas."""
from sqlalchemy import func
counts = {STATUS_PENDING: 0, STATUS_SENT: 0, STATUS_FAILED: 0}
for modelo in (EfcSyncOutbox, EfcFileOutbox):
rows = (
db.query(modelo.status, func.count())
.filter(modelo.tenant_id == tenant_id, modelo.company_id == company_id)
.group_by(modelo.status)
.all()
)
for estado, n in rows:
counts[estado] = counts.get(estado, 0) + n
return {
"pending": counts.get(STATUS_PENDING, 0),
"sent": counts.get(STATUS_SENT, 0),
"failed": counts.get(STATUS_FAILED, 0),
}
def find_expediente_gaps(db: Session, limit: int = 200) -> list:
"""Expedientes (no borrados) SIN ninguna fila de outbox que los referencie.
Nunca se encolaron: expedientes creados **antes** de activar la integración, o un crash. Se
re-encolan para no perder la réplica.
Los ``failed`` **no son huecos** —existen como fila, son visibles y reintentables desde el
tablero—, así que la fila los excluye por estar presente, no por su estado. Corre sin contexto
de tenant (beat); cada expediente lleva el suyo.
"""
from sqlalchemy import exists
ya_encolado = exists().where(EfcSyncOutbox.expediente_ref == Case.id)
return (
db.query(Case)
.filter(Case.deleted_at.is_(None), ~ya_encolado)
.order_by(Case.id.desc())
.limit(limit)
.all()
)

View File

@@ -0,0 +1,39 @@
"""La llave de almacenamiento del expediente en EFC.
Vive en el carril y no en el módulo del expediente a propósito: el expediente (``crm.cases``) es
del CRM y no sabe nada de EFC; esto es exclusivamente cómo EFC nombra su carpeta.
El generador de folios NO está aquí. Es ``crm/common/folios.py::next_folio``, que ya reserva el
consecutivo mensual por ``(tenant, company, entidad, periodo)`` con bloqueo de fila. El carril lo
consume, no lo reimplementa.
"""
from __future__ import annotations
# Longitud de ``Pedimento.pedimento_app`` en EFC (api/customs/models.py). El token se guarda ahí.
PEDIMENTO_APP_MAX = 25
def storage_token(company_id: int, folio: str) -> str:
"""``CRM-{company_id}-{folio}`` — la llave del pedimento provisional en EFC.
Empieza con letras, así que es imposible que colisione con la llave de un pedimento real, que
es ``^\\d{2}-\\d{2}-\\d{4}-\\d{7}$``. El ``company_id`` va dentro porque el puente con EFC es
tenant → organización 1:1 pero un tenant tiene N companies: sin él, dos companies del mismo
tenant generarían el mismo ``EXP2026-08-001`` y chocarían en el ``unique_together`` de EFC.
PRESUPUESTO DE CARACTERES: ``CRM-`` (4) + company + ``-`` (1) + ``EXP2026-08-001`` (14) = 19 +
los dígitos del company. En los 25 de ``pedimento_app`` caben hasta **6 dígitos** de company,
no 7 como decía la primera versión de este docstring: con 7 salen 26 y el insert del lado de
EFC reventaría. Un consecutivo de 4 dígitos (mes con más de 999 expedientes) gasta uno más.
Se valida en vez de truncar: un token recortado apuntaría a la carpeta de OTRO expediente y
los documentos se mezclarían en silencio, que es peor que fallar aquí.
"""
token = f"CRM-{company_id}-{folio}"
if len(token) > PEDIMENTO_APP_MAX:
raise ValueError(
f"storage_token de {len(token)} caracteres excede los {PEDIMENTO_APP_MAX} de "
f"pedimento_app en EFC: {token!r}. Revisa el largo del company_id o del consecutivo."
)
return token

View File

@@ -0,0 +1,122 @@
"""Tareas Celery del carril CRM Agentes de Carga -> EFC.
- ``deliver_outbox_row`` / ``sweep_outbox``: expedientes (alta del provisional y completado).
- ``deliver_file_outbox_row`` / ``sweep_file_outbox``: archivos.
- ``sweep_expediente_gaps``: reconciliación de expedientes que nunca se encolaron.
**La trampa de RLS, que es lo que más fácil se pasa por alto.** ``core/celery_app.py`` materializa el
contexto desde los headers ``rls_tenant_id`` / ``rls_company_id``. Por tanto:
- Las tareas **por fila** se despachan siempre con esos headers.
- Los **barridos corren sin contexto de tenant**: leen los ids pendientes con una sesión sin scope y
despachan una tarea hija por fila con sus propios headers. Si un barrido abriera una sesión con
scope e iterara, o no vería nada o se saltaría el aislamiento.
**Sin ``autoretry_for``, ``retry_backoff`` ni ``max_retries``**: duplicarían el mecanismo de
reintento que ya está en el cliente (3 intentos con backoff lineal) y en el barrido (cada 120 s
hasta ``MAX_ATTEMPTS``).
"""
import logging
from core.celery_app import celery_app
from core.config import settings
from core.database import scoped_core_db
from . import service
from .models import STATUS_PENDING, EfcFileOutbox, EfcSyncOutbox
logger = logging.getLogger(__name__)
# ── Expedientes ─────────────────────────────────────────────────────────────
@celery_app.task(name="expediente_gateway.deliver_outbox_row")
def deliver_outbox_row(outbox_id: int, tenant_id: int, company_id: int) -> None:
with scoped_core_db(tenant_id, company_id) as db:
row = db.query(EfcSyncOutbox).filter(EfcSyncOutbox.id == outbox_id).first()
if row is None:
logger.warning(
"expediente_gateway: outbox_id=%s no encontrado (tenant=%s)", outbox_id, tenant_id
)
return
service.deliver_row(db, row)
@celery_app.task(name="expediente_gateway.sweep_outbox")
def sweep_outbox(limit: int = 100) -> int:
"""Re-despacha filas pendientes de expediente. Sin contexto de tenant: cada fila lleva el suyo."""
with scoped_core_db() as db:
rows = (
db.query(EfcSyncOutbox.id, EfcSyncOutbox.tenant_id, EfcSyncOutbox.company_id)
.filter(EfcSyncOutbox.status == STATUS_PENDING)
.order_by(EfcSyncOutbox.created_at.asc())
.limit(limit)
.all()
)
for rid, tid, cid in rows:
deliver_outbox_row.apply_async(
args=[rid, tid, cid],
headers={"rls_tenant_id": str(tid), "rls_company_id": str(cid) if cid is not None else ""},
)
if rows:
logger.info("expediente_gateway: sweep (expedientes) re-despachó %s filas pendientes", len(rows))
return len(rows)
# ── Archivos ────────────────────────────────────────────────────────────────
@celery_app.task(name="expediente_gateway.deliver_file_outbox_row")
def deliver_file_outbox_row(outbox_id: int, tenant_id: int, company_id: int) -> None:
with scoped_core_db(tenant_id, company_id) as db:
row = db.query(EfcFileOutbox).filter(EfcFileOutbox.id == outbox_id).first()
if row is None:
logger.warning(
"expediente_gateway: file outbox_id=%s no encontrado (tenant=%s)", outbox_id, tenant_id
)
return
service.deliver_file_row(db, row)
@celery_app.task(name="expediente_gateway.sweep_file_outbox")
def sweep_file_outbox(limit: int = 100) -> int:
"""Re-despacha archivos pendientes (EFC o el broker caídos cuando el usuario subió el archivo)."""
with scoped_core_db() as db:
rows = (
db.query(EfcFileOutbox.id, EfcFileOutbox.tenant_id, EfcFileOutbox.company_id)
.filter(EfcFileOutbox.status == STATUS_PENDING)
.order_by(EfcFileOutbox.created_at.asc())
.limit(limit)
.all()
)
for rid, tid, cid in rows:
deliver_file_outbox_row.apply_async(
args=[rid, tid, cid],
headers={"rls_tenant_id": str(tid), "rls_company_id": str(cid) if cid is not None else ""},
)
if rows:
logger.info("expediente_gateway: sweep (archivos) re-despachó %s archivos pendientes", len(rows))
return len(rows)
# ── Reconciliación de huecos ────────────────────────────────────────────────
@celery_app.task(name="expediente_gateway.sweep_expediente_gaps")
def sweep_expediente_gaps(limit: int = 200) -> int:
"""Detecta expedientes que nunca se encolaron a EFC y los re-encola.
No-op si la integración está apagada.
"""
if not settings.EFC_API_URL:
return 0
n = 0
with scoped_core_db() as db:
gaps = service.find_expediente_gaps(db, limit=limit)
for expediente in gaps:
service.replicate_expediente_best_effort(db, expediente)
n += 1
if n:
db.commit()
if n:
logger.info("expediente_gateway: sweep de huecos re-encoló %s expedientes", n)
return n