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)
387 lines
15 KiB
Python
387 lines
15 KiB
Python
"""Pruebas de la máquina de reintentos del outbox hacia EFC.
|
|
|
|
Lo que se fija aquí es la capa 2 de las tres del carril: el worker **nunca lanza**, registra el
|
|
fallo en la propia fila, y decide reintentar o rendirse por el campo ``retryable`` —nunca parseando
|
|
el texto del error—.
|
|
"""
|
|
|
|
import pytest
|
|
|
|
from api.v1.modules.crm.expediente_gateway import service as gateway
|
|
from api.v1.modules.crm.expediente_gateway.models import (
|
|
FILE_KIND_DOCUMENTO,
|
|
KIND_EXPEDIENTE,
|
|
MAX_ATTEMPTS,
|
|
SOURCE_CRM_DOCUMENTS,
|
|
SOURCE_OPS_SHIPMENT_DOCUMENTS,
|
|
STATUS_FAILED,
|
|
STATUS_PENDING,
|
|
STATUS_SENT,
|
|
EfcFileOutbox,
|
|
EfcSyncOutbox,
|
|
)
|
|
from api.v1.modules.crm.expedientes import service as expedientes_service
|
|
from api.v1.modules.crm.service_requests import service as sr_service
|
|
from api.v1.modules.crm.service_requests.dto import ServiceRequestCreate
|
|
from core.efc_client import EfcClientError
|
|
from tests.conftest import COMPANY_ID, TENANT_ID
|
|
|
|
OTRO_TENANT = 99
|
|
|
|
|
|
@pytest.fixture()
|
|
def efc_encendido(monkeypatch):
|
|
"""Enciende la integración y evita que el encolado toque el broker o el tenant real."""
|
|
from core.config import settings
|
|
|
|
monkeypatch.setattr(settings, "EFC_API_URL", "https://efc.example.test/", raising=False)
|
|
monkeypatch.setattr(gateway, "_dispatch_delivery", lambda *a, **k: None)
|
|
monkeypatch.setattr(gateway, "_dispatch_file_delivery", lambda *a, **k: None)
|
|
monkeypatch.setattr(gateway, "_tenant_slug", lambda tid: ("temex", "TEMEX"))
|
|
return settings
|
|
|
|
|
|
@pytest.fixture()
|
|
def efc_apagado(monkeypatch):
|
|
from core.config import settings
|
|
|
|
monkeypatch.setattr(settings, "EFC_API_URL", "", raising=False)
|
|
return settings
|
|
|
|
|
|
def _expediente(db):
|
|
solicitud = sr_service.create_service_request(
|
|
db, ServiceRequestCreate(operation_type="importacion"), TENANT_ID, COMPANY_ID, "user-1"
|
|
)
|
|
return expedientes_service.find_by_service_request(db, solicitud.id, TENANT_ID, COMPANY_ID)
|
|
|
|
|
|
def _fila_sync(db, expediente, **kwargs):
|
|
row = EfcSyncOutbox(
|
|
kind=kwargs.pop("kind", KIND_EXPEDIENTE),
|
|
payload=kwargs.pop("payload", {"folio": expediente.folio}),
|
|
expediente_ref=expediente.id,
|
|
status=kwargs.pop("status", STATUS_PENDING),
|
|
tenant_id=kwargs.pop("tenant_id", TENANT_ID),
|
|
company_id=kwargs.pop("company_id", COMPANY_ID),
|
|
**kwargs,
|
|
)
|
|
db.add(row)
|
|
db.commit()
|
|
return row
|
|
|
|
|
|
def _fila_archivo(db, expediente, **kwargs):
|
|
row = EfcFileOutbox(
|
|
kind=kwargs.pop("kind", FILE_KIND_DOCUMENTO),
|
|
s3_key=kwargs.pop("s3_key", "tenants/1/companies/1/expedientes/1/guia.pdf"),
|
|
file_name=kwargs.pop("file_name", "guia.pdf"),
|
|
content_type="application/pdf",
|
|
efc_tipo=kwargs.pop("efc_tipo", "MBL"),
|
|
source_table=kwargs.pop("source_table", SOURCE_CRM_DOCUMENTS),
|
|
source_id=kwargs.pop("source_id", 1),
|
|
crm_document_ref=kwargs.pop("crm_document_ref", "CRMDOC-1-1"),
|
|
expediente_ref=expediente.id,
|
|
status=kwargs.pop("status", STATUS_PENDING),
|
|
tenant_id=kwargs.pop("tenant_id", TENANT_ID),
|
|
company_id=kwargs.pop("company_id", COMPANY_ID),
|
|
**kwargs,
|
|
)
|
|
db.add(row)
|
|
db.commit()
|
|
return row
|
|
|
|
|
|
# ── _register_failure ────────────────────────────────────────────────────────
|
|
|
|
def test_un_fallo_retryable_suma_un_intento_y_deja_la_fila_pendiente(db, efc_encendido):
|
|
expediente = _expediente(db)
|
|
row = _fila_sync(db, expediente)
|
|
|
|
gateway._register_failure(db, row, EfcClientError("EFC no responde", retryable=True), True)
|
|
|
|
assert row.attempts == 1
|
|
assert row.status == STATUS_PENDING
|
|
assert "EFC no responde" in row.last_error
|
|
|
|
|
|
def test_un_fallo_no_retryable_marca_failed_de_inmediato(db, efc_encendido):
|
|
"""Un 400 no mejora insistiendo: reintentarlo ocho veces solo retrasa que alguien lo vea."""
|
|
expediente = _expediente(db)
|
|
row = _fila_sync(db, expediente)
|
|
|
|
gateway._register_failure(db, row, EfcClientError("tipo inválido", retryable=False), False)
|
|
|
|
assert row.attempts == 1
|
|
assert row.status == STATUS_FAILED
|
|
|
|
|
|
def test_al_llegar_a_max_attempts_la_fila_queda_failed(db, efc_encendido):
|
|
expediente = _expediente(db)
|
|
row = _fila_sync(db, expediente, attempts=MAX_ATTEMPTS - 1)
|
|
|
|
gateway._register_failure(db, row, EfcClientError("otra vez", retryable=True), True)
|
|
|
|
assert row.attempts == MAX_ATTEMPTS
|
|
assert row.status == STATUS_FAILED
|
|
|
|
|
|
def test_el_ultimo_error_se_trunca_a_2000_caracteres(db, efc_encendido):
|
|
"""``last_error`` es Text, pero un traceback de 5000 caracteres por fila llena la tabla de ruido."""
|
|
expediente = _expediente(db)
|
|
row = _fila_sync(db, expediente)
|
|
|
|
gateway._register_failure(db, row, Exception("x" * 5000), True)
|
|
|
|
assert len(row.last_error) == 2000
|
|
|
|
|
|
# ── deliver_row: no propaga ──────────────────────────────────────────────────
|
|
|
|
class _ClienteQueRevienta:
|
|
is_configured = True
|
|
|
|
def __init__(self, exc):
|
|
self._exc = exc
|
|
self.llamadas = 0
|
|
|
|
def ingest_expediente(self, payload):
|
|
self.llamadas += 1
|
|
raise self._exc
|
|
|
|
def completar_expediente(self, folio, payload):
|
|
self.llamadas += 1
|
|
raise self._exc
|
|
|
|
|
|
def test_deliver_row_no_propaga_la_excepcion_de_efc(db, efc_encendido, monkeypatch):
|
|
"""Si esto propagara, un EFC caído mataría al worker y se perdería la cola entera."""
|
|
expediente = _expediente(db)
|
|
row = _fila_sync(db, expediente)
|
|
monkeypatch.setattr(gateway, "_resolve_org_id", lambda c, t: "org-1")
|
|
|
|
gateway.deliver_row(db, row, _ClienteQueRevienta(EfcClientError("caído", retryable=True)))
|
|
|
|
assert row.status == STATUS_PENDING
|
|
assert row.attempts == 1
|
|
|
|
|
|
def test_deliver_row_tampoco_propaga_una_excepcion_inesperada(db, efc_encendido, monkeypatch):
|
|
expediente = _expediente(db)
|
|
row = _fila_sync(db, expediente)
|
|
monkeypatch.setattr(gateway, "_resolve_org_id", lambda c, t: "org-1")
|
|
|
|
gateway.deliver_row(db, row, _ClienteQueRevienta(RuntimeError("algo raro")))
|
|
|
|
assert row.attempts == 1
|
|
# Una excepción inesperada se trata como transitoria: no se sabe que sea permanente.
|
|
assert row.status == STATUS_PENDING
|
|
|
|
|
|
def test_una_fila_ya_enviada_no_vuelve_a_llamar_a_efc(db, efc_encendido):
|
|
"""Segunda guarda de idempotencia. Sin ella, un re-despacho duplicaría el expediente en EFC."""
|
|
expediente = _expediente(db)
|
|
row = _fila_sync(db, expediente, status=STATUS_SENT)
|
|
cliente = _ClienteQueRevienta(EfcClientError("no debería llamarse"))
|
|
|
|
gateway.deliver_row(db, row, cliente)
|
|
|
|
assert cliente.llamadas == 0
|
|
assert row.status == STATUS_SENT
|
|
|
|
|
|
# ── _ya_entregado: la ambigüedad de las dos secuencias ───────────────────────
|
|
|
|
def test_no_se_encola_dos_veces_el_mismo_archivo(db, efc_encendido):
|
|
expediente = _expediente(db)
|
|
_fila_archivo(db, expediente, source_id=7, status=STATUS_SENT)
|
|
|
|
assert gateway._ya_entregado(db, SOURCE_CRM_DOCUMENTS, 7, FILE_KIND_DOCUMENTO) is True
|
|
|
|
row = gateway.enqueue_file_best_effort(
|
|
db, kind=FILE_KIND_DOCUMENTO, s3_key="k", file_name="f.pdf", content_type=None,
|
|
efc_tipo="MBL", source_table=SOURCE_CRM_DOCUMENTS, source_id=7,
|
|
crm_document_ref="CRMDOC-1-7", expediente_ref=expediente.id,
|
|
tenant_id=TENANT_ID, company_id=COMPANY_ID,
|
|
)
|
|
assert row is None
|
|
|
|
|
|
def test_el_mismo_id_en_otra_tabla_de_origen_SI_se_encola(db, efc_encendido):
|
|
"""La prueba de la ambigüedad de las dos secuencias.
|
|
|
|
``crm.documents.id = 7`` y ``ops.shipment_documents.id = 7`` son documentos DISTINTOS. Sin
|
|
``source_table`` en la guarda, entregar el primero haría que el segundo se saltara para
|
|
siempre — y nadie vería un error.
|
|
"""
|
|
expediente = _expediente(db)
|
|
_fila_archivo(db, expediente, source_id=7, source_table=SOURCE_CRM_DOCUMENTS, status=STATUS_SENT)
|
|
|
|
assert gateway._ya_entregado(db, SOURCE_OPS_SHIPMENT_DOCUMENTS, 7, FILE_KIND_DOCUMENTO) is False
|
|
|
|
row = gateway.enqueue_file_best_effort(
|
|
db, kind=FILE_KIND_DOCUMENTO, s3_key="k", file_name="f.pdf", content_type=None,
|
|
efc_tipo="MBL", source_table=SOURCE_OPS_SHIPMENT_DOCUMENTS, source_id=7,
|
|
crm_document_ref="SHPDOC-1-7", expediente_ref=expediente.id,
|
|
tenant_id=TENANT_ID, company_id=COMPANY_ID,
|
|
)
|
|
assert row is not None
|
|
assert row.source_table == SOURCE_OPS_SHIPMENT_DOCUMENTS
|
|
|
|
|
|
# ── retry ────────────────────────────────────────────────────────────────────
|
|
|
|
def test_retry_resetea_la_fila_y_la_re_despacha(db, efc_encendido, monkeypatch):
|
|
despachos = []
|
|
monkeypatch.setattr(
|
|
gateway, "_dispatch_file_delivery", lambda oid, t, c: despachos.append((oid, t, c))
|
|
)
|
|
expediente = _expediente(db)
|
|
row = _fila_archivo(db, expediente, status=STATUS_FAILED, attempts=MAX_ATTEMPTS,
|
|
last_error="se acabaron los intentos")
|
|
|
|
ok = gateway.retry_outbox_row(db, row.id, TENANT_ID, COMPANY_ID, "file")
|
|
|
|
assert ok is True
|
|
assert row.status == STATUS_PENDING
|
|
assert row.attempts == 0
|
|
assert row.last_error is None
|
|
assert despachos == [(row.id, TENANT_ID, COMPANY_ID)]
|
|
|
|
|
|
def test_retry_de_otro_tenant_devuelve_false(db, efc_encendido):
|
|
"""Devuelve False y el llamador lo traduce a 404: un 200 le haría creer al frontend que se
|
|
reencoló algo que ni siquiera es suyo."""
|
|
expediente = _expediente(db)
|
|
row = _fila_archivo(db, expediente, status=STATUS_FAILED)
|
|
|
|
assert gateway.retry_outbox_row(db, row.id, OTRO_TENANT, COMPANY_ID, "file") is False
|
|
assert row.status == STATUS_FAILED # intacta
|
|
|
|
|
|
def test_retry_de_una_fila_inexistente_devuelve_false(db, efc_encendido):
|
|
assert gateway.retry_outbox_row(db, 999999, TENANT_ID, COMPANY_ID, "file") is False
|
|
|
|
|
|
# ── métricas y listado ───────────────────────────────────────────────────────
|
|
|
|
def test_las_metricas_cuentan_por_status_sumando_las_dos_tablas(db, efc_encendido):
|
|
expediente = _expediente(db)
|
|
_fila_sync(db, expediente, status=STATUS_SENT)
|
|
_fila_archivo(db, expediente, source_id=1, status=STATUS_PENDING)
|
|
_fila_archivo(db, expediente, source_id=2, status=STATUS_FAILED)
|
|
_fila_archivo(db, expediente, source_id=3, status=STATUS_FAILED)
|
|
|
|
metricas = gateway.outbox_metrics(db, TENANT_ID, COMPANY_ID)
|
|
|
|
# El alta del expediente encoló su propia fila pendiente al crearse la solicitud.
|
|
assert metricas["failed"] == 2
|
|
assert metricas["sent"] == 1
|
|
assert metricas["pending"] >= 1
|
|
|
|
|
|
def test_las_metricas_no_ven_las_filas_de_otro_tenant(db, efc_encendido):
|
|
expediente = _expediente(db)
|
|
_fila_archivo(db, expediente, source_id=5, status=STATUS_FAILED, tenant_id=OTRO_TENANT)
|
|
|
|
assert gateway.outbox_metrics(db, TENANT_ID, COMPANY_ID)["failed"] == 0
|
|
|
|
|
|
def test_el_listado_marca_de_que_tabla_viene_cada_fila(db, efc_encendido):
|
|
expediente = _expediente(db)
|
|
_fila_archivo(db, expediente, source_id=1)
|
|
|
|
filas = gateway.list_outbox(db, TENANT_ID, COMPANY_ID)
|
|
tablas = {f["tabla"] for f in filas}
|
|
assert tablas == {"sync", "file"}
|
|
|
|
solo_archivos = gateway.list_outbox(db, TENANT_ID, COMPANY_ID, tipo="file")
|
|
assert {f["tabla"] for f in solo_archivos} == {"file"}
|
|
|
|
|
|
# ── huecos ───────────────────────────────────────────────────────────────────
|
|
|
|
def test_find_expediente_gaps_encuentra_los_que_no_tienen_fila(db, efc_apagado):
|
|
"""Con EFC apagado no se encola nada, así que todos los expedientes son huecos.
|
|
|
|
Es exactamente el caso que el barrido cubre: lo creado ANTES de activar la integración.
|
|
"""
|
|
_expediente(db)
|
|
_expediente(db)
|
|
|
|
huecos = gateway.find_expediente_gaps(db)
|
|
assert len(huecos) == 2
|
|
|
|
|
|
def test_un_expediente_con_fila_failed_NO_es_un_hueco(db, efc_encendido):
|
|
"""Un ``failed`` existe como fila: es visible en el tablero y reintentable a mano.
|
|
|
|
Tratarlo como hueco lo re-encolaría en cada barrido y escondería el fallo.
|
|
"""
|
|
expediente = _expediente(db)
|
|
fila = db.query(EfcSyncOutbox).filter(EfcSyncOutbox.expediente_ref == expediente.id).first()
|
|
assert fila is not None
|
|
fila.status = STATUS_FAILED
|
|
db.commit()
|
|
|
|
assert gateway.find_expediente_gaps(db) == []
|
|
|
|
|
|
# ── best-effort ──────────────────────────────────────────────────────────────
|
|
|
|
def test_con_efc_apagado_no_se_encola_nada(db, efc_apagado):
|
|
"""``EFC_API_URL`` vacía apaga el carril entero. El CRM sigue funcionando igual."""
|
|
expediente = _expediente(db)
|
|
|
|
assert db.query(EfcSyncOutbox).count() == 0
|
|
|
|
fila = gateway.enqueue_file_best_effort(
|
|
db, kind=FILE_KIND_DOCUMENTO, s3_key="k", file_name="f.pdf", content_type=None,
|
|
efc_tipo="MBL", source_table=SOURCE_CRM_DOCUMENTS, source_id=1,
|
|
crm_document_ref="CRMDOC-1-1", expediente_ref=expediente.id,
|
|
tenant_id=TENANT_ID, company_id=COMPANY_ID,
|
|
)
|
|
assert fila is None
|
|
assert db.query(EfcFileOutbox).count() == 0
|
|
|
|
|
|
def test_con_efc_encendido_crear_una_solicitud_encola_su_expediente(db, efc_encendido):
|
|
expediente = _expediente(db)
|
|
|
|
filas = db.query(EfcSyncOutbox).filter(EfcSyncOutbox.expediente_ref == expediente.id).all()
|
|
assert len(filas) == 1
|
|
assert filas[0].kind == KIND_EXPEDIENTE
|
|
assert filas[0].status == STATUS_PENDING
|
|
assert filas[0].payload["folio"] == expediente.folio
|
|
assert filas[0].payload["storage_token"] == expediente.efc_storage_token
|
|
|
|
|
|
def test_si_el_encolado_revienta_la_operacion_local_no_se_rompe(db, efc_encendido, monkeypatch):
|
|
"""La integración NUNCA puede tumbar el alta de una solicitud del usuario."""
|
|
def _revienta(*a, **k):
|
|
raise RuntimeError("la tabla del outbox no existe")
|
|
|
|
monkeypatch.setattr(gateway, "_expediente_ya_encolado", _revienta)
|
|
|
|
solicitud = sr_service.create_service_request(
|
|
db, ServiceRequestCreate(operation_type="importacion"), TENANT_ID, COMPANY_ID, "user-1"
|
|
)
|
|
|
|
assert solicitud.id is not None
|
|
assert expedientes_service.find_by_service_request(db, solicitud.id, TENANT_ID, COMPANY_ID) is not None
|
|
|
|
|
|
def test_no_se_encola_dos_veces_el_mismo_expediente(db, efc_encendido):
|
|
"""Primera guarda: ``ensure`` es idempotente y no debe generar una segunda réplica."""
|
|
expediente = _expediente(db)
|
|
solicitud_id = expediente.service_request_id
|
|
|
|
expedientes_service.ensure_expediente(db, solicitud_id, TENANT_ID, COMPANY_ID, "user-1")
|
|
expedientes_service.ensure_expediente(db, solicitud_id, TENANT_ID, COMPANY_ID, "user-1")
|
|
|
|
filas = db.query(EfcSyncOutbox).filter(
|
|
EfcSyncOutbox.expediente_ref == expediente.id,
|
|
EfcSyncOutbox.kind == KIND_EXPEDIENTE,
|
|
).all()
|
|
assert len(filas) == 1
|