"""FlowDeck — Agent router (v4.10.0). /api/agent/* Agents, conversations, SSE run, audit & rollback, skills, triggers, tools. """ from __future__ import annotations import asyncio import json import logging from datetime import UTC from fastapi import APIRouter, HTTPException, Request from fastapi.responses import StreamingResponse from app.auth.session import get_current_user from app.config import settings from app.db import get_conn from app.services import skill_gallery from app.services.agent_engine import AgentEngine, undo_action from app.services.llm_client import PROVIDER_MODELS, PROVIDERS, LLMClient from app.services.llm_config import ( delete_user_llm_key, fetch_provider_models, get_llm_config, get_user_llm_key, list_user_llm_keys, mark_llm_config_verified, mark_user_llm_key_verified, provider_info, set_llm_config, upsert_user_llm_key, ) from app.services.tool_registry import ToolRegistry logger = logging.getLogger(__name__) router = APIRouter(tags=["agent"], prefix="/api/agent") async def agent_scheduler(interval_seconds: int = 60): """Background loop that fires scheduled custom agents (trigger_type=schedule). Runs inside the FastAPI lifespan. A trigger's config_json may use a cron-like string ("* * * * *") — for simplicity we fire at most once per poll interval when the trigger is active and not currently running. """ from datetime import datetime, timedelta logger.info("FlowDeck Agent scheduler started") while True: try: await asyncio.sleep(interval_seconds) with get_conn() as conn: triggers = conn.execute( "SELECT * FROM agent_triggers WHERE trigger_type='schedule' AND is_active=1" ).fetchall() now = datetime.now(UTC).replace(tzinfo=None) for trig in triggers: last = trig["last_fired_at"] if last: try: last_dt = datetime.fromisoformat(last) except ValueError: last_dt = None if last_dt and now - last_dt < timedelta(seconds=interval_seconds): continue # fire: run the agent's system_instructions through the engine agent_row = conn.execute("SELECT * FROM agents WHERE id=?", (trig["agent_id"],)).fetchone() if not agent_row: continue user_id = agent_row["created_by"] cur = conn.execute( "INSERT INTO agent_conversations (agent_id, user_id, title, context_json) VALUES (?,?,?,?)", (agent_row["id"], user_id, f"Scheduled: {agent_row['name']}", json.dumps({"workspace_id": agent_row["workspace_id"], "trigger": "schedule"})), ) conv_id = cur.lastrowid conn.execute( "UPDATE agent_triggers SET last_fired_at=? WHERE id=?", (now.isoformat(), trig["id"]), ) conn.commit() objective = (agent_row["system_instructions"] or "").strip() or \ f"Exécute l'agent planifié « {agent_row['name']} »." # best-effort fire-and-forget (offline/mock path is safe & fast) try: engine = AgentEngine(user_id, workspace_id=agent_row["workspace_id"]) async for _ev in engine.run(conv_id, objective, model=agent_row["model"]): pass except Exception as exc: # noqa: BLE001 logger.warning("Scheduled agent %s failed: %s", agent_row["id"], exc) except asyncio.CancelledError: raise except Exception: # noqa: BLE001 logger.exception("Agent scheduler tick failed") async def _current_user_id(request: Request) -> int: """A14 : plus de fallback sur la row `admin` — 401 sans session.""" user = await get_current_user(request) if not user or not user.get("id"): raise HTTPException(status_code=401, detail="Authentication required") return user["id"] async def _workspace_id(request: Request) -> int | None: user = await get_current_user(request) if user and user.get("workspace_id"): return user["workspace_id"] try: ws = int(request.query_params.get("workspace_id", 0) or 0) return ws or None except (TypeError, ValueError): return None async def _current_admin(request: Request) -> dict: """A14 : session obligatoire, puis admin. L'ancien fallback « row admin » laissait un anonymous diriger `PATCH /api/agent/providers` (et donc le `ping()` vers un `api_base` de son choix = SSRF).""" user = await get_current_user(request) if not user: raise HTTPException(status_code=401, detail="Authentication required") if not user.get("is_admin"): from app.db import get_conn as _gc with _gc() as conn: row = conn.execute("SELECT is_admin FROM users WHERE id=?", (user.get("id"),)).fetchone() if not row or not row["is_admin"]: raise HTTPException(status_code=403, detail="Accès administrateur requis") return user def _default_agent(conn, user_id: int) -> dict: row = conn.execute( "SELECT * FROM agents WHERE agent_type='personal' ORDER BY id LIMIT 1" ).fetchone() if row: return dict(row) cur = conn.execute( "INSERT INTO agents (workspace_id, name, agent_type, model, created_by) VALUES (?, 'FlowDeck Agent', 'personal', 'gpt-4o', ?)", (None, user_id), ) conn.commit() return dict(conn.execute("SELECT * FROM agents WHERE id=?", (cur.lastrowid,)).fetchone()) # ── Agents ── @router.get("") async def list_agents(request: Request): user_id = await _current_user_id(request) ws = await _workspace_id(request) with get_conn() as conn: _default_agent(conn, user_id) rows = conn.execute("SELECT * FROM agents WHERE workspace_id IS ? OR workspace_id=? ORDER BY agent_type, name", (ws, ws)).fetchall() return {"agents": [dict(r) for r in rows]} @router.post("") async def create_agent(request: Request): user_id = await _current_user_id(request) ws = await _workspace_id(request) body = await request.json() if request.headers.get("content-type") else {} name = (body.get("name") or "").strip() or "Custom Agent" with get_conn() as conn: try: cur = conn.execute( """INSERT INTO agents (workspace_id, name, icon, agent_type, description, system_instructions, model, scope_json, trigger_json, approval_mode, created_by) VALUES (?,?,?,?,?,?,?,?,?,?,?)""", (ws, name, body.get("icon", "🤖"), body.get("agent_type", "custom"), body.get("description", ""), body.get("system_instructions", ""), body.get("model", "gpt-4o"), json.dumps(body.get("scope", {})), json.dumps(body.get("trigger", {})), body.get("approval_mode", "auto"), user_id), ) conn.commit() except Exception as exc: # noqa: BLE001 raise HTTPException(status_code=409, detail=f"Impossible de créer l'agent: {exc}") from exc return {"id": cur.lastrowid, "name": name, "status": "created"} # ── Conversations ── @router.get("/conversations") async def list_conversations(request: Request): user_id = await _current_user_id(request) with get_conn() as conn: rows = conn.execute( "SELECT * FROM agent_conversations WHERE user_id=? ORDER BY updated_at DESC", (user_id,), ).fetchall() return {"conversations": [dict(r) for r in rows]} @router.post("/conversations") async def create_conversation(request: Request): user_id = await _current_user_id(request) body = await request.json() if request.headers.get("content-type") else {} ws = await _workspace_id(request) with get_conn() as conn: agent = _default_agent(conn, user_id) cur = conn.execute( """INSERT INTO agent_conversations (agent_id, user_id, title, context_json, provider, model) VALUES (?,?,?,?,?,?)""", (agent["id"], user_id, body.get("title", "New conversation"), json.dumps({"workspace_id": ws}), body.get("provider", ""), body.get("model", "")), ) conn.commit() return {"id": cur.lastrowid, "title": body.get("title", "New conversation"), "status": "created"} @router.get("/conversations/{conversation_id}") def get_conversation(request: Request, conversation_id: int): with get_conn() as conn: conv = conn.execute("SELECT * FROM agent_conversations WHERE id=?", (conversation_id,)).fetchone() if not conv: raise HTTPException(status_code=404, detail="Conversation introuvable") messages = conn.execute( "SELECT * FROM agent_messages WHERE conversation_id=? ORDER BY created_at, id", (conversation_id,), ).fetchall() return {"conversation": dict(conv), "messages": [dict(m) for m in messages]} @router.delete("/conversations/{conversation_id}") def delete_conversation(request: Request, conversation_id: int): with get_conn() as conn: if not conn.execute("SELECT id FROM agent_conversations WHERE id=?", (conversation_id,)).fetchone(): raise HTTPException(status_code=404, detail="Conversation introuvable") conn.execute("DELETE FROM agent_conversations WHERE id=?", (conversation_id,)) conn.commit() return {"id": conversation_id, "status": "deleted"} @router.patch("/conversations/{conversation_id}") async def patch_conversation(request: Request, conversation_id: int): """Update a conversation's title / provider / model (slash-command support).""" user_id = await _current_user_id(request) body = await request.json() if request.headers.get("content-type") else {} with get_conn() as conn: conv = conn.execute( "SELECT id FROM agent_conversations WHERE id=? AND user_id=?", (conversation_id, user_id), ).fetchone() if not conv: raise HTTPException(status_code=404, detail="Conversation introuvable") sets, params = ["updated_at=CURRENT_TIMESTAMP"], [] for col in ("title", "provider", "model"): if body.get(col) is not None: sets.append(f"{col}=?") params.append(str(body[col])) params.append(conversation_id) conn.execute(f"UPDATE agent_conversations SET {', '.join(sets)} WHERE id=?", params) conn.commit() return {"id": conversation_id, "status": "updated"} # ── Run (SSE) ── @router.post("/conversations/{conversation_id}/run") async def run_conversation(request: Request, conversation_id: int): user_id = await _current_user_id(request) ws = await _workspace_id(request) body = await request.json() if request.headers.get("content-type") else {} objective = (body.get("message") or "").strip() if not objective: raise HTTPException(status_code=400, detail="message est requis") with get_conn() as conn: conv = conn.execute("SELECT * FROM agent_conversations WHERE id=?", (conversation_id,)).fetchone() if not conv: raise HTTPException(status_code=404, detail="Conversation introuvable") eff_provider = body.get("provider") or conv["provider"] or None eff_model = body.get("model") or conv["model"] or None # Persist the selection so the same provider/model is reused next time. if body.get("provider") is not None or body.get("model") is not None: conn.execute( "UPDATE agent_conversations SET provider=?, model=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (body.get("provider", conv["provider"] or ""), body.get("model", conv["model"] or ""), conversation_id), ) conn.commit() engine = AgentEngine(user_id, workspace_id=ws) if eff_provider: user_key = get_user_llm_key(user_id, eff_provider) if user_key and user_key.get("api_key"): engine.llm = LLMClient(provider=eff_provider, api_key=user_key["api_key"], api_base=user_key.get("api_base") or None) else: engine.llm = LLMClient(provider=eff_provider) async def event_stream(): async for ev in engine.run( conversation_id, objective, model=eff_model, mentions=body.get("mentions"), files=body.get("files"), skill_id=body.get("skill_id"), skill_ids=body.get("skill_ids"), extra_context=body.get("context"), ): yield f"data: {json.dumps(ev, ensure_ascii=False)}\n\n" return StreamingResponse( event_stream(), media_type="text/event-stream", headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"}, ) @router.post("/generate") async def agent_generate(request: Request): """Headless text generation (NO tools) for content actions on documents. Used by the Agent panel's contextual actions ("Résumer ce document", "Traduire cette page", "Proposer des améliorations", "AI meeting note"). Because no tool schema is offered, the model answers with plain text based on the provided document context instead of issuing search_workspace / tools. """ user_id = await _current_user_id(request) body = await request.json() if request.headers.get("content-type") else {} prompt = (body.get("prompt") or "").strip() if not prompt: raise HTTPException(status_code=400, detail="prompt est requis") provider = (body.get("provider") or "").strip().lower() or None model = (body.get("model") or "").strip() or None context = (body.get("context") or "") or "" llm = LLMClient() if provider: if provider not in PROVIDERS: raise HTTPException(status_code=400, detail=f"Provider inconnu: {provider}") user_key = get_user_llm_key(user_id, provider) if user_key and user_key.get("api_key"): llm = LLMClient(provider=provider, api_key=user_key["api_key"], api_base=(user_key.get("api_base") or "").strip() or None) else: llm = LLMClient(provider=provider) offline = llm.provider == "offline" or not llm._has_credentials() system = ( "Tu es l'assistant d'édition de FlowDeck (contenu). " "Réponds UNIQUEMENT avec le texte demandé, en Markdown léger " "(paragraphes, listes à puces, titres si utile). " "N'ajoute aucun préambule, aucun commentaire, aucun bloc de code autour du texte. " "Fonde ta réponse sur le contexte du document fourni." ) user_content = prompt if context and context.strip(): user_content += "\n\n# Contexte du document\n" + context.strip()[:20000] messages = [ {"role": "system", "content": system}, {"role": "user", "content": user_content}, ] try: resp = await llm.complete(messages, model=model or None, tools=None, stream=False) except Exception as exc: # noqa: BLE001 — surface connectivity errors return {"ok": False, "error": str(exc), "offline": offline, "model": model or ""} return { "ok": True, "text": (resp.text or "").strip(), "model": resp.model or model or "", "offline": offline, } # ── AI Writing Assist (v5.9.0) ── @router.post("/writing") async def agent_writing(request: Request): """Headless writing actions for the editor (v5.9.0). `action` is one of: write, summarize, translate, continue, autocomplete. Returns plain Markdown (`text`) plus `model`/`offline` metadata. The `properties` action is served by `/writing/properties`. """ from app.services.ai_writing import WRITING_ACTIONS, AIWritingService user_id = await _current_user_id(request) body = await request.json() if request.headers.get("content-type") else {} action = (body.get("action") or "").strip().lower() if not action: raise HTTPException(status_code=400, detail="action est requis") if action == "properties": raise HTTPException(status_code=400, detail="Utilisez /writing/properties pour l'action properties") if action not in WRITING_ACTIONS: raise HTTPException(status_code=400, detail=f"Action inconnue: {action}") prompt = (body.get("prompt") or "").strip() if action == "write" and not prompt and not (body.get("title") or "").strip(): raise HTTPException(status_code=400, detail="prompt est requis pour l'action write") service = AIWritingService( user_id=user_id, provider=(body.get("provider") or "").strip() or None, model=(body.get("model") or "").strip() or None, ) result = await service.run( action, prompt=prompt, context=(body.get("context") or ""), target_language=(body.get("target_language") or "English"), prefix=(body.get("prefix") or ""), title=(body.get("title") or ""), ) return result @router.post("/writing/properties") async def agent_writing_properties(request: Request): """Suggest values for a collection page's properties (v5.9.0). Body: `{title, content?, properties: [{name, type}], provider?, model?}`. Returns `{ok, suggestions: {property_name: value}, model, offline}`. """ from app.services.ai_writing import AIWritingService user_id = await _current_user_id(request) body = await request.json() if request.headers.get("content-type") else {} properties = body.get("properties") or [] if not isinstance(properties, list) or not properties: raise HTTPException(status_code=400, detail="properties est requis (liste non vide)") service = AIWritingService( user_id=user_id, provider=(body.get("provider") or "").strip() or None, model=(body.get("model") or "").strip() or None, ) suggestions = await service.suggest_properties( context=(body.get("content") or ""), title=(body.get("title") or ""), properties=properties, ) return { "ok": True, "suggestions": suggestions, "model": service.model or "", "offline": service._offline_hint(), } # ── Audit & rollback ── @router.get("/conversations/{conversation_id}/actions") def list_actions(request: Request, conversation_id: int): with get_conn() as conn: rows = conn.execute( "SELECT * FROM agent_actions WHERE conversation_id=? ORDER BY created_at, id", (conversation_id,), ).fetchall() return {"actions": [dict(r) for r in rows]} @router.post("/actions/{action_id}/undo") def undo(request: Request, action_id: int): try: undo_action(action_id) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc except Exception as exc: # noqa: BLE001 raise HTTPException(status_code=500, detail=f"Rollback échoué: {exc}") from exc return {"id": action_id, "status": "reverted"} # ── Skills ── @router.get("/skills") async def list_skills(request: Request): ws = await _workspace_id(request) with get_conn() as conn: rows = conn.execute("SELECT * FROM agent_skills WHERE workspace_id IS ? OR workspace_id=? ORDER BY name", (ws, ws)).fetchall() return {"skills": [dict(r) for r in rows]} @router.post("/skills") async def create_skill(request: Request): user_id = await _current_user_id(request) ws = await _workspace_id(request) body = await request.json() if request.headers.get("content-type") else {} name = (body.get("name") or "").strip() if not name: raise HTTPException(status_code=400, detail="name est requis") with get_conn() as conn: try: cur = conn.execute( """INSERT INTO agent_skills (workspace_id, name, description, prompt_template, allowed_tools_json, created_by) VALUES (?,?,?,?,?,?)""", (ws, name, body.get("description", ""), body.get("prompt_template", ""), json.dumps(body.get("allowed_tools", [])), user_id), ) conn.commit() except Exception as exc: # noqa: BLE001 raise HTTPException(status_code=409, detail=f"Skill existe déjà: {exc}") from exc return {"id": cur.lastrowid, "name": name, "status": "created"} @router.post("/skills/{skill_id}/apply") async def apply_skill(request: Request, skill_id: int): """Create a conversation pre-loaded with a skill, ready to run.""" user_id = await _current_user_id(request) ws = await _workspace_id(request) with get_conn() as conn: skill = conn.execute("SELECT * FROM agent_skills WHERE id=?", (skill_id,)).fetchone() if not skill: raise HTTPException(status_code=404, detail="Skill introuvable") agent = _default_agent(conn, user_id) cur = conn.execute( """INSERT INTO agent_conversations (agent_id, user_id, title, context_json) VALUES (?,?,?,?)""", (agent["id"], user_id, skill["name"], json.dumps({"workspace_id": ws, "skill_id": skill_id})), ) conn.commit() return {"conversation_id": cur.lastrowid, "skill": skill["name"], "status": "ready"} # ── Skill marketplace (v6.6.0, Agent phase 5) ── # Galerie de presets + export/import portable — même implémentation que # l'API publique (/api/v2/skills/*), via app.services.skill_gallery. @router.get("/skills/gallery") def skills_gallery(request: Request): presets = skill_gallery.list_gallery() return {"gallery": presets, "total": len(presets), "install": "POST /api/agent/skills/gallery/{slug}/install"} @router.post("/skills/gallery/{slug}/install") async def install_gallery_skill(request: Request, slug: str): user_id = await _current_user_id(request) ws = await _workspace_id(request) preset = skill_gallery.get_gallery(slug) if not preset: raise HTTPException(status_code=404, detail=f"Skill inconnue dans la galerie: {slug}") body = await request.json() if request.headers.get("content-type") else {} try: row, created = skill_gallery.upsert_skill( skill_gallery.parse_payload(preset), workspace_id=ws, created_by=user_id, overwrite=bool(body.get("overwrite", True)), ) except ValueError as exc: raise HTTPException(status_code=409, detail=str(exc)) from exc return {"slug": slug, "id": row.get("id"), "name": row.get("name"), "status": "installed" if created else "updated", "skill": row} @router.post("/skills/import") async def import_skill(request: Request): """Importe un skill portable (JSON exporté depuis une autre instance).""" user_id = await _current_user_id(request) ws = await _workspace_id(request) body = await request.json() if request.headers.get("content-type") else {} payload = body.get("payload") if isinstance(body.get("payload"), dict) else body try: fields = skill_gallery.parse_payload(payload) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc try: row, created = skill_gallery.upsert_skill( fields, workspace_id=ws, created_by=user_id, overwrite=bool(body.get("overwrite")), ) except ValueError as exc: raise HTTPException(status_code=409, detail=str(exc)) from exc return {"id": row.get("id"), "name": fields["name"], "status": "imported" if created else "updated", "skill": row} @router.get("/skills/{skill_id}/export") def export_skill(request: Request, skill_id: int): """Document JSON portable — à rejouer sur /api/agent/skills/import.""" with get_conn() as conn: row = conn.execute("SELECT * FROM agent_skills WHERE id=?", (skill_id,)).fetchone() if not row: raise HTTPException(status_code=404, detail="Skill introuvable") return skill_gallery.export_skill(row) @router.delete("/skills/{skill_id}") def delete_skill(request: Request, skill_id: int): with get_conn() as conn: row = conn.execute("SELECT name FROM agent_skills WHERE id=?", (skill_id,)).fetchone() if not row: raise HTTPException(status_code=404, detail="Skill introuvable") conn.execute("DELETE FROM agent_skills WHERE id=?", (skill_id,)) conn.commit() return {"id": skill_id, "status": "deleted"} # ── Mentions (commande @ / +) & feedback (boutons 👍 / 👎) ── @router.get("/mentions") async def list_mentions(request: Request, q: str = ""): """Éléments mentionnables dans le panneau agent (commande « @ » / bouton « + »). Retourne des sections d'objets FlowDeck que l'utilisateur peut épingler au contexte du chat. La première section privilégie les **fichiers de l'espace de travail courant ouvert** (`workspace_id`), puis viennent les collections, les pages de collection et les autres documents. Chaque élément porte un jeton ``token`` (``document:``, ``collection:``, ``page:``) que ContextBuilder résout en contenu réel pour le LLM. Les entrées identiques (même type + même titre, ex. plusieurs pages « Sans titre ») sont **dédupliquées** pour ne pas afficher le même fichier plusieurs fois. """ query = (q or "").strip()[:60] like = f"%{query}%" limit = 8 # L'espace courant ouvert est fourni par le frontend (localWsId) et doit primer # sur l'espace par défaut de l'utilisateur pour cette liste contextuelle. try: ws = int(request.query_params.get("workspace_id", 0) or 0) or None except (TypeError, ValueError): ws = None if not ws: ws = await _workspace_id(request) def dedupe(items: list[dict]) -> list[dict]: seen: set = set() out: list[dict] = [] for it in items: key = (it["type"], (it.get("label") or "").strip().lower()) if key in seen: continue seen.add(key) out.append(it) return out sections: list[dict] = [] ws_name = None with get_conn() as conn: if ws: row = conn.execute("SELECT name FROM workspaces WHERE id=?", (ws,)).fetchone() ws_name = row["name"] if row else None # 1) Fichiers de l'espace de travail courant ouvert — token document: if ws: sql_ws = ( "SELECT p.id, p.title FROM pages p " "WHERE p.workspace_id=? AND p.deleted_at IS NULL AND (? = '' OR p.title LIKE ?) " "ORDER BY p.updated_at DESC LIMIT ?" ) files = [ {"type": "document", "icon": "📄", "id": r["id"], "label": r["title"] or "Sans titre", "sub": ws_name or "Espace courant", "token": f"document:{r['id']}"} for r in conn.execute(sql_ws, (ws, query, like, limit)).fetchall() ] sections.append({"key": "files", "label": (ws_name or "Fichiers de l'espace"), "items": dedupe(files)}) # 2) Collections (bases de données) — token collection: colls = [] for r in conn.execute( "SELECT id, name, icon FROM collections WHERE (? = '' OR name LIKE ?) " "ORDER BY updated_at DESC LIMIT ?", (query, like, limit) ).fetchall(): colls.append({"type": "collection", "icon": r["icon"] or "🗄️", "id": r["id"], "label": r["name"], "sub": "Base de données", "token": f"collection:{r['id']}"}) if colls: sections.append({"key": "collections", "label": "Collections / bases", "items": dedupe(colls)}) # 3) Pages de collection — token page: pages = [] for r in conn.execute( "SELECT cp.id, cp.title, cp.icon, c.name AS coll_name " "FROM collection_pages cp LEFT JOIN collections c ON c.id=cp.collection_id " "WHERE (? = '' OR cp.title LIKE ?) ORDER BY cp.updated_at DESC LIMIT ?", (query, like, limit) ).fetchall(): pages.append({ "type": "page", "icon": r["icon"] if str(r["icon"]).startswith(("📄", "📑", "🗂️", "📁")) else "📑", "id": r["id"], "label": r["title"] or "Sans titre", "sub": f"Page · {r['coll_name']}" if r["coll_name"] else "Page", "token": f"page:{r['id']}", }) if pages: sections.append({"key": "pages", "label": "Pages de collections", "items": dedupe(pages)}) # 4) Autres documents éditeur hors de l'espace courant — token document: if ws: other = [] for r in conn.execute( "SELECT p.id, p.title, w.name AS ws_name " "FROM pages p LEFT JOIN workspaces w ON w.id=p.workspace_id " "WHERE p.deleted_at IS NULL AND (p.workspace_id IS ? OR p.workspace_id IS NULL) " "AND (? = '' OR p.title LIKE ?) ORDER BY p.updated_at DESC LIMIT ?", (ws, query, like, limit) ).fetchall(): other.append({"type": "document", "icon": "📄", "id": r["id"], "label": r["title"] or "Sans titre", "sub": r["ws_name"] or "", "token": f"document:{r['id']}"}) dedup_other = dedupe(other) if dedup_other: sections.append({"key": "other", "label": "Autres documents", "items": dedup_other}) else: # Pas de workspace courant → liste les documents éditeur (comme avant) docs = [] for r in conn.execute( "SELECT p.id, p.title, w.name AS ws_name FROM pages p " "LEFT JOIN workspaces w ON w.id=p.workspace_id " "WHERE p.deleted_at IS NULL AND (? = '' OR p.title LIKE ?) " "ORDER BY p.updated_at DESC LIMIT ?", (query, like, limit) ).fetchall(): docs.append({"type": "document", "icon": "📄", "id": r["id"], "label": r["title"] or "Sans titre", "sub": r["ws_name"] or "", "token": f"document:{r['id']}"}) if docs: sections.insert(0, {"key": "files", "label": "Documents", "items": dedupe(docs)}) return {"workspace_id": ws, "sections": sections} @router.post("/feedback") async def add_feedback(request: Request): """Enregistre le retour (👍 / 👎) porté sur une réponse de l'agent.""" user_id = await _current_user_id(request) body = await request.json() if request.headers.get("content-type") else {} rating = (body.get("rating") or "").strip().lower() if rating not in ("up", "down"): raise HTTPException(status_code=400, detail="rating doit être 'up' ou 'down'") conversation_id = body.get("conversation_id") message_id = body.get("message_id") snippet = (body.get("content") or "")[:1000].strip() comment = (body.get("comment") or "").strip()[:2000] try: conversation_id = int(conversation_id) if conversation_id not in (None, "") else None except (TypeError, ValueError): conversation_id = None try: message_id = int(message_id) if message_id not in (None, "") else None except (TypeError, ValueError): message_id = None with get_conn() as conn: cur = conn.execute( """INSERT INTO agent_feedback (conversation_id, message_id, user_id, rating, snippet, comment) VALUES (?,?,?,?,?,?)""", (conversation_id, message_id, user_id, rating, snippet, comment), ) conn.commit() return {"id": cur.lastrowid, "rating": rating, "status": "recorded"} # ── Triggers & tools ── @router.post("/{agent_id}/trigger") async def trigger_agent(request: Request, agent_id: int): """Manually fire a custom agent: create a conversation and run it with the agent's instructions as the objective (falls back to a generic prompt).""" user_id = await _current_user_id(request) ws = await _workspace_id(request) with get_conn() as conn: agent = conn.execute("SELECT * FROM agents WHERE id=?", (agent_id,)).fetchone() if not agent: raise HTTPException(status_code=404, detail="Agent introuvable") cur = conn.execute( """INSERT INTO agent_conversations (agent_id, user_id, title, context_json) VALUES (?,?,?,?)""", (agent_id, user_id, f"Run: {agent['name']}", json.dumps({"workspace_id": ws})), ) conv_id = cur.lastrowid conn.commit() objective = (agent["system_instructions"] or "").strip() or f"Exécute l'agent « {agent['name']} »." engine = AgentEngine(user_id, workspace_id=ws) async def event_stream(): async for ev in engine.run(conv_id, objective, model=agent["model"]): yield f"data: {json.dumps(ev, ensure_ascii=False)}\n\n" return StreamingResponse( event_stream(), media_type="text/event-stream", headers={"Cache-Control": "no-cache"}, ) @router.get("/tools") def list_tools(request: Request): registry = ToolRegistry() tools = registry.schema() return {"tools": tools} @router.get("/providers") async def list_providers(request: Request): """Expose the LLM providers list + current config to the UI. Every provider known by FlowDeck is returned (the Settings UI lets the user configure any of them). Each provider also carries activation flags used by the Agent panel to offer only the providers that are *configured and functional* for the current user: - ``configured`` : a usable credential exists (own key OR the workspace default provider carries a key, OR no key is required). - ``verified`` : the last connection test / model fetch succeeded. - ``functional`` : the provider is ready to chat (verified, or `offline`). """ user_id = await _current_user_id(request) llm = LLMClient() cfg = get_llm_config() keys = list_user_llm_keys(user_id) keys_map = {k["provider"]: k for k in keys} default_provider = (cfg.get("provider") or "offline").lower() global_has_key = bool(cfg.get("api_key")) and default_provider != "offline" global_verified = bool(cfg.get("verified")) and global_has_key provs = provider_info() for p in provs: k = keys_map.get(p["id"]) if k and k["models"]: for m in k["models"]: if m not in p["models"]: p["models"].append(m) user_has_key = bool(k and k["has_key"]) is_default = p["id"] == default_provider configured = bool(k and k["has_key"]) or (is_default and global_has_key) \ or not p["requires_key"] user_verified = bool(k and k.get("verified")) verified = (user_verified and bool(k and k["has_key"])) or \ (is_default and global_has_key and global_verified) if not p["requires_key"] and not (k and k["has_key"]) and not (is_default and global_has_key): # ollama / offline: no key needed, but only offline is usable as-is; # local providers must still pass a connection test. verified = bool(user_verified or (is_default and global_verified)) \ if p["id"] != "offline" else True p["has_key"] = user_has_key p["configured"] = bool(configured) p["source"] = ("user" if user_has_key else "global" if is_default and global_has_key else "open") p["verified"] = bool(verified) p["verified_model"] = (k.get("verified_model") if user_has_key and k else cfg.get("verified_model") if is_default and global_has_key else "") or "" p["last_error"] = (k.get("last_error") if user_has_key and k else cfg.get("last_error") if is_default and global_has_key else "") or "" p["functional"] = p["id"] == "offline" or bool(verified) return { "provider": llm.provider, "model": llm.default_model, "available": await llm.is_available(), "max_iterations": settings.agent_max_iterations, "api_base": cfg["api_base"], "has_api_key": bool(cfg["api_key"]), "default_provider": default_provider, "default_verified": bool(cfg.get("verified")), "providers": provs, "keys": keys, } # ── Per-user provider API keys (non-admin: any authenticated user) ── @router.get("/keys") async def list_llm_keys(request: Request): """The user's saved provider keys + API keys (masked).""" user_id = await _current_user_id(request) return {"keys": list_user_llm_keys(user_id)} @router.put("/keys/{llm_provider}") async def save_llm_key(request: Request, llm_provider: str): """Upsert a provider key for the current user (masked in responses).""" user_id = await _current_user_id(request) provider = llm_provider.lower() if provider not in PROVIDERS: raise HTTPException(status_code=400, detail=f"Provider inconnu: {provider}") body = await request.json() if request.headers.get("content-type") else {} api_base_raw = body.get("api_base") raw = upsert_user_llm_key( user_id, provider, api_key=(body.get("api_key") or "").strip(), # None = keep the stored base, "" = reset to the provider default. api_base=api_base_raw.strip() if isinstance(api_base_raw, str) else None, default_model=(body.get("default_model") or "").strip(), models=body.get("models"), ) stored = bool(raw.get("api_key")) return {"status": "saved", "provider": provider, "has_key": stored} @router.delete("/keys/{llm_provider}") async def delete_llm_key(request: Request, llm_provider: str): """Remove a saved provider key for the current user.""" user_id = await _current_user_id(request) provider = llm_provider.lower() if provider not in PROVIDERS: raise HTTPException(status_code=400, detail=f"Provider inconnu: {provider}") delete_user_llm_key(user_id, provider) return {"status": "deleted", "provider": provider} @router.post("/keys/{llm_provider}/test") async def test_user_llm_key(request: Request, llm_provider: str): """Non-admin: verify one of the user's own provider keys (no mock fallback). On success the provider is flagged ``verified`` so it can be offered in the Agent panel; on failure the stored error is kept for display in Settings. """ user_id = await _current_user_id(request) provider = llm_provider.lower() if provider not in PROVIDERS: raise HTTPException(status_code=400, detail=f"Provider inconnu: {provider}") body = await request.json() if request.headers.get("content-type") else {} stored = get_user_llm_key(user_id, provider) api_key = (body.get("api_key") or "").strip() api_base = (body.get("api_base") or "").strip() if not api_key and stored and stored.get("api_key"): api_key = stored["api_key"] if not api_base: api_base = (stored.get("api_base") or "").strip() if stored else "" _, provider_default = PROVIDERS.get(provider, (None, "gpt-4o")) candidates = PROVIDER_MODELS.get(provider) or [provider_default] model = (body.get("model") or "").strip() or \ ((stored or {}).get("default_model") or "") or \ next((m for m in candidates if m), provider_default) # Only record the result against the *stored* credential when it is the one # being tested (avoids flagging a stored key from a test run on a typed key). can_record = bool(stored) and ( not (stored or {}).get("api_key") or api_key == (stored or {}).get("api_key") ) llm = LLMClient(provider=provider, api_key=api_key, api_base=api_base or None) try: resp = await llm.ping(model=model or None) except Exception as exc: # noqa: BLE001 — surface connectivity errors if can_record: mark_user_llm_key_verified(user_id, provider, False, error=str(exc)) return {"ok": False, "provider": provider, "error": str(exc), "verified": False} if not stored and provider in ("ollama", "offline"): stored = upsert_user_llm_key(user_id, provider, api_key="", api_base=api_base) can_record = True if can_record: mark_user_llm_key_verified( user_id, provider, True, model=resp.model or model or "", ) return { "ok": True, "provider": provider, "model": resp.model or model or "", "reply": (resp.text or "").strip()[:200], "verified": bool(can_record), } @router.post("/keys/{llm_provider}/models") async def fetch_llm_models(request: Request, llm_provider: str): """Fetch the live model list from a provider. Falls back to the user's stored key when no key is supplied in the body. A successful fetch proves connectivity, so when it used the *stored* key the provider is flagged ``verified`` (functional) for the Agent panel. """ user_id = await _current_user_id(request) provider = llm_provider.lower() if provider not in PROVIDERS: raise HTTPException(status_code=400, detail=f"Provider inconnu: {provider}") body = await request.json() if request.headers.get("content-type") else {} api_key = (body.get("api_key") or "").strip() api_base = (body.get("api_base") or "").strip() used_stored = False if not api_key: stored = get_user_llm_key(user_id, provider) if stored and stored.get("api_key"): api_key = stored["api_key"] api_base = api_base or (stored.get("api_base") or "").strip() used_stored = True try: models = await fetch_provider_models(provider, api_key=api_key, api_base=api_base) if used_stored and models: mark_user_llm_key_verified(user_id, provider, True, model=(body.get("default_model") or "").strip()) return {"ok": True, "provider": provider, "models": models} except Exception as exc: # noqa: BLE001 — surface connectivity errors return {"ok": False, "provider": provider, "error": str(exc)} def _check_api_base(value: str) -> str: """A14 : `api_base` doit être une URL http(s) sans identifiants. ponytail: les hôtes PRIVÉS restent acceptés — le provider par défaut du produit est `http://localhost:11434/v1` (Ollama, `llm_client.PROVIDERS`) et le verrou nommé par l'audit (un anonymous qui oriente le `ping()` du serveur) est neutralisé par `_current_admin` (401 sans session / 403 non admin). Pour verrouiller plus tard : allowlist des providers locaux ou un settings `llm_allow_private=false`. """ url = (value or "").strip() if not url: return "" from urllib.parse import urlparse parsed = urlparse(url) if parsed.scheme not in ("http", "https") or not parsed.netloc: raise HTTPException(status_code=400, detail=f"api_base invalide: {url!r}") if parsed.username or parsed.password: raise HTTPException(status_code=400, detail="api_base ne doit pas contenir d'identifiants") return url @router.patch("/providers") async def update_provider_config(request: Request): await _current_admin(request) body = await request.json() if request.headers.get("content-type") else {} provider = (body.get("provider") or "").strip().lower() if provider and provider not in PROVIDERS: raise HTTPException(status_code=400, detail=f"Provider inconnu: {provider}") _ = set_llm_config( provider=provider or None, model=(body.get("model") or "").strip() or None, api_key=body.get("api_key"), api_base=_check_api_base(body.get("api_base") or "") or None, clear_keys=(provider == "offline"), ) llm = LLMClient() return { "status": "saved", "provider": llm.provider, "model": llm.default_model, "available": await llm.is_available(), "api_base": (body.get("api_base") or "").strip() or "", "has_api_key": bool(llm.api_key), } @router.post("/providers/test") async def test_provider_config(request: Request): """Admin: verify a provider is reachable (no mock fallback). A successful test flags the workspace default provider as ``verified`` so it becomes available (functional) for every user in the Agent panel. """ await _current_admin(request) body = await request.json() if request.headers.get("content-type") else {} provider = (body.get("provider") or "").strip().lower() or None if provider and provider not in PROVIDERS: raise HTTPException(status_code=400, detail=f"Provider inconnu: {provider}") llm = LLMClient( provider=provider, api_key=body.get("api_key"), api_base=_check_api_base(body.get("api_base") or "") or None, ) try: resp = await llm.ping(model=(body.get("model") or "").strip() or None) except Exception as exc: # noqa: BLE001 — surface real connectivity errors cfg = get_llm_config() if (provider or cfg.get("provider")) == cfg.get("provider"): mark_llm_config_verified(False, error=str(exc)) return {"ok": False, "error": str(exc), "verified": False} cfg = get_llm_config() if (provider or cfg.get("provider")) == cfg.get("provider"): mark_llm_config_verified(True, model=resp.model or llm.default_model or "") return { "ok": True, "model": resp.model or (body.get("model") or "").strip() or llm.default_model, "reply": (resp.text or "").strip()[:200], "verified": True, } # ── Agents by id (registered LAST so static routes /tools, /skills, … win) ── @router.get("/{agent_id}") def get_agent(request: Request, agent_id: int): with get_conn() as conn: row = conn.execute("SELECT * FROM agents WHERE id=?", (agent_id,)).fetchone() if not row: raise HTTPException(status_code=404, detail="Agent introuvable") return dict(row) @router.put("/{agent_id}") async def update_agent(request: Request, agent_id: int): body = await request.json() if request.headers.get("content-type") else {} with get_conn() as conn: existing = conn.execute("SELECT * FROM agents WHERE id=?", (agent_id,)).fetchone() if not existing: raise HTTPException(status_code=404, detail="Agent introuvable") sets, params = [], [] for col in ("name", "icon", "description", "system_instructions", "model", "approval_mode", "is_active"): if col in body: sets.append(f"{col}=?") params.append(body[col]) if "scope" in body: sets.append("scope_json=?") params.append(json.dumps(body["scope"])) if "trigger" in body: sets.append("trigger_json=?") params.append(json.dumps(body["trigger"])) if sets: params.append(agent_id) conn.execute(f"UPDATE agents SET {', '.join(sets)} WHERE id=?", params) conn.commit() return {"id": agent_id, "status": "updated"} @router.delete("/{agent_id}") def delete_agent(request: Request, agent_id: int): with get_conn() as conn: existing = conn.execute("SELECT * FROM agents WHERE id=?", (agent_id,)).fetchone() if not existing: raise HTTPException(status_code=404, detail="Agent introuvable") conn.execute("DELETE FROM agents WHERE id=?", (agent_id,)) conn.commit() return {"id": agent_id, "status": "deleted"}