dead_letter_queue
Cola de mensajes fallidos (DLQ) para capturar tareas del pipeline KYC que exceden el numero maximo de reintentos, como inferencia facial fallida, OCR timeout o errores de liveness detection. Permite analisis post-mortem de fallos recurrentes y re-procesamiento manual o automatizado de tareas recuperables.
When to use
Usa esta skill cuando trabajes con el worker_pool_agent y necesites implementar o configurar el manejo de tareas fallidas en el pipeline de verificacion de identidad. Aplica cuando tareas criticas (face_match, OCR, liveness) fallan repetidamente y necesitas capturarlas en lugar de perderlas.
Instructions
Definir la configuracion de reintentos maximos por tipo de tarea en la configuracion de Celery:
# backend/modules/worker_pool/celery_config.py
TASK_MAX_RETRIES = {
"face_match": 3,
"ocr_extraction": 5,
"liveness_detection": 3,
"doc_processing": 4,
}
Crear el modelo de base de datos para almacenar tareas en la DLQ:
# backend/modules/worker_pool/models/dead_letter.py
class DeadLetterEntry(Base):
__tablename__ = "dead_letter_queue"
id = Column(UUID, primary_key=True, default=uuid4)
task_name = Column(String, nullable=False)
task_id = Column(String, unique=True, nullable=False)
args = Column(JSON)
kwargs = Column(JSON)
exception = Column(Text)
traceback = Column(Text)
retry_count = Column(Integer)
created_at = Column(DateTime, default=datetime.utcnow)
status = Column(String, default="pending") # pending, reprocessed, discarded
Implementar el handler de task_failure que envia tareas agotadas a la DLQ:
from celery.signals import task_failure
@task_failure.connect
def handle_task_failure(sender, task_id, exception, args, kwargs, traceback, einfo, **kw):
if sender.request.retries >= TASK_MAX_RETRIES.get(sender.name, 3):
DeadLetterEntry.create(
task_name=sender.name,
task_id=task_id,
args=args,
kwargs=kwargs,
exception=str(exception),
traceback=str(einfo),
retry_count=sender.request.retries,
)
Crear endpoint API para consultar y gestionar la DLQ:
# backend/api/routes/dead_letter.py
@router.get("/dlq/entries")
async def list_dlq_entries(status: str = "pending", limit: int = 50):
return await DeadLetterEntry.filter(status=status).limit(limit).all()
@router.post("/dlq/entries/{entry_id}/reprocess")
async def reprocess_entry(entry_id: UUID):
entry = await DeadLetterEntry.get(id=entry_id)
celery_app.send_task(entry.task_name, args=entry.args, kwargs=entry.kwargs)
entry.status = "reprocessed"
await entry.save()
Implementar un job periodico que analice patrones de fallo en la DLQ:
@celery_app.task(name="dlq_analysis")
def analyze_dlq_patterns():
recent = DeadLetterEntry.filter(
created_at__gte=datetime.utcnow() - timedelta(hours=1)
).all()
failure_counts = Counter(e.task_name for e in recent)
for task_name, count in failure_counts.items():
if count > ALERT_THRESHOLD:
send_alert(f"DLQ: {task_name} tiene {count} fallos en la ultima hora")
Configurar alertas y metricas de la DLQ para monitoreo:
DLQ_METRICS = {
"dlq_entries_total": Counter("dlq_entries_total", "Total DLQ entries", ["task_name"]),
"dlq_reprocessed_total": Counter("dlq_reprocessed_total", "Reprocessed entries", ["task_name"]),
}
Implementar politica de retencion para limpiar entradas antiguas de la DLQ:
@celery_app.task(name="dlq_cleanup")
def cleanup_old_entries(retention_days: int = 30):
cutoff = datetime.utcnow() - timedelta(days=retention_days)
DeadLetterEntry.filter(created_at__lt=cutoff, status="discarded").delete()
Notes
- Las entradas de la DLQ deben anonimizar datos biometricos y personales segun GDPR/LOPD; almacenar solo referencias de sesion, nunca imagenes ni embeddings faciales.
- El re-procesamiento manual desde la DLQ debe pasar por las mismas validaciones antifraude que el flujo original para evitar bypass de seguridad.
- Monitorear el crecimiento de la DLQ como indicador de salud del sistema: un incremento sostenido indica problemas sistemicos en el pipeline que requieren atencion inmediata.
1---2name: dead-letter-queue3description: Cola de mensajes fallidos para capturar y reprocesar tareas KYC que exceden reintentos maximos4---56# dead_letter_queue78Cola de mensajes fallidos (DLQ) para capturar tareas del pipeline KYC que exceden el numero maximo de reintentos, como inferencia facial fallida, OCR timeout o errores de liveness detection. Permite analisis post-mortem de fallos recurrentes y re-procesamiento manual o automatizado de tareas recuperables.910## When to use1112Usa esta skill cuando trabajes con el **worker_pool_agent** y necesites implementar o configurar el manejo de tareas fallidas en el pipeline de verificacion de identidad. Aplica cuando tareas criticas (face_match, OCR, liveness) fallan repetidamente y necesitas capturarlas en lugar de perderlas.1314## Instructions15161. Definir la configuracion de reintentos maximos por tipo de tarea en la configuracion de Celery:17 ```python18 # backend/modules/worker_pool/celery_config.py19 TASK_MAX_RETRIES = {20 "face_match": 3,21 "ocr_extraction": 5,22 "liveness_detection": 3,23 "doc_processing": 4,24 }25 ```26272. Crear el modelo de base de datos para almacenar tareas en la DLQ:28 ```python29 # backend/modules/worker_pool/models/dead_letter.py30 class DeadLetterEntry(Base):31 __tablename__ = "dead_letter_queue"32 id = Column(UUID, primary_key=True, default=uuid4)33 task_name = Column(String, nullable=False)34 task_id = Column(String, unique=True, nullable=False)35 args = Column(JSON)36 kwargs = Column(JSON)37 exception = Column(Text)38 traceback = Column(Text)39 retry_count = Column(Integer)40 created_at = Column(DateTime, default=datetime.utcnow)41 status = Column(String, default="pending") # pending, reprocessed, discarded42 ```43443. Implementar el handler de task_failure que envia tareas agotadas a la DLQ:45 ```python46 from celery.signals import task_failure4748 @task_failure.connect49 def handle_task_failure(sender, task_id, exception, args, kwargs, traceback, einfo, **kw):50 if sender.request.retries >= TASK_MAX_RETRIES.get(sender.name, 3):51 DeadLetterEntry.create(52 task_name=sender.name,53 task_id=task_id,54 args=args,55 kwargs=kwargs,56 exception=str(exception),57 traceback=str(einfo),58 retry_count=sender.request.retries,59 )60 ```61624. Crear endpoint API para consultar y gestionar la DLQ:63 ```python64 # backend/api/routes/dead_letter.py65 @router.get("/dlq/entries")66 async def list_dlq_entries(status: str = "pending", limit: int = 50):67 return await DeadLetterEntry.filter(status=status).limit(limit).all()6869 @router.post("/dlq/entries/{entry_id}/reprocess")70 async def reprocess_entry(entry_id: UUID):71 entry = await DeadLetterEntry.get(id=entry_id)72 celery_app.send_task(entry.task_name, args=entry.args, kwargs=entry.kwargs)73 entry.status = "reprocessed"74 await entry.save()75 ```76775. Implementar un job periodico que analice patrones de fallo en la DLQ:78 ```python79 @celery_app.task(name="dlq_analysis")80 def analyze_dlq_patterns():81 recent = DeadLetterEntry.filter(82 created_at__gte=datetime.utcnow() - timedelta(hours=1)83 ).all()84 failure_counts = Counter(e.task_name for e in recent)85 for task_name, count in failure_counts.items():86 if count > ALERT_THRESHOLD:87 send_alert(f"DLQ: {task_name} tiene {count} fallos en la ultima hora")88 ```89906. Configurar alertas y metricas de la DLQ para monitoreo:91 ```python92 DLQ_METRICS = {93 "dlq_entries_total": Counter("dlq_entries_total", "Total DLQ entries", ["task_name"]),94 "dlq_reprocessed_total": Counter("dlq_reprocessed_total", "Reprocessed entries", ["task_name"]),95 }96 ```97987. Implementar politica de retencion para limpiar entradas antiguas de la DLQ:99 ```python100 @celery_app.task(name="dlq_cleanup")101 def cleanup_old_entries(retention_days: int = 30):102 cutoff = datetime.utcnow() - timedelta(days=retention_days)103 DeadLetterEntry.filter(created_at__lt=cutoff, status="discarded").delete()104 ```105106## Notes107108- Las entradas de la DLQ deben anonimizar datos biometricos y personales segun GDPR/LOPD; almacenar solo referencias de sesion, nunca imagenes ni embeddings faciales.109- El re-procesamiento manual desde la DLQ debe pasar por las mismas validaciones antifraude que el flujo original para evitar bypass de seguridad.110- Monitorear el crecimiento de la DLQ como indicador de salud del sistema: un incremento sostenido indica problemas sistemicos en el pipeline que requieren atencion inmediata.