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)
448 lines
17 KiB
Python
448 lines
17 KiB
Python
"""Servicio del expediente del CRM.
|
|
|
|
Funciones libres que reciben ``db, tenant_id, company_id``, como el resto de los módulos del repo.
|
|
"""
|
|
|
|
import hashlib
|
|
import uuid
|
|
from datetime import datetime, timezone
|
|
|
|
from fastapi import HTTPException, UploadFile, status
|
|
from sqlalchemy.orm import Session
|
|
|
|
from core.s3_keys import expediente_document_key
|
|
from core.storage_s3 import put_object_bytes
|
|
|
|
from ..documents.models import Document
|
|
from ..expediente_gateway import service as gateway
|
|
from ..expediente_gateway.models import FILE_KIND_DOCUMENTO, SOURCE_CRM_DOCUMENTS
|
|
from ..service_requests.models import ServiceRequest
|
|
from ..uploads.routes import leer_acotado, validar_extension
|
|
from .doc_types import is_valid_doc_type
|
|
from .dto import ExpedienteCompleteInput
|
|
from .folio import next_folio, storage_token
|
|
from .models import Expediente
|
|
|
|
|
|
def _get_service_request(db: Session, service_request_id: int, tenant_id: int, company_id: int) -> ServiceRequest:
|
|
obj = (
|
|
db.query(ServiceRequest)
|
|
.filter(
|
|
ServiceRequest.id == service_request_id,
|
|
ServiceRequest.tenant_id == tenant_id,
|
|
ServiceRequest.company_id == company_id,
|
|
ServiceRequest.deleted_at.is_(None),
|
|
)
|
|
.first()
|
|
)
|
|
if not obj:
|
|
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Solicitud no encontrada")
|
|
return obj
|
|
|
|
|
|
def get_expediente(db: Session, expediente_id: int, tenant_id: int, company_id: int) -> Expediente:
|
|
obj = (
|
|
db.query(Expediente)
|
|
.filter(
|
|
Expediente.id == expediente_id,
|
|
Expediente.tenant_id == tenant_id,
|
|
Expediente.company_id == company_id,
|
|
Expediente.deleted_at.is_(None),
|
|
)
|
|
.first()
|
|
)
|
|
if not obj:
|
|
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Expediente no encontrado")
|
|
return obj
|
|
|
|
|
|
def list_expedientes(
|
|
db: Session,
|
|
tenant_id: int,
|
|
company_id: int,
|
|
service_request_id: int | None = None,
|
|
account_id: int | None = None,
|
|
exp_status: str | None = None,
|
|
) -> list[Expediente]:
|
|
query = db.query(Expediente).filter(
|
|
Expediente.tenant_id == tenant_id,
|
|
Expediente.company_id == company_id,
|
|
Expediente.deleted_at.is_(None),
|
|
)
|
|
if service_request_id is not None:
|
|
query = query.filter(Expediente.service_request_id == service_request_id)
|
|
if account_id is not None:
|
|
query = query.filter(Expediente.account_id == account_id)
|
|
if exp_status:
|
|
query = query.filter(Expediente.status == exp_status)
|
|
return query.order_by(Expediente.created_at.desc()).all()
|
|
|
|
|
|
def find_by_service_request(
|
|
db: Session, service_request_id: int, tenant_id: int, company_id: int
|
|
) -> Expediente | None:
|
|
return (
|
|
db.query(Expediente)
|
|
.filter(
|
|
Expediente.service_request_id == service_request_id,
|
|
Expediente.tenant_id == tenant_id,
|
|
Expediente.company_id == company_id,
|
|
Expediente.deleted_at.is_(None),
|
|
)
|
|
.first()
|
|
)
|
|
|
|
|
|
def ensure_expediente_for_service_request(
|
|
db: Session,
|
|
service_request_id: int,
|
|
tenant_id: int,
|
|
company_id: int,
|
|
user_id: str | None = None,
|
|
account_id: int | None = None,
|
|
) -> Expediente:
|
|
"""Devuelve el expediente de esa solicitud, creándolo si todavía no existe. **Idempotente.**
|
|
|
|
**No hace commit**: hace ``flush`` sobre la sesión que recibe. La razón es que sus dos
|
|
llamadores —``create_service_request`` y ``create_from_opportunity``— lo invocan ANTES de su
|
|
propio commit, dentro de la misma transacción. Si esto commiteara por su cuenta, un fallo
|
|
posterior en el alta de la solicitud dejaría un expediente huérfano con su folio ya quemado.
|
|
Es la misma disciplina de ``next_folio`` y la del asignador de folios de Anexo22.
|
|
"""
|
|
existing = find_by_service_request(db, service_request_id, tenant_id, company_id)
|
|
if existing is not None:
|
|
return existing
|
|
|
|
folio, year, month, sequence = next_folio(db, tenant_id, company_id)
|
|
expediente = Expediente(
|
|
folio=folio,
|
|
period_year=year,
|
|
period_month=month,
|
|
sequence=sequence,
|
|
service_request_id=service_request_id,
|
|
account_id=account_id,
|
|
status="abierto",
|
|
efc_storage_token=storage_token(company_id, folio),
|
|
efc_link_state="PENDING",
|
|
tenant_id=tenant_id,
|
|
company_id=company_id,
|
|
created_by=user_id,
|
|
updated_by=user_id,
|
|
)
|
|
db.add(expediente)
|
|
db.flush()
|
|
# Réplica hacia EFC: best-effort y en la MISMA transacción. Si EFC está apagado
|
|
# (``EFC_API_URL`` vacía) esto es un no-op y el expediente vive igual, solo en el CRM.
|
|
gateway.replicate_expediente_best_effort(db, expediente)
|
|
return expediente
|
|
|
|
|
|
def ensure_expediente(
|
|
db: Session,
|
|
service_request_id: int,
|
|
tenant_id: int,
|
|
company_id: int,
|
|
user_id: str | None = None,
|
|
) -> Expediente:
|
|
"""Variante de cara al usuario: valida la solicitud, asegura el expediente y **sí** commitea.
|
|
|
|
La usa el endpoint ``POST /expedientes/ensure``, donde la transacción empieza y termina aquí.
|
|
"""
|
|
solicitud = _get_service_request(db, service_request_id, tenant_id, company_id)
|
|
expediente = ensure_expediente_for_service_request(
|
|
db, solicitud.id, tenant_id, company_id, user_id, account_id=solicitud.account_id
|
|
)
|
|
db.commit()
|
|
db.refresh(expediente)
|
|
return expediente
|
|
|
|
|
|
def _campos_para_efc(campos: dict) -> dict:
|
|
"""Traduce los campos del expediente al vocabulario del contrato de EFC.
|
|
|
|
Las fechas van en ISO porque el payload del outbox se serializa a JSON, y un ``date`` de Python
|
|
no es serializable: sin esto la fila se encolaría bien y **fallaría al entregar**, que es el peor
|
|
momento para descubrirlo.
|
|
"""
|
|
salida = {}
|
|
for clave, valor in campos.items():
|
|
salida[clave] = valor.isoformat() if hasattr(valor, "isoformat") else valor
|
|
return salida
|
|
|
|
|
|
def complete_expediente(
|
|
db: Session,
|
|
expediente_id: int,
|
|
payload: ExpedienteCompleteInput,
|
|
tenant_id: int,
|
|
company_id: int,
|
|
user_id: str | None = None,
|
|
) -> Expediente:
|
|
"""Registra en el CRM la data aduanera real de un expediente.
|
|
|
|
Solo toca la fila del CRM. Completar el pedimento provisional del lado de EFC es una operación
|
|
aparte, del carril del gateway, porque puede fallar por causas de EFC —un pedimento real que ya
|
|
existe con esa llave— y eso no debe impedir que el CRM guarde lo que el usuario capturó.
|
|
"""
|
|
expediente = get_expediente(db, expediente_id, tenant_id, company_id)
|
|
if expediente.status == "completado":
|
|
raise HTTPException(
|
|
status_code=status.HTTP_409_CONFLICT, detail="El expediente ya está completado"
|
|
)
|
|
|
|
campos = payload.model_dump(exclude_unset=True)
|
|
for field, value in campos.items():
|
|
setattr(expediente, field, value)
|
|
expediente.status = "completado"
|
|
expediente.updated_by = user_id
|
|
|
|
# El completado del provisional en EFC va por el outbox, no inline: puede devolver 409 si allá
|
|
# ya existe ese pedimento real, y eso no debe impedir que el CRM guarde lo que se capturó.
|
|
gateway.enqueue_completar_best_effort(db, expediente, _campos_para_efc(campos))
|
|
|
|
db.commit()
|
|
db.refresh(expediente)
|
|
return expediente
|
|
|
|
|
|
def delete_expediente(db: Session, expediente_id: int, tenant_id: int, company_id: int) -> None:
|
|
"""Baja lógica. No propaga nada a EFC: el expediente electrónico se conserva."""
|
|
expediente = get_expediente(db, expediente_id, tenant_id, company_id)
|
|
expediente.deleted_at = datetime.now(timezone.utc)
|
|
db.commit()
|
|
|
|
|
|
# ══ Documentos del expediente (fase 7) ══════════════════════════════════════
|
|
|
|
async def attach_document(
|
|
db: Session,
|
|
expediente_id: int,
|
|
file: UploadFile,
|
|
doc_type: str,
|
|
tenant_id: int,
|
|
company_id: int,
|
|
name: str | None = None,
|
|
user_id: str | None = None,
|
|
) -> Document:
|
|
"""Subida de UN paso: guarda el archivo, crea el documento y encola su entrega a EFC.
|
|
|
|
El orden importa y no es negociable:
|
|
|
|
1. Validar tipo y extensión **antes** de tocar el almacén, para no dejar un objeto huérfano.
|
|
2. Leer el archivo **acotado**: nunca ``await file.read()`` completo (ver ``leer_acotado``).
|
|
3. Escribir a MinIO.
|
|
4. Crear la fila del documento y la del outbox **en una sola transacción**, de modo que no
|
|
pueda existir un documento sin su intención de entrega ni al revés.
|
|
5. Despachar best-effort y responder 201 **pase lo que pase con EFC**.
|
|
|
|
El paso 5 es lo que hace que un EFC caído no le cueste su trabajo al usuario: la subida tiene
|
|
éxito y el documento queda en «Pendiente de enviar» hasta que el barrido lo entregue.
|
|
"""
|
|
expediente = get_expediente(db, expediente_id, tenant_id, company_id)
|
|
|
|
if not is_valid_doc_type(doc_type):
|
|
raise HTTPException(
|
|
status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
|
|
detail="Ese tipo de documento no está en el catálogo.",
|
|
)
|
|
validar_extension(file.filename)
|
|
contenido = await leer_acotado(file)
|
|
|
|
key = expediente_document_key(
|
|
tenant_id, company_id, expediente.id, uuid.uuid4().hex, file.filename or "archivo"
|
|
)
|
|
content_type = file.content_type or "application/octet-stream"
|
|
put_object_bytes(key, contenido, content_type=content_type)
|
|
|
|
documento = Document(
|
|
doc_type=doc_type,
|
|
name=name or file.filename or "archivo",
|
|
file_key=key,
|
|
content_type=content_type,
|
|
size_bytes=len(contenido),
|
|
expediente_id=expediente.id,
|
|
efc_sync_state="PENDING",
|
|
content_sha256=hashlib.sha256(contenido).hexdigest(),
|
|
tenant_id=tenant_id,
|
|
company_id=company_id,
|
|
uploaded_by=user_id,
|
|
)
|
|
db.add(documento)
|
|
db.flush()
|
|
|
|
documento.efc_document_ref = f"CRMDOC-{company_id}-{documento.id}"
|
|
fila = gateway.enqueue_file_best_effort(
|
|
db,
|
|
kind=FILE_KIND_DOCUMENTO,
|
|
s3_key=key,
|
|
file_name=file.filename or "archivo",
|
|
content_type=content_type,
|
|
efc_tipo=doc_type,
|
|
source_table=SOURCE_CRM_DOCUMENTS,
|
|
source_id=documento.id,
|
|
crm_document_ref=documento.efc_document_ref,
|
|
expediente_ref=expediente.id,
|
|
tenant_id=tenant_id,
|
|
company_id=company_id,
|
|
delete_local=True,
|
|
)
|
|
db.commit()
|
|
db.refresh(documento)
|
|
|
|
if fila is not None:
|
|
gateway._dispatch_file_delivery(fila.id, tenant_id, company_id)
|
|
return documento
|
|
|
|
|
|
def get_expediente_document(
|
|
db: Session, expediente_id: int, document_id: int, tenant_id: int, company_id: int
|
|
) -> Document:
|
|
"""El documento, validando que pertenezca a ESE expediente, tenant y company.
|
|
|
|
La validación de pertenencia va **antes** de tocar EFC. Sin ella, un ``document_id`` que
|
|
coincidiera leería el expediente de otro tenant — y peor: el ``organizacion_id`` con el que el
|
|
proxy pregunta se deriva del expediente, así que el CRM iría a preguntarle a la organización de
|
|
otro cliente.
|
|
"""
|
|
expediente = get_expediente(db, expediente_id, tenant_id, company_id)
|
|
documento = (
|
|
db.query(Document)
|
|
.filter(
|
|
Document.id == document_id,
|
|
Document.expediente_id == expediente.id,
|
|
Document.tenant_id == tenant_id,
|
|
Document.company_id == company_id,
|
|
Document.deleted_at.is_(None),
|
|
)
|
|
.first()
|
|
)
|
|
if not documento:
|
|
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Documento no encontrado")
|
|
return documento
|
|
|
|
|
|
def list_expediente_documents(
|
|
db: Session, expediente_id: int, tenant_id: int, company_id: int
|
|
) -> list[Document]:
|
|
expediente = get_expediente(db, expediente_id, tenant_id, company_id)
|
|
return (
|
|
db.query(Document)
|
|
.filter(
|
|
Document.expediente_id == expediente.id,
|
|
Document.tenant_id == tenant_id,
|
|
Document.company_id == company_id,
|
|
Document.deleted_at.is_(None),
|
|
)
|
|
.order_by(Document.created_at.desc())
|
|
.all()
|
|
)
|
|
|
|
|
|
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**.
|
|
|
|
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():
|
|
try:
|
|
async for chunk in upstream.aiter_bytes():
|
|
yield chunk
|
|
finally:
|
|
await upstream.aclose()
|
|
await client.aclose()
|
|
|
|
return _generador()
|
|
|
|
|
|
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.
|
|
|
|
Devuelve ``(iterador, content_type, filename)``. Dos guardas, y las dos importan:
|
|
|
|
1. La pertenencia se valida **antes** de tocar EFC (ver ``get_expediente_document``).
|
|
2. Se manda el ``organizacion_id`` del expediente, para que la verificación de EFC también
|
|
dispare y no baste con acertar un id de documento.
|
|
"""
|
|
from core.efc_client import EfcClientError, efc_client
|
|
|
|
expediente = get_expediente(db, expediente_id, tenant_id, company_id)
|
|
documento = get_expediente_document(db, expediente_id, document_id, tenant_id, company_id)
|
|
|
|
if not documento.efc_document_id:
|
|
raise HTTPException(
|
|
status_code=status.HTTP_409_CONFLICT,
|
|
detail="El documento todavía no llegó al expediente electrónico.",
|
|
)
|
|
if not efc_client.is_configured:
|
|
raise EfcClientError("Integración con EFC no configurada.", retryable=False)
|
|
|
|
return (
|
|
await _abrir_upstream(
|
|
efc_client.download_url(documento.efc_document_id),
|
|
efc_client.auth_headers,
|
|
{"organizacion_id": str(expediente.efc_organizacion_id or "")},
|
|
efc_client.verify_ssl,
|
|
efc_client.upload_timeout_s,
|
|
),
|
|
documento.content_type or "application/octet-stream",
|
|
documento.name or "documento",
|
|
)
|
|
|
|
|
|
def detach_document(
|
|
db: Session, expediente_id: int, document_id: int, tenant_id: int, company_id: int
|
|
) -> None:
|
|
"""Baja lógica en el CRM. **NO destruye el documento en EFC.**
|
|
|
|
El gateway de Anexo22 nunca llama al DELETE de EFC —verificado: ese método no existe en su
|
|
cliente—, y ``record.Document`` en EFC no tiene vigencia ni purga, así que la política implícita
|
|
del sistema es **conservar**. Un documento que mañana puede ser parte del expediente de un
|
|
pedimento real es riesgo de retención fiscal: desaparece de la vista del CRM y sigue en EFC.
|
|
"""
|
|
documento = get_expediente_document(db, expediente_id, document_id, tenant_id, company_id)
|
|
documento.deleted_at = datetime.now(timezone.utc)
|
|
db.commit()
|