Introduce an admin portal (React + Nginx), WebSocket routing, and API versioning middleware with `/api/v1/` prefix deprecation. Add master API key authentication, new Prometheus metrics for AI token consumption and active WebSockets, and extend S3 config with a public endpoint URL. Update test paths and fixtures to align with the new routing structure.
252 lines
9.4 KiB
Python
252 lines
9.4 KiB
Python
"""
|
|
Router WebSocket — suivi temps réel du pipeline de traitement d'images.
|
|
|
|
Endpoints :
|
|
- WS /ws/pipeline/{image_id}?token=<api_key> → événements d'un pipeline
|
|
- WS /ws/admin/monitor?token=<admin_api_key> → monitoring admin global
|
|
"""
|
|
import json
|
|
import logging
|
|
from typing import Any
|
|
|
|
from fastapi import APIRouter, WebSocket, WebSocketDisconnect, status
|
|
from sqlalchemy import select
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from app.database import AsyncSessionLocal, get_db
|
|
from app.dependencies.auth import hash_api_key
|
|
from app.metrics import hub_active_websockets
|
|
from app.models.client import APIClient
|
|
from app.models.image import Image, ProcessingStatus
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
router = APIRouter(tags=["WebSocket"])
|
|
|
|
|
|
# ─────────────────────────────────────────────────────────────
|
|
# Helpers
|
|
# ─────────────────────────────────────────────────────────────
|
|
|
|
async def _authenticate_ws(websocket: WebSocket, db: AsyncSession) -> APIClient | None:
|
|
"""
|
|
Authentifie une connexion WebSocket via le query param `token`.
|
|
Retourne le client ou None si invalide.
|
|
"""
|
|
token = websocket.query_params.get("token")
|
|
if not token:
|
|
return None
|
|
|
|
# Vérification Master Key
|
|
from app.config import settings
|
|
if settings.ADMIN_API_KEY and token == settings.ADMIN_API_KEY:
|
|
return APIClient(
|
|
id="admin-master",
|
|
name="Imago Master Admin",
|
|
scopes=["admin", "images:read", "images:write", "ai:use"],
|
|
plan="premium",
|
|
)
|
|
|
|
key_hash = hash_api_key(token)
|
|
result = await db.execute(
|
|
select(APIClient).where(APIClient.api_key_hash == key_hash)
|
|
)
|
|
client = result.scalar_one_or_none()
|
|
if client and client.is_active:
|
|
return client
|
|
return None
|
|
|
|
|
|
async def _get_image(image_id: int, db: AsyncSession) -> Image | None:
|
|
"""Charge une image depuis la BDD."""
|
|
result = await db.execute(select(Image).where(Image.id == image_id))
|
|
return result.scalar_one_or_none()
|
|
|
|
|
|
# ─────────────────────────────────────────────────────────────
|
|
# WS /ws/pipeline/{image_id} — suivi d'un pipeline
|
|
# ─────────────────────────────────────────────────────────────
|
|
|
|
from fastapi import Depends
|
|
|
|
@router.websocket("/ws/pipeline/{image_id}")
|
|
async def ws_pipeline(
|
|
websocket: WebSocket,
|
|
image_id: int,
|
|
db: AsyncSession = Depends(get_db)
|
|
) -> None:
|
|
"""
|
|
WebSocket de suivi temps réel d'un pipeline image.
|
|
|
|
- Authentification via query param `token`
|
|
- Vérifie que l'image appartient au client
|
|
- Envoie le buffer de reconnexion puis les événements live
|
|
- Ferme après pipeline.done ou pipeline.error
|
|
"""
|
|
# ── Authentification ──────────────────────────────────────
|
|
client = await _authenticate_ws(websocket, db)
|
|
if client is None:
|
|
await websocket.close(code=4001, reason="Token manquant ou invalide")
|
|
return
|
|
|
|
# ── Vérification propriété de l'image ─────────────────────
|
|
image = await _get_image(image_id, db)
|
|
|
|
if image is None:
|
|
await websocket.accept()
|
|
await websocket.close(code=4004, reason="Image introuvable")
|
|
return
|
|
|
|
# Admin peut voir toutes les images, sinon vérifier ownership
|
|
if not client.has_scope("admin") and image.client_id != client.id:
|
|
await websocket.close(code=4003, reason="Accès interdit")
|
|
return
|
|
|
|
# ── Accepter la connexion ─────────────────────────────────
|
|
await websocket.accept()
|
|
hub_active_websockets.inc()
|
|
|
|
try:
|
|
# ── Image déjà terminée → message synthétique ─────────
|
|
if image.processing_status == ProcessingStatus.DONE:
|
|
await websocket.send_json({
|
|
"event": "pipeline.done",
|
|
"image_id": image_id,
|
|
"status": "done",
|
|
"synthetic": True,
|
|
})
|
|
return
|
|
|
|
if image.processing_status == ProcessingStatus.ERROR:
|
|
await websocket.send_json({
|
|
"event": "pipeline.error",
|
|
"image_id": image_id,
|
|
"error": image.processing_error or "Erreur inconnue",
|
|
"synthetic": True,
|
|
})
|
|
return
|
|
|
|
# ── Récupérer le buffer de reconnexion depuis Redis ───
|
|
redis = getattr(websocket.app.state, "redis", None)
|
|
if redis is not None:
|
|
try:
|
|
buffer_key = f"pipeline:buffer:{image_id}"
|
|
buffered = await redis.lrange(buffer_key, 0, -1)
|
|
for raw_event in buffered:
|
|
try:
|
|
event_data = json.loads(raw_event)
|
|
await websocket.send_json(event_data)
|
|
except (json.JSONDecodeError, Exception):
|
|
pass
|
|
except Exception as e:
|
|
logger.warning("ws.buffer_read_error", extra={"error": str(e)})
|
|
|
|
# ── S'abonner au channel Redis et écouter les événements
|
|
if redis is not None:
|
|
pubsub = redis.pubsub()
|
|
try:
|
|
await pubsub.subscribe(f"pipeline:{image_id}")
|
|
|
|
async for message in pubsub.listen():
|
|
if message["type"] != "message":
|
|
continue
|
|
|
|
try:
|
|
data = json.loads(message["data"])
|
|
except (json.JSONDecodeError, TypeError):
|
|
continue
|
|
|
|
await websocket.send_json(data)
|
|
|
|
# Fermer après pipeline.done ou pipeline.error
|
|
event_type = data.get("event", "")
|
|
if event_type in ("pipeline.done", "pipeline.error"):
|
|
break
|
|
finally:
|
|
await pubsub.unsubscribe(f"pipeline:{image_id}")
|
|
await pubsub.close()
|
|
else:
|
|
# Pas de Redis — envoyer un message d'info et fermer
|
|
await websocket.send_json({
|
|
"event": "error",
|
|
"message": "Redis indisponible — utilisez le polling GET /images/{id}/status",
|
|
})
|
|
|
|
except WebSocketDisconnect:
|
|
logger.info("ws.client_disconnected", extra={
|
|
"image_id": image_id,
|
|
"client_id": client.id,
|
|
})
|
|
except Exception as e:
|
|
logger.error("ws.unexpected_error", extra={
|
|
"image_id": image_id,
|
|
"error": str(e),
|
|
})
|
|
finally:
|
|
hub_active_websockets.dec()
|
|
|
|
|
|
# ─────────────────────────────────────────────────────────────
|
|
# WS /ws/admin/monitor — monitoring admin global
|
|
# ─────────────────────────────────────────────────────────────
|
|
|
|
@router.websocket("/ws/admin/monitor")
|
|
async def ws_admin_monitor(
|
|
websocket: WebSocket,
|
|
db: AsyncSession = Depends(get_db)
|
|
) -> None:
|
|
"""
|
|
WebSocket admin pour surveiller tous les pipelines en temps réel.
|
|
|
|
Nécessite le scope `admin`. Pousse un événement à chaque démarrage
|
|
ou fin de pipeline sur n'importe quelle image.
|
|
"""
|
|
# ── Authentification ──────────────────────────────────────
|
|
client = await _authenticate_ws(websocket, db)
|
|
if client is None:
|
|
await websocket.close(code=4001, reason="Token manquant ou invalide")
|
|
return
|
|
|
|
if not client.has_scope("admin"):
|
|
await websocket.close(code=4003, reason="Scope admin requis")
|
|
return
|
|
|
|
# ── Accepter la connexion ─────────────────────────────────
|
|
await websocket.accept()
|
|
hub_active_websockets.inc()
|
|
|
|
try:
|
|
redis = getattr(websocket.app.state, "redis", None)
|
|
if redis is None:
|
|
await websocket.send_json({
|
|
"event": "error",
|
|
"message": "Redis indisponible",
|
|
})
|
|
return
|
|
|
|
pubsub = redis.pubsub()
|
|
try:
|
|
await pubsub.subscribe("pipeline:admin")
|
|
|
|
async for message in pubsub.listen():
|
|
if message["type"] != "message":
|
|
continue
|
|
|
|
try:
|
|
data = json.loads(message["data"])
|
|
except (json.JSONDecodeError, TypeError):
|
|
continue
|
|
|
|
await websocket.send_json(data)
|
|
|
|
finally:
|
|
await pubsub.unsubscribe("pipeline:admin")
|
|
await pubsub.close()
|
|
|
|
except WebSocketDisconnect:
|
|
logger.info("ws.admin_disconnected", extra={"client_id": client.id})
|
|
except Exception as e:
|
|
logger.error("ws.admin_error", extra={"error": str(e)})
|
|
finally:
|
|
hub_active_websockets.dec()
|