feature/validaciones-clarion-csv-transportistas

This commit is contained in:
hreyes
2026-03-06 13:13:46 -07:00
parent dfcf4aad6b
commit 3cd88c199b
18 changed files with 1115 additions and 7 deletions

View File

@@ -44,6 +44,9 @@ from api.v1.modules.a76.layouts_csv.drivers.template_config import (
from api.v1.modules.a76.layouts_csv.trailers.template_config import (
TEMPLATE_COLUMNS as TRAILERS_TEMPLATE_COLUMNS,
)
from api.v1.modules.a76.layouts_csv.transportistas.template_config import (
TEMPLATE_COLUMNS as TRANSPORTERS_TEMPLATE_COLUMNS,
)
def _canonicals_from_columns(cols: Optional[List[Dict]]) -> List[str]:
@@ -109,6 +112,9 @@ def _build_registry() -> Dict[str, List[str]]:
# trailers
registry["trailers"] = _canonicals_from_columns(TRAILERS_TEMPLATE_COLUMNS.get("trailers"))
# transporters (transportistas)
registry["transporters"] = _canonicals_from_columns(TRANSPORTERS_TEMPLATE_COLUMNS.get("transporters"))
return registry
@@ -128,6 +134,7 @@ TEMPLATE_FILENAMES: Dict[str, str] = {
"transports": "EstructuraCatTransportes.csv",
"drivers": "EstructuraCatConductor.csv",
"trailers": "EstructuraCatTrailers.csv",
"transporters": "EstructuraCatTransportistas.csv",
"imp_temp_header": "EstructuraEncFacImpoTemp.csv",
"imp_temp_details": "EstructuraParFacImpoTempAF.csv",
"imp_def_header": "EstructuraEncFacImpoDef.csv",

View File

@@ -0,0 +1 @@
# layouts_csv.transportistas: CSV import for Transportistas (carriers catalog). Paridad Clarion.

View File

@@ -0,0 +1 @@
# transportistas common: fk_loader, common_validators, mappers

View File

@@ -0,0 +1,143 @@
"""
Validadores reutilizables para import CSV de transportistas.
Paridad Clarion: VALIDACIONES_TRANSPORTISTAS (estado, país, estado-país), desfase Col S.
"""
from typing import Dict, Any, Optional, Set, Tuple
# Max lengths from Transporter model (a76.transporter)
MAX_LEN = {
"transporter_key": 23,
"name": 256,
"short_name": 10,
"responsible": 100,
"rfc": 30,
"streets": 100,
"postal_code": 15,
"city": 30,
"state": 30,
"country": 3,
"loader_code": 9,
"caat_code": 49,
"transport_code": 8,
"transport_interface_type": 20,
"ftp_server": 200,
"ftp_user": 200,
"ftp_password": 100,
"ftp_directory": 1000,
}
def check_required(row: Dict[str, Any], col: str, max_len: int, line_num: int) -> Optional[Dict[str, Any]]:
val = (row.get(col) or "").strip()
if not val:
return {"line": line_num, "col": col, "msg": "Requerido"}
if max_len and len(val) > max_len:
return {"line": line_num, "col": col, "msg": f"Máximo {max_len} caracteres"}
return None
def check_max_length(
row: Dict[str, Any],
col: str,
max_len: int,
line_num: int,
required: bool = False,
) -> Optional[Dict[str, Any]]:
val = (row.get(col) or "").strip()
if not val:
if required:
return {"line": line_num, "col": col, "msg": "Requerido"}
return None
if len(val) > max_len:
return {"line": line_num, "col": col, "msg": f"Máximo {max_len} caracteres"}
return None
def check_desfase_transportistas(row: Dict[str, Any], line_num: int) -> Optional[Dict[str, Any]]:
"""Si COL_EXTRA (columna 19 / Col S) tiene valor → advertencia de desfase (Clarion, no bloqueante)."""
val = (row.get("COL_EXTRA") or "").strip()
if not val:
return None
return {
"line": line_num,
"col": "COL_EXTRA",
"msg": "Advertencia: Podría existir un desfase en esta línea.",
"solution": "Revisar esta línea del archivo CSV y verificar cada campo esté en la posición correcta.",
"warning": True,
}
def check_pais_catalog_transportistas(
row: Dict[str, Any],
line_num: int,
valid_country_ame: Optional[Set[str]] = None,
) -> Optional[Dict[str, Any]]:
"""Col J (PAIS): Si no vacío, debe ser clave americana (2 chars) y existir en catálogo (GPaises.Pais_Ame)."""
val = (row.get("PAIS") or "").strip()
if not val or valid_country_ame is None:
return None
val_upper = val.upper()
if len(val) > 2:
return {
"line": line_num,
"col": "PAIS",
"msg": f"Error: (Celda J{line_num}) El Pais: {val} Es Incorrecto",
"solution": "Capturar en columna J un Pais en Clave Americana(US = Estados Unidos, MX = Mexico, ES = España, etc).",
}
if val_upper in valid_country_ame:
return None
return {
"line": line_num,
"col": "PAIS",
"msg": f"Error: (Col. J) El Pais: {val} es incorrecto.",
"solution": "Capturar en columna J el Pais del Transportista en Clave Americana.",
}
def check_estado_catalog_transportistas(
row: Dict[str, Any],
line_num: int,
state_descriptions_upper: Optional[Set[str]] = None,
) -> Optional[Dict[str, Any]]:
"""Col I (ESTADO): Si no vacío, debe existir en catálogo (GEstados por descripción). Si hay estado, PAIS no puede estar vacío."""
val = (row.get("ESTADO") or "").strip()
if not val or state_descriptions_upper is None:
return None
val_upper = val.upper()
if val_upper not in state_descriptions_upper:
return {
"line": line_num,
"col": "ESTADO",
"msg": f"Error: (Celda I{line_num}) El Estado: {val} es incorrecto.",
"solution": "Capturar en columna I el Estado del Transportista en Clave Americana o Nombre Completo.",
}
pais = (row.get("PAIS") or "").strip()
if not pais:
return {
"line": line_num,
"col": "PAIS",
"msg": f"Error: (Celda I{line_num}) El Estado: {val} No esta ligado a ningun Pais.",
"solution": "Capturar en columna J un Pais en Clave Americana(US = Estados Unidos, MX = Mexico, ES = España, etc).",
}
return None
def check_estado_pais_consistency_transportistas(
row: Dict[str, Any],
line_num: int,
state_country_set: Optional[Set[Tuple[str, str]]] = None,
) -> Optional[Dict[str, Any]]:
"""Si hay ESTADO y PAIS, validar que el estado pertenezca al país (Clarion GEstados-GPaises)."""
estado = (row.get("ESTADO") or "").strip()
pais = (row.get("PAIS") or "").strip().upper()
if not estado or not pais or state_country_set is None:
return None
key = (pais, estado.upper())
if key in state_country_set:
return None
return {
"line": line_num,
"col": "ESTADO",
"msg": f"Error: (Col. I) EL Estado: {estado} no pertenece al Pais: {pais}.",
"solution": "Capturar en columna I un Estado que pertenesca al Pais de la columna J.",
}

View File

@@ -0,0 +1,79 @@
"""
Carga de conjuntos FK para validación de import CSV de transportistas.
Clarion: GTransportista (ClaveTrans), GPaises (Pais_Ame), GEstados (Descripcion), relación Estado-País.
"""
from typing import Set, Tuple, Optional
import logging
from core.database import CoreSessionLocal
logger = logging.getLogger(__name__)
def load_transportistas_fk_sets(
tenant_id: Optional[int] = None,
company_id: Optional[int] = None,
) -> Tuple[
Set[str],
Set[str],
Set[str],
Set[Tuple[str, str]],
]:
"""
Carga conjuntos para validación CSV de transportistas (paridad Clarion).
Devuelve:
- existing_transporter_keys: claves de transportistas existentes (tenant/company) en mayúsculas
- valid_country_ame: claves americana de países (GPaises.Pais_Ame), mayúsculas
- state_descriptions_upper: descripciones de estados en mayúsculas (GEstados)
- state_country_set: set de (ame_key_pais, description_estado_upper) para validar estado pertenece a país
"""
existing_transporter_keys: Set[str] = set()
valid_country_ame: Set[str] = set()
state_descriptions_upper: Set[str] = set()
state_country_set: Set[Tuple[str, str]] = set()
try:
with CoreSessionLocal() as session:
from api.v1.modules.a76.transportation.transporters.models import Transporter
from api.v1.modules.public.reference_data.countries.models import Country
from api.v1.modules.public.reference_data.states.models import State
if tenant_id is not None and company_id is not None:
for row in (
session.query(Transporter.transporter_key)
.filter(
Transporter.tenant_id == tenant_id,
Transporter.company_id == company_id,
)
.all()
):
if row[0] and (row[0] or "").strip():
existing_transporter_keys.add((row[0] or "").strip().upper())
for row in session.query(Country.ame_key).all():
if row[0]:
valid_country_ame.add((row[0] or "").strip().upper())
for state in session.query(State).all():
desc = (state.description or "").strip()
if desc:
state_descriptions_upper.add(desc.upper())
country = (
session.query(Country)
.filter(Country.m3_key == state.m3_key)
.first()
)
if country and (country.ame_key or "").strip():
state_country_set.add(
((country.ame_key or "").strip().upper(), desc.upper())
)
except Exception as e:
logger.warning("Transportistas import: could not load FK sets: %s", e)
return (
existing_transporter_keys,
valid_country_ame,
state_descriptions_upper,
state_country_set,
)

View File

@@ -0,0 +1,95 @@
"""
Mapeo fila CSV → datos para Transporter.
Paridad Clarion LLENA_TRANSPORTISTAS / VALIDA_PARCIAL: en actualización, campo vacío en CSV usa valor existente.
UPPER aplicado según Clarion: NOMBRE CORTO, CODIGO POSTAL, PAIS, CODIGO CARGADOR, CODIGO CAAT, CODIGO TRANS, etc.
"""
from typing import Dict, Any, Optional, Union
from .common_validators import MAX_LEN
def _str_or_none(val: Any, max_len: Optional[int] = None, upper: bool = False) -> Optional[str]:
if val is None:
return None
s = str(val).strip()
if not s:
return None
if upper:
s = s.upper()
if max_len and len(s) > max_len:
return s[:max_len]
return s
def _get_existing_val(existing: Union[Any, Dict[str, Any]], key: str) -> Any:
if existing is None:
return None
if isinstance(existing, dict):
return existing.get(key)
return getattr(existing, key, None)
def row_to_transporter_data(row_norm: Dict[str, Any]) -> Dict[str, Any]:
"""Build dict suitable for TransporterCreateDTO from normalized row (Clarion QueCSV → GenTra)."""
transporter_key = _str_or_none(row_norm.get("CLAVE TRANSPORTISTA"), MAX_LEN["transporter_key"])
if not transporter_key:
return {}
return {
"transporter_key": transporter_key,
"name": _str_or_none(row_norm.get("NOMBRE"), MAX_LEN["name"]),
"short_name": _str_or_none(row_norm.get("NOMBRE CORTO"), MAX_LEN["short_name"], upper=True),
"responsible": _str_or_none(row_norm.get("RESPONSABLE"), MAX_LEN["responsible"]),
"rfc": _str_or_none(row_norm.get("RFC"), MAX_LEN["rfc"]),
"streets": _str_or_none(row_norm.get("CALLES"), MAX_LEN["streets"]),
"postal_code": _str_or_none(row_norm.get("CODIGO POSTAL"), MAX_LEN["postal_code"], upper=True),
"city": _str_or_none(row_norm.get("CIUDAD"), MAX_LEN["city"]),
"state": _str_or_none(row_norm.get("ESTADO"), MAX_LEN["state"]),
"country": _str_or_none(row_norm.get("PAIS"), MAX_LEN["country"], upper=True),
"loader_code": _str_or_none(row_norm.get("CODIGO CARGADOR"), MAX_LEN["loader_code"], upper=True),
"caat_code": _str_or_none(row_norm.get("CODIGO CAAT"), MAX_LEN["caat_code"], upper=True),
"transport_code": _str_or_none(row_norm.get("CODIGO TRANS"), MAX_LEN["transport_code"], upper=True),
"transport_interface_type": _str_or_none(row_norm.get("TIPO INTERFASE TRANS"), MAX_LEN["transport_interface_type"]),
"ftp_server": _str_or_none(row_norm.get("SERVIDOR FTP"), MAX_LEN["ftp_server"]),
"ftp_user": _str_or_none(row_norm.get("USUARIO FTP"), MAX_LEN["ftp_user"]),
"ftp_password": _str_or_none(row_norm.get("CLAVE ACCESO FTP"), MAX_LEN["ftp_password"]),
"ftp_directory": _str_or_none(row_norm.get("DIRECTORIO FTP"), MAX_LEN["ftp_directory"]),
}
def row_to_transporter_data_for_update(
row_norm: Dict[str, Any],
existing_transporter: Union[Any, Dict[str, Any]],
) -> Dict[str, Any]:
"""
Build dict for TransporterUpdateDTO: CSV value if non-empty, else existing (Clarion VALIDA_PARCIAL).
"""
transporter_key = _str_or_none(row_norm.get("CLAVE TRANSPORTISTA"), MAX_LEN["transporter_key"]) or _get_existing_val(existing_transporter, "transporter_key")
if not transporter_key:
return {}
def _csv_or_existing(csv_key: str, dto_key: str, max_len: Optional[int] = None, upper: bool = False):
v = _str_or_none(row_norm.get(csv_key), max_len, upper=upper)
if v is not None and v != "":
return v
return _get_existing_val(existing_transporter, dto_key)
return {
"transporter_key": transporter_key,
"name": _csv_or_existing("NOMBRE", "name", MAX_LEN["name"]),
"short_name": _csv_or_existing("NOMBRE CORTO", "short_name", MAX_LEN["short_name"], upper=True),
"responsible": _csv_or_existing("RESPONSABLE", "responsible", MAX_LEN["responsible"]),
"rfc": _csv_or_existing("RFC", "rfc", MAX_LEN["rfc"]),
"streets": _csv_or_existing("CALLES", "streets", MAX_LEN["streets"]),
"postal_code": _csv_or_existing("CODIGO POSTAL", "postal_code", MAX_LEN["postal_code"], upper=True),
"city": _csv_or_existing("CIUDAD", "city", MAX_LEN["city"]),
"state": _csv_or_existing("ESTADO", "state", MAX_LEN["state"]),
"country": _csv_or_existing("PAIS", "country", MAX_LEN["country"], upper=True),
"loader_code": _csv_or_existing("CODIGO CARGADOR", "loader_code", MAX_LEN["loader_code"], upper=True),
"caat_code": _csv_or_existing("CODIGO CAAT", "caat_code", MAX_LEN["caat_code"], upper=True),
"transport_code": _csv_or_existing("CODIGO TRANS", "transport_code", MAX_LEN["transport_code"], upper=True),
"transport_interface_type": _csv_or_existing("TIPO INTERFASE TRANS", "transport_interface_type", MAX_LEN["transport_interface_type"]),
"ftp_server": _csv_or_existing("SERVIDOR FTP", "ftp_server", MAX_LEN["ftp_server"]),
"ftp_user": _csv_or_existing("USUARIO FTP", "ftp_user", MAX_LEN["ftp_user"]),
"ftp_password": _csv_or_existing("CLAVE ACCESO FTP", "ftp_password", MAX_LEN["ftp_password"]),
"ftp_directory": _csv_or_existing("DIRECTORIO FTP", "ftp_directory", MAX_LEN["ftp_directory"]),
}

View File

@@ -0,0 +1,190 @@
"""
Rutas de importación CSV para Transportistas.
Flujo: upload → scan → status (polling) → commit.
"""
import base64
import json
import logging
import os
import threading
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.paths import layout_path
from core.security import get_current_user, validate_access_to_resource
from .schemas import ImportJobResponse
from .tasks import (
scan_file,
run_scan_sync,
run_commit_sync,
TRP_IMPORT_FILE_PREFIX,
TRP_IMPORT_META_PREFIX,
TRP_IMPORT_STATUS_PREFIX,
TRP_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"),
actualizar: bool = Query(False, description="Modo Agregar/Actualizar (ACT); si True, validación parcial para claves existentes"),
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
try:
tenant_id = validate_access_to_resource(db, company_id, current_user)
except Exception as e:
logger.error("Transportistas import: access validation failed: %s", 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": "transporters",
"actualizar": actualizar,
}
try:
r = _get_redis()
r.set(
f"{TRP_IMPORT_FILE_PREFIX}{job_id}",
base64.b64encode(contents),
ex=TRP_IMPORT_REDIS_TTL,
)
r.set(
f"{TRP_IMPORT_META_PREFIX}{job_id}",
json.dumps(meta_data).encode("utf-8"),
ex=TRP_IMPORT_REDIS_TTL,
)
except Exception as e:
logger.error("Transportistas import: Redis store error: %s", e)
raise HTTPException(status_code=500, detail="No se pudo encolar el archivo.")
try:
upload_dir = layout_path("imports", "temp")
os.makedirs(upload_dir, exist_ok=True)
with open(os.path.join(upload_dir, f"trp_{job_id}.csv"), "wb") as f:
f.write(contents)
with open(os.path.join(upload_dir, f"trp_{job_id}.meta.json"), "w") as f:
json.dump(meta_data, f)
except Exception as e:
logger.warning("Transportistas import: local file save failed: %s", e)
scan_file.apply_async(args=[job_id], task_id=job_id)
def run_scan_background():
try:
run_scan_sync(job_id)
except Exception as e:
logger.exception("Transportistas import: background scan failed: %s", e)
threading.Thread(target=run_scan_background, daemon=True).start()
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):
try:
r = _get_redis()
raw = r.get(f"{TRP_IMPORT_STATUS_PREFIX}{job_id}")
if raw:
data = json.loads(raw.decode("utf-8"))
return data
except Exception as e:
logger.debug("Transportistas import: could not read status from Redis: %s", e)
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}
result = getattr(task_result, "result", None)
if isinstance(result, dict) and result.get("status") in ("finished", "warning"):
return result
logger.warning("Transportistas 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):
try:
r = _get_redis()
r.set(
f"{TRP_IMPORT_STATUS_PREFIX}{job_id}",
json.dumps({"status": "processing", "message": "Insertando..."}).encode("utf-8"),
ex=TRP_IMPORT_REDIS_TTL,
)
except Exception as e:
logger.debug("Transportistas import: could not write processing status: %s", e)
def run_commit_background():
try:
run_commit_sync(job_id)
except Exception as e:
logger.exception("Transportistas import: background commit failed: %s", e)
threading.Thread(target=run_commit_background, daemon=True).start()
return {
"status": "committing",
"message": "Inserción iniciada.",
"commit_job_id": job_id,
}

View File

@@ -0,0 +1,22 @@
from pydantic import BaseModel
from typing import Optional
class ImportJobResponse(BaseModel):
job_id: str
status: str
message: str
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
updated: Optional[int] = 0
skipped_invalid: Optional[int] = 0
skipped_duplicate: Optional[int] = 0
skipped_details: Optional[list] = None

View File

@@ -0,0 +1,344 @@
"""
Tareas Celery para importación CSV de Transportistas.
Flujo: scan_file (validación) → commit.
Usa layouts_csv.common (storage, normalize, meta, responses, csv_reader).
"""
import json
import logging
import os
from typing import Dict, Any, Optional, List, Set
from core.celery_app import celery_app
from core.database import CoreSessionLocal
from ..common import storage as common_storage
from ..common import normalize as common_normalize
from ..common import meta as common_meta
from ..common import responses as common_responses
from ..common import csv_reader as common_csv_reader
from .template_config import row_from_template
from .validators import validate_row_transporter, validate_row_transporter_desfase
from .common.mappers import row_to_transporter_data, row_to_transporter_data_for_update
from .common.fk_loader import load_transportistas_fk_sets
logger = logging.getLogger(__name__)
JOB_TYPE = "trp"
# Para routes.py (deben coincidir con common_storage.storage_keys(JOB_TYPE, job_id))
TRP_IMPORT_FILE_PREFIX = "trp_import_file:"
TRP_IMPORT_META_PREFIX = "trp_import_meta:"
TRP_IMPORT_ERROR_LINES_PREFIX = "trp_import_error_lines:"
TRP_IMPORT_STATUS_PREFIX = "trp_import_status:"
TRP_IMPORT_REDIS_TTL = common_storage.IMPORT_REDIS_TTL
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 _do_scan(job_id: str, progress_callback: Optional[Any] = None) -> Dict[str, Any]:
file_path = common_storage.ensure_file_from_redis(JOB_TYPE, job_id, "Transportistas import")
if not file_path:
return {"status": "failed", "error": "Archivo no encontrado (expirado o no subido). Sube de nuevo."}
common_storage.ensure_meta_from_redis(JOB_TYPE, job_id, file_path, "Transportistas import")
error_path = common_storage.error_path_for_job(JOB_TYPE, job_id)
try:
total_rows = common_csv_reader.count_csv_rows(file_path)
except Exception as e:
return {"status": "failed", "error": str(e)}
try:
tenant_id, company_id = common_meta.require_tenant_context(file_path)
except ValueError as e:
return {"status": "failed", "error": str(e)}
meta = common_meta.load_meta(file_path) or {}
actualizar = meta.get("actualizar", False)
(
existing_transporter_keys,
valid_country_ame,
state_descriptions_upper,
state_country_set,
) = load_transportistas_fk_sets(tenant_id, company_id)
if actualizar:
pass # existing_transporter_keys ya cargado
else:
existing_transporter_keys = set()
error_count = 0
processed_rows = 0
errors_detail: List[Dict[str, Any]] = []
error_lines_list: List[int] = []
try:
with open(error_path, "w", encoding="utf-8") as f_err:
for i, row in common_csv_reader.iter_csv_rows(file_path):
if progress_callback and i % 500 == 0:
progress_callback(i, total_rows, error_count)
row_norm = row_from_template(row, common_normalize.normalize_header)
_ = validate_row_transporter_desfase(row_norm, i)
err = validate_row_transporter(
row_norm,
i,
actualizar=actualizar,
existing_transporter_keys=existing_transporter_keys,
valid_country_ame=valid_country_ame,
state_descriptions_upper=state_descriptions_upper,
state_country_set=state_country_set,
)
if err:
error_count += 1
error_lines_list.append(err["line"])
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
if error_lines_list:
common_storage.store_error_lines(JOB_TYPE, job_id, error_lines_list)
except Exception as e:
logger.error("Transportistas import scan failed: %s", e)
return {"status": "failed", "error": str(e)}
return common_responses.scan_result(job_id, processed_rows, error_count, errors_detail)
def run_scan_sync(job_id: str) -> Dict[str, Any]:
result = _do_scan(job_id, progress_callback=None)
try:
r = _get_redis()
r.set(
f"{TRP_IMPORT_STATUS_PREFIX}{job_id}",
json.dumps(result).encode("utf-8"),
ex=TRP_IMPORT_REDIS_TTL,
)
except Exception as e:
logger.warning("Transportistas import: failed to store scan status in Redis: %s", e)
return result
@celery_app.task(bind=True)
def scan_file(self, job_id: str, config: str = None):
logger.info("Transportistas import: starting scan for job %s", job_id)
def on_progress(current: int, total: int, errors: int) -> None:
self.update_state(state="PROGRESS", meta={"current": current, "total": total, "errors": errors})
result = _do_scan(job_id, progress_callback=on_progress)
try:
r = _get_redis()
r.set(
f"{TRP_IMPORT_STATUS_PREFIX}{job_id}",
json.dumps(result).encode("utf-8"),
ex=TRP_IMPORT_REDIS_TTL,
)
except Exception as e:
logger.warning("Transportistas import: failed to store scan status in Redis: %s", e)
return result
def _do_commit(job_id: str) -> Dict[str, Any]:
file_path = common_storage.ensure_file_from_redis(JOB_TYPE, job_id, "Transportistas import")
if not file_path:
alt_path = common_storage.file_path_for_job(JOB_TYPE, job_id)
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:
common_storage.ensure_meta_from_redis(JOB_TYPE, job_id, file_path, "Transportistas import")
error_path = common_storage.error_path_for_job(JOB_TYPE, job_id)
error_lines = common_storage.get_error_lines(JOB_TYPE, job_id, error_path)
try:
tenant_id, company_id = common_meta.require_tenant_context(file_path)
except ValueError as e:
return {"status": "failed", "error": str(e)}
meta = common_meta.load_meta(file_path) or {}
actualizar = meta.get("actualizar", False)
(
existing_transporter_keys,
valid_country_ame,
state_descriptions_upper,
state_country_set,
) = load_transportistas_fk_sets(tenant_id, company_id)
from api.v1.modules.a76.transportation.transporters.services import TransporterService
from api.v1.modules.a76.transportation.transporters.dto import TransporterCreateDTO, TransporterUpdateDTO
inserted_count = 0
updated_count = 0
skipped_invalid = 0
skipped_duplicate = 0
skipped_details: List[Dict[str, Any]] = []
seen_keys_in_file: Dict[str, int] = {}
meta_path = common_meta.get_meta_path(file_path)
try:
with CoreSessionLocal() as session:
for i, row in common_csv_reader.iter_csv_rows(file_path):
if i in error_lines:
continue
row_norm = row_from_template(row, common_normalize.normalize_header)
err = validate_row_transporter(
row_norm,
i,
actualizar=actualizar,
existing_transporter_keys=existing_transporter_keys,
valid_country_ame=valid_country_ame,
state_descriptions_upper=state_descriptions_upper,
state_country_set=state_country_set,
)
if err:
skipped_invalid += 1
tk = (row_norm.get("CLAVE TRANSPORTISTA") or "").strip()[:23] or "-"
skipped_details.append({
"line": i,
"transporter_key": tk,
"reason": f"{err.get('col', '')}: {err.get('msg', '')}",
})
continue
tk = (row_norm.get("CLAVE TRANSPORTISTA") or "").strip()[:23] or ""
if not tk:
skipped_invalid += 1
continue
tk_upper = tk.upper()
if tk_upper in seen_keys_in_file:
skipped_duplicate += 1
skipped_details.append({
"line": i,
"transporter_key": tk,
"reason": "Clave duplicada en el archivo (se usa la primera)",
})
continue
seen_keys_in_file[tk_upper] = i
existing = TransporterService.get_by_id(session, tk, tenant_id, company_id)
if not existing:
existing = TransporterService.get_by_id_ignore_case(session, tk, tenant_id, company_id)
try:
if existing:
if actualizar:
data = row_to_transporter_data_for_update(row_norm, existing)
else:
data = row_to_transporter_data(row_norm)
if not data or not data.get("transporter_key"):
skipped_invalid += 1
continue
update_data = TransporterUpdateDTO(**{k: v for k, v in data.items() if k != "transporter_key"})
TransporterService.update(session, existing.transporter_key, tenant_id, update_data, company_id)
updated_count += 1
else:
data = row_to_transporter_data(row_norm)
if not data or not data.get("transporter_key"):
skipped_invalid += 1
continue
create_data = TransporterCreateDTO(**data)
TransporterService.create(session, create_data, tenant_id, company_id)
inserted_count += 1
except Exception as db_err:
session.rollback()
skipped_invalid += 1
skipped_details.append({
"line": i,
"transporter_key": tk,
"reason": str(db_err),
})
continue
try:
session.commit()
except Exception as db_err:
session.rollback()
logger.error("Transportistas import DB error: %s", db_err)
return {"status": "failed", "error": str(db_err)}
except Exception as e:
logger.exception("Transportistas import task failed")
return {"status": "failed", "error": str(e)}
common_storage.cleanup_import_job(
JOB_TYPE, job_id,
file_path=file_path,
error_path=error_path,
meta_path=meta_path,
)
try:
r = _get_redis()
r.delete(f"{TRP_IMPORT_STATUS_PREFIX}{job_id}")
except Exception as e:
logger.warning("Transportistas import: failed to delete status key: %s", e)
total_ok = inserted_count + updated_count
if total_ok == 0 and (skipped_invalid + skipped_duplicate) > 0:
return {
"status": "warning",
"inserted": inserted_count,
"updated": updated_count,
"skipped_invalid": skipped_invalid,
"skipped_duplicate": skipped_duplicate,
"skipped_details": skipped_details,
"message": f"No se insertaron registros. {skipped_invalid + skipped_duplicate} rechazados.",
}
if total_ok == 0:
return {
"status": "failed",
"error": "No hay registros válidos en el archivo CSV",
"inserted": 0,
"updated": 0,
"skipped_invalid": skipped_invalid,
"skipped_duplicate": skipped_duplicate,
"skipped_details": skipped_details,
}
return {
"status": "finished",
"inserted": inserted_count,
"updated": updated_count,
"skipped_invalid": skipped_invalid,
"skipped_duplicate": skipped_duplicate,
"skipped_details": skipped_details,
}
def run_commit_sync(job_id: str) -> Dict[str, Any]:
result = _do_commit(job_id)
try:
r = _get_redis()
r.set(
f"{TRP_IMPORT_STATUS_PREFIX}{job_id}",
json.dumps(result).encode("utf-8"),
ex=TRP_IMPORT_REDIS_TTL,
)
except Exception as e:
logger.warning("Transportistas import: failed to store commit status in Redis: %s", e)
return result
@celery_app.task(bind=True)
def insert_valid_rows(self, job_id: str):
logger.info("Transportistas import: starting commit for job %s", job_id)
result = _do_commit(job_id)
try:
r = _get_redis()
r.set(
f"{TRP_IMPORT_STATUS_PREFIX}{job_id}",
json.dumps(result).encode("utf-8"),
ex=TRP_IMPORT_REDIS_TTL,
)
except Exception as e:
logger.warning("Transportistas import: failed to store commit status in Redis: %s", e)
return result

View File

@@ -0,0 +1,55 @@
"""
Configuración de plantilla CSV para Transportistas (EstructuraCatTransportistas).
Mapeo Clarion: Col A = CLAVE TRANSPORTISTA, B = NOMBRE, ... R = DIRECTORIO FTP, S = desfase.
"""
from typing import Dict, List, Any
TEMPLATE_COLUMNS: Dict[str, List[Dict[str, Any]]] = {
"transporters": [
{"canonical": "CLAVE TRANSPORTISTA", "aliases": ["TRANSPORTISTA", "CLAVE TRANS", "CARRIER KEY"]},
{"canonical": "NOMBRE", "aliases": ["NOMBRE TRANSPORTISTA", "NAME"]},
{"canonical": "NOMBRE CORTO", "aliases": ["NOMBRE CORTO TRANS", "SHORT NAME"]},
{"canonical": "RESPONSABLE", "aliases": ["RESPONSABLE TRANS"]},
{"canonical": "RFC", "aliases": ["RFC TRANS"]},
{"canonical": "CALLES", "aliases": ["STREETS", "DIRECCION"]},
{"canonical": "CODIGO POSTAL", "aliases": ["CODIGO POSTAL TRANS", "POSTAL CODE", "CP"]},
{"canonical": "CIUDAD", "aliases": ["CITY", "CIUDAD TRANS"]},
{"canonical": "ESTADO", "aliases": ["STATE", "ESTADO TRANS"]},
{"canonical": "PAIS", "aliases": ["COUNTRY", "PAIS TRANS"]},
{"canonical": "CODIGO CARGADOR", "aliases": ["COD CARGADOR", "LOADER CODE"]},
{"canonical": "CODIGO CAAT", "aliases": ["CAAT", "CODIGO CAAT TRANS"]},
{"canonical": "CODIGO TRANS", "aliases": ["COD TRANS", "TRANSPORT CODE", "SCAC"]},
{"canonical": "TIPO INTERFASE TRANS", "aliases": ["TIPO INTERFASE", "INTERFACE TYPE"]},
{"canonical": "SERVIDOR FTP", "aliases": ["FTP SERVER", "SERVIDOR FTP TRANS"]},
{"canonical": "USUARIO FTP", "aliases": ["FTP USER", "USUARIO FTP TRANS"]},
{"canonical": "CLAVE ACCESO FTP", "aliases": ["CLAVE FTP", "FTP PASSWORD", "PASSWORD FTP"]},
{"canonical": "DIRECTORIO FTP", "aliases": ["FTP DIRECTORY", "DIRECTORIO FTP TRANS"]},
{"canonical": "COL_EXTRA", "aliases": ["COLUMNA S", "COL S", "COL 19"]},
],
}
def build_normalized_lookup(normalize_header_fn) -> Dict[str, str]:
cols = TEMPLATE_COLUMNS.get("transporters")
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]:
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

@@ -0,0 +1,3 @@
from .create import validate_row_transporter, validate_row_transporter_desfase
__all__ = ["validate_row_transporter", "validate_row_transporter_desfase"]

View File

@@ -0,0 +1,94 @@
"""
Validaciones comunes de fila para import CSV de transportistas.
Paridad Clarion: VALIDACIONES_TRANSPORTISTAS, VALIDA_TODA_TRANSPORTISTAS, VALIDA_PARCIAL_TRANSPORTISTAS.
"""
from typing import Dict, Any, Optional, Set, Tuple
from ..common.common_validators import (
MAX_LEN,
check_required,
check_max_length,
check_pais_catalog_transportistas,
check_estado_catalog_transportistas,
check_estado_pais_consistency_transportistas,
)
def validate_row_transporter_required(row: Dict[str, Any], line_num: int) -> Optional[Dict[str, Any]]:
"""Col A (CLAVE TRANSPORTISTA) obligatoria."""
return check_required(row, "CLAVE TRANSPORTISTA", MAX_LEN["transporter_key"], line_num)
def validate_row_transporter_required_full(
row: Dict[str, Any],
line_num: int,
) -> Optional[Dict[str, Any]]:
"""VALIDA_TODA: Col A y Col B (NOMBRE) obligatorios cuando se agrega o se reemplaza."""
err = validate_row_transporter_required(row, line_num)
if err:
return err
err = check_max_length(
row, "NOMBRE", MAX_LEN["name"], line_num, required=True
)
if err:
return err
return None
def validaciones_transportistas(
row: Dict[str, Any],
line_num: int,
valid_country_ame: Optional[Set[str]] = None,
state_descriptions_upper: Optional[Set[str]] = None,
state_country_set: Optional[Set[Tuple[str, str]]] = None,
) -> Optional[Dict[str, Any]]:
"""
VALIDACIONES_TRANSPORTISTAS: estado, país, estado-país (Clarion).
"""
err = check_estado_catalog_transportistas(row, line_num, state_descriptions_upper)
if err:
return err
err = check_pais_catalog_transportistas(row, line_num, valid_country_ame)
if err:
return err
err = check_estado_pais_consistency_transportistas(row, line_num, state_country_set)
if err:
return err
return None
def valida_toda_transportistas(
row: Dict[str, Any],
line_num: int,
valid_country_ame: Optional[Set[str]] = None,
state_descriptions_upper: Optional[Set[str]] = None,
state_country_set: Optional[Set[Tuple[str, str]]] = None,
) -> Optional[Dict[str, Any]]:
"""VALIDA_TODA: obligatorios A y B + VALIDACIONES_TRANSPORTISTAS (usado para agregar nuevo o reemplazar)."""
err = validate_row_transporter_required_full(row, line_num)
if err:
return err
return validaciones_transportistas(
row,
line_num,
valid_country_ame=valid_country_ame,
state_descriptions_upper=state_descriptions_upper,
state_country_set=state_country_set,
)
def valida_parcial_transportistas(
row: Dict[str, Any],
line_num: int,
valid_country_ame: Optional[Set[str]] = None,
state_descriptions_upper: Optional[Set[str]] = None,
state_country_set: Optional[Set[Tuple[str, str]]] = None,
) -> Optional[Dict[str, Any]]:
"""VALIDA_PARCIAL_TRANSPORTISTAS: solo VALIDACIONES_TRANSPORTISTAS (actualizar registro existente)."""
return validaciones_transportistas(
row,
line_num,
valid_country_ame=valid_country_ame,
state_descriptions_upper=state_descriptions_upper,
state_country_set=state_country_set,
)

View File

@@ -0,0 +1,60 @@
"""
Punto de entrada de validación para import de una fila transportista.
Actualizar = merge: si no existe la clave se crea (valida_toda), si existe se actualizan solo campos enviados (valida_parcial).
Reemplazar = reescribir: siempre valida_toda; en commit se sustituye el registro completo o se agrega.
"""
from typing import Dict, Any, Optional, Set, Tuple
from .common import (
validate_row_transporter_required,
valida_toda_transportistas,
valida_parcial_transportistas,
)
from ..common.common_validators import check_desfase_transportistas
def validate_row_transporter(
row: Dict[str, Any],
line_num: int,
actualizar: bool = False,
existing_transporter_keys: Optional[Set[str]] = None,
valid_country_ame: Optional[Set[str]] = None,
state_descriptions_upper: Optional[Set[str]] = None,
state_country_set: Optional[Set[Tuple[str, str]]] = None,
) -> Optional[Dict[str, Any]]:
"""
Valida una fila de CSV de transportistas.
1. CLAVE TRANSPORTISTA (Col A) vacía → error.
2. Si actualizar y clave existe en catálogo → valida_parcial (solo validaciones de catálogo; campos vacíos = mantener actual).
3. Si reemplazar o (actualizar y clave no existe) → valida_toda (A y B obligatorios + validaciones).
"""
err = validate_row_transporter_required(row, line_num)
if err:
return err
existing = existing_transporter_keys or set()
clave = (row.get("CLAVE TRANSPORTISTA") or "").strip().upper()
use_partial = actualizar and bool(clave and clave in existing)
if use_partial:
err = valida_parcial_transportistas(
row,
line_num,
valid_country_ame=valid_country_ame,
state_descriptions_upper=state_descriptions_upper,
state_country_set=state_country_set,
)
else:
err = valida_toda_transportistas(
row,
line_num,
valid_country_ame=valid_country_ame,
state_descriptions_upper=state_descriptions_upper,
state_country_set=state_country_set,
)
return err
def validate_row_transporter_desfase(row: Dict[str, Any], line_num: int) -> Optional[Dict[str, Any]]:
"""Advertencia de desfase si COL_EXTRA (Col S) tiene valor. No bloqueante."""
return check_desfase_transportistas(row, line_num)

View File

@@ -1,11 +1,19 @@
from fastapi import APIRouter
from api.v1.common.tenant_crud_routes import TenantCRUDRoutes
from .dto import TransporterCreateDTO, TransporterResponseDTO, TransporterUpdateDTO
from .services import TransporterService
from api.v1.modules.a76.layouts_csv.transportistas.routes import router as imports_router
# Create router using TenantCRUDRoutes factory
# Note: transporter_key is a string (not int) and is used as the primary key
router = TenantCRUDRoutes(
# Main router: transporters CRUD + CSV imports
router = APIRouter()
# CSV import (upload → scan → status → commit)
router.include_router(imports_router, prefix="/transporters/imports", tags=["a76 / transporters / csv_import"])
# CRUD routes
crud_router = TenantCRUDRoutes(
service=TransporterService,
create_schema=TransporterCreateDTO,
update_schema=TransporterUpdateDTO,
@@ -20,3 +28,4 @@ router = TenantCRUDRoutes(
default_page_size=50,
max_page_size=100,
).router
router.include_router(crud_router)

View File

@@ -46,6 +46,7 @@ celery_app.conf.update(
"api.v1.modules.a76.layouts_csv.vehicles.tasks",
"api.v1.modules.a76.layouts_csv.drivers.tasks",
"api.v1.modules.a76.layouts_csv.trailers.tasks",
"api.v1.modules.a76.layouts_csv.transportistas.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",