FlowDeck CI / lint (push) Successful in 43s
FlowDeck CI / test (push) Successful in 4m2s
FlowDeck CI / lint (pull_request) Successful in 42s
FlowDeck CI / test (pull_request) Successful in 4m3s
FlowDeck CI / docker (push) Successful in 1m2s
FlowDeck CI / docker (pull_request) Successful in 35s
tests/conftest.py: mutate the settings singleton (instead of rebinding) so DB + backup dir are isolated per test -> pytest-xdist safe. Real backup tests (snapshot/prune/admin API) and OAuth mock tests (Gitea/GitHub/link) replace the previous skips. init_db() now also creates webhook_subscriptions (full schema without the FastAPI lifespan). ruff check is clean; .eslintrc.json migrated to eslint.config.mjs (flat config). CI: lint job (ruff + eslint), parallel tests (-n auto), run on every branch push. VERSION 5.11.1.
1029 lines
44 KiB
Python
1029 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, HTTPException, Request
|
|
from fastapi.responses import StreamingResponse
|
|
|
|
from app.auth.session import get_current_user
|
|
from app.config import settings
|
|
from app.db import get_conn
|
|
from app.services.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.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}") from exc
|
|
return {"id": cur.lastrowid, "name": name, "status": "created"}
|
|
|
|
|
|
# ── Conversations ──
|
|
|
|
|
|
@router.get("/conversations")
|
|
async def list_conversations(request: Request):
|
|
user_id = await _current_user_id(request)
|
|
with get_conn() as conn:
|
|
rows = conn.execute(
|
|
"SELECT * FROM agent_conversations WHERE user_id=? ORDER BY updated_at DESC",
|
|
(user_id,),
|
|
).fetchall()
|
|
return {"conversations": [dict(r) for r in rows]}
|
|
|
|
|
|
@router.post("/conversations")
|
|
async def create_conversation(request: Request):
|
|
user_id = await _current_user_id(request)
|
|
body = await request.json() if request.headers.get("content-type") else {}
|
|
ws = await _workspace_id(request)
|
|
with get_conn() as conn:
|
|
agent = _default_agent(conn, user_id)
|
|
cur = conn.execute(
|
|
"""INSERT INTO agent_conversations (agent_id, user_id, title, context_json, provider, model)
|
|
VALUES (?,?,?,?,?,?)""",
|
|
(agent["id"], user_id, body.get("title", "New conversation"),
|
|
json.dumps({"workspace_id": ws}),
|
|
body.get("provider", ""), body.get("model", "")),
|
|
)
|
|
conn.commit()
|
|
return {"id": cur.lastrowid, "title": body.get("title", "New conversation"), "status": "created"}
|
|
|
|
|
|
@router.get("/conversations/{conversation_id}")
|
|
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 WRITING_ACTIONS, AIWritingService
|
|
|
|
user_id = await _current_user_id(request)
|
|
body = await request.json() if request.headers.get("content-type") else {}
|
|
action = (body.get("action") or "").strip().lower()
|
|
if not action:
|
|
raise HTTPException(status_code=400, detail="action est requis")
|
|
if action == "properties":
|
|
raise HTTPException(status_code=400, detail="Utilisez /writing/properties pour l'action properties")
|
|
if action not in WRITING_ACTIONS:
|
|
raise HTTPException(status_code=400, detail=f"Action inconnue: {action}")
|
|
|
|
prompt = (body.get("prompt") or "").strip()
|
|
if action == "write" and not prompt and not (body.get("title") or "").strip():
|
|
raise HTTPException(status_code=400, detail="prompt est requis pour l'action write")
|
|
|
|
service = AIWritingService(
|
|
user_id=user_id,
|
|
provider=(body.get("provider") or "").strip() or None,
|
|
model=(body.get("model") or "").strip() or None,
|
|
)
|
|
result = await service.run(
|
|
action,
|
|
prompt=prompt,
|
|
context=(body.get("context") or ""),
|
|
target_language=(body.get("target_language") or "English"),
|
|
prefix=(body.get("prefix") or ""),
|
|
title=(body.get("title") or ""),
|
|
)
|
|
return result
|
|
|
|
|
|
@router.post("/writing/properties")
|
|
async def agent_writing_properties(request: Request):
|
|
"""Suggest values for a collection page's properties (v5.9.0).
|
|
|
|
Body: `{title, content?, properties: [{name, type}], provider?, model?}`.
|
|
Returns `{ok, suggestions: {property_name: value}, model, offline}`.
|
|
"""
|
|
from app.services.ai_writing import AIWritingService
|
|
|
|
user_id = await _current_user_id(request)
|
|
body = await request.json() if request.headers.get("content-type") else {}
|
|
properties = body.get("properties") or []
|
|
if not isinstance(properties, list) or not properties:
|
|
raise HTTPException(status_code=400, detail="properties est requis (liste non vide)")
|
|
|
|
service = AIWritingService(
|
|
user_id=user_id,
|
|
provider=(body.get("provider") or "").strip() or None,
|
|
model=(body.get("model") or "").strip() or None,
|
|
)
|
|
suggestions = await service.suggest_properties(
|
|
context=(body.get("content") or ""),
|
|
title=(body.get("title") or ""),
|
|
properties=properties,
|
|
)
|
|
return {
|
|
"ok": True,
|
|
"suggestions": suggestions,
|
|
"model": service.model or "",
|
|
"offline": service._offline_hint(),
|
|
}
|
|
|
|
|
|
# ── Audit & rollback ──
|
|
|
|
|
|
@router.get("/conversations/{conversation_id}/actions")
|
|
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)) from exc
|
|
except Exception as exc: # noqa: BLE001
|
|
raise HTTPException(status_code=500, detail=f"Rollback échoué: {exc}") from exc
|
|
return {"id": action_id, "status": "reverted"}
|
|
|
|
|
|
# ── Skills ──
|
|
|
|
|
|
@router.get("/skills")
|
|
async def list_skills(request: Request):
|
|
ws = await _workspace_id(request)
|
|
with get_conn() as conn:
|
|
rows = conn.execute("SELECT * FROM agent_skills WHERE workspace_id IS ? OR workspace_id=? ORDER BY name", (ws, ws)).fetchall()
|
|
return {"skills": [dict(r) for r in rows]}
|
|
|
|
|
|
@router.post("/skills")
|
|
async def create_skill(request: Request):
|
|
user_id = await _current_user_id(request)
|
|
ws = await _workspace_id(request)
|
|
body = await request.json() if request.headers.get("content-type") else {}
|
|
name = (body.get("name") or "").strip()
|
|
if not name:
|
|
raise HTTPException(status_code=400, detail="name est requis")
|
|
with get_conn() as conn:
|
|
try:
|
|
cur = conn.execute(
|
|
"""INSERT INTO agent_skills (workspace_id, name, description, prompt_template, allowed_tools_json, created_by)
|
|
VALUES (?,?,?,?,?,?)""",
|
|
(ws, name, body.get("description", ""), body.get("prompt_template", ""),
|
|
json.dumps(body.get("allowed_tools", [])), user_id),
|
|
)
|
|
conn.commit()
|
|
except Exception as exc: # noqa: BLE001
|
|
raise HTTPException(status_code=409, detail=f"Skill existe déjà: {exc}") from exc
|
|
return {"id": cur.lastrowid, "name": name, "status": "created"}
|
|
|
|
|
|
@router.post("/skills/{skill_id}/apply")
|
|
async def apply_skill(request: Request, skill_id: int):
|
|
"""Create a conversation pre-loaded with a skill, ready to run."""
|
|
user_id = await _current_user_id(request)
|
|
ws = await _workspace_id(request)
|
|
with get_conn() as conn:
|
|
skill = conn.execute("SELECT * FROM agent_skills WHERE id=?", (skill_id,)).fetchone()
|
|
if not skill:
|
|
raise HTTPException(status_code=404, detail="Skill introuvable")
|
|
agent = _default_agent(conn, user_id)
|
|
cur = conn.execute(
|
|
"""INSERT INTO agent_conversations (agent_id, user_id, title, context_json)
|
|
VALUES (?,?,?,?)""",
|
|
(agent["id"], user_id, skill["name"], json.dumps({"workspace_id": ws, "skill_id": skill_id})),
|
|
)
|
|
conn.commit()
|
|
return {"conversation_id": cur.lastrowid, "skill": skill["name"], "status": "ready"}
|
|
|
|
|
|
# ── 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"}
|