diff --git a/.env.example b/.env.example index ec6c76d..56c0eb0 100644 --- a/.env.example +++ b/.env.example @@ -1,76 +1,85 @@ -# ============================================================ +# ═══════════════════════════════════════════════════════════════ # Imago — Configuration -# Copier ce fichier en .env et remplir les valeurs -# ============================================================ +# ═══════════════════════════════════════════════════════════════ +# Copiez ce fichier en .env et ajustez les valeurs. +# cp .env.example .env -# Application -APP_NAME="Imago" -APP_VERSION="1.0.0" -DEBUG=true -SECRET_KEY="changez-moi-en-production-avec-une-cle-aleatoire-longue" +# ── Application ────────────────────────────────────────────── +APP_NAME=Imago +APP_VERSION=2.0.0 +DEBUG=false +SECRET_KEY=changez-moi-par-une-valeur-longue-et-aleatoire -# AI — Configuration -AI_ENABLED=true - -# AI — Provider (gemini/openrouter) -AI_PROVIDER="openrouter" - -# Serveur +# ── Serveur ────────────────────────────────────────────────── HOST=0.0.0.0 PORT=8000 -# Base de données -DATABASE_URL="postgresql+asyncpg://imago:imago@db:5432/imago" -# Modifiez les valeurs ci-dessus si vous utilisez une instance externe ou locale. -# Pour SQLite (développement local sans Docker): -# DATABASE_URL="sqlite+aiosqlite:///./data/imago.db" +# ── Base de données ────────────────────────────────────────── +# SQLite (dev): sqlite+aiosqlite:///./data/imago.db +# PostgreSQL: postgresql+asyncpg://user:pass@host:5432/imago +DATABASE_URL=sqlite+aiosqlite:///./data/imago.db -# Redis (ARQ Worker) -REDIS_URL="redis://redis:6379/0" - -# Stockage des fichiers -UPLOAD_DIR="./data/uploads" -THUMBNAILS_DIR="./data/thumbnails" +# ── Stockage ───────────────────────────────────────────────── +# local (fichiers sur disque) ou s3 (MinIO / AWS S3 / Cloudflare R2) +STORAGE_BACKEND=local +UPLOAD_DIR=./data/uploads +THUMBNAILS_DIR=./data/thumbnails MAX_UPLOAD_SIZE_MB=50 -# AI — Google Gemini -GEMINI_API_KEY="AIza..." -GEMINI_MODEL="gemini-3.1-pro-preview" +# ── Stockage S3 (si STORAGE_BACKEND=s3) ────────────────────── +S3_BUCKET=imago +S3_PREFIX= +S3_REGION=us-east-1 +S3_ENDPOINT_URL=http://localhost:9000 +S3_ENDPOINT_PUBLIC=http://localhost:9000 +S3_ACCESS_KEY=minioadmin +S3_SECRET_KEY=minioadmin +SIGNED_URL_SECRET=changez-moi-en-prod + +# ── AI ─────────────────────────────────────────────────────── +AI_ENABLED=true +AI_PROVIDER=openrouter + +# Google Gemini +GEMINI_API_KEY= +GEMINI_MODEL=gemini-3.1-pro-preview GEMINI_MAX_TOKENS=1024 -# AI - Openrouter -# model name : mistralai/mistral-small-3.1-24b-instruct:free -# model name : google/gemini-2.0-flash-001 -OPENROUTER_API_KEY="..." -OPENROUTER_MODEL="qwen/qwen2.5-vl-72b-instruct" +# OpenRouter +OPENROUTER_API_KEY= +OPENROUTER_MODEL=qwen/qwen2.5-vl-72b-instruct -# AI — Comportement +# AI Comportement AI_TAGS_MIN=5 AI_TAGS_MAX=10 -AI_DESCRIPTION_LANGUAGE="français" +AI_DESCRIPTION_LANGUAGE=francais AI_CACHE_DAYS=30 +AI_REQUEST_TIMEOUT=60 +AI_MAX_RETRIES=2 -# OCR +# ── OCR ────────────────────────────────────────────────────── OCR_ENABLED=true -TESSERACT_CMD="/usr/bin/tesseract" -OCR_LANGUAGES="fra+eng" +TESSERACT_CMD=/usr/bin/tesseract +OCR_LANGUAGES=fra+eng -# CORS -CORS_ORIGINS=["http://localhost:3000","http://localhost:8080","http://localhost:5173"] +# ── CORS ───────────────────────────────────────────────────── +# Format JSON : ["http://localhost:3000", "http://localhost:5173"] +CORS_ORIGINS=["http://localhost:3000", "http://localhost:5173"] -# Authentification -ADMIN_API_KEY="" -JWT_SECRET_KEY="changez-moi-jwt-secret-en-production" -JWT_ALGORITHM="HS256" +# ── Authentification ───────────────────────────────────────── +ADMIN_API_KEY= +JWT_SECRET_KEY=changez-moi-jwt-secret +JWT_ALGORITHM=HS256 -# Rate Limiting (requêtes par minute — legacy) -RATE_LIMIT_UPLOAD=10 -RATE_LIMIT_AI=20 - -# Rate Limiting par plan (requêtes par heure) +# ── Rate Limiting ──────────────────────────────────────────── RATE_LIMIT_FREE_UPLOAD=20 RATE_LIMIT_FREE_AI=50 RATE_LIMIT_STANDARD_UPLOAD=100 RATE_LIMIT_STANDARD_AI=200 RATE_LIMIT_PREMIUM_UPLOAD=500 RATE_LIMIT_PREMIUM_AI=1000 +# Redis pour persistence des compteurs (optionnel, mémoire par défaut) +# Format: redis://host:***@RATE_LIMIT_STORAGE_URL= + +# ── Redis + ARQ Worker ─────────────────────────────────────── +REDIS_URL=redis://localhost:***@PIPELINE_TIMEOUT=300 diff --git a/CHANGELOG.md b/CHANGELOG.md index 4950c01..95fddf1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,7 +4,45 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). -## [Unreleased] +## [2.0.0] - 2026-06-22 + +### Security +- API Key hashing now uses SHA-256 + pepper (SECRET_KEY) with timing-safe comparison via `secrets.compare_digest()` +- CORS validation rejects `["*"]` when `allow_credentials=True` (Pydantic model_validator) +- Scope validation: `ClientCreate.scopes` rejects invalid values via `field_validator` +- Master key comparison uses constant-time compare + +### Added +- **Dead Letter Queue (DLQ)**: Failed ARQ jobs pushed to Redis DLQ after `PIPELINE_MAX_RETRIES` attempts. Admin endpoints to list/retry/clear dead jobs (`GET/POST/DELETE /admin/api/queue/dead`) +- **Circuit Breaker AI**: All AI calls wrapped with `asyncio.wait_for()` + exponential backoff retry (configurable via `AI_REQUEST_TIMEOUT`, `AI_MAX_RETRIES`) +- **Request ID Middleware**: `X-Request-ID` header injected in all responses + bound to structlog context +- **Rate Limiting per Plan**: Dynamic rate limits based on client plan (free/standard/premium) via ContextVar +- **Redis Rate Limit Storage**: Optional persistence via `RATE_LIMIT_STORAGE_URL` +- **Worker script**: `worker.py` for standalone ARQ worker launch +- **Docker Compose stack**: PostgreSQL 16 + Redis 7 + MinIO + API + Worker +- **`.env.example`**: All configuration variables documented +- **Grafana Dashboard**: `docs/grafana-dashboard.json` with 10 panels (images, tokens, pipeline, storage, WebSocket, errors) +- **Database indexes**: Composite indexes `(client_id, uploaded_at)` and `(client_id, processing_status)` on images table + +### Changed +- `APP_VERSION` aligned to `2.0.0` in config (was `1.0.0`) +- `delete_files()` now fully async — uses `await backend.delete()` instead of fire-and-forget `ensure_future()` +- Upload flow: files deleted on DB commit failure (no orphaned S3 files) +- Gemini client: cached with TTL — recreated if `GEMINI_API_KEY` changes +- `parse_redis_url()` extracted to `app/workers/redis_client.py` (DRY — was duplicated in main.py and image_worker.py) +- Rate limit key function encodes plan for per-plan bucket isolation +- `_FallbackArqPool` moved to `app/workers/arq_fallback.py` (avoids circular imports) +- ARQ worker: configurable `PIPELINE_TIMEOUT` and `PIPELINE_MAX_RETRIES`, `on_shutdown` handler +- Pipeline: `push_dead_job()` on final retry failure marks image as `ERROR` with DLQ metadata +- Admin panel: VITE_API_BASE_URL changed from internal Docker hostname to `http://localhost:8000` + +### Fixed +- ARQ fallback now returns explicit warning in `UploadResponse.message` instead of silently ignoring +- Removed obsolete comment in `ai_vision.py` about `StorageBackend.read` +- Removed signal handler instability in test fixtures +- Fixed Docker Compose admin connectivity (browser couldn't resolve internal `backend` hostname) + +## [1.0.0] - Previous ### Added (Phase 3: DX) - **WebSockets / Real-Time**: diff --git a/app/config.py b/app/config.py index 783bf96..a5208e5 100644 --- a/app/config.py +++ b/app/config.py @@ -4,14 +4,14 @@ Configuration centralisée — chargée depuis .env from pathlib import Path from typing import List from pydantic_settings import BaseSettings -from pydantic import field_validator +from pydantic import field_validator, model_validator import json class Settings(BaseSettings): # Application APP_NAME: str = "Imago" - APP_VERSION: str = "1.0.0" + APP_VERSION: str = "2.0.0" API_V1_SUNSET_DATE: str = "" DEBUG: bool = False SECRET_KEY: str = "changez-moi" @@ -44,9 +44,13 @@ class Settings(BaseSettings): # AI — Comportement AI_TAGS_MIN: int = 5 AI_TAGS_MAX: int = 10 - AI_DESCRIPTION_LANGUAGE: str = "français" + AI_DESCRIPTION_LANGUAGE: str = "fran\u00e7ais" AI_CACHE_DAYS: int = 30 + # AI — Résilience + AI_REQUEST_TIMEOUT: int = 60 + AI_MAX_RETRIES: int = 2 + # OCR OCR_ENABLED: bool = True TESSERACT_CMD: str = "/usr/bin/tesseract" @@ -59,6 +63,7 @@ class Settings(BaseSettings): ADMIN_API_KEY: str = "" JWT_SECRET_KEY: str = "changez-moi-jwt-secret" JWT_ALGORITHM: str = "HS256" + API_KEY_HASH_ROUNDS: int = 12 # Rate limiting — global (legacy) RATE_LIMIT_UPLOAD: int = 10 @@ -72,24 +77,25 @@ class Settings(BaseSettings): RATE_LIMIT_PREMIUM_UPLOAD: int = 500 RATE_LIMIT_PREMIUM_AI: int = 1000 - # Redis + ARQ Worker - REDIS_URL: str = "redis://localhost:6379" - WORKER_MAX_JOBS: int = 10 - WORKER_JOB_TIMEOUT: int = 180 - WORKER_MAX_TRIES: int = 3 - AI_STEP_TIMEOUT: int = 120 - OCR_STEP_TIMEOUT: int = 30 + # Rate limiting — stockage Redis (optionnel, mémoire par défaut) + RATE_LIMIT_STORAGE_URL: str = "" - # Storage Backend - STORAGE_BACKEND: str = "local" # "local" | "s3" - S3_BUCKET: str = "" + # Redis + ARQ Worker + REDIS_URL: str = "redis://localhost:6379/0" + PIPELINE_TIMEOUT: int = 300 # secondes — timeout max d'un job pipeline + PIPELINE_MAX_RETRIES: int = 3 # tentatives avant Dead Letter Queue + REDIS_POOL_MAX_CONNECTIONS: int = 20 + + # Storage backends + STORAGE_BACKEND: str = "local" + S3_BUCKET: str = "imago" + S3_PREFIX: str = "" S3_REGION: str = "us-east-1" - S3_ENDPOINT_URL: str = "" # vide = AWS, sinon MinIO/R2 - S3_ENDPOINT_PUBLIC: str = "" # URL utilisée pour les liens présignés (ex: http://localhost:9000) + S3_ENDPOINT_URL: str = "" S3_ACCESS_KEY: str = "" S3_SECRET_KEY: str = "" - S3_PREFIX: str = "imago" - SIGNED_URL_SECRET: str = "changez-moi-signed-url" + S3_ENDPOINT_PUBLIC: str = "" + SIGNED_URL_SECRET: str = "changez-moi-signed-url-secret" @field_validator("CORS_ORIGINS", mode="before") @classmethod @@ -101,6 +107,16 @@ class Settings(BaseSettings): return [v] return v + @model_validator(mode="after") + def validate_cors_credentials(self): + """Empêche allow_credentials=True avec origins=['*'].""" + if self.CORS_ORIGINS == ["*"] or "*" in self.CORS_ORIGINS: + raise ValueError( + "CORS_ORIGINS ne peut pas contenir '*' quand " + "allow_credentials=True. Spécifiez les origines explicitement." + ) + return self + @property def upload_path(self) -> Path: p = Path(self.UPLOAD_DIR) diff --git a/app/database.py b/app/database.py index e28ef9a..47ef8d2 100644 --- a/app/database.py +++ b/app/database.py @@ -44,8 +44,8 @@ async def get_db() -> AsyncSession: async def init_db(): """Crée toutes les tables et initialise un client par défaut si nécessaire.""" import secrets - import hashlib from sqlalchemy import select + from app.dependencies.auth import hash_api_key from app.models.client import APIClient, ClientPlan async with engine.begin() as conn: @@ -59,7 +59,7 @@ async def init_db(): if result.scalar_one_or_none() is None: # Table vide -> Création du client bootstrap raw_key = secrets.token_urlsafe(32) - key_hash = hashlib.sha256(raw_key.encode("utf-8")).hexdigest() + key_hash = hash_api_key(raw_key) bootstrap_client = APIClient( name="Default Admin", diff --git a/app/dependencies/auth.py b/app/dependencies/auth.py index ddb96b5..6012b35 100644 --- a/app/dependencies/auth.py +++ b/app/dependencies/auth.py @@ -1,27 +1,42 @@ """ Dépendances FastAPI — authentification par API Key + vérification de scopes. +Sécurité : +- Hashage SHA-256 avec pepper (SECRET_KEY) pour le stockage +- Comparaison timing-safe via secrets.compare_digest() +- Injection du plan client dans le ContextVar pour le rate limiting dynamique +- Les clés API sont des tokens 256-bit aléatoires (secrets.token_urlsafe(32)) + Usage dans les routers : client = Depends(get_current_client) _ = Depends(require_scope("images:read")) """ import hashlib import logging +import secrets as sec from typing import Callable from fastapi import Depends, Header, HTTPException, Request, status from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession +from app.config import settings from app.database import get_db from app.models.client import APIClient +from app.middleware import _current_client_plan logger = logging.getLogger(__name__) def hash_api_key(api_key: str) -> str: - """Hash SHA-256 d'une clé API — fonction utilitaire réutilisable.""" - return hashlib.sha256(api_key.encode("utf-8")).hexdigest() + """ + Hash une clé API avec SHA-256 + pepper (SECRET_KEY). + + Le pepper empêche les attaques par rainbow table même si la BDD fuit. + Les clés API sont des tokens aléatoires 256-bit → SHA-256 est suffisant. + """ + peppered = f"{api_key}:{settings.SECRET_KEY}" + return hashlib.sha256(peppered.encode("utf-8")).hexdigest() async def verify_api_key( @@ -41,16 +56,17 @@ async def verify_api_key( """ Vérifie la clé API fournie dans le header Authorization ou X-API-Key. Injecte client_id et client_plan dans request.state pour le rate limiter. + Injecte le plan dans le ContextVar pour le rate limiting dynamique. Raises: HTTPException 401: clé absente, invalide ou client inactif. """ raw_key = None - # ── 1. Tentative avec Authorization: Bearer ──────── + # ── 1. Tentative avec Authorization: Bearer *** ──────── if authorization and authorization.startswith("Bearer "): raw_key = authorization[7:].strip() - + # ── 2. Tentative avec X-API-Key ────────────────────────── if not raw_key and x_api_key: raw_key = x_api_key.strip() @@ -63,15 +79,21 @@ async def verify_api_key( ) # ── 1.5. Vérification Master Key (ADMIN_API_KEY) ─────── - from app.config import settings - if settings.ADMIN_API_KEY and raw_key == settings.ADMIN_API_KEY: - # Retourne un client virtuel avec tous les droits - return APIClient( - id="admin-master", - name="Imago Master Admin", - scopes=["admin", "images:read", "images:write", "ai:use"], - plan="premium" - ) + # Comparaison timing-safe de la master key + if settings.ADMIN_API_KEY: + master_hash = hash_api_key(settings.ADMIN_API_KEY) + candidate_hash = hash_api_key(raw_key) + if sec.compare_digest(candidate_hash, master_hash): + client = APIClient( + id="admin-master", + name="Imago Master Admin", + scopes=["admin", "images:read", "images:write", "ai:use"], + plan="premium", + ) + request.state.client_id = client.id + request.state.client_plan = "premium" + _current_client_plan.set("premium") + return client # ── Lookup par hash ─────────────────────────────────────── key_hash = hash_api_key(raw_key) @@ -81,18 +103,18 @@ async def verify_api_key( client = result.scalar_one_or_none() if client is None: - logger.warning("Tentative d'authentification avec une clé invalide") + logger.warning("auth.invalid_key") raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, - detail="Authentification requise", + detail="Clé API invalide", headers={"WWW-Authenticate": "Bearer"}, ) if not client.is_active: - logger.warning("Tentative d'authentification avec un client inactif: %s", client.id) + logger.warning("auth.inactive_client", extra={"client_id": client.id}) raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, - detail="Authentification requise", + detail="Client désactivé", headers={"WWW-Authenticate": "Bearer"}, ) @@ -100,6 +122,9 @@ async def verify_api_key( request.state.client_id = client.id request.state.client_plan = client.plan.value if client.plan else "free" + # Injecter dans le ContextVar pour le rate limiting dynamique + _current_client_plan.set(client.plan.value if client.plan else "free") + return client @@ -120,12 +145,12 @@ def require_scope(scope: str) -> Callable: ) -> APIClient: if not client.has_scope(scope): logger.warning( - "Client %s (%s) a tenté d'accéder au scope '%s' sans autorisation", - client.id, client.name, scope, + "auth.scope_denied", + extra={"client_id": client.id, "scope": scope}, ) raise HTTPException( status_code=status.HTTP_403_FORBIDDEN, - detail="Permission insuffisante", + detail=f"Permission insuffisante : scope '{scope}' requis", ) return client diff --git a/app/main.py b/app/main.py index f73b531..e151b7d 100644 --- a/app/main.py +++ b/app/main.py @@ -3,7 +3,7 @@ Imago — Application principale FastAPI """ import logging from contextlib import asynccontextmanager -from fastapi import FastAPI, Request +from fastapi import FastAPI, Request, APIRouter from fastapi.middleware.cors import CORSMiddleware from fastapi.staticfiles import StaticFiles from slowapi import _rate_limit_exceeded_handler @@ -16,7 +16,8 @@ from app.routers import images_router, ai_router, auth_router, files_router, ws_ from app.middleware import limiter from app.middleware.logging_middleware import LoggingMiddleware from app.middleware.versioning import APIVersioningMiddleware -from app.workers.redis_client import get_redis_pool, close_redis_pool +from app.workers.redis_client import get_redis_pool, close_redis_pool, parse_redis_url +from app.workers.arq_fallback import _FallbackArqPool # Configure le logging structuré dès l'import configure_logging(debug=settings.DEBUG) @@ -38,41 +39,10 @@ except ImportError: def _arq_redis_settings() -> "RedisSettings": - """Parse REDIS_URL en RedisSettings ARQ.""" - url = settings.REDIS_URL - if url.startswith("redis://"): - url = url[8:] - elif url.startswith("rediss://"): - url = url[9:] - - password = None - host = "localhost" - port = 6379 - database = 0 - - if "@" in url: - auth_part, url = url.rsplit("@", 1) - if ":" in auth_part: - password = auth_part.split(":", 1)[1] - else: - password = auth_part - - if "/" in url: - host_port, db_str = url.split("/", 1) - if db_str: - database = int(db_str) - else: - host_port = url - - if ":" in host_port: - host, port_str = host_port.rsplit(":", 1) - if port_str: - port = int(port_str) - else: - host = host_port - + """Parse REDIS_URL en RedisSettings ARQ via le parser centralisé.""" + host, port, password, database = parse_redis_url() return RedisSettings( - host=host or "localhost", + host=host, port=port, password=password, database=database, @@ -85,11 +55,9 @@ def _arq_redis_settings() -> "RedisSettings": @asynccontextmanager async def lifespan(app: FastAPI): - # Création des répertoires de données settings.upload_path settings.thumbnails_path - # Initialisation de la base de données (création des tables) await init_db() active_model = settings.OPENROUTER_MODEL if settings.AI_PROVIDER == "openrouter" else settings.GEMINI_MODEL @@ -101,7 +69,6 @@ async def lifespan(app: FastAPI): "ocr_enabled": settings.OCR_ENABLED, }) - # Initialisation Redis + ARQ pool try: app.state.redis = await get_redis_pool() logger.info("startup.redis_connected", extra={"url": settings.REDIS_URL}) @@ -122,24 +89,12 @@ async def lifespan(app: FastAPI): yield - # Fermeture propre if hasattr(app.state, "arq_pool") and hasattr(app.state.arq_pool, "close"): await app.state.arq_pool.close() await close_redis_pool() logger.info("shutdown.complete") -class _FallbackArqPool: - """Fallback quand Redis/ARQ n'est pas disponible.""" - - async def enqueue_job(self, *args, **kwargs): - logger.warning("arq.fallback_enqueue", extra={"args": str(args)}) - return None - - async def close(self): - pass - - # ───────────────────────────────────────────────────────────── # Application # ───────────────────────────────────────────────────────────── @@ -152,31 +107,24 @@ app = FastAPI( Backend de gestion d'images et fonctionnalités AI pour l'interface Shaarli. -> **Note de versionnement :** L'API principale est désormais servie sous le préfixe `/api/v1/`. Les anciennes routes sans ce préfixe sont dépréciées. +> **Note de versionnement :** L'API principale est désormais servie sous le préfixe `/api/v1/`. ### Fonctionnalités -- 📸 **Upload et stockage d'images** avec génération de thumbnails -- 🔍 **Extraction EXIF** automatique (appareil, GPS, paramètres de prise de vue) -- 📝 **OCR** — extraction de texte depuis les images (Tesseract) -- 🤖 **Vision AI** — description et classification par tags (Gemini) -- 🔗 **Résumé d'URL** — scraping + résumé AI de pages web -- ✅ **Rédaction de tâches** — génération structurée via AI -- 📋 **File de tâches ARQ** — pipeline persistant avec retry automatique -- 📊 **Métriques Prometheus** — /metrics endpoint - -### Pipeline de traitement -Chaque image uploadée est automatiquement traitée via ARQ (Redis) : -`EXIF → OCR → Vision AI → stockage BDD` +- Upload et stockage d'images avec génération de thumbnails +- Extraction EXIF automatique (appareil, GPS, paramètres de prise de vue) +- OCR — extraction de texte depuis les images (Tesseract) +- Vision AI — description et classification par tags +- Résumé d'URL — scraping + résumé AI de pages web +- Rédaction de tâches — génération structurée via AI +- File de tâches ARQ — pipeline persistant avec retry automatique +- Métriques Prometheus — /metrics endpoint """, lifespan=lifespan, docs_url="/docs", redoc_url="/redoc", ) -# ───────────────────────────────────────────────────────────── -# Middleware CORS -# ───────────────────────────────────────────────────────────── - +# ── Middleware CORS ───────────────────────────────────────── app.add_middleware( CORSMiddleware, allow_origins=settings.CORS_ORIGINS, @@ -185,29 +133,21 @@ app.add_middleware( allow_headers=["*"], ) -# ───────────────────────────────────────────────────────────── -# Middleware de Versionnement API -# ───────────────────────────────────────────────────────────── +# ── Request ID (trace_id) ─────────────────────────────────── +from app.middleware.request_id import RequestIDMiddleware +app.add_middleware(RequestIDMiddleware) +# ── API Versioning ────────────────────────────────────────── app.add_middleware(APIVersioningMiddleware) -# ───────────────────────────────────────────────────────────── -# Middleware Logging HTTP -# ───────────────────────────────────────────────────────────── - +# ── Logging HTTP ──────────────────────────────────────────── app.add_middleware(LoggingMiddleware) -# ───────────────────────────────────────────────────────────── -# Rate Limiting (slowapi) -# ───────────────────────────────────────────────────────────── - +# ── Rate Limiting ─────────────────────────────────────────── app.state.limiter = limiter app.add_exception_handler(RateLimitExceeded, _rate_limit_exceeded_handler) -# ───────────────────────────────────────────────────────────── -# Prometheus Metrics -# ───────────────────────────────────────────────────────────── - +# ── Prometheus Metrics ────────────────────────────────────── if _prometheus_available: Instrumentator( should_group_status_codes=True, @@ -215,39 +155,25 @@ if _prometheus_available: excluded_handlers=["/health", "/health/detailed", "/metrics"], ).instrument(app).expose(app, endpoint="/metrics", tags=["Observabilité"]) -# ───────────────────────────────────────────────────────────── -# Fichiers statiques / URLs signées -# ───────────────────────────────────────────────────────────── - +# ── Fichiers statiques ────────────────────────────────────── if settings.STORAGE_BACKEND == "local": app.include_router(files_router) app.mount("/static/uploads", StaticFiles(directory=str(settings.upload_path)), name="uploads") app.mount("/static/thumbnails", StaticFiles(directory=str(settings.thumbnails_path)), name="thumbnails") -# ───────────────────────────────────────────────────────────── -# Routers Versionnés (/api/v1/) -# ───────────────────────────────────────────────────────────── -from fastapi import APIRouter - +# ── Routers versionnés (/api/v1/) ─────────────────────────── api_v1_router = APIRouter(prefix="/api/v1") api_v1_router.include_router(images_router) api_v1_router.include_router(ai_router) api_v1_router.include_router(auth_router) - -# Mount the versioned router on the app app.include_router(api_v1_router) -# ───────────────────────────────────────────────────────────── -# Routers Non-Versionnés (WebSockets, Fichiers, Santé, Metrics) -# ───────────────────────────────────────────────────────────── - +# ── Routers non-versionnés ────────────────────────────────── app.include_router(ws_router) app.include_router(admin_router) -# ───────────────────────────────────────────────────────────── -# Routes utilitaires -# ───────────────────────────────────────────────────────────── +# ── Routes utilitaires ────────────────────────────────────── @app.get("/", tags=["Santé"]) async def root(): @@ -266,7 +192,6 @@ async def health(): (settings.AI_PROVIDER == "openrouter" and bool(settings.OPENROUTER_API_KEY)) ) active_model = settings.OPENROUTER_MODEL if settings.AI_PROVIDER == "openrouter" else settings.GEMINI_MODEL - return { "status": "healthy", "ai_enabled": settings.AI_ENABLED, @@ -282,14 +207,12 @@ async def health_detailed(request: Request): """Endpoint de santé détaillé pour monitoring avancé.""" checks = {} - # 1. Backend check checks["backend"] = { "status": "ok", "version": settings.APP_VERSION, - "uptime": "N/A" # Optionnel: ajouter un tracker d'uptime si nécessaire + "uptime": "N/A", } - # 2. Database check try: from app.database import AsyncSessionLocal from sqlalchemy import text @@ -299,7 +222,6 @@ async def health_detailed(request: Request): except Exception as e: checks["database"] = {"status": "error", "error": str(e)} - # 3. Redis check redis = getattr(request.app.state, "redis", None) if redis: try: @@ -310,29 +232,17 @@ async def health_detailed(request: Request): else: checks["redis"] = {"status": "not_configured"} - # 4. Worker & Queue check (ARQ) arq_pool = getattr(request.app.state, "arq_pool", None) if arq_pool and not isinstance(arq_pool, _FallbackArqPool): try: - # On essaye de voir le nombre de jobs en attente queued_jobs = await arq_pool.zcard("arq:queue") - checks["queue"] = { - "status": "ok", - "pending_jobs": queued_jobs - } - - # Pour le worker, on vérifie la clé de health check déposée par ARQ - # Par défaut c'est :health-check - hc_key = f"{settings.S3_PREFIX or 'arq'}:queue:health-check" - # On ré-essaye avec la queue par défaut si le préfixe est différent - if not await redis.exists(hc_key): - hc_key = "arq:queue:health-check" - - worker_active = await redis.exists(hc_key) + checks["queue"] = {"status": "ok", "pending_jobs": queued_jobs} + hc_key = "arq:queue:health-check" + worker_active = await redis.exists(hc_key) if redis else False checks["worker"] = { "status": "ok" if worker_active else "error", "active_workers": 1 if worker_active else 0, - "detail": "Worker opérationnel" if worker_active else "Aucun worker détecté" + "detail": "Worker opérationnel" if worker_active else "Aucun worker détecté", } except Exception as e: checks["queue"] = {"status": "error", "error": str(e)} @@ -341,33 +251,19 @@ async def health_detailed(request: Request): checks["queue"] = {"status": "fallback"} checks["worker"] = {"status": "error", "detail": "ARQ non disponible"} - # 5. MinIO / Storage check from app.services.storage_backend import get_storage_backend, S3Storage storage = get_storage_backend() try: if isinstance(storage, S3Storage): - # On vérifie la connexion réelle à MinIO en listant les objets (ou head bucket) - # Pour faire simple/rapide on vérifie juste si on peut accéder au bucket - # via l'interface du storage qui a déjà les sessions asynchrones - exists = await storage.exists(".health-check") # path arbitraire - checks["minio"] = { - "status": "ok", - "backend": "s3", - "bucket": settings.S3_BUCKET - } + await storage.exists(".health-check") + checks["minio"] = {"status": "ok", "backend": "s3", "bucket": settings.S3_BUCKET} else: - checks["minio"] = { - "status": "ok", - "backend": "local", - "path": settings.UPLOAD_DIR - } + checks["minio"] = {"status": "ok", "backend": "local", "path": settings.UPLOAD_DIR} except Exception as e: checks["minio"] = {"status": "error", "error": str(e)} - # 6. OCR check (Tesseract) checks["tesseract"] = {"status": "ok" if settings.OCR_ENABLED else "disabled"} - # 7. AI Provider check ai_configured = ( (settings.AI_PROVIDER == "gemini" and bool(settings.GEMINI_API_KEY)) or (settings.AI_PROVIDER == "openrouter" and bool(settings.OPENROUTER_API_KEY)) @@ -375,16 +271,12 @@ async def health_detailed(request: Request): checks["ai"] = { "status": "ok" if ai_configured else "warning", "provider": settings.AI_PROVIDER, - "configured": ai_configured + "configured": ai_configured, } overall = "healthy" if all( c.get("status") in ("ok", "enabled", "disabled", "not_configured", "fallback") - for k, c in checks.items() if k != "ai" # On tolère AI warning + for k, c in checks.items() if k != "ai" ) else "degraded" - return { - "status": overall, - "checks": checks, - "version": settings.APP_VERSION, - } + return {"status": overall, "checks": checks, "version": settings.APP_VERSION} diff --git a/app/middleware/__init__.py b/app/middleware/__init__.py index d8cc5e2..18a3a43 100644 --- a/app/middleware/__init__.py +++ b/app/middleware/__init__.py @@ -8,8 +8,12 @@ Limites par plan (par heure) : - free : 20 uploads, 50 AI - standard : 100 uploads, 200 AI - premium : 500 uploads, 1000 AI + +Supporte le stockage Redis pour la persistence des compteurs. """ import logging +from contextvars import ContextVar + from slowapi import Limiter from slowapi.util import get_remote_address from starlette.requests import Request @@ -18,21 +22,31 @@ from app.config import settings logger = logging.getLogger(__name__) +# ContextVar pour transmettre le plan du client au rate limiter dynamique +_current_client_plan: ContextVar[str] = ContextVar("current_client_plan", default="free") + def _get_client_id_from_request(request: Request) -> str: """ Extrait le client_id depuis la state de la requête. + Encode le plan pour isoler les buckets par plan. Fallback vers l'IP si le client n'est pas encore authentifié. """ - # Le client_id est injecté par le middleware ou la dépendance auth client_id = getattr(request.state, "client_id", None) + plan = getattr(request.state, "client_plan", "free") if client_id: - return str(client_id) + return f"{client_id}:{plan}" return get_remote_address(request) -# Instance globale du limiter -limiter = Limiter(key_func=_get_client_id_from_request) +# Stockage Redis pour les compteurs (si configuré), sinon mémoire +if settings.RATE_LIMIT_STORAGE_URL: + limiter = Limiter( + key_func=_get_client_id_from_request, + storage_uri=settings.RATE_LIMIT_STORAGE_URL, + ) +else: + limiter = Limiter(key_func=_get_client_id_from_request) def get_upload_rate_limit(plan: str) -> str: @@ -55,11 +69,13 @@ def get_ai_rate_limit(plan: str) -> str: return limits.get(plan, limits["free"]) -def upload_rate_limit_key(request: Request) -> str: - """Clé dynamique pour le rate limiting des uploads.""" - return _get_client_id_from_request(request) +def dynamic_upload_limit() -> str: + """Callable pour slowapi — lit le plan depuis le ContextVar.""" + plan = _current_client_plan.get() + return get_upload_rate_limit(plan) -def ai_rate_limit_key(request: Request) -> str: - """Clé dynamique pour le rate limiting des endpoints AI.""" - return _get_client_id_from_request(request) +def dynamic_ai_limit() -> str: + """Callable pour slowapi — lit le plan depuis le ContextVar.""" + plan = _current_client_plan.get() + return get_ai_rate_limit(plan) diff --git a/app/models/image.py b/app/models/image.py index 7663bb9..5c5902b 100644 --- a/app/models/image.py +++ b/app/models/image.py @@ -5,7 +5,7 @@ import enum from datetime import datetime, timezone from sqlalchemy import ( Column, Integer, String, Text, DateTime, - JSON, Float, Enum as SAEnum, BigInteger, Boolean, ForeignKey + JSON, Float, Enum as SAEnum, BigInteger, Boolean, ForeignKey, Index ) from sqlalchemy.orm import relationship from app.database import Base @@ -31,11 +31,11 @@ class Image(Base): # ── Fichier ─────────────────────────────────────────────── original_name = Column(String(512), nullable=False) - filename = Column(String(512), nullable=False) # nom sur disque (uuid-based) + filename = Column(String(512), nullable=False) file_path = Column(String(1024), nullable=False) thumbnail_path = Column(String(1024)) mime_type = Column(String(128)) - file_size = Column(BigInteger) # bytes + file_size = Column(BigInteger) width = Column(Integer) height = Column(Integer) uploaded_at = Column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc)) @@ -52,18 +52,18 @@ class Image(Base): processing_done_at = Column(DateTime(timezone=True)) # ── Métadonnées EXIF ────────────────────────────────────── - exif_raw = Column(JSON) # dict complet brut - exif_make = Column(String(256)) # Appareil — fabricant - exif_model = Column(String(256)) # Appareil — modèle + exif_raw = Column(JSON) + exif_make = Column(String(256)) + exif_model = Column(String(256)) exif_lens = Column(String(256)) - exif_taken_at = Column(DateTime(timezone=True)) # DateTimeOriginal EXIF + exif_taken_at = Column(DateTime(timezone=True)) exif_gps_lat = Column(Float) exif_gps_lon = Column(Float) exif_altitude = Column(Float) exif_iso = Column(Integer) - exif_aperture = Column(String(32)) # ex: "f/2.8" - exif_shutter = Column(String(32)) # ex: "1/250" - exif_focal = Column(String(32)) # ex: "50mm" + exif_aperture = Column(String(32)) + exif_shutter = Column(String(32)) + exif_focal = Column(String(32)) exif_flash = Column(Boolean) exif_orientation = Column(Integer) exif_software = Column(String(256)) @@ -71,18 +71,26 @@ class Image(Base): # ── OCR ─────────────────────────────────────────────────── ocr_text = Column(Text) ocr_language = Column(String(64)) - ocr_confidence = Column(Float) # 0.0 – 1.0 + ocr_confidence = Column(Float) ocr_has_text = Column(Boolean, default=False) # ── AI Vision ───────────────────────────────────────────── ai_description = Column(Text) - ai_tags = Column(JSON) # ["nature", "paysage", ...] - ai_confidence = Column(Float) # score de confiance global + ai_tags = Column(JSON) + ai_confidence = Column(Float) ai_model_used = Column(String(128)) ai_processed_at = Column(DateTime(timezone=True)) ai_prompt_tokens = Column(Integer) ai_output_tokens = Column(Integer) + # ── Indexes composites ──────────────────────────────────── + __table_args__ = ( + # Pour les listings paginés par client (requête la plus fréquente) + Index("ix_images_client_uploaded", "client_id", "uploaded_at"), + # Pour le filtrage par statut + client + Index("ix_images_client_status", "client_id", "processing_status"), + ) + def __repr__(self): return f"" diff --git a/app/routers/admin.py b/app/routers/admin.py index 69b30f1..aa57e4d 100644 --- a/app/routers/admin.py +++ b/app/routers/admin.py @@ -14,6 +14,7 @@ router = APIRouter(prefix="/admin/api", tags=["Admin"]) PROJECT_ROOT = Path(__file__).resolve().parent.parent.parent ALLOWED_DOCS = { "README.md": PROJECT_ROOT / "README.md", + "ROADMAP.md": PROJECT_ROOT / "docs" / "ROADMAP.md", "USER_GUIDE.md": PROJECT_ROOT / "docs" / "USER_GUIDE.md", "ARCHITECTURE.md": PROJECT_ROOT / "docs" / "ARCHITECTURE.md", "API_GUIDE.md": PROJECT_ROOT / "docs" / "API_GUIDE.md", @@ -82,21 +83,60 @@ async def get_queue_status( request: Request, _: APIClient = Depends(require_scope("admin")) ) -> Dict[str, Any]: - """Get ARQ queue status using the redis connection.""" + """Get ARQ queue status including Dead Letter Queue.""" redis = request.app.state.redis - - # Note: A real implementation would parse the ARQ keys here. - # For now, we return basic statistics by probing redis directly with arq known queues. - # Count pending jobs in arq:queue + from app.workers.dead_letter import get_dead_job_count + pending_count = await redis.llen("arq:queue") if hasattr(redis, "llen") else 0 - # In ARQ, active jobs are harder to count without querying the worker sets, - # so we return pending count. - + dead_count = await get_dead_job_count(redis) + return { "pending_jobs": pending_count, - "status": "active" + "dead_jobs": dead_count, + "status": "active", } + +# ───────────────────────────────────────────────────────────── +# Dead Letter Queue +# ───────────────────────────────────────────────────────────── + +@router.get("/queue/dead") +async def list_dead_jobs( + request: Request, + _: APIClient = Depends(require_scope("admin")), + limit: int = 50, +) -> list[dict]: + """List dead jobs from the Dead Letter Queue.""" + from app.workers.dead_letter import get_dead_jobs + redis = request.app.state.redis + return await get_dead_jobs(redis, limit=limit) + + +@router.post("/queue/dead/{index}/retry") +async def retry_dead_job( + request: Request, + index: int, + _: APIClient = Depends(require_scope("admin")), +) -> dict: + """Retry a dead job by its index (0 = most recent).""" + from app.workers.dead_letter import retry_dead_job + redis = request.app.state.redis + arq_pool = request.app.state.arq_pool + return await retry_dead_job(redis, arq_pool, index) + + +@router.delete("/queue/dead") +async def clear_dead_jobs( + request: Request, + _: APIClient = Depends(require_scope("admin")), +) -> Dict[str, Any]: + """Clear all dead jobs from the DLQ.""" + from app.workers.dead_letter import clear_dead_jobs + redis = request.app.state.redis + count = await clear_dead_jobs(redis) + return {"cleared": count} + @router.post("/clients/{client_id}/toggle") async def toggle_client( client_id: str, diff --git a/app/routers/images.py b/app/routers/images.py index 6455772..a7e718e 100644 --- a/app/routers/images.py +++ b/app/routers/images.py @@ -23,8 +23,9 @@ from app.schemas import ( TagsResponse, ReprocessResponse, ) from app.services import storage -from app.middleware import limiter, get_upload_rate_limit +from app.middleware import limiter, dynamic_upload_limit from app.workers.image_worker import QUEUE_STANDARD, QUEUE_PREMIUM +from app.workers.arq_fallback import is_fallback_pool from app.metrics import ( hub_images_uploaded, hub_images_deleted, hub_storage_used_bytes, hub_arq_jobs_enqueued ) @@ -57,12 +58,6 @@ async def get_image_or_404( return image -def _dynamic_upload_limit(key: str) -> str: - """Retourne la limite dynamique basée sur le plan du client.""" - # On parse le plan depuis la clé ou le state — fallback free - return get_upload_rate_limit("free") - - # ───────────────────────────────────────────────────────────── # UPLOAD # ───────────────────────────────────────────────────────────── @@ -75,7 +70,7 @@ def _dynamic_upload_limit(key: str) -> str: description="Upload une image, lance automatiquement le pipeline AI (EXIF + OCR + Vision).", dependencies=[Depends(require_scope("images:write"))], ) -@limiter.limit("500/hour") +@limiter.limit(dynamic_upload_limit) async def upload_image( request: Request, file: UploadFile = File(...), @@ -102,23 +97,49 @@ async def upload_image( file_size = file_data.get("file_size", 0) client.storage_used_bytes = (client.storage_used_bytes or 0) + file_size - await db.commit() - await db.refresh(image) + try: + await db.commit() + await db.refresh(image) + except Exception: + # Rollback : supprimer les fichiers uploadés car la BDD n'a pas persisté + await storage.delete_files( + file_data["file_path"], + file_data.get("thumbnail_path"), + ) + raise HTTPException( + status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, + detail="Échec de l'enregistrement — veuillez réessayer", + ) - # Enqueue dans ARQ + # Enqueue dans ARQ — détecter le fallback arq_pool = request.app.state.arq_pool - await arq_pool.enqueue_job("process_image_task", image.id, str(client.id)) + job_enqueued = True + if is_fallback_pool(arq_pool): + job_enqueued = False + logger.warning("upload.arq_fallback", extra={ + "image_id": image.id, + "client_id": client.id, + }) + else: + await arq_pool.enqueue_job("process_image_task", image.id, str(client.id)) # Metrics hub_images_uploaded.labels(client_id=client.id).inc() hub_storage_used_bytes.labels(client_id=client.id).set(client.storage_used_bytes) - hub_arq_jobs_enqueued.labels(queue="default").inc() + if job_enqueued: + hub_arq_jobs_enqueued.labels(queue="default").inc() return UploadResponse( id=image.id, uuid=image.uuid, original_name=image.original_name, status=image.processing_status, + message=( + "Image uploadée — traitement AI en cours" + if job_enqueued + else "Image uploadée — le pipeline AI est indisponible (Redis/ARQ down). " + "Le traitement sera effectué automatiquement dès le rétablissement." + ), ) @@ -397,7 +418,7 @@ async def get_all_tags( description="Reprocess une image existante (utile après changement de modèle AI).", dependencies=[Depends(require_scope("images:write"))], ) -@limiter.limit("500/hour") +@limiter.limit(dynamic_upload_limit) async def reprocess_image( request: Request, image_id: int, @@ -444,8 +465,8 @@ async def delete_image( file_size = image.file_size or 0 client.storage_used_bytes = max(0, (client.storage_used_bytes or 0) - file_size) - # Suppression des fichiers sur disque - storage.delete_files(image.file_path, image.thumbnail_path) + # Suppression des fichiers — maintenant async et déterministe + await storage.delete_files(image.file_path, image.thumbnail_path) await db.delete(image) await db.commit() @@ -503,4 +524,3 @@ async def get_thumbnail_url( backend = get_storage_backend() url = await backend.get_signed_url(image.thumbnail_path, expires_in=expires_in) return {"url": url, "expires_in": expires_in} - diff --git a/app/schemas/auth.py b/app/schemas/auth.py index a1e566a..e6350bc 100644 --- a/app/schemas/auth.py +++ b/app/schemas/auth.py @@ -2,10 +2,19 @@ Schémas Pydantic — authentification et gestion des clients API """ from datetime import datetime -from typing import List, Optional -from pydantic import BaseModel, ConfigDict, Field +from typing import List, Optional, Set +from pydantic import BaseModel, ConfigDict, Field, field_validator from app.models.client import ClientPlan +# Scopes valides — toute valeur hors de cet ensemble est rejetée +VALID_SCOPES: Set[str] = { + "images:read", + "images:write", + "images:delete", + "ai:use", + "admin", +} + # ───────────────────────────────────────────────────────────── # Requêtes @@ -20,6 +29,17 @@ class ClientCreate(BaseModel): ) plan: ClientPlan = Field(default=ClientPlan.FREE, description="Plan tarifaire") + @field_validator("scopes") + @classmethod + def validate_scopes(cls, v): + invalid = set(v) - VALID_SCOPES + if invalid: + raise ValueError( + f"Scopes invalides : {invalid}. " + f"Scopes valides : {sorted(VALID_SCOPES)}" + ) + return v + class ClientUpdate(BaseModel): """Modifier un client API existant.""" diff --git a/app/services/ai_vision.py b/app/services/ai_vision.py index dc919cc..be9a5f7 100644 --- a/app/services/ai_vision.py +++ b/app/services/ai_vision.py @@ -1,11 +1,17 @@ """ -Service AI Vision — description, classification et tags via Google Gemini ou OpenRouter +Service AI Vision — description, classification et tags via Google Gemini ou OpenRouter. + +Résilience : +- Timeout configurable (AI_REQUEST_TIMEOUT) +- Retry avec backoff exponentiel (AI_MAX_RETRIES) +- Client Gemini non-singleton (recréé si la clé change) """ import asyncio import json import logging import re import base64 +import time import httpx import io from pathlib import Path @@ -19,14 +25,18 @@ from app.services.storage_backend import get_storage_backend logger = logging.getLogger(__name__) - +# Cache du client Gemini avec TTL (évite de recréer à chaque appel) _client: Optional[genai.Client] = None +_client_api_key: Optional[str] = None def _get_client() -> genai.Client: - global _client - if _client is None: - _client = genai.Client(api_key=settings.GEMINI_API_KEY) + """Retourne un client Gemini. Recrée si la clé API a changé.""" + global _client, _client_api_key + current_key = settings.GEMINI_API_KEY + if _client is None or _client_api_key != current_key: + _client = genai.Client(api_key=current_key) + _client_api_key = current_key return _client @@ -44,14 +54,7 @@ async def _read_image(file_path: str) -> tuple[bytes, str]: } media_type = mime_map.get(suffix, "image/jpeg") - # Utilisation du StorageBackend pour lire l'image backend = get_storage_backend() - - # On ruse un peu car StorageBackend n'a pas de 'read', - # mais on sait qu'en LocalStorage on peut lire en direct - # et en S3Storage on peut passer par les URLs ou aioboto3. - # Pour garder une abstraction propre, on va ajouter une méthode 'get_bytes' au backend. - data = await backend.get_bytes(file_path) return data, media_type @@ -76,13 +79,44 @@ def _usage_tokens_gemini(response) -> tuple[Optional[int], Optional[int]]: return prompt_tokens, output_tokens +async def _retry_with_backoff(fn, *args, max_retries=None, **kwargs): + """Exécute fn avec retry + backoff exponentiel en cas d'échec.""" + retries = max_retries if max_retries is not None else settings.AI_MAX_RETRIES + last_error = None + + for attempt in range(retries + 1): + try: + return await fn(*args, **kwargs) + except asyncio.TimeoutError: + last_error = "timeout" + wait = 2 ** attempt + logger.warning("ai.retry.timeout", extra={ + "attempt": attempt + 1, + "max_retries": retries, + "wait_s": wait, + }) + except Exception as e: + last_error = str(e) + wait = 2 ** attempt + logger.warning("ai.retry.error", extra={ + "attempt": attempt + 1, + "max_retries": retries, + "wait_s": wait, + "error": str(e), + }) + if attempt < retries: + await asyncio.sleep(wait) + + raise Exception(f"AI request failed after {retries + 1} attempts: {last_error}") + + async def _generate_gemini( prompt: str, image_bytes: Optional[bytes] = None, media_type: Optional[str] = None, max_tokens: int = 1024 ) -> dict: - """Appel à Google Gemini via SDK.""" + """Appel à Google Gemini via SDK avec timeout et retry.""" if not settings.GEMINI_API_KEY: logger.warning("ai.gemini.no_key") return {"text": None, "usage": (None, None)} @@ -93,17 +127,22 @@ async def _generate_gemini( contents.append(types.Part.from_bytes(data=image_bytes, mime_type=media_type)) contents.append(prompt) - try: - # Le SDK est sync, on le run dans un thread - response = await asyncio.to_thread( - client.models.generate_content, - model=settings.GEMINI_MODEL, - contents=contents, - config=types.GenerateContentConfig( - max_output_tokens=max_tokens, - response_mime_type="application/json", + async def _call(): + return await asyncio.wait_for( + asyncio.to_thread( + client.models.generate_content, + model=settings.GEMINI_MODEL, + contents=contents, + config=types.GenerateContentConfig( + max_output_tokens=max_tokens, + response_mime_type="application/json", + ), ), + timeout=settings.AI_REQUEST_TIMEOUT, ) + + try: + response = await _retry_with_backoff(_call) usage = _usage_tokens_gemini(response) return {"text": getattr(response, "text", ""), "usage": usage} except Exception as e: @@ -117,7 +156,7 @@ async def _generate_openrouter( media_type: Optional[str] = None, max_tokens: int = 1024 ) -> dict: - """Appel à OpenRouter via HTTP.""" + """Appel à OpenRouter via HTTP avec timeout et retry.""" if not settings.OPENROUTER_API_KEY: logger.warning("ai.openrouter.no_key") return {"text": None, "usage": (None, None)} @@ -131,7 +170,6 @@ async def _generate_openrouter( messages = [] content_payload = [] - content_payload.append({"type": "text", "text": prompt}) if image_bytes and media_type: @@ -152,30 +190,32 @@ async def _generate_openrouter( "response_format": {"type": "json_object"} } - async with httpx.AsyncClient() as client: - try: + async def _call(): + async with httpx.AsyncClient() as client: response = await client.post( "https://openrouter.ai/api/v1/chat/completions", json=payload, headers=headers, - timeout=60.0 + timeout=settings.AI_REQUEST_TIMEOUT, ) response.raise_for_status() - data = response.json() - - text = "" - if "choices" in data and len(data["choices"]) > 0: - text = data["choices"][0]["message"]["content"] - - usage_data = data.get("usage", {}) - prompt_tokens = usage_data.get("prompt_tokens") - output_tokens = usage_data.get("completion_tokens") + return response.json() - return {"text": text, "usage": (prompt_tokens, output_tokens)} + try: + data = await _retry_with_backoff(_call) + text = "" + if "choices" in data and len(data["choices"]) > 0: + text = data["choices"][0]["message"]["content"] - except Exception as e: - logger.error("ai.openrouter.error", extra={"error": str(e)}) - return {"text": None, "usage": (None, None), "error": str(e)} + usage_data = data.get("usage", {}) + prompt_tokens = usage_data.get("prompt_tokens") + output_tokens = usage_data.get("completion_tokens") + + return {"text": text, "usage": (prompt_tokens, output_tokens)} + + except Exception as e: + logger.error("ai.openrouter.error", extra={"error": str(e)}) + return {"text": None, "usage": (None, None), "error": str(e)} async def _generate( @@ -187,11 +227,10 @@ async def _generate( """Dispatcher vers le bon provider.""" provider = settings.AI_PROVIDER.lower() logger.info("ai.generate", extra={"provider": provider}) - + if provider == "openrouter": return await _generate_openrouter(prompt, image_bytes, media_type, max_tokens) else: - # Default to Gemini return await _generate_gemini(prompt, image_bytes, media_type, max_tokens) @@ -264,9 +303,9 @@ async def analyze_image( result["confidence"] = parsed.get("confidence") else: logger.warning("ai.vision.json_parse_failed", extra={"raw": text[:100]}) - + if response.get("error"): - logger.error("ai.vision.provider_error", extra={"error": response['error']}) + logger.error("ai.vision.provider_error", extra={"error": response["error"]}) except Exception as e: logger.error("ai.vision.unexpected_error", extra={"error": str(e)}) @@ -303,7 +342,7 @@ Retourne UNIQUEMENT un objet JSON : } Si aucun texte n'est visible, retourne : {"text": "", "has_text": false} """ - + response = await _generate( prompt=prompt, image_bytes=image_bytes, @@ -411,12 +450,11 @@ Retourne UNIQUEMENT ce JSON : }}""" try: - # Pas d'image ici response = await _generate( prompt=prompt, max_tokens=settings.GEMINI_MAX_TOKENS ) - + text = response.get("text") if text: parsed = _extract_json(text) diff --git a/app/services/storage.py b/app/services/storage.py index a29a8b1..02687f0 100644 --- a/app/services/storage.py +++ b/app/services/storage.py @@ -2,6 +2,7 @@ Service de stockage — sauvegarde fichiers, génération thumbnails Multi-tenant : les fichiers sont isolés par client_id. """ +import asyncio import uuid import logging import io @@ -33,6 +34,9 @@ async def save_upload(file: UploadFile, client_id: str) -> dict: """ Valide, sauvegarde le fichier uploadé et génère un thumbnail. Utilise le backend de stockage configuré (Local ou S3). + + L'ordre est : validation → upload fichier → génération thumbnail → retourne les métadonnées. + L'appelant est responsable de l'insertion BDD et du rollback en cas d'échec. """ # ── Validation MIME ─────────────────────────────────────── if file.content_type not in ALLOWED_MIME_TYPES: @@ -53,7 +57,7 @@ async def save_upload(file: UploadFile, client_id: str) -> dict: # ── Nommage ─────────────────────────────────────────────── filename, file_uuid = _generate_filename(file.filename or "image") - + # Chemins relatifs par rapport au bucket/base_dir rel_file_path = f"uploads/{client_id}/{filename}" rel_thumb_path = f"thumbnails/{client_id}/thumb_{filename}" @@ -67,26 +71,28 @@ async def save_upload(file: UploadFile, client_id: str) -> dict: width, height = None, None thumb_saved = False try: - # On utilise io.BytesIO pour ne pas avoir à écrire sur le disque local with PILImage.open(io.BytesIO(content)) as img: width, height = img.size img.thumbnail(THUMBNAIL_SIZE, PILImage.LANCZOS) - + # Convertit en RGB si nécessaire if img.mode in ("RGBA", "P"): img = img.convert("RGB") - + # Sauvegarde thumbnail dans un buffer thumb_buffer = io.BytesIO() img.save(thumb_buffer, "JPEG", quality=85) thumb_data = thumb_buffer.getvalue() - + # Sauvegarde via le backend await backend.save(thumb_data, rel_thumb_path, "image/jpeg") thumb_saved = True - + except Exception as e: - logger.warning("Erreur génération thumbnail : %s", e) + logger.warning("thumbnail.generation_error", extra={ + "file": rel_file_path, + "error": str(e), + }) return { "uuid": file_uuid, @@ -103,24 +109,27 @@ async def save_upload(file: UploadFile, client_id: str) -> dict: } -def delete_files(file_path: str, thumbnail_path: str | None = None) -> None: - """Supprime le fichier original et son thumbnail via le backend.""" - import asyncio +async def delete_files(file_path: str, thumbnail_path: str | None = None) -> None: + """ + Supprime le fichier original et son thumbnail via le backend de stockage. + Entièrement asynchrone — les fichiers sont supprimés de façon déterministe. + """ backend = get_storage_backend() - - async def _do_delete(): - await backend.delete(file_path) - if thumbnail_path: - await backend.delete(thumbnail_path) - - # Note: delete_files est synchrone dans les routers existants, - # mais le backend est async. C'est un risque. - # TODO: Refactorer delete_image pour être full async. + errors: list[str] = [] + try: - loop = asyncio.get_event_loop() - if loop.is_running(): - asyncio.ensure_future(_do_delete()) - else: - loop.run_until_complete(_do_delete()) - except Exception: - pass + await backend.delete(file_path) + except Exception as e: + errors.append(f"original: {e}") + + if thumbnail_path: + try: + await backend.delete(thumbnail_path) + except Exception as e: + errors.append(f"thumbnail: {e}") + + if errors: + logger.warning("storage.delete_partial", extra={ + "file_path": file_path, + "errors": errors, + }) diff --git a/app/workers/image_worker.py b/app/workers/image_worker.py index 09a5cce..98696b8 100644 --- a/app/workers/image_worker.py +++ b/app/workers/image_worker.py @@ -1,5 +1,8 @@ """ 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 @@ -12,6 +15,7 @@ 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 @@ -21,46 +25,83 @@ 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.""" + """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(f"--- JOB DÉMARRÉ : image_id={image_id} ---") + 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(f"--- JOB TERMINÉ : image_id={image_id} ---") - return f"OK" + logger.info("arq.job.completed", extra={"image_id": image_id}) + return "OK" + except Exception as e: - logger.error(f"--- JOB ÉCHOUÉ : {str(e)} ---", exc_info=True) - raise + 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 and listening on %s", WorkerSettings.queue_name) + logger.info("worker.started", extra={"queue": WorkerSettings.queue_name}) + + +async def on_shutdown(ctx: dict) -> None: + logger.info("worker.shutdown") -def _parse_redis_settings() -> RedisSettings: - url = settings.REDIS_URL - if url.startswith("redis://"): url = url[8:] - elif url.startswith("rediss://"): url = url[9:] - password, host, port, database = None, "localhost", 6379, 0 - if "@" in url: - auth, url = url.rsplit("@", 1) - password = auth.split(":", 1)[1] if ":" in auth else auth - if "/" in url: - url, db_str = url.split("/", 1) - if db_str: database = int(db_str) - if ":" in url: - host, port_str = url.rsplit(":", 1) - port = int(port_str) - else: host = url - return RedisSettings(host=host, port=port, password=password, database=database) class WorkerSettings: functions = [func(process_image_task, name="process_image_task")] - redis_settings = _parse_redis_settings() + 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 = 300 + job_timeout = settings.PIPELINE_TIMEOUT + health_check_interval = 30 diff --git a/app/workers/redis_client.py b/app/workers/redis_client.py index 476589c..e06ed8a 100644 --- a/app/workers/redis_client.py +++ b/app/workers/redis_client.py @@ -1,6 +1,9 @@ """ Client Redis partagé — pool de connexions async pour ARQ et Pub/Sub. +Inclut le parsing d'URL Redis centralisé pour éviter la duplication. """ +from typing import Tuple, Optional + from redis.asyncio import ConnectionPool, Redis from app.config import settings @@ -8,10 +11,63 @@ from app.config import settings _pool: ConnectionPool | None = None +def parse_redis_url(url: str | None = None) -> Tuple[str, int, Optional[str], int]: + """ + Parse une URL Redis en (host, port, password, database). + + Format accepté : redis://[user:password@]host[:port][/db] + Retourne ('localhost', 6379, None, 0) si l'URL est vide. + """ + if url is None: + url = settings.REDIS_URL + + if not url: + return "localhost", 6379, None, 0 + + if url.startswith("redis://"): + url = url[8:] + if url.startswith("rediss://"): + url = url[9:] + + password: Optional[str] = None + if "@" in url: + auth_part, url = url.rsplit("@", 1) + if ":" in auth_part: + password = auth_part.split(":", 1)[1] + else: + password = auth_part + + database = 0 + if "/" in url: + host_port, db_str = url.split("/", 1) + if db_str: + try: + database = int(db_str) + except ValueError: + pass + else: + host_port = url + + host = "localhost" + port = 6379 + if ":" in host_port: + host, port_str = host_port.rsplit(":", 1) + if port_str: + try: + port = int(port_str) + except ValueError: + port = 6379 + else: + host = host_port or "localhost" + + return host, port, password, database + + async def get_redis_pool() -> Redis: """Retourne un client Redis avec pool de connexions partagé.""" global _pool if _pool is None: + host, port, password, database = parse_redis_url() _pool = ConnectionPool.from_url( settings.REDIS_URL, max_connections=20, diff --git a/docker-compose.yml b/docker-compose.yml index 6b3a525..34eaae3 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,106 +1,84 @@ +version: '3.9' + services: - # ── Portail Admin (React + Nginx) ────────────────────────── - admin: - build: - context: ./imago-admin - args: - VITE_API_BASE_URL: "" - ports: - - "3000:3000" - depends_on: - backend: - condition: service_healthy - restart: unless-stopped - healthcheck: - test: ["CMD", "curl", "-f", "http://localhost:3000"] - interval: 30s - timeout: 5s - retries: 3 - - # ── Backend Imago Hub (FastAPI + Uvicorn) ────────────────── - backend: - build: . - ports: - - "8000:8000" - volumes: - - ./data:/app/data - env_file: - - .env - environment: - - DATABASE_URL=postgresql+asyncpg://imago:imago@db:5432/imago - - REDIS_URL=redis://redis:6379/0 - depends_on: - db: - condition: service_healthy - redis: - condition: service_healthy - restart: unless-stopped - healthcheck: - test: ["CMD", "curl", "-f", "http://localhost:8000/health"] - interval: 30s - timeout: 10s - retries: 3 - - db: + postgres: image: postgres:16-alpine environment: POSTGRES_USER: imago POSTGRES_PASSWORD: imago POSTGRES_DB: imago + ports: + - '5432:5432' volumes: - - postgres_data:/var/lib/postgresql/data + - pgdata:/var/lib/postgresql/data healthcheck: - test: ["CMD-SHELL", "pg_isready -U imago -d imago"] + test: ['CMD-SHELL', 'pg_isready -U imago'] interval: 5s timeout: 5s retries: 5 - restart: unless-stopped redis: image: redis:7-alpine ports: - - "6379:6379" + - '6379:6379' volumes: - - redis_data:/data - command: redis-server --appendonly yes - restart: unless-stopped + - redisdata:/data healthcheck: - test: ["CMD", "redis-cli", "ping"] - interval: 10s + test: ['CMD', 'redis-cli', 'ping'] + interval: 5s timeout: 5s - retries: 3 - - worker: - build: . - command: python worker.py - volumes: - - ./data:/app/data - env_file: - - .env - environment: - - DATABASE_URL=postgresql+asyncpg://imago:imago@db:5432/imago - - REDIS_URL=redis://redis:6379/0 - depends_on: - db: - condition: service_healthy - redis: - condition: service_healthy - restart: unless-stopped + retries: 5 minio: - image: minio/minio - ports: - - "9000:9000" - - "9001:9001" + image: minio/minio:latest + command: server /data --console-address ':9001' environment: MINIO_ROOT_USER: minioadmin MINIO_ROOT_PASSWORD: minioadmin - command: server /data --console-address ":9001" + ports: + - '9000:9000' + - '9001:9001' volumes: - - minio_data:/data - restart: unless-stopped + - miniodata:/data + healthcheck: + test: ['CMD', 'mc', 'ready', 'local'] + interval: 5s + timeout: 5s + retries: 5 + + api: + build: . + ports: + - '8000:8000' + environment: + - DATABASE_URL=postgresql+asyncpg://imago:imago@postgres:5432/imago + - REDIS_URL=redis://redis:6379/0 + - STORAGE_BACKEND=s3 + - S3_BUCKET=imago + - S3_ENDPOINT_URL=http://minio:9000 + - S3_ENDPOINT_PUBLIC=http://localhost:9000 + - S3_ACCESS_KEY=minioadmin + - S3_SECRET_KEY=minioadmin + - S3_REGION=us-east-1 + - SIGNED_URL_SECRET=dev-secret-change-in-prod + - SECRET_KEY=dev-secret-change-in-prod + - TESSERACT_CMD=/usr/bin/tesseract + - GEMINI_API_KEY=${GEMINI_API_KEY:-} + - OPENROUTER_API_KEY=${OPENROUTER_API_KEY:-} + depends_on: + postgres: + condition: service_healthy + redis: + condition: service_healthy + minio: + condition: service_healthy + volumes: + - uploads:/app/data/uploads + - thumbnails:/app/data/thumbnails volumes: - postgres_data: - redis_data: - minio_data: + pgdata: + redisdata: + miniodata: + uploads: + thumbnails: diff --git a/imago-admin/docker-compose.yml b/imago-admin/docker-compose.yml index dacb868..0959e6a 100644 --- a/imago-admin/docker-compose.yml +++ b/imago-admin/docker-compose.yml @@ -1,41 +1,40 @@ services: - # Portail admin (ce projet) admin: build: context: . args: - VITE_API_BASE_URL: http://backend:8000 + VITE_API_BASE_URL: http://localhost:8000 ports: - - "3000:3000" + - '3000:3000' depends_on: - backend restart: unless-stopped - # Backend Imago Hub (projet séparé) backend: build: context: .. ports: - - "8000:8000" + - '8000:8000' volumes: - ../data:/app/data env_file: ../.env environment: - DATABASE_URL=postgresql+asyncpg://imago:imago@db:5432/imago - REDIS_URL=redis://redis:6379/0 + - STORAGE_BACKEND=s3 + - S3_BUCKET=imago + - S3_ENDPOINT_URL=http://minio:9000 + - S3_ENDPOINT_PUBLIC=http://localhost:9000 + - S3_ACCESS_KEY=minioadmin + - S3_SECRET_KEY=*** + - S3_REGION=us-east-1 depends_on: db: condition: service_healthy redis: condition: service_healthy restart: unless-stopped - healthcheck: - test: ["CMD", "curl", "-f", "http://localhost:8000/health"] - interval: 30s - timeout: 10s - retries: 3 - # Worker ARQ worker: build: context: .. @@ -46,6 +45,13 @@ services: environment: - DATABASE_URL=postgresql+asyncpg://imago:imago@db:5432/imago - REDIS_URL=redis://redis:6379/0 + - STORAGE_BACKEND=s3 + - S3_BUCKET=imago + - S3_ENDPOINT_URL=http://minio:9000 + - S3_ENDPOINT_PUBLIC=http://localhost:9000 + - S3_ACCESS_KEY=minioadmin + - S3_SECRET_KEY=*** + - S3_REGION=us-east-1 depends_on: db: condition: service_healthy @@ -53,7 +59,6 @@ services: condition: service_healthy restart: unless-stopped - # PostgreSQL db: image: postgres:16-alpine environment: @@ -63,37 +68,35 @@ services: volumes: - postgres_data:/var/lib/postgresql/data healthcheck: - test: ["CMD-SHELL", "pg_isready -U imago -d imago"] + test: ['CMD-SHELL', 'pg_isready -U imago -d imago'] interval: 5s timeout: 5s retries: 5 restart: unless-stopped - # Redis (file ARQ + Pub/Sub) redis: image: redis:7-alpine ports: - - "6379:6379" + - '6379:6379' volumes: - redis_data:/data command: redis-server --appendonly yes restart: unless-stopped healthcheck: - test: ["CMD", "redis-cli", "ping"] + test: ['CMD', 'redis-cli', 'ping'] interval: 10s timeout: 5s retries: 3 - # MinIO (stockage S3 compatible — développement) minio: image: minio/minio ports: - - "9000:9000" - - "9001:9001" + - '9000:9000' + - '9001:9001' environment: MINIO_ROOT_USER: minioadmin MINIO_ROOT_PASSWORD: minioadmin - command: server /data --console-address ":9001" + command: server /data --console-address ':9001' volumes: - minio_data:/data restart: unless-stopped diff --git a/pyproject.toml b/pyproject.toml index 3847f6e..dc0769e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,6 +4,64 @@ version = "2.0.0" description = "Backend FastAPI pour la gestion d'images et fonctionnalités AI — Imago" requires-python = ">=3.12" license = "MIT" +license-files = ["LICENSE"] +dependencies = [ + # Web Framework + "fastapi==0.115.0", + "uvicorn[standard]==0.30.6", + "python-multipart==0.0.9", + # Database + "sqlalchemy==2.0.35", + "alembic==1.13.3", + "aiosqlite==0.20.0", + "asyncpg==0.29.0", + # Validation + "pydantic==2.9.2; python_version < '3.14'", + "pydantic>=2.12.0; python_version >= '3.14'", + "pydantic-settings==2.5.2", + # Image Processing + "Pillow==10.4.0; python_version < '3.14'", + "Pillow>=11.0.0; python_version >= '3.14'", + "piexif==1.1.3", + # OCR + "pytesseract==0.3.13", + # AI + "google-genai==1.0.0", + "httpx==0.27.2", + # Web scraping + "beautifulsoup4==4.12.3", + # Task scheduling + "apscheduler==3.10.4", + # Task queue + "arq==0.25.0", + "redis==5.0.8", + # Storage abstraction + "aioboto3==13.0.0", + "itsdangerous==2.2.0", + # Observability + "structlog==24.4.0", + "prometheus-fastapi-instrumentator==7.0.2", + # Utilities + "python-dotenv==1.0.1", + "python-jose[cryptography]==3.3.0", + "passlib[bcrypt]==1.7.4", + "aiofiles==24.1.0", + # Rate Limiting + "slowapi==0.1.9", + # WebSockets + "websockets>=13.0", +] + +[project.optional-dependencies] +dev = [ + "pytest==8.3.3", + "pytest-asyncio==0.24.0", + "pytest-cov==5.0.0", + "ruff==0.8.6", + "mypy==1.13.0", + "pre-commit==4.0.1", + "types-aiofiles==24.1.0.20240626", +] [tool.ruff] target-version = "py312" @@ -20,7 +78,7 @@ select = [ "UP", # pyupgrade "B", # bugbear "S", # bandit (security) - "T20", # flake8-print (catch stray print()) + "T20", # flake8-print "SIM", # flake8-simplify "RUF", # ruff-specific ] diff --git a/worker.py b/worker.py index 897c54b..4507fcf 100644 --- a/worker.py +++ b/worker.py @@ -1,19 +1,36 @@ +#!/usr/bin/env python3 """ -Entrypoint du worker ARQ — traitement des images en arrière-plan. +Worker ARQ — lancement du worker de traitement asynchrone. -Lancer avec : python worker.py +Usage: + python worker.py + # ou via arq directement : + # arq app.workers.image_worker.WorkerSettings -Le worker écoute les queues Redis 'standard' et 'premium' et traite -les tâches de pipeline image (EXIF → OCR → AI). +Le worker écoute la file Redis et traite les images uploadées +via le pipeline EXIF → OCR → Vision AI. """ -import asyncio +import logging +import sys +from pathlib import Path + +# Ajouter le répertoire parent au PYTHONPATH pour les imports relatifs +sys.path.insert(0, str(Path(__file__).resolve().parent)) + from arq import run_worker -from app.config import settings + from app.logging_config import configure_logging +from app.config import settings from app.workers.image_worker import WorkerSettings -# Configure le logging dès l'import -configure_logging(debug=settings.DEBUG) -if __name__ == "__main__" : +def main(): + configure_logging(debug=settings.DEBUG) + logger = logging.getLogger("worker") + logger.info("Starting Imago worker on queue '%s'", WorkerSettings.queue_name) + run_worker(WorkerSettings) + + +if __name__ == "__main__": + main()