diff --git a/backend/main.py b/backend/main.py index 29dbe1d..673fbc2 100644 --- a/backend/main.py +++ b/backend/main.py @@ -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() \ No newline at end of file + 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() \ No newline at end of file diff --git a/backend/orchestrator.py b/backend/orchestrator.py new file mode 100644 index 0000000..01c1cee --- /dev/null +++ b/backend/orchestrator.py @@ -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)}") \ No newline at end of file