T2026-05-030
This commit is contained in:
@@ -1,3 +1,4 @@
|
||||
from api.organization.models import Organizacion
|
||||
from celery import group
|
||||
from celery import shared_task, group
|
||||
from api.customs.models import *
|
||||
@@ -8,6 +9,11 @@ import requests
|
||||
from config.settings import SERVICE_API_URL_V2
|
||||
from datetime import datetime
|
||||
import json
|
||||
import logging
|
||||
import uuid
|
||||
# este solo fue para pruebas personales, lo dejo por si en un futuro lo requiero
|
||||
TEST_ORG_ID = uuid.UUID('defc7848-4f39-4d67-9dba-5bb445248d23')
|
||||
logger = logging.getLogger('api.customs.microservice_v2')
|
||||
|
||||
def credenciales_to_dict(credenciales):
|
||||
if not credenciales:
|
||||
@@ -132,7 +138,7 @@ def procesar_edocs_pedimento(pedimento_id):
|
||||
}
|
||||
|
||||
response = requests.post(
|
||||
f"{SERVICE_API_URL_V2}/services/download/edoc/",
|
||||
f"{SERVICE_API_URL_V2}/services/download/all/edocs/",
|
||||
data=json.dumps(payload),
|
||||
headers={"Content-Type": "application/json"}
|
||||
)
|
||||
@@ -277,27 +283,40 @@ def procesar_remesas(organizacion_id):
|
||||
pedimentos = Pedimento.objects.filter(organizacion_id=organizacion_id)
|
||||
|
||||
for pedimento in pedimentos:
|
||||
if not pedimento.documents.filter(document_type=3).exists(): # Tipo 3: Remesa
|
||||
# Convertir el pedimento a JSON usando el serializer
|
||||
logger.info(f"pedimento >>>> {pedimento}")
|
||||
try:
|
||||
# if pedimento.documents.filter(document_type=3).exists(): # Remesa ya descargada
|
||||
# logger.info(f"Pedimento {pedimento.pedimento} ya tiene remesa descargada, omitiendo.")
|
||||
# continue
|
||||
|
||||
pedimento_dict = pedimento_to_dict(pedimento)
|
||||
credenciales = Vucem.objects.filter(id=CredencialesImportador.objects.filter(rfc=pedimento.contribuyente).first().vucem.id).first()
|
||||
|
||||
credencial_importador = CredencialesImportador.objects.filter(rfc=pedimento.contribuyente).first()
|
||||
if not credencial_importador:
|
||||
logger.warning(f"Sin credenciales para RFC {pedimento.contribuyente} (pedimento {pedimento.pedimento}), omitiendo.")
|
||||
continue
|
||||
|
||||
credenciales = Vucem.objects.filter(id=credencial_importador.vucem.id).first()
|
||||
if not credenciales:
|
||||
logger.warning(f"Credencial Vucem no encontrada para pedimento {pedimento.pedimento}, omitiendo.")
|
||||
continue
|
||||
|
||||
credenciales_dict = credenciales_to_dict(credenciales)
|
||||
|
||||
|
||||
payload = {
|
||||
"pedimento": pedimento_dict,
|
||||
"credencial": credenciales_dict
|
||||
}
|
||||
|
||||
|
||||
|
||||
response = requests.post(
|
||||
f"{SERVICE_API_URL_V2}/services/remesas",
|
||||
f"{SERVICE_API_URL_V2}/services/remesas/",
|
||||
data=json.dumps(payload),
|
||||
headers={"Content-Type": "application/json"}
|
||||
)
|
||||
# Aquí puedes continuar con el resto de tu lógica
|
||||
logger.info(f"Servicio enviado para pedimento {pedimento.pedimento} — status {response.status_code}")
|
||||
|
||||
print(f"Servicio enviado para pedimento {pedimento.pedimento}")
|
||||
except Exception as e:
|
||||
logger.error(f"Error procesando remesa para pedimento {pedimento.pedimento}: {e}", exc_info=True)
|
||||
|
||||
@shared_task
|
||||
def procesar_coves(organizacion_id):
|
||||
@@ -522,6 +541,34 @@ def ejecutar_todos_por_organizacion(organizacion_id):
|
||||
procesar_pedimentos_completos.delay(organizacion_id)
|
||||
procesar_remesas.delay(organizacion_id)
|
||||
|
||||
def ejecutar_basicos_organizacion(organizacion_id):
|
||||
# solo coves y e documents, si es necesario ya en un futuro se agregan los de partidas, pedimento completo y esas madres
|
||||
procesar_coves.delay(organizacion_id)
|
||||
procesar_acuse_coves.delay(organizacion_id)
|
||||
procesar_edocs.delay(organizacion_id)
|
||||
procesar_acuses.delay(organizacion_id)
|
||||
# procesar_partidas.delay(organizacion_id)
|
||||
# procesar_pedimentos_completos.delay(organizacion_id)
|
||||
# procesar_remesas.delay(organizacion_id)
|
||||
|
||||
@shared_task
|
||||
def process_organization_batch(org_id):
|
||||
"""
|
||||
Procesa todos los tipos de documentos pendientes para una organización.
|
||||
"""
|
||||
ejecutar_basicos_organizacion(org_id)
|
||||
|
||||
@shared_task
|
||||
def process_all_organizations():
|
||||
"""
|
||||
Envía una tarea por organización activa a la cola org_processing.
|
||||
"""
|
||||
active_orgs = Organizacion.objects.filter(is_active=True, is_verified=True)
|
||||
|
||||
for org in active_orgs:
|
||||
process_organization_batch.apply_async(
|
||||
args=[org.id],
|
||||
queue='org_processing'
|
||||
)
|
||||
return f"Dispatched {active_orgs.count()} organizations"
|
||||
|
||||
|
||||
Reference in New Issue
Block a user