From d671558a5915b8fe1c2f0ddb7919a8708400d208 Mon Sep 17 00:00:00 2001 From: marcos Date: Mon, 10 Aug 2026 07:35:25 -0600 Subject: [PATCH] fix(crm): el error de EFC al descargar un documento llega al usuario, y con su codigo El proxy de descarga del expediente comprobaba el status de EFC DENTRO del generador de streaming. Starlette manda la linea de estado en cuanto construye la StreamingResponse -antes de pedir el primer trozo-, asi que cuando el generador descubria el 404 la respuesta ya habia salido con 200: el usuario recibia un archivo vacio o cortado en vez del error. La traduccion asimetrica que fija el ticket -404 pasa como 404, todo lo demas como 502- estaba escrita pero nunca llegaba a ocurrir. Ahora la conexion se abre y su status se valida ANTES de construir la respuesta, con client.send(peticion, stream=True). El streaming se conserva intacto: lo unico que se adelanta es la cabecera de EFC, que por protocolo llega antes que el cuerpo. El cliente httpx se sigue creando fuera de todo `async with` y cerrandose en el finally, que era la parte que si estaba bien. De paso, un fallo de red contra EFC (timeout, conexion rechazada) daba un 500 opaco y dejaba el cliente httpx sin cerrar; ahora sale 502 y se cierra. Es el mismo camino de codigo que el ticket manda cubrir. Se agregan las pruebas del proxy que el ticket exigia y que no existian: streaming real -el cliente sigue vivo mientras se consumen los trozos y se cierra al terminar-, la pertenencia validada antes de tocar EFC, el organizacion_id del expediente en la peticion, que ninguna URL de almacenamiento se filtre al cuerpo ni a los headers, y la traduccion completa de errores. Todo contra httpx.MockTransport, sin red. stream_document y su ruta pasan a async por el cambio; la prueba del 409 en test_expediente_upload gana su await, sin el cual la corrutina no se ejecutaba y la prueba pasaba sin probar nada. Refs: T2026-08-046 (fase 7) --- .../api/v1/modules/crm/expedientes/routes.py | 8 +- .../api/v1/modules/crm/expedientes/service.py | 63 +++- backend/tests/test_expediente_download.py | 332 ++++++++++++++++++ backend/tests/test_expediente_upload.py | 2 +- 4 files changed, 383 insertions(+), 22 deletions(-) create mode 100644 backend/tests/test_expediente_download.py diff --git a/backend/api/v1/modules/crm/expedientes/routes.py b/backend/api/v1/modules/crm/expedientes/routes.py index 014ae03..803e3f8 100644 --- a/backend/api/v1/modules/crm/expedientes/routes.py +++ b/backend/api/v1/modules/crm/expedientes/routes.py @@ -129,7 +129,7 @@ async def upload_expediente_document( @router.get("/expedientes/{expediente_id}/documentos/{document_id}/archivo") -def download_expediente_document( +async def download_expediente_document( expediente_id: int, document_id: int, company_id: int = Query(..., description="Company ID"), @@ -141,10 +141,14 @@ def download_expediente_document( La traducción de errores es asimétrica **a propósito**: cualquier ``EfcClientError`` sale como 502 —es un fallo de la integración, no del usuario—, salvo un 404 de EFC, que sale como 404 porque significa que ese documento realmente no está. + + Es ``async`` porque el servicio abre la conexión con EFC y **comprueba su status antes** de que + esta función devuelva la ``StreamingResponse``: una vez devuelta, el status ya se envió y + traducir el error sería tarde (ver ``service._abrir_upstream``). """ tenant_id = current_user["tenant_id"] try: - iterador, content_type, filename = service.stream_document( + iterador, content_type, filename = await service.stream_document( db, expediente_id, document_id, tenant_id, company_id ) except EfcClientError as exc: diff --git a/backend/api/v1/modules/crm/expedientes/service.py b/backend/api/v1/modules/crm/expedientes/service.py index e85b374..181ca60 100644 --- a/backend/api/v1/modules/crm/expedientes/service.py +++ b/backend/api/v1/modules/crm/expedientes/service.py @@ -338,39 +338,64 @@ def list_expediente_documents( ) -def _iter_upstream(url: str, headers: dict, params: dict, verify: bool, timeout_s: float): - """Generador async que hace de proxy del archivo de EFC hacia el navegador. +async def _abrir_upstream(url: str, headers: dict, params: dict, verify: bool, timeout_s: float): + """Abre la descarga contra EFC y **valida el status antes de devolver el generador**. - **El ``AsyncClient`` se crea DENTRO del generador y se cierra en ``finally``.** Si se creara en - un ``async with`` de fuera, ese bloque cerraría el cliente antes de que empiece el streaming - —FastAPI consume el generador después de devolver la respuesta— y la descarga moriría a medias. + Que la validación ocurra aquí y no dentro del generador **no es cosmético**: Starlette manda la + línea de estado en cuanto construye la ``StreamingResponse``, o sea **antes** de pedir el primer + trozo. Un 404 de EFC detectado ya dentro del generador llega tarde —la respuesta salió con 200— + y el usuario recibe un archivo cortado en lugar del error. Detectándolo antes, la traducción + 404→404 / resto→502 que fija el ticket llega de verdad al cliente. + + **El ``AsyncClient`` se crea aquí y se cierra en el ``finally`` del generador**, no en un + ``async with`` de fuera: FastAPI consume el generador después de que esta función retorna, así + que un cierre anticipado mataría la descarga a media transferencia. EFC nunca entrega una URL de MinIO, y reescribir el host de una URL ya firmada invalida su SigV4, así que el CRM tiene que hacer de segundo proxy: no hay atajo. """ import httpx + client = httpx.AsyncClient(verify=verify, timeout=timeout_s) + try: + peticion = client.build_request("GET", url, headers=headers, params=params) + upstream = await client.send(peticion, stream=True) + except httpx.HTTPError as exc: + # Un fallo de red es de la integración, no del usuario: 502, nunca un 500 opaco. Y el + # cliente se cierra aquí porque todavía no hay generador que lo haga. + await client.aclose() + raise HTTPException( + status_code=status.HTTP_502_BAD_GATEWAY, + detail="No se pudo obtener el archivo del expediente electrónico.", + ) from exc + except BaseException: + await client.aclose() + raise + + if upstream.status_code >= 400: + codigo = upstream.status_code + await upstream.aread() + await upstream.aclose() + await client.aclose() + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND + if codigo == 404 + else status.HTTP_502_BAD_GATEWAY, + detail="No se pudo obtener el archivo del expediente electrónico.", + ) + async def _generador(): - client = httpx.AsyncClient(verify=verify, timeout=timeout_s) try: - async with client.stream("GET", url, headers=headers, params=params) as upstream: - if upstream.status_code >= 400: - await upstream.aread() - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND - if upstream.status_code == 404 - else status.HTTP_502_BAD_GATEWAY, - detail="No se pudo obtener el archivo del expediente electrónico.", - ) - async for chunk in upstream.aiter_bytes(): - yield chunk + async for chunk in upstream.aiter_bytes(): + yield chunk finally: + await upstream.aclose() await client.aclose() return _generador() -def stream_document( +async def stream_document( db: Session, expediente_id: int, document_id: int, tenant_id: int, company_id: int ) -> tuple: """Prepara la descarga de un documento del expediente. @@ -395,7 +420,7 @@ def stream_document( raise EfcClientError("Integración con EFC no configurada.", retryable=False) return ( - _iter_upstream( + await _abrir_upstream( efc_client.download_url(documento.efc_document_id), efc_client.auth_headers, {"organizacion_id": str(expediente.efc_organizacion_id or "")}, diff --git a/backend/tests/test_expediente_download.py b/backend/tests/test_expediente_download.py new file mode 100644 index 0000000..5513848 --- /dev/null +++ b/backend/tests/test_expediente_download.py @@ -0,0 +1,332 @@ +"""Pruebas del proxy de descarga de documentos del expediente. + +EFC nunca entrega una URL de MinIO —reescribir el host de una URL ya firmada invalida su SigV4— así +que el CRM tiene que hacer de segundo proxy. Lo que se fija aquí: + +* La pertenencia se valida **antes** de tocar EFC, y el ``organizacion_id`` que se manda es el del + expediente, no uno que venga del cliente. +* La descarga es **streaming de verdad**: el cliente httpx sigue vivo mientras se consumen los + trozos y se cierra al terminar, pase lo que pase. +* La respuesta **no** filtra ninguna URL de almacenamiento, ni en el cuerpo ni en los headers. +* La traducción de errores es la asimétrica del ticket: un 404 de EFC sale 404, todo lo demás 502. + +Todo va contra ``httpx.MockTransport``: **ninguna de estas pruebas toca la red.** +""" + +import httpx +import pytest +from fastapi import HTTPException + +from api.v1.modules.crm.expediente_gateway import service as gateway +from api.v1.modules.crm.expedientes import routes as expedientes_routes +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 tests.conftest import COMPANY_ID, TENANT_ID + +OTRO_TENANT = 99 +ORGANIZACION = "11111111-2222-3333-4444-555555555555" +DOC_EFC = "9001" +PARTES = [b"%PDF-1.4 parte-1;", b"parte-2;", b"parte-3"] +ARCHIVO = b"".join(PARTES) + +# Un header de almacenamiento que EFC podría llegar a mandar. Está aquí para que la prueba muerda: +# si alguien "simplificara" el proxy reenviando los headers del upstream tal cual, o devolviendo un +# redirect a la URL firmada, la fuga se vería en la respuesta del CRM. +URL_ALMACEN = "https://minio.efc.interno/bucket/doc.pdf?X-Amz-Signature=abc123" + +USUARIO = {"tenant_id": TENANT_ID, "sub": "user-1"} + + +class _ClienteEspia(httpx.AsyncClient): + """``AsyncClient`` con el transporte simulado inyectado, que además anota si lo cerraron. + + El proxy construye su cliente **dentro** del generador, así que la única forma de alcanzarlo + desde una prueba es sustituir la clase. + """ + + transporte = None + instancias: list["_ClienteEspia"] = [] + + def __init__(self, *args, **kwargs): + kwargs.pop("verify", None) # con transporte simulado no hay TLS que verificar + kwargs["transport"] = _ClienteEspia.transporte + super().__init__(*args, **kwargs) + self.cerrado = False + _ClienteEspia.instancias.append(self) + + async def aclose(self): + self.cerrado = True + await super().aclose() + + +@pytest.fixture() +def entorno(db, monkeypatch): + """Un expediente con un documento que **ya** vive en EFC, y EFC simulado.""" + from core.config import settings + from core.efc_client import efc_client + + 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")) + monkeypatch.setattr(efc_client, "base_url", "https://efc.example.test", raising=False) + monkeypatch.setattr(efc_client, "api_key", "llave-de-prueba", raising=False) + + subidos: dict = {} + monkeypatch.setattr( + expedientes_service, "put_object_bytes", + lambda key, body, content_type="application/octet-stream": subidos.update({key: body}), + ) + + solicitud = sr_service.create_service_request( + db, ServiceRequestCreate(operation_type="importacion"), TENANT_ID, COMPANY_ID, "user-1" + ) + expediente = expedientes_service.find_by_service_request( + db, solicitud.id, TENANT_ID, COMPANY_ID + ) + expediente.efc_organizacion_id = ORGANIZACION + + documento = _documento_en_efc(db, expediente) + + peticiones: list[httpx.Request] = [] + _ClienteEspia.instancias = [] + monkeypatch.setattr(httpx, "AsyncClient", _ClienteEspia) + + def usar(handler): + def _registrar(request): + peticiones.append(request) + return handler(request) + + _ClienteEspia.transporte = httpx.MockTransport(_registrar) + + usar(_efc_entrega_el_archivo) + + return { + "db": db, + "expediente": expediente, + "documento": documento, + "peticiones": peticiones, + "usar": usar, + } + + +def _documento_en_efc(db, expediente): + """Documento ya replicado: el proxy solo entra en juego cuando hay ``efc_document_id``.""" + from api.v1.modules.crm.documents.models import Document + + documento = Document( + tenant_id=TENANT_ID, + company_id=COMPANY_ID, + expediente_id=expediente.id, + name="guia.pdf", + doc_type="MBL", + content_type="application/pdf", + efc_document_id=DOC_EFC, + efc_sync_state="SYNCED", + ) + db.add(documento) + db.commit() + db.refresh(documento) + return documento + + +def _efc_entrega_el_archivo(request): + async def _por_partes(): + for parte in PARTES: + yield parte + + return httpx.Response( + 200, + content=_por_partes(), + headers={ + "Content-Type": "application/pdf", + "X-Storage-Url": URL_ALMACEN, + "Content-Disposition": 'attachment; filename="interno-de-efc.pdf"', + }, + ) + + +async def _descargar(entorno, **kwargs): + return await expedientes_service.stream_document( + entorno["db"], + kwargs.pop("expediente_id", entorno["expediente"].id), + kwargs.pop("document_id", entorno["documento"].id), + kwargs.pop("tenant_id", TENANT_ID), + kwargs.pop("company_id", COMPANY_ID), + ) + + +async def _por_la_ruta(entorno, **kwargs): + return await expedientes_routes.download_expediente_document( + kwargs.pop("expediente_id", entorno["expediente"].id), + kwargs.pop("document_id", entorno["documento"].id), + company_id=kwargs.pop("company_id", COMPANY_ID), + current_user=kwargs.pop("current_user", USUARIO), + db=entorno["db"], + ) + + +# ── Streaming ──────────────────────────────────────────────────────────────── + +@pytest.mark.asyncio +async def test_el_proxy_entrega_los_bytes_tal_cual_llegan_de_efc(entorno): + iterador, content_type, filename = await _descargar(entorno) + + assert b"".join([chunk async for chunk in iterador]) == ARCHIVO + assert content_type == "application/pdf" + assert filename == "guia.pdf" + + +@pytest.mark.asyncio +async def test_el_cliente_sigue_vivo_mientras_se_consumen_los_trozos_y_se_cierra_al_final(entorno): + """**El defecto que el ``finally`` del proxy cubre.** + + Si el cliente se creara en un ``async with`` de fuera del generador, ese bloque lo cerraría en + cuanto la función devuelve —FastAPI pide los trozos *después*— y la descarga moriría a medias + sin error visible. Aquí se comprueba lo contrario: abierto durante, cerrado después. + """ + iterador, _, _ = await _descargar(entorno) + generador = iterador.__aiter__() + + primero = await generador.__anext__() + assert primero + cliente = _ClienteEspia.instancias[-1] + assert cliente.cerrado is False, "el cliente se cerró antes de terminar de leer el archivo" + + async for _ in generador: + pass + assert cliente.cerrado is True + + +@pytest.mark.asyncio +async def test_el_cliente_se_cierra_aunque_efc_conteste_error(entorno): + entorno["usar"](lambda request: httpx.Response(500, text="boom")) + + with pytest.raises(HTTPException): + await _descargar(entorno) + + assert _ClienteEspia.instancias[-1].cerrado is True + + +# ── Las dos guardas de pertenencia ─────────────────────────────────────────── + +@pytest.mark.asyncio +async def test_la_pertenencia_se_valida_ANTES_de_tocar_efc(entorno): + """Sin esto, acertar un ``document_id`` bastaría para leer el expediente de otro cliente — y + peor: el ``organizacion_id`` se deriva del expediente, así que el proxy iría a preguntarle a la + organización ajena.""" + with pytest.raises(HTTPException) as exc: + await _descargar(entorno, tenant_id=OTRO_TENANT) + + assert exc.value.status_code == 404 + assert entorno["peticiones"] == [], "se llamó a EFC antes de validar la pertenencia" + + +@pytest.mark.asyncio +async def test_un_documento_de_otra_company_no_se_descarga(entorno): + with pytest.raises(HTTPException) as exc: + await _descargar(entorno, company_id=777) + + assert exc.value.status_code == 404 + assert entorno["peticiones"] == [] + + +@pytest.mark.asyncio +async def test_el_proxy_manda_el_organizacion_id_del_expediente(entorno): + iterador, _, _ = await _descargar(entorno) + async for _ in iterador: + pass + + peticion = entorno["peticiones"][0] + assert peticion.url.params["organizacion_id"] == ORGANIZACION + assert DOC_EFC in str(peticion.url.path) + assert peticion.headers["X-Api-Key"] == "llave-de-prueba" + + +# ── Nada de URLs de almacenamiento hacia afuera ────────────────────────────── + +@pytest.mark.asyncio +async def test_la_respuesta_no_expone_ninguna_url_de_almacenamiento(entorno): + """El proxy existe justamente para no entregar la URL firmada. Ni redirect, ni header + reenviado, ni URL en el cuerpo.""" + respuesta = await _por_la_ruta(entorno) + + assert respuesta.status_code == 200 + cabeceras = {k.lower(): v for k, v in respuesta.headers.items()} + assert "location" not in cabeceras + assert not any("x-amz" in v.lower() or "minio" in v.lower() for v in cabeceras.values()) + assert "x-storage-url" not in cabeceras + + cuerpo = b"".join([chunk async for chunk in respuesta.body_iterator]) + assert cuerpo == ARCHIVO + assert b"X-Amz-Signature" not in cuerpo + + +@pytest.mark.asyncio +async def test_el_nombre_que_ve_el_usuario_es_el_del_crm_no_el_de_efc(entorno): + respuesta = await _por_la_ruta(entorno) + assert 'filename="guia.pdf"' in respuesta.headers["content-disposition"] + + +# ── Traducción de errores: 404 pasa, el resto es 502 ───────────────────────── + +@pytest.mark.asyncio +async def test_un_404_de_efc_llega_al_usuario_como_404(entorno): + """Y **llega**: la traducción tiene que ocurrir antes de que la respuesta empiece a salir, o el + navegador recibiría un 200 con el cuerpo cortado.""" + entorno["usar"](lambda request: httpx.Response(404, json={"detail": "no está"})) + + with pytest.raises(HTTPException) as exc: + await _por_la_ruta(entorno) + + assert exc.value.status_code == 404 + assert exc.value.detail == "No se pudo obtener el archivo del expediente electrónico." + + +@pytest.mark.asyncio +@pytest.mark.parametrize("codigo", [400, 403, 500, 503]) +async def test_cualquier_otro_error_de_efc_sale_como_502(entorno, codigo): + entorno["usar"](lambda request: httpx.Response(codigo, text="boom")) + + with pytest.raises(HTTPException) as exc: + await _por_la_ruta(entorno) + + assert exc.value.status_code == 502 + + +@pytest.mark.asyncio +async def test_un_timeout_contra_efc_sale_como_502_y_no_como_500(entorno): + """Un fallo de red es de la integración, no del usuario: 502. Dejarlo escapar sería un 500 + opaco y un cliente httpx sin cerrar.""" + def _revienta(request): + raise httpx.ConnectTimeout("EFC no responde") + + entorno["usar"](_revienta) + + with pytest.raises(HTTPException) as exc: + await _por_la_ruta(entorno) + + assert exc.value.status_code == 502 + assert _ClienteEspia.instancias[-1].cerrado is True + + +@pytest.mark.asyncio +async def test_sin_integracion_configurada_la_ruta_responde_502(entorno, monkeypatch): + from core.efc_client import efc_client + + monkeypatch.setattr(efc_client, "base_url", "", raising=False) + + with pytest.raises(HTTPException) as exc: + await _por_la_ruta(entorno) + + assert exc.value.status_code == 502 + assert entorno["peticiones"] == [] + + +@pytest.mark.asyncio +async def test_un_documento_de_otro_tenant_por_la_ruta_da_404(entorno): + with pytest.raises(HTTPException) as exc: + await _por_la_ruta(entorno, current_user={"tenant_id": OTRO_TENANT, "sub": "user-2"}) + + assert exc.value.status_code == 404 diff --git a/backend/tests/test_expediente_upload.py b/backend/tests/test_expediente_upload.py index 40650cf..f9d9b1f 100644 --- a/backend/tests/test_expediente_upload.py +++ b/backend/tests/test_expediente_upload.py @@ -260,7 +260,7 @@ async def test_descargar_un_documento_que_aun_no_llego_a_efc_da_409(entorno): documento = await _adjuntar(entorno) with pytest.raises(HTTPException) as exc: - expedientes_service.stream_document( + await expedientes_service.stream_document( entorno["db"], entorno["expediente"].id, documento.id, TENANT_ID, COMPANY_ID ) assert exc.value.status_code == 409