"""Cliente HTTP hacia EFC (carril de integración CRM Agentes de Carga -> EFC). Clon del cliente del gateway de Anexo22, que es el carril de referencia ya en producción (``anexo22/backend/api/v1/modules/pedimentos/pedimento_gateway/client.py``), con las rutas cambiadas a ``.../integrations/crm/...``. No es una reinterpretación: el molde de reintentos, el corte en 4xx y la forma del error se conservan tal cual. **Síncrono a propósito** (``httpx.Client``): el consumidor es el worker de Celery que drena el outbox, que corre en contexto sync. El único punto async del carril es el proxy de descarga de cara al usuario, y ése no usa este cliente. Todos los endpoints se autentican con el header ``X-Api-Key`` (``settings.EFC_API_KEY`` == ``CRM_INTEGRATION_API_KEY`` del lado de EFC). """ import logging import time from typing import Any, Optional import httpx from core.config import settings logger = logging.getLogger(__name__) # Rutas de EFC (prefijo /api/v1/, ver config/urls.py del backend de EFC). _PATH_ORG_BUSCAR = "/api/v1/organization/integrations/crm/organizaciones/" _PATH_ORG_RESOLVER = "/api/v1/organization/integrations/crm/organizaciones/resolver/" _PATH_EXPEDIENTE = "/api/v1/customs/integrations/crm/expedientes/" _PATH_EXPEDIENTE_COMPLETAR = "/api/v1/customs/integrations/crm/expedientes/{folio}/completar/" _PATH_EXPEDIENTE_DETALLE = "/api/v1/customs/integrations/crm/expedientes/{folio}/" _PATH_DOCS = "/api/v1/record/integrations/crm/documentos/" _PATH_DOCS_LIST = "/api/v1/record/integrations/crm/documentos/list/" _PATH_DOC_DESCARGAR = "/api/v1/record/integrations/crm/documentos/{doc_id}/descargar/" _PATH_DOC_ELIMINAR = "/api/v1/record/integrations/crm/documentos/{doc_id}/eliminar/" _PATH_DOC_REEMPLAZAR = "/api/v1/record/integrations/crm/documentos/{doc_id}/reemplazar/" class EfcClientError(Exception): """Error de comunicación con EFC. ``retryable=True`` marca fallos transitorios (5xx/timeout/red) que el worker debe reintentar. ``status_code``/``code`` exponen la respuesta de EFC para que el worker pueda ramificar (p. ej. 404 ``expediente_no_encontrado`` → ensure-then-upload). El worker decide **por el campo ``retryable``**, nunca parseando el mensaje: un texto de error cambia con cualquier refactor del otro repo y con él se caería la política de reintentos sin que nada se vea roto. """ def __init__(self, message: str, status_code: Optional[int] = None, code: Optional[str] = None, retryable: bool = False): super().__init__(message) self.status_code = status_code self.code = code self.retryable = retryable class EfcClient: def __init__( self, base_url: Optional[str] = None, api_key: Optional[str] = None, timeout_ms: Optional[int] = None, upload_timeout_ms: Optional[int] = None, verify_ssl: Optional[bool] = None, retries: int = 2, transport: Optional[httpx.BaseTransport] = None, ): self.base_url = (base_url if base_url is not None else settings.EFC_API_URL).rstrip("/") self.api_key = api_key if api_key is not None else settings.EFC_API_KEY self.timeout_s = max(0.1, float(timeout_ms or settings.EFC_API_TIMEOUT_MS) / 1000.0) # Timeout aparte para las subidas: los 8 s de los metadatos no alcanzan para un archivo de # 25 MB. Tiene que quedar POR DEBAJO del proxy_read_timeout del nginx de EFC — si el CRM # esperara más, vería un 504 opaco sin saber si el documento entró. Fallando primero de este # lado, el reintento con el mismo efc_document_ref es limpio. self.upload_timeout_s = max( 0.1, float(upload_timeout_ms or settings.EFC_UPLOAD_TIMEOUT_MS) / 1000.0 ) self.verify_ssl = settings.EFC_API_VERIFY_SSL if verify_ssl is None else verify_ssl self.retries = max(0, int(retries)) self.transport = transport # mTLS (scaffolding pre-prod): si hay CA se usa para verificar; si hay par cert/key se # presenta como certificado de cliente. Vacío = TLS normal. self._ca_path = settings.EFC_MTLS_CA_PATH or "" self._cert_path = settings.EFC_MTLS_CERT_PATH or "" self._key_path = settings.EFC_MTLS_KEY_PATH or "" def _client_kwargs(self, timeout_s: Optional[float] = None) -> dict: """``verify``/``cert`` para httpx según config mTLS (o TLS normal si no hay mTLS).""" verify = self._ca_path if self._ca_path else self.verify_ssl kwargs = { "timeout": timeout_s or self.timeout_s, "verify": verify, "transport": self.transport, } if self._cert_path and self._key_path: kwargs["cert"] = (self._cert_path, self._key_path) return kwargs @property def is_configured(self) -> bool: """``False`` = integración deshabilitada (best-effort): sin URL o sin key.""" return bool(self.base_url and self.api_key) # ── HTTP interno ────────────────────────────────────────────────────────── def _request(self, method: str, path: str, *, json: Any = None, params: dict = None, files: dict = None, data: dict = None, stream: bool = False, timeout_s: Optional[float] = None): if not self.is_configured: raise EfcClientError("EFC no configurado (EFC_API_URL/EFC_API_KEY vacíos).", retryable=False) url = f"{self.base_url}{path}" headers = {"X-Api-Key": self.api_key} last_error: Optional[Exception] = None for attempt in range(self.retries + 1): try: client = httpx.Client(**self._client_kwargs(timeout_s)) try: response = client.request(method, url, headers=headers, json=json, params=params, files=files, data=data) except Exception: client.close() raise if 200 <= response.status_code < 300: if stream: # El caller lee response.content y cierra el cliente. return response, client client.close() return response # 5xx: transitorio, reintentar. if response.status_code >= 500 and attempt < self.retries: client.close() time.sleep(0.15 * (attempt + 1)) continue # 4xx u otro: no reintentar. Extraer code/mensaje de EFC. code, message = _parse_error_body(response) client.close() raise EfcClientError( message or f"EFC respondió {response.status_code}", status_code=response.status_code, code=code, retryable=response.status_code >= 500, ) except (httpx.TimeoutException, httpx.NetworkError) as exc: last_error = exc if attempt >= self.retries: break time.sleep(0.15 * (attempt + 1)) except EfcClientError: raise except Exception as exc: # noqa: BLE001 — cualquier fallo inesperado es no-retryable raise EfcClientError(str(exc), retryable=False) from exc raise EfcClientError( f"EFC inaccesible tras {self.retries + 1} intentos: {last_error}", retryable=True, ) # ── Organización ────────────────────────────────────────────────────────── def buscar_organizaciones(self, q: str) -> list: """Búsqueda por texto, para el alta manual desde una pantalla de administración.""" return self._request("GET", _PATH_ORG_BUSCAR, params={"q": q}).json() def resolve_organizacion(self, tenant_slug: str, tenant_name: Optional[str] = None) -> dict: payload = {"tenant_slug": tenant_slug} if tenant_name: payload["tenant_name"] = tenant_name return self._request("POST", _PATH_ORG_RESOLVER, json=payload).json() # ── Expediente (pedimento provisional en EFC) ────────────────────────────── def ingest_expediente(self, payload: dict) -> dict: """Crea el pedimento provisional del expediente. 201 si es nuevo, 200 si ya existía.""" return self._request("POST", _PATH_EXPEDIENTE, json=payload).json() def completar_expediente(self, folio: str, payload: dict) -> dict: """Completa un provisional con la data aduanera real. Ningún archivo se mueve.""" path = _PATH_EXPEDIENTE_COMPLETAR.format(folio=folio) return self._request("POST", path, json=payload).json() def get_expediente(self, folio: str, organizacion_id: str) -> dict: path = _PATH_EXPEDIENTE_DETALLE.format(folio=folio) return self._request("GET", path, params={"organizacion_id": str(organizacion_id)}).json() # ── Documentos ──────────────────────────────────────────────────────────── def upload_documento(self, organizacion_id: str, crm_company_id: int, crm_expediente_id: int, tipo: str, filename: str, content: bytes, content_type: str = "application/octet-stream", crm_document_ref: Optional[str] = None) -> dict: """Sube un documento al expediente. Multipart, nunca base64. ``crm_document_ref`` es **el handle de NUESTRO registro de origen**, y es lo que después permite recuperar el archivo sin guardar de este lado ningún identificador de EFC. EFC lo guarda junto al documento con un UNIQUE parcial por ``(organizacion, ref)``, así que una entrega repetida devuelve el documento que ya existía (200) en vez de crear otro (201). Es la cuarta capa de idempotencia del carril y la única que garantiza la base: la entrega la hace un worker con reintentos, así que un timeout ambiguo —EFC commiteó y contestó tarde— duplicaría el documento sin esto. """ files = {"file": (filename, content, content_type)} data = { "organizacion_id": str(organizacion_id), "crm_company_id": str(int(crm_company_id)), "crm_expediente_id": str(int(crm_expediente_id)), "tipo": tipo, } if crm_document_ref: data["crm_document_ref"] = crm_document_ref return self._request( "POST", _PATH_DOCS, files=files, data=data, timeout_s=self.upload_timeout_s ).json() def list_documentos(self, organizacion_id: str, crm_expediente_id: int, *, tipo: Optional[str] = None, crm_document_ref: Optional[str] = None) -> list: """Documentos del expediente. ``tipo`` y ``crm_document_ref`` son filtros OPCIONALES. Preguntar por ``crm_document_ref`` es lo que permite recuperar ``efc_document_id`` si el cache local se perdió, sin guardar identificadores ajenos como handle. """ params = { "organizacion_id": str(organizacion_id), "crm_expediente_id": str(int(crm_expediente_id)), } if tipo: params["tipo"] = tipo if crm_document_ref: params["crm_document_ref"] = crm_document_ref return self._request("GET", _PATH_DOCS_LIST, params=params).json() def replace_documento(self, organizacion_id: str, doc_id: str, filename: str, content: bytes, content_type: str = "application/octet-stream") -> dict: """Sustituye el CONTENIDO de un documento conservando su fila (mismo id, mismo tipo). EFC sube primero y borra el viejo al final, así que una subida fallida deja el anterior intacto y descargable. """ files = {"file": (filename, content, content_type)} data = {"organizacion_id": str(organizacion_id)} path = _PATH_DOC_REEMPLAZAR.format(doc_id=doc_id) return self._request( "PUT", path, files=files, data=data, timeout_s=self.upload_timeout_s ).json() def download_documento(self, organizacion_id: str, doc_id: str) -> tuple[bytes, str]: path = _PATH_DOC_DESCARGAR.format(doc_id=doc_id) params = {"organizacion_id": str(organizacion_id)} response, client = self._request("GET", path, params=params, stream=True) try: content = response.content filename = _filename_from_response(response, default=str(doc_id)) finally: client.close() return content, filename def download_url(self, doc_id: str) -> str: """URL absoluta del endpoint de descarga de EFC, para el proxy async de la fase 7. El proxy no puede usar ``download_documento``: éste es síncrono y bufferiza el archivo entero. Lo que necesita es la URL y el header, y hace su propio streaming. """ return f"{self.base_url}{_PATH_DOC_DESCARGAR.format(doc_id=doc_id)}" @property def auth_headers(self) -> dict: """El header de autenticación, para el proxy async que no pasa por ``_request``.""" return {"X-Api-Key": self.api_key} def _parse_error_body(response) -> tuple[Optional[str], Optional[str]]: """Extrae ``(code, message)`` del cuerpo de error estructurado de EFC (``{"error": {"code", "message"}}``) sin reventar si no es JSON.""" try: body = response.json() except Exception: return None, None err = body.get("error") if isinstance(body, dict) else None if isinstance(err, dict): return err.get("code"), err.get("message") return None, None def _filename_from_response(response, default: str) -> str: disp = response.headers.get("Content-Disposition", "") if "filename=" in disp: return disp.split("filename=")[-1].strip().strip('"') or default return default # Instancia por defecto (lee settings). Los tests inyectan su propio transport construyendo # EfcClient(transport=httpx.MockTransport(...)). efc_client = EfcClient()