Files
flowdeck/app/routers/agent.py
T
bruno f09de98406
FlowDeck CI / test (push) Failing after 25s
FlowDeck CI / docker (push) Skipped
feat(v5.9.0): AI Writing Assist - slash /ai, autocompletion, AI properties
- Service app/services/ai_writing.py: 6 actions sans outils (write, summarize,
  translate, continue, autocomplete, properties) + replis deterministes offline
- Endpoints POST /api/agent/writing et /api/agent/writing/properties
- Editeur: groupe slash AI (Write/Summarize/Translate/Continue) via E.aiSlash
- Autocompletion inline AIAC (suggestion ~900ms, Tab accepte, Escape rejette)
- Database: bouton AI fill (suggestions de proprietes Status/Priority/Resume)
- Tests tests/test_ai_writing.py (+29) ; version 5.9.0 (VERSION, main.py)
- CHANGELOG + ROADMAP v5.9.0 completes
2026-09-10 11:24:21 -04:00

1023 lines
44 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
from fastapi import APIRouter, Request, HTTPException
from fastapi.responses import StreamingResponse
from app.db import get_conn
from app.config import settings
from app.auth.session import get_current_user
from app.services.agent_engine import AgentEngine, undo_action
from app.services.llm_client import LLMClient, PROVIDERS, PROVIDER_MODELS
from app.services.llm_config import (
get_llm_config, set_llm_config, provider_info,
get_user_llm_key, list_user_llm_keys, upsert_user_llm_key,
delete_user_llm_key, fetch_provider_models,
mark_llm_config_verified, mark_user_llm_key_verified,
)
from app.services.tool_registry import ToolRegistry
from app.services.permission_manager import PermissionManager
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.utcnow()
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 | None:
user = await get_current_user(request)
if user and user.get("id"):
return user["id"]
with get_conn() as conn:
row = conn.execute("SELECT id FROM users WHERE login='admin' ORDER BY id LIMIT 1").fetchone()
return row["id"] if row else None
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:
"""Require an admin session. Falls back to the single admin row, matching
the agent router's unauthenticated convention (single-user deployments)."""
user = await get_current_user(request)
if user:
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
with get_conn() as conn:
row = conn.execute("SELECT * FROM users WHERE login='admin' ORDER BY id LIMIT 1").fetchone()
if not row or not row["is_admin"]:
raise HTTPException(status_code=403, detail="Accès administrateur requis")
return dict(row)
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}")
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}")
async 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}")
async 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 AIWritingService, WRITING_ACTIONS
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")
async 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")
async def undo(request: Request, action_id: int):
try:
undo_action(action_id)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc))
except Exception as exc: # noqa: BLE001
raise HTTPException(status_code=500, detail=f"Rollback échoué: {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}")
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"}
# ── 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:<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 = 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:<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")
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")
async 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 {}
raw = upsert_user_llm_key(
user_id,
provider,
api_key=(body.get("api_key") or "").strip(),
api_base=(body.get("api_base") or "").strip(),
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)}
@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=(body.get("api_base") or "").strip() 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=(body.get("api_base") or "").strip() 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}")
async 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}")
async 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"}