""" Rutas de importación CSV para Clases de Materiales. Flujo: 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 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, CLS_IMPORT_META_PREFIX, CLS_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 (ACT): validación parcial si la clase existe"), siempre_toda: bool = Query(False, description="Forzar siempre validación completa"), 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(f"Classes 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": "material_classes", "actualizar": actualizar, "siempre_toda": siempre_toda, } try: common_storage.store_import_file( JOB_TYPE, job_id, contents, meta_data, tenant_id=int(tenant_id), company_id=company_id, ttl=CLS_IMPORT_REDIS_TTL, log_label="Classes import", ) except common_storage.ImportStoreError as e: logger.error(f"Classes import: store error: {e}") raise HTTPException(status_code=500, detail="No se pudo encolar el archivo.") except Exception as e: logger.error(f"Classes 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="classes_scan_file", task_group="layouts_csv", task_origin="a76/layouts_csv/classes/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): 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) # Si Celery devolvió el resultado como string (p. ej. JSON), parsear y devolver como scan si aplica if isinstance(result, str): try: parsed = json.loads(result) if isinstance(parsed, dict) and ( parsed.get("status") == "waiting_confirmation" or (parsed.get("job_id") and "total_rows" in parsed) ): return parsed if isinstance(parsed, dict) and parsed.get("status") in ("finished", "warning"): return normalize_commit_status_payload(parsed) except (json.JSONDecodeError, TypeError): pass # Si el resultado tiene forma de escaneo (waiting_confirmation), devolverlo para que el front muestre el modal if isinstance(result, dict) and ( result.get("status") == "waiting_confirmation" or (result.get("job_id") and "total_rows" in result) ): return result logger.warning("Classes 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), ): 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=f"{CLS_IMPORT_META_PREFIX}{job_id}", celery_task=insert_valid_rows, task_name="classes_insert_valid_rows", task_origin="a76/layouts_csv/classes/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("cls", job_id)