- app/services/plugins.py + migration 35 : table plugins (slug, name,
description, enabled) pré-remplie avec 3 modules câblés — web-tools,
web-clipper, automations ; ligne absente = activé (défaut sûr)
- automations OFF → dépendance FastAPI posée à l'include_router dans main.py
(aucun router touché) → toutes les routes /workspace/automations* refusées +
garde de tick du scheduler en arrière-plan
- web-clipper OFF → GET /extensions et tout /api/v2/web-clipper/* refusés
- web-tools OFF → web_search et fetch_url retirés du schéma ET de execute()
via ToolRegistry._all() : le LLM ne les voit plus
- UI rendue côté serveur : global Jinja plugin_enabled(slug) — nav
« Extensions » / « Automations » en {% if %} (absentes du DOM), sections
conditionnées en x-show dans settings.html
- menu + : l'entrée « Add plugins » devient vivante (fini disabled:true) —
liste des 3 plugins avec bascule, GET/PATCH /api/agent/plugins[/slug]
(slug inconnu → 404, 401 sans session)
- tests : tests/test_v758_plugins.py (10 tests) — routes refusées (302 hors
/api, 404 JSON pour /api*), outils retirés, nav disparue, persistance,
câblage ; assertions disabled:true == 0 dans les tests des phases 1/3/4/5/7
- livraison : VERSION + app/main = 7.58.0, OpenAPI 525 chemins, CHANGELOG,
ROADMAP phase 8 cochée (menu + complet), avenant phase 8 (docs)
1348 lines
56 KiB
Python
1348 lines
56 KiB
Python
"""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
|
||
import secrets
|
||
from datetime import UTC
|
||
from urllib.parse import urlparse
|
||
|
||
from fastapi import APIRouter, Body, HTTPException, Request
|
||
from fastapi.responses import JSONResponse, RedirectResponse, StreamingResponse
|
||
|
||
from app.auth.session import get_current_user
|
||
from app.config import settings
|
||
from app.db import get_conn
|
||
from app.services import connectors, oauth_connectors, plugins, 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")
|
||
|
||
|
||
def _current_user_id(request: Request) -> int:
|
||
"""A14 : plus de fallback sur la row `admin` — 401 sans session."""
|
||
user = get_current_user(request)
|
||
if not user or not user.get("id"):
|
||
raise HTTPException(status_code=401, detail="Authentication required")
|
||
return user["id"]
|
||
|
||
|
||
def _workspace_id(request: Request) -> int | None:
|
||
user = 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
|
||
|
||
|
||
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 = 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("")
|
||
def list_agents(request: Request):
|
||
user_id = _current_user_id(request)
|
||
ws = _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("")
|
||
def create_agent(request: Request, body: dict = Body(default={})):
|
||
user_id = _current_user_id(request)
|
||
ws = _workspace_id(request)
|
||
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")
|
||
def list_conversations(request: Request):
|
||
user_id = _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")
|
||
def create_conversation(request: Request, body: dict = Body(default={})):
|
||
user_id = _current_user_id(request)
|
||
ws = _workspace_id(request)
|
||
with get_conn() as conn:
|
||
agent = _default_agent(conn, user_id)
|
||
mem_default = 1 if settings.agent_memory_default else 0
|
||
cur = conn.execute(
|
||
"""INSERT INTO agent_conversations (agent_id, user_id, title, context_json, provider, model, memory_enabled)
|
||
VALUES (?,?,?,?,?,?,?)""",
|
||
(agent["id"], user_id, body.get("title", "New conversation"),
|
||
json.dumps({"workspace_id": ws}),
|
||
body.get("provider", ""), body.get("model", ""), mem_default),
|
||
)
|
||
conn.commit()
|
||
return {"id": cur.lastrowid, "title": body.get("title", "New conversation"),
|
||
"status": "created", "memory_enabled": mem_default}
|
||
|
||
|
||
@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}")
|
||
def patch_conversation(request: Request, conversation_id: int, body: dict = Body(default={})):
|
||
"""Update a conversation's title / provider / model (slash-command support)."""
|
||
user_id = _current_user_id(request)
|
||
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]))
|
||
# v7.54.0 — toggle « Mémoire » de la conversation (0/1).
|
||
if body.get("memory_enabled") is not None:
|
||
sets.append("memory_enabled=?")
|
||
params.append(1 if body["memory_enabled"] else 0)
|
||
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, body: dict = Body(default={})):
|
||
user_id = _current_user_id(request)
|
||
ws = _workspace_id(request)
|
||
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 = _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 = _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 = _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")
|
||
def list_skills(request: Request):
|
||
ws = _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")
|
||
def create_skill(request: Request, body: dict = Body(default={})):
|
||
user_id = _current_user_id(request)
|
||
ws = _workspace_id(request)
|
||
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.patch("/skills/{skill_id}")
|
||
def update_skill(request: Request, skill_id: int, body: dict = Body(default={})):
|
||
"""Édition d'un skill enregistré (v7.52.0 — « Gérer les compétences »).
|
||
|
||
Seuls les champs présents dans ``body`` sont mis à jour ; la liste des
|
||
colonnes est figée ici (jamais interpolée depuis la requête).
|
||
"""
|
||
_current_user_id(request)
|
||
fields: dict = {}
|
||
if "name" in body:
|
||
name = str(body.get("name") or "").strip()
|
||
if not name:
|
||
raise HTTPException(status_code=400, detail="name est requis")
|
||
fields["name"] = name
|
||
if "description" in body:
|
||
fields["description"] = str(body.get("description") or "")
|
||
if "prompt_template" in body:
|
||
fields["prompt_template"] = str(body.get("prompt_template") or "")
|
||
if "allowed_tools" in body:
|
||
fields["allowed_tools_json"] = json.dumps(body.get("allowed_tools") or [])
|
||
if not fields:
|
||
raise HTTPException(status_code=400, detail="Aucun champ à mettre à jour")
|
||
cols = ", ".join(f"{k}=?" for k in fields)
|
||
with get_conn() as conn:
|
||
if not conn.execute("SELECT id FROM agent_skills WHERE id=?", (skill_id,)).fetchone():
|
||
raise HTTPException(status_code=404, detail="Skill introuvable")
|
||
try:
|
||
conn.execute(
|
||
f"UPDATE agent_skills SET {cols} WHERE id=?",
|
||
(*fields.values(), skill_id),
|
||
)
|
||
conn.commit()
|
||
except Exception as exc: # noqa: BLE001 — contrainte UNIQUE(name)
|
||
raise HTTPException(status_code=409, detail=f"Skill existe déjà: {exc}") from exc
|
||
return {"id": skill_id, "status": "updated", "fields": sorted(fields)}
|
||
|
||
|
||
@router.post("/skills/{skill_id}/apply")
|
||
def apply_skill(request: Request, skill_id: int):
|
||
"""Create a conversation pre-loaded with a skill, ready to run."""
|
||
user_id = _current_user_id(request)
|
||
ws = _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")
|
||
def install_gallery_skill(request: Request, slug: str, body: dict = Body(default={})):
|
||
user_id = _current_user_id(request)
|
||
ws = _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}")
|
||
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")
|
||
def import_skill(request: Request, body: dict = Body(default={})):
|
||
"""Importe un skill portable (JSON exporté depuis une autre instance)."""
|
||
user_id = _current_user_id(request)
|
||
ws = _workspace_id(request)
|
||
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"}
|
||
|
||
|
||
# ── Connecteurs (v7.55.0) ──
|
||
# ponytail : catalogue partagé au même titre que les skills (`list_skills`
|
||
# n'a pas non plus de filtre par créateur) — la clé n'est jamais renvoyée.
|
||
|
||
|
||
@router.get("/connectors")
|
||
def list_connectors_route(request: Request):
|
||
"""Catalogue : natifs (statut dérivé de la config) + personnels."""
|
||
user_id = _current_user_id(request)
|
||
return {"connectors": connectors.list_connectors(user_id)}
|
||
|
||
|
||
@router.get("/plugins")
|
||
def list_plugins_route(request: Request):
|
||
"""Catalogue des plugins de l'instance (menu + → « Add plugins »)."""
|
||
_current_user_id(request)
|
||
return {"plugins": plugins.list_plugins()}
|
||
|
||
|
||
@router.patch("/plugins/{slug}")
|
||
def set_plugin_route(request: Request, slug: str, body: dict = Body(default={})):
|
||
"""Bascule un plugin : l'effet est réel (routes/UI/outils), pas un drapeau."""
|
||
_current_user_id(request)
|
||
try:
|
||
return plugins.set_enabled(slug, bool(body.get("enabled")))
|
||
except ValueError as exc:
|
||
raise HTTPException(status_code=404, detail=str(exc)) from exc
|
||
|
||
|
||
@router.post("/connectors")
|
||
def create_connectors_route(request: Request, body: dict = Body(default={})):
|
||
"""Ajoute un connecteur personnalisé — l'URL est validée (garde SSRF)."""
|
||
_current_user_id(request)
|
||
try:
|
||
return connectors.create_connector(
|
||
body.get("name") or "", body.get("url") or "", str(body.get("secret") or ""),
|
||
kind=str(body.get("kind") or "custom"),
|
||
auth=str(body.get("auth") or "bearer"),
|
||
)
|
||
except ValueError as exc:
|
||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||
|
||
|
||
@router.patch("/connectors/{connector_id}")
|
||
def update_connectors_route(request: Request, connector_id: int, body: dict = Body(default={})):
|
||
_current_user_id(request)
|
||
try:
|
||
row = connectors.update_connector(
|
||
connector_id,
|
||
name=body.get("name"),
|
||
enabled=body.get("enabled"),
|
||
secret=body.get("secret"),
|
||
)
|
||
except ValueError as exc:
|
||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||
if row is None:
|
||
raise HTTPException(status_code=404, detail="Connecteur introuvable")
|
||
return row
|
||
|
||
|
||
@router.delete("/connectors/{connector_id}")
|
||
def delete_connectors_route(request: Request, connector_id: int):
|
||
_current_user_id(request)
|
||
if not connectors.delete_connector(connector_id):
|
||
raise HTTPException(status_code=404, detail="Connecteur introuvable")
|
||
return {"id": connector_id, "status": "deleted"}
|
||
|
||
|
||
@router.post("/connectors/probe")
|
||
async def probe_connectors_route(request: Request, body: dict = Body(default={})):
|
||
"""Teste un connecteur (id, nom ou kind natif) et persiste le résultat."""
|
||
target = str(body.get("connector") or "").strip()
|
||
if not target:
|
||
raise HTTPException(status_code=400, detail="connector est requis")
|
||
user_id = _current_user_id(request)
|
||
try:
|
||
return await connectors.probe(target, user_id)
|
||
except ValueError as exc:
|
||
raise HTTPException(status_code=404, detail=str(exc)) from exc
|
||
|
||
|
||
# ── Connecteurs OAuth : Google / Microsoft 365 (v7.56.0) ──
|
||
# Flow : POST authorize (cookies state/verifier/next) → fournisseur →
|
||
# GET callback (échange du code, tokens chiffrés) → redirection vers la page
|
||
# d'origine. Le GET est en SAFE_METHODS : pas de CSRF (comme /auth/callback).
|
||
|
||
|
||
def _oauth_kind(kind: str) -> str:
|
||
if kind not in oauth_connectors.OAUTH_KINDS:
|
||
raise HTTPException(status_code=404, detail="Connecteur OAuth inconnu")
|
||
return kind
|
||
|
||
|
||
def _oauth_next(request: Request, kind: str, key: str, default: str = "/") -> str:
|
||
"""Chemin de retour same-origin (cookie `key` pour le callback)."""
|
||
raw = request.cookies.get(key, "") or default
|
||
path = urlparse(raw).path if raw.startswith("http") else raw
|
||
if not path.startswith("/") or path.startswith("//"):
|
||
return default
|
||
return path
|
||
|
||
|
||
@router.get("/connectors/oauth/{kind}/status")
|
||
def connector_oauth_status(request: Request, kind: str):
|
||
"""État du connecteur OAuth pour l'utilisateur courant."""
|
||
_oauth_kind(kind)
|
||
user_id = _current_user_id(request)
|
||
try:
|
||
oauth_connectors._client(kind)
|
||
configured, detail = True, ""
|
||
except ValueError as exc:
|
||
configured, detail = False, str(exc)
|
||
tok = oauth_connectors.tokens(kind, user_id)
|
||
return {
|
||
"kind": kind, "configured": configured, "configured_detail": detail,
|
||
"connected": bool(tok.get("access_token")), "scope": tok.get("scope", ""),
|
||
"expires_at": tok.get("expires_at", 0),
|
||
}
|
||
|
||
|
||
@router.post("/connectors/oauth/{kind}/authorize")
|
||
def connector_oauth_authorize(request: Request, kind: str):
|
||
"""URL d'autorisation (PKCE) + cookies d'état éphémères (10 min)."""
|
||
_oauth_kind(kind)
|
||
_current_user_id(request)
|
||
try:
|
||
flow = oauth_connectors.begin(kind)
|
||
except ValueError as exc:
|
||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||
referer = request.headers.get("referer") or ""
|
||
next_path = urlparse(referer).path if referer else "/"
|
||
if not next_path.startswith("/") or next_path.startswith("//"):
|
||
next_path = "/"
|
||
secure = request.url.scheme == "https"
|
||
resp = JSONResponse({"url": flow["url"], "kind": kind})
|
||
for key, value in (
|
||
(f"fd_oauth_state_{kind}", flow["state"]),
|
||
(f"fd_oauth_verifier_{kind}", flow["verifier"]),
|
||
(f"fd_oauth_next_{kind}", next_path),
|
||
):
|
||
resp.set_cookie(key, value, max_age=600, httponly=True, samesite="lax", secure=secure)
|
||
return resp
|
||
|
||
|
||
@router.get("/connectors/oauth/{kind}/callback")
|
||
async def connector_oauth_callback(
|
||
request: Request, kind: str, code: str = "", state: str = "", error: str = ""
|
||
):
|
||
"""Reçoit le code du fournisseur, échange, stocke, renvoie sur la page."""
|
||
_oauth_kind(kind)
|
||
next_path = _oauth_next(request, kind, f"fd_oauth_next_{kind}")
|
||
expected = request.cookies.get(f"fd_oauth_state_{kind}", "")
|
||
verifier = request.cookies.get(f"fd_oauth_verifier_{kind}", "")
|
||
|
||
def _back(flag: str) -> RedirectResponse:
|
||
sep = "&" if "?" in next_path else "?"
|
||
resp = RedirectResponse(f"{next_path}{sep}{flag}", status_code=302)
|
||
for key in (f"fd_oauth_state_{kind}", f"fd_oauth_verifier_{kind}",
|
||
f"fd_oauth_next_{kind}"):
|
||
resp.delete_cookie(key, samesite="lax")
|
||
return resp
|
||
|
||
if error:
|
||
return _back("oauth_error=" + error[:80])
|
||
if not code or not expected or not secrets.compare_digest(expected, state):
|
||
return _back("oauth_error=state")
|
||
user = get_current_user(request)
|
||
if not user or not user.get("id"):
|
||
return _back("oauth_error=session")
|
||
try:
|
||
tokens = await oauth_connectors.complete(kind, code, verifier)
|
||
except ValueError as exc:
|
||
return _back("oauth_error=" + str(exc)[:80].replace(" ", "+"))
|
||
oauth_connectors.save(kind, int(user["id"]), tokens)
|
||
logger.info("Connector OAuth connected: %s (user #%s)", kind, user["id"])
|
||
return _back("oauth=connected")
|
||
|
||
|
||
@router.post("/connectors/oauth/{kind}/disconnect")
|
||
def connector_oauth_disconnect(request: Request, kind: str):
|
||
_oauth_kind(kind)
|
||
user_id = _current_user_id(request)
|
||
removed = oauth_connectors.clear(kind, user_id)
|
||
return {"kind": kind, "status": "disconnected" if removed else "not-connected"}
|
||
|
||
|
||
# ── Mentions (commande @ / +) & feedback (boutons 👍 / 👎) ──
|
||
|
||
|
||
@router.get("/mentions")
|
||
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:<id>``, ``collection:<id>``, ``page:<id>``) 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 = _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:<id>
|
||
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:<id>
|
||
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:<id>
|
||
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:<id>
|
||
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")
|
||
def add_feedback(request: Request, body: dict = Body(default={})):
|
||
"""Enregistre le retour (👍 / 👎) porté sur une réponse de l'agent."""
|
||
user_id = _current_user_id(request)
|
||
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 = _current_user_id(request)
|
||
ws = _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 = _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")
|
||
def list_llm_keys(request: Request):
|
||
"""The user's saved provider keys + API keys (masked)."""
|
||
user_id = _current_user_id(request)
|
||
return {"keys": list_user_llm_keys(user_id)}
|
||
|
||
|
||
@router.put("/keys/{llm_provider}")
|
||
def save_llm_key(request: Request, llm_provider: str, body: dict = Body(default={})):
|
||
"""Upsert a provider key for the current user (masked in responses)."""
|
||
user_id = _current_user_id(request)
|
||
provider = llm_provider.lower()
|
||
if provider not in PROVIDERS:
|
||
raise HTTPException(status_code=400, detail=f"Provider inconnu: {provider}")
|
||
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}")
|
||
def delete_llm_key(request: Request, llm_provider: str):
|
||
"""Remove a saved provider key for the current user."""
|
||
user_id = _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 = _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 = _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):
|
||
_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.
|
||
"""
|
||
_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}")
|
||
def update_agent(request: Request, agent_id: int, body: dict = Body(default={})):
|
||
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"}
|