fix/forzar el procesamiento de un pedimento cargado por datastage
This commit is contained in:
169
tasks.py
169
tasks.py
@@ -2,7 +2,7 @@ from celery import Celery
|
||||
from celery_app import celery_app
|
||||
import asyncio
|
||||
import logging
|
||||
from typing import Dict, Any
|
||||
from typing import Dict, Any, List
|
||||
from contextlib import asynccontextmanager
|
||||
from controllers.RESTController import rest_controller
|
||||
from controllers.SOAPController import soap_controller
|
||||
@@ -158,6 +158,173 @@ def pedimento_completo_task(self, request_data: Dict[str, Any]):
|
||||
|
||||
return run_async_task(_execute_pedimento_completo)
|
||||
|
||||
@celery_app.task(bind=True, name='tasks.multi_pedimento_completo_task')
|
||||
def multi_pedimento_completo_task(self, pedimentos: List[str], organizacion: str):
|
||||
"""
|
||||
Tarea asíncrona para procesar MÚLTIPLES pedimentos completos.
|
||||
|
||||
Args:
|
||||
pedimentos: Lista de IDs de pedimentos a procesar
|
||||
organizacion: ID de la organización
|
||||
"""
|
||||
import time
|
||||
from datetime import datetime
|
||||
|
||||
start_time = time.time()
|
||||
|
||||
results = {
|
||||
"total": len(pedimentos),
|
||||
"successful": [],
|
||||
"failed": [],
|
||||
"started_at": datetime.utcnow().isoformat()
|
||||
}
|
||||
|
||||
total = len(pedimentos)
|
||||
|
||||
for idx, pedimento_id in enumerate(pedimentos, 1):
|
||||
try:
|
||||
# Actualizar progreso (igual que en tus otras tareas)
|
||||
self.update_state(
|
||||
state='PROGRESS',
|
||||
meta={
|
||||
'status': f'Procesando pedimento {idx}/{total}',
|
||||
'current': idx,
|
||||
'total': total,
|
||||
'current_pedimento': pedimento_id,
|
||||
'percentage': round((idx / total) * 100, 2)
|
||||
}
|
||||
)
|
||||
|
||||
logger.info(f"[MULTI] Procesando pedimento {idx}/{total}: {pedimento_id}")
|
||||
|
||||
# Preparar datos exactamente como lo espera la tarea individual
|
||||
request_data = {
|
||||
"pedimento": pedimento_id,
|
||||
"organizacion": organizacion
|
||||
}
|
||||
|
||||
# Reutilizar la lógica de la tarea individual
|
||||
# Esto ejecuta el mismo código que tu endpoint individual
|
||||
async def _execute():
|
||||
return await _execute_pedimento_completo_logic(request_data)
|
||||
|
||||
result = run_async_task(_execute)
|
||||
|
||||
results["successful"].append({
|
||||
"pedimento_id": pedimento_id,
|
||||
"result": result
|
||||
})
|
||||
|
||||
logger.info(f"[MULTI] Pedimento {pedimento_id} procesado exitosamente")
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[MULTI] Error procesando pedimento {pedimento_id}: {e}")
|
||||
results["failed"].append({
|
||||
"pedimento_id": pedimento_id,
|
||||
"error": str(e)
|
||||
})
|
||||
|
||||
elapsed_time = time.time() - start_time
|
||||
|
||||
results["completed_at"] = datetime.utcnow().isoformat()
|
||||
results["elapsed_seconds"] = round(elapsed_time, 2)
|
||||
results["success_count"] = len(results["successful"])
|
||||
results["failed_count"] = len(results["failed"])
|
||||
|
||||
return results
|
||||
|
||||
async def _execute_pedimento_completo_logic(request_data: dict) -> dict:
|
||||
"""
|
||||
Lógica compartida para procesar un pedimento completo.
|
||||
Esta es la misma lógica que usa tu endpoint individual.
|
||||
"""
|
||||
operation_name = "pedimento_completo"
|
||||
service_data = None
|
||||
|
||||
try:
|
||||
logger.info(f"Procesando pedimento completo - Pedimento: {request_data['pedimento']}")
|
||||
|
||||
# Validar datos de entrada
|
||||
await _validate_request_data(request_data)
|
||||
|
||||
# Obtener servicio
|
||||
service_data = await _get_pedimento_service(
|
||||
pedimento_id=request_data['pedimento'],
|
||||
service_type=3,
|
||||
operation_name=operation_name
|
||||
)
|
||||
|
||||
# Actualizar estado a "En proceso"
|
||||
update_success = await _update_service_status(
|
||||
service_data['id'], ESTADO_EN_PROCESO, service_data, operation_name
|
||||
)
|
||||
|
||||
if not update_success:
|
||||
raise Exception("Error al actualizar estado del servicio")
|
||||
|
||||
# Obtener credenciales VUCEM
|
||||
contribuyente_id = service_data.get('pedimento', {}).get('contribuyente', '')
|
||||
if not contribuyente_id:
|
||||
await _update_service_status(service_data['id'], ESTADO_ERROR, service_data, operation_name)
|
||||
raise Exception("ID de contribuyente no encontrado")
|
||||
|
||||
credentials = await _get_vucem_credentials(contribuyente_id, operation_name)
|
||||
|
||||
# Procesar petición SOAP
|
||||
soap_response = await get_soap_pedimento_completo(
|
||||
credenciales=credentials,
|
||||
response_service=service_data,
|
||||
soap_controller=soap_controller
|
||||
)
|
||||
|
||||
if not soap_response:
|
||||
raise Exception("Error en la petición SOAP")
|
||||
|
||||
# Actualizar datos del pedimento
|
||||
xml_content = soap_response.get('xml_content', {})
|
||||
if xml_content:
|
||||
update_content = {k: v for k, v in xml_content.items() if k != 'identificadores_ed'}
|
||||
update_content['existe_expediente'] = True
|
||||
await rest_controller.put_pedimento(
|
||||
service_data['pedimento']['id'],
|
||||
update_content
|
||||
)
|
||||
|
||||
# Procesar COVEs
|
||||
coves = xml_content.get('coves', [])
|
||||
if coves:
|
||||
await _post_coves(
|
||||
response_service=service_data,
|
||||
coves=coves
|
||||
)
|
||||
|
||||
# Procesar documentos digitalizados
|
||||
identificadores_ed = xml_content.get('identificadores_ed', [])
|
||||
if identificadores_ed:
|
||||
await _post_edocuments(
|
||||
response_service=service_data,
|
||||
identificadores_ed=identificadores_ed
|
||||
)
|
||||
|
||||
# Finalizar servicio
|
||||
await _update_service_status(service_data['id'], ESTADO_FINALIZADO, service_data, operation_name)
|
||||
|
||||
return {
|
||||
"success": True,
|
||||
"pedimento_id": request_data['pedimento'],
|
||||
"message": "Pedimento procesado exitosamente"
|
||||
}
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error en pedimento {request_data['pedimento']}: {e}")
|
||||
if service_data:
|
||||
await _update_service_status(service_data['id'], ESTADO_ERROR, service_data, operation_name)
|
||||
return {
|
||||
"success": False,
|
||||
"pedimento_id": request_data['pedimento'],
|
||||
"error": str(e)
|
||||
}
|
||||
|
||||
@celery_app.task(bind=True)
|
||||
def partidas_task(self, **kwargs):
|
||||
"""Tarea asíncrona para obtener partidas"""
|
||||
|
||||
Reference in New Issue
Block a user