108 lines
3.8 KiB
Python
108 lines
3.8 KiB
Python
"""
|
|
Worker ARQ — traitement asynchrone des images via Redis.
|
|
|
|
Utilise le parser d'URL Redis centralisé de app.workers.redis_client.
|
|
Inclut la Dead Letter Queue pour les jobs échoués après N tentatives.
|
|
"""
|
|
import logging
|
|
import asyncio
|
|
from datetime import datetime, timezone
|
|
|
|
from arq import cron, func
|
|
from arq.connections import RedisSettings
|
|
|
|
from app.config import settings
|
|
from app.database import AsyncSessionLocal
|
|
from app.models.image import Image, ProcessingStatus
|
|
from app.services.pipeline import process_image_pipeline
|
|
from app.workers.redis_client import parse_redis_url
|
|
from sqlalchemy import select
|
|
|
|
# Préfixes ARQ
|
|
QUEUE_STANDARD = "standard"
|
|
QUEUE_PREMIUM = "premium"
|
|
DEFAULT_QUEUE_NAME = "arq:queue"
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
async def process_image_task(ctx: dict, image_id: int, client_id: str) -> str:
|
|
"""Tâche ARQ : traite une image. Pousse en DLQ après échecs répétés."""
|
|
job_try = ctx.get("job_try", 1)
|
|
redis = ctx.get("redis")
|
|
max_retries = getattr(settings, "PIPELINE_MAX_RETRIES", 3)
|
|
|
|
logger.info(
|
|
"arq.job.started",
|
|
extra={"image_id": image_id, "client_id": client_id, "attempt": job_try},
|
|
)
|
|
|
|
async with AsyncSessionLocal() as db:
|
|
try:
|
|
result = await db.execute(select(Image).where(Image.id == image_id))
|
|
image = result.scalar_one_or_none()
|
|
if image is None:
|
|
logger.warning("arq.job.image_not_found", extra={"image_id": image_id})
|
|
return "SKIPPED: image not found"
|
|
if image.processing_status == ProcessingStatus.DONE:
|
|
logger.info("arq.job.already_done", extra={"image_id": image_id})
|
|
return "SKIPPED: already done"
|
|
|
|
await process_image_pipeline(image_id, db, redis=redis)
|
|
logger.info("arq.job.completed", extra={"image_id": image_id})
|
|
return "OK"
|
|
|
|
except Exception as e:
|
|
logger.error(
|
|
"arq.job.failed",
|
|
extra={"image_id": image_id, "error": str(e), "attempt": job_try},
|
|
exc_info=True,
|
|
)
|
|
|
|
# Si c'est la dernière tentative, pousser en DLQ
|
|
if job_try >= max_retries:
|
|
from app.workers.dead_letter import push_dead_job
|
|
await push_dead_job(
|
|
redis=redis,
|
|
job_name="process_image_task",
|
|
args=[image_id, client_id],
|
|
kwargs={},
|
|
error=str(e),
|
|
attempts=job_try,
|
|
)
|
|
# Marquer l'image en erreur permanente
|
|
async with AsyncSessionLocal() as db2:
|
|
result2 = await db2.execute(select(Image).where(Image.id == image_id))
|
|
img = result2.scalar_one_or_none()
|
|
if img:
|
|
img.processing_status = ProcessingStatus.ERROR
|
|
img.processing_error = f"[DLQ] Échec après {job_try} tentatives: {str(e)[:200]}"
|
|
await db2.commit()
|
|
return f"DLQ: failed after {job_try} attempts"
|
|
|
|
raise # Re-raise pour que ARQ retente
|
|
|
|
|
|
async def on_startup(ctx: dict) -> None:
|
|
logger.info("worker.started", extra={"queue": WorkerSettings.queue_name})
|
|
|
|
|
|
async def on_shutdown(ctx: dict) -> None:
|
|
logger.info("worker.shutdown")
|
|
|
|
|
|
class WorkerSettings:
|
|
functions = [func(process_image_task, name="process_image_task")]
|
|
redis_settings = RedisSettings(
|
|
host=parse_redis_url()[0],
|
|
port=parse_redis_url()[1],
|
|
password=parse_redis_url()[2],
|
|
database=parse_redis_url()[3],
|
|
)
|
|
queue_name = DEFAULT_QUEUE_NAME
|
|
on_startup = on_startup
|
|
on_shutdown = on_shutdown
|
|
max_jobs = 10
|
|
job_timeout = settings.PIPELINE_TIMEOUT
|
|
health_check_interval = 30
|