Objectif : une anomalie injectée dans
process_signal_datadéclenche en < 60s un work_order avec RCA structuré, broadcasté en temps réel via WebSocket. Bloquant pour scènes 2 et 3 de la démo.
Scope. backend/core/ws_manager.py :
- Classe
WSManageravecconnections: set[WebSocket] await connect(ws),disconnect(ws)await broadcast(event_type: str, payload: dict)→ envoie JSON à tous les sockets, ignore les sockets fermés- Singleton module-level
ws_manager = WSManager()
Events utilisés (contrat aligné avec frontend, cf. docs/planning/ALIGNMENT.md).
| Event type | Payload | Émis par |
|---|---|---|
anomaly_detected |
{cell_id, signal_def_id, value, threshold, work_order_id, time} |
Sentinel |
agent_start |
{agent, turn_id} |
orchestrator |
agent_end |
{agent, turn_id, finish_reason} |
orchestrator |
tool_call_started |
{agent, tool_name, args, turn_id} |
Investigator / Q&A / WO Gen |
tool_call_completed |
{agent, tool_name, duration_ms, turn_id} |
idem |
agent_handoff |
{from_agent, to_agent, reason, turn_id} (cf. M4.6) |
orchestrator on ask_* tool call |
thinking_delta |
{agent, content, turn_id} (cf. M4.5) |
Investigator (extended thinking) |
ui_render |
{agent, component, props, turn_id} (cf. M2.9) |
tout agent appelant un render_* |
rca_ready |
{work_order_id, rca_summary, confidence, turn_id} |
Investigator |
work_order_ready |
{work_order_id} |
Work Order Generator |
turn_id. UUID v4 généré par l'orchestrateur à chaque agent_start.
Corrèle tous les events d'un même tour agentique côté frontend (Activity Feed,
Agent Inspector). Stocker dans une ContextVar Python pour ne pas le passer
explicitement à chaque ws_manager.broadcast().
Sérialisation. JSON une ligne par event. Pas d'event error — erreurs via
HTTP status ou {type: "done", error: "..."} final côté /agent/chat.
✅ DÉCIDÉ — un seul topic global. Le frontend filtre côté client par
cell_iddans le payload. Pour la démo on a 1 site, 5 cells, 1 opérateur connecté au max — le coup d'étoffer en rooms par cell n'apporte rien.
Acceptance.
-
wscat ws://localhost:8000/api/v1/events→ reçoit les events broadcastés - Endpoint
WS /api/v1/eventsenregistré dansmain.py
Scope. backend/agents/sentinel.py :
async def sentinel_loop(): bouclewhile Trueavecasyncio.sleep(30)- À chaque tick :
- Liste les cells qui ont une
equipment_kbaveconboarding_complete=true - Pour chaque cell : récupère les 5 dernières minutes de
process_signal_datapour les signaux référencés danskb.thresholds - Compare la dernière valeur vs
kb.thresholds.<signal>.alert - Si dépassement nouveau (pas déjà un work_order ouvert pour ce cell+signal dans
les 30 dernières minutes) :
- INSERT
work_order(status='detected', generated_by_agent=true, trigger_anomaly_time=value_time, triggered_by_signal_def_id=..., title="Anomalie détectée — <signal>", priority='high') ws_manager.broadcast("anomaly_detected", ...)- Broadcast
ui_renderavecrender_alert_banner(cf. M2.9) pour afficher le banner inline dans le dashboard/chat :ws_manager.broadcast("ui_render", {agent: "sentinel", component: "render_alert_banner", props: {severity, cell_id, message, anomaly_id: work_order_id}, turn_id: ...}) asyncio.create_task(run_investigator(work_order_id))
- INSERT
- Liste les cells qui ont une
Démarrage. Lancer dans le lifespan de main.py :
sentinel_task = asyncio.create_task(sentinel_loop())
yield
sentinel_task.cancel()
✅ DÉCIDÉ — cells sans KB skipées silencieusement. Au startup du
sentinel_loop, log INFO une fois la liste des cells surveillées vs ignorées. Pas de log répétitif à chaque tick. Côté frontend, le badge "Surveillé par ARIA" sur la card cell rend le statut évident.
✅ DÉCIDÉ — debounce via query DB. À chaque tick, avant d'INSERT un nouveau work_order, faire :
SELECT 1 FROM work_order WHERE cell_id = $1 AND triggered_by_signal_def_id = $2 AND created_at > NOW() - INTERVAL '30 minutes' AND status NOT IN ('completed','cancelled') LIMIT 1;Si row → skip. La DB est la source de vérité, survit aux restarts du Sentinel, et un humain qui ferme manuellement le WO débloque la surveillance immédiatement.
Acceptance.
- Simulateur monte vibration P-02 à 3.4 mm/s → 1 work_order créé en < 35s
- Pas de work_order doublon si la valeur reste élevée pendant 5 minutes
Scope. backend/agents/investigator.py :
async def run_investigator(work_order_id: int) -> None- Charge le contexte initial : work_order + cell + signal_def + anomaly_time
- Construit user message : "Anomalie détectée sur <cell.name>. Signal <signal.name> a atteint à (seuil alerte: ). Investigue."
- System prompt : "Tu es un expert maintenance industrielle. Une anomalie a été
détectée. Utilise les tools disponibles pour investiguer librement — tu décides
quoi consulter, dans quel ordre. Produis ensuite un RCA structuré au format JSON :
{root_cause, confidence, contributing_factors, similar_past_failure, recommended_action}." - Boucle agent (cf. pattern
technical.md§2.2 mais avecMCPClient) :tools_schema = await mcp_client.get_tools_schema() + UI_TOOLS + [SUBMIT_RCA_TOOL, ASK_KB_BUILDER_TOOL](UI_TOOLS cf. M2.9, ASK_KB_BUILDER_TOOL cf. M4.6)- while not end_turn :
messages.create(...)avec model agent + tools- pour chaque
tool_useblock :- si
tool_name.startswith("render_")→ broadcastui_render+ tool_result "rendered" (cf. M2.9) - si
tool_name == "ask_kb_builder"→ spawn mini-session KB Builder (cf. M3.5) + tool_result - si
tool_name == "submit_rca"→ capture args, break loop - sinon → broadcast
tool_call_started→mcp_client.call_tool()→ broadcasttool_call_completed
- si
- append assistant + tool_results à messages
- À end_turn : extract le bloc JSON RCA du dernier message texte
- UPDATE
work_order SET rca_summary=..., status='analyzed' - INSERT
failure_history(cell_id, failure_time, failure_mode, root_cause, signal_patterns, work_order_id) ws_manager.broadcast("rca_ready", ...)asyncio.create_task(run_work_order_generator(work_order_id))
✅ DÉCIDÉ — tool
submit_rca(Option A). Plus robuste que parser du markdown. Déclaré inline côté Python (pas via FastMCP server) :SUBMIT_RCA_TOOL = { "name": "submit_rca", "description": "À appeler une fois pour soumettre le RCA final.", "input_schema": {"type": "object", "properties": { "root_cause": {"type": "string"}, "confidence": {"type": "number"}, "contributing_factors": {"type": "array", "items": {"type": "string"}}, "similar_past_failure": {"type": "string"}, "recommended_action": {"type": "string"} }, "required": ["root_cause", "confidence", "recommended_action"]} } tools = await mcp_client.get_tools_schema() + [SUBMIT_RCA_TOOL]Quand l'agent appelle
submit_rca, on ne fait PAS de tool_result back — on stoppe la boucle, persiste le RCA, et break.
✅ DÉCIDÉ —
submit_rcaest local au module Investigator. Déclaré inline dansagents/investigator.pyet concaténé manuellement àtools_schema. Le MCP server n'expose JAMAISsubmit_rca. Même pattern poursubmit_work_order(Investigator only) et tout futur tool spécifique à un agent. Règle générale : tool d'output structuré = local agent ; tool de lecture/écriture DB partagée = MCP.
Acceptance.
- Anomalie injectée → work_order avec
rca_summarynon-null en < 60s -
failure_historycontient une nouvelle entrée - Frontend voit les
tool_call_startedevents streamer en live
Scope. Modifier backend/main.py lifespan pour :
- Démarrer
sentinel_taskau startup - L'annuler au shutdown
- Logger les exceptions du
sentinel_loopsans laisser la task mourir silencieusement (wrappertry/exceptavec re-raise contrôlé)
Acceptance.
-
docker compose up→ log "Sentinel started" visible -
docker compose down→ log "Sentinel cancelled"
Scope. Activer thinking sur l'agent loop Investigator. C'est le seul argument
visible "pourquoi Opus 4.7 vs Sonnet" pour les juges.
response = await anthropic.messages.create(
model=model_for("agent"),
thinking={"type": "enabled", "budget_tokens": 10000},
system=INVESTIGATOR_SYSTEM,
messages=messages,
tools=tools_schema,
stream=True,
)Streaming. À chaque chunk thinking_delta reçu du SDK Anthropic, broadcast :
await ws_manager.broadcast("thinking_delta", {
"agent": "investigator",
"content": chunk.thinking_delta.text,
"turn_id": turn_id,
})Périmètre. Activé uniquement sur Investigator. Budget 10k tokens ≈ 5¢ par run avec Opus 4.7, négligeable. Les autres agents n'en ont pas besoin (et économie coûts).
Pourquoi critique. Le frontend (M8.5 Agent Inspector) streame le thinking en live dans un panel dédié → le juge voit Opus réfléchir. Sans ça, "pourquoi Opus 4.7 ?" n'a pas de réponse visuelle. 25% de la note.
Acceptance.
-
thinking_deltaevents streamed pendant un run Investigator - Frontend M8.5 affiche le thinking en live (vérifié J6)
- Latence end-to-end Investigator reste < 60s avec thinking activé
Bloque. Frontend M8.5, prix "Opus 4.7 Use".
Scope. Remplacer le pipeline scripted (asyncio.create_task(run_work_order_generator)
en fin de M4.3) par des handoffs décidés par l'agent via tool call. Sans ça, le
"multi-agent" est juste un workflow Python aux yeux des juges.
Décision — garder les deux chemins.
- Pipeline scripted RCA → WO Gen reste en place (chemin garanti pour la démo)
- En plus, l'Investigator a des tools
ask_*qu'il peut choisir d'appeler en cours de raisonnement (chemin wow factor)
Tools à déclarer (locaux, pas via FastMCP).
ASK_KB_BUILDER_TOOL = {
"name": "ask_kb_builder",
"description": "Consulte le KB Builder pour un détail constructeur absent de la "
"KB courante (ex: torque max d'un boulon, réf pièce introuvable).",
"input_schema": {"type": "object", "properties": {
"question": {"type": "string"},
"cell_id": {"type": "integer"}
}, "required": ["question", "cell_id"]}
}Handler. Quand Investigator appelle ask_kb_builder :
ws_manager.broadcast("agent_handoff", {from: "investigator", to: "kb_builder", reason: args.question, turn_id})ws_manager.broadcast("agent_start", {agent: "kb_builder", turn_id: new_turn_id})- Spawn une mini-session KB Builder (Messages API loop, system prompt spécialisé "réponds factuellement à un collègue agent, format JSON")
ws_manager.broadcast("agent_end", {agent: "kb_builder", turn_id, finish_reason: "end_turn"})- Retour comme
tool_resultà Investigator
Symétrique côté Q&A (M5.2/M5.4). Déclarer ask_investigator pour que Q&A
délègue les questions diagnostiques pointues.
Scénarios démo (au moins 2 handoffs visibles).
- Scène 3 (Investigation) : Investigator → KB Builder pour chercher une réf pièce
- Scène 5 (Q&A) : Q&A → Investigator pour une question diagnostique poussée
Acceptance.
- Investigator peut appeler
ask_kb_builderet reçoit une réponse structurée - Event
agent_handoffvisible dans le WS stream - Scénario démo P-02 déclenche au moins 1 handoff dynamique
Bloque. Pitch "Best Managed Agents" crédible.
Scope. L'Investigator charge failure_history du cell dans son contexte initial
et produit un RCA visiblement plus rapide/précis au 2e diagnostic similaire.
Implem.
- Début du run :
past_failures = await mcp_client.call_tool("get_failure_history", {cell_id, limit: 5}) - Inject dans le system prompt : "Pannes précédentes de cet équipement : {past_failures}.
Si le pattern actuel matche une panne passée, cite-la explicitement dans
similar_past_failure." - Le tool
submit_rca(cf. M4.3) accepte déjàsimilar_past_failure→ rien à changer côté schema
Scène démo dédiée. POST /api/v1/demo/trigger-memory-scene :
- INSERT une fausse
failure_historydatée de 3 mois (pattern P-02 similaire) - Trigger une anomalie P-02 actuelle
- L'Investigator match → RCA cite la panne passée → scène flex "l'agent apprend"
Priorité. P1. Skip si serré J6 PM. Si skipé, l'essentiel reste : Investigator
utilise failure_history en contexte (sans la scène scénarisée).
Acceptance.
- Investigator cite la panne passée dans son RCA si pattern matche
- Endpoint demo
/trigger-memory-scenerejouable autant de fois que voulu
- Scène 2 (anomalie live) et Scène 3 (RCA) de la démo
- M5 (le Work Order Generator est triggered par l'Investigator)
- M1 (
work_ordercolonnes) - M2 (tools, MCPClient, M2.9 UI tools)
- M3 (KB doit exister pour que Sentinel ait des seuils)