feature/tabla-bitacora-logs-task

This commit is contained in:
2026-03-24 13:07:35 -06:00
parent 5b2962d53b
commit b5b99a5ddb
50 changed files with 2289 additions and 430 deletions

View File

@@ -6,6 +6,7 @@ from sqlalchemy.orm import Session
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 .task import process_export_invoice_task
@@ -25,8 +26,16 @@ def trigger_invoice_process(
"""
tenant_id = validate_access_to_resource(db, company_id, current_user)
task = process_export_invoice_task.apply_async(
args=[invoice_id, str(tenant_id), str(company_id)]
task = track_and_dispatch(
db=db,
task=process_export_invoice_task,
tenant_id=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="process_export_invoice_task",
task_group="invoices",
task_origin="a76/invoices/exports/process",
args=[invoice_id, str(tenant_id), str(company_id)],
)
return {"task_id": task.id}

View File

@@ -6,6 +6,7 @@ from sqlalchemy.orm import Session
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 .task import revert_invoice_task
@@ -33,8 +34,16 @@ def trigger_invoice_revert(
or "system"
)
task = revert_invoice_task.apply_async(
args=[invoice_id, str(tenant_id), str(company_id), str(cancelled_by)]
task = track_and_dispatch(
db=db,
task=revert_invoice_task,
tenant_id=tenant_id,
company_id=company_id,
requested_by_user=str(cancelled_by),
task_name="revert_export_invoice_task",
task_group="invoices",
task_origin="a76/invoices/exports/revert",
args=[invoice_id, str(tenant_id), str(company_id), str(cancelled_by)],
)
return {"task_id": task.id}

View File

@@ -6,6 +6,7 @@ from sqlalchemy.orm import Session
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 api.v1.modules.a76.invoices.models import InvoiceHeader, OperationType
from .task import process_invoice_task
@@ -34,12 +35,28 @@ def trigger_invoice_process(
raise HTTPException(status_code=404, detail=f"Factura {invoice_id} no encontrada.")
if invoice.operation_type == OperationType.EXP:
task = process_export_invoice_task.apply_async(
args=[invoice_id, str(tenant_id), str(company_id)]
task = track_and_dispatch(
db=db,
task=process_export_invoice_task,
tenant_id=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="process_export_invoice_task",
task_group="invoices",
task_origin="a76/invoices/exports/process",
args=[invoice_id, str(tenant_id), str(company_id)],
)
else:
task = process_invoice_task.apply_async(
args=[invoice_id, str(tenant_id), str(company_id)]
task = track_and_dispatch(
db=db,
task=process_invoice_task,
tenant_id=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="process_invoice_task",
task_group="invoices",
task_origin="a76/invoices/imports/process",
args=[invoice_id, str(tenant_id), str(company_id)],
)
return {"task_id": task.id}

View File

@@ -6,6 +6,7 @@ from sqlalchemy.orm import Session
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 api.v1.modules.a76.invoices.models import InvoiceHeader, OperationType
from .task import revert_invoice_task as revert_import_invoice_task
@@ -40,12 +41,28 @@ def trigger_invoice_revert(
raise HTTPException(status_code=404, detail=f"Factura {invoice_id} no encontrada.")
if invoice.operation_type == OperationType.EXP:
task = revert_export_invoice_task.apply_async(
args=[invoice_id, str(tenant_id), str(company_id), str(cancelled_by)]
task = track_and_dispatch(
db=db,
task=revert_export_invoice_task,
tenant_id=tenant_id,
company_id=company_id,
requested_by_user=str(cancelled_by),
task_name="revert_export_invoice_task",
task_group="invoices",
task_origin="a76/invoices/exports/revert",
args=[invoice_id, str(tenant_id), str(company_id), str(cancelled_by)],
)
else:
task = revert_import_invoice_task.apply_async(
args=[invoice_id, str(tenant_id), str(company_id), str(cancelled_by)]
task = track_and_dispatch(
db=db,
task=revert_import_invoice_task,
tenant_id=tenant_id,
company_id=company_id,
requested_by_user=str(cancelled_by),
task_name="revert_invoice_task",
task_group="invoices",
task_origin="a76/invoices/imports/revert",
args=[invoice_id, str(tenant_id), str(company_id), str(cancelled_by)],
)
return {"task_id": task.id}

View File

@@ -16,6 +16,7 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
from core.security import get_current_user, validate_access_to_resource
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .schemas import ImportJobResponse
from .tasks import (
@@ -26,6 +27,7 @@ from .tasks import (
BOM_IMPORT_REDIS_TTL,
)
from ..common.error_csv import download_scan_errors_csv_stream
from ..common.track_commit_dispatch import dispatch_tracked_layouts_csv_commit
router = APIRouter()
logger = logging.getLogger(__name__)
@@ -89,7 +91,18 @@ async def upload_import_file(
except Exception as e:
logger.warning(f"BOMs import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
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="boms_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/boms/upload",
args=[job_id],
task_id=job_id,
)
return ImportJobResponse(
job_id=job_id,
@@ -144,12 +157,26 @@ async def get_import_status(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str):
task = insert_valid_rows.delay(job_id)
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,
meta_redis_key=f"{BOM_IMPORT_META_PREFIX}{job_id}",
celery_task=insert_valid_rows,
task_name="boms_insert_valid_rows",
task_origin="a76/layouts_csv/boms/commit",
args=[job_id],
)
return {
"status": "committing",
"message": "Inserción iniciada.",
"commit_job_id": task.id,
"commit_job_id": commit_id,
}

View File

@@ -16,7 +16,9 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
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 .schemas import ImportJobResponse, CommitRequest
from .tasks import scan_file, insert_valid_rows, JOB_TYPE, CRREG_IMPORT_REDIS_TTL
from ..common import storage as common_storage
@@ -92,7 +94,20 @@ async def upload_import_file(
except Exception as e:
logger.warning("Cambio régimen/Regularización import: local file save failed: %s", e)
scan_file.apply_async(args=[job_id, model_target, footer_config], task_id=job_id)
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="cambio_regimen_regularizacion_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/cambio_regimen_regularizacion/upload",
args=[job_id, model_target, footer_config],
task_id=job_id,
)
return ImportJobResponse(
job_id=job_id,
@@ -148,13 +163,29 @@ async def get_import_status(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str, body: CommitRequest):
async def commit_import_job(
job_id: str,
body: CommitRequest,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
"""Usuario confirma; se encola la tarea de commit (por ahora sin inserción real)."""
task = insert_valid_rows.delay(job_id, body.model_target)
_, meta_key, _ = common_storage.storage_keys(JOB_TYPE, job_id)
r = _get_redis()
commit_id = dispatch_tracked_layouts_csv_commit(
db=db,
current_user=current_user,
redis_client=r,
meta_redis_key=meta_key,
celery_task=insert_valid_rows,
task_name="cambio_regimen_regularizacion_insert_valid_rows",
task_origin="a76/layouts_csv/cambio_regimen_regularizacion/commit",
args=[job_id, body.model_target],
)
return {
"status": "committing",
"message": "Proceso de commit iniciado.",
"commit_job_id": task.id,
"commit_job_id": commit_id,
}

View File

@@ -16,7 +16,9 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
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 .schemas import ImportJobResponse
from .tasks import (
scan_file,
@@ -93,7 +95,20 @@ async def upload_import_file(
except Exception as e:
logger.warning(f"Classes import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
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,
@@ -169,12 +184,26 @@ async def get_import_status(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str):
task = insert_valid_rows.delay(job_id)
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,
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": task.id,
"commit_job_id": commit_id,
}

View File

@@ -16,7 +16,9 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
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 .schemas import ImportJobResponse
from .tasks import (
scan_file,
@@ -92,7 +94,20 @@ async def upload_import_file(
except Exception as e:
logger.warning(f"CP import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
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="clients_and_providers_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/clients_and_providers/upload",
args=[job_id],
task_id=job_id,
)
return ImportJobResponse(
job_id=job_id,
@@ -151,15 +166,29 @@ async def get_import_status(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str):
async def commit_import_job(
job_id: str,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
"""
Fase 2: Usuario confirma; se encola la inserción de filas válidas.
"""
task = insert_valid_rows.delay(job_id)
r = _get_redis()
commit_id = dispatch_tracked_layouts_csv_commit(
db=db,
current_user=current_user,
redis_client=r,
meta_redis_key=f"{CP_IMPORT_META_PREFIX}{job_id}",
celery_task=insert_valid_rows,
task_name="clients_and_providers_insert_valid_rows",
task_origin="a76/layouts_csv/clients_and_providers/commit",
args=[job_id],
)
return {
"status": "committing",
"message": "Inserción iniciada.",
"commit_job_id": task.id,
"commit_job_id": commit_id,
}

View File

@@ -0,0 +1,58 @@
"""Encolar `insert_valid_rows` con registro en core.task_runs y validación de acceso."""
from __future__ import annotations
import json
from typing import Any
from uuid import uuid4
from celery import Task
from fastapi import HTTPException
from sqlalchemy.orm import Session
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from core.security import validate_access_to_resource
def dispatch_tracked_layouts_csv_commit(
*,
db: Session,
current_user: dict[str, Any],
redis_client: Any,
meta_redis_key: str,
celery_task: Task,
task_name: str,
task_origin: str,
args: list[Any],
) -> str:
"""
Lee tenant_id / company_id del meta en Redis, valida acceso y despacha la tarea Celery
con un task_id nuevo (commit) para que el polling use commit_job_id.
"""
raw = redis_client.get(meta_redis_key)
if not raw:
raise HTTPException(status_code=404, detail="Job no encontrado o expiró")
meta = json.loads(raw.decode("utf-8"))
tenant_id = int(meta["tenant_id"])
company_id = meta.get("company_id")
if company_id is not None:
try:
validate_access_to_resource(db, int(company_id), current_user)
except Exception:
raise HTTPException(status_code=403, detail="Sin acceso a este job") from None
commit_id = str(uuid4())
track_and_dispatch(
db=db,
task=celery_task,
tenant_id=tenant_id,
company_id=int(company_id) if company_id is not None else None,
requested_by_user=current_user.get("preferred_username")
or current_user.get("email")
or current_user.get("sub"),
task_name=task_name,
task_group="layouts_csv",
task_origin=task_origin,
args=args,
task_id=commit_id,
)
return commit_id

View File

@@ -16,7 +16,9 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
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 .schemas import ImportJobResponse
from .tasks import (
scan_file,
@@ -92,7 +94,20 @@ async def upload_import_file(
except Exception as e:
logger.warning(f"CB import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
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="customs_brokers_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/customs_brokers/upload",
args=[job_id],
task_id=job_id,
)
return ImportJobResponse(
job_id=job_id,
@@ -151,15 +166,29 @@ async def get_import_status(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str):
async def commit_import_job(
job_id: str,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
"""
Fase 2: Usuario confirma; se encola la inserción de filas válidas.
"""
task = insert_valid_rows.delay(job_id)
r = _get_redis()
commit_id = dispatch_tracked_layouts_csv_commit(
db=db,
current_user=current_user,
redis_client=r,
meta_redis_key=f"{CB_IMPORT_META_PREFIX}{job_id}",
celery_task=insert_valid_rows,
task_name="customs_brokers_insert_valid_rows",
task_origin="a76/layouts_csv/customs_brokers/commit",
args=[job_id],
)
return {
"status": "committing",
"message": "Inserción iniciada.",
"commit_job_id": task.id,
"commit_job_id": commit_id,
}

View File

@@ -16,7 +16,9 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
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 .schemas import ImportJobResponse
from .tasks import (
scan_file,
@@ -103,7 +105,20 @@ async def upload_import_file(
except Exception as e:
logger.warning(f"ER import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
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="exchange_rate_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/exchange_rate/upload",
args=[job_id],
task_id=job_id,
)
return ImportJobResponse(
job_id=job_id,
@@ -161,15 +176,29 @@ async def get_import_status(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str):
async def commit_import_job(
job_id: str,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
"""
Fase 2: Usuario confirma; se encola la inserción de filas válidas.
"""
task = insert_valid_rows.delay(job_id)
r = _get_redis()
commit_id = dispatch_tracked_layouts_csv_commit(
db=db,
current_user=current_user,
redis_client=r,
meta_redis_key=f"{ER_IMPORT_META_PREFIX}{job_id}",
celery_task=insert_valid_rows,
task_name="exchange_rate_insert_valid_rows",
task_origin="a76/layouts_csv/exchange_rate/commit",
args=[job_id],
)
return {
"status": "committing",
"message": "Inserción iniciada.",
"commit_job_id": task.id,
"commit_job_id": commit_id,
}

View File

@@ -19,10 +19,12 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
from core.security import get_current_user, validate_access_to_resource
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .schemas import ImportJobResponse, CommitRequest
from .tasks import scan_file, insert_valid_rows, JOB_TYPE, EXP_IMPORT_REDIS_TTL
from ..common import storage as common_storage
from ..common.track_commit_dispatch import dispatch_tracked_layouts_csv_commit
router = APIRouter()
logger = logging.getLogger(__name__)
@@ -99,7 +101,18 @@ async def upload_import_file(
except Exception as e:
logger.warning("Exportación import: local file save failed: %s", e)
scan_file.apply_async(args=[job_id, model_target, footer_config], task_id=job_id)
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="exportacion_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/exportacion/upload",
args=[job_id, model_target, footer_config],
task_id=job_id,
)
return ImportJobResponse(
job_id=job_id,
@@ -155,13 +168,29 @@ async def get_import_status(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str, body: CommitRequest):
async def commit_import_job(
job_id: str,
body: CommitRequest,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
"""Usuario confirma; se encola la tarea de commit (por ahora sin inserción real)."""
task = insert_valid_rows.delay(job_id, body.model_target)
_, meta_key, _ = common_storage.storage_keys(JOB_TYPE, job_id)
r = _get_redis()
commit_id = dispatch_tracked_layouts_csv_commit(
db=db,
current_user=current_user,
redis_client=r,
meta_redis_key=meta_key,
celery_task=insert_valid_rows,
task_name="exportacion_insert_valid_rows",
task_origin="a76/layouts_csv/exportacion/commit",
args=[job_id, body.model_target],
)
return {
"status": "committing",
"message": "Proceso de commit iniciado.",
"commit_job_id": task.id,
"commit_job_id": commit_id,
}

View File

@@ -16,6 +16,7 @@ from core.config import settings
from core.database import get_core_db
from core.paths import layout_path
from core.security import get_current_user, validate_access_to_resource
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .tasks import (
scan_file,
@@ -26,6 +27,7 @@ from .tasks import (
)
from .schemas import ImportJobResponse, ImportJobStatus, CommitRequest
from ..common import storage as common_storage
from ..common.track_commit_dispatch import dispatch_tracked_layouts_csv_commit
router = APIRouter()
logger = logging.getLogger(__name__)
@@ -111,7 +113,18 @@ async def upload_import_file(
logger.warning(f"Local file save failed (worker will use Redis): {e}")
# Trigger Celery Task (Async). Worker loads file from Redis.
scan_file.apply_async(args=[job_id, model_target, footer_config], task_id=job_id)
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="facturas_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/facturas/upload",
args=[job_id, model_target, footer_config],
task_id=job_id,
)
return ImportJobResponse(
job_id=job_id,
@@ -229,14 +242,29 @@ async def download_scan_errors_csv(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str, body: CommitRequest):
async def commit_import_job(
job_id: str,
body: CommitRequest,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
"""
Step 2: User confirms import. Trigger bulk insert.
"""
task = insert_valid_rows.delay(job_id, body.model_target)
r = _get_redis()
commit_id = dispatch_tracked_layouts_csv_commit(
db=db,
current_user=current_user,
redis_client=r,
meta_redis_key=f"{IMPORT_META_KEY_PREFIX}{job_id}",
celery_task=insert_valid_rows,
task_name="facturas_insert_valid_rows",
task_origin="a76/layouts_csv/facturas/commit",
args=[job_id, body.model_target],
)
return {
"status": "committing",
"message": "Bulk insert started.",
"commit_job_id": task.id
"commit_job_id": commit_id,
}

View File

@@ -16,6 +16,7 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
from core.security import get_current_user, validate_access_to_resource
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .schemas import ImportJobResponse
from .tasks import (
@@ -26,6 +27,7 @@ from .tasks import (
PART_IMPORT_REDIS_TTL,
)
from ..common.error_csv import download_scan_errors_csv_stream
from ..common.track_commit_dispatch import dispatch_tracked_layouts_csv_commit
router = APIRouter()
logger = logging.getLogger(__name__)
@@ -93,7 +95,18 @@ async def upload_import_file(
except Exception as e:
logger.warning(f"Parts import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
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="parts_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/parts/upload",
args=[job_id],
task_id=job_id,
)
return ImportJobResponse(
job_id=job_id,
@@ -169,12 +182,26 @@ async def get_import_status(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str):
task = insert_valid_rows.delay(job_id)
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,
meta_redis_key=f"{PART_IMPORT_META_PREFIX}{job_id}",
celery_task=insert_valid_rows,
task_name="parts_insert_valid_rows",
task_origin="a76/layouts_csv/parts/commit",
args=[job_id],
)
return {
"status": "committing",
"message": "Inserción iniciada.",
"commit_job_id": task.id,
"commit_job_id": commit_id,
}

View File

@@ -16,7 +16,9 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
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 .schemas import ImportJobResponse
from .tasks import (
scan_file,
@@ -100,7 +102,20 @@ async def upload_import_file(
except Exception as e:
logger.warning(f"Pedimentos import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
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="pedmientos_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/pedmientos/upload",
args=[job_id],
task_id=job_id,
)
return ImportJobResponse(
job_id=job_id,
@@ -158,15 +173,30 @@ async def get_import_status(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str):
async def commit_import_job(
job_id: str,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
"""
Fase 2: Usuario confirma; se encola la inserción de filas válidas.
"""
task = insert_valid_rows.delay(job_id)
_, meta_key, _ = common_storage.storage_keys(PED_JOB_TYPE, job_id)
r = _get_redis()
commit_id = dispatch_tracked_layouts_csv_commit(
db=db,
current_user=current_user,
redis_client=r,
meta_redis_key=meta_key,
celery_task=insert_valid_rows,
task_name="pedmientos_insert_valid_rows",
task_origin="a76/layouts_csv/pedmientos/commit",
args=[job_id],
)
return {
"status": "committing",
"message": "Inserción iniciada.",
"commit_job_id": task.id,
"commit_job_id": commit_id,
}

View File

@@ -17,6 +17,7 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
from core.security import get_current_user, validate_access_to_resource
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .schemas import ImportJobResponse
from .tasks import (
@@ -94,7 +95,20 @@ async def upload_import_file(
except Exception as e:
logger.warning(f"Trailers import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
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="trailers_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/trailers/upload",
args=[job_id],
task_id=job_id,
)
def run_scan_background():
try:

View File

@@ -17,6 +17,7 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
from core.security import get_current_user, validate_access_to_resource
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .schemas import ImportJobResponse
from .tasks import (
@@ -94,7 +95,21 @@ async def upload_import_file(
except Exception as e:
logger.warning("Transportistas import: local file save failed: %s", e)
scan_file.apply_async(args=[job_id], task_id=job_id)
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="transportistas_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/transportistas/upload",
args=[job_id],
task_id=job_id,
)
def run_scan_background():
try:
run_scan_sync(job_id)

View File

@@ -16,7 +16,9 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
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 .schemas import ImportJobResponse
from .tasks import (
scan_file,
@@ -95,7 +97,20 @@ async def upload_import_file(
except Exception as e:
logger.warning(f"FA import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
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="us_tariff_fractions_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/us_tariff_fractions/upload",
args=[job_id],
task_id=job_id,
)
return ImportJobResponse(
job_id=job_id,
@@ -153,15 +168,29 @@ async def get_import_status(job_id: str):
@router.post("/{job_id}/commit")
async def commit_import_job(job_id: str):
async def commit_import_job(
job_id: str,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
"""
Fase 2: Usuario confirma; se encola la inserción de filas válidas.
"""
task = insert_valid_rows.delay(job_id)
r = _get_redis()
commit_id = dispatch_tracked_layouts_csv_commit(
db=db,
current_user=current_user,
redis_client=r,
meta_redis_key=f"{FA_IMPORT_META_PREFIX}{job_id}",
celery_task=insert_valid_rows,
task_name="us_tariff_fractions_insert_valid_rows",
task_origin="a76/layouts_csv/us_tariff_fractions/commit",
args=[job_id],
)
return {
"status": "committing",
"message": "Inserción iniciada.",
"commit_job_id": task.id,
"commit_job_id": commit_id,
}

View File

@@ -17,6 +17,7 @@ from core.celery_app import celery_app
from core.database import get_core_db
from core.paths import layout_path
from core.security import get_current_user, validate_access_to_resource
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .schemas import ImportJobResponse
from .tasks import (
@@ -94,7 +95,18 @@ async def upload_import_file(
except Exception as e:
logger.warning(f"Vehicles import: local file save failed: {e}")
scan_file.apply_async(args=[job_id], task_id=job_id)
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="vehicles_scan_file",
task_group="layouts_csv",
task_origin="a76/layouts_csv/vehicles/upload",
args=[job_id],
task_id=job_id,
)
# Fallback sin worker: ejecutar scan en un hilo y guardar resultado en Redis
def run_scan_background():
try:

View File

@@ -6,6 +6,7 @@ from celery.result import AsyncResult
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 .task import generar_pdf_aviso_consolidado_exp_async
router = APIRouter()
@@ -41,6 +42,16 @@ async def trigger_descarga_aviso_consolidado_exp(
current_user: Dict[str, Any] = Depends(get_current_user),
db: Session = Depends(get_core_db)
):
validate_access_to_resource(db, company_id, current_user)
task = generar_pdf_aviso_consolidado_exp_async.delay(invoice_id, company_id)
tenant_id = validate_access_to_resource(db, company_id, current_user)
task = track_and_dispatch(
db=db,
task=generar_pdf_aviso_consolidado_exp_async,
tenant_id=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="generar_pdf_aviso_consolidado_exp_async",
task_group="reports",
task_origin="a76/reports/exportacion/aviso_consolidado",
args=[invoice_id, company_id],
)
return {"task_id": task.id, "message": "Generación iniciada"}

View File

@@ -5,7 +5,8 @@ from sqlalchemy.orm import Session
from typing import Dict, Any
from core.database import get_core_db as get_db
from core.security import get_current_user
from core.security import get_current_user, validate_access_to_resource
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .service import FIFOAssignmentService
from .task import generate_descarga_pdf_task # FORCE RELOAD
from celery.result import AsyncResult
@@ -40,6 +41,7 @@ def run_fifo_assignment(
async def trigger_descarga_generation(
invoice_id: int,
company_id: int,
db: Session = Depends(get_db),
current_user: Any = Depends(get_current_user)
):
"""
@@ -47,8 +49,18 @@ async def trigger_descarga_generation(
Retorna el task_id para polling.
"""
try:
# Lanza la tarea de Celery
task = generate_descarga_pdf_task.delay(invoice_id, company_id)
tenant_id = validate_access_to_resource(db, company_id, current_user)
task = track_and_dispatch(
db=db,
task=generate_descarga_pdf_task,
tenant_id=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="generate_descarga_pdf_task",
task_group="reports",
task_origin="a76/reports/exportacion/descargo",
args=[invoice_id, company_id],
)
return {"task_id": task.id, "status": "processing"}
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))

View File

@@ -1,8 +1,11 @@
from typing import Dict, Any
from fastapi import APIRouter, Depends, Body
from sqlalchemy.orm import Session
from celery.result import AsyncResult
from core.celery_app import celery_app
from core.security import get_current_user
from core.database import get_core_db
from core.security import get_current_user, get_tenant_from_token
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .task import generar_transmission_file_async
from .schemas import Mainx30GenerationRequest
@@ -35,9 +38,19 @@ async def get_task_status(
@router.post("/generate")
async def trigger_generation(
request: Mainx30GenerationRequest,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user)
):
tenant_id = current_user.get("tenant_id")
# Pass request as dict to Celery task
task = generar_transmission_file_async.delay(request.model_dump(), tenant_id)
tenant_id = get_tenant_from_token(current_user) or current_user.get("tenant_id")
task = track_and_dispatch(
db=db,
task=generar_transmission_file_async,
tenant_id=int(tenant_id),
company_id=getattr(request, "company_id", None),
requested_by_user=current_user.get("preferred_username") or current_user.get("email") or current_user.get("sub"),
task_name="generar_transmission_file_async",
task_group="reports",
task_origin="a76/reports/exportacion/transmission/MAINX30",
args=[request.model_dump(), tenant_id],
)
return {"task_id": task.id, "message": "Generación iniciada"}

View File

@@ -6,6 +6,7 @@ from celery.result import AsyncResult
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 .mex.service import ConsolidadoImportacionMexService
from .task import generar_pdf_consolidado_async
@@ -43,6 +44,16 @@ async def trigger_descarga_consolidado(
current_user: Dict[str, Any] = Depends(get_current_user),
db: Session = Depends(get_core_db)
):
validate_access_to_resource(db, company_id, current_user)
task = generar_pdf_consolidado_async.delay(invoice_id, company_id)
tenant_id = validate_access_to_resource(db, company_id, current_user)
task = track_and_dispatch(
db=db,
task=generar_pdf_consolidado_async,
tenant_id=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="generar_pdf_consolidado_async",
task_group="reports",
task_origin="a76/reports/importacion/consolidados",
args=[invoice_id, company_id],
)
return {"task_id": task.id, "message": "Generación iniciada"}

View File

@@ -6,6 +6,7 @@ from celery.result import AsyncResult
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 .mex.service import FacturaImportacionMexService
from .task import generar_pdf_factura_async
@@ -45,6 +46,16 @@ async def trigger_descarga_factura(
current_user: Dict[str, Any] = Depends(get_current_user),
db: Session = Depends(get_core_db)
):
validate_access_to_resource(db, company_id, current_user)
task = generar_pdf_factura_async.delay(invoice_id, company_id, invoice_type, currency_code)
tenant_id = validate_access_to_resource(db, company_id, current_user)
task = track_and_dispatch(
db=db,
task=generar_pdf_factura_async,
tenant_id=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="generar_pdf_factura_async",
task_group="reports",
task_origin="a76/reports/importacion/facturas",
args=[invoice_id, company_id, invoice_type, currency_code],
)
return {"task_id": task.id, "message": "Generación iniciada"}

View File

@@ -3,6 +3,7 @@ from fastapi import APIRouter, Depends, Query, Response, HTTPException
from sqlalchemy.orm import Session
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 .service import PackingListService
router = APIRouter()
@@ -43,7 +44,16 @@ async def trigger_download_packing_list(
current_user: Dict[str, Any] = Depends(get_current_user),
db: Session = Depends(get_core_db)
):
validate_access_to_resource(db, company_id, current_user)
task = generar_packing_list_async.delay(invoice_id, company_id)
tenant_id = validate_access_to_resource(db, company_id, current_user)
task = track_and_dispatch(
db=db,
task=generar_packing_list_async,
tenant_id=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="generar_packing_list_async",
task_group="reports",
task_origin="a76/reports/importacion/packing_list",
args=[invoice_id, company_id],
)
return {"task_id": task.id, "message": "Generación iniciada"}

View File

@@ -1,8 +1,11 @@
from typing import Dict, Any
from fastapi import APIRouter, Depends, Body
from sqlalchemy.orm import Session
from celery.result import AsyncResult
from core.celery_app import celery_app
from core.security import get_current_user
from core.database import get_core_db
from core.security import get_current_user, get_tenant_from_token
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .task import generar_transmission_definitiva_async
from .schemas import Mainx30GenerationRequest
@@ -34,8 +37,19 @@ async def get_task_status(
@router.post("/generate")
async def trigger_generation(
request: Mainx30GenerationRequest,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user)
):
tenant_id = current_user.get("tenant_id")
task = generar_transmission_definitiva_async.delay(request.model_dump(), tenant_id)
tenant_id = get_tenant_from_token(current_user) or current_user.get("tenant_id")
task = track_and_dispatch(
db=db,
task=generar_transmission_definitiva_async,
tenant_id=int(tenant_id),
company_id=getattr(request, "company_id", None),
requested_by_user=current_user.get("preferred_username") or current_user.get("email") or current_user.get("sub"),
task_name="generar_transmission_definitiva_async",
task_group="reports",
task_origin="a76/reports/importacion/transmission/definitive/MAINX30",
args=[request.model_dump(), tenant_id],
)
return {"task_id": task.id, "message": "Generación Definitiva iniciada"}

View File

@@ -1,8 +1,11 @@
from typing import Dict, Any
from fastapi import APIRouter, Depends, Body
from sqlalchemy.orm import Session
from celery.result import AsyncResult
from core.celery_app import celery_app
from core.security import get_current_user
from core.database import get_core_db
from core.security import get_current_user, get_tenant_from_token
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .task import generar_transmission_temporal_async
from .schemas import Mainx30GenerationRequest
@@ -35,9 +38,19 @@ async def get_task_status(
@router.post("/generate")
async def trigger_generation(
request: Mainx30GenerationRequest,
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user)
):
tenant_id = current_user.get("tenant_id")
# Pass request as dict to Celery task
task = generar_transmission_temporal_async.delay(request.model_dump(), tenant_id)
tenant_id = get_tenant_from_token(current_user) or current_user.get("tenant_id")
task = track_and_dispatch(
db=db,
task=generar_transmission_temporal_async,
tenant_id=int(tenant_id),
company_id=getattr(request, "company_id", None),
requested_by_user=current_user.get("preferred_username") or current_user.get("email") or current_user.get("sub"),
task_name="generar_transmission_temporal_async",
task_group="reports",
task_origin="a76/reports/importacion/transmission/temporal/MAINX30",
args=[request.model_dump(), tenant_id],
)
return {"task_id": task.id, "message": "Generación Temporal iniciada"}

View File

@@ -7,6 +7,8 @@ from core.celery_app import celery_app
from celery.result import AsyncResult
from typing import Dict, Any
import logging
from core.security import get_current_user, get_tenant_from_token
from api.v1.modules.core.tasks_tracking import track_and_dispatch
logger = logging.getLogger(__name__)
@@ -15,9 +17,20 @@ router = APIRouter(tags=["WINSAAI"])
@router.post("/generate", response_model=WinsaaiResponse)
async def trigger_winsaai_generation(
request: WinsaaiGenerationRequest,
db: Session = Depends(get_core_db)
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
task = generate_winsaai_task.delay(request.invoice_ids, request.is_temporal)
tenant_id = get_tenant_from_token(current_user) or current_user.get("tenant_id")
task = track_and_dispatch(
db=db,
task=generate_winsaai_task,
tenant_id=int(tenant_id),
requested_by_user=current_user.get("preferred_username") or current_user.get("email") or current_user.get("sub"),
task_name="generate_winsaai_task",
task_group="reports",
task_origin="a76/reports/importacion/winsaai/invoices",
args=[request.invoice_ids, request.is_temporal],
)
return WinsaaiResponse(
task_id=task.id,
status="PENDING",

View File

@@ -7,6 +7,8 @@ from core.celery_app import celery_app
from celery.result import AsyncResult
from typing import Dict, Any
import logging
from core.security import get_current_user, get_tenant_from_token
from api.v1.modules.core.tasks_tracking import track_and_dispatch
logger = logging.getLogger(__name__)
@@ -15,12 +17,19 @@ router = APIRouter(tags=["WINSAAI-Pedimentos"])
@router.post("/generate", response_model=WinsaaiResponse)
async def trigger_pedimentos_generation(
request: PedimentosWinsaaiGenerationRequest,
db: Session = Depends(get_core_db)
db: Session = Depends(get_core_db),
current_user: Dict[str, Any] = Depends(get_current_user),
):
task = generate_pedimentos_winsaai_task.delay(
request.pedimento_ids,
request.is_temporal,
request.is_by_class
tenant_id = get_tenant_from_token(current_user) or current_user.get("tenant_id")
task = track_and_dispatch(
db=db,
task=generate_pedimentos_winsaai_task,
tenant_id=int(tenant_id),
requested_by_user=current_user.get("preferred_username") or current_user.get("email") or current_user.get("sub"),
task_name="generate_pedimentos_winsaai_task",
task_group="reports",
task_origin="a76/reports/importacion/winsaai/pedimentos",
args=[request.pedimento_ids, request.is_temporal, request.is_by_class],
)
return WinsaaiResponse(
task_id=task.id,

View File

@@ -4,7 +4,8 @@ from sqlalchemy.orm import Session
from typing import List, Union
from core.database import get_core_db
from core.security import get_current_user
from core.security import get_current_user, get_tenant_from_token
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .schemas import (
ImportTemporaryFilter,
ImportDefinitiveFilter,
@@ -700,6 +701,7 @@ async def get_all_movements(
)
def generate_invoice_report_async(
filters: AllMovementsFilter,
db: Session = Depends(get_core_db),
current_user: dict = Depends(get_current_user)
):
"""
@@ -715,7 +717,25 @@ def generate_invoice_report_async(
user_email = current_user.get('email')
# Trigger task
task = generate_invoice_movements_async.delay(filter_data, user_email)
tenant_id = get_tenant_from_token(current_user)
if not tenant_id:
tenant_id = current_user.get("tenant_id")
if not tenant_id:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="Tenant ID not found in token",
)
task = track_and_dispatch(
db=db,
task=generate_invoice_movements_async,
tenant_id=int(tenant_id),
company_id=filters.company_id,
requested_by_user=current_user.get("preferred_username") or current_user.get("email") or current_user.get("sub"),
task_name="generate_invoice_movements_async",
task_group="reports",
task_origin="a76/reports/movements/invoices",
args=[filter_data, user_email],
)
return {"task_id": task.id}

View File

@@ -11,6 +11,7 @@ from sqlalchemy.orm import Session
from core.database import get_core_db
from core.security import get_current_user
from api.v1.modules.core.tasks_tracking import track_and_dispatch
from .schemas import SaldosFilter
logger = logging.getLogger(__name__)
@@ -51,7 +52,17 @@ def generate_saldos_report_async(
filter_data = filters.model_dump()
user_email = current_user.get("email")
task = generate_saldos_temporales_async.delay(filter_data, user_email)
task = track_and_dispatch(
db=db,
task=generate_saldos_temporales_async,
tenant_id=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="generate_saldos_temporales_async",
task_group="reports",
task_origin="a76/reports/movements/saldos",
args=[filter_data, user_email],
)
return {"task_id": task.id}

View File

@@ -6,6 +6,7 @@ from .user_tenant.routes import router as user_tenant_router
from .users.routes import router as users_router
from .dashboard.routes import router as dashboard_router
from .help_center.routes import router as help_center_router
from .tasks_tracking.routes import router as tasks_tracking_router
from fastapi import APIRouter
router = APIRouter()
@@ -18,3 +19,4 @@ router.include_router(licenses_router, prefix="/core", tags=["core / licenses"])
router.include_router(permissions_router, prefix="/core", tags=["core / permissions"])
router.include_router(dashboard_router, prefix="/core", tags=["core / dashboard"])
router.include_router(help_center_router, prefix="/core", tags=["core / help-center"])
router.include_router(tasks_tracking_router, prefix="/core", tags=["core / tasks"])

View File

@@ -0,0 +1,4 @@
from .dispatch import track_and_dispatch
from .service import TaskTrackerService
__all__ = ["track_and_dispatch", "TaskTrackerService"]

View File

@@ -0,0 +1,36 @@
from typing import Any
from celery import Task
from sqlalchemy.orm import Session
from .service import TaskTrackerService
def track_and_dispatch(
*,
db: Session,
task: Task,
tenant_id: int,
task_name: str,
task_group: str,
company_id: int | None = None,
requested_by_user: str | None = None,
task_origin: str | None = None,
args: list[Any] | None = None,
kwargs: dict[str, Any] | None = None,
task_id: str | None = None,
meta_payload: dict[str, Any] | None = None,
):
celery_task = task.apply_async(args=args or [], kwargs=kwargs or {}, task_id=task_id)
tracker = TaskTrackerService(db)
tracker.register_dispatch(
task_id=celery_task.id,
tenant_id=tenant_id,
company_id=company_id,
requested_by_user=requested_by_user,
task_name=task_name,
task_group=task_group,
task_origin=task_origin,
meta_payload=meta_payload,
)
return celery_task

View File

@@ -0,0 +1,62 @@
import enum
from datetime import datetime
from sqlalchemy import DateTime, ForeignKey, Index, Integer, String, Text, func
from sqlalchemy.dialects.postgresql import JSONB
from sqlalchemy.orm import Mapped, mapped_column
from core.database import Base
class TaskStatus(str, enum.Enum):
PENDING = "pending"
ACTIVE = "active"
COMPLETED = "completed"
FAILED = "failed"
class TaskRun(Base):
__tablename__ = "task_runs"
__table_args__ = (
Index("ix_core_task_runs_task_id", "task_id", unique=True),
Index("ix_core_task_runs_tenant_updated", "tenant_id", "updated_at"),
Index("ix_core_task_runs_tenant_status_updated", "tenant_id", "status", "updated_at"),
Index("ix_core_task_runs_tenant_group_updated", "tenant_id", "task_group", "updated_at"),
Index("ix_core_task_runs_tenant_company_updated", "tenant_id", "company_id", "updated_at"),
{"schema": "core"},
)
id: Mapped[int] = mapped_column(Integer, primary_key=True)
task_id: Mapped[str] = mapped_column(String(255), nullable=False)
tenant_id: Mapped[int] = mapped_column(
Integer, ForeignKey("core.tenants.id"), nullable=False, index=True
)
company_id: Mapped[int | None] = mapped_column(
Integer, ForeignKey("a76.company.id"), nullable=True, index=True
)
requested_by_user: Mapped[str | None] = mapped_column(String(255), nullable=True)
task_name: Mapped[str] = mapped_column(String(255), nullable=False)
task_group: Mapped[str] = mapped_column(String(100), nullable=False)
task_origin: Mapped[str | None] = mapped_column(String(255), nullable=True)
status: Mapped[TaskStatus] = mapped_column(
String(20), nullable=False, default=TaskStatus.PENDING.value
)
celery_state_raw: Mapped[str] = mapped_column(String(30), nullable=False, default="PENDING")
progress_current: Mapped[int | None] = mapped_column(Integer, nullable=True)
progress_total: Mapped[int | None] = mapped_column(Integer, nullable=True)
progress_percent: Mapped[float | None] = mapped_column(nullable=True)
progress_message: Mapped[str | None] = mapped_column(String(500), nullable=True)
retries: Mapped[int | None] = mapped_column(Integer, nullable=True)
exception_type: Mapped[str | None] = mapped_column(String(255), nullable=True)
exception_message: Mapped[str | None] = mapped_column(Text, nullable=True)
traceback_excerpt: Mapped[str | None] = mapped_column(Text, nullable=True)
result_summary: Mapped[dict | None] = mapped_column(JSONB, nullable=True)
meta_payload: Mapped[dict | None] = mapped_column(JSONB, nullable=True)
started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now()
)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now(), onupdate=func.now()
)

View File

@@ -0,0 +1,150 @@
from typing import Any
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy.orm import Session
from core.database import get_core_db
from core.security import get_current_user, get_tenant_from_token
from .models import TaskRun, TaskStatus
from .schemas import TaskCatalogsResponse, TaskRunDetail, TaskRunListItem, TaskRunsResponse, TaskSyncRequest
from .service import TaskTrackerService
router = APIRouter()
def _map_row(row: TaskRun) -> TaskRunListItem:
pct = row.progress_percent
if row.status == TaskStatus.COMPLETED.value and (pct is None or pct == 0):
pct = 100.0
return TaskRunListItem(
task_id=row.task_id,
task_name=row.task_name,
task_group=row.task_group,
task_origin=row.task_origin,
status=row.status,
celery_state_raw=row.celery_state_raw,
progress={
"current": row.progress_current,
"total": row.progress_total,
"percent": pct,
"message": row.progress_message,
},
retries=row.retries,
error=(
{"type": row.exception_type, "message": row.exception_message}
if row.exception_type or row.exception_message
else None
),
tenant_id=row.tenant_id,
company_id=row.company_id,
requested_by_user=row.requested_by_user,
started_at=row.started_at,
finished_at=row.finished_at,
created_at=row.created_at,
updated_at=row.updated_at,
)
@router.get("/tasks", response_model=TaskRunsResponse)
def list_tasks(
page: int = Query(1, ge=1),
page_size: int = Query(50, ge=1, le=200),
status: list[str] | None = Query(None),
task_group: list[str] | None = Query(None),
task_name: list[str] | None = Query(None),
company_id: int | None = Query(None),
search: str | None = Query(None),
sync_active: bool = Query(False),
current_user: dict[str, Any] = Depends(get_current_user),
db: Session = Depends(get_core_db),
):
tenant_id = get_tenant_from_token(current_user)
if not tenant_id:
raise HTTPException(status_code=400, detail="Tenant ID not found in token")
tracker = TaskTrackerService(db)
if sync_active:
tracker.sync_active_tasks(tenant_id=tenant_id)
rows, total = tracker.list_tasks(
tenant_id=tenant_id,
page=page,
page_size=page_size,
status=status,
task_group=task_group,
task_name=task_name,
company_id=company_id,
search=search,
)
items = [_map_row(row) for row in rows]
return TaskRunsResponse(
items=items,
total=total,
page=page,
page_size=page_size,
has_next=(page * page_size) < total,
)
@router.get("/tasks/{task_id}", response_model=TaskRunDetail)
def get_task_detail(
task_id: str,
sync: bool = Query(True),
current_user: dict[str, Any] = Depends(get_current_user),
db: Session = Depends(get_core_db),
):
tenant_id = get_tenant_from_token(current_user)
if not tenant_id:
raise HTTPException(status_code=400, detail="Tenant ID not found in token")
row = db.query(TaskRun).filter(TaskRun.task_id == task_id, TaskRun.tenant_id == tenant_id).first()
if not row:
raise HTTPException(status_code=404, detail="Task not found")
tracker = TaskTrackerService(db)
if sync:
row = tracker.sync_task(row)
item = _map_row(row)
return TaskRunDetail(
**item.model_dump(),
traceback_excerpt=row.traceback_excerpt,
result_summary=row.result_summary,
meta_payload=row.meta_payload,
)
@router.post("/tasks/sync")
def sync_tasks(
body: TaskSyncRequest,
current_user: dict[str, Any] = Depends(get_current_user),
db: Session = Depends(get_core_db),
):
tenant_id = get_tenant_from_token(current_user)
if not tenant_id:
raise HTTPException(status_code=400, detail="Tenant ID not found in token")
tracker = TaskTrackerService(db)
updated = tracker.sync_active_tasks(tenant_id=tenant_id, task_ids=body.task_ids)
return {"updated": updated}
@router.get("/tasks/catalogs", response_model=TaskCatalogsResponse)
def get_catalogs(
current_user: dict[str, Any] = Depends(get_current_user),
db: Session = Depends(get_core_db),
):
tenant_id = get_tenant_from_token(current_user)
if not tenant_id:
raise HTTPException(status_code=400, detail="Tenant ID not found in token")
groups = (
db.query(TaskRun.task_group).filter(TaskRun.tenant_id == tenant_id).distinct().order_by(TaskRun.task_group).all()
)
names = db.query(TaskRun.task_name).filter(TaskRun.tenant_id == tenant_id).distinct().order_by(TaskRun.task_name).all()
statuses = db.query(TaskRun.status).filter(TaskRun.tenant_id == tenant_id).distinct().order_by(TaskRun.status).all()
return TaskCatalogsResponse(
task_groups=[g[0] for g in groups if g[0]],
task_names=[n[0] for n in names if n[0]],
statuses=[s[0] for s in statuses if s[0]],
)

View File

@@ -0,0 +1,59 @@
from datetime import datetime
from typing import Any
from pydantic import BaseModel
class TaskProgress(BaseModel):
current: int | None = None
total: int | None = None
percent: float | None = None
message: str | None = None
class TaskError(BaseModel):
type: str | None = None
message: str | None = None
class TaskRunListItem(BaseModel):
task_id: str
task_name: str
task_group: str
task_origin: str | None = None
status: str
celery_state_raw: str
progress: TaskProgress
retries: int | None = None
error: TaskError | None = None
tenant_id: int
company_id: int | None = None
requested_by_user: str | None = None
started_at: datetime | None = None
finished_at: datetime | None = None
created_at: datetime
updated_at: datetime
class TaskRunDetail(TaskRunListItem):
traceback_excerpt: str | None = None
result_summary: dict[str, Any] | None = None
meta_payload: dict[str, Any] | None = None
class TaskRunsResponse(BaseModel):
items: list[TaskRunListItem]
total: int
page: int
page_size: int
has_next: bool
class TaskSyncRequest(BaseModel):
task_ids: list[str] | None = None
class TaskCatalogsResponse(BaseModel):
task_groups: list[str]
task_names: list[str]
statuses: list[str]

View File

@@ -0,0 +1,242 @@
from datetime import datetime, timezone
from typing import Any
from celery.result import AsyncResult
from sqlalchemy import asc, desc, func, or_
from sqlalchemy.orm import Session
from core.celery_app import celery_app
from .models import TaskRun, TaskStatus
# No usar valid_rows/total_rows del resultado final: en SUCCESS distorsiona el % (errores de scan).
_CELERY_TERMINAL_RAW = frozenset({"SUCCESS", "FAILURE", "REVOKED", "REJECTED"})
def _progress_from_row_count_dict(d: dict[str, Any]) -> tuple[int, int] | None:
total = d.get("total_rows")
if not isinstance(total, (int, float)) or total <= 0:
return None
valid = d.get("valid_rows")
if isinstance(valid, (int, float)):
return int(valid), int(total)
processed = d.get("processed_rows")
if isinstance(processed, (int, float)):
return int(processed), int(total)
return None
def normalize_celery_state(state: str | None) -> TaskStatus:
raw = (state or "PENDING").upper()
if raw in {"STARTED", "PROGRESS", "PROCESSING", "RETRY"}:
return TaskStatus.ACTIVE
if raw == "SUCCESS":
return TaskStatus.COMPLETED
if raw in {"FAILURE", "REVOKED", "REJECTED"}:
return TaskStatus.FAILED
return TaskStatus.PENDING
def _extract_progress(result: AsyncResult, raw_state_upper: str) -> tuple[int | None, int | None, float | None, str | None]:
payload = result.info if isinstance(result.info, dict) else {}
current = payload.get("current")
total = payload.get("total")
message = payload.get("status")
if current is None and isinstance(result.result, dict):
current = result.result.get("current")
if total is None and isinstance(result.result, dict):
total = result.result.get("total")
if message is None and isinstance(result.result, dict):
message = result.result.get("status")
percent = None
if isinstance(current, (int, float)) and isinstance(total, (int, float)) and total > 0:
percent = min(100.0, max(0.0, (float(current) / float(total)) * 100.0))
raw = (raw_state_upper or "PENDING").upper()
if percent is None and raw not in _CELERY_TERMINAL_RAW:
for src in (payload, result.result if isinstance(result.result, dict) else {}):
if not isinstance(src, dict):
continue
pair = _progress_from_row_count_dict(src)
if pair is None:
continue
cur_i, tot_i = pair
current, total = cur_i, tot_i
percent = min(100.0, max(0.0, (float(cur_i) / float(tot_i)) * 100.0))
break
return (
int(current) if isinstance(current, (int, float)) else None,
int(total) if isinstance(total, (int, float)) else None,
percent,
str(message) if message is not None else None,
)
def _extract_failure(result: AsyncResult) -> tuple[str | None, str | None, str | None]:
exception_type = None
exception_message = None
traceback_excerpt = None
err = result.result
if isinstance(err, Exception):
exception_type = type(err).__name__
exception_message = str(err)
elif isinstance(err, dict):
exception_type = err.get("exc_type")
exception_message = err.get("exc_message") or err.get("error") or err.get("message")
elif err is not None:
exception_message = str(err)
tb = getattr(result, "traceback", None)
if isinstance(tb, str):
traceback_excerpt = tb[-4000:]
return exception_type, exception_message, traceback_excerpt
class TaskTrackerService:
def __init__(self, db: Session):
self.db = db
def register_dispatch(
self,
*,
task_id: str,
tenant_id: int,
task_name: str,
task_group: str,
company_id: int | None = None,
requested_by_user: str | None = None,
task_origin: str | None = None,
meta_payload: dict[str, Any] | None = None,
) -> TaskRun:
current = self.db.query(TaskRun).filter(TaskRun.task_id == task_id).first()
if current:
return current
row = TaskRun(
task_id=task_id,
tenant_id=tenant_id,
company_id=company_id,
requested_by_user=requested_by_user,
task_name=task_name,
task_group=task_group,
task_origin=task_origin,
status=TaskStatus.PENDING.value,
celery_state_raw="PENDING",
progress_current=0,
progress_total=100,
progress_percent=0.0,
progress_message="Queued",
retries=0,
meta_payload=meta_payload,
started_at=None,
finished_at=None,
)
self.db.add(row)
self.db.commit()
self.db.refresh(row)
return row
def sync_task(self, task_run: TaskRun) -> TaskRun:
async_result = celery_app.AsyncResult(task_run.task_id)
raw_state = (async_result.state or "PENDING").upper()
normalized = normalize_celery_state(raw_state)
now = datetime.now(timezone.utc)
task_run.celery_state_raw = raw_state
task_run.status = normalized.value
task_run.retries = int(getattr(async_result, "retries", 0) or 0)
current, total, percent, message = _extract_progress(async_result, raw_state)
if current is not None:
task_run.progress_current = current
if total is not None:
task_run.progress_total = total
if percent is not None:
task_run.progress_percent = percent
if message:
task_run.progress_message = message
if normalized == TaskStatus.ACTIVE and task_run.started_at is None:
task_run.started_at = now
if normalized == TaskStatus.COMPLETED:
if task_run.started_at is None:
task_run.started_at = now
task_run.finished_at = now
task_run.exception_type = None
task_run.exception_message = None
task_run.traceback_excerpt = None
if isinstance(async_result.result, dict):
task_run.result_summary = async_result.result
else:
task_run.result_summary = {"result": str(async_result.result)}
# register_dispatch seeds progress_percent=0; Celery SUCCESS often has no current/total meta
if percent is None:
task_run.progress_percent = 100.0
if normalized == TaskStatus.FAILED:
if task_run.started_at is None:
task_run.started_at = now
task_run.finished_at = now
etype, emsg, tb = _extract_failure(async_result)
task_run.exception_type = etype
task_run.exception_message = emsg
task_run.traceback_excerpt = tb
self.db.add(task_run)
self.db.commit()
self.db.refresh(task_run)
return task_run
def sync_active_tasks(self, tenant_id: int, task_ids: list[str] | None = None) -> int:
query = self.db.query(TaskRun).filter(
TaskRun.tenant_id == tenant_id, TaskRun.status.in_([TaskStatus.PENDING.value, TaskStatus.ACTIVE.value])
)
if task_ids:
query = query.filter(TaskRun.task_id.in_(task_ids))
rows = query.limit(200).all()
for row in rows:
self.sync_task(row)
return len(rows)
def list_tasks(
self,
*,
tenant_id: int,
page: int,
page_size: int,
status: list[str] | None = None,
task_group: list[str] | None = None,
task_name: list[str] | None = None,
company_id: int | None = None,
search: str | None = None,
order: str = "desc",
) -> tuple[list[TaskRun], int]:
query = self.db.query(TaskRun).filter(TaskRun.tenant_id == tenant_id)
if status:
query = query.filter(TaskRun.status.in_(status))
if task_group:
query = query.filter(TaskRun.task_group.in_(task_group))
if task_name:
query = query.filter(TaskRun.task_name.in_(task_name))
if company_id is not None:
# Tareas sin company_id (p. ej. reportes solo por tenant) deben seguir visibles
query = query.filter(or_(TaskRun.company_id == company_id, TaskRun.company_id.is_(None)))
if search:
term = f"%{search}%"
query = query.filter(
(TaskRun.task_id.ilike(term))
| (TaskRun.task_name.ilike(term))
| (TaskRun.task_origin.ilike(term))
| (TaskRun.exception_message.ilike(term))
)
total = query.with_entities(func.count(TaskRun.id)).scalar() or 0
order_expr = asc(TaskRun.updated_at) if order == "asc" else desc(TaskRun.updated_at)
items = query.order_by(order_expr).offset((page - 1) * page_size).limit(page_size).all()
return items, total

View File

@@ -1,4 +1,3 @@
# TODO: Revisar las claves de pedimento que no estan en la tabla de pedimento_codes
seed = [
("A4", "DFI", "E"),
("A4", "DFI", "I"),
@@ -33,7 +32,7 @@ seed = [
("H1", "EXD", "E"),
("H8", "EXD", "E"),
("I1", "EXD", "E"),
# ("J1", "EXD", "E"), # No existe en la tabla de pedimento_codes
("J1", "EXD", "E"),
("J2", "EXD", "E"),
("K1", "EXD", "E"),
("K2", "EXD", "E"),

View File

@@ -1,5 +1,9 @@
seed = [
("A1", "IMPORTACION O EXPORTACION DEFINITIVA."),
(
"A2",
"IMPORTACION TEMPORAL DE BIENES DISTINTOS A LOS DE ACTIVO FIJO POR PARTE DE EMPRESAS CON PITEX.",
),
("A3", "REGULARIZACION DE MERCANCIAS (IMPORTACION DEFINITIVA)."),
("A4", "INTRODUCCION PARA DEPOSITO FISCAL (AGD)."),
("A5", "INTRODUCCION A DEPOSITO FISCAL EN LOCAL AUTORIZADO."),
@@ -7,6 +11,22 @@ seed = [
"A6",
"IMPORTACIÓN TEMPORAL DE BIENES DE ACTIVO FIJO POR PARTE DE EMPRESAS CON PITEX.",
),
(
"A7",
"IMPORTACION TEMPORAL DE MERCANCIAS PARA RETORNO Y REEXPORTACION.",
),
(
"A8",
"IMPORTACION TEMPORAL DE MERCANCIAS E INSUMOS ASOCIADOS A PROCESOS INDUSTRIALES.",
),
(
"A9",
"IMPORTACION TEMPORAL DE MERCANCIAS BAJO ESQUEMAS DE CONTROL ADUANERO.",
),
(
"AA",
"IMPORTACION TEMPORAL DE MERCANCIAS EN OPERACIONES DIVERSAS.",
),
(
"AD",
"IMPORTACIÓN TEMPORAL DE MERCANCIAS DESTINADAS A CONVENCIONES Y CONGRESOS INTERNACIONALES (ARTICULO 106, FRACCION III, INCISO A) DE LA LEY).",
@@ -59,6 +79,10 @@ seed = [
"C1",
"IMPORTACION DEFINITIVA A LA FRANJA FRONTERIZA NORTE Y REGION FRONTERIZA AL AMPARO DEL „DECRETO DE LA FRANJA O REGION FRONTERIZA“ (DOF 24/12/2008 Y SUS POSTERIORES MODIFICACIONES).",
),
(
"C2",
"IMPORTACION DEFINITIVA AL INTERIOR DEL PAIS DE MERCANCIAS PROCEDENTES DE LA FRANJA FRONTERIZA NORTE Y REGION FRONTERIZA.",
),
("C3", "EXTRACCION DE DEPOSITO FISCAL DE FRANJA O REGION FRONTERIZA (AGD)."),
("CT", "PEDIMENTO COMPLEMENTARIO."),
("D1", "RETORNO POR SUSTITUCION."),
@@ -101,6 +125,10 @@ seed = [
),
("GC", "GLOBAL COMPLEMENTARIO."),
("H1", "RETORNO DE MERCANCIAS EN SU MISMO ESTADO."),
(
"H3",
"IMPORTACION TEMPORAL DE ACTIVO FIJO POR PARTE DE MAQUILADORAS.",
),
("H8", "RETORNO DE ENVASES."),
(
"I1",
@@ -110,6 +138,14 @@ seed = [
"IN",
"IMPORTACION TEMPORAL DE BIENES QUE SERAN SUJETOS A TRANSFORMACION, ELABORACION O REPARACION (IMMEX).",
),
(
"J1",
"EXPORTACION DEFINITIVA DE INSUMOS ELABORADOS O TRANSFORMADOS EN RECINTO FISCALIZADO.",
),
(
"J2",
"EXPORTACION DEFINITIVA DE MAQUINARIA Y EQUIPO ELABORADO O TRANSFORMADO EN RECINTO FISCALIZADO.",
),
(
"J3",
"RETORNO Y EXPORTACION DE INSUMOS ELABORADOS O TRANSFORMADOS EN RECINTO FISCALIZADO.",