Files
plantillas-proyectos/backend/api/v1/modules/a76/layouts_csv/pedmientos/routes.py

195 lines
6.5 KiB
Python

"""
Rutas de importación CSV para Pedimentos.
Mismo flujo que customs_brokers/imports: upload → scan → status (polling) → commit.
"""
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, Optional
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 api.v1.modules.core.tasks_tracking import track_and_dispatch
from ..common.track_commit_dispatch import dispatch_tracked_layouts_csv_commit
from ..common.responses import normalize_commit_status_payload
from .schemas import ImportJobResponse
from .tasks import (
scan_file,
insert_valid_rows,
JOB_TYPE as PED_JOB_TYPE,
PED_IMPORT_REDIS_TTL,
)
from ..common import storage as common_storage
from ..common.error_csv import download_scan_errors_csv_stream
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 actualizar: validación parcial si el pedimento existe"),
dateFormat: Optional[str] = Query(None, description="Formato de fecha: dd/mm/yyyy, mm/dd/yyyy, yyyy-mm-dd"),
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, ["csv_upload.process"])
except Exception as e:
logger.error(f"Pedimentos 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: Dict[str, Any] = {
"tenant_id": tenant_id,
"company_id": company_id,
"user_id": current_user.get("id"),
"template_id": "pedimentos",
"actualizar": actualizar,
}
if dateFormat:
meta_data["dateFormat"] = dateFormat
try:
common_storage.store_import_file(
PED_JOB_TYPE,
job_id,
contents,
meta_data,
tenant_id=int(tenant_id),
company_id=company_id,
ttl=PED_IMPORT_REDIS_TTL,
log_label="Pedimentos import",
)
except common_storage.ImportStoreError as e:
logger.error(f"Pedimentos import: store error: {e}")
raise HTTPException(status_code=500, detail="No se pudo encolar el archivo.")
except Exception as e:
logger.error(f"Pedimentos import: Redis store error: {e}")
raise HTTPException(status_code=500, detail="No se pudo encolar el archivo.")
track_and_dispatch(
db=db,
task=scan_file,
tenant_id=int(tenant_id),
company_id=company_id,
requested_by_user=current_user.get("preferred_username")
or current_user.get("email")
or current_user.get("sub"),
task_name="pedmientos_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/pedmientos/upload",
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 normalize_commit_status_payload(result)
return {"status": "finished", "result": result}
result = getattr(task_result, "result", None)
if isinstance(result, dict) and result.get("status") in ("finished", "warning"):
return normalize_commit_status_payload(result)
logger.warning("Pedimentos 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,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
"""
Fase 2: Usuario confirma; se encola la inserción de filas válidas.
"""
_, meta_key, _ = common_storage.storage_keys(PED_JOB_TYPE, job_id)
r = _get_redis()
commit_id = dispatch_tracked_layouts_csv_commit(
db=db,
current_user=current_user,
redis_client=r,
layout_job_id=job_id,
meta_redis_key=meta_key,
celery_task=insert_valid_rows,
task_name="pedmientos_insert_valid_rows",
task_origin="a76/layouts_csv/pedmientos/commit",
args=[job_id],
)
return {
"status": "committing",
"message": "Inserción iniciada.",
"commit_job_id": commit_id,
}
@router.get("/{job_id}/errors/scan-csv")
async def download_scan_errors_csv(job_id: str):
return download_scan_errors_csv_stream(PED_JOB_TYPE, job_id)