Files
flowdeck/app/routers/agent.py
T
2026-09-05 09:58:04 -04:00

424 lines
17 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
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
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)
VALUES (?,?,?,?)""",
(agent["id"], user_id, body.get("title", "New conversation"),
json.dumps({"workspace_id": ws})),
)
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"}
# ── 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")
engine = AgentEngine(user_id, workspace_id=ws)
provider = body.get("provider") or None
llm = LLMClient(provider=provider) if provider else None
if llm:
engine.llm = llm
async def event_stream():
async for ev in engine.run(
conversation_id, objective,
model=body.get("model"),
mentions=body.get("mentions"),
files=body.get("files"),
skill_id=body.get("skill_id"),
):
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"},
)
# ── 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"}
# ── 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 configured LLM provider + availability to the UI."""
llm = LLMClient()
return {
"provider": llm.provider,
"model": llm.default_model,
"available": await llm.is_available(),
"max_iterations": settings.agent_max_iterations,
}
# ── 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"}