Files
Imago/app/workers/dead_letter.py
T
bruno caefc1b1a0
CI / Lint & Format (push) Failing after 11s
CI / Tests (push) Has been skipped
CI / Security Scan (push) Failing after 7s
CI / Docker Build (push) Has been skipped
Add configuration files and comprehensive project documentation
- Add .env.backup with full application configuration template
- Create RequestIDMiddleware for trace ID injection in logs
- Implement ARQ fallback pool for Redis unavailability
- Add Dead Letter Queue module for failed job handling
- Include project ROADMAP with audit and 27 applied corrections
- Add Grafana monitoring dashboard with 10 panels
2026-06-22 14:50:05 -04:00

140 lines
3.8 KiB
Python

"""
Dead Letter Queue — jobs ARQ échoués après N tentatives.
Stocke les jobs échoués dans Redis pour inspection et retry manuel.
Expose des endpoints admin pour lister et réessayer les jobs morts.
"""
import json
import logging
from datetime import datetime, timezone
from typing import Any
from app.config import settings
from app.metrics import hub_arq_jobs_failed
logger = logging.getLogger(__name__)
DLQ_KEY = "arq:dead:jobs"
DLQ_MAX_JOBS = 1000 # Garder max 1000 jobs morts
async def push_dead_job(
redis: Any,
job_name: str,
args: list,
kwargs: dict,
error: str,
attempts: int,
) -> None:
"""
Pousse un job échoué dans la Dead Letter Queue Redis.
Format stocké : JSON avec job_name, args, kwargs, error, attempts, failed_at.
"""
if redis is None:
logger.warning("dlq.no_redis", extra={"job_name": job_name})
return
try:
entry = {
"job_name": job_name,
"args": [str(a) for a in (args or [])],
"kwargs": {k: str(v) for k, v in (kwargs or {}).items()},
"error": str(error)[:500],
"attempts": attempts,
"failed_at": datetime.now(timezone.utc).isoformat(),
}
await redis.lpush(DLQ_KEY, json.dumps(entry))
await redis.ltrim(DLQ_KEY, 0, DLQ_MAX_JOBS - 1)
hub_arq_jobs_failed.inc()
logger.warning("dlq.pushed", extra={
"job_name": job_name,
"attempts": attempts,
"error": str(error)[:100],
})
except Exception as e:
logger.error("dlq.push_error", extra={"error": str(e)})
async def get_dead_jobs(redis: Any, limit: int = 50) -> list[dict]:
"""Récupère les N derniers jobs morts."""
if redis is None:
return []
try:
raw = await redis.lrange(DLQ_KEY, 0, limit - 1)
return [json.loads(r) for r in raw]
except Exception:
return []
async def get_dead_job_count(redis: Any) -> int:
"""Nombre de jobs dans la DLQ."""
if redis is None:
return 0
try:
return await redis.llen(DLQ_KEY)
except Exception:
return 0
async def retry_dead_job(
redis: Any,
arq_pool: Any,
index: int,
) -> dict:
"""
Retente un job mort par son index dans la DLQ (0 = le plus récent).
Supprime l'entrée de la DLQ après remise en file.
"""
from app.workers.arq_fallback import is_fallback_pool
if redis is None:
return {"error": "Redis indisponible"}
if is_fallback_pool(arq_pool):
return {"error": "ARQ indisponible (fallback)"}
try:
raw = await redis.lindex(DLQ_KEY, index)
if raw is None:
return {"error": f"Index {index} hors limites"}
entry = json.loads(raw)
job_name = entry["job_name"]
kwargs = {k: v for k, v in entry.get("kwargs", {}).items()}
# Convertir les args en int si nécessaire (image_id, client_id)
args = []
for a in entry.get("args", []):
try:
args.append(int(a))
except (ValueError, TypeError):
args.append(a)
await arq_pool.enqueue_job(job_name, *args, **kwargs)
# Supprimer l'entrée de la DLQ
await redis.lset(DLQ_KEY, index, "__RETRIED__")
await redis.lrem(DLQ_KEY, 0, "__RETRIED__")
logger.info("dlq.retried", extra={"job_name": job_name, "args": args})
return {"status": "retried", "job_name": job_name, "args": args}
except Exception as e:
logger.error("dlq.retry_error", extra={"error": str(e)})
return {"error": str(e)}
async def clear_dead_jobs(redis: Any) -> int:
"""Vide la DLQ. Retourne le nombre de jobs supprimés."""
if redis is None:
return 0
try:
count = await redis.llen(DLQ_KEY)
await redis.delete(DLQ_KEY)
return count
except Exception:
return 0