Files
IdeA/crates/application/src/orchestrator/service.rs
Blomios 11dc8d7ba3 feat(skill-awareness): T1→T5 — manifeste de skills + outil MCP idea_skill_read
Surface les skills assignés à un agent « à la MCP » : description sur l'entité
Skill (+ effective_description fallback), section « # Skills disponibles » haute
altitude dans le convention-file (mode MCP), outil read-only idea_skill_read
(résolution projet→global), use case ReadSkill (port SkillStore existant), et
câblage au composition root. Dump legacy du corps complet conservé en mode
sans-MCP (zéro régression). Rétro-compat index.json (serde default).

Tests : domain+application+infrastructure 1212 passed / 0 failed (23 ajoutés).
Reste : T6 (champ description front) + T7 (e2e après rebuild AppImage).

NB topologie : commit réalisé par l'orchestrateur car l'agent Git était
injoignable (bug de livraison cold-start) ; à faire relire/rebaser par Git.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-21 18:30:47 +02:00

2359 lines
107 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

//! [`OrchestratorService`] — dispatches a validated [`OrchestratorCommand`] to the
//! existing agent/terminal use cases (ARCHITECTURE §14.3).
//!
//! The orchestrator agent never spawns a process itself: IdeA is the single source
//! of truth for the agent lifecycle. This service is the application-layer seam
//! that turns a request into the *same* calls the UI makes:
//!
//! - `spawn_agent` → [`CreateAgentFromScratch`] (if unknown) then [`LaunchAgent`],
//! - `stop_agent` → resolve the agent's live session, then [`CloseTerminal`],
//! - `update_agent_context` → [`UpdateAgentContext`].
//!
//! It talks **only** to use cases and ports ([`ProfileStore`], [`TerminalSessions`]):
//! no filesystem watching, no JSON, no process spawning here — those are the
//! infrastructure adapter's job. That keeps this fully unit-testable with fakes.
use std::collections::HashMap;
use std::sync::{Arc, Mutex as StdMutex};
use std::time::{Duration, Instant};
use tokio::sync::Mutex as AsyncMutex;
use domain::conversation::{ConversationParty, ConversationRegistry, SessionRef, WaitForGraph};
use domain::conversation_log::{ConversationTurn, TurnId, TurnRole};
use domain::input::{InputMediator, InputSource, SubmitConfig};
use domain::mailbox::{Ticket, TicketId, TurnResolution};
use domain::ports::{Clock, EventBus, ProfileStore, PtyHandle};
use domain::project::ProjectPath;
use domain::{
AgentId, AgentProfile, DomainEvent, OrchestratorCommand, OrchestratorVisibility, ProfileId,
Project,
};
use crate::conversation::RecordTurn;
use crate::memory::HarvestMemoryFromTurn;
use crate::agent::{
drain_with_readiness, CreateAgentFromScratch, CreateAgentInput, LaunchAgent, LaunchAgentInput,
ListAgents, ListAgentsInput, McpRuntime, ReattachDecision, UpdateAgentContext,
UpdateAgentContextInput, CODEX_SUBMIT_DELAY_MS,
};
use crate::error::AppError;
use crate::orchestrator::{
ProposeContext, ProposeContextInput, ProposeOutcome, ReadContext, ReadContextInput, ReadMemory,
ReadMemoryInput, WriteMemory, WriteMemoryInput,
};
use crate::skill::{CreateSkill, CreateSkillInput, ReadSkill, ReadSkillInput};
use crate::terminal::{CloseTerminal, CloseTerminalInput, StructuredSessions, TerminalSessions};
/// Default terminal geometry for an orchestrator-launched agent cell. The UI
/// resizes the PTY to the real cell size on attach; these are sane starting rows
/// /cols so the spawn never fails on a zero-sized terminal.
const DEFAULT_ROWS: u16 = 24;
/// See [`DEFAULT_ROWS`].
const DEFAULT_COLS: u16 = 80;
/// Submit defaults for delegated prompts, after applying profile-specific
/// compatibility fallbacks for existing saved profiles.
fn submit_config_for_profile(profile: &AgentProfile) -> SubmitConfig {
let delay_ms = profile.submit_delay_ms.or_else(|| {
matches!(
profile.structured_adapter,
Some(domain::profile::StructuredAdapter::Codex)
)
.then_some(CODEX_SUBMIT_DELAY_MS)
});
SubmitConfig::new(profile.submit_sequence.clone(), delay_ms)
}
/// Bound on the synchronous inter-agent rendezvous (`agent.message` → `AskAgent`).
///
/// A target agent's turn can be long (reasoning + tool use), so the cap is
/// generous; on expiry [`send_blocking`] returns [`domain::ports::AgentSessionError::Timeout`]
/// **without killing the session**, so the requester can retry. Internal and
/// intentionally not yet config-exposed (it may become a per-project setting
/// without changing the contract).
///
/// TEMPORAIRE (demande utilisateur 2026-06-13) : la borne de 300 s coupait des tours
/// délégués légitimement très longs (gros refactors + tests). On la relève donc à 24 h
/// (≈ « pas de limite ») le temps de concevoir un vrai signal de vivacité (savoir si
/// l'agent travaille encore) qui remplacera ce timeout brut. NE PAS considérer comme
/// définitif : à reconvertir en borne courte + heartbeat once that liveness signal exists.
const ASK_AGENT_TIMEOUT: Duration = Duration::from_secs(24 * 60 * 60);
/// Borne d'attente **en file** pour acquérir le verrou de tour d'un agent (A0,
/// cadrage v5 §4).
///
/// La sérialisation FIFO par agent (« 1 agent = 1 employé : un tour à la fois »)
/// fait qu'un `ask` concurrent vers la **même** cible patiente derrière le tour en
/// cours. Le timeout de tour ([`ASK_AGENT_TIMEOUT`]) borne **le tour lui-même**,
/// **pas** cette attente : on lui donne donc son propre plafond, généreux mais fini,
/// pour qu'un `ask` ne reste jamais bloqué indéfiniment si la file est longue
/// (inanition). À l'expiration, l'`ask` renvoie un timeout typé (réutilise le
/// **même** type que le timeout de tour, cf. [`AppError::Process`] via
/// [`domain::ports::AgentSessionError::Timeout`]) ; le tour en cours n'est pas
/// affecté. Largeur = un tour complet + sa propre file ⇒ on autorise deux tours
/// pleins d'attente.
///
/// TEMPORAIRE (cf. [`ASK_AGENT_TIMEOUT`]) : relevé à 48 h (≈ deux tours « sans limite »)
/// en cohérence avec la levée temporaire de la borne de tour, le temps d'un vrai
/// signal de vivacité.
const ASK_QUEUE_WAIT_CAP: Duration = Duration::from_secs(48 * 60 * 60);
/// Borne de tour effective (lot 2, timeouts pilotés par profil) : le
/// `turn_timeout_ms` du profil de la cible **quand il existe**, sinon le défaut
/// historique [`ASK_AGENT_TIMEOUT`] (zéro régression pour un profil sans `liveness` /
/// sans `turn_timeout_ms`). Pure et testable sans wiring de service.
#[must_use]
fn resolve_turn_timeout(turn_timeout_ms: Option<u32>) -> Duration {
turn_timeout_ms.map_or(ASK_AGENT_TIMEOUT, |ms| Duration::from_millis(u64::from(ms)))
}
/// Garde RAII de **fin de tour** : ramène la cible `Idle` quoi qu'il arrive (succès,
/// erreur, timeout, **et surtout future abandonné/dropped**) tant qu'elle n'a pas été
/// désarmée.
///
/// Cause racine du blocage `Busy` à vie : l'agent passe `Idle→Busy` dès l'`enqueue`
/// (qui démarre le tour) mais ne revenait `Idle` que sur la **branche succès** du
/// `select!`/`match`. Si le futur `ask_agent` est **dropped** (le demandeur a interrompu
/// son appel), AUCUNE branche ne s'exécute ⇒ l'agent reste `Busy`, et toutes les
/// délégations suivantes s'empilent derrière ce tour fantôme sans jamais être livrées.
///
/// Le garde tient les `Arc` nécessaires pour, à son `Drop`, faire à la fois
/// `cancel_head(agent, ticket)` (retire le ticket fantôme de la FIFO — idempotent et
/// positionnel : no-op si le head a déjà changé) **et** `mark_idle(agent)` (libère la
/// FIFO — idempotent). Sur **succès**, on appelle [`BusyTurnGuard::disarm`] AVANT de
/// retourner : le `mark_idle` « propre » déjà présent reste en place et on évite tout
/// double `cancel_head` d'un ticket déjà résolu.
struct BusyTurnGuard {
input: Arc<dyn InputMediator>,
mailbox: Arc<dyn domain::mailbox::AgentMailbox>,
agent: AgentId,
ticket: TicketId,
armed: bool,
}
impl BusyTurnGuard {
/// Arme le garde juste après l'`enqueue` (qui a démarré le tour ⇒ cible `Busy`).
fn new(
input: Arc<dyn InputMediator>,
mailbox: Arc<dyn domain::mailbox::AgentMailbox>,
agent: AgentId,
ticket: TicketId,
) -> Self {
Self {
input,
mailbox,
agent,
ticket,
armed: true,
}
}
/// Désarme le garde (branche succès) : le `Drop` devient un no-op. À appeler AVANT
/// de retourner le contenu, pour préserver le `mark_idle` propre déjà fait et ne pas
/// re-`cancel_head` un ticket déjà résolu.
fn disarm(mut self) {
self.armed = false;
}
}
impl Drop for BusyTurnGuard {
fn drop(&mut self) {
if self.armed {
// Ordre : retirer le ticket fantôme de la FIFO PUIS libérer le busy state.
// `cancel_head` est positionnel/idempotent (no-op si le head n'est plus ce
// ticket) et `mark_idle` est idempotent (no-op si déjà Idle) ⇒ sûr quel que
// soit le chemin (erreur, timeout, drop).
self.mailbox.cancel_head(self.agent, self.ticket);
self.input.mark_idle(self.agent);
// Rendezvous beacon (diagnostics) : le garde a libéré une cible restée Busy
// (chemin erreur / timeout / futur ask abandonné). Si ce beacon apparaît
// sans « ask resolved », la cible n'a jamais répondu — c'est le scénario de
// blocage à diagnostiquer.
crate::diag!(
"[rendezvous] busy-guard freed target agent {} (ticket {})",
self.agent,
self.ticket,
);
}
}
}
/// Fournit les faits OS/runtime (exe + endpoint projet) pour écrire la déclaration MCP
/// réelle quand l'orchestrateur (re)lance une cible sur le chemin `ask`. Implémenté dans
/// app-tauri (seul détenteur de current_exe/$APPIMAGE/mcp_endpoint).
///
/// La couche `application` ne connaît que ce **port** : elle ne calcule jamais le chemin
/// de l'exécutable ni l'endpoint loopback (ces faits vivent dans `app-tauri`, cadrage v5
/// §0.3 / §7). Seules les **chaînes** d'un [`McpRuntime`] traversent la frontière.
pub trait McpRuntimeProvider: Send + Sync {
/// `agent_id` = la cible relancée = le `--requester` (c'est elle qui appellera idea_reply).
/// `None` ⇒ dégrade vers la déclaration minimale (jamais d'échec de lancement).
fn runtime_for(&self, project: &Project, agent_id: AgentId) -> Option<McpRuntime>;
}
/// Fournit le use case [`RecordTurn`] **lié au project root** de la délégation en
/// cours (lot P6b).
///
/// L'[`OrchestratorService`] est **unique et partagé par tous les projets ouverts**
/// (un seul `Arc` au composition root), alors que le log/handoff conversationnel est
/// **par project root** (`<root>/.ideai/conversations/`, comme la mémoire). Les
/// adapters `Fs*` de P6a fixent leur racine à la construction et leur port ne porte
/// pas le root par appel ; on ne peut donc pas figer un `RecordTurn` global. Ce
/// **port** lève la tension : `ask_agent` connaît le `project.root` du tour et demande
/// au provider de matérialiser un `RecordTurn` ciblant le **bon** dossier (les adapters
/// `Fs*` ne font que des jointures de chemin, leur construction est triviale).
///
/// `None` ⇒ aucune persistance (best-effort absente) : zéro régression pour les call
/// sites/tests qui ne le branchent pas. Implémenté dans app-tauri (seul détenteur des
/// adapters `Fs*`).
pub trait RecordTurnProvider: Send + Sync {
/// Construit le [`RecordTurn`] dont le log/handoff ciblent `root`. Appelé une fois
/// par checkpoint best-effort ; `None` ⇒ on saute silencieusement la persistance.
fn record_turn_for(&self, root: &ProjectPath) -> Option<Arc<RecordTurn>>;
}
/// Dispatches validated orchestrator commands to the agent/terminal use cases.
pub struct OrchestratorService {
create_agent: Arc<CreateAgentFromScratch>,
launch_agent: Arc<LaunchAgent>,
list_agents: Arc<ListAgents>,
close_terminal: Arc<CloseTerminal>,
update_context: Arc<UpdateAgentContext>,
create_skill: Arc<CreateSkill>,
profiles: Arc<dyn ProfileStore>,
sessions: Arc<TerminalSessions>,
/// Médiateur d'entrée (cadrage C3 §5.2) — point de convergence unique de l'entrée
/// d'un agent. La cible d'un `agent.message`/`idea_ask_agent` y reçoit un ticket
/// (`enqueue`) dont on **attend** la résolution (`idea_reply` ⇒
/// [`OrchestratorCommand::Reply`]) ; son impl écrit aussi le tour dans le PTY de la
/// cible (livraison sérialisée, plus d'écriture ad hoc ici). Injecté via
/// [`Self::with_input_mediator`] ; `None` ⇒ `AskAgent`/`Reply` non servis (call
/// sites/tests legacy restent verts).
input: Option<Arc<dyn InputMediator>>,
/// Mailbox sous-jacent du médiateur, pour `resolve`/`resolve_ticket`/`cancel_head`
/// (corrélation par ticket). C'est le **même** moteur de corrélation que celui que
/// `input` enveloppe ; injecté ensemble via [`Self::with_input_mediator`].
mailbox: Option<Arc<dyn domain::mailbox::AgentMailbox>>,
/// Registre des conversations par paire (cadrage C3 §5.2) — résout paresseusement
/// le fil `A↔B` (ou `User↔B`) d'un `ask`, sépare strictement les contextes.
/// Injecté via [`Self::with_conversations`] ; `None` ⇒ on retombe sur un routage
/// par agent sans matérialisation de fil (legacy).
conversations: Option<Arc<dyn ConversationRegistry>>,
/// Graphe d'attente inter-agents (cadrage C3 §6) — arête posée à l'`enqueue` d'un
/// `ask` A→B, retirée au reply/timeout (RAII via le garde de tour). Sert à
/// **refuser** une délégation ré-entrante (A→B→…→A) avant deadlock.
wait_for: StdMutex<WaitForGraph>,
/// Bus d'événements pour publier [`DomainEvent::AgentReplied`] à l'issue d'un
/// `ask` réussi (§17.4). Injecté via [`Self::with_events`] ; `None` ⇒ pas de
/// publication (l'`ask` fonctionne quand même).
events: Option<Arc<dyn EventBus>>,
/// Verrous de **tour par agent** (A0, cadrage v5 §4) — sérialisation FIFO des
/// `ask` vers une **même** cible : « 1 agent = 1 employé, un tour à la fois ».
///
/// Clé = `AgentId` de la **cible** ; valeur = un `tokio::Mutex` d'unité dont le
/// garde est tenu **à travers** le `.await` de `send_blocking` (d'où le mutex
/// **async**, pas `std`). La `HashMap` elle-même n'est tenue que le temps du
/// get-or-create (jamais à travers un `.await`), donc protégée par un mutex
/// **synchrone** `std` — pas de garde gardé en travers d'un point de suspension.
///
/// Le registre vit **ici** (règle applicative d'orchestration, frontière
/// hexagonale, cadrage v5 §5) et **non** dans [`StructuredSessions`] : sérialiser
/// les tours est une responsabilité de l'orchestrateur, pas du stockage de
/// sessions (SRP) ; le couplage reste minimal. Verrou **par agent** ⇒ deux `ask`
/// vers des agents **différents** n'entrent jamais en contention (un map d'`Arc`
/// distincts).
///
/// Croissance bornée en pratique au nombre d'agents du projet ; une entrée
/// morte ne coûte qu'un `Arc<Mutex<()>>` vide (pas de session, pas de process).
ask_locks: StdMutex<HashMap<AgentId, Arc<AsyncMutex<()>>>>,
/// Fournisseur des faits OS/runtime (exe + endpoint) pour écrire la déclaration
/// MCP **réelle** quand `ensure_live_pty` (re)lance une cible sur le chemin `ask`
/// (B-3). Injecté au câblage via [`Self::with_mcp_runtime_provider`] depuis
/// app-tauri ; `None` ⇒ on conserve la déclaration minimale (`mcp_runtime: None`),
/// donc zéro régression pour les call sites/tests qui ne le branchent pas.
mcp_runtime_provider: Option<Arc<dyn McpRuntimeProvider>>,
/// FileGuard-mediated context/memory use cases (cadrage C7). Injected via
/// [`Self::with_context_guard`] ; `None` ⇒ les commandes `context.*`/`memory.*`
/// renvoient une erreur typée (call sites/tests legacy restent verts).
context_guard: Option<Arc<ContextGuardUseCases>>,
/// Provider du use case [`RecordTurn`] **par project root** (lot P6b), pour
/// persister **best-effort** le Prompt et la Response de chaque paire déléguée dans
/// le bon `conversationId`. Injecté avec son horloge via [`Self::with_record_turn`] ;
/// `None` ⇒ aucune persistance (zéro régression pour les call sites/tests legacy).
/// Un échec de persistance ne transforme **jamais** un succès de délégation en erreur.
record_turn: Option<Arc<dyn RecordTurnProvider>>,
/// Horloge millis (port [`Clock`]) pour estampiller `at_ms` des tours persistés —
/// injectée avec [`Self::record_turn`]. La couche `application` reste **pure** : pas
/// de `SystemTime::now()` brut ici, le temps vient du port injecté au composition root.
clock: Option<Arc<dyn Clock>>,
/// Registre des **sessions IA structurées** vivantes (§17.5), pour router un `ask`
/// vers une cible à `structured_adapter` directement sur sa session
/// ([`drain_with_readiness`](crate::agent::drain_with_readiness)) — le `Final`
/// déterministe réveille le `pending` **et** marque l'agent `Idle` (chantier
/// readiness/heartbeat, lot 1, fix du blocage `Busy`). Injecté via
/// [`Self::with_structured`] ; `None` ⇒ on conserve **uniquement** le chemin
/// PTY/mailbox legacy (zéro régression pour les call sites/tests qui ne le branchent
/// pas).
structured: Option<Arc<StructuredSessions>>,
/// Harvester d'auto-mémoire (Lot E1), branché **après** le checkpoint Response
/// best-effort d'un `ask` réussi : parse les blocs ` ```idea-memory ` de la
/// réponse et persiste les notes valides via le `MemoryStore`. Injecté via
/// [`Self::with_memory_harvest`] ; `None` ⇒ aucun harvest (zéro régression). Un
/// échec de harvest ne transforme **jamais** un succès de délégation en erreur.
memory_harvest: Option<Arc<HarvestMemoryFromTurn>>,
/// Use case de lecture d'un skill **par nom** (`idea_skill_read`, feature
/// « skills à la MCP »). Compose le port `SkillStore` existant (aucun nouveau
/// port). Injecté via [`Self::with_read_skill`] ; `None` ⇒ la commande
/// `skill.read` renvoie une erreur typée (call sites/tests legacy restent verts).
read_skill: Option<Arc<ReadSkill>>,
}
/// Bundle des quatre use cases C7 sous [`domain::fileguard::FileGuard`], injectés
/// ensemble dans l'[`OrchestratorService`] (cadrage C7). Regroupés pour garder la
/// signature de [`OrchestratorService::with_context_guard`] simple (un seul `Arc`).
pub struct ContextGuardUseCases {
/// Lecture d'un contexte `.md` IdeA sous read-lease.
pub read_context: Arc<ReadContext>,
/// Proposition/écriture d'un contexte `.md` IdeA sous le garde.
pub propose_context: Arc<ProposeContext>,
/// Lecture mémoire sous read-lease.
pub read_memory: Arc<ReadMemory>,
/// Écriture mémoire sous write-lease.
pub write_memory: Arc<WriteMemory>,
}
/// Outcome of dispatching a command — a short, human-readable success summary the
/// infrastructure adapter folds into the JSON response file.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OrchestratorOutcome {
/// One-line description of what IdeA did (e.g. `"launched agent dev-backend"`).
pub detail: String,
/// The target agent's **reply content** for a synchronous `agent.message`
/// (`AskAgent`, §17.4). `Some(content)` carries the turn's `Final`; every other
/// command leaves it `None`. Additive field ⇒ existing call sites/tests are
/// untouched (they only read [`Self::detail`]).
pub reply: Option<String>,
}
impl OrchestratorService {
/// Builds the service from the use cases and ports it dispatches to.
#[must_use]
#[allow(clippy::too_many_arguments)]
pub fn new(
create_agent: Arc<CreateAgentFromScratch>,
launch_agent: Arc<LaunchAgent>,
list_agents: Arc<ListAgents>,
close_terminal: Arc<CloseTerminal>,
update_context: Arc<UpdateAgentContext>,
create_skill: Arc<CreateSkill>,
profiles: Arc<dyn ProfileStore>,
sessions: Arc<TerminalSessions>,
) -> Self {
Self {
create_agent,
launch_agent,
list_agents,
close_terminal,
update_context,
create_skill,
profiles,
sessions,
input: None,
mailbox: None,
conversations: None,
wait_for: StdMutex::new(WaitForGraph::new()),
events: None,
ask_locks: StdMutex::new(HashMap::new()),
mcp_runtime_provider: None,
context_guard: None,
record_turn: None,
clock: None,
structured: None,
memory_harvest: None,
read_skill: None,
}
}
/// Branche le use case [`ReadSkill`] (`skill.read`/`idea_skill_read`). Builder
/// additif (signature de [`Self::new`] inchangée) ; `None` ⇒ la commande
/// renvoie une erreur typée « not configured ».
#[must_use]
pub fn with_read_skill(mut self, read_skill: Arc<ReadSkill>) -> Self {
self.read_skill = Some(read_skill);
self
}
/// Branche le registre des [`StructuredSessions`] (§17.5) pour router un `ask`
/// vers une cible **structurée** sur sa propre session. Builder additif (signature
/// de [`Self::new`] inchangée) ; `None` ⇒ chemin PTY/mailbox legacy uniquement.
#[must_use]
pub fn with_structured(mut self, structured: Arc<StructuredSessions>) -> Self {
self.structured = Some(structured);
self
}
/// Branche les use cases C7 (`context.*`/`memory.*`) sous
/// [`domain::fileguard::FileGuard`]. Builder additif (signature de [`Self::new`]
/// inchangée).
#[must_use]
pub fn with_context_guard(mut self, guard: Arc<ContextGuardUseCases>) -> Self {
self.context_guard = Some(guard);
self
}
/// Returns the per-agent **turn lock**, creating it on first use.
///
/// Get-or-create under the synchronous map mutex (held only for this lookup,
/// never across an `.await`); the returned `Arc<AsyncMutex<()>>` is the lock the
/// caller acquires (and holds across `send_blocking`) to serialise turns for
/// `agent_id`. Two different agents get two distinct `Arc`s ⇒ no contention.
fn ask_lock_for(&self, agent_id: &AgentId) -> Arc<AsyncMutex<()>> {
let mut locks = self
.ask_locks
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
Arc::clone(locks.entry(*agent_id).or_default())
}
/// Branche le **médiateur d'entrée** (cadrage C3 §5.2) pour servir
/// `agent.message`/[`OrchestratorCommand::AskAgent`] et
/// `agent.reply`/[`OrchestratorCommand::Reply`]. Le `mailbox` est le moteur de
/// corrélation **sous-jacent** au médiateur (le même `InMemoryMailbox` que
/// `MediatedInbox` enveloppe) : on l'injecte ensemble pour pouvoir `resolve`/
/// `resolve_ticket`/`cancel_head` un ticket. Builder additif : signature de
/// [`Self::new`] **inchangée** (les tests/call sites legacy restent verts).
#[must_use]
pub fn with_input_mediator(
mut self,
input: Arc<dyn InputMediator>,
mailbox: Arc<dyn domain::mailbox::AgentMailbox>,
) -> Self {
self.input = Some(input);
self.mailbox = Some(mailbox);
self
}
/// Branche le [`ConversationRegistry`] (cadrage C3 §5.2) pour résoudre
/// paresseusement le fil `A↔B` (ou `User↔B`) d'un `ask` et séparer les contextes.
/// Builder additif (signature de [`Self::new`] inchangée).
#[must_use]
pub fn with_conversations(mut self, conversations: Arc<dyn ConversationRegistry>) -> Self {
self.conversations = Some(conversations);
self
}
/// Branche l'[`EventBus`] pour publier [`DomainEvent::AgentReplied`] après un
/// `ask` réussi (§17.4). Builder additif (cf. [`Self::with_structured`]).
#[must_use]
pub fn with_events(mut self, events: Arc<dyn EventBus>) -> Self {
self.events = Some(events);
self
}
/// Branche le [`McpRuntimeProvider`] (app-tauri) pour que les (re)lancements
/// issus du chemin `ask` (`ensure_live_pty`) écrivent la déclaration MCP **réelle**
/// (endpoint + exe + requester) au lieu de la minimale — sans quoi le pont MCP
/// n'est jamais spawné et la cible ne peut pas appeler `idea_reply` (timeout).
/// Builder additif : signature de [`Self::new`] **inchangée**.
#[must_use]
pub fn with_mcp_runtime_provider(mut self, provider: Arc<dyn McpRuntimeProvider>) -> Self {
self.mcp_runtime_provider = Some(provider);
self
}
/// Branche le provider de [`RecordTurn`] **par project root** + son horloge (lot
/// P6b) pour persister **best-effort** le Prompt et la Response de chaque paire
/// déléguée. Builder additif : signature de [`Self::new`] **inchangée** (les
/// tests/call sites legacy restent verts ; `None` ⇒ aucune persistance, donc aucune
/// régression). Un échec de persistance ne casse jamais la délégation live.
#[must_use]
pub fn with_record_turn(
mut self,
record_turn: Arc<dyn RecordTurnProvider>,
clock: Arc<dyn Clock>,
) -> Self {
self.record_turn = Some(record_turn);
self.clock = Some(clock);
self
}
/// Branche le harvester d'auto-mémoire (Lot E1), appelé **après** le checkpoint
/// `RecordTurn` de la réponse d'un `ask` réussi. `None` (défaut) ⇒ aucun harvest.
/// Le `MemoryStore` prend le `root` par appel, donc un seul use case sert tous
/// les projets (pas de provider par root nécessaire, contrairement au handoff).
#[must_use]
pub fn with_memory_harvest(mut self, harvester: Arc<HarvestMemoryFromTurn>) -> Self {
self.memory_harvest = Some(harvester);
self
}
/// Persiste **best-effort** un tour (`Prompt`/`Response`) dans `conversation`.
///
/// No-op silencieux quand le provider/horloge ne sont pas câblés, quand le provider
/// ne rend pas de [`RecordTurn`] pour ce root, ou quand l'`append`/`save` échoue : la
/// persistance ne doit **jamais** transformer un succès de délégation en erreur, ni
/// paniquer (contrat P6b). N'ajoute pas de latence inutile (un seul `await` borné par
/// les adapters `Fs*`, déjà sérialisés par le verrou de tour de la cible).
async fn record_turn_best_effort(
&self,
root: &ProjectPath,
conversation: domain::conversation::ConversationId,
source: InputSource,
role: TurnRole,
text: String,
) {
let (Some(provider), Some(clock)) = (&self.record_turn, &self.clock) else {
return;
};
let Some(record) = provider.record_turn_for(root) else {
return;
};
let at_ms = u64::try_from(clock.now_millis()).unwrap_or(0);
let turn = ConversationTurn::new(
TurnId::new_random(),
conversation,
at_ms,
source,
role,
text,
);
// Best-effort : un échec de persistance est avalé (le contrat P6b interdit qu'il
// remonte). Pas de framework de log dans `application` ; on reste cohérent avec
// les autres effets best-effort du service (cf. publication `AgentReplied`).
let _ = record.record(conversation, turn).await;
}
/// Récolte **best-effort** les notes d'auto-mémoire (Lot E1) d'une réponse.
///
/// No-op silencieux quand le harvester n'est pas câblé. Ne parse **que** le texte
/// du tour courant (jamais le log complet) et n'agit que sur un `Response`. Tout
/// échec parse/store reste dans l'`outcome` (jamais propagé) : un succès de
/// délégation n'est **jamais** transformé en erreur. Branché **après** le
/// checkpoint `RecordTurn` de la réponse (append/handoff déjà effectués).
async fn harvest_memory_best_effort(&self, root: &ProjectPath, text: &str) {
let Some(harvester) = &self.memory_harvest else {
return;
};
let outcome = harvester.harvest(root, TurnRole::Response, text).await;
if !outcome.saved.is_empty() || !outcome.errors.is_empty() {
crate::diag!(
"[harvest] auto-memory: found={} saved={} skipped={} errors={}",
outcome.found,
outcome.saved.len(),
outcome.skipped,
outcome.errors.len(),
);
}
}
/// Dispatches a validated command against `project`.
///
/// # Errors
/// Propagates the underlying use-case [`AppError`] (e.g. unknown profile,
/// unknown agent, PTY failure). For `spawn_agent` a *known* agent is launched
/// directly; an *unknown* one is created from scratch first.
pub async fn dispatch(
&self,
project: &Project,
command: OrchestratorCommand,
) -> Result<OrchestratorOutcome, AppError> {
match command {
OrchestratorCommand::SpawnAgent {
name,
profile,
context,
visibility,
} => {
self.spawn_agent(project, name, profile, context, visibility)
.await
}
OrchestratorCommand::AskAgent {
target,
task,
requester,
} => self.ask_agent(project, target, task, requester).await,
OrchestratorCommand::Reply {
from,
ticket,
result,
} => self.reply(from, ticket, result),
OrchestratorCommand::ListAgents => self.list_agents(project).await,
OrchestratorCommand::StopAgent { name } => self.stop_agent(project, name).await,
OrchestratorCommand::UpdateAgentContext { name, context } => {
self.update_agent_context(project, name, context).await
}
OrchestratorCommand::CreateSkill {
name,
content,
scope,
} => self.create_skill(project, name, content, scope).await,
OrchestratorCommand::ReadContext { target, requester } => {
self.read_context(project, target, requester).await
}
OrchestratorCommand::ProposeContext {
target,
content,
requester,
} => {
self.propose_context(project, target, content, requester)
.await
}
OrchestratorCommand::ReadMemory { slug, requester } => {
self.read_memory(project, slug, requester).await
}
OrchestratorCommand::ReadSkill { name, requester } => {
self.read_skill(project, name, requester).await
}
OrchestratorCommand::WriteMemory {
slug,
content,
requester,
} => self.write_memory(project, slug, content, requester).await,
}
}
/// Returns the injected C7 use cases, or a typed error when unwired.
fn require_context_guard(&self) -> Result<&ContextGuardUseCases, AppError> {
self.context_guard.as_deref().ok_or_else(|| {
AppError::Invalid("FileGuard context/memory tools are not configured".to_owned())
})
}
/// `context.read` → reads an IdeA-owned context under a shared read-lease; the
/// body is returned inline in the outcome's `reply`.
async fn read_context(
&self,
project: &Project,
target: Option<String>,
requester: ConversationParty,
) -> Result<OrchestratorOutcome, AppError> {
let md = self
.require_context_guard()?
.read_context
.execute(ReadContextInput {
project: project.clone(),
target: target.clone(),
requester,
})
.await?;
Ok(OrchestratorOutcome {
detail: format!("read {} context", target.as_deref().unwrap_or("project")),
reply: Some(md.into_string()),
})
}
/// `context.propose` → direct write (agent ctx / orchestrator on global) or a
/// materialised proposal (non-orchestrator on global).
async fn propose_context(
&self,
project: &Project,
target: Option<String>,
content: String,
requester: ConversationParty,
) -> Result<OrchestratorOutcome, AppError> {
let outcome = self
.require_context_guard()?
.propose_context
.execute(ProposeContextInput {
project: project.clone(),
target: target.clone(),
content,
requester,
})
.await?;
let detail = match outcome {
ProposeOutcome::Written => {
format!("wrote {} context", target.as_deref().unwrap_or("project"))
}
ProposeOutcome::Proposed { path } => {
format!("filed proposal for project context at {path}")
}
};
Ok(OrchestratorOutcome {
detail,
reply: None,
})
}
/// `memory.read` → reads a note (or the index) under a shared read-lease; the
/// content is returned inline in the outcome's `reply`.
async fn read_memory(
&self,
project: &Project,
slug: Option<String>,
requester: ConversationParty,
) -> Result<OrchestratorOutcome, AppError> {
let content = self
.require_context_guard()?
.read_memory
.execute(ReadMemoryInput {
project: project.clone(),
slug: slug.clone(),
requester,
})
.await?;
Ok(OrchestratorOutcome {
detail: format!("read memory {}", slug.as_deref().unwrap_or("index")),
reply: Some(content),
})
}
/// `skill.read` → resolves a skill by name (project scope first, then global)
/// and returns its Markdown body inline in the outcome's `reply`. Read-only:
/// `requester` is accepted for symmetry with the other read tools but skill
/// reads take no [`domain::fileguard::FileGuard`] lease.
async fn read_skill(
&self,
project: &Project,
name: String,
_requester: ConversationParty,
) -> Result<OrchestratorOutcome, AppError> {
let read_skill = self.read_skill.as_deref().ok_or_else(|| {
AppError::Invalid("the idea_skill_read tool is not configured".to_owned())
})?;
let md = read_skill
.execute(ReadSkillInput {
name: name.clone(),
project_root: project.root.clone(),
})
.await?;
Ok(OrchestratorOutcome {
detail: format!("read skill {name}"),
reply: Some(md.into_string()),
})
}
/// `memory.write` → writes a note under an exclusive write-lease.
async fn write_memory(
&self,
project: &Project,
slug: String,
content: String,
requester: ConversationParty,
) -> Result<OrchestratorOutcome, AppError> {
self.require_context_guard()?
.write_memory
.execute(WriteMemoryInput {
project: project.clone(),
slug: slug.clone(),
content,
requester,
})
.await?;
Ok(OrchestratorOutcome {
detail: format!("wrote memory {slug}"),
reply: None,
})
}
/// `spawn_agent`: create the agent if the manifest doesn't already hold one by
/// that name, then launch it (which publishes `AgentLaunched` → the UI opens a
/// cell + the Agents tab).
async fn spawn_agent(
&self,
project: &Project,
name: String,
profile: Option<String>,
context: Option<String>,
visibility: OrchestratorVisibility,
) -> Result<OrchestratorOutcome, AppError> {
let existing = self.find_agent_id_by_name(project, &name).await?;
let agent_id = match existing {
Some(id) => id,
None => {
let profile = profile.as_deref().ok_or_else(|| {
AppError::Invalid("profile is required to create an agent".to_owned())
})?;
let profile_id = self.resolve_profile(profile).await?;
let created = self
.create_agent
.execute(CreateAgentInput {
project: project.clone(),
name: name.clone(),
profile_id,
initial_content: context,
})
.await?;
created.agent.id
}
};
if let Some(session_id) = self.sessions.session_for_agent(&agent_id) {
match visibility {
OrchestratorVisibility::Background => {
return Ok(OrchestratorOutcome {
detail: format!("agent {name} already running in background"),
reply: None,
});
}
OrchestratorVisibility::Visible { node_id } => {
// R0a — même règle que le garde de `LaunchAgent` (cadrage v5 §3.2,
// Trou A) : un `spawn` est un **lancement neuf** (pas de
// conversation portée), donc ré-attacher à la **même** cellule
// hôte est légitime (rebind de vue), mais viser un **autre** node
// pour un agent singleton déjà vivant est un second lancement ⇒
// refus `AgentAlreadyRunning`.
let host_node = self.sessions.node_for_agent(&agent_id);
match ReattachDecision::resolve(Some(node_id), host_node, None) {
ReattachDecision::Rebind { node_id } => {
let session = self
.sessions
.rebind_agent_node(&agent_id, node_id)
.ok_or_else(|| {
AppError::NotFound(format!(
"running session {session_id} for agent {name}"
))
})?;
return Ok(OrchestratorOutcome {
detail: format!(
"attached agent {name} to cell {}",
session.node_id
),
reply: None,
});
}
ReattachDecision::Refuse { node_id } => {
return Err(AppError::AgentAlreadyRunning { agent_id, node_id });
}
// A `Visible` spawn always carries a node, so the decision is
// never `Idempotent`; treat it as a same-node rebind for safety.
ReattachDecision::Idempotent => {
return Ok(OrchestratorOutcome {
detail: format!("agent {name} already running"),
reply: None,
});
}
}
}
}
}
let node_id = match visibility {
OrchestratorVisibility::Background => None,
OrchestratorVisibility::Visible { node_id } => Some(node_id),
};
self.launch_agent
.execute(LaunchAgentInput {
project: project.clone(),
agent_id,
rows: DEFAULT_ROWS,
cols: DEFAULT_COLS,
node_id,
conversation_id: None,
// Orchestrator-driven launch (inside `application`): no OS/runtime
// facts to inject here; the real MCP declaration is written when the
// agent is (re)launched through the app-tauri composition root.
mcp_runtime: None,
})
.await?;
Ok(OrchestratorOutcome {
detail: match visibility {
OrchestratorVisibility::Background => {
format!("launched agent {name} in background")
}
OrchestratorVisibility::Visible { node_id } => {
format!("launched agent {name} in cell {node_id}")
}
},
reply: None,
})
}
/// `agent.message` / `idea_ask_agent`: the **inter-agent delegation rendezvous**
/// (Option 1 « Terminal + MCP », lot B-3).
///
/// The target's human-facing view is now a **raw native terminal** (PTY REPL), and
/// delegation flows through the terminal's single FIFO input plus the MCP mailbox:
///
/// 1. Resolve the target by name and acquire its **per-agent turn lock** so two
/// `ask`s for the same target serialise FIFO (1 agent = 1 employee).
/// 2. Ensure the target is **live in the PTY registry** — reusing its terminal if
/// it is already running, otherwise launching it in the background (a normal
/// PTH launch: a live PTY *is* the channel now, not an error as before).
/// 3. **Enqueue a ticket** in the [`AgentMailbox`] (registering the reply slot)
/// **then write** the task into the target's terminal, prefixed with the asking
/// agent + ticket id so the target knows to answer via `idea_reply`.
/// 4. **Await** the [`domain::mailbox::PendingReply`] bounded by [`ASK_AGENT_TIMEOUT`]:
/// the target's later `idea_reply(result)` lands in [`Self::reply`] →
/// `mailbox.resolve`, waking this await. On timeout the ticket is retired from
/// the head ([`AgentMailbox::cancel_head`]) — **the target stays alive** — and a
/// typed timeout is returned (retry possible).
/// 5. Return the reply as [`OrchestratorOutcome::reply`] and publish
/// [`DomainEvent::AgentReplied`].
///
/// # Errors
/// - [`AppError::NotFound`] if the target agent is unknown;
/// - [`AppError::Invalid`] if the mailbox/PTY channel is not wired;
/// - [`AppError::Process`] on a launch/PTY-write failure, or on the await timeout
/// (turn timeout *or* queue-wait timeout — same typed error).
async fn ask_agent(
&self,
project: &Project,
target: String,
task: String,
requester: Option<AgentId>,
) -> Result<OrchestratorOutcome, AppError> {
let (input, mailbox) = match (&self.input, &self.mailbox) {
(Some(i), Some(m)) => (i, m),
_ => {
return Err(AppError::Invalid(
"la messagerie inter-agents (idea_ask_agent) n'est pas disponible : \
médiateur d'entrée non câblé"
.to_owned(),
))
}
};
let agent = self
.find_agent_by_name(project, &target)
.await?
.ok_or_else(|| AppError::NotFound(format!("agent {target}")))?;
let agent_id = agent.id;
// F2 — garde profil : refuser **immédiatement** une cible dont le profil ne
// sait pas consommer le pont `idea_*` matérialisé via `.mcp.json`, plutôt que
// de laisser le round-trip échouer en timeout muet (300s).
self.guard_mcp_bridge_supported(&agent.profile_id, &target)
.await?;
// Détection de cycle (cadrage C3 §6) : si l'ask vient d'un **agent** A vers la
// cible B, refuser AVANT tout enqueue si poser l'arête A→B fermerait un cycle
// d'attente (B attend déjà …→A). Pur, sans I/O ⇒ jamais de deadlock.
if let Some(from) = requester {
let cycles = {
let g = self
.wait_for
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
g.would_cycle(from, agent_id)
};
if cycles {
crate::diag!(
"[rendezvous] cycle refused: requester={from} -> target={target} \
(agent {agent_id}) (deadlock évité, aucun enqueue)",
);
return Err(AppError::Invalid(format!(
"délégation ré-entrante refusée : demander à l'agent '{target}' créerait \
un cycle d'attente inter-agents (deadlock évité)"
)));
}
}
// Résoudre paresseusement le **fil** de l'ask : A↔B si un agent demande, sinon
// User↔B. La session vivante est désormais keyée par conversation (lève
// `session-registry-agent-ambiguity`).
let conversation_id = self.resolve_conversation(requester, agent_id);
// Sérialisation FIFO **par agent** (A0) : verrou de tour de la **cible**, tenu
// pour TOUT le tour (enqueue → réponse). RAII : tombe sur chaque early-return.
let lock = self.ask_lock_for(&agent_id);
// Diagnostics : tracer l'attente du verrou de tour. Un « turn-lock waiting » sans
// « turn-lock acquired » qui suit signe un tour précédent jamais relâché (la cible
// reste Busy) — le verrou n'est libéré qu'à la fin du tour courant.
crate::diag!(
"[rendezvous] turn-lock waiting: requester={} -> target={target} (agent {agent_id}) \
wait_cap_ms={}",
requester.map_or_else(|| "user".to_owned(), |a| a.to_string()),
ASK_QUEUE_WAIT_CAP.as_millis(),
);
let lock_wait = Instant::now();
let _turn = match tokio::time::timeout(ASK_QUEUE_WAIT_CAP, lock.lock_owned()).await {
Ok(guard) => {
crate::diag!(
"[rendezvous] turn-lock acquired: target={target} (agent {agent_id}) \
waited_ms={}",
lock_wait.elapsed().as_millis(),
);
guard
}
Err(_elapsed) => {
crate::diag!(
"[rendezvous] turn-lock TIMEOUT: target={target} (agent {agent_id}) \
waited_ms={} (tour précédent jamais relâché ?)",
lock_wait.elapsed().as_millis(),
);
return Err(AppError::from(domain::ports::AgentSessionError::Timeout));
}
};
// Poser l'arête d'attente A→B (retirée en fin de tour par le RAII `_edge`).
let _edge = requester.map(|from| WaitEdgeGuard::new(self, from, agent_id));
// ── Chemin **structuré** (readiness/heartbeat lot 1) ──────────────────────
// Une cible à `structured_adapter` n'a **pas** de PTY : `ensure_live_pty`
// échouerait, et le tour ne pourrait se débloquer que par un `idea_reply`
// explicite (cause racine du blocage `Busy`). Quand le registre structuré est
// câblé et que la cible a une session vivante, on draine son tour via
// `drain_with_readiness` : le `Final` déterministe réveille le `pending` (valeur
// de retour) **et** marque l'agent `Idle` (`mark_idle`). On enregistre tout de
// même un ticket dans la FIFO pour la comptabilité busy et pour préserver
// `idea_reply` comme signal **alternatif** (premier arrivé gagne).
if let Some(structured) = self.structured.as_ref() {
// Lot 1b — auto-lancement d'une cible **froide** structurée : si le registre
// est câblé, qu'aucune session ne vit encore pour la cible, mais que son
// profil porte un `structured_adapter`, on **démarre** sa session via le
// launcher (qui route §17.4 vers `launch_structured` et l'insère dans CE
// même registre, avec la conf MCP matérialisée) plutôt que de tomber dans
// `ensure_live_pty` — chemin PTY qui échouerait pour une cible sans PTY.
// Une cible SANS `structured_adapter` (agent PTY/TUI legacy) conserve le
// chemin `ensure_live_pty` ci-dessous (zéro régression).
if let Some(session) = self
.ensure_structured_session(project, agent_id, &agent.profile_id, structured)
.await?
{
return self
.ask_structured(
project,
agent_id,
&target,
conversation_id,
requester,
task,
session.as_ref(),
)
.await;
}
}
// 1. Garantir la cible vivante en PTY pour CE fil ; lier sa session à la
// conversation, et brancher son handle d'entrée sur le médiateur (livraison).
let (handle, cold_launch) = self
.ensure_live_pty(project, agent_id, conversation_id, &target)
.await?;
// Arm prompt-ready detection (C5) with the target profile's literal marker, so a
// return-to-prompt frees the turn (the other OR signal being `idea_reply`).
let (prompt_pattern, submit, has_mcp) =
self.prompt_and_submit_for_agent(project, agent_id).await;
// Gate cold-launch : un agent froid n'est pas encore à son prompt. On diffère le
// 1er tour s'il existe un signal pour le libérer — soit le prompt-ready watcher
// (pattern non vide), soit la connexion du pont MCP de l'agent
// (`InputMediator::release_cold_start`, déclenchée par l'McpServer). Sans aucun
// des deux ⇒ pas de gate (livraison immédiate, sinon blocage indéfini).
let gate_cold_start =
cold_launch && (prompt_pattern.as_ref().is_some_and(|p| !p.is_empty()) || has_mcp);
// Diagnostics : décision de gate du premier tour. `gate_cold_start=false` sur une
// cible froide SANS signal de libération (ni prompt-ready ni MCP) livrerait
// immédiatement ; un `true` diffère la livraison jusqu'au signal.
crate::diag!(
"[rendezvous] gate decision: target={target} (agent {agent_id}) cold_launch={cold_launch} \
has_prompt_pattern={} has_mcp={has_mcp} gate_cold_start={gate_cold_start}",
prompt_pattern.as_ref().is_some_and(|p| !p.is_empty()),
);
if gate_cold_start {
input.mark_starting(agent_id);
}
input.bind_handle_with_prompt(agent_id, handle.clone(), prompt_pattern, submit);
// 2. Enregistrer le ticket (slot de réponse) + livrer le tour via le médiateur
// (écriture sérialisée dans le PTY — plus d'écriture ad hoc ici). Le ticket
// porte la source (Human/Agent) et la conversation cible.
let requester_label = self.requester_label(project, requester).await;
let ticket_id = TicketId::new_random();
// Checkpoint Prompt (P6b, best-effort) : persister l'invite AVANT que `task` ne
// soit déplacé dans le `Ticket`. Source = origine de la requête (agent demandeur
// `Some(from)` ⇒ Agent, sinon Humain) — **même** `InputSource` que le ticket.
let prompt_source = match requester {
Some(from) => InputSource::agent(from),
None => InputSource::Human,
};
self.record_turn_best_effort(
&project.root,
conversation_id,
prompt_source,
TurnRole::Prompt,
task.clone(),
)
.await;
let ticket = match requester {
Some(from) => {
Ticket::from_agent(ticket_id, from, conversation_id, requester_label, task)
}
None => Ticket::from_human(ticket_id, conversation_id, requester_label, task),
};
// Timeout de tour piloté par profil (lot 2) : `turn_timeout_ms` de la cible si
// défini, sinon le défaut [`ASK_AGENT_TIMEOUT`]. Arme aussi le seuil de stall sur
// le médiateur AVANT l'enqueue (consommé au start_turn).
let turn_timeout = self.turn_timeout_for(project, agent_id).await;
let pending = input.enqueue(agent_id, ticket);
// Rendezvous beacon (diagnostics) : l'ask est désormais en attente du
// `idea_reply` (ou prompt-ready) de la cible. Si la cible termine son tour en
// texte SANS appeler `idea_reply`, ce beacon « ask started » n'aura pas de
// « ask resolved » correspondant avant l'expiration du `turn_timeout` — la
// signature exacte du blocage Main→cible.
let started = Instant::now();
crate::diag!(
"[rendezvous] ask started: requester={} -> target={target} (agent {agent_id}) \
conversation={conversation_id} ticket={ticket_id} cold_launch={cold_launch} \
gate_cold_start={gate_cold_start} turn_timeout_ms={}",
requester.map_or_else(|| "user".to_owned(), |a| a.to_string()),
turn_timeout.as_millis(),
);
// Garde RAII de fin de tour, armé JUSTE après l'enqueue (la cible est maintenant
// `Busy`). Quel que soit le chemin de sortie — erreur, timeout, ou **futur
// abandonné (drop)** — son `Drop` ramène la cible `Idle` et retire le ticket
// fantôme de la FIFO. C'est le fix de la cause racine (cf. [`BusyTurnGuard`]).
let busy_guard =
BusyTurnGuard::new(Arc::clone(input), Arc::clone(mailbox), agent_id, ticket_id);
// Delivery is the mediator's responsibility (`InputMediator::enqueue` writes the
// turn into the bound handle). The service no longer writes the PTY directly —
// no ad-hoc `[IdeA · tâche …]` line here, no `\r` band-aid (cadrage C3 §5.1).
// 3. Attendre la réponse, bornée. Timeout/canal fermé ⇒ le garde retire le ticket
// (cible laissée vivante) et la ramène `Idle` au Drop ; renvoie une erreur typée.
match tokio::time::timeout(turn_timeout, pending).await {
Ok(Ok(TurnResolution::Replied(result))) => {
// Checkpoint Response (P6b, best-effort) : persister la réponse AVANT de
// déplacer `result` dans `reply_outcome`. Source = la **cible** (c'est
// elle qui a rendu le tour) ; même conversation que le Prompt.
self.record_turn_best_effort(
&project.root,
conversation_id,
InputSource::agent(agent_id),
TurnRole::Response,
result.clone(),
)
.await;
// Auto-memory harvest (Lot E1), APRÈS l'append/handoff de la réponse :
// parse les blocs ` ```idea-memory ` et persiste les notes valides.
// Best-effort strict — n'altère jamais ce succès de délégation.
self.harvest_memory_best_effort(&project.root, &result)
.await;
// Succès : désarmer le garde AVANT de retourner. Le `mark_idle` propre
// sur cette branche est porté par le médiateur (prompt-ready / idea_reply
// qui a résolu le `pending`) ; on ne veut ni re-`cancel_head` un ticket
// déjà résolu, ni libérer un busy state qui ne nous appartient plus.
busy_guard.disarm();
crate::diag!(
"[rendezvous] ask resolved: target={target} (agent {agent_id}) \
ticket={ticket_id} after_ms={} reply_len={}",
started.elapsed().as_millis(),
result.len(),
);
Ok(self.reply_outcome(agent_id, &target, result))
}
// La cible est revenue à son prompt SANS `idea_reply` : la fenêtre de grâce a
// expiré et le médiateur a complété le tour « sans réponse ». Le ticket est
// déjà retiré (`complete_without_reply`) et la cible déjà `Idle` (prompt-ready
// `mark_idle`) ⇒ on **désarme** le garde (pas de `cancel_head`/`mark_idle`
// redondant ni de beacon « busy-guard freed » trompeur) et on renvoie une
// erreur typée claire/retryable, en ~G au lieu du timeout long.
Ok(Ok(TurnResolution::ReturnedToPromptNoReply)) => {
busy_guard.disarm();
crate::diag!(
"[rendezvous] ask returned-to-prompt-no-reply: target={target} \
(agent {agent_id}) ticket={ticket_id} after_ms={}",
started.elapsed().as_millis(),
);
Err(AppError::TargetReturnedNoReply(target))
}
// Erreur / timeout : on laisse le garde faire `cancel_head` + `mark_idle` au
// Drop (retrait des `cancel_head` redondants — `cancel_head` reste idempotent).
Ok(Err(_cancelled)) => {
crate::diag!(
"[rendezvous] ask channel-closed: target={target} (agent {agent_id}) \
ticket={ticket_id} after_ms={}",
started.elapsed().as_millis(),
);
Err(AppError::Process(format!(
"agent {target} : canal de réponse fermé avant un résultat"
)))
}
Err(_elapsed) => {
crate::diag!(
"[rendezvous] ask TIMEOUT: target={target} (agent {agent_id}) \
ticket={ticket_id} after_ms={} (la cible n'a jamais appelé idea_reply \
ni atteint son prompt-ready dans le turn_timeout)",
started.elapsed().as_millis(),
);
Err(AppError::from(domain::ports::AgentSessionError::Timeout))
}
}
}
/// Chemin `ask` **structuré** (readiness/heartbeat lot 1) : la cible a un
/// `structured_adapter` et une session vivante ([`StructuredSessions`]). On pilote
/// son tour **directement** sur le port [`domain::ports::AgentSession`] et on le
/// draine via [`drain_with_readiness`] : le [`domain::ports::ReplyEvent::Final`]
/// déterministe rend le contenu (réveille le `pending`) **et** marque l'agent `Idle`
/// (`mark_idle`), sans dépendre d'un `idea_reply` explicite ni d'un PTY (que la cible
/// n'a pas). C'est le fix de la cause racine du blocage `Busy`.
///
/// On enregistre tout de même un ticket dans la FIFO (`enqueue`) pour la comptabilité
/// busy et pour **préserver `idea_reply` comme signal alternatif** : on attend la
/// **première** des deux issues (réponse de la session OU résolution du `pending`).
/// Le checkpoint Prompt/Response best-effort est conservé à l'identique du chemin PTY.
#[allow(clippy::too_many_arguments)]
async fn ask_structured(
&self,
project: &Project,
agent_id: AgentId,
target: &str,
conversation_id: domain::conversation::ConversationId,
requester: Option<AgentId>,
task: String,
session: &dyn domain::ports::AgentSession,
) -> Result<OrchestratorOutcome, AppError> {
let (input, mailbox) = match (&self.input, &self.mailbox) {
(Some(i), Some(m)) => (i, m),
_ => {
return Err(AppError::Invalid(
"la messagerie inter-agents (idea_ask_agent) n'est pas disponible : \
médiateur d'entrée non câblé"
.to_owned(),
))
}
};
// Checkpoint Prompt (best-effort), AVANT de déplacer `task` dans le ticket.
let prompt_source = match requester {
Some(from) => InputSource::agent(from),
None => InputSource::Human,
};
self.record_turn_best_effort(
&project.root,
conversation_id,
prompt_source,
TurnRole::Prompt,
task.clone(),
)
.await;
// Ticket dans la FIFO : comptabilité busy + `idea_reply` comme signal alternatif.
let requester_label = self.requester_label(project, requester).await;
let ticket_id = TicketId::new_random();
let ticket = match requester {
Some(from) => Ticket::from_agent(
ticket_id,
from,
conversation_id,
requester_label,
task.clone(),
),
None => Ticket::from_human(ticket_id, conversation_id, requester_label, task.clone()),
};
// Timeout de tour piloté par profil (lot 2) + armement du seuil de stall, AVANT
// l'enqueue qui démarre le tour (le médiateur arme alors sa fenêtre de vivacité).
let turn_timeout = self.turn_timeout_for(project, agent_id).await;
let pending = input.enqueue(agent_id, ticket);
// Garde RAII de fin de tour, armé JUSTE après l'enqueue (cible `Busy`). Couvre
// toutes les sorties — `return Err` des bras du select, OU **futur abandonné
// (drop)** — en ramenant la cible `Idle` au Drop (cf. [`BusyTurnGuard`]). C'est
// le fix de la cause racine du blocage `Busy` à vie.
let busy_guard =
BusyTurnGuard::new(Arc::clone(input), Arc::clone(mailbox), agent_id, ticket_id);
// Rendezvous beacon (chemin structuré) : équivalent du « ask started » du chemin
// PTY. La cible n'a pas de PTY ; le tour se débloque sur le `Final` de sa session
// OU sur un `idea_reply`. Sans l'un des deux avant `turn_timeout`, le drain expire.
let started = Instant::now();
crate::diag!(
"[rendezvous] ask started (structured): requester={} -> target={target} \
(agent {agent_id}) conversation={conversation_id} ticket={ticket_id} \
turn_timeout_ms={}",
requester.map_or_else(|| "user".to_owned(), |a| a.to_string()),
turn_timeout.as_millis(),
);
// Drainer le tour structuré (le `Final` ⇒ contenu + `mark_idle`), borné par le
// **même** garde-fou que le chemin PTY (profil ou défaut). On attend la
// **première** issue : le tour structuré OU un `idea_reply` explicite.
let drain =
drain_with_readiness(session, &task, Some(turn_timeout), input.as_ref(), agent_id);
let result = tokio::select! {
biased;
// Issue déterministe : la session a rendu son `Final`.
drained = drain => match drained {
Ok(content) => {
crate::diag!(
"[rendezvous] ask resolved (structured/final): target={target} \
(agent {agent_id}) ticket={ticket_id} after_ms={} reply_len={}",
started.elapsed().as_millis(),
content.len(),
);
content
}
// Erreur de drain : le garde fait `cancel_head` + `mark_idle` au Drop.
Err(err) => {
crate::diag!(
"[rendezvous] ask drain-error (structured): target={target} \
(agent {agent_id}) ticket={ticket_id} after_ms={} err={err}",
started.elapsed().as_millis(),
);
return Err(AppError::from(err));
}
},
// Issue alternative : un `idea_reply` explicite a résolu le ticket d'abord.
replied = pending => match replied {
Ok(TurnResolution::Replied(content)) => {
// La session draine encore en arrière-plan ; on s'assure que la FIFO
// avance même si le `Final` n'est pas encore observé.
input.mark_idle(agent_id);
crate::diag!(
"[rendezvous] ask resolved (structured/reply): target={target} \
(agent {agent_id}) ticket={ticket_id} after_ms={} reply_len={}",
started.elapsed().as_millis(),
content.len(),
);
content
}
// Retour au prompt sans `idea_reply` : très rare sur le chemin structuré
// (la cible n'a pas de PTY/watcher), mais on traite la résolution comme
// sur le chemin PTY — erreur typée claire/retryable, garde désarmé.
Ok(TurnResolution::ReturnedToPromptNoReply) => {
input.mark_idle(agent_id);
busy_guard.disarm();
crate::diag!(
"[rendezvous] ask returned-to-prompt-no-reply (structured): \
target={target} (agent {agent_id}) ticket={ticket_id} after_ms={}",
started.elapsed().as_millis(),
);
return Err(AppError::TargetReturnedNoReply(target.to_owned()));
}
// Canal fermé : le garde fait `cancel_head` + `mark_idle` au Drop.
Err(_cancelled) => {
crate::diag!(
"[rendezvous] ask channel-closed (structured): target={target} \
(agent {agent_id}) ticket={ticket_id} after_ms={}",
started.elapsed().as_millis(),
);
return Err(AppError::Process(format!(
"agent {target} : canal de réponse fermé avant un résultat"
)));
}
},
};
// Succès : désarmer le garde (le `mark_idle` propre est déjà porté par le `Final`
// de la session ou par le bras `replied`) — pas de double `cancel_head`.
busy_guard.disarm();
// Checkpoint Response (best-effort), AVANT de déplacer `result`.
self.record_turn_best_effort(
&project.root,
conversation_id,
InputSource::agent(agent_id),
TurnRole::Response,
result.clone(),
)
.await;
// Auto-memory harvest (Lot E1), après l'append/handoff de la réponse.
// Best-effort strict — un échec n'altère jamais ce succès de délégation.
self.harvest_memory_best_effort(&project.root, &result)
.await;
Ok(self.reply_outcome(agent_id, target, result))
}
/// `SubmitHumanInput` (cadrage C4 §5.3) — the **human** Envoyer path.
///
/// The operator types into IdeA's mediated input; this resolves the `User↔Agent`
/// thread, ensures the target is live in the PTY registry, binds its handle on the
/// mediator, and **enqueues** a `Ticket::from_human` into the **same FIFO** the
/// inter-agent delegations use (`InputMediator::enqueue`) — so a human submit and a
/// delegation serialise on the same agent («1 agent = 1 employee»).
///
/// Unlike [`Self::ask_agent`], it is **fire-and-forget**: the human watches the
/// terminal for the answer, so we do **not** await the [`PendingReply`] (the reply
/// slot is registered and simply left to resolve/expire on its own). The busy
/// event is emitted at the mediator's source (Idle→Busy on the starting enqueue).
///
/// # Errors
/// - [`AppError::Invalid`] if the input mediator is not wired;
/// - [`AppError::NotFound`] if the target agent is unknown;
/// - [`AppError::Process`] on a launch/PTY failure while ensuring the live session.
pub async fn submit_human_input(
&self,
project: &Project,
agent_id: AgentId,
text: String,
) -> Result<OrchestratorOutcome, AppError> {
let input = self.input.as_ref().ok_or_else(|| {
AppError::Invalid(
"l'entrée médiée (submit_agent_input) n'est pas disponible : \
médiateur d'entrée non câblé"
.to_owned(),
)
})?;
// Display label for error/launch messages; the agent must exist in the
// manifest. A human submit to an unknown id is a NotFound, never a panic.
let target = self
.find_name_by_agent_id(project, agent_id)
.await
.ok_or_else(|| AppError::NotFound(format!("agent {agent_id}")))?;
let target = target.as_str();
// User↔Agent thread (no requester ⇒ left = User). Same lazy resolution as ask.
let conversation_id = self.resolve_conversation(None, agent_id);
// Ensure the target is live for this thread and bind its input handle on the
// mediator (delivery path). Same call the ask path uses.
let (handle, cold_launch) = self
.ensure_live_pty(project, agent_id, conversation_id, target)
.await?;
let (prompt_pattern, submit, has_mcp) =
self.prompt_and_submit_for_agent(project, agent_id).await;
// Gate cold-launch : un agent froid n'est pas encore à son prompt. On diffère le
// 1er tour s'il existe un signal pour le libérer — soit le prompt-ready watcher
// (pattern non vide), soit la connexion du pont MCP de l'agent
// (`InputMediator::release_cold_start`, déclenchée par l'McpServer). Sans aucun
// des deux ⇒ pas de gate (livraison immédiate, sinon blocage indéfini).
let gate_cold_start =
cold_launch && (prompt_pattern.as_ref().is_some_and(|p| !p.is_empty()) || has_mcp);
if gate_cold_start {
input.mark_starting(agent_id);
}
input.bind_handle_with_prompt(agent_id, handle, prompt_pattern, submit);
// Enqueue a human-sourced ticket in the SAME FIFO as delegations. Fire-and-
// forget: we drop the PendingReply (the human reads the terminal). The
// mediator emits AgentBusyChanged at the source on a starting turn.
let ticket = Ticket::from_human(TicketId::new_random(), conversation_id, "vous", text);
let _pending = input.enqueue(agent_id, ticket);
Ok(OrchestratorOutcome {
detail: format!("submitted human input to agent {target}"),
reply: None,
})
}
/// Interrupt (cadrage C4 §5.3) — the **Interrompre** path (Échap/stop).
///
/// Resolves the target by name and calls [`InputMediator::preempt`], which signals
/// the running turn to stop (best-effort interrupt byte to the agent's bound PTY
/// handle). It is **not** an enqueue and resolves **no** ticket — a pending caller
/// is never silently answered. Idempotent: interrupting an idle agent is a no-op.
///
/// # Errors
/// - [`AppError::Invalid`] if the input mediator is not wired;
/// - [`AppError::NotFound`] if the target agent is unknown.
pub async fn interrupt_agent(
&self,
project: &Project,
agent_id: AgentId,
) -> Result<OrchestratorOutcome, AppError> {
let input = self.input.as_ref().ok_or_else(|| {
AppError::Invalid(
"l'interruption (interrupt_agent) n'est pas disponible : \
médiateur d'entrée non câblé"
.to_owned(),
)
})?;
// Confirm the agent exists (typed NotFound rather than a silent no-op on a
// bogus id). The manifest lookup also keeps the contract symmetric with submit.
if self
.find_name_by_agent_id(project, agent_id)
.await
.is_none()
{
return Err(AppError::NotFound(format!("agent {agent_id}")));
}
input.preempt(agent_id);
Ok(OrchestratorOutcome {
detail: format!("interrupted agent {agent_id}"),
reply: None,
})
}
/// **Ack de livraison** (ARCHITECTURE §20.3) — best-effort, observabilité only.
///
/// The frontend write-portal calls this (via the `delegation_delivered` Tauri
/// command) once it has **physically written** a delegation `ticket` into the agent's
/// native PTY, to distinguish "queued because the human line was busy" from "written"
/// in logs/observability. It **does not** change correlation: the requester's `ask`
/// is still woken by `idea_reply`→`resolve_ticket`, and the per-turn timeout/cycle
/// guards are untouched. It is a pure log no-op when nothing is wired, so a missing
/// mediator/registry never breaks the rendezvous.
pub fn note_delegation_delivered(
&self,
project: &Project,
agent_id: AgentId,
ticket: TicketId,
) {
// Observability beacon only: never mutates the mailbox/busy state nor resolves
// the ticket (that stays `idea_reply`-driven). Kept best-effort by construction —
// a pure, infallible hook so a missing mediator/registry never breaks the
// rendezvous.
let _ = project;
crate::diag!(
"[delivery] front ack received: agent={agent_id} ticket={ticket} \
(write-portal reports text+submit writes completed)"
);
}
/// Libère le premier tour différé d'un agent **lancé à froid** quand son pont MCP se
/// connecte (readiness de démarrage). Pont entre l'McpServer (adapter entrant) et le
/// médiateur d'entrée. No-op si aucun médiateur n'est câblé ou si rien n'est différé.
pub fn release_agent_cold_start(&self, agent: domain::AgentId) {
if let Some(input) = &self.input {
input.release_cold_start(agent);
}
}
/// Déclare si une **cellule terminal frontend** est montée pour `agent` (write-portal
/// actif). Pont entre le write-portal (adapter UI) et le médiateur : un agent avec
/// cellule reçoit ses tours via l'événement `DelegationReady` (le front écrit) ; un
/// agent **headless** (délégué en arrière-plan, sans cellule) voit le médiateur écrire
/// lui-même la tâche dans son PTY — sinon le tour est perdu. No-op sans médiateur.
pub fn set_agent_front_attached(&self, agent: domain::AgentId, attached: bool) {
crate::diag!("[delivery] front attachment changed: agent={agent} attached={attached}");
if let Some(input) = &self.input {
input.set_front_attached(agent, attached);
} else {
crate::diag!(
"[delivery] front attachment ignored because input mediator is not wired: \
agent={agent} attached={attached}"
);
}
}
/// Resolves the conversation thread id for an ask: `A↔B` when an agent requests,
/// else `User↔B` (cadrage C3 §5.2). Without a wired registry, falls back to a
/// stable per-agent id derived from the target (legacy routing — never panics).
fn resolve_conversation(
&self,
requester: Option<AgentId>,
target: AgentId,
) -> domain::conversation::ConversationId {
let left = match requester {
Some(from) => ConversationParty::agent(from),
None => ConversationParty::User,
};
let right = ConversationParty::agent(target);
match &self.conversations {
Some(reg) => reg.resolve(left, right).id,
// Repli pur déterministe partagé avec `LaunchAgent` (ARCHITECTURE §19.7,
// lot P8a) : la même paire dérive la même clé de conversation des deux
// côtés (sauvegarde du handoff ici, dérivation côté cellule là-bas).
None => domain::conversation::ConversationId::for_pair(left, right),
}
}
/// `agent.reply` / `idea_reply`: the target agent renders the result of the task
/// it is currently processing (Option 1, lot B-4).
///
/// Positional correlation: `from` is the **emitting** agent (its identity comes
/// from the MCP handshake, not from a model-managed id), so the result resolves
/// the ticket at the **head** of *that agent's* mailbox queue — the task it is
/// working on. ACK only: no inline payload, no `AgentReplied` here (that belongs
/// to the asking side's `ask_agent`).
///
/// # Errors
/// - [`AppError::Invalid`] if the mailbox is not wired;
/// - [`AppError::Invalid`] if `from` has no in-flight request (a reply with no
/// matching ask) — typed, never a panic.
fn reply(
&self,
from: AgentId,
ticket: Option<TicketId>,
result: String,
) -> Result<OrchestratorOutcome, AppError> {
// Diagnostics : un `idea_reply` est entré, AVANT toute corrélation. Longueur
// seulement (jamais le corps). Si ce beacon apparaît mais que `correlated`/
// `UNMATCHED` ne suit pas, la file n'est pas câblée.
crate::diag!(
"[rendezvous] idea_reply received: from agent {from} ticket={} result_len={}",
ticket.map_or_else(|| "head".to_owned(), |t| t.to_string()),
result.len(),
);
let mailbox = self.mailbox.as_ref().ok_or_else(|| {
crate::diag!(
"[rendezvous] idea_reply DROPPED: from agent {from} (file inter-agents non câblée)",
);
AppError::Invalid(
"idea_reply n'est pas disponible : file inter-agents non câblée".to_owned(),
)
})?;
// Corrélation par ticket quand l'agent l'a renvoyé (déterministe, multi-fil) ;
// sinon repli sur la tête de file de l'émetteur (compat agents mono-fil).
let correlation = match ticket {
Some(ticket_id) => mailbox.resolve_ticket(from, ticket_id, result),
None => mailbox.resolve(from, result),
};
// Rendezvous beacon (diagnostics) : un `idea_reply` est arrivé. Tracer s'il a
// corrélé à un ask en vol — un échec ici (« no matching ask ») signe une
// délégation déjà expirée/abandonnée ou un ticket erroné côté cible.
match &correlation {
Ok(()) => crate::diag!(
"[rendezvous] idea_reply correlated: from agent {from} ticket={}",
ticket.map_or_else(|| "head".to_owned(), |t| t.to_string()),
),
Err(e) => crate::diag!(
"[rendezvous] idea_reply UNMATCHED: from agent {from} ticket={} err={e}",
ticket.map_or_else(|| "head".to_owned(), |t| t.to_string()),
),
}
correlation.map_err(|e| AppError::Invalid(e.to_string()))?;
// Explicit «end-of-turn» signal (cadrage §6, lot C5): an `idea_reply` means the
// emitting agent `from` finished its delegated task ⇒ mark it Idle so its FIFO
// advances to the next queued ticket. This is the deterministic OR signal that
// pairs with prompt-ready detection; whichever fires first frees the turn. No-op
// (and no spurious event) when the mediator is absent or `from` was already idle.
if let Some(input) = self.input.as_ref() {
input.mark_idle(from);
}
Ok(OrchestratorOutcome {
detail: format!("reply from agent {from} delivered"),
reply: None,
})
}
/// Ensures the target agent has a **live PTY session**, returning its handle.
///
/// Reuses the running terminal when present; otherwise launches the agent in the
/// background (a normal PTH launch — a live PTY is the delegation channel). After
/// a launch the handle is resolved from the registry; a missing handle is a
/// [`AppError::Process`] (the launch did not register a PTY session, e.g. a profile
/// IdeA cannot drive as a terminal).
///
/// Le booléen renvoyé indique un **lancement à froid** (`true` quand l'agent n'avait
/// pas de session vivante et vient d'être démarré) : l'appelant l'utilise pour
/// *gater* le premier tour sur le prompt-ready (cf. [`InputMediator::mark_starting`]),
/// car un agent froid n'est pas encore à son prompt. `false` ⇒ session réutilisée
/// (agent déjà chaud) ⇒ livraison immédiate, comportement inchangé.
async fn ensure_live_pty(
&self,
project: &Project,
agent_id: AgentId,
conversation_id: domain::conversation::ConversationId,
target: &str,
) -> Result<(PtyHandle, bool), AppError> {
// «1 session vivante / conversation» (cadrage C3 §5.2) : on cherche d'abord la
// session du **fil**, puis on retombe sur la session de l'agent (compat : un
// agent mono-fil dont la session n'a pas encore été liée à sa conversation).
let existing = self
.sessions
.session_for(conversation_id)
.or_else(|| self.sessions.session_for_agent(&agent_id));
if let Some(session_id) = existing {
if let Some(handle) = self.sessions.handle(&session_id) {
// (Re)lier le fil à cette session vivante (idempotent). Réutilisation
// d'une session déjà chaude ⇒ pas de gate de démarrage à froid.
self.bind_conversation_session(conversation_id, session_id);
return Ok((handle, false));
}
}
// Dead target: launch it in the background (PTY). On injecte ici la déclaration
// MCP **réelle** via le [`McpRuntimeProvider`] câblé (app-tauri détient l'exe et
// l'endpoint) — c'est ce qui permet au pont MCP de la cible de se spawner et donc
// à la cible d'appeler `idea_reply`. Provider absent (ou `runtime_for` → `None`)
// ⇒ déclaration minimale comme avant (dégradation gracieuse).
self.launch_agent
.execute(LaunchAgentInput {
project: project.clone(),
agent_id,
rows: DEFAULT_ROWS,
cols: DEFAULT_COLS,
node_id: None,
conversation_id: None,
mcp_runtime: self
.mcp_runtime_provider
.as_ref()
.and_then(|p| p.runtime_for(project, agent_id)),
})
.await?;
let session_id = self.sessions.session_for_agent(&agent_id).ok_or_else(|| {
AppError::Process(format!(
"agent {target} n'a pas de session terminal vivante après lancement"
))
})?;
// Lier la session fraîchement lancée à CE fil (registre terminal + registre de
// conversations) ⇒ un prochain ask sur le même fil la réutilise.
self.bind_conversation_session(conversation_id, session_id);
let handle = self.sessions.handle(&session_id).ok_or_else(|| {
AppError::Process(format!(
"handle PTY de l'agent {target} introuvable après lancement"
))
})?;
// Lancement à froid : `true` ⇒ l'appelant gatera le premier tour sur le prompt.
Ok((handle, true))
}
/// Lot 1b — garantit une **session structurée vivante** pour la cible d'un `ask`.
///
/// Retourne :
/// - `Ok(Some(session))` si la cible a déjà une session vivante, **ou** si son
/// profil porte un `structured_adapter` et qu'on vient de la (re)lancer ;
/// - `Ok(None)` si la cible n'est **pas** structurée (profil sans `structured_adapter`,
/// ou profil introuvable) ⇒ l'appelant retombe sur le chemin PTY `ensure_live_pty`.
///
/// L'auto-lancement réutilise le **même** [`LaunchAgent`] que le chemin PTV
/// ([`Self::ensure_live_pty`]) avec le **même** `mcp_runtime` matérialisé : pour un
/// profil structuré, le launcher route §17.4 vers `launch_structured`, démarre la
/// session via la fabrique et l'insère dans le registre [`StructuredSessions`]
/// **partagé** (le même `Arc` que `self.structured`, câblé au composition root).
/// La conf MCP est donc matérialisée comme pour une cellule chat lancée à la main,
/// si bien que la cible voit les outils `idea_*` pour répondre.
async fn ensure_structured_session(
&self,
project: &Project,
agent_id: AgentId,
profile_id: &ProfileId,
structured: &Arc<StructuredSessions>,
) -> Result<Option<Arc<dyn domain::ports::AgentSession>>, AppError> {
// Cible déjà chaude : route directe (aucun lancement).
if let Some(session) = structured.session_for_agent(&agent_id) {
return Ok(Some(session));
}
// Cible froide : ne (re)lancer que si le profil sait être piloté en mode
// structuré. Sinon (agent PTY/TUI legacy, ou profil introuvable) ⇒ chemin PTY.
let is_structured = self
.profiles
.list()
.await?
.into_iter()
.find(|p| &p.id == profile_id)
.is_some_and(|p| p.is_selectable());
if !is_structured {
return Ok(None);
}
// Démarrer la session via le launcher (route §17.4 → `launch_structured`,
// insère dans CE registre). Mêmes faits MCP que le chemin PTY pour que le pont
// `idea_*` de la cible se branche. `conversation_id: None` ⇒ le launcher dérive
// l'id de paire (User↔agent) ou réutilise celui de la cellule (P8a), comme pour
// un lancement direct utilisateur.
self.launch_agent
.execute(LaunchAgentInput {
project: project.clone(),
agent_id,
rows: DEFAULT_ROWS,
cols: DEFAULT_COLS,
node_id: None,
conversation_id: None,
mcp_runtime: self
.mcp_runtime_provider
.as_ref()
.and_then(|p| p.runtime_for(project, agent_id)),
})
.await?;
// Le launcher a inséré la session dans le registre partagé : la relire.
structured
.session_for_agent(&agent_id)
.map(Some)
.ok_or_else(|| {
AppError::Process(format!(
"agent {agent_id} : aucune session structurée vivante après lancement"
))
})
}
/// Binds `conversation` to `session` in both the terminal registry (fast
/// `session_for`) and the [`ConversationRegistry`] (domain thread state), when the
/// latter is wired. Idempotent.
fn bind_conversation_session(
&self,
conversation: domain::conversation::ConversationId,
session: domain::SessionId,
) {
self.sessions.bind_conversation(conversation, session);
if let Some(reg) = &self.conversations {
reg.bind_session(conversation, SessionRef::new(session));
}
}
/// Resolves a human-friendly label for the **requesting** agent to prefix the
/// delegated task with. Best-effort: there is no requester id threaded into
/// [`OrchestratorCommand::AskAgent`] today, so this falls back to a stable
/// `"un autre agent"` label. Kept as a seam so a future requester-aware ask can
/// surface the real name without touching the call sites.
async fn requester_label(&self, project: &Project, requester: Option<AgentId>) -> String {
match requester {
None => "un autre agent".to_owned(),
Some(from) => self
.find_name_by_agent_id(project, from)
.await
.unwrap_or_else(|| "un autre agent".to_owned()),
}
}
/// Best-effort display name for an agent id (manifest lookup). `None` when the id
/// is unknown — the caller falls back to a generic label.
async fn find_name_by_agent_id(&self, project: &Project, id: AgentId) -> Option<String> {
self.list_agents
.execute(ListAgentsInput {
project: project.clone(),
})
.await
.ok()?
.agents
.into_iter()
.find(|a| a.id == id)
.map(|a| a.name)
}
/// Builds the success outcome of an `ask` and publishes [`DomainEvent::AgentReplied`]
/// (best-effort: a missing event bus never fails the rendezvous).
fn reply_outcome(
&self,
agent_id: domain::AgentId,
target: &str,
content: String,
) -> OrchestratorOutcome {
if let Some(events) = &self.events {
events.publish(DomainEvent::AgentReplied {
agent_id,
reply_len: content.len(),
});
}
OrchestratorOutcome {
detail: format!("agent {target} replied ({} bytes)", content.len()),
reply: Some(content),
}
}
/// `list_agents`: discovery — return the project's agents exactly as the UI
/// reads them from the manifest, via the **same** [`ListAgents`] use case.
///
/// The list is serialised as a JSON array into [`OrchestratorOutcome::reply`]
/// (the existing inline-payload channel, also used by `ask`), with a one-line
/// count in [`OrchestratorOutcome::detail`]. Each element carries the agent's
/// `id`, `name`, `contextPath`, `profileId`, `origin`, `synchronized` and
/// `skills` (camelCase, the [`domain::Agent`] serde shape).
///
/// # Errors
/// Propagates [`AppError`] from the use case (manifest load / invariant) or a
/// serialisation failure ([`AppError::Invalid`]).
async fn list_agents(&self, project: &Project) -> Result<OrchestratorOutcome, AppError> {
let listed = self
.list_agents
.execute(ListAgentsInput {
project: project.clone(),
})
.await?;
let reply = serde_json::to_string(&listed.agents)
.map_err(|e| AppError::Invalid(format!("failed to serialise agent list: {e}")))?;
Ok(OrchestratorOutcome {
detail: format!("listed {} agent(s)", listed.agents.len()),
reply: Some(reply),
})
}
/// `stop_agent`: translate the agent name → its live session → `CloseTerminal`.
async fn stop_agent(
&self,
project: &Project,
name: String,
) -> Result<OrchestratorOutcome, AppError> {
let agent_id = self
.find_agent_id_by_name(project, &name)
.await?
.ok_or_else(|| AppError::NotFound(format!("agent {name}")))?;
let session_id = self
.sessions
.session_for_agent(&agent_id)
.ok_or_else(|| AppError::NotFound(format!("running session for agent {name}")))?;
self.close_terminal
.execute(CloseTerminalInput { session_id })
.await?;
Ok(OrchestratorOutcome {
detail: format!("stopped agent {name}"),
reply: None,
})
}
/// `update_agent_context`: overwrite the agent's `.md` body.
async fn update_agent_context(
&self,
project: &Project,
name: String,
context: String,
) -> Result<OrchestratorOutcome, AppError> {
let agent_id = self
.find_agent_id_by_name(project, &name)
.await?
.ok_or_else(|| AppError::NotFound(format!("agent {name}")))?;
self.update_context
.execute(UpdateAgentContextInput {
project: project.clone(),
agent_id,
content: context,
})
.await?;
Ok(OrchestratorOutcome {
detail: format!("updated context for agent {name}"),
reply: None,
})
}
/// `create_skill`: create a reusable skill in the requested scope — the same
/// path the UI's "New skill" action takes. For [`SkillScope::Project`] the
/// skill lands under `<root>/.ideai/skills/`; for [`SkillScope::Global`] the
/// project root is ignored by the store.
async fn create_skill(
&self,
project: &Project,
name: String,
content: String,
scope: domain::SkillScope,
) -> Result<OrchestratorOutcome, AppError> {
let created = self
.create_skill
.execute(CreateSkillInput {
name: name.clone(),
// `idea_create_skill` carries no description argument; the affordance
// falls back to the body's first line (see `effective_description`).
description: None,
content,
scope,
project_root: project.root.clone(),
})
.await?;
Ok(OrchestratorOutcome {
detail: format!("created skill {} ({:?})", created.skill.name, scope),
reply: None,
})
}
/// Finds an agent id by display name (case-insensitive) in the project manifest.
async fn find_agent_id_by_name(
&self,
project: &Project,
name: &str,
) -> Result<Option<domain::AgentId>, AppError> {
let listed = self
.list_agents
.execute(ListAgentsInput {
project: project.clone(),
})
.await?;
Ok(listed
.agents
.into_iter()
.find(|a| a.name.eq_ignore_ascii_case(name))
.map(|a| a.id))
}
/// Finds the full [`domain::Agent`] by display name (case-insensitive) in the
/// project manifest. Variante de [`Self::find_agent_id_by_name`] qui conserve
/// l'agent entier (notamment son `profile_id`), nécessaire à la garde F2.
async fn find_agent_by_name(
&self,
project: &Project,
name: &str,
) -> Result<Option<domain::Agent>, AppError> {
let listed = self
.list_agents
.execute(ListAgentsInput {
project: project.clone(),
})
.await?;
Ok(listed
.agents
.into_iter()
.find(|a| a.name.eq_ignore_ascii_case(name)))
}
/// Garde F2 : vérifie que le profil de la cible **sait consommer** le pont
/// `idea_*` matérialisé par IdeA, et renvoie sinon une [`AppError::Invalid`]
/// **immédiate** (au lieu d'un timeout 300s muet sur le round-trip).
///
/// **Critère retenu** (le plus robuste aujourd'hui) : le pont est honoré ssi le
/// profil porte une capacité MCP en stratégie `ConfigFile` ciblant `.mcp.json`
/// **ET** que son adaptateur structuré est `Claude`. En effet IdeA matérialise le
/// serveur MCP sous forme d'un fichier `.mcp.json` dans le run dir, ce que **seul**
/// Claude Code lit réellement ; Codex déclare pourtant la même stratégie
/// `ConfigFile(.mcp.json)` mais lit en pratique `~/.codex/config.toml` ⇒ le pont
/// n'est jamais branché et la cible ne peut pas appeler `idea_reply`. On exige donc
/// l'adaptateur `Claude` plutôt qu'une simple présence de capacité MCP, ce qui
/// exclut Codex de fait et reste valable pour tout futur profil non-Claude.
///
/// Profil introuvable ⇒ on **n'interdit pas** (laisse le flux suivre son cours
/// comme avant) : la garde ne fait que transformer un échec connu en erreur typée.
async fn guard_mcp_bridge_supported(
&self,
profile_id: &ProfileId,
target: &str,
) -> Result<(), AppError> {
let Some(profile) = self
.profiles
.list()
.await?
.into_iter()
.find(|p| &p.id == profile_id)
else {
return Ok(());
};
// Source de vérité UNIQUE (domaine) : un profil supporte le pont `idea_*` ssi
// IdeA matérialise réellement sa config MCP pour la CLI qu'il pilote — Claude
// via `.mcp.json`, Codex via `config.toml`/`CODEX_HOME`. Tout autre couple
// (y compris MCP absent) ⇒ repli fichier, pont non branché.
if profile.materializes_idea_bridge() {
return Ok(());
}
Err(AppError::Invalid(format!(
"la cible '{target}' (profil '{}', adaptateur {:?}) ne supporte pas encore le \
pont idea_* : la délégation inter-agents passe par un serveur MCP qu'IdeA \
matérialise pour la CLI cible, et ce profil ne déclare pas un couple \
(adaptateur × stratégie MCP) pris en charge. Cible un agent dont le profil \
expose le pont MCP (Claude ou Codex).",
profile.name, profile.structured_adapter
)))
}
/// Resolves the target agent profile's **prompt-ready pattern** (§6, lot C5) **and**
/// its **submit config** (`submit_sequence`/`submit_delay_ms`, ARCHITECTURE §20.3) in
/// a single profile lookup. Both are carried into `bind_handle_with_prompt`: the
/// pattern arms prompt detection, the submit config is echoed on the next
/// `DelegationReady` so the frontend write-portal knows how to submit. Returns
/// `(None, default)` when the agent, its profile, or the field is absent — the safe
/// fallback (no pattern ⇒ Idle only via explicit signal/timeout; no submit ⇒ the
/// front applies its own `"\r"`/~60 ms default).
async fn prompt_and_submit_for_agent(
&self,
project: &Project,
agent_id: AgentId,
) -> (Option<String>, SubmitConfig, bool) {
let Some(agent) = self
.list_agents
.execute(ListAgentsInput {
project: project.clone(),
})
.await
.ok()
.and_then(|out| out.agents.into_iter().find(|a| a.id == agent_id))
else {
return (None, SubmitConfig::default(), false);
};
let Some(profile) = self
.profiles
.list()
.await
.ok()
.and_then(|ps| ps.into_iter().find(|p| p.id == agent.profile_id))
else {
return (None, SubmitConfig::default(), false);
};
let submit = submit_config_for_profile(&profile);
// 3e élément : le profil cible déclare-t-il un pont MCP ? Si oui, sa connexion
// (initialize) servira de signal de readiness de démarrage pour libérer un 1er
// tour différé — d'où le gate cold-launch même sans `prompt_ready_pattern`.
(profile.prompt_ready_pattern, submit, profile.mcp.is_some())
}
/// Résout la [`domain::profile::LivenessStrategy`] du profil de la cible (lot 2) :
/// le seuil de stagnation (`stall_after_ms`) et le timeout de tour
/// (`turn_timeout_ms`). `None`/absent ⇒ `(None, None)` : pas de détection de stall et
/// repli sur les bornes par défaut codées en dur (zéro régression pour un profil sans
/// bloc `liveness`).
async fn liveness_for_agent(
&self,
project: &Project,
agent_id: AgentId,
) -> (Option<u32>, Option<u32>) {
let Some(agent) = self
.list_agents
.execute(ListAgentsInput {
project: project.clone(),
})
.await
.ok()
.and_then(|out| out.agents.into_iter().find(|a| a.id == agent_id))
else {
return (None, None);
};
let Some(profile) = self
.profiles
.list()
.await
.ok()
.and_then(|ps| ps.into_iter().find(|p| p.id == agent.profile_id))
else {
return (None, None);
};
match profile.liveness {
Some(l) => (l.stall_after_ms, l.turn_timeout_ms),
None => (None, None),
}
}
/// Borne de tour effective pour un `ask` : le `turn_timeout_ms` du profil de la cible
/// **quand il existe**, sinon le défaut historique [`ASK_AGENT_TIMEOUT`] (lot 2,
/// timeouts pilotés par profil). Arme aussi le seuil de stagnation sur le médiateur
/// avant l'enqueue (battement/`sweep_stalled`).
async fn turn_timeout_for(&self, project: &Project, agent_id: AgentId) -> Duration {
let (stall_after_ms, turn_timeout_ms) = self.liveness_for_agent(project, agent_id).await;
if let Some(input) = &self.input {
input.set_stall_threshold(agent_id, stall_after_ms);
}
resolve_turn_timeout(turn_timeout_ms)
}
/// Resolves a human-friendly profile reference (slug like `claude-code`,
/// command like `claude`, or display name like `Claude Code`) to a configured
/// [`ProfileId`]. Matching is universal — never hard-coded to one AI — by
/// scanning the configured profiles' command and name.
///
/// # Errors
/// [`AppError::NotFound`] when no configured profile matches.
async fn resolve_profile(&self, reference: &str) -> Result<ProfileId, AppError> {
let needle = normalise(reference);
let profiles = self.profiles.list().await?;
profiles
.into_iter()
.find(|p| {
normalise(&p.command) == needle
|| normalise(&p.name) == needle
|| p.id.to_string() == reference
})
.map(|p| p.id)
.ok_or_else(|| AppError::NotFound(format!("profile matching '{reference}'")))
}
}
/// RAII guard that posts a wait-for edge `from → to` for the duration of an ask and
/// removes it on drop (reply, timeout, or any early return) — the same discipline as
/// the per-agent turn lock. Holds a raw pointer-free borrow via the shared mutex on
/// the service's [`WaitForGraph`]; constructed only inside `ask_agent` where the
/// service outlives the guard.
struct WaitEdgeGuard<'a> {
graph: &'a StdMutex<WaitForGraph>,
from: AgentId,
to: AgentId,
}
impl<'a> WaitEdgeGuard<'a> {
fn new(service: &'a OrchestratorService, from: AgentId, to: AgentId) -> Self {
{
let mut g = service
.wait_for
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
g.add_edge(from, to);
}
Self {
graph: &service.wait_for,
from,
to,
}
}
}
impl Drop for WaitEdgeGuard<'_> {
fn drop(&mut self) {
let mut g = self
.graph
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
g.remove_edge(self.from, self.to);
}
}
/// Normalises a profile reference for tolerant matching: lowercased, with spaces,
/// dashes and underscores stripped (`"Claude Code"`, `"claude-code"`, `"claude"`
/// → comparable forms; `claude` ⊂ ... handled by the command match above).
fn normalise(s: &str) -> String {
s.chars()
.filter(|c| c.is_ascii_alphanumeric())
.map(|c| c.to_ascii_lowercase())
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
use domain::profile::{AgentProfile, ContextInjection, StructuredAdapter};
use domain::ProfileId;
// (c) — timeouts pilotés par profil (lot 2) : `turn_timeout_ms` prime sur le défaut.
#[test]
fn profile_turn_timeout_overrides_default() {
// Seuil de profil défini ⇒ il prime sur ASK_AGENT_TIMEOUT.
assert_eq!(
resolve_turn_timeout(Some(120_000)),
Duration::from_millis(120_000)
);
// Absent (profil sans liveness / sans turn_timeout_ms) ⇒ repli sur le défaut.
assert_eq!(resolve_turn_timeout(None), ASK_AGENT_TIMEOUT);
}
fn profile(id: u128, name: &str, command: &str) -> AgentProfile {
AgentProfile::new(
ProfileId::from_uuid(uuid::Uuid::from_u128(id)),
name,
command,
Vec::new(),
ContextInjection::convention_file("CLAUDE.md").unwrap(),
None,
"{agentRunDir}",
None,
)
.unwrap()
}
#[test]
fn normalise_makes_slug_command_and_name_comparable() {
assert_eq!(normalise("Claude Code"), "claudecode");
assert_eq!(normalise("claude-code"), "claudecode");
assert_eq!(normalise("claude_code"), "claudecode");
}
#[test]
fn resolve_matches_by_command_name_or_id() {
// We exercise the pure matching predicate the same way `resolve_profile`
// does, without standing up the whole service/ports.
let p = profile(1, "Claude Code", "claude");
let by_command = normalise("claude") == normalise(&p.command);
let by_name = normalise("claude-code") == normalise(&p.name);
assert!(by_command);
assert!(by_name);
assert_eq!(p.id.to_string(), p.id.to_string());
}
#[test]
fn codex_submit_config_gets_conservative_delay_when_profile_omits_it() {
let p = profile(2, "OpenAI Codex CLI", "codex")
.with_structured_adapter(StructuredAdapter::Codex);
let submit = submit_config_for_profile(&p);
assert_eq!(submit.sequence, None);
assert_eq!(submit.delay_ms, Some(CODEX_SUBMIT_DELAY_MS));
}
#[test]
fn explicit_profile_submit_delay_is_preserved() {
let p = profile(3, "OpenAI Codex CLI", "codex")
.with_structured_adapter(StructuredAdapter::Codex)
.with_submit_delay_ms(900);
let submit = submit_config_for_profile(&p);
assert_eq!(submit.delay_ms, Some(900));
}
// --- BusyTurnGuard (RAII de fin de tour) -------------------------------
//
// Fakes minimaux pour observer ce que le garde appelle à son Drop : un médiateur
// qui enregistre les `mark_idle` et un mailbox qui enregistre les `cancel_head`.
use std::sync::Mutex as TestMutex;
#[derive(Default)]
struct SpyMediator {
idled: TestMutex<Vec<AgentId>>,
}
impl domain::input::InputMediator for SpyMediator {
fn enqueue(
&self,
_agent: AgentId,
_ticket: domain::mailbox::Ticket,
) -> domain::mailbox::PendingReply {
domain::mailbox::PendingReply::new(Box::pin(async {
Err(domain::mailbox::MailboxError::Cancelled)
}))
}
fn preempt(&self, _agent: AgentId) {}
fn mark_idle(&self, agent: AgentId) {
self.idled.lock().unwrap().push(agent);
}
fn busy_state(&self, _agent: AgentId) -> domain::input::AgentBusyState {
domain::input::AgentBusyState::Idle
}
}
#[derive(Default)]
struct SpyMailbox {
cancelled: TestMutex<Vec<(AgentId, TicketId)>>,
}
impl domain::mailbox::AgentMailbox for SpyMailbox {
fn enqueue(
&self,
_agent: AgentId,
_ticket: domain::mailbox::Ticket,
) -> domain::mailbox::PendingReply {
domain::mailbox::PendingReply::new(Box::pin(async {
Err(domain::mailbox::MailboxError::Cancelled)
}))
}
fn resolve(
&self,
_agent: AgentId,
_result: String,
) -> Result<(), domain::mailbox::MailboxError> {
Ok(())
}
fn cancel_head(&self, agent: AgentId, ticket_id: TicketId) {
self.cancelled.lock().unwrap().push((agent, ticket_id));
}
}
fn aid(n: u128) -> AgentId {
AgentId::from_uuid(uuid::Uuid::from_u128(n))
}
fn tid(n: u128) -> TicketId {
TicketId::from_uuid(uuid::Uuid::from_u128(n))
}
/// Drop d'un garde **armé** ⇒ `cancel_head` + `mark_idle` sur la cible (c'est le
/// comportement qui débloque un agent resté `Busy` sur un futur abandonné).
#[test]
fn armed_guard_drop_cancels_head_and_marks_idle() {
let med = Arc::new(SpyMediator::default());
let mb = Arc::new(SpyMailbox::default());
{
let _g = BusyTurnGuard::new(
Arc::clone(&med) as Arc<dyn domain::input::InputMediator>,
Arc::clone(&mb) as Arc<dyn domain::mailbox::AgentMailbox>,
aid(1),
tid(7),
);
} // Drop ici.
assert_eq!(
*med.idled.lock().unwrap(),
vec![aid(1)],
"mark_idle au Drop"
);
assert_eq!(
*mb.cancelled.lock().unwrap(),
vec![(aid(1), tid(7))],
"cancel_head au Drop"
);
}
/// Garde **désarmé** (branche succès) ⇒ son Drop est un no-op : pas de `mark_idle`
/// parasite, surtout pas de `cancel_head` d'un ticket déjà résolu.
#[test]
fn disarmed_guard_drop_is_a_noop() {
let med = Arc::new(SpyMediator::default());
let mb = Arc::new(SpyMailbox::default());
let g = BusyTurnGuard::new(
Arc::clone(&med) as Arc<dyn domain::input::InputMediator>,
Arc::clone(&mb) as Arc<dyn domain::mailbox::AgentMailbox>,
aid(1),
tid(7),
);
g.disarm();
assert!(
med.idled.lock().unwrap().is_empty(),
"pas de mark_idle après disarm"
);
assert!(
mb.cancelled.lock().unwrap().is_empty(),
"pas de cancel_head après disarm"
);
}
}