From b5b99a5ddb4cad09b09ab33ba4cc676c3e5c35b5 Mon Sep 17 00:00:00 2001 From: hreyes Date: Tue, 24 Mar 2026 13:07:35 -0600 Subject: [PATCH 1/2] feature/tabla-bitacora-logs-task --- .../c1a2b3d4e5f6_create_core_task_runs.py | 85 ++++ ...9_backfill_full_code_pedimento_regimens.py | 82 ++++ .../a76/invoices/exports/process/routes.py | 13 +- .../a76/invoices/exports/revert/routes.py | 13 +- .../a76/invoices/imports/process/routes.py | 25 +- .../a76/invoices/imports/revert/routes.py | 25 +- .../v1/modules/a76/layouts_csv/boms/routes.py | 35 +- .../cambio_regimen_regularizacion/routes.py | 39 +- .../modules/a76/layouts_csv/classes/routes.py | 37 +- .../clients_and_providers/routes.py | 37 +- .../common/track_commit_dispatch.py | 58 +++ .../a76/layouts_csv/customs_brokers/routes.py | 37 +- .../a76/layouts_csv/exchange_rate/routes.py | 37 +- .../a76/layouts_csv/exportacion/routes.py | 37 +- .../a76/layouts_csv/facturas/routes.py | 38 +- .../modules/a76/layouts_csv/parts/routes.py | 35 +- .../a76/layouts_csv/pedmientos/routes.py | 38 +- .../a76/layouts_csv/trailers/routes.py | 16 +- .../a76/layouts_csv/transportistas/routes.py | 17 +- .../layouts_csv/us_tariff_fractions/routes.py | 37 +- .../a76/layouts_csv/vehicles/routes.py | 14 +- .../exportacion/aviso_consolidado/routes.py | 15 +- .../reports/exportacion/descargo/routes.py | 18 +- .../transmission/MAINX30/routes.py | 21 +- .../importacion/consolidados/routes.py | 15 +- .../reports/importacion/facturas/routes.py | 15 +- .../importacion/packing_list/routes.py | 16 +- .../transmission/definitive/MAINX30/routes.py | 20 +- .../transmission/temporal/MAINX30/routes.py | 21 +- .../importacion/winsaai/invoices/routes.py | 17 +- .../importacion/winsaai/pedimentos/routes.py | 19 +- .../a76/reports/movements/invoices/routes.py | 24 +- .../a76/reports/movements/saldos/routes.py | 13 +- backend/api/v1/modules/core/router.py | 2 + .../modules/core/tasks_tracking/__init__.py | 4 + .../modules/core/tasks_tracking/dispatch.py | 36 ++ .../v1/modules/core/tasks_tracking/models.py | 62 +++ .../v1/modules/core/tasks_tracking/routes.py | 150 +++++++ .../v1/modules/core/tasks_tracking/schemas.py | 59 +++ .../v1/modules/core/tasks_tracking/service.py | 242 +++++++++++ .../code_pedimento_regimens/seed.py | 3 +- .../reference_data/pedimento_codes/seed.py | 36 ++ frontend/messages/en.json | 4 +- frontend/messages/es.json | 4 +- frontend/src/lib/api/dashboard/a76/tasks.ts | 78 ++++ .../routes/dashboard/audit_logs/+page.svelte | 398 +++--------------- .../src/routes/dashboard/audit_logs/+page.ts | 7 + .../dashboard/audit_logs/bitacora-tab.svelte | 321 ++++++++++++++ .../dashboard/audit_logs/tasks-tab.svelte | 338 +++++++++++++++ frontend/src/routes/dashboard/tasks/+page.ts | 6 + 50 files changed, 2289 insertions(+), 430 deletions(-) create mode 100644 backend/alembic/versions/c1a2b3d4e5f6_create_core_task_runs.py create mode 100644 backend/alembic/versions/d4e5f6a7b8c9_backfill_full_code_pedimento_regimens.py create mode 100644 backend/api/v1/modules/a76/layouts_csv/common/track_commit_dispatch.py create mode 100644 backend/api/v1/modules/core/tasks_tracking/__init__.py create mode 100644 backend/api/v1/modules/core/tasks_tracking/dispatch.py create mode 100644 backend/api/v1/modules/core/tasks_tracking/models.py create mode 100644 backend/api/v1/modules/core/tasks_tracking/routes.py create mode 100644 backend/api/v1/modules/core/tasks_tracking/schemas.py create mode 100644 backend/api/v1/modules/core/tasks_tracking/service.py create mode 100644 frontend/src/lib/api/dashboard/a76/tasks.ts create mode 100644 frontend/src/routes/dashboard/audit_logs/+page.ts create mode 100644 frontend/src/routes/dashboard/audit_logs/bitacora-tab.svelte create mode 100644 frontend/src/routes/dashboard/audit_logs/tasks-tab.svelte create mode 100644 frontend/src/routes/dashboard/tasks/+page.ts diff --git a/backend/alembic/versions/c1a2b3d4e5f6_create_core_task_runs.py b/backend/alembic/versions/c1a2b3d4e5f6_create_core_task_runs.py new file mode 100644 index 00000000..dae5b1c5 --- /dev/null +++ b/backend/alembic/versions/c1a2b3d4e5f6_create_core_task_runs.py @@ -0,0 +1,85 @@ +"""create_core_task_runs + +Revision ID: c1a2b3d4e5f6 +Revises: bccb7f8986c7 +Create Date: 2026-03-24 11:30:00.000000 +""" + +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa +from sqlalchemy.dialects import postgresql + + +# revision identifiers, used by Alembic. +revision: str = "c1a2b3d4e5f6" +down_revision: Union[str, Sequence[str], None] = "bccb7f8986c7" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.create_table( + "task_runs", + sa.Column("id", sa.Integer(), nullable=False), + sa.Column("task_id", sa.String(length=255), nullable=False), + sa.Column("tenant_id", sa.Integer(), nullable=False), + sa.Column("company_id", sa.Integer(), nullable=True), + sa.Column("requested_by_user", sa.String(length=255), nullable=True), + sa.Column("task_name", sa.String(length=255), nullable=False), + sa.Column("task_group", sa.String(length=100), nullable=False), + sa.Column("task_origin", sa.String(length=255), nullable=True), + sa.Column("status", sa.String(length=20), nullable=False), + sa.Column("celery_state_raw", sa.String(length=30), nullable=False), + sa.Column("progress_current", sa.Integer(), nullable=True), + sa.Column("progress_total", sa.Integer(), nullable=True), + sa.Column("progress_percent", sa.Float(), nullable=True), + sa.Column("progress_message", sa.String(length=500), nullable=True), + sa.Column("retries", sa.Integer(), nullable=True), + sa.Column("exception_type", sa.String(length=255), nullable=True), + sa.Column("exception_message", sa.Text(), nullable=True), + sa.Column("traceback_excerpt", sa.Text(), nullable=True), + sa.Column("result_summary", postgresql.JSONB(astext_type=sa.Text()), nullable=True), + sa.Column("meta_payload", postgresql.JSONB(astext_type=sa.Text()), nullable=True), + sa.Column("started_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.text("now()"), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), server_default=sa.text("now()"), nullable=False), + sa.ForeignKeyConstraint(["company_id"], ["a76.company.id"]), + sa.ForeignKeyConstraint(["tenant_id"], ["core.tenants.id"]), + sa.PrimaryKeyConstraint("id"), + schema="core", + ) + op.create_index("ix_core_task_runs_task_id", "task_runs", ["task_id"], unique=True, schema="core") + op.create_index("ix_core_task_runs_tenant_updated", "task_runs", ["tenant_id", "updated_at"], unique=False, schema="core") + op.create_index( + "ix_core_task_runs_tenant_status_updated", + "task_runs", + ["tenant_id", "status", "updated_at"], + unique=False, + schema="core", + ) + op.create_index( + "ix_core_task_runs_tenant_group_updated", + "task_runs", + ["tenant_id", "task_group", "updated_at"], + unique=False, + schema="core", + ) + op.create_index( + "ix_core_task_runs_tenant_company_updated", + "task_runs", + ["tenant_id", "company_id", "updated_at"], + unique=False, + schema="core", + ) + + +def downgrade() -> None: + op.drop_index("ix_core_task_runs_tenant_company_updated", table_name="task_runs", schema="core") + op.drop_index("ix_core_task_runs_tenant_group_updated", table_name="task_runs", schema="core") + op.drop_index("ix_core_task_runs_tenant_status_updated", table_name="task_runs", schema="core") + op.drop_index("ix_core_task_runs_tenant_updated", table_name="task_runs", schema="core") + op.drop_index("ix_core_task_runs_task_id", table_name="task_runs", schema="core") + op.drop_table("task_runs", schema="core") diff --git a/backend/alembic/versions/d4e5f6a7b8c9_backfill_full_code_pedimento_regimens.py b/backend/alembic/versions/d4e5f6a7b8c9_backfill_full_code_pedimento_regimens.py new file mode 100644 index 00000000..c28f7ba5 --- /dev/null +++ b/backend/alembic/versions/d4e5f6a7b8c9_backfill_full_code_pedimento_regimens.py @@ -0,0 +1,82 @@ +"""backfill_full_code_pedimento_regimens + +Revision ID: d4e5f6a7b8c9 +Revises: c1a2b3d4e5f6 +Create Date: 2026-03-24 13:10:00.000000 +""" + +from typing import Sequence, Union + +from alembic import op +from api.v1.modules.public.reference_data.code_pedimento_regimens.seed import ( + seed as code_pedimento_regimens_seed_full, +) +from api.v1.modules.public.reference_data.pedimento_codes.seed import ( + seed as pedimento_codes_seed, +) + + +# revision identifiers, used by Alembic. +revision: str = "d4e5f6a7b8c9" +down_revision: Union[str, Sequence[str], None] = "c1a2b3d4e5f6" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def _format_value(val: str | None) -> str: + if val is None or str(val).strip() == "" or str(val).upper() == "NONE": + return "NULL" + return f"'{str(val).replace(chr(39), chr(39) * 2)}'" + + +def upgrade() -> None: + """Backfill full code_pedimento_regimens data without editing old migrations.""" + + pedimento_desc_by_code = {code: desc for code, desc in pedimento_codes_seed} + required_pedimento_codes = sorted( + {pedimento_code for pedimento_code, _regimen_code, _type_code in code_pedimento_regimens_seed_full} + ) + + values_missing_codes = ", ".join( + [ + f"({_format_value(code)}, {_format_value(pedimento_desc_by_code.get(code, f'AUTO-GENERATED FOR FK ({code})'))})" + for code in required_pedimento_codes + ] + ) + op.execute( + f""" + INSERT INTO public.pedimento_codes (code, description) + VALUES {values_missing_codes} + ON CONFLICT (code) DO NOTHING; + """ + ) + + values_relations = ", ".join( + [ + f"({_format_value(pedimento_code)}, {_format_value(regimen_code)}, {_format_value(type_code)})" + for pedimento_code, regimen_code, type_code in code_pedimento_regimens_seed_full + ] + ) + op.execute( + f""" + INSERT INTO public.code_pedimento_regimens (pedimento_code, regimen_code, type_code) + SELECT src.pedimento_code, src.regimen_code, src.type_code + FROM (VALUES {values_relations}) AS src(pedimento_code, regimen_code, type_code) + JOIN public.pedimento_codes pc + ON pc.code = src.pedimento_code + JOIN public.pedimento_regimens pr + ON pr.code = src.regimen_code + WHERE NOT EXISTS ( + SELECT 1 + FROM public.code_pedimento_regimens cpr + WHERE cpr.pedimento_code = src.pedimento_code + AND cpr.regimen_code = src.regimen_code + AND COALESCE(cpr.type_code, '') = COALESCE(src.type_code, '') + ); + """ + ) + + +def downgrade() -> None: + """Preserve seeded data on downgrade.""" + diff --git a/backend/api/v1/modules/a76/invoices/exports/process/routes.py b/backend/api/v1/modules/a76/invoices/exports/process/routes.py index 1f8a2297..2bd193bd 100644 --- a/backend/api/v1/modules/a76/invoices/exports/process/routes.py +++ b/backend/api/v1/modules/a76/invoices/exports/process/routes.py @@ -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} diff --git a/backend/api/v1/modules/a76/invoices/exports/revert/routes.py b/backend/api/v1/modules/a76/invoices/exports/revert/routes.py index 1419eb96..ad50c3a7 100644 --- a/backend/api/v1/modules/a76/invoices/exports/revert/routes.py +++ b/backend/api/v1/modules/a76/invoices/exports/revert/routes.py @@ -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} diff --git a/backend/api/v1/modules/a76/invoices/imports/process/routes.py b/backend/api/v1/modules/a76/invoices/imports/process/routes.py index 66a00121..675b7bea 100644 --- a/backend/api/v1/modules/a76/invoices/imports/process/routes.py +++ b/backend/api/v1/modules/a76/invoices/imports/process/routes.py @@ -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} diff --git a/backend/api/v1/modules/a76/invoices/imports/revert/routes.py b/backend/api/v1/modules/a76/invoices/imports/revert/routes.py index 8990f34e..a1e60c9e 100644 --- a/backend/api/v1/modules/a76/invoices/imports/revert/routes.py +++ b/backend/api/v1/modules/a76/invoices/imports/revert/routes.py @@ -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} diff --git a/backend/api/v1/modules/a76/layouts_csv/boms/routes.py b/backend/api/v1/modules/a76/layouts_csv/boms/routes.py index 38fb20ed..3d33f075 100644 --- a/backend/api/v1/modules/a76/layouts_csv/boms/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/boms/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/cambio_regimen_regularizacion/routes.py b/backend/api/v1/modules/a76/layouts_csv/cambio_regimen_regularizacion/routes.py index d84e32c9..e646e107 100644 --- a/backend/api/v1/modules/a76/layouts_csv/cambio_regimen_regularizacion/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/cambio_regimen_regularizacion/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/classes/routes.py b/backend/api/v1/modules/a76/layouts_csv/classes/routes.py index ff123ccf..24801833 100644 --- a/backend/api/v1/modules/a76/layouts_csv/classes/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/classes/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/clients_and_providers/routes.py b/backend/api/v1/modules/a76/layouts_csv/clients_and_providers/routes.py index f1d082d8..59c99ef7 100644 --- a/backend/api/v1/modules/a76/layouts_csv/clients_and_providers/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/clients_and_providers/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/common/track_commit_dispatch.py b/backend/api/v1/modules/a76/layouts_csv/common/track_commit_dispatch.py new file mode 100644 index 00000000..b8479a0a --- /dev/null +++ b/backend/api/v1/modules/a76/layouts_csv/common/track_commit_dispatch.py @@ -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 diff --git a/backend/api/v1/modules/a76/layouts_csv/customs_brokers/routes.py b/backend/api/v1/modules/a76/layouts_csv/customs_brokers/routes.py index 5c9c32b8..ab884351 100644 --- a/backend/api/v1/modules/a76/layouts_csv/customs_brokers/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/customs_brokers/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/exchange_rate/routes.py b/backend/api/v1/modules/a76/layouts_csv/exchange_rate/routes.py index 838887ae..ff55abf5 100644 --- a/backend/api/v1/modules/a76/layouts_csv/exchange_rate/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/exchange_rate/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/exportacion/routes.py b/backend/api/v1/modules/a76/layouts_csv/exportacion/routes.py index 953d2860..82c2af93 100644 --- a/backend/api/v1/modules/a76/layouts_csv/exportacion/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/exportacion/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/facturas/routes.py b/backend/api/v1/modules/a76/layouts_csv/facturas/routes.py index 488c5871..18b39eee 100644 --- a/backend/api/v1/modules/a76/layouts_csv/facturas/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/facturas/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/parts/routes.py b/backend/api/v1/modules/a76/layouts_csv/parts/routes.py index 71c2d418..d4a5af18 100644 --- a/backend/api/v1/modules/a76/layouts_csv/parts/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/parts/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/pedmientos/routes.py b/backend/api/v1/modules/a76/layouts_csv/pedmientos/routes.py index 5eeb8890..36f601fe 100644 --- a/backend/api/v1/modules/a76/layouts_csv/pedmientos/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/pedmientos/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/trailers/routes.py b/backend/api/v1/modules/a76/layouts_csv/trailers/routes.py index 71c80a26..f1543f78 100644 --- a/backend/api/v1/modules/a76/layouts_csv/trailers/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/trailers/routes.py @@ -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: diff --git a/backend/api/v1/modules/a76/layouts_csv/transportistas/routes.py b/backend/api/v1/modules/a76/layouts_csv/transportistas/routes.py index 7cf41b14..18671ddb 100644 --- a/backend/api/v1/modules/a76/layouts_csv/transportistas/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/transportistas/routes.py @@ -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) diff --git a/backend/api/v1/modules/a76/layouts_csv/us_tariff_fractions/routes.py b/backend/api/v1/modules/a76/layouts_csv/us_tariff_fractions/routes.py index 1d23acd1..322ca183 100644 --- a/backend/api/v1/modules/a76/layouts_csv/us_tariff_fractions/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/us_tariff_fractions/routes.py @@ -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, } diff --git a/backend/api/v1/modules/a76/layouts_csv/vehicles/routes.py b/backend/api/v1/modules/a76/layouts_csv/vehicles/routes.py index 3b13656c..da98d135 100644 --- a/backend/api/v1/modules/a76/layouts_csv/vehicles/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/vehicles/routes.py @@ -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: diff --git a/backend/api/v1/modules/a76/reports/exportacion/aviso_consolidado/routes.py b/backend/api/v1/modules/a76/reports/exportacion/aviso_consolidado/routes.py index 09cf878f..5b804b3a 100644 --- a/backend/api/v1/modules/a76/reports/exportacion/aviso_consolidado/routes.py +++ b/backend/api/v1/modules/a76/reports/exportacion/aviso_consolidado/routes.py @@ -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"} diff --git a/backend/api/v1/modules/a76/reports/exportacion/descargo/routes.py b/backend/api/v1/modules/a76/reports/exportacion/descargo/routes.py index 2a3878f3..8882ed5f 100644 --- a/backend/api/v1/modules/a76/reports/exportacion/descargo/routes.py +++ b/backend/api/v1/modules/a76/reports/exportacion/descargo/routes.py @@ -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)) diff --git a/backend/api/v1/modules/a76/reports/exportacion/transmission/MAINX30/routes.py b/backend/api/v1/modules/a76/reports/exportacion/transmission/MAINX30/routes.py index 76a916ed..e53fd95f 100644 --- a/backend/api/v1/modules/a76/reports/exportacion/transmission/MAINX30/routes.py +++ b/backend/api/v1/modules/a76/reports/exportacion/transmission/MAINX30/routes.py @@ -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"} diff --git a/backend/api/v1/modules/a76/reports/importacion/consolidados/routes.py b/backend/api/v1/modules/a76/reports/importacion/consolidados/routes.py index 5254cfe4..629805ea 100644 --- a/backend/api/v1/modules/a76/reports/importacion/consolidados/routes.py +++ b/backend/api/v1/modules/a76/reports/importacion/consolidados/routes.py @@ -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"} diff --git a/backend/api/v1/modules/a76/reports/importacion/facturas/routes.py b/backend/api/v1/modules/a76/reports/importacion/facturas/routes.py index 12595fda..b602b15c 100644 --- a/backend/api/v1/modules/a76/reports/importacion/facturas/routes.py +++ b/backend/api/v1/modules/a76/reports/importacion/facturas/routes.py @@ -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"} \ No newline at end of file diff --git a/backend/api/v1/modules/a76/reports/importacion/packing_list/routes.py b/backend/api/v1/modules/a76/reports/importacion/packing_list/routes.py index 54c4dacf..45869290 100644 --- a/backend/api/v1/modules/a76/reports/importacion/packing_list/routes.py +++ b/backend/api/v1/modules/a76/reports/importacion/packing_list/routes.py @@ -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"} diff --git a/backend/api/v1/modules/a76/reports/importacion/transmission/definitive/MAINX30/routes.py b/backend/api/v1/modules/a76/reports/importacion/transmission/definitive/MAINX30/routes.py index 833ab452..74f95712 100644 --- a/backend/api/v1/modules/a76/reports/importacion/transmission/definitive/MAINX30/routes.py +++ b/backend/api/v1/modules/a76/reports/importacion/transmission/definitive/MAINX30/routes.py @@ -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"} diff --git a/backend/api/v1/modules/a76/reports/importacion/transmission/temporal/MAINX30/routes.py b/backend/api/v1/modules/a76/reports/importacion/transmission/temporal/MAINX30/routes.py index e8dbcf51..a2453785 100644 --- a/backend/api/v1/modules/a76/reports/importacion/transmission/temporal/MAINX30/routes.py +++ b/backend/api/v1/modules/a76/reports/importacion/transmission/temporal/MAINX30/routes.py @@ -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"} diff --git a/backend/api/v1/modules/a76/reports/importacion/winsaai/invoices/routes.py b/backend/api/v1/modules/a76/reports/importacion/winsaai/invoices/routes.py index 5d2062ef..bd83d6d3 100644 --- a/backend/api/v1/modules/a76/reports/importacion/winsaai/invoices/routes.py +++ b/backend/api/v1/modules/a76/reports/importacion/winsaai/invoices/routes.py @@ -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", diff --git a/backend/api/v1/modules/a76/reports/importacion/winsaai/pedimentos/routes.py b/backend/api/v1/modules/a76/reports/importacion/winsaai/pedimentos/routes.py index 995177d2..4c274f1f 100644 --- a/backend/api/v1/modules/a76/reports/importacion/winsaai/pedimentos/routes.py +++ b/backend/api/v1/modules/a76/reports/importacion/winsaai/pedimentos/routes.py @@ -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, diff --git a/backend/api/v1/modules/a76/reports/movements/invoices/routes.py b/backend/api/v1/modules/a76/reports/movements/invoices/routes.py index 351e4d9e..9d71b705 100644 --- a/backend/api/v1/modules/a76/reports/movements/invoices/routes.py +++ b/backend/api/v1/modules/a76/reports/movements/invoices/routes.py @@ -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} diff --git a/backend/api/v1/modules/a76/reports/movements/saldos/routes.py b/backend/api/v1/modules/a76/reports/movements/saldos/routes.py index 2d7341cd..489480bc 100644 --- a/backend/api/v1/modules/a76/reports/movements/saldos/routes.py +++ b/backend/api/v1/modules/a76/reports/movements/saldos/routes.py @@ -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} diff --git a/backend/api/v1/modules/core/router.py b/backend/api/v1/modules/core/router.py index be36fde0..8e3361a6 100644 --- a/backend/api/v1/modules/core/router.py +++ b/backend/api/v1/modules/core/router.py @@ -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"]) diff --git a/backend/api/v1/modules/core/tasks_tracking/__init__.py b/backend/api/v1/modules/core/tasks_tracking/__init__.py new file mode 100644 index 00000000..3ffbd045 --- /dev/null +++ b/backend/api/v1/modules/core/tasks_tracking/__init__.py @@ -0,0 +1,4 @@ +from .dispatch import track_and_dispatch +from .service import TaskTrackerService + +__all__ = ["track_and_dispatch", "TaskTrackerService"] diff --git a/backend/api/v1/modules/core/tasks_tracking/dispatch.py b/backend/api/v1/modules/core/tasks_tracking/dispatch.py new file mode 100644 index 00000000..f37632b6 --- /dev/null +++ b/backend/api/v1/modules/core/tasks_tracking/dispatch.py @@ -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 diff --git a/backend/api/v1/modules/core/tasks_tracking/models.py b/backend/api/v1/modules/core/tasks_tracking/models.py new file mode 100644 index 00000000..dc86a33a --- /dev/null +++ b/backend/api/v1/modules/core/tasks_tracking/models.py @@ -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() + ) diff --git a/backend/api/v1/modules/core/tasks_tracking/routes.py b/backend/api/v1/modules/core/tasks_tracking/routes.py new file mode 100644 index 00000000..57395ea4 --- /dev/null +++ b/backend/api/v1/modules/core/tasks_tracking/routes.py @@ -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]], + ) diff --git a/backend/api/v1/modules/core/tasks_tracking/schemas.py b/backend/api/v1/modules/core/tasks_tracking/schemas.py new file mode 100644 index 00000000..9af50cb4 --- /dev/null +++ b/backend/api/v1/modules/core/tasks_tracking/schemas.py @@ -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] diff --git a/backend/api/v1/modules/core/tasks_tracking/service.py b/backend/api/v1/modules/core/tasks_tracking/service.py new file mode 100644 index 00000000..4d98672a --- /dev/null +++ b/backend/api/v1/modules/core/tasks_tracking/service.py @@ -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 diff --git a/backend/api/v1/modules/public/reference_data/code_pedimento_regimens/seed.py b/backend/api/v1/modules/public/reference_data/code_pedimento_regimens/seed.py index 76e5e1e9..83e2b16c 100644 --- a/backend/api/v1/modules/public/reference_data/code_pedimento_regimens/seed.py +++ b/backend/api/v1/modules/public/reference_data/code_pedimento_regimens/seed.py @@ -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"), diff --git a/backend/api/v1/modules/public/reference_data/pedimento_codes/seed.py b/backend/api/v1/modules/public/reference_data/pedimento_codes/seed.py index 1afbcea6..48971d9c 100644 --- a/backend/api/v1/modules/public/reference_data/pedimento_codes/seed.py +++ b/backend/api/v1/modules/public/reference_data/pedimento_codes/seed.py @@ -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.", diff --git a/frontend/messages/en.json b/frontend/messages/en.json index 51be78a8..a4985a02 100644 --- a/frontend/messages/en.json +++ b/frontend/messages/en.json @@ -116,7 +116,9 @@ "customs_brokers": "Customs Brokers", "audit_logs": "Audit Logs", "audit_logs_title": "Audit Logs", - "audit_logs_description": "Detailed audit trail of system operations", + "audit_logs_description": "Audit trail of operations and background task (Celery) status.", + "audit_logs_tab_bitacora": "Audit trail", + "audit_logs_tab_tasks": "Background tasks", "client_provider_type": { "client_indicator": "C", "provider_indicator": "P", diff --git a/frontend/messages/es.json b/frontend/messages/es.json index b1857839..53e27d56 100644 --- a/frontend/messages/es.json +++ b/frontend/messages/es.json @@ -116,7 +116,9 @@ "customs_brokers": "Agentes Aduanales", "audit_logs": "Bitácora", "audit_logs_title": "Bitácora de Movimientos", - "audit_logs_description": "Auditoría detallada de operaciones del sistema", + "audit_logs_description": "Auditoría de operaciones y seguimiento de tareas en segundo plano (Celery).", + "audit_logs_tab_bitacora": "Bitácora", + "audit_logs_tab_tasks": "Tareas en segundo plano", "client_provider_type": { "client_indicator": "C", "provider_indicator": "P", diff --git a/frontend/src/lib/api/dashboard/a76/tasks.ts b/frontend/src/lib/api/dashboard/a76/tasks.ts new file mode 100644 index 00000000..1bd7b194 --- /dev/null +++ b/frontend/src/lib/api/dashboard/a76/tasks.ts @@ -0,0 +1,78 @@ +import { api } from '$lib/api'; + +export type UnifiedTaskStatus = 'pending' | 'active' | 'completed' | 'failed'; + +export interface UnifiedTask { + task_id: string; + task_name: string; + task_group: string; + task_origin?: string | null; + status: UnifiedTaskStatus; + celery_state_raw: string; + progress: { + current?: number | null; + total?: number | null; + percent?: number | null; + message?: string | null; + }; + retries?: number | null; + error?: { + type?: string | null; + message?: string | null; + } | null; + tenant_id: number; + company_id?: number | null; + requested_by_user?: string | null; + started_at?: string | null; + finished_at?: string | null; + created_at: string; + updated_at: string; +} + +export interface UnifiedTaskDetail extends UnifiedTask { + traceback_excerpt?: string | null; + result_summary?: Record | null; + meta_payload?: Record | null; +} + +export interface UnifiedTaskListResponse { + items: UnifiedTask[]; + total: number; + page: number; + page_size: number; + has_next: boolean; +} + +export interface UnifiedTaskCatalogs { + task_groups: string[]; + task_names: string[]; + statuses: string[]; +} + +export const coreTasksApi = { + list: (params?: { + page?: number; + page_size?: number; + status?: UnifiedTaskStatus[]; + task_group?: string[]; + task_name?: string[]; + company_id?: number; + search?: string; + sync_active?: boolean; + }) => { + const qs = new URLSearchParams(); + if (params?.page) qs.set('page', String(params.page)); + if (params?.page_size) qs.set('page_size', String(params.page_size)); + if (params?.company_id != null) qs.set('company_id', String(params.company_id)); + if (params?.search) qs.set('search', params.search); + if (params?.sync_active != null) qs.set('sync_active', String(params.sync_active)); + (params?.status || []).forEach((v) => qs.append('status', v)); + (params?.task_group || []).forEach((v) => qs.append('task_group', v)); + (params?.task_name || []).forEach((v) => qs.append('task_name', v)); + return api.get(`/v1/core/tasks?${qs.toString()}`); + }, + get: (taskId: string, sync = true) => + api.get(`/v1/core/tasks/${encodeURIComponent(taskId)}?sync=${String(sync)}`), + sync: (taskIds?: string[]) => api.post<{ updated: number }>(`/v1/core/tasks/sync`, { task_ids: taskIds }), + catalogs: () => api.get(`/v1/core/tasks/catalogs`) +}; diff --git a/frontend/src/routes/dashboard/audit_logs/+page.svelte b/frontend/src/routes/dashboard/audit_logs/+page.svelte index 8968b41b..7de90e27 100644 --- a/frontend/src/routes/dashboard/audit_logs/+page.svelte +++ b/frontend/src/routes/dashboard/audit_logs/+page.svelte @@ -1,346 +1,78 @@ -
-
+
+

{m['sidebar.audit_logs_title']()}

{m['sidebar.audit_logs_description']()}

-
- -
- -
- - - Filtros - - -
-
- -
- - -
-
+ + + + + {m['sidebar.audit_logs_tab_bitacora']()} + + + + {m['sidebar.audit_logs_tab_tasks']()} + + -
- - -
- -
- - -
- -
- - -
- -
- - -
-
-
- -
-
-
-
- - - - -
-
- Registros - Total: {total} registros encontrados -
-
-
- - {#if error} -
- Error al cargar datos: {error} -
- {:else} -
- - - - ID - Referencia - Procedimiento - Movimiento - Usuario - Fecha - Hora - - - - {#if loading && logs.length === 0} - - Cargando... - - {:else if logs.length === 0} - - No se encontraron registros - - {:else} - {#each logs as log} - - {log.spec_id} - {log.reference} - {log.procedure} - {log.movement} - {log.username} - {formatDate(log.date, log.timestamp)} - {formatTime(log.time, log.timestamp)} - - {/each} - {/if} - - - - -
- {#if loading && logs.length > 0} -
- - Cargando más registros... -
- {/if} -
-
+ + {#if tabValue === 'bitacora'} + {/if} -
-
+ + + {#if tabValue === 'tasks'} + + {/if} + +
diff --git a/frontend/src/routes/dashboard/audit_logs/+page.ts b/frontend/src/routes/dashboard/audit_logs/+page.ts new file mode 100644 index 00000000..1bf476b2 --- /dev/null +++ b/frontend/src/routes/dashboard/audit_logs/+page.ts @@ -0,0 +1,7 @@ +import type { PageLoad } from './$types'; + +export const load: PageLoad = ({ url }) => { + const tab = url.searchParams.get('tab'); + const initialTab = tab === 'tasks' ? 'tasks' : 'bitacora'; + return { initialTab }; +}; diff --git a/frontend/src/routes/dashboard/audit_logs/bitacora-tab.svelte b/frontend/src/routes/dashboard/audit_logs/bitacora-tab.svelte new file mode 100644 index 00000000..da1e2bd7 --- /dev/null +++ b/frontend/src/routes/dashboard/audit_logs/bitacora-tab.svelte @@ -0,0 +1,321 @@ + + +
+
+ +
+ +
+ + + Filtros + + +
+
+ +
+ + +
+
+ +
+ + +
+ +
+ + +
+ +
+ + +
+ +
+ + +
+
+
+ +
+
+
+
+ + + +
+
+ Registros + Total: {total} registros encontrados +
+
+
+ + {#if error} +
+ Error al cargar datos: {error} +
+ {:else} +
+ + + + ID + Referencia + Procedimiento + Movimiento + Usuario + Fecha + Hora + + + + {#if loading && logs.length === 0} + + + Cargando... + + + {:else if logs.length === 0} + + + No se encontraron registros + + + {:else} + {#each logs as log} + + {log.spec_id} + + {log.reference} + + {log.procedure} + {log.movement} + {log.username} + {formatDate(log.date, log.timestamp)} + {formatTime(log.time, log.timestamp)} + + {/each} + {/if} + + + +
+ {#if loading && logs.length > 0} +
+ + Cargando más registros... +
+ {/if} +
+
+ {/if} +
+
+
diff --git a/frontend/src/routes/dashboard/audit_logs/tasks-tab.svelte b/frontend/src/routes/dashboard/audit_logs/tasks-tab.svelte new file mode 100644 index 00000000..43c94c52 --- /dev/null +++ b/frontend/src/routes/dashboard/audit_logs/tasks-tab.svelte @@ -0,0 +1,338 @@ + + +
+
+ +
+ +
+ + + + + Todos + En cola + En progreso + Completadas + Fallidas + + + +
+ + {#if error} +
+ {error} +
+ {/if} + +
+ + + + + + + + + + + + + + {#if tasks.length === 0 && !loading} + + + + {:else if tasks.length === 0 && loading} + + + + {:else} + {#each tasks as task} + void openDetail(task)} + > + + + + + + + + + {/each} + {/if} + +
Task IDTipoEmpresaEstadoProgresoReintentosActualizado
Sin tareas registradas
Cargando...
{task.task_id}{task.task_group} / {task.task_name}{task.company_id ?? '—'} + {statusLabel(task.status)} + ({task.celery_state_raw}) + + {formatPercent(task)} + {task.progress?.message ? ` · ${task.progress.message}` : ''} + {task.retries ?? 0}{new Date(task.updated_at).toLocaleString()}
+
+ +
+
Total: {total}
+
+ + Página {page} + +
+
+
+ + { + if (!open) closeDetail(); + }} +> + + + Detalle de tarea + + {#if detailLoading} + Sincronizando estado con Celery… + {:else if selectedDetail} + {selectedDetail.task_id} + {:else} + Cargando… + {/if} + + + + {#if detailLoading} +
+ Sincronizando… +
+ {/if} + + {#if detailError} +

{detailError}

+ {/if} + + {#if selectedDetail} +
+
+ Task: + {selectedDetail.task_id} +
+
+ Estado: + {statusLabel(selectedDetail.status)} + ({selectedDetail.celery_state_raw}) +
+
+ Progreso: + {formatPercent(selectedDetail)} + {#if selectedDetail.progress?.message} + — {selectedDetail.progress.message} + {/if} +
+
Origen: {selectedDetail.task_origin || '—'}
+
Solicitante: {selectedDetail.requested_by_user || '—'}
+
+ Error: + {selectedDetail.error?.type || '—'} — {selectedDetail.error?.message || '—'} +
+ {#if selectedDetail.started_at} +
Inicio: {new Date(selectedDetail.started_at).toLocaleString()}
+ {/if} + {#if selectedDetail.finished_at} +
Fin: {new Date(selectedDetail.finished_at).toLocaleString()}
+ {/if} + {#if selectedDetail.traceback_excerpt} +
+ Traceback (extracto) +
{selectedDetail.traceback_excerpt}
+
+ {/if} + {#if selectedDetail.result_summary && Object.keys(selectedDetail.result_summary).length > 0} +
+ Resultado (resumen) +
{JSON.stringify(
+								selectedDetail.result_summary,
+								null,
+								2
+							)}
+
+ {/if} +
+ {/if} + + + + +
+
diff --git a/frontend/src/routes/dashboard/tasks/+page.ts b/frontend/src/routes/dashboard/tasks/+page.ts new file mode 100644 index 00000000..2239e38f --- /dev/null +++ b/frontend/src/routes/dashboard/tasks/+page.ts @@ -0,0 +1,6 @@ +import { redirect } from '@sveltejs/kit'; +import type { PageLoad } from './$types'; + +export const load: PageLoad = () => { + throw redirect(303, '/dashboard/audit_logs?tab=tasks'); +}; From 3d1c73565f16733f0cf1667eb23dfe49af6823e1 Mon Sep 17 00:00:00 2001 From: hreyes Date: Tue, 24 Mar 2026 13:15:15 -0600 Subject: [PATCH 2/2] fix/un-alembic --- .../c1a2b3d4e5f6_create_core_task_runs.py | 59 ++++++++++++- ...9_backfill_full_code_pedimento_regimens.py | 82 ------------------- 2 files changed, 58 insertions(+), 83 deletions(-) delete mode 100644 backend/alembic/versions/d4e5f6a7b8c9_backfill_full_code_pedimento_regimens.py diff --git a/backend/alembic/versions/c1a2b3d4e5f6_create_core_task_runs.py b/backend/alembic/versions/c1a2b3d4e5f6_create_core_task_runs.py index dae5b1c5..2af79dcf 100644 --- a/backend/alembic/versions/c1a2b3d4e5f6_create_core_task_runs.py +++ b/backend/alembic/versions/c1a2b3d4e5f6_create_core_task_runs.py @@ -1,4 +1,4 @@ -"""create_core_task_runs +"""create_core_task_runs and backfill code_pedimento_regimens Revision ID: c1a2b3d4e5f6 Revises: bccb7f8986c7 @@ -11,6 +11,13 @@ from alembic import op import sqlalchemy as sa from sqlalchemy.dialects import postgresql +from api.v1.modules.public.reference_data.code_pedimento_regimens.seed import ( + seed as code_pedimento_regimens_seed_full, +) +from api.v1.modules.public.reference_data.pedimento_codes.seed import ( + seed as pedimento_codes_seed, +) + # revision identifiers, used by Alembic. revision: str = "c1a2b3d4e5f6" @@ -19,6 +26,12 @@ branch_labels: Union[str, Sequence[str], None] = None depends_on: Union[str, Sequence[str], None] = None +def _format_value(val: str | None) -> str: + if val is None or str(val).strip() == "" or str(val).upper() == "NONE": + return "NULL" + return f"'{str(val).replace(chr(39), chr(39) * 2)}'" + + def upgrade() -> None: op.create_table( "task_runs", @@ -75,6 +88,50 @@ def upgrade() -> None: schema="core", ) + pedimento_desc_by_code = {code: desc for code, desc in pedimento_codes_seed} + required_pedimento_codes = sorted( + {pedimento_code for pedimento_code, _regimen_code, _type_code in code_pedimento_regimens_seed_full} + ) + + values_missing_codes = ", ".join( + [ + f"({_format_value(code)}, {_format_value(pedimento_desc_by_code.get(code, f'AUTO-GENERATED FOR FK ({code})'))})" + for code in required_pedimento_codes + ] + ) + op.execute( + f""" + INSERT INTO public.pedimento_codes (code, description) + VALUES {values_missing_codes} + ON CONFLICT (code) DO NOTHING; + """ + ) + + values_relations = ", ".join( + [ + f"({_format_value(pedimento_code)}, {_format_value(regimen_code)}, {_format_value(type_code)})" + for pedimento_code, regimen_code, type_code in code_pedimento_regimens_seed_full + ] + ) + op.execute( + f""" + INSERT INTO public.code_pedimento_regimens (pedimento_code, regimen_code, type_code) + SELECT src.pedimento_code, src.regimen_code, src.type_code + FROM (VALUES {values_relations}) AS src(pedimento_code, regimen_code, type_code) + JOIN public.pedimento_codes pc + ON pc.code = src.pedimento_code + JOIN public.pedimento_regimens pr + ON pr.code = src.regimen_code + WHERE NOT EXISTS ( + SELECT 1 + FROM public.code_pedimento_regimens cpr + WHERE cpr.pedimento_code = src.pedimento_code + AND cpr.regimen_code = src.regimen_code + AND COALESCE(cpr.type_code, '') = COALESCE(src.type_code, '') + ); + """ + ) + def downgrade() -> None: op.drop_index("ix_core_task_runs_tenant_company_updated", table_name="task_runs", schema="core") diff --git a/backend/alembic/versions/d4e5f6a7b8c9_backfill_full_code_pedimento_regimens.py b/backend/alembic/versions/d4e5f6a7b8c9_backfill_full_code_pedimento_regimens.py deleted file mode 100644 index c28f7ba5..00000000 --- a/backend/alembic/versions/d4e5f6a7b8c9_backfill_full_code_pedimento_regimens.py +++ /dev/null @@ -1,82 +0,0 @@ -"""backfill_full_code_pedimento_regimens - -Revision ID: d4e5f6a7b8c9 -Revises: c1a2b3d4e5f6 -Create Date: 2026-03-24 13:10:00.000000 -""" - -from typing import Sequence, Union - -from alembic import op -from api.v1.modules.public.reference_data.code_pedimento_regimens.seed import ( - seed as code_pedimento_regimens_seed_full, -) -from api.v1.modules.public.reference_data.pedimento_codes.seed import ( - seed as pedimento_codes_seed, -) - - -# revision identifiers, used by Alembic. -revision: str = "d4e5f6a7b8c9" -down_revision: Union[str, Sequence[str], None] = "c1a2b3d4e5f6" -branch_labels: Union[str, Sequence[str], None] = None -depends_on: Union[str, Sequence[str], None] = None - - -def _format_value(val: str | None) -> str: - if val is None or str(val).strip() == "" or str(val).upper() == "NONE": - return "NULL" - return f"'{str(val).replace(chr(39), chr(39) * 2)}'" - - -def upgrade() -> None: - """Backfill full code_pedimento_regimens data without editing old migrations.""" - - pedimento_desc_by_code = {code: desc for code, desc in pedimento_codes_seed} - required_pedimento_codes = sorted( - {pedimento_code for pedimento_code, _regimen_code, _type_code in code_pedimento_regimens_seed_full} - ) - - values_missing_codes = ", ".join( - [ - f"({_format_value(code)}, {_format_value(pedimento_desc_by_code.get(code, f'AUTO-GENERATED FOR FK ({code})'))})" - for code in required_pedimento_codes - ] - ) - op.execute( - f""" - INSERT INTO public.pedimento_codes (code, description) - VALUES {values_missing_codes} - ON CONFLICT (code) DO NOTHING; - """ - ) - - values_relations = ", ".join( - [ - f"({_format_value(pedimento_code)}, {_format_value(regimen_code)}, {_format_value(type_code)})" - for pedimento_code, regimen_code, type_code in code_pedimento_regimens_seed_full - ] - ) - op.execute( - f""" - INSERT INTO public.code_pedimento_regimens (pedimento_code, regimen_code, type_code) - SELECT src.pedimento_code, src.regimen_code, src.type_code - FROM (VALUES {values_relations}) AS src(pedimento_code, regimen_code, type_code) - JOIN public.pedimento_codes pc - ON pc.code = src.pedimento_code - JOIN public.pedimento_regimens pr - ON pr.code = src.regimen_code - WHERE NOT EXISTS ( - SELECT 1 - FROM public.code_pedimento_regimens cpr - WHERE cpr.pedimento_code = src.pedimento_code - AND cpr.regimen_code = src.regimen_code - AND COALESCE(cpr.type_code, '') = COALESCE(src.type_code, '') - ); - """ - ) - - -def downgrade() -> None: - """Preserve seeded data on downgrade.""" -