Files
flowdeck/app/services/agent_engine.py
T
bruno ba363eaee9
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
feat(v5.2.0): finalize Infrastructure & Polish (tests isolation, xdist, lint, CI)
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.
2026-09-11 23:36:53 -04:00

457 lines
24 KiB
Python

"""FlowDeck — AgentEngine: ReAct orchestrator (v4.14.0).
`objective → comprehension → context → reasoning ↔ action → result`.
The engine drives the LLM (which only emits tool intentions), gates each call
through PermissionManager, executes it via ToolRegistry, journals every action
to `agent_actions` with an undo snapshot, and yields a stream of SSE events so
the UI can render reasoning + actions live. Since v4.14.0 a freshly created
conversation is automatically renamed with a descriptive title derived from its
content so the history stays easy to browse.
"""
from __future__ import annotations
import asyncio
import json
import logging
import re
from app.config import settings
from app.db import get_conn
from app.services.context_builder import ContextBuilder
from app.services.llm_client import LLMClient
from app.services.permission_manager import PermissionManager
from app.services.tool_registry import ToolRegistry
logger = logging.getLogger(__name__)
MAX_ITERATIONS = 12
# Compact in-app guide so the LLM can answer « comment faire… ? » questions even
# when no document is attached to the conversation (generic help / onboarding).
APP_GUIDE = """## Guide de l'utilisateur FlowDeck (sert à répondre aux questions « comment … ? »)
- **Pages** : le contenu est organisé en blocs (paragraphes, titres, listes, to-do, tableaux, images, formules, bases embarquées). La barre latérale liste les pages récentes, favoris, agents, partagées et publiées.
- **Documents & espaces de travail** : un « document » est une page éditeur (type Notion) qui vit dans un espace de travail. Pour créer un document dans un espace : appelle `read_workspaces` (reprends le `workspace_name` ou l'id exact), puis `create_document`. Pour modifier un document existant : `read_document` puis `write_blocks` (blocs et/ou titre). Pour supprimer : `delete_document` (corbeille). `search_workspace` retrouve aussi les documents et les espaces par titre.
- **Format des blocs** (pour `write_blocks`) : chaque bloc est un objet `{"type": "...", "content": "texte"}`. Le champ du texte s'appelle **`content`** (jamais `text`). Un script / code s'écrit dans un bloc `{"type": "code", "content": "...", "language": "powershell"}`. Les titres sont `heading_1`, `heading_2`, `heading_3`, `heading_4`. Autres types : paragraph, bulleted_list, numbered_list, to_do, quote, divider, toggle, callout.
- **Collections (bases de données)** : des ensembles de pages structurées avec des propriétés (texte, nombre, sélection, dates…). Chaque collection peut avoir plusieurs vues : tableau, board (kanban), calendrier, galerie, liste, timeline, graphique, formulaire, carte, flux, gantt. Ajouter une propriété ou une vue = outils add_property / create_view.
- **Créer du contenu** : « crée une collection X », « crée une page », « ajoute une propriété Statut à la collection Y » sont des actions que l'agent peut exécuter directement avec ses outils.
- **Espaces de travail** : FlowDeck gère des espaces locaux et des dépôts Gitea/GitHub (pages privées dans un dépôt, issues reliées via read_gitea_issues). On change d'espace depuis le menu en bas à gauche (« Switch workspace »).
- **Recherche** : la commande Ctrl+K / la barre de recherche du haut permet de retrouver pages et collections.
- **Corbeille & Bibliothèque** : les pages supprimées vont dans la Corbeille ; Favoris / Récents / Partagés / Publiés se consultent dans la Bibliothèque.
- **Réglages** : Paramètres (en bas à gauche → Settings) pour le compte, les notifications, les tags, les intégrations et la section « Agent & IA » (clés API, fournisseurs, modèle global).
- **Agent IA** : ouvrable via le bouton 🤖 en bas à droite ou la section « Agents » du sidebar. On peut lui parler de la page ouverte, ou lui poser des questions générales sur l'utilisation de l'application.
Quand la question est générale (« comment créer un kanban ? », « où sont mes favoris ? »), réponds de façon concise et guidée à partir de ces informations, sans inventer de fonctionnalités absentes."""
# Deterministic auto-title heuristics (used when no real LLM is configured, and
# as a fallback when the generated title is unusable). Ordered by priority: the
# first matching intent wins.
_TITLE_INTENTS = (
("Création", ("créer", "crée", "crées", "création", "nouveau", "nouvelle",
"create", "creation")),
("Ajout", ("ajouter", "ajoute", "ajout d", "ajoutons", "add")),
("Renommage", ("renommer", "renomme", "renommage", "rename")),
("Suppression", ("supprimer", "supprime", "suppression", "delete")),
("Déplacement", ("déplacer", "déplace", "déplacement", "move")),
("Mise à jour", ("modifier", "modifie", "modification", "mets à jour",
"met à jour", "mettre à jour", "update", "éditer")),
("Analyse", ("analyser", "analyse")),
("Résumé", ("résumer", "résume", "résumé", "resume")),
("Traduction", ("traduire", "traduis", "traduit", "traduction", "translate")),
("Planification", ("planifier", "planifie", "préparer", "prépare", "organiser",
"organise", "sprint", "agenda")),
("Recherche", ("chercher", "cherche", "rechercher", "recherche", "trouver",
"trouve", "liste", "lister", "search", "find")),
)
_TITLE_TYPES = (
("collection", "collection", ("collection", "base de données", "database", "db")),
("propriété", "propriété", ("propriété", "property")),
("vue", "vue", (" vue", "view")),
("board", "board", ("board", "kanban")),
("sprint", "sprint", ("sprint",)),
("document", "document", ("document", "note de réunion", "compte-rendu", "compte rendu")),
("tâche", "tâche", ("tâche", "task", "tache")),
("issue", "issue", ("issue",)),
)
class AgentEngine:
def __init__(self, user_id: int, workspace_id: int | None = None,
llm: LLMClient | None = None):
self.user_id = user_id
self.workspace_id = workspace_id
self.llm = llm or LLMClient()
self.tools = ToolRegistry()
self.ctx = ContextBuilder(user_id, workspace_id)
self.perms = PermissionManager(user_id)
self._tokens = 0
# ── Helpers ──
def _load_agent(self, conversation_id: int) -> dict:
with get_conn() as conn:
row = conn.execute(
"SELECT a.* FROM agents a JOIN agent_conversations c ON c.agent_id=a.id WHERE c.id=?",
(conversation_id,),
).fetchone()
if not row:
row = {"id": None, "name": "FlowDeck Agent", "icon": "🤖", "agent_type": "personal",
"system_instructions": "", "model": settings.llm_model,
"scope_json": "{}", "approval_mode": "auto"}
return dict(row)
def _build_system_prompt(self, agent: dict, skills=None) -> str:
if skills is None:
skills = []
if isinstance(skills, dict):
skills = [skills]
lines = [
"Tu es FlowDeck Agent, un agent IA qui réalise des tâches dans le workspace FlowDeck.",
"Tu réfléchis (reasoning) puis agis en appelant les outils disponibles.",
"Appelle UN ou PLUSIEURS outils pour atteindre l'objectif, puis conclus avec une réponse finale.",
"N'invente jamais d'IDs : utilise ceux fournis dans le contexte.",
f"Workspace courant : {self.workspace_id}.",
]
if agent.get("system_instructions"):
lines.append(f"\nInstructions de l'agent {agent.get('name','')}:\n{agent['system_instructions']}")
# Plusieurs skills peuvent être appliqués au même post : chacun injecte
# son prompt dans les instructions système.
for skill in skills:
if skill:
lines.append(f"\nSkill appliquée « {skill.get('name','')} »:\n{skill.get('prompt_template','')}")
# L'utilisateur peut poser des questions d'aide sans contexte de document ;
# le guide intégré permet d'y répondre (aucun outil requis).
lines.append("\n" + APP_GUIDE)
return "\n".join(lines)
# ── Main run (async generator of SSE events) ──
async def run(self, conversation_id: int, objective: str, *, model: str | None = None,
mentions: list[str] | None = None, files: list[dict] | None = None,
skill_id: int | None = None, skill_ids: list[int] | None = None,
extra_context: str | None = None):
agent = self._load_agent(conversation_id)
scope = json.loads(agent.get("scope_json") or "{}")
approval_mode = agent.get("approval_mode") or "auto"
model = model or agent.get("model") or settings.llm_model
ids = list(skill_ids or [])
if skill_id and skill_id not in ids:
ids.append(skill_id)
skills = [s for s in (self._load_skill(i, scope) for i in ids) if s]
system = self._build_system_prompt(agent, skills)
context = self.ctx.build(mentions=mentions, files=files)
if extra_context and extra_context.strip():
context += "\n\n## Document / contexte fourni par l'utilisateur\n" + extra_context.strip()
messages = [
{"role": "system", "content": system},
{"role": "user", "content": f"{objective}\n\n# Contexte\n{context}"},
]
self._persist_message(conversation_id, "user", objective)
self._update_conversation(conversation_id, status="running")
# Update the history title right away (before the run finishes) and
# refine it once we have the final answer (_autotitle below).
try:
suggested = self._suggest_title(objective, None)
if suggested:
self._update_conversation(conversation_id, title=suggested[:80])
except Exception: # noqa: BLE001 — never break a run because of the title
logger.exception("Auto-title failed for conversation #%s", conversation_id)
tool_schema = self.tools.schema(scope)
final_text = None
used_model = model or "" # peut être ajusté par un repli de modèle (404/410)
try:
for _step in range(settings.agent_max_iterations or MAX_ITERATIONS):
if self._tokens >= settings.agent_max_tokens_budget:
yield self._event("error", {"message": "Budget de tokens dépassé"})
break
response = await asyncio.wait_for(
self.llm.complete(messages, model=model, tools=tool_schema, stream=True),
timeout=settings.agent_run_timeout_seconds,
)
self._tokens += response.usage.get("total_tokens", 0) or 0
if getattr(response, "notice", ""):
yield self._event("notice", {"message": response.notice})
used_model = getattr(response, "model", "") or used_model
if response.text and response.text.strip():
yield self._event("reasoning", {"content": response.text})
if not response.tool_calls:
messages.append({"role": "assistant", "content": response.text or ""})
final_text = response.text or self._no_tool_message(response)
yield self._event("final", {"content": final_text})
break
# L'API de chat exige que le message assistant qui *annonce* les appels
# d'outils porte les `tool_calls` (avec id), puis que chaque résultat
# d'outil soit fourni avec le `tool_call_id` correspondant. Sans cela
# la passe suivante est refusée par le fournisseur (et l'agent retombait
# silencieusement sur le mock hors-ligne).
tool_specs = []
for idx, call in enumerate(response.tool_calls):
call_id = call.get("id") or f"call_{conversation_id}_{idx}_{self._tokens}"
tool_specs.append({
"id": call_id,
"type": "function",
"function": {
"name": call["name"],
"arguments": call.get("arguments_raw")
or json.dumps(call.get("arguments") or {}, ensure_ascii=False),
},
})
assistant_msg = {"role": "assistant", "content": response.text or ""}
assistant_msg["tool_calls"] = tool_specs
messages.append(assistant_msg)
for idx, call in enumerate(response.tool_calls):
tool, args = call["name"], call.get("arguments") or {}
call_id = tool_specs[idx]["id"]
denied = False
try:
self.perms.assert_can(tool, args, self.workspace_id, approval_mode)
except Exception as exc: # permission / approval guard
detail = self._exc_detail(exc)
yield self._event("action", {"tool": tool, "status": "error", "detail": detail})
self._log_action(conversation_id, tool, args, {}, "error", detail=detail)
messages.append({
"role": "tool", "tool_call_id": call_id,
"content": json.dumps({"status": "error", "message": f"Permission refusée: {detail}"}, ensure_ascii=False),
})
denied = True
if not denied:
result = await self.tools.execute(tool, args, user_id=self.user_id)
if result.status == "success":
yield self._event("action", {
"tool": tool, "status": result.status,
"target_type": result.target_type, "target_id": result.target_id,
"message": result.message,
})
self._log_action(conversation_id, tool, args, result.data, "success",
target_type=result.target_type, target_id=result.target_id,
undo=result.undo)
messages.append({
"role": "tool", "tool_call_id": call_id,
"content": json.dumps({"status": "ok", "result": result.data, "target_id": result.target_id}, ensure_ascii=False),
})
else:
yield self._event("action", {"tool": tool, "status": "error", "detail": result.message})
self._log_action(conversation_id, tool, args, {}, "error", detail=result.message)
messages.append({
"role": "tool", "tool_call_id": call_id,
"content": json.dumps({"status": "error", "message": result.message}, ensure_ascii=False),
})
if final_text is None:
final_text = "Objectif traité. Consultez le journal des actions pour le détail."
yield self._event("final", {"content": final_text})
self._persist_message(conversation_id, "assistant", final_text,
model=used_model, tokens=self._tokens)
await self._autotitle(conversation_id, objective, final_text)
except Exception as exc: # noqa: BLE001
logger.exception("AgentEngine run failed")
yield self._event("error", {"message": f"Erreur interne: {exc}"})
finally:
self._update_conversation(conversation_id, status="idle")
# ── Skills ──
def _load_skill(self, skill_id: int, scope: dict | None) -> dict | None:
with get_conn() as conn:
row = conn.execute("SELECT * FROM agent_skills WHERE id=?", (skill_id,)).fetchone()
if not row:
return None
skill = dict(row)
allowed = json.loads(skill.get("allowed_tools_json") or "[]")
if allowed:
skill["prompt_template"] = (skill.get("prompt_template") or "") + \
"\nOutils autorisés: " + ", ".join(allowed)
return skill
# ── Audit & persistence ──
def _log_action(self, conversation_id, tool, args, result, status,
*, target_type="", target_id=None, undo=None, detail=""):
with get_conn() as conn:
conn.execute(
"""INSERT INTO agent_actions
(conversation_id, tool_name, target_type, target_id, payload_json,
result_json, status, undo_snapshot_json, executed_by)
VALUES (?,?,?,?,?,?,?,?,?)""",
(conversation_id, tool, target_type,
str(target_id) if target_id is not None else None,
json.dumps(args, ensure_ascii=False),
json.dumps(result, ensure_ascii=False, default=str),
status,
json.dumps(undo or {}, ensure_ascii=False),
self.user_id),
)
conn.commit()
def _persist_message(self, conversation_id, role, content, *, model="", tokens=0):
with get_conn() as conn:
conn.execute(
"INSERT INTO agent_messages (conversation_id, role, content, model, tokens_used) VALUES (?,?,?,?,?)",
(conversation_id, role, content, model, tokens),
)
conn.commit()
def _update_conversation(self, conversation_id, *, status=None, title=None):
with get_conn() as conn:
sets, params = ["updated_at=CURRENT_TIMESTAMP"], []
if status:
sets.append("status=?")
params.append(status)
if title:
sets.append("title=?")
params.append(title)
params.append(conversation_id)
conn.execute(f"UPDATE agent_conversations SET {', '.join(sets)} WHERE id=?", params)
conn.commit()
# ── Misc ──
@staticmethod
def _event(etype: str, data: dict) -> dict:
return {"type": etype, **data}
@staticmethod
def _no_tool_message(response) -> str:
return "Je n'ai pas d'action à proposer pour cet objectif. Posez-moi une question plus précise ou demandez-moi de créer un élément."
@staticmethod
def _exc_detail(exc: Exception) -> str:
detail = getattr(exc, "detail", None)
return detail if isinstance(detail, str) else str(exc)
# ── Auto-title (v4.14.0) ──
async def _autotitle(self, conversation_id: int, objective: str, final_text: str | None):
"""Rename the conversation with a descriptive title derived from the
*latest* user request. Runs after every AI call so the history list
always reflects the current topic and stays easy to browse."""
try:
suggested = self._suggest_title(objective, final_text)
if suggested:
self._update_conversation(conversation_id, title=suggested[:80])
except Exception: # noqa: BLE001 — never break a run because of the title
logger.exception("Auto-title failed for conversation #%s", conversation_id)
@classmethod
def _suggest_title(cls, objective: str | None, final_text: str | None) -> str:
"""Produce a short descriptive title from the user objective (offline-safe)."""
text = (objective or "").split("\n# Contexte", 1)[0].strip() or (final_text or "").strip()
if not text:
return "Conversation"
# Strip the composer prefixes ("Contexte « X »", "Skill « Y »") that the
# frontend prepends before the real user text.
text = re.sub(
r"(?:Contexte\s*«[^»]*»|Skill\s*«[^»]*»|Skill\s+«[^»]*»)(?:\s*[,;\n.])+\s*",
"", text,
).strip()
if not text:
return "Conversation"
low = text.lower()
intent = None
for label, words in _TITLE_INTENTS:
if any(w in low for w in words):
intent = label
break
type_label = next(
(t for t, _noun, words in _TITLE_TYPES if any(w in low for w in words)), None
)
has_workspace = any(w in low for w in ("workspace", "espace de travail"))
quotes = [q.strip() for q in re.findall(r'[«"]([^«»"]{1,80})[»"]', text) if q.strip()]
subject = quotes[0] if quotes else None
ws = quotes[-1] if (has_workspace and len(quotes) > 1) else None
def _clean(s: str) -> str:
return re.sub(r"\s+", " ", s).strip(" .;:-")
if not subject and type_label:
noun = next((n for t, n, _w in _TITLE_TYPES if t == type_label), type_label)
m = re.search(
rf"\b{noun}\b\s*(?:nomm[ée]e?\s+|appel[ée]e?\s+|intitul[ée]e?\s+)?"
r'[«"]?\s*([A-Za-zÀ-ÿ0-9][A-Za-zÀ-ÿ0-9_ \-]{1,60}?)\s*[»"]?',
text, re.IGNORECASE,
)
if m:
subject = m.group(1).strip()
if intent and subject:
core = f"{intent} {type_label or 'élément'} « {subject} »" \
if type_label else f"{intent} « {subject} »"
if ws:
core += f" (dans {ws})"
return _clean(core)
generic = _clean(text)
return generic[:70] if generic else "Conversation"
def undo_action(action_id: int) -> bool:
"""Reverse a single agent action using its stored undo snapshot.
Returns True on success. Marks the action row `reverted`.
"""
with get_conn() as conn:
action = conn.execute("SELECT * FROM agent_actions WHERE id=?", (action_id,)).fetchone()
if not action:
raise ValueError(f"Action #{action_id} introuvable")
undo = json.loads(action["undo_snapshot_json"] or "{}")
op, table, rid = undo.get("action"), undo.get("table"), undo.get("id")
if not op or not table or rid is None:
raise ValueError(f"Action #{action_id} n'a pas de snapshot annulable")
if op == "delete":
conn.execute(f"DELETE FROM {table} WHERE id=?", (rid,))
elif op == "softdelete":
conn.execute(f"UPDATE {table} SET deleted_at=NULL WHERE id=?", (rid,))
elif op == "insert":
snapshot = undo.get("snapshot")
if not snapshot:
raise ValueError("Snapshot manquant pour insert")
cols = ", ".join(snapshot.keys())
ph = ", ".join("?" for _ in snapshot)
conn.execute(f"INSERT INTO {table} ({cols}) VALUES ({ph})", list(snapshot.values()))
elif op == "update":
snapshot = undo.get("snapshot")
if table == "collection_pages":
conn.execute(
"UPDATE collection_pages SET property_values_json=?, title=?, updated_at=CURRENT_TIMESTAMP WHERE id=?",
(json.dumps(snapshot, ensure_ascii=False), undo.get("title", ""), rid),
)
elif table == "pages":
if isinstance(snapshot, dict):
conn.execute(
"UPDATE pages SET content=?, content_format=?, title=?, updated_at=CURRENT_TIMESTAMP WHERE id=?",
(snapshot.get("content", ""),
snapshot.get("content_format", "markdown"),
snapshot.get("title", ""), rid),
)
else:
conn.execute("UPDATE pages SET content=?, updated_at=CURRENT_TIMESTAMP WHERE id=?",
(snapshot, rid))
else:
raise ValueError(f"Table non gérée pour rollback: {table}")
else:
raise ValueError(f"Opération de rollback inconnue: {op}")
conn.execute("UPDATE agent_actions SET status='reverted' WHERE id=?", (action_id,))
conn.commit()
return True