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)
This commit is contained in:
2026-08-10 07:35:25 -06:00
parent d97cdd2f73
commit d671558a59
4 changed files with 383 additions and 22 deletions

View File

@@ -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:

View File

@@ -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 "")},

View File

@@ -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

View File

@@ -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