From c565696242513b44676d79b2d07f41eb84e88ad3 Mon Sep 17 00:00:00 2001 From: hreyes Date: Tue, 24 Mar 2026 16:58:48 -0600 Subject: [PATCH] feature/limpieza-periodica --- .../v1/modules/a76/layouts_csv/boms/routes.py | 1 + .../cambio_regimen_regularizacion/routes.py | 1 + .../modules/a76/layouts_csv/classes/routes.py | 1 + .../clients_and_providers/routes.py | 1 + .../layouts_csv/common/import_redis_keys.py | 69 +++++++ .../common/track_commit_dispatch.py | 2 + .../modules/a76/layouts_csv/common/victor.py | 172 ++++++++++++++++++ .../a76/layouts_csv/customs_brokers/routes.py | 1 + .../a76/layouts_csv/exchange_rate/routes.py | 1 + .../a76/layouts_csv/exportacion/routes.py | 1 + .../a76/layouts_csv/facturas/routes.py | 1 + .../modules/a76/layouts_csv/parts/routes.py | 1 + .../a76/layouts_csv/pedmientos/routes.py | 1 + .../layouts_csv/us_tariff_fractions/routes.py | 1 + backend/core/celery_app.py | 5 + docker-compose.prod.yml | 31 ++++ scripts/run_layout_import_janitor.sh | 8 + 17 files changed, 298 insertions(+) create mode 100644 backend/api/v1/modules/a76/layouts_csv/common/import_redis_keys.py create mode 100644 backend/api/v1/modules/a76/layouts_csv/common/victor.py create mode 100644 scripts/run_layout_import_janitor.sh 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 3d33f075..f99183d5 100644 --- a/backend/api/v1/modules/a76/layouts_csv/boms/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/boms/routes.py @@ -167,6 +167,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=f"{BOM_IMPORT_META_PREFIX}{job_id}", celery_task=insert_valid_rows, task_name="boms_insert_valid_rows", 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 e646e107..9f8a4ccd 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 @@ -176,6 +176,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=meta_key, celery_task=insert_valid_rows, task_name="cambio_regimen_regularizacion_insert_valid_rows", 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 24801833..53464802 100644 --- a/backend/api/v1/modules/a76/layouts_csv/classes/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/classes/routes.py @@ -194,6 +194,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=f"{CLS_IMPORT_META_PREFIX}{job_id}", celery_task=insert_valid_rows, task_name="classes_insert_valid_rows", 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 59c99ef7..fcc092dd 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 @@ -179,6 +179,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=f"{CP_IMPORT_META_PREFIX}{job_id}", celery_task=insert_valid_rows, task_name="clients_and_providers_insert_valid_rows", diff --git a/backend/api/v1/modules/a76/layouts_csv/common/import_redis_keys.py b/backend/api/v1/modules/a76/layouts_csv/common/import_redis_keys.py new file mode 100644 index 00000000..8d763d1d --- /dev/null +++ b/backend/api/v1/modules/a76/layouts_csv/common/import_redis_keys.py @@ -0,0 +1,69 @@ +""" +Claves Redis por job de import CSV (layouts_csv): las de storage_keys más extras por módulo. +Usado por la limpieza periódica (victor "el senior de limpieza") para saber si un job sigue vivo en Redis. +""" +from __future__ import annotations + +from typing import Any, List + +from . import storage as common_storage + +# Prefijo en nombre de archivo `{prefix}_{uuid}.csv` -> job_type (vacío = solo `{uuid}.csv`) +FILE_PREFIX_TO_JOB_TYPE: dict[str, str] = { + "veh": "veh", + "fa": "fa", + "trl": "trl", + "trp": "trp", + "ped": "ped", + "part": "part", + "er": "er", + "cb": "cb", + "drv": "drv", + "cls": "cls", + "cp": "cp", + "bom": "bom", + "exp": "exp", + "crreg": "crreg", +} + + +def all_redis_keys_for_layout_import(job_type: str, job_id: str) -> List[str]: + """Todas las claves Redis que pueden existir para un job (file, meta, error_lines, extras).""" + file_k, meta_k, err_k = common_storage.storage_keys(job_type, job_id) + keys: List[str] = [file_k, meta_k, err_k] + if job_type == "veh": + keys.append(f"veh_import_status:{job_id}") + elif job_type == "trl": + keys.append(f"trl_import_status:{job_id}") + elif job_type == "trp": + keys.append(f"trp_import_status:{job_id}") + elif job_type == "drv": + keys.append(f"drv_import_status:{job_id}") + keys.append(f"drv_import_transporter_map:{job_id}") + return keys + + +def any_layout_import_key_exists(redis_client: Any, job_type: str, job_id: str) -> bool: + for k in all_redis_keys_for_layout_import(job_type, job_id): + if redis_client.exists(k): + return True + return False + + +def delete_extra_layout_import_keys_from_redis(redis_client: Any, job_type: str, job_id: str) -> None: + """Borra claves que delete_import_from_redis no cubre (status, mapas).""" + extras: list[str] = [] + if job_type == "veh": + extras.append(f"veh_import_status:{job_id}") + elif job_type == "trl": + extras.append(f"trl_import_status:{job_id}") + elif job_type == "trp": + extras.append(f"trp_import_status:{job_id}") + elif job_type == "drv": + extras.append(f"drv_import_status:{job_id}") + extras.append(f"drv_import_transporter_map:{job_id}") + if extras: + try: + redis_client.delete(*extras) + except Exception: + pass 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 index b8479a0a..296906c4 100644 --- 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 @@ -18,6 +18,7 @@ def dispatch_tracked_layouts_csv_commit( db: Session, current_user: dict[str, Any], redis_client: Any, + layout_job_id: str, meta_redis_key: str, celery_task: Task, task_name: str, @@ -54,5 +55,6 @@ def dispatch_tracked_layouts_csv_commit( task_origin=task_origin, args=args, task_id=commit_id, + meta_payload={"layout_import_job_id": layout_job_id}, ) return commit_id diff --git a/backend/api/v1/modules/a76/layouts_csv/common/victor.py b/backend/api/v1/modules/a76/layouts_csv/common/victor.py new file mode 100644 index 00000000..eab7d2bf --- /dev/null +++ b/backend/api/v1/modules/a76/layouts_csv/common/victor.py @@ -0,0 +1,172 @@ +""" +Limpieza periódica de archivos huérfanos de import CSV (layouts/imports/temp y errors). +""" +from __future__ import annotations + +import logging +import os +import re +import time +from typing import Optional, Tuple + +from sqlalchemy import or_ + +from api.v1.modules.core.tasks_tracking.models import TaskRun, TaskStatus +from core.celery_app import celery_app +from core.database import CoreSessionLocal + +from . import storage as common_storage +from .import_redis_keys import ( + FILE_PREFIX_TO_JOB_TYPE, + any_layout_import_key_exists, + delete_extra_layout_import_keys_from_redis, +) + +logger = logging.getLogger(__name__) + +# Evita carreras con uploads/commits recién creados (segundos) +MIN_ORPHAN_FILE_AGE_SEC = int(os.getenv("LAYOUT_ORPHAN_MIN_AGE_SEC", "600")) + +_UUID_RE = r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}" + + +def _get_redis(): + import redis + + url = os.getenv("VALKEY_URL", os.getenv("REDIS_URL", "redis://valkey:6379/0")) + return redis.Redis.from_url(url, decode_responses=False) + + +def parse_temp_csv_filename(name: str) -> Optional[Tuple[str, str]]: + """Devuelve (job_type, job_id) o None si no es un CSV de import conocido.""" + if not name.endswith(".csv"): + return None + base = name[:-4] + m = re.match(rf"^({_UUID_RE})$", base, re.I) + if m: + return "", m.group(1) + m = re.match(rf"^([a-z0-9]+)_({_UUID_RE})$", base, re.I) + if not m: + return None + prefix, jid = m.group(1).lower(), m.group(2) + jt = FILE_PREFIX_TO_JOB_TYPE.get(prefix) + if jt is None: + return None + return jt, jid + + +def parse_errors_jsonl_filename(name: str) -> Optional[Tuple[str, str]]: + if not name.endswith(".jsonl"): + return None + base = name[:-6] + m = re.match(rf"^({_UUID_RE})$", base, re.I) + if m: + return "", m.group(1) + m = re.match(rf"^([a-z0-9]+)_({_UUID_RE})$", base, re.I) + if not m: + return None + prefix, jid = m.group(1).lower(), m.group(2) + jt = FILE_PREFIX_TO_JOB_TYPE.get(prefix) + if jt is None: + return None + return jt, jid + + +def _layout_import_protected_by_task_run(db, job_id: str) -> bool: + row = ( + db.query(TaskRun) + .filter( + TaskRun.task_group == "layouts_csv", + TaskRun.status.in_([TaskStatus.PENDING.value, TaskStatus.ACTIVE.value]), + or_( + TaskRun.task_id == job_id, + TaskRun.meta_payload.contains({"layout_import_job_id": job_id}), + ), + ) + .first() + ) + return row is not None + + +def _reference_mtime(job_type: str, job_id: str) -> Optional[float]: + csv_path = common_storage.file_path_for_job(job_type, job_id) + err_path = common_storage.error_path_for_job(job_type, job_id) + mtimes = [] + if os.path.isfile(csv_path): + mtimes.append(os.path.getmtime(csv_path)) + if os.path.isfile(err_path): + mtimes.append(os.path.getmtime(err_path)) + if not mtimes: + return None + return max(mtimes) + + +def _try_cleanup_layout_import_job(r, db, job_type: str, job_id: str) -> bool: + if any_layout_import_key_exists(r, job_type, job_id): + return False + if _layout_import_protected_by_task_run(db, job_id): + return False + mtime = _reference_mtime(job_type, job_id) + if mtime is None: + return False + if time.time() - mtime < MIN_ORPHAN_FILE_AGE_SEC: + return False + + file_path = common_storage.file_path_for_job(job_type, job_id) + error_path = common_storage.error_path_for_job(job_type, job_id) + meta_path = None + if os.path.isfile(file_path): + meta_path = file_path.replace(".csv", ".meta.json") + + common_storage.cleanup_import_job( + job_type, + job_id, + file_path=file_path if os.path.isfile(file_path) else None, + error_path=error_path if os.path.isfile(error_path) else None, + meta_path=meta_path if meta_path and os.path.isfile(meta_path) else None, + ) + delete_extra_layout_import_keys_from_redis(r, job_type, job_id) + return True + + +def run_orphan_layout_import_cleanup() -> dict: + """Barrido síncrono; devuelve contadores para logs/resultado Celery.""" + removed_jobs: list[str] = [] + r = _get_redis() + db = CoreSessionLocal() + try: + temp_dir = common_storage.upload_dir() + os.makedirs(temp_dir, exist_ok=True) + for name in os.listdir(temp_dir): + parsed = parse_temp_csv_filename(name) + if not parsed: + continue + job_type, job_id = parsed + if _try_cleanup_layout_import_job(r, db, job_type, job_id): + removed_jobs.append(f"{job_type or 'invoice'}:{job_id}") + + err_dir = common_storage.error_dir() + if os.path.isdir(err_dir): + for name in os.listdir(err_dir): + parsed = parse_errors_jsonl_filename(name) + if not parsed: + continue + job_type, job_id = parsed + csv_path = common_storage.file_path_for_job(job_type, job_id) + if os.path.isfile(csv_path): + continue + if _try_cleanup_layout_import_job(r, db, job_type, job_id): + rid = f"{job_type or 'invoice'}:{job_id}:errors_only" + if rid not in removed_jobs: + removed_jobs.append(rid) + finally: + db.close() + + if removed_jobs: + logger.info("layout import cleanup removed %s job(s): %s", len(removed_jobs), removed_jobs[:20]) + return {"removed_count": len(removed_jobs), "removed": removed_jobs} + + +@celery_app.task(name="cleanup_orphan_layout_imports") +def cleanup_orphan_layout_imports(): + return run_orphan_layout_import_cleanup() 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 ab884351..f0cebaaf 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 @@ -179,6 +179,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=f"{CB_IMPORT_META_PREFIX}{job_id}", celery_task=insert_valid_rows, task_name="customs_brokers_insert_valid_rows", 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 ff55abf5..db570733 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 @@ -189,6 +189,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=f"{ER_IMPORT_META_PREFIX}{job_id}", celery_task=insert_valid_rows, task_name="exchange_rate_insert_valid_rows", 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 82c2af93..dca65485 100644 --- a/backend/api/v1/modules/a76/layouts_csv/exportacion/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/exportacion/routes.py @@ -181,6 +181,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=meta_key, celery_task=insert_valid_rows, task_name="exportacion_insert_valid_rows", 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 18b39eee..c9ac5285 100644 --- a/backend/api/v1/modules/a76/layouts_csv/facturas/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/facturas/routes.py @@ -256,6 +256,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=f"{IMPORT_META_KEY_PREFIX}{job_id}", celery_task=insert_valid_rows, task_name="facturas_insert_valid_rows", 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 d4a5af18..5184073e 100644 --- a/backend/api/v1/modules/a76/layouts_csv/parts/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/parts/routes.py @@ -192,6 +192,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=f"{PART_IMPORT_META_PREFIX}{job_id}", celery_task=insert_valid_rows, task_name="parts_insert_valid_rows", 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 36f601fe..eda3ea86 100644 --- a/backend/api/v1/modules/a76/layouts_csv/pedmientos/routes.py +++ b/backend/api/v1/modules/a76/layouts_csv/pedmientos/routes.py @@ -187,6 +187,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=meta_key, celery_task=insert_valid_rows, task_name="pedmientos_insert_valid_rows", 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 322ca183..afb8c3bb 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 @@ -181,6 +181,7 @@ async def commit_import_job( db=db, current_user=current_user, redis_client=r, + layout_job_id=job_id, meta_redis_key=f"{FA_IMPORT_META_PREFIX}{job_id}", celery_task=insert_valid_rows, task_name="us_tariff_fractions_insert_valid_rows", diff --git a/backend/core/celery_app.py b/backend/core/celery_app.py index 7137ff77..d7cae18a 100644 --- a/backend/core/celery_app.py +++ b/backend/core/celery_app.py @@ -66,6 +66,7 @@ celery_app.conf.update( "api.v1.modules.a76.invoices.imports.revert.task", "api.v1.modules.a76.invoices.exports.process.task", "api.v1.modules.a76.invoices.exports.revert.task", + "api.v1.modules.a76.layouts_csv.common.victor", ] # Ruta al módulo donde están las tareas ) @@ -84,6 +85,10 @@ celery_app.conf.beat_schedule = { "task": "sync_from_hub_task", "schedule": 60.0, # Run every 60 seconds }, + "cleanup-orphan-layout-imports-hourly": { + "task": "cleanup_orphan_layout_imports", + "schedule": 3600.0, + }, } if __name__ == "__main__": diff --git a/docker-compose.prod.yml b/docker-compose.prod.yml index 938efb43..e823b73e 100644 --- a/docker-compose.prod.yml +++ b/docker-compose.prod.yml @@ -187,6 +187,7 @@ services: condition: service_healthy volumes: - backend_uploads:/app/uploads + - backend_layouts:/app/layouts - ./scripts/backend-entrypoint.sh:/entrypoint.sh:ro networks: - backend-net @@ -233,9 +234,37 @@ services: depends_on: - backend - valkey + volumes: + - backend_layouts:/app/layouts networks: - backend-net + celery_beat: + image: dev.aduanasoft.com/anexo76/backend:latest + container_name: celery_beat + command: celery -A core.celery_app beat --loglevel=info + environment: + - VALKEY_URL=${VALKEY_URL:-redis://valkey:6379/0} + - CENTRAL_SERVER_URL=${CENTRAL_SERVER_URL:-""} + - SYNC_SECRET_TOKEN=${SYNC_SECRET_TOKEN:-change-this-sync-token-in-production} + - SPOKE_URLS=${SPOKE_URLS:-""} + - CORE_DB_HOST=${CORE_DB_HOST:-postgres-a76} + - CORE_DB_PORT=${CORE_DB_PORT:-5432} + - CORE_DB_NAME=${CORE_DB_NAME:-anexo76_core} + - CORE_DB_USER=${CORE_DB_USER:-postgres} + - CORE_DB_PASSWORD=${POSTGRES_APP_PASSWORD:-postgres} + - SITAR_API_URL=${SITAR_API_URL} + - SITAR_API_USER=${SITAR_API_USER} + - SITAR_API_PASSWORD=${SITAR_API_PASSWORD} + depends_on: + - backend + - valkey + volumes: + - backend_layouts:/app/layouts + networks: + - backend-net + restart: unless-stopped + valkey: image: valkey/valkey:7.2 container_name: valkey @@ -305,6 +334,8 @@ volumes: driver: local backend_uploads: driver: local + backend_layouts: + driver: local networks: backend-net: diff --git a/scripts/run_layout_import_janitor.sh b/scripts/run_layout_import_janitor.sh new file mode 100644 index 00000000..f4c99089 --- /dev/null +++ b/scripts/run_layout_import_janitor.sh @@ -0,0 +1,8 @@ +#!/usr/bin/env bash +# Ejecuta una pasada del janitor de layouts CSV (tarea Celery cleanup_orphan_layout_imports). +# Requiere stack levantado: docker compose up -d worker valkey postgres-a76 +set -euo pipefail +ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +cd "$ROOT" +CONTAINER="${CELERY_WORKER_CONTAINER:-worker}" +docker compose exec -T "$CONTAINER" celery -A core.celery_app call cleanup_orphan_layout_imports