Compare commits

...
2 Commits
Author SHA1 Message Date
xavier 8a32c16ff4 feat: implémentation de l'orchestrateur asynchrone multi-agents (prompteur/travailleurs/concaténeur) et API de synchronisation des modèles
Build and Push Docker Image / build-and-push (push) Successful in 42s
2026-10-07 15:39:24 +02:00
xavier 4ce11d0bfe fix: mise à jour de l'identifiant du modèle Gemini vers gemini-3.5-flash-lite
Build and Push Docker Image / build-and-push (push) Successful in 40s
2026-10-07 15:17:03 +02:00
3 changed files with 154 additions and 8 deletions

No files matched your search

+4 -4
View File
@@ -10,20 +10,20 @@ def get_ai_response(messages: list, settings: SystemSettings) -> str:
formatted_messages = [{"role": msg.role, "content": msg.content} for msg in messages]
try:
# 1. Priorité absolue : Test Gemini direct (Rapide, sans intermédiaire)
# 1. Priorité absolue : Test Gemini direct
if settings.gemini_api_key:
client = OpenAI(
base_url="https://generativelanguage.googleapis.com/v1beta/openai/",
api_key=settings.gemini_api_key
)
# Utilisation de l'identifiant standard API pour la version Flash
model = "gemini-1.5-flash"
# Correction : Utilisation exacte du modèle 3.5 Flash-Lite
model = "gemini-3.5-flash-lite"
response = client.chat.completions.create(model=model, messages=formatted_messages)
# 2. Test OpenRouter (Fallback)
elif settings.openrouter_api_key:
client = OpenAI(base_url="https://openrouter.ai/api/v1", api_key=settings.openrouter_api_key)
model = "google/gemini-1.5-pro" # Identifiant mis à jour et valide
model = "google/gemini-3.5-flash-lite"
response = client.chat.completions.create(
model=model,
messages=formatted_messages,
+23 -4
View File
@@ -10,7 +10,8 @@ from .database import engine, Base, get_db
from .auth import get_password_hash, generate_totp_secret, get_totp_uri, verify_password, verify_totp, create_access_token, verify_token
from .schemas import AdminCreate, LoginRequest, ProjectCreate, ProjectResponse, ProjectRename, MessageCreate, MessageResponse
from .models import User, Project, Message, SystemSettings
from .ai import get_ai_response
from .orchestrator import run_orchestrator, sync_providers_models
from .models import User, Project, Message, SystemSettings, AIModel
# Création des tables dans la base de données
Base.metadata.create_all(bind=engine)
@@ -151,7 +152,8 @@ def get_messages(project_id: int, db: Session = Depends(get_db), current_user: U
return db.query(Message).filter(Message.project_id == project_id).order_by(Message.created_at.asc()).all()
@app.post("/api/projects/{project_id}/messages", response_model=List[MessageResponse])
def create_message(project_id: int, message: MessageCreate, db: Session = Depends(get_db), current_user: User = Depends(get_current_user)):
async def create_message(project_id: int, message: MessageCreate, db: Session = Depends(get_db), current_user: User = Depends(get_current_user)):
"""Ajoute un message, lance l'Orchestrateur asynchrone et retourne l'historique."""
project = db.query(Project).filter(Project.id == project_id, Project.user_id == current_user.id).first()
if not project: raise HTTPException(status_code=404, detail="Projet introuvable")
@@ -162,10 +164,27 @@ def create_message(project_id: int, message: MessageCreate, db: Session = Depend
history = db.query(Message).filter(Message.project_id == project_id).order_by(Message.created_at.asc()).all()
settings = db.query(SystemSettings).first()
ai_response_text = get_ai_response(history, settings)
# Configuration temporaire (sera remplacée par les choix du Frontend à la Phase 2)
# Ex pour tester le pipeline complet : {"workers": ["gemini-3.5-flash-lite", "google/gemini-1.5-pro"]}
orchestrator_config = {"workers": ["gemini-3.5-flash-lite"]}
ai_response_text = await run_orchestrator(history, settings, orchestrator_config)
ai_message = Message(role="assistant", content=ai_response_text, project_id=project_id)
db.add(ai_message)
db.commit()
return db.query(Message).filter(Message.project_id == project_id).order_by(Message.created_at.asc()).all()
return db.query(Message).filter(Message.project_id == project_id).order_by(Message.created_at.asc()).all()
@app.post("/api/models/sync")
async def trigger_model_sync(db: Session = Depends(get_db), current_user: User = Depends(get_current_user)):
"""Déclenche manuellement l'extraction et la mise à jour des modèles depuis les API."""
if not current_user.is_admin:
raise HTTPException(status_code=403, detail="Accès réservé aux administrateurs.")
settings = db.query(SystemSettings).first()
return await sync_providers_models(db, settings)
@app.get("/api/models")
def get_models(db: Session = Depends(get_db), current_user: User = Depends(get_current_user)):
"""Renvoie la liste complète des modèles stockés en base de données."""
return db.query(AIModel).order_by(AIModel.name.asc()).all()
+127
View File
@@ -0,0 +1,127 @@
import asyncio
import httpx
from openai import AsyncOpenAI
from fastapi import HTTPException
from sqlalchemy.orm import Session
from .models import SystemSettings, AIModel
from .logger import system_logger
from datetime import datetime, timezone
# --- PARTIE 1 : EXTRACTION DES MODÈLES ---
async def sync_providers_models(db: Session, settings: SystemSettings):
"""Extrait et met à jour les modèles depuis les fournisseurs configurés."""
added_or_updated = 0
# httpx.AsyncClient permet des requêtes non-bloquantes (ultra rapide)
async with httpx.AsyncClient() as client:
# 1. OpenRouter (Exemple principal pour l'extraction massive)
if settings.openrouter_api_key:
try:
response = await client.get("https://openrouter.ai/api/v1/models")
if response.status_code == 200:
for item in response.json().get("data", []):
model_id = item["id"]
existing = db.query(AIModel).filter(AIModel.model_id == model_id).first()
pricing = item.get("pricing", {})
# Conversion en coût pour 1 Million de tokens
p_prompt = float(pricing.get("prompt", 0)) * 1000000 if pricing.get("prompt") else 0.0
p_comp = float(pricing.get("completion", 0)) * 1000000 if pricing.get("completion") else 0.0
if existing:
existing.pricing_prompt = p_prompt
existing.pricing_completion = p_comp
existing.last_updated = datetime.now(timezone.utc)
else:
new_model = AIModel(
provider="openrouter",
model_id=model_id,
name=item["name"],
context_length=item.get("context_length", 0),
pricing_prompt=p_prompt,
pricing_completion=p_comp
)
db.add(new_model)
added_or_updated += 1
except Exception as e:
system_logger.error(f"Erreur Sync OpenRouter: {e}")
db.commit()
return {"status": "success", "models_processed": added_or_updated}
# --- PARTIE 2 : MOTEUR MULTI-AGENTS ---
def get_client_for_model(model_id: str, settings: SystemSettings):
"""Retourne le client AsyncOpenAI approprié selon le modèle sélectionné."""
if "gemini" in model_id.lower() and settings.gemini_api_key:
return AsyncOpenAI(base_url="https://generativelanguage.googleapis.com/v1beta/openai/", api_key=settings.gemini_api_key)
elif settings.openrouter_api_key:
return AsyncOpenAI(base_url="https://openrouter.ai/api/v1", api_key=settings.openrouter_api_key)
raise ValueError(f"Aucun fournisseur configuré pour {model_id}")
async def ask_agent(client, model_id, messages, is_openrouter=False):
"""Appel asynchrone à un modèle IA."""
kwargs = {"model": model_id, "messages": messages}
if is_openrouter:
kwargs["extra_headers"] = {"HTTP-Referer": "https://aethas38.duckdns.org", "X-Title": "AETHAS38 Orchestrator"}
response = await client.chat.completions.create(**kwargs)
return response.choices[0].message.content
async def run_orchestrator(history: list, settings: SystemSettings, config: dict) -> str:
"""
Gère la logique : 1 Prompteur -> N Travailleurs -> 1 Concaténeur.
Le dictionnaire 'config' proviendra de l'interface graphique.
"""
workers = config.get("workers", [])
if not workers:
workers = ["gemini-3.5-flash-lite"] # Fallback de sécurité
user_prompt = history[-1].content
formatted_history = [{"role": msg.role, "content": msg.content} for msg in history[:-1]]
# SCÉNARIO 1 : Un seul travailleur (Pas besoin de prompteur/concaténeur)
if len(workers) == 1:
worker_model = workers[0]
client = get_client_for_model(worker_model, settings)
messages = formatted_history + [{"role": "user", "content": user_prompt}]
return await ask_agent(client, worker_model, messages, "openrouter" in worker_model.lower())
# SCÉNARIO 2 : Multi-Travailleurs (Le pipeline complet)
try:
# Étape 1 : Le Prompteur améliore la requête
prompter_model = config.get("prompter", "gemini-3.5-flash-lite")
p_client = get_client_for_model(prompter_model, settings)
p_messages = [{"role": "system", "content": "Tu es un expert en Prompt Engineering. Optimise la requête de l'utilisateur pour qu'elle soit claire, directive et parfaite pour des IAs de génération. Retourne UNIQUEMENT le prompt optimisé."}]
p_messages.append({"role": "user", "content": user_prompt})
system_logger.info("Démarrage du Prompteur...")
optimized_prompt = await ask_agent(p_client, prompter_model, p_messages)
# Étape 2 : Les Travailleurs en parallèle (Magie de l'Asynchrone)
system_logger.info(f"Lancement de {len(workers)} travailleurs en parallèle...")
w_tasks = []
for w_model in workers:
w_client = get_client_for_model(w_model, settings)
w_messages = formatted_history + [{"role": "user", "content": optimized_prompt}]
# On stocke les tâches sans les attendre immédiatement
w_tasks.append(ask_agent(w_client, w_model, w_messages, "openrouter" in w_model.lower()))
# 'gather' exécute toutes les requêtes en même temps !
workers_responses = await asyncio.gather(*w_tasks, return_exceptions=True)
# Étape 3 : Le Concaténeur synthétise
concat_model = config.get("concatenator", "gemini-3.5-flash-lite")
c_client = get_client_for_model(concat_model, settings)
synthesis_prompt = f"Voici la requête initiale : {user_prompt}\n\nVoici les réponses de {len(workers)} experts IA différents :\n"
for i, resp in enumerate(workers_responses):
synthesis_prompt += f"--- EXPERT {i+1} ---\n{resp if not isinstance(resp, Exception) else 'Erreur de génération'}\n\n"
synthesis_prompt += "Fais une synthèse finale parfaite, complète et structurée de ces réponses, en gardant le meilleur de chacune."
system_logger.info("Démarrage du Concaténeur...")
c_messages = [{"role": "user", "content": synthesis_prompt}]
return await ask_agent(c_client, concat_model, c_messages)
except Exception as e:
system_logger.error(f"Erreur Pipeline Multi-Agents: {e}")
raise HTTPException(status_code=502, detail=f"Échec de l'orchestration : {str(e)}")