Files
Imago/app/workers/image_worker.py
T

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