feature/csv-client-provider

This commit is contained in:
hreyes
2026-02-23 17:20:22 -06:00
parent e8e3774758
commit 4a8bea4966
10 changed files with 807 additions and 7 deletions

View File

@@ -0,0 +1,2 @@
# CSV import flow for Clientes y Proveedores (Client Providers).
# Replicates the same two-phase flow as customs_brokers/imports: upload → scan → waiting_confirmation → commit.

View File

@@ -0,0 +1,161 @@
"""
Rutas de importación CSV para Clientes y Proveedores.
Mismo flujo que a76.imports: upload → scan → status (polling) → commit.
"""
import base64
import json
import logging
import os
from uuid import uuid4
from fastapi import APIRouter, File, HTTPException, Query, UploadFile, Depends
from sqlalchemy.orm import Session
from typing import Dict, Any
from core.celery_app import celery_app
from core.database import get_core_db
from core.security import get_current_user, validate_access_to_resource
from .schemas import ImportJobResponse
from .tasks import (
scan_file,
insert_valid_rows,
CP_IMPORT_FILE_PREFIX,
CP_IMPORT_META_PREFIX,
CP_IMPORT_REDIS_TTL,
)
router = APIRouter()
logger = logging.getLogger(__name__)
def _get_redis():
import redis
url = os.getenv("VALKEY_URL", os.getenv("REDIS_URL", "redis://valkey:6379/0"))
return redis.Redis.from_url(url, decode_responses=False)
@router.post("/upload", response_model=ImportJobResponse)
async def upload_import_file(
file: UploadFile = File(...),
company_id: int = Query(..., description="Company ID"),
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
"""
Fase 1: Subir CSV, guardar en Redis, encolar tarea de escaneo.
"""
try:
tenant_id = validate_access_to_resource(db, company_id, current_user)
except Exception as e:
logger.error(f"CP import: access validation failed: {e}")
raise HTTPException(status_code=403, detail="Invalid company access")
if not file.filename or not file.filename.lower().endswith(".csv"):
raise HTTPException(status_code=400, detail="Solo se permiten archivos .csv")
job_id = str(uuid4())
contents = await file.read()
meta_data = {
"tenant_id": tenant_id,
"company_id": company_id,
"user_id": current_user.get("id"),
"template_id": "client_providers",
}
try:
r = _get_redis()
r.set(
f"{CP_IMPORT_FILE_PREFIX}{job_id}",
base64.b64encode(contents),
ex=CP_IMPORT_REDIS_TTL,
)
r.set(
f"{CP_IMPORT_META_PREFIX}{job_id}",
json.dumps(meta_data).encode("utf-8"),
ex=CP_IMPORT_REDIS_TTL,
)
except Exception as e:
logger.error(f"CP import: Redis store error: {e}")
raise HTTPException(status_code=500, detail="No se pudo encolar el archivo.")
try:
upload_dir = os.path.join(os.getcwd(), "uploads", "temp")
os.makedirs(upload_dir, exist_ok=True)
with open(os.path.join(upload_dir, f"cp_{job_id}.csv"), "wb") as f:
f.write(contents)
with open(os.path.join(upload_dir, f"cp_{job_id}.meta.json"), "w") as f:
json.dump(meta_data, f)
except Exception as e:
logger.warning(f"CP import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
return ImportJobResponse(
job_id=job_id,
status="queued",
message="Archivo subido. Escaneo iniciado.",
)
@router.get("/{job_id}/status")
async def get_import_status(job_id: str):
"""
Polling: estado del escaneo o del commit.
"""
task_result = celery_app.AsyncResult(job_id)
if task_result.state == "PENDING":
return {"status": "processing", "progress": 0}
if task_result.state == "PROGRESS":
info = (task_result.info or {})
return {
"status": "processing",
"progress": info.get("current", 0),
"total": info.get("total", 0),
}
if task_result.state == "SUCCESS":
result = task_result.result
if isinstance(result, dict) and "status" in result:
return result
return {"status": "finished", "result": result}
# A veces Celery tiene el result disponible pero state aún no es SUCCESS; si el result es éxito, devolverlo
result = getattr(task_result, "result", None)
if isinstance(result, dict) and result.get("status") in ("finished", "warning"):
return result
logger.warning("CP import task %s failed: state=%s", job_id, task_result.state)
err_msg = None
tb = getattr(task_result, "traceback", None)
if tb and isinstance(tb, str):
lines = [l.strip() for l in tb.strip().split("\n") if l.strip()]
if lines:
err_msg = lines[-1]
if not err_msg:
try:
exc = task_result.get(propagate=False)
if exc is not None:
err_msg = str(exc)
except Exception:
pass
if not err_msg and result is not None:
if not isinstance(result, dict):
err_msg = str(result)
elif result.get("error") or result.get("message"):
err_msg = result.get("error") or result.get("message")
return {"status": "failed", "error": err_msg or "Task failed"}
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str):
"""
Fase 2: Usuario confirma; se encola la inserción de filas válidas.
"""
task = insert_valid_rows.delay(job_id)
return {
"status": "committing",
"message": "Inserción iniciada.",
"commit_job_id": task.id,
}

View File

@@ -0,0 +1,25 @@
from pydantic import BaseModel
from typing import Optional
class ImportJobResponse(BaseModel):
job_id: str
status: str
message: str
class CommitRequest(BaseModel):
pass # no body needed for single model
class ImportJobStatus(BaseModel):
status: str
job_id: str
total_rows: Optional[int] = 0
error_count: Optional[int] = 0
valid_rows: Optional[int] = 0
error: Optional[str] = None
inserted: Optional[int] = 0
skipped_invalid: Optional[int] = 0
skipped_missing_fk: Optional[int] = 0
skipped_details: Optional[list] = None

View File

@@ -0,0 +1,500 @@
"""
Tareas Celery para importación CSV de Clientes y Proveedores.
Flujo en dos fases: scan_file (validación) → insert_valid_rows (commit).
"""
import os
import base64
import csv
import json
import logging
import re
import unicodedata
from typing import Dict, Any, Optional, List
from core.celery_app import celery_app
from core.database import CoreSessionLocal
from .template_config import row_from_template
from api.v1.modules.a76.clients_and_providers.models import (
ClientProvider,
ClientProviderAddress,
ClientOrProviderEnum,
)
logger = logging.getLogger(__name__)
# Redis keys (prefijo propio para no colisionar con cb_ ni a76.imports)
CP_IMPORT_FILE_PREFIX = "cp_import_file:"
CP_IMPORT_META_PREFIX = "cp_import_meta:"
CP_IMPORT_ERROR_LINES_PREFIX = "cp_import_error_lines:"
CP_IMPORT_REDIS_TTL = 3600 # 1 hour
def _get_redis():
import redis
url = os.getenv("VALKEY_URL", os.getenv("REDIS_URL", "redis://valkey:6379/0"))
return redis.Redis.from_url(url, decode_responses=False)
def _worker_upload_dir() -> str:
return os.path.join(os.getcwd(), "uploads", "temp")
def _ensure_worker_has_file_from_redis(job_id: str) -> Optional[str]:
r = _get_redis()
data = r.get(f"{CP_IMPORT_FILE_PREFIX}{job_id}")
if not data:
return None
try:
raw = base64.b64decode(data)
except Exception as e:
logger.warning(f"CP import: failed to decode file from Redis: {e}")
return None
upload_dir = _worker_upload_dir()
os.makedirs(upload_dir, exist_ok=True)
file_path = os.path.join(upload_dir, f"cp_{job_id}.csv")
with open(file_path, "wb") as f:
f.write(raw)
return file_path
def _ensure_worker_has_meta_from_redis(job_id: str, file_path: str) -> bool:
r = _get_redis()
data = r.get(f"{CP_IMPORT_META_PREFIX}{job_id}")
if not data:
return False
try:
meta = json.loads(data.decode("utf-8"))
except Exception as e:
logger.warning(f"CP import: failed to decode meta from Redis: {e}")
return False
meta_path = file_path.replace(".csv", ".meta.json")
with open(meta_path, "w", encoding="utf-8") as f:
json.dump(meta, f)
return True
def _delete_import_from_redis(job_id: str) -> None:
try:
r = _get_redis()
r.delete(
f"{CP_IMPORT_FILE_PREFIX}{job_id}",
f"{CP_IMPORT_META_PREFIX}{job_id}",
f"{CP_IMPORT_ERROR_LINES_PREFIX}{job_id}",
)
except Exception as e:
logger.warning(f"CP import: failed to delete Redis keys: {e}")
def normalize_header(name: Optional[str]) -> str:
if not name:
return ""
name = unicodedata.normalize("NFKD", str(name)).upper()
name = "".join(ch for ch in name if not unicodedata.combining(ch))
name = re.sub(r"[^A-Z0-9]+", " ", name)
return re.sub(r"\s+", " ", name).strip()
def _parse_client_or_provider(val: Optional[str]) -> Optional[ClientOrProviderEnum]:
"""Mapea valor CSV a ClientOrProviderEnum. Retorna None si no reconocido."""
if not val or not str(val).strip():
return None
v = str(val).strip().lower()
if v in ("client", "cliente", "c"):
return ClientOrProviderEnum.CLIENT
if v in ("provider", "proveedor", "p"):
return ClientOrProviderEnum.PROVIDER
if v in ("both", "ambos", "b", "cliente y proveedor"):
return ClientOrProviderEnum.BOTH
return None
def _parse_active(val: Optional[str]) -> bool:
"""Interpreta ACTIVO: 1/true/si/yes -> True, 0/false/no -> False. Default True."""
if not val or not str(val).strip():
return True
v = str(val).strip().lower()
if v in ("1", "true", "si", "", "yes", "s", "x"):
return True
if v in ("0", "false", "no", "n"):
return False
return True
def _validate_row_client_provider(row: Dict[str, Any], line_num: int) -> Optional[Dict[str, Any]]:
"""Valida una fila para Cliente/Proveedor. Retorna error dict o None."""
rfc = (row.get("RFC") or "").strip()
if not rfc:
return {"line": line_num, "col": "RFC", "msg": "Requerido"}
if len(rfc) > 30:
return {"line": line_num, "col": "RFC", "msg": "Máximo 30 caracteres"}
tipo_raw = (row.get("TIPO") or "").strip()
if tipo_raw and _parse_client_or_provider(tipo_raw) is None:
return {
"line": line_num,
"col": "TIPO",
"msg": "Valor no válido. Use Cliente, Proveedor o Ambos.",
}
name = (row.get("NOMBRE") or "").strip()
if len(name) > 256:
return {"line": line_num, "col": "NOMBRE", "msg": "Máximo 256 caracteres"}
short_name = (row.get("SHORT_NAME") or "").strip()
if short_name and len(short_name) > 10:
return {"line": line_num, "col": "SHORT_NAME", "msg": "Máximo 10 caracteres"}
curp = (row.get("CURP") or "").strip()
if curp and len(curp) > 19:
return {"line": line_num, "col": "CURP", "msg": "Máximo 19 caracteres"}
return None
@celery_app.task(bind=True)
def scan_file(self, job_id: str, config: str = None):
"""
Fase 1: Leer CSV, validar filas, escribir errores en JSONL.
Devuelve waiting_confirmation con total_rows, error_count, valid_rows, errors.
"""
logger.info(f"CP import: starting scan for job {job_id}")
file_path = _ensure_worker_has_file_from_redis(job_id)
if not file_path:
return {"status": "failed", "error": "Archivo no encontrado (expirado o no subido). Sube de nuevo."}
_ensure_worker_has_meta_from_redis(job_id, file_path)
error_dir = os.path.join(os.path.dirname(file_path).replace("temp", "errors"), "")
os.makedirs(error_dir, exist_ok=True)
error_path = os.path.join(error_dir, f"cp_{job_id}.jsonl")
total_rows = 0
try:
with open(file_path, "r", encoding="utf-8-sig") as f:
total_rows = sum(1 for _ in f) - 1
except Exception as e:
return {"status": "failed", "error": str(e)}
meta_path = file_path.replace(".csv", ".meta.json")
meta = {}
if os.path.exists(meta_path):
try:
with open(meta_path, "r", encoding="utf-8") as f:
meta = json.load(f) or {}
except Exception as e:
logger.warning(f"CP import: failed to read meta: {e}")
tenant_id = meta.get("tenant_id")
company_id = meta.get("company_id")
if not tenant_id or not company_id:
return {"status": "failed", "error": "Falta contexto (tenant/company)"}
error_count = 0
processed_rows = 0
errors_detail: List[Dict[str, Any]] = []
try:
with open(file_path, "r", encoding="utf-8-sig") as f_in, open(
error_path, "w", encoding="utf-8"
) as f_err:
sample = f_in.read(2048)
f_in.seek(0)
try:
dialect = csv.Sniffer().sniff(sample, delimiters=",;\t")
except Exception:
dialect = "excel"
reader = csv.DictReader(f_in, dialect=dialect)
for i, row in enumerate(reader, start=1):
if i % 500 == 0:
self.update_state(
state="PROGRESS",
meta={"current": i, "total": total_rows, "errors": error_count},
)
row_norm = row_from_template(row, normalize_header)
err = _validate_row_client_provider(row_norm, i)
if err:
error_count += 1
f_err.write(json.dumps(err) + "\n")
if len(errors_detail) < 500:
errors_detail.append(
{"line": err["line"], "col": err.get("col", ""), "msg": err.get("msg", "")}
)
processed_rows += 1
except Exception as e:
logger.error(f"CP import scan failed: {e}")
return {"status": "failed", "error": str(e)}
error_lines_list = []
try:
if os.path.exists(error_path):
with open(error_path, "r", encoding="utf-8") as f:
for line in f:
try:
err = json.loads(line)
if "line" in err:
error_lines_list.append(err["line"])
except Exception:
pass
if error_lines_list:
r = _get_redis()
r.set(
f"{CP_IMPORT_ERROR_LINES_PREFIX}{job_id}",
json.dumps(error_lines_list).encode("utf-8"),
ex=CP_IMPORT_REDIS_TTL,
)
except Exception as e:
logger.warning(f"CP import: failed to store error lines in Redis: {e}")
return {
"status": "waiting_confirmation",
"job_id": job_id,
"total_rows": processed_rows,
"error_count": error_count,
"valid_rows": processed_rows - error_count,
"errors": errors_detail,
}
def _str_or_none(val: Any, max_len: Optional[int] = None) -> Optional[str]:
if val is None:
return None
s = str(val).strip()
if not s:
return None
if max_len and len(s) > max_len:
return s[:max_len]
return s
@celery_app.task(bind=True)
def insert_valid_rows(self, job_id: str):
"""
Fase 2: Re-leer CSV, omitir filas con error, insertar/actualizar ClientProvider.
Upsert por (tenant_id, company_id, rfc).
"""
logger.info(f"CP import: starting commit for job {job_id}")
file_path = _ensure_worker_has_file_from_redis(job_id)
if not file_path:
alt_path = os.path.join(_worker_upload_dir(), f"cp_{job_id}.csv")
if not os.path.exists(alt_path):
return {
"status": "failed",
"error": "Archivo no encontrado (expirado). Sube y confirma de nuevo.",
}
file_path = alt_path
else:
_ensure_worker_has_meta_from_redis(job_id, file_path)
base_dir = os.path.dirname(file_path)
error_dir = base_dir.replace("temp", "errors")
error_path = os.path.join(error_dir, f"cp_{job_id}.jsonl")
error_lines = set()
try:
r = _get_redis()
raw = r.get(f"{CP_IMPORT_ERROR_LINES_PREFIX}{job_id}")
if raw:
error_lines = set(json.loads(raw.decode("utf-8")))
except Exception as e:
logger.debug(f"CP import: could not load error lines from Redis: {e}")
if not error_lines and os.path.exists(error_path):
with open(error_path, "r", encoding="utf-8") as f:
for line in f:
try:
err = json.loads(line)
error_lines.add(err["line"])
except Exception:
pass
meta_path = file_path.replace(".csv", ".meta.json")
tenant_id = None
company_id = None
meta = {}
if os.path.exists(meta_path):
try:
with open(meta_path, "r", encoding="utf-8") as f:
meta = json.load(f) or {}
tenant_id = meta.get("tenant_id")
company_id = meta.get("company_id")
except Exception:
pass
if not tenant_id or not company_id:
return {"status": "failed", "error": "Falta contexto (tenant/company)"}
inserted_count = 0
skipped_invalid = 0
skipped_details: List[Dict[str, Any]] = []
response = None
try:
with CoreSessionLocal() as session:
# Cargar existentes por (tenant_id, company_id, rfc); rfc puede ser None en BD, usamos '' como key
existing_by_rfc: Dict[str, ClientProvider] = {}
for cp in (
session.query(ClientProvider)
.filter(
ClientProvider.tenant_id == tenant_id,
ClientProvider.company_id == company_id,
)
.all()
):
key = (cp.rfc or "").strip()
existing_by_rfc[key] = cp
with open(file_path, "r", encoding="utf-8-sig") as f:
sample = f.read(2048)
f.seek(0)
try:
dialect = csv.Sniffer().sniff(sample, delimiters=",;\t")
except Exception:
dialect = "excel"
reader = csv.DictReader(f, dialect=dialect)
for i, row in enumerate(reader, start=1):
if i in error_lines:
continue
row_norm = row_from_template(row, normalize_header)
err = _validate_row_client_provider(row_norm, i)
if err:
skipped_invalid += 1
skipped_details.append(
{
"line": i,
"reason": f"{err.get('col', '')}: {err.get('msg', '')}",
}
)
continue
rfc = _str_or_none(row_norm.get("RFC"), 30)
if not rfc:
skipped_invalid += 1
skipped_details.append({"line": i, "reason": "RFC requerido"})
continue
client_or_provider = _parse_client_or_provider(row_norm.get("TIPO"))
if client_or_provider is None:
client_or_provider = ClientOrProviderEnum.BOTH
is_active = _parse_active(row_norm.get("ACTIVO"))
existing = existing_by_rfc.get(rfc)
if existing:
existing.name = _str_or_none(row_norm.get("NOMBRE"), 256)
existing.short_name = _str_or_none(row_norm.get("SHORT_NAME"), 10)
existing.curp = _str_or_none(row_norm.get("CURP"), 19)
existing.client_or_provider = client_or_provider
existing.responsible = _str_or_none(row_norm.get("RESPONSABLE"), 80)
existing.position = _str_or_none(row_norm.get("POSICION"), 30)
existing.incoterm = _str_or_none(row_norm.get("INCOTERM"), 19)
existing.is_active = is_active
session.add(existing)
inserted_count += 1
else:
new_cp = ClientProvider(
tenant_id=tenant_id,
company_id=company_id,
rfc=rfc,
name=_str_or_none(row_norm.get("NOMBRE"), 256),
short_name=_str_or_none(row_norm.get("SHORT_NAME"), 10),
curp=_str_or_none(row_norm.get("CURP"), 19),
client_or_provider=client_or_provider,
responsible=_str_or_none(row_norm.get("RESPONSABLE"), 80),
position=_str_or_none(row_norm.get("POSICION"), 30),
incoterm=_str_or_none(row_norm.get("INCOTERM"), 19),
is_active=is_active,
)
session.add(new_cp)
session.flush()
existing_by_rfc[rfc] = new_cp
inserted_count += 1
# Opcional: crear dirección si hay email/teléfono/dirección
email = _str_or_none(row_norm.get("EMAIL"), 100)
phone = _str_or_none(row_norm.get("TELEFONO"), 30)
address_str = _str_or_none(row_norm.get("DIRECCION"), 100)
if email or phone or address_str:
addr = ClientProviderAddress(
client_id=new_cp.id,
tenant_id=tenant_id,
company_id=company_id,
streets=address_str,
postal_code=_str_or_none(row_norm.get("CODIGO POSTAL"), 15),
city=_str_or_none(row_norm.get("CIUDAD"), 30),
state=_str_or_none(row_norm.get("ESTADO"), 30),
country=_str_or_none(row_norm.get("PAIS"), 3),
phone=phone,
email=email,
contact=_str_or_none(row_norm.get("CONTACTO"), 50),
)
session.add(addr)
try:
session.commit()
except Exception as db_err:
session.rollback()
logger.error(f"CP import DB error: {db_err}")
return {"status": "failed", "error": str(db_err)}
total_skipped = skipped_invalid
if inserted_count == 0 and total_skipped > 0:
response = {
"status": "warning",
"inserted": 0,
"skipped_invalid": skipped_invalid,
"skipped_missing_fk": 0,
"skipped_details": skipped_details,
"message": f"No se insertaron registros. {total_skipped} rechazados.",
}
elif inserted_count == 0:
response = {
"status": "failed",
"error": "No hay registros válidos en el archivo CSV",
"inserted": 0,
"skipped_invalid": skipped_invalid,
"skipped_missing_fk": 0,
"skipped_details": skipped_details,
}
else:
response = {
"status": "finished",
"inserted": inserted_count,
"skipped_invalid": skipped_invalid,
"skipped_missing_fk": 0,
"skipped_details": skipped_details,
}
except Exception as e:
logger.error(f"CP import task failed: {e}")
import traceback
logger.error(traceback.format_exc())
return {"status": "failed", "error": str(e)}
try:
if file_path and os.path.exists(file_path):
os.remove(file_path)
if os.path.exists(error_path):
os.remove(error_path)
meta_path = file_path.replace(".csv", ".meta.json")
if os.path.exists(meta_path):
os.remove(meta_path)
_delete_import_from_redis(job_id)
except Exception as cleanup_err:
logger.warning(f"CP import cleanup failed: {cleanup_err}")
if response is None:
response = {
"status": "failed",
"error": "Error inesperado",
"inserted": 0,
"skipped_invalid": skipped_invalid,
"skipped_missing_fk": 0,
"skipped_details": skipped_details,
}
return response

View File

@@ -0,0 +1,56 @@
"""
Configuración de plantilla CSV para Clientes y Proveedores (EstructuraCatClienteProv.xls).
Solo se leen columnas definidas aquí; el resto se ignora.
Definir cabeceras según la primera fila del XLS oficial (frontend/static/csv/EstructuraCatClienteProv.xls).
"""
from typing import Dict, List, Any, Optional
TEMPLATE_COLUMNS: Dict[str, List[Dict[str, Any]]] = {
"client_providers": [
{"canonical": "NOMBRE", "aliases": ["RAZON SOCIAL", "NAME", "RAZON SOCIAL O NOMBRE"]},
{"canonical": "RFC", "aliases": ["TAX_ID", "TAXID", "IDENTIFICADOR FISCAL", "IDENTIFICACION FISCAL"]},
{"canonical": "TIPO", "aliases": ["CLIENT_OR_PROVIDER", "TIPO ENTIDAD", "CLIENTE O PROVEEDOR"]},
{"canonical": "EMAIL", "aliases": ["CORREO", "E-MAIL", "CORREO ELECTRONICO"]},
{"canonical": "SHORT_NAME", "aliases": ["CLAVE", "CLAVE CORTA", "NOMBRE CORTO", "SIGLAS"]},
{"canonical": "CURP", "aliases": []},
{"canonical": "TELEFONO", "aliases": ["PHONE", "TEL", "TELEFONO CONTACTO"]},
{"canonical": "DIRECCION", "aliases": ["DOMICILIO", "DIRECCION FISCAL", "CALLE"]},
{"canonical": "CODIGO POSTAL", "aliases": ["CODIGOPOSTAL", "CP", "C.P."]},
{"canonical": "CIUDAD", "aliases": ["MUNICIPIO"]},
{"canonical": "ESTADO", "aliases": []},
{"canonical": "PAIS", "aliases": ["COUNTRY"]},
{"canonical": "CONTACTO", "aliases": ["CONTACT", "PERSONA CONTACTO"]},
{"canonical": "RESPONSABLE", "aliases": ["RESPONSABLE AREA"]},
{"canonical": "POSICION", "aliases": ["CARGO", "PUESTO"]},
{"canonical": "INCOTERM", "aliases": []},
{"canonical": "ACTIVO", "aliases": ["IS_ACTIVE", "ACTIVE", "ESTADO ACTIVO"]},
],
}
def build_normalized_lookup(normalize_header_fn) -> Dict[str, str]:
"""normalized_header -> canonical_name para plantilla client_providers."""
cols = TEMPLATE_COLUMNS.get("client_providers")
if not cols:
return {}
lookup: Dict[str, str] = {}
for item in cols:
canonical = item["canonical"]
lookup[normalize_header_fn(canonical)] = canonical
for alias in item.get("aliases") or []:
lookup[normalize_header_fn(alias)] = canonical
return lookup
def row_from_template(row: Dict[str, Any], normalize_header_fn) -> Dict[str, Any]:
"""Fila CSV con solo columnas de la plantilla, en nombres canónicos."""
lookup = build_normalized_lookup(normalize_header_fn)
if not lookup:
return {normalize_header_fn(k): v for k, v in row.items()}
out: Dict[str, Any] = {}
for csv_header, value in row.items():
key_norm = normalize_header_fn(csv_header)
if key_norm in lookup:
out[lookup[key_norm]] = value
return out

View File

@@ -20,10 +20,14 @@ from .dto import (
)
from .service import ClientProviderService
from .models import ClientProvider
from .imports.routes import router as imports_router
# Create main router to add custom endpoints
router = APIRouter(prefix="/clients-providers")
# CSV import (mismo flujo que customs_brokers/imports: upload → scan → commit)
router.include_router(imports_router, prefix="/imports", tags=["clients_and_providers / csv_import"])
@router.get("/", response_model=ClientProviderPaginatedResponseDTO)
async def get_clients_and_providers(

View File

@@ -121,6 +121,11 @@ async def get_import_status(job_id: str):
return result
return {"status": "finished", "result": result}
# A veces Celery tiene el result disponible pero state aún no es SUCCESS; si el result es éxito, devolverlo
result = getattr(task_result, "result", None)
if isinstance(result, dict) and result.get("status") in ("finished", "warning"):
return result
logger.warning("CB import task %s failed: state=%s", job_id, task_result.state)
err_msg = None
tb = getattr(task_result, "traceback", None)
@@ -135,11 +140,10 @@ async def get_import_status(job_id: str):
err_msg = str(exc)
except Exception:
pass
if not err_msg:
result = getattr(task_result, "result", None)
if result is not None and not isinstance(result, dict):
if not err_msg and result is not None:
if not isinstance(result, dict):
err_msg = str(result)
elif isinstance(result, dict) and (result.get("error") or result.get("message")):
elif result.get("error") or result.get("message"):
err_msg = result.get("error") or result.get("message")
return {"status": "failed", "error": err_msg or "Task failed"}

View File

@@ -27,6 +27,7 @@ celery_app = Celery(
"api.v1.modules.a76.reports.exportacion.descargo.task",
"api.v1.modules.a76.imports.tasks",
"api.v1.modules.a76.customs_brokers.imports.tasks",
"api.v1.modules.a76.clients_and_providers.imports.tasks",
"api.v1.modules.a76.reports.exportacion.transmission.MAINX30.task",
"api.v1.modules.a76.reports.importacion.transmission.temporal.MAINX30.task",
"api.v1.modules.a76.reports.importacion.transmission.definitive.MAINX30.task"

View File

@@ -389,6 +389,21 @@ export const api = {
api.post(`/v1/a76/customs-brokers/imports/${jobId}/commit`, {})
},
// CSV import for Clientes y Proveedores (flujo propio en clients_and_providers/imports)
clientProviderImports: {
upload: (file: File, companyId: number) => {
const formData = new FormData();
formData.append('file', file);
return fetchApi(
`/v1/a76/clients-providers/imports/upload?company_id=${companyId}`,
{ method: 'POST', body: formData }
);
},
status: (jobId: string) => api.get(`/v1/a76/clients-providers/imports/${jobId}/status`),
commit: (jobId: string) =>
api.post(`/v1/a76/clients-providers/imports/${jobId}/commit`, {})
},
// Generic request for custom needs (like file uploads)
request: <T = any>(endpoint: string, options: RequestInit = {}) => fetchApi<T>(endpoint, options)
};

View File

@@ -26,6 +26,8 @@
let showResultModal = $state(false);
// Cuando es true, usamos API de importación de Agentes Aduanales (customs_brokers/imports)
let useCustomsBrokerImport = $state(false);
// Cuando es true, usamos API de importación de Clientes y Proveedores (clients_and_providers/imports)
let useClientProviderImport = $state(false);
// Initialize settings for all tabs upfront to avoid reactivity loops
let allSettings = $state<Record<string, any>>(() => {
@@ -45,6 +47,7 @@
activeModelTarget = config.modelTarget || null;
scanResults = null;
useCustomsBrokerImport = config.id === 'customs_brokers';
useClientProviderImport = config.id === 'clients_providers';
const companyId = companyStore.activeCompany?.id || 1;
@@ -66,6 +69,24 @@
return;
}
if (useClientProviderImport) {
try {
const res = await api.clientProviderImports.upload(file, companyId);
if (res.data?.job_id) {
currentJobId = res.data.job_id;
pollStatus();
} else {
toast.error(res.error || 'Error al subir el archivo');
isUploading = false;
}
} catch (e) {
console.error('Upload exception', e);
toast.error('Error inesperado al subir el archivo');
isUploading = false;
}
return;
}
const currentSettings = allSettings[activeTab] || {};
const footerConfig = { ...currentSettings };
if (activeTab === 'importacion') {
@@ -102,6 +123,8 @@
try {
const res = useCustomsBrokerImport
? await api.customsBrokerImports.status(currentJobId)
: useClientProviderImport
? await api.clientProviderImports.status(currentJobId)
: await api.imports.status(currentJobId);
console.log('Poll response', res);
if (res.error && !res.data) {
@@ -116,7 +139,14 @@
toast.success('Escaneo completado. Revisa los resultados.');
isUploading = false;
} else if (res.data?.status === 'failed' || res.data?.status === 'FAILURE') {
toast.error('Error en el procesamiento: ' + (res.data.error || 'Error desconocido'));
const errRaw = res.data.error;
const errText =
typeof errRaw === 'string'
? errRaw.includes('finished') && errRaw.includes('inserted')
? 'La importación pudo completarse. Revisa el listado de registros.'
: errRaw
: (errRaw?.message ?? 'Error desconocido');
toast.error('Error en el procesamiento: ' + errText);
isUploading = false;
currentJobId = null;
scanResults = null;
@@ -239,6 +269,8 @@
isUploading = true;
const res = useCustomsBrokerImport
? await api.customsBrokerImports.commit(currentJobId)
: useClientProviderImport
? await api.clientProviderImports.commit(currentJobId)
: await api.imports.commit(currentJobId, activeModelTarget || '');
if (res.data?.commit_job_id) {
currentJobId = res.data.commit_job_id;