import os from datetime import datetime from decimal import Decimal import csv import json import logging import re import unicodedata from celery import shared_task from typing import Dict, Any, Optional from core.database import CoreSessionLocal # Models are imported inside tasks to avoid circular dependencies and mapper initialization issues in the API process # We'll need schemas for validation # from api.v1.modules.a76.invoices.schemas import InvoiceHeaderCreate # But for Phase 1 we use a lighter check logger = logging.getLogger(__name__) class ForeignKeyValidator: def __init__(self, session, tenant_id, company_id): self.session = session self.tenant_id = tenant_id self.company_id = company_id self.cache = {} # {(model_name, value): bool} def check_exists(self, model, value, field_name="id", is_public=False): if value is None: return True # Assume optional if None, or let DB handle not-null key = (model.__name__, value) if key in self.cache: return self.cache[key] query = self.session.query(getattr(model, field_name)).filter(getattr(model, field_name) == value) if not is_public: query = query.filter(model.tenant_id == self.tenant_id, model.company_id == self.company_id) exists = query.first() is not None self.cache[key] = exists return exists @shared_task(bind=True) def scan_file(self, job_id: str, file_path: str, model_target: str, config: str = None): """ Pass 1: Read CSV, Validate types, Write Errors to JSONL. """ logger.info(f"Starting scan for job {job_id} target {model_target}") # 1. Setup Error Log error_path = file_path.replace("temp", "errors").replace(".csv", ".jsonl") os.makedirs(os.path.dirname(error_path), exist_ok=True) total_rows = 0 error_count = 0 processed_rows = 0 # 2. Count Total (Quick Pass) or just estimate # For better progress, we can get file line count first try: with open(file_path, 'r', encoding='utf-8-sig') as f: total_rows = sum(1 for _ in f) - 1 # Minus header except Exception as e: return {"status": "failed", "error": f"Cannot read file: {e}"} footer_config = parse_footer_config(config) date_format = footer_config.get("dateFormat") # Validate and set default date_format if not provided if not date_format: date_format = "yyyy-mm-dd" # Default to ISO format logger.info(f"No date_format specified in config, using default: {date_format}") try: with open(file_path, 'r', encoding='utf-8-sig') as f_in, \ open(error_path, 'w', encoding='utf-8') as f_err: # Detect Delimiter sample = f_in.read(2048) f_in.seek(0) try: dialect = csv.Sniffer().sniff(sample, delimiters=",;\t") except: dialect = 'excel' reader = csv.DictReader(f_in, dialect=dialect) for i, row in enumerate(reader, start=1): # Check for Progress Update if i % 1000 == 0: self.update_state(state='PROGRESS', meta={ 'current': i, 'total': total_rows, 'errors': error_count }) # Validation (Phase 1: Minimal) row_norm = normalize_row(row) errors = validate_row_phase_1(row_norm, model_target, i, date_format) if errors: error_count += 1 # Write simple JSON error f_err.write(json.dumps(errors) + "\n") processed_rows += 1 except Exception as e: logger.error(f"Scan failed: {e}") return {"status": "failed", "error": str(e)} # 4. Result return { "status": "waiting_confirmation", "job_id": job_id, "total_rows": processed_rows, "error_count": error_count, "valid_rows": processed_rows - error_count, "error_file": error_path } def validate_row_phase_1( row: Dict[str, Any], target: str, line_num: int, date_format: Optional[str], ) -> Dict[str, Any]: """ Minimal validation: Unique IDs and Dates. Target: 'invoice_header' or 'invoice_details' """ errors = {} # A. Invoice Header if target == 'invoice_header': # 1. Unique ID if not row.get('NUMERO FACTURA') and not row.get('NUM FACTURA') and not row.get('ID'): return {"line": line_num, "col": "NUMERO FACTURA", "msg": "Requerido"} # 2. Date Format date_str = row.get('FECHA FACTURA') if date_str: if not is_valid_date(date_str, date_format): expected = display_date_format(date_format) return { "line": line_num, "col": "FECHA FACTURA", "msg": f"Formato inválido ({expected})", } else: return {"line": line_num, "col": "FECHA FACTURA", "msg": "Requerido"} # B. Invoice Details (Parts) elif target == 'invoice_details': # 1. Line Number if not row.get('LINEA'): return {"line": line_num, "col": "LINEA", "msg": "Requerido"} # 2. Parent Link (Invoice Number) if not (row.get('NUMERO FACTURA') or row.get('NUM FACTURA') or row.get('FACTURA')): return {"line": line_num, "col": "NUMERO FACTURA", "msg": "Requerido"} # 2. Parent Link (Simplified for now, we assume parent exists or is in same batch) # In a real scenario, we'd check if the invoice exists. pass return errors if errors else None def parse_footer_config(config: Optional[str]) -> Dict[str, Any]: if not config: return {} try: if isinstance(config, str): return json.loads(config) if isinstance(config, dict): return config except Exception: return {} return {} def display_date_format(date_format: Optional[str]) -> str: if not date_format: return "YYYY-MM-DD" return date_format.upper() def parse_date(date_text: Optional[str], date_format: Optional[str]) -> Optional[datetime.date]: if not date_text: return None candidates = [] fmt_map = { "dd/mm/yyyy": "%d/%m/%Y", "mm/dd/yyyy": "%m/%d/%Y", "yyyy-mm-dd": "%Y-%m-%d", } if date_format and date_format in fmt_map: candidates.append(fmt_map[date_format]) candidates.extend(["%Y-%m-%d", "%d/%m/%Y", "%m/%d/%Y"]) for fmt in candidates: try: return datetime.strptime(str(date_text).strip(), fmt).date() except ValueError: continue return None def is_valid_date(date_text: Optional[str], date_format: Optional[str]) -> bool: return parse_date(date_text, date_format) is not None 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 normalize_row(row: Dict[str, Any]) -> Dict[str, Any]: return {normalize_header(k): v for k, v in row.items()} def parse_int(value: Any) -> Optional[int]: if value is None: return None text = str(value).strip() if not text: return None try: return int(text) except ValueError: return None def parse_decimal(value: Any) -> Optional[Decimal]: if value is None: return None text = str(value).strip() if not text: return None text = text.replace(",", "") try: return Decimal(text) except Exception: return None def parse_currency(value: Optional[str], currency_type: Optional[str]): from api.v1.modules.a76.invoices.models import Currency if value: normalized = normalize_header(value) if normalized in {"MN", "M N", "NACIONAL", "LOCAL", "PESOS", "PESO"}: return Currency.LOCAL if normalized in {"ME", "M E", "EXTRANJERA", "EXTRANJERO", "FOREIGN", "USD", "DOLAR", "DOLARES"}: return Currency.FOREIGN if "MANUAL" in normalized: return Currency.MANUAL if currency_type and str(currency_type).strip().upper() == "MXN": return Currency.LOCAL if currency_type: return Currency.FOREIGN return Currency.MANUAL def parse_weight_unit(value: Optional[str]): from api.v1.modules.a76.invoices.models import WeightUnit if not value: return None normalized = normalize_header(value) if normalized in {"KG", "KGS", "KILOS", "KILOGRAMOS"}: return WeightUnit.KGS if normalized in {"LB", "LBS", "LIBRAS"}: return WeightUnit.LBS return None def resolve_tenant_fk_id( session: CoreSessionLocal, model, value: Optional[int], tenant_id: int, company_id: int, cache: Dict[int, Optional[int]], ) -> Optional[int]: if value is None: return None if value in cache: return cache[value] exists = ( session.query(model.id) .filter( model.id == value, model.tenant_id == tenant_id, model.company_id == company_id, ) .scalar() ) cache[value] = value if exists is not None else None return cache[value] def resolve_public_code( session: CoreSessionLocal, model, column, value: Optional[str], cache: Dict[str, Optional[str]], ) -> Optional[str]: if not value: return None normalized = str(value).strip().upper() if not normalized: return None if normalized in cache: return cache[normalized] exists = session.query(column).filter(column == normalized).scalar() cache[normalized] = normalized if exists is not None else None return cache[normalized] @shared_task(bind=True) def insert_valid_rows(self, job_id: str, model_target: str): """ Pass 2: Re-read CSV, Skip Errors, Bulk Insert. """ logger.info(f"Starting Commit for {job_id} target {model_target}") try: from api.v1.modules.a76.invoices.models import ( InvoiceHeader, InvoiceComplianceMx, InvoiceFinancials, InvoiceLogistics, InvoiceSalesDetails, OperationType, WeightUnit, ) from api.v1.modules.a76.clients_and_providers.models import ClientProvider from api.v1.modules.a76.customs_brokers.models import CustomsBroker from api.v1.modules.public.reference_data.currency_types.models import CurrencyType from api.v1.modules.public.reference_data.pedimento_regimens.models import RegimenPedimento from api.v1.modules.public.reference_data.code_pedimento_regimens.models import CodePedimentoRegimen from api.v1.modules.public.reference_data.pedimento_codes.models import PedimentoCode from api.v1.modules.public.reference_data.invoice_types.models import InvoiceType from api.v1.modules.public.reference_data.customs_sections.models import CustomsSection from api.v1.modules.a76.items.models import LineItem from api.v1.modules.a76.items.line_financials.models import LineFinancial from api.v1.modules.a76.items.line_quantities.models import LineQuantity from api.v1.modules.a76.items.line_customs.models import LineCustom from api.v1.modules.a76.items.line_descriptions.models import LineDescription from api.v1.modules.a76.parts.models import Part upload_dir = os.path.join(os.getcwd(), "uploads", "temp") file_path = os.path.join(upload_dir, f"{job_id}.csv") error_path = file_path.replace("temp", "errors").replace(".csv", ".jsonl") # 1. Load Error Line Numbers error_lines = set() 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) error_lines.add(err['line']) except: pass # Load Metadata (Context) meta_path = file_path.replace("temp", "temp").replace(".csv", ".meta.json") tenant_id = None company_id = None footer_config = {} if os.path.exists(meta_path): try: with open(meta_path, 'r') as f: meta = json.load(f) tenant_id = meta.get('tenant_id') company_id = meta.get('company_id') operation_type_raw = meta.get('operation_type', 'imp') footer_config = parse_footer_config(meta.get('footer_config')) except: pass if not tenant_id or not company_id: return {"status": "failed", "error": "Missing context (tenant/company)"} # 2. Re-read and Map # Initialize counters outside the session block so they're accessible later headers_to_insert = [] details_to_insert = [] skipped_invalid = 0 skipped_missing_invoice = 0 skipped_missing_fk = 0 skipped_fk_details = [] inserted_count = 0 response = None # Will be set inside the session block date_format = footer_config.get("dateFormat") # Validate and set default date_format if not provided if not date_format: date_format = "yyyy-mm-dd" # Default to ISO format logger.info(f"No date_format specified in config, using default: {date_format}") else: logger.info(f"Using date_format from config: {date_format}") # Default types from config or fallback op_type_value = OperationType(meta.get('operation_type', 'imp').lower()) inv_type_value = footer_config.get('invoice_type', 'TEM') logger.info(f"Processing CSV with operation_type={op_type_value}, invoice_type={inv_type_value}, date_format={date_format}") with CoreSessionLocal() as session: invoice_id_cache = {} cleared_invoices = set() # Track invoices where we've already cleared items in this job provider_cache: Dict[int, Optional[int]] = {} sold_to_cache: Dict[int, Optional[int]] = {} shipped_to_cache: Dict[int, Optional[int]] = {} broker_cache: Dict[int, Optional[int]] = {} regimen_cache: Dict[str, Optional[str]] = {} currency_type_cache: Dict[str, Optional[str]] = {} customs_section_cache: Dict[str, Optional[str]] = {} part_cache: Dict[str, Optional[int]] = {} validator = ForeignKeyValidator(session, tenant_id, company_id) with open(file_path, 'r', encoding='utf-8-sig') as f: # Detect Delimiter sample = f.read(2048) f.seek(0) try: dialect = csv.Sniffer().sniff(sample, delimiters=",;\t") except: dialect = 'excel' reader = csv.DictReader(f, dialect=dialect) for i, row in enumerate(reader, start=1): if i in error_lines: continue row_norm = normalize_row(row) # Mapping Logic if model_target == 'invoice_header': invoice_number = (row_norm.get('NUMERO FACTURA') or row_norm.get('NUM FACTURA') or row_norm.get('FACTURA') or '').strip() invoice_date = parse_date(row_norm.get('FECHA FACTURA') or row_norm.get('FECHA'), date_format) if not invoice_number or not invoice_date: skipped_invalid += 1 logger.debug(f"Row {i}: Skipped - missing invoice_number or invalid invoice_date. " f"Invoice: {invoice_number}, Date: {row_norm.get('FECHA FACTURA') or row_norm.get('FECHA')}") continue # --- NEW: Foreign Key Validations --- # 1. Invoice Type (Public) if not validator.check_exists(InvoiceType, inv_type_value, field_name="key", is_public=True): skipped_missing_fk += 1 reason = f"Tipo de factura '{inv_type_value}' no existe" skipped_fk_details.append({"line": i, "invoice": invoice_number, "reason": reason}) logger.warning(f"Row {i} (Invoice {invoice_number}): {reason}") continue # 2. Client/Provider (Tenant) provider_id = parse_int(row_norm.get('CLAVE PROVEEDOR')) if provider_id and not validator.check_exists(ClientProvider, provider_id): skipped_missing_fk += 1 reason = f"Proveedor ID '{provider_id}' no existe" skipped_fk_details.append({"line": i, "invoice": invoice_number, "reason": reason}) logger.warning(f"Row {i} (Invoice {invoice_number}): {reason}") continue # 3. Customs Broker (Tenant) broker_id = parse_int(row_norm.get('AGENTE ADUANAL')) if broker_id and not validator.check_exists(CustomsBroker, broker_id): skipped_missing_fk += 1 reason = f"Agente Aduanal ID '{broker_id}' no existe" skipped_fk_details.append({"line": i, "invoice": invoice_number, "reason": reason}) logger.warning(f"Row {i} (Invoice {invoice_number}): {reason}") continue # --- 4. Check for Existing Invoice (Upsert Logic) --- existing_header = None if invoice_number: existing_header = ( session.query(InvoiceHeader) .filter( InvoiceHeader.tenant_id == tenant_id, InvoiceHeader.company_id == company_id, InvoiceHeader.invoice_number == invoice_number, InvoiceHeader.invoice_type == inv_type_value ) .first() ) if existing_header: # UPDATE existing header header = existing_header header.invoice_date = invoice_date header.operation_type = op_type_value header.is_updated = True # Mark as updated header.updated_date = datetime.utcnow() header.document_type = resolve_public_code( session, RegimenPedimento, RegimenPedimento.code, (row_norm.get('REGIMEN') or row_norm.get('CLAVEDOCUMENTO')), regimen_cache, ) header.project_number = (row_norm.get('NUM PROYECTO') or row_norm.get('NUMPROYECTO') or None) header.purchase_order = (row_norm.get('ORDEN COMPRA') or row_norm.get('ORDENCOMPRA') or None) header.alternate_invoice = (row_norm.get('FACTURA ALTERNA') or None) header.invoice_ref = (row_norm.get('FACTURA EXPO REF') or row_norm.get('FACTURAEXPOREF') or None) header.emission_date = parse_date(row_norm.get('FECHA EMISION'), date_format) header.observation_es = (row_norm.get('OBSERVACIONES E') or None) header.observation_en = (row_norm.get('OBSERVACIONES I') or None) logger.info(f"Row {i}: Updating existing invoice {invoice_number}") # Clean up related data that will be re-inserted/updated # Note: compliance, financials, logistics are 1-to-1 relationships and will be updated by assignment below # but we might want to be explicit if ORM doesn't handle replace well. # SQLAlchemy relationship assignment usually handles 1-to-1 updates correctly. else: # CREATE new header header = InvoiceHeader( invoice_number=invoice_number, invoice_date=invoice_date, operation_type=op_type_value, is_updated=False, system="CSV", capture_date=datetime.utcnow(), invoice_type=inv_type_value, document_type=resolve_public_code( session, RegimenPedimento, RegimenPedimento.code, (row_norm.get('REGIMEN') or row_norm.get('CLAVEDOCUMENTO')), regimen_cache, ), project_number=(row_norm.get('NUM PROYECTO') or row_norm.get('NUMPROYECTO') or None), purchase_order=(row_norm.get('ORDEN COMPRA') or row_norm.get('ORDENCOMPRA') or None), alternate_invoice=(row_norm.get('FACTURA ALTERNA') or None), invoice_ref=(row_norm.get('FACTURA EXPO REF') or row_norm.get('FACTURAEXPOREF') or None), emission_date=parse_date(row_norm.get('FECHA EMISION'), date_format), observation_es=(row_norm.get('OBSERVACIONES E') or None), observation_en=(row_norm.get('OBSERVACIONES I') or None), tenant_id=tenant_id, company_id=company_id, ) compliance = InvoiceComplianceMx( remesa=parse_int(row_norm.get('REMESA')), aduana=resolve_public_code( session, CustomsSection, CustomsSection.customs_code, row_norm.get('ADUANA DE CRUCE'), customs_section_cache, ), provider_id=resolve_tenant_fk_id( session, ClientProvider, parse_int(row_norm.get('CLAVE PROVEEDOR')), tenant_id, company_id, provider_cache, ), sold_to_id=resolve_tenant_fk_id( session, ClientProvider, parse_int(row_norm.get('CLAVE VENDIDO A')), tenant_id, company_id, sold_to_cache, ), shipped_to_id=resolve_tenant_fk_id( session, ClientProvider, parse_int(row_norm.get('CLAVE ENVIADO A')), tenant_id, company_id, shipped_to_cache, ), customs_broker_id=resolve_tenant_fk_id( session, CustomsBroker, parse_int(row_norm.get('AGENTE ADUANAL')), tenant_id, company_id, broker_cache, ), edocument=(row_norm.get('E DOCUMENT') or None), vucem_operation_num=(row_norm.get('NUM OPERACION') or None), tenant_id=tenant_id, company_id=company_id, ) financials_currency_type = resolve_public_code( session, CurrencyType, CurrencyType.code, row_norm.get('CLAVE MONEDA'), currency_type_cache, ) financials = InvoiceFinancials( currency=parse_currency(row_norm.get('TIPO MONEDA'), financials_currency_type), currency_type=financials_currency_type, exchange_rate=parse_decimal(row_norm.get('TIPO DE CAMBIO')), freight=parse_decimal(row_norm.get('FLETES')), insurance_value=parse_decimal(row_norm.get('VALOR SEGUROS')), insurance=parse_decimal(row_norm.get('SEGUROS')), packaging=parse_decimal(row_norm.get('EMBALAJES')), other_increments=parse_decimal(row_norm.get('OTROS INCREMENTABLES')), tenant_id=tenant_id, company_id=company_id, ) weight_type = parse_weight_unit(row_norm.get('TIPO PESO')) logistics = None if weight_type or row_norm.get('TIPO TRANSPORTE') or row_norm.get('NUMERO TRANSPORTE'): logistics = InvoiceLogistics( carrier_id=(row_norm.get('CLAVE TRANSPORTISTA') or None), driver_name=(row_norm.get('NOMBRE CONDUCTOR') or None), transport_type=str(row_norm.get('TIPO TRANSPORTE') or "none").lower(), transport_num=(row_norm.get('NUMERO TRANSPORTE') or None), weight_type=weight_type or WeightUnit.KGS, seal_number=(row_norm.get('PRECINTO') or None), incoterm=(row_norm.get('CLAVE INCOTERM') or None), entry_exit_date=parse_date(row_norm.get('FECHA EMISION'), date_format), tenant_id=tenant_id, company_id=company_id, ) header.compliance_mx = compliance header.financials = financials if logistics: header.logistics = logistics headers_to_insert.append(header) elif model_target == 'invoice_details': invoice_number = (row_norm.get('NUMERO FACTURA') or row_norm.get('NUM FACTURA') or '').strip() if not invoice_number: skipped_invalid += 1 continue if invoice_number in invoice_id_cache: invoice_id = invoice_id_cache[invoice_number] else: invoice_id = ( session.query(InvoiceHeader.id) .filter( InvoiceHeader.tenant_id == tenant_id, InvoiceHeader.company_id == company_id, InvoiceHeader.invoice_number == invoice_number, ) .scalar() ) invoice_id_cache[invoice_number] = invoice_id if not invoice_id: logger.warning( "Invoice not found for details row %s (invoice_number=%s)", i, invoice_number, ) skipped_missing_invoice += 1 continue # --- Prevent Duplicates: Clear existing items for this invoice (Once per job) --- if invoice_id not in cleared_invoices: logger.info(f"Clearing existing details for Invoice {invoice_number} (ID: {invoice_id}) to prevent duplicates") # 1. Delete Items (Cascades to LineItem, LineFinancial, etc. if DB configured, check models) # Checking Item model, we usually need to be careful. # Assuming Cascade delete is set up on FKs or we rely on ORM cascade if using relationships. # Here we use bulk delete. session.query(Item).filter(Item.invoice_id == invoice_id).delete(synchronize_session=False) # 2. Delete InvoiceSalesDetails session.query(InvoiceSalesDetails).filter(InvoiceSalesDetails.invoice_id == invoice_id).delete(synchronize_session=False) cleared_invoices.add(invoice_id) # --- NEW LOGIC: Expanded Anexo 76 Structure --- # A. Find/Cache Part part_num = (row_norm.get('NUMPARTE') or row_norm.get('NUMERO PARTE') or '').strip() part_id = None if part_num: part_id = part_cache.get(part_num) if part_id is None: p = session.query(Part.id).filter( Part.part_number == part_num, Part.tenant_id == tenant_id, Part.company_id == company_id ).first() if p: part_id = p.id part_cache[part_num] = part_id line_num_val = (row_norm.get('LINEA') or row_norm.get('RENGLON') or row_norm.get('PARTIDA')) line_num = parse_int(line_num_val) or (len(details_to_insert) + 1) # 1. Parent Item item = Item( invoice_id=invoice_id, tenant_id=tenant_id, company_id=company_id, item_type="N", # Default to Normal system_origin="CSV" ) session.add(item) session.flush() # Need item.id # 2. Main Line line = LineItem( item_id=item.id, line_number=line_num, part_number=part_id, tenant_id=tenant_id, company_id=company_id ) session.add(line) session.flush() # Need line.id # 3. Financial Data price = parse_decimal(row_norm.get('PRECIO UNITARIO') or row_norm.get('PRECIOUNITARIO')) val_com = parse_decimal(row_norm.get('VALOR COMERCIAL') or row_norm.get('VALORCOMERCIAL')) qty = parse_decimal(row_norm.get('CANTIDAD')) session.add(LineFinancial( item_line_id=line.id, unit_price=price, commercial_value=val_com or (price * qty if price and qty else None), )) # 4. Quantities if qty: session.add(LineQuantity( item_line_id=line.id, quantity=qty, )) # 5. Customs/Fraction origin = row_norm.get('PAIS ORIGEN') or row_norm.get('PAISORIGEN') fraction = row_norm.get('FRACCION') if origin or fraction: session.add(LineCustom( item_line_id=line.id, fraction=fraction, origin_country=origin, )) # 6. Description desc = row_norm.get('DESCRIPCION') if desc: session.add(LineDescription( item_line_id=line.id, description_spanish=desc, )) # 7. Legacy Sales Details (For specific audit/UI fields) detail = InvoiceSalesDetails( invoice_id=invoice_id, line_number=line_num, sales_order=(row_norm.get('ORDEN DE COMPRA') or row_norm.get('ORDENCOMPRA') or None), line_bundles=parse_int(row_norm.get('CANTIDAD BULTOS') or row_norm.get('CANTIDADBULTOS')), tenant_id=tenant_id, company_id=company_id, ) session.add(detail) details_to_insert.append(item) # Use as counter/ref # 3. Bulk Insert (ORM Transaction) try: if model_target == 'invoice_header': if headers_to_insert: logger.info(f"Attempting to commit {len(headers_to_insert)} headers") session.add_all(headers_to_insert) session.commit() inserted_count = len(headers_to_insert) logger.info(f"Headers commit successful. Inserted: {inserted_count}") else: logger.warning(f"No headers to insert for job {job_id}") else: if details_to_insert: logger.info(f"Attempting to commit {len(details_to_insert)} items and related data") session.commit() # Everything was already added with session.add() inserted_count = len(details_to_insert) logger.info(f"Details commit successful. Inserted: {inserted_count}") else: logger.warning(f"No details to insert for job {job_id}") except Exception as db_err: session.rollback() logger.error(f"DB Error during {model_target} commit: {db_err}") import traceback logger.error(traceback.format_exc()) return {"status": "failed", "error": str(db_err)} # 4. Determine final status and prepare response (inside session block to access variables) total_skipped = skipped_invalid + skipped_missing_fk + skipped_missing_invoice # Log summary logger.info(f"Job {job_id} completed. Inserted: {inserted_count}, Skipped: {total_skipped} " f"(invalid: {skipped_invalid}, missing_fk: {skipped_missing_fk}, missing_invoice: {skipped_missing_invoice})") # Prepare response based on results if inserted_count == 0: if total_skipped > 0: logger.warning(f"No valid records to insert for job {job_id}. All {total_skipped} records were rejected.") response = { "status": "warning", "inserted": 0, "skipped_invalid": skipped_invalid, "skipped_missing_invoice": skipped_missing_invoice, "skipped_missing_fk": skipped_missing_fk, "skipped_details": skipped_fk_details, "message": f"No se insertaron registros. {total_skipped} fueron rechazados." } else: logger.error(f"No valid records found in CSV for job {job_id}") response = { "status": "failed", "error": "No hay registros válidos en el archivo CSV", "inserted": 0, "skipped_invalid": skipped_invalid, "skipped_missing_invoice": skipped_missing_invoice, "skipped_missing_fk": skipped_missing_fk, "skipped_details": skipped_fk_details } else: # Success case - at least some records were inserted response = { "status": "finished", "inserted": inserted_count, "skipped_invalid": skipped_invalid, "skipped_missing_invoice": skipped_missing_invoice, "skipped_missing_fk": skipped_missing_fk, "skipped_details": skipped_fk_details } except Exception as e: logger.error(f"Task failed: {e}") import traceback logger.error(traceback.format_exc()) return {"status": "failed", "error": str(e)} # 5. Cleanup try: if os.path.exists(file_path): os.remove(file_path) if os.path.exists(error_path): os.remove(error_path) except: logger.warning("Failed to cleanup temp files") # Ensure response is defined (fallback in case of unexpected errors) if response is None: logger.error(f"Unexpected error: response not set for job {job_id}") response = { "status": "failed", "error": "Error inesperado durante el procesamiento", "inserted": 0, "skipped_invalid": skipped_invalid, "skipped_missing_invoice": skipped_missing_invoice, "skipped_missing_fk": skipped_missing_fk, "skipped_details": skipped_fk_details } return response