""" 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, 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, 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) def _assert_classes_csv_job_access( db: Session, job_id: str, current_user: Dict[str, Any], ) -> None: from api.v1.modules.core.tasks_tracking.models import TaskRun company_id: Optional[int] = None row = db.query(TaskRun).filter(TaskRun.task_id == job_id).first() if row is not None and row.company_id is not None: company_id = int(row.company_id) r = _get_redis() if company_id is None: raw = r.get(f"{CLS_IMPORT_META_PREFIX}{job_id}") if raw: meta = json.loads(raw.decode("utf-8")) cid = meta.get("company_id") if cid is not None: company_id = int(cid) if company_id is None and row is not None and row.meta_payload: layout_jid = row.meta_payload.get("layout_import_job_id") if layout_jid: raw2 = r.get(f"{CLS_IMPORT_META_PREFIX}{layout_jid}") if raw2: meta2 = json.loads(raw2.decode("utf-8")) cid2 = meta2.get("company_id") if cid2 is not None: company_id = int(cid2) if company_id is None: raise HTTPException(status_code=403, detail="Sin acceso a este job") try: validate_access_to_resource( db, company_id, current_user, ["csv_upload.process"] ) except HTTPException: raise except Exception as e: logger.error("Classes import: job access validation failed: %s", e) raise HTTPException(status_code=403, detail="Sin acceso a este job") from None @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, ["csv_upload.process"]) 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, db: Session = Depends(get_core_db), current_user: Dict[str, Any] = Depends(get_current_user), ): _assert_classes_csv_job_access(db, job_id, current_user) 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], required_permissions=["csv_upload.process"], ) 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, db: Session = Depends(get_core_db), current_user: Dict[str, Any] = Depends(get_current_user), ): _assert_classes_csv_job_access(db, job_id, current_user) return download_scan_errors_csv_stream("cls", job_id)