agent conversation fix
This commit is contained in:
@ -22,7 +22,7 @@ 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::mailbox::{Ticket, TicketId};
|
||||
use domain::ports::{Clock, EventBus, ProfileStore, PtyHandle};
|
||||
use domain::project::ProjectPath;
|
||||
use domain::{
|
||||
@ -315,6 +315,11 @@ pub struct OrchestratorService {
|
||||
/// `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>,
|
||||
/// Snapshot opérationnel des mêmes arêtes d'attente, exposé au chemin d'arrêt
|
||||
/// utilisateur : si A est stoppé pendant qu'il attend B, IdeA stoppe aussi B.
|
||||
/// Séparé de [`WaitForGraph`] qui reste un objet domaine minimal de détection de
|
||||
/// cycle, sans API de traversal.
|
||||
active_waits: StdMutex<Vec<(AgentId, AgentId)>>,
|
||||
/// 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).
|
||||
@ -459,6 +464,7 @@ impl OrchestratorService {
|
||||
mailbox: None,
|
||||
conversations: None,
|
||||
wait_for: StdMutex::new(WaitForGraph::new()),
|
||||
active_waits: StdMutex::new(Vec::new()),
|
||||
events: None,
|
||||
ask_locks: StdMutex::new(HashMap::new()),
|
||||
mcp_runtime_provider: None,
|
||||
@ -537,6 +543,46 @@ impl OrchestratorService {
|
||||
Arc::clone(locks.entry(*agent_id).or_default())
|
||||
}
|
||||
|
||||
/// Returns the transitive set of agents currently waited on by `agent`.
|
||||
///
|
||||
/// Used by user-driven cancellation: stopping A while A waits on B should also
|
||||
/// stop B, and then any agent B itself waits on. The snapshot is best-effort and
|
||||
/// lock-bounded; callers perform the actual interruption/stop outside the mutex.
|
||||
#[must_use]
|
||||
pub fn active_wait_dependencies(&self, agent: AgentId) -> Vec<AgentId> {
|
||||
let edges = self
|
||||
.active_waits
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner)
|
||||
.clone();
|
||||
Self::active_wait_dependencies_from_edges(&edges, agent)
|
||||
}
|
||||
|
||||
/// Returns the transitive dependencies of `agent` in an active wait-edge snapshot.
|
||||
#[must_use]
|
||||
fn active_wait_dependencies_from_edges(
|
||||
edges: &[(AgentId, AgentId)],
|
||||
agent: AgentId,
|
||||
) -> Vec<AgentId> {
|
||||
let mut out = Vec::new();
|
||||
let mut stack: Vec<AgentId> = edges
|
||||
.iter()
|
||||
.filter_map(|(from, to)| (*from == agent).then_some(*to))
|
||||
.collect();
|
||||
while let Some(next) = stack.pop() {
|
||||
if out.contains(&next) {
|
||||
continue;
|
||||
}
|
||||
out.push(next);
|
||||
stack.extend(
|
||||
edges
|
||||
.iter()
|
||||
.filter_map(|(from, to)| (*from == next).then_some(*to)),
|
||||
);
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
/// 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
|
||||
@ -1203,32 +1249,28 @@ impl OrchestratorService {
|
||||
})
|
||||
}
|
||||
|
||||
/// `agent.message` / `idea_ask_agent`: the **inter-agent delegation rendezvous**
|
||||
/// (Option 1 « Terminal + MCP », lot B-3).
|
||||
/// `agent.message` / `idea_ask_agent`: the **inter-agent delegation rendezvous**.
|
||||
///
|
||||
/// 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:
|
||||
/// The MCP tool and the file watcher both enter here, but they are only entrypoints
|
||||
/// into the same application orchestration path:
|
||||
///
|
||||
/// 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
|
||||
/// 2. Resolve the requester/target conversation and reject wait-for cycles before
|
||||
/// enqueueing anything.
|
||||
/// 3. Ensure the target has a structured/headless [`domain::ports::AgentSession`],
|
||||
/// launching it through the profile adapter if it is cold.
|
||||
/// 4. Drive the turn directly through [`domain::ports::AgentSession::send`] and
|
||||
/// drain until the structured `Final` is captured. The target does **not** see
|
||||
/// a ticket and does **not** call `idea_reply`.
|
||||
/// 5. Return the captured final answer 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
|
||||
/// - [`AppError::Invalid`] if structured/headless orchestration is not wired, or if
|
||||
/// the target profile cannot be driven as a structured session;
|
||||
/// - [`AppError::Process`] on launch/session failure, or on the await timeout
|
||||
/// (turn timeout *or* queue-wait timeout — same typed error).
|
||||
async fn ask_agent(
|
||||
&self,
|
||||
@ -1237,29 +1279,12 @@ impl OrchestratorService {
|
||||
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.
|
||||
@ -1323,234 +1348,37 @@ impl OrchestratorService {
|
||||
// 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?;
|
||||
// Résout la config de soumission du profil cible + s'il déclare un pont MCP.
|
||||
let (submit, has_mcp) = self.submit_and_mcp_for_agent(project, agent_id).await;
|
||||
// Gate cold-launch : un agent froid n'est pas encore prêt à recevoir son 1er tour.
|
||||
// On le diffère s'il existe un signal pour le libérer — la connexion du pont MCP
|
||||
// de l'agent (`InputMediator::release_cold_start`, déclenchée par l'McpServer sur
|
||||
// `initialize`). C'est désormais le **seul** signal de readiness (le watcher
|
||||
// prompt-ready PTY a été supprimé). Sans pont MCP ⇒ pas de gate (livraison
|
||||
// immédiate, sinon blocage indéfini).
|
||||
let gate_cold_start = cold_launch && has_mcp;
|
||||
// Diagnostics : décision de gate du premier tour. `gate_cold_start=false` sur une
|
||||
// cible froide SANS pont MCP livrerait immédiatement ; un `true` diffère la
|
||||
// livraison jusqu'au signal MCP-initialize.
|
||||
crate::diag!(
|
||||
"[rendezvous] gate decision: target={target} (agent {agent_id}) cold_launch={cold_launch} \
|
||||
has_mcp={has_mcp} gate_cold_start={gate_cold_start}",
|
||||
);
|
||||
if gate_cold_start {
|
||||
input.mark_starting(agent_id);
|
||||
}
|
||||
input.bind_handle_with_submit(agent_id, handle.clone(), 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,
|
||||
// ── Chemin **structuré/headless uniquement** ──────────────────────────────
|
||||
// La conversation inter-agent ne passe plus par le PTY ni par le rendez-vous
|
||||
// `idea_reply` MCP. La cible doit être pilotable via `AgentSession::send`; son
|
||||
// `Final` est la seule réponse normale du tour. Les profils legacy PTY/TUI sont
|
||||
// refusés au lieu de retomber sur l'ancien chemin MCP fragile.
|
||||
let structured = self.structured.as_ref().ok_or_else(|| {
|
||||
AppError::Invalid(
|
||||
"la conversation inter-agent headless n'est pas disponible : registre \
|
||||
de sessions structurées non câblé"
|
||||
.to_owned(),
|
||||
)
|
||||
})?;
|
||||
let session = self
|
||||
.ensure_structured_session(project, agent_id, &agent.profile_id, structured)
|
||||
.await?
|
||||
.ok_or_else(|| {
|
||||
AppError::Invalid(format!(
|
||||
"la cible '{target}' ne peut pas recevoir de conversation inter-agent : \
|
||||
son profil ne déclare pas d'adaptateur structured/headless"
|
||||
))
|
||||
})?;
|
||||
self.ask_structured(
|
||||
project,
|
||||
agent_id,
|
||||
&target,
|
||||
conversation_id,
|
||||
prompt_source,
|
||||
TurnRole::Prompt,
|
||||
task.clone(),
|
||||
requester,
|
||||
task,
|
||||
session.as_ref(),
|
||||
)
|
||||
.await;
|
||||
// Live-state (lot LS3) : distiller l'intent AVANT que `task` ne soit déplacé
|
||||
// dans le `Ticket`. Posé sur la cible après l'enqueue (délégation acceptée).
|
||||
let working_intent = Self::distill_intent(&task);
|
||||
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);
|
||||
// Auto-update live-state (lot LS3), best-effort : la cible passe `Working` sur
|
||||
// cette transition d'`ask` acceptée. N'altère jamais le succès de la délégation.
|
||||
self.mark_target_working_best_effort(&project.root, agent_id, ticket_id, working_intent)
|
||||
.await;
|
||||
// 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 par la **fenêtre d'inactivité réarmable**
|
||||
// (signe de vie = sonde de transcript) au lieu d'un timeout plat : un long
|
||||
// tour unique qui progresse n'est plus coupé à `turn_timeout`. Vrai silence
|
||||
// ⇒ timeout typé (comme avant) ; plafond atteint malgré progrès ⇒ erreur
|
||||
// typée distincte. Canal fermé ⇒ le garde retire le ticket au Drop.
|
||||
let verdict = self
|
||||
.run_ask_with_watchdog(pending, turn_timeout, &project.root, agent_id, &target, started)
|
||||
.await;
|
||||
match verdict {
|
||||
WatchdogOutcome::Resolved(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;
|
||||
// Auto-update live-state (lot LS3), best-effort : la cible vient de rendre
|
||||
// son résultat ⇒ `Done` + `last_delegation` = ticket résolu.
|
||||
self.mark_target_done_best_effort(&project.root, agent_id, ticket_id)
|
||||
.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. On repasse
|
||||
// aussi la live-state de la cible à `Done` (best-effort) : le tour est conclu
|
||||
// sans réponse, donc `idea_workstate_read` ne doit pas la voir `Working`
|
||||
// (busy fantôme) comme sur la branche succès.
|
||||
WatchdogOutcome::Resolved(Ok(TurnResolution::ReturnedToPromptNoReply)) => {
|
||||
busy_guard.disarm();
|
||||
self.mark_target_done_best_effort(&project.root, agent_id, ticket_id)
|
||||
.await;
|
||||
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))
|
||||
}
|
||||
// Plafond absolu atteint alors que la cible **progresse encore** (sonde de vie
|
||||
// croissante) : verdict DISTINCT, non un faux timeout muet. Le garde fait
|
||||
// `cancel_head` + `mark_idle` au Drop ; on réconcilie aussi la live-state à
|
||||
// `Done` pour qu'un abandon ne laisse pas un busy fantôme.
|
||||
WatchdogOutcome::CeilingActive => {
|
||||
self.mark_target_done_best_effort(&project.root, agent_id, ticket_id)
|
||||
.await;
|
||||
crate::diag!(
|
||||
"[rendezvous] ask CEILING (still active): target={target} \
|
||||
(agent {agent_id}) ticket={ticket_id} after_ms={}",
|
||||
started.elapsed().as_millis(),
|
||||
);
|
||||
Err(AppError::TargetCeilingActive(target))
|
||||
}
|
||||
// Erreur / timeout : on laisse le garde faire `cancel_head` + `mark_idle` au
|
||||
// Drop (retrait des `cancel_head` redondants — `cancel_head` reste idempotent).
|
||||
WatchdogOutcome::Resolved(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"
|
||||
)))
|
||||
}
|
||||
WatchdogOutcome::NoReply => {
|
||||
// Vrai silence (aucun progrès sur une fenêtre) : sémantique identique à
|
||||
// l'ancien timeout plat. On réconcilie la live-state à `Done` pour qu'un
|
||||
// abandon ne laisse pas la cible `Working` (busy fantôme) ; le garde fait
|
||||
// `cancel_head` + `mark_idle` au Drop.
|
||||
self.mark_target_done_best_effort(&project.root, agent_id, ticket_id)
|
||||
.await;
|
||||
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, et aucun progrès observé sur la fenêtre)",
|
||||
started.elapsed().as_millis(),
|
||||
);
|
||||
Err(AppError::from(domain::ports::AgentSessionError::Timeout))
|
||||
}
|
||||
}
|
||||
.await
|
||||
}
|
||||
|
||||
/// Chemin `ask` **structuré** (readiness/heartbeat lot 1) : la cible a un
|
||||
@ -1562,8 +1390,8 @@ impl OrchestratorService {
|
||||
/// 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`).
|
||||
/// busy et l'exclusion mutuelle, mais il n'est plus une source de réponse : le ticket
|
||||
/// est retiré quand le [`domain::ports::ReplyEvent::Final`] structured a été reçu.
|
||||
/// Le checkpoint Prompt/Response best-effort est conservé à l'identique du chemin PTY.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
async fn ask_structured(
|
||||
@ -1601,7 +1429,8 @@ impl OrchestratorService {
|
||||
)
|
||||
.await;
|
||||
|
||||
// Ticket dans la FIFO : comptabilité busy + `idea_reply` comme signal alternatif.
|
||||
// Ticket dans la FIFO : comptabilité busy + exclusion mutuelle. La réponse ne
|
||||
// viendra plus de cette mailbox, seulement du `Final` structured/headless.
|
||||
let requester_label = self.requester_label(project, requester).await;
|
||||
let ticket_id = TicketId::new_random();
|
||||
let ticket = match requester {
|
||||
@ -1617,7 +1446,7 @@ impl OrchestratorService {
|
||||
// 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);
|
||||
let _pending = input.enqueue_silent(agent_id, ticket);
|
||||
// Auto-update live-state (lot LS3), best-effort : la cible passe `Working` sur
|
||||
// cette transition d'`ask` acceptée (chemin structuré). `task` est encore vivant
|
||||
// ici (utilisé par le drain plus bas), on le distille directement.
|
||||
@ -1636,8 +1465,8 @@ impl OrchestratorService {
|
||||
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.
|
||||
// PTY. La cible n'a pas de PTY ; le tour se débloque uniquement sur le `Final`
|
||||
// de sa session structured/headless.
|
||||
let started = Instant::now();
|
||||
crate::diag!(
|
||||
"[rendezvous] ask started (structured): requester={} -> target={target} \
|
||||
@ -1651,35 +1480,12 @@ impl OrchestratorService {
|
||||
// est **non borné** (`None`) : la borne de tour est désormais portée par la
|
||||
// **fenêtre d'inactivité réarmable** autour de l'attente (signe de vie = sonde de
|
||||
// transcript), pas par un timeout plat interne — un long tour unique qui progresse
|
||||
// n'est plus coupé à `turn_timeout`. On attend la **première** issue : le tour
|
||||
// structuré OU un `idea_reply` explicite.
|
||||
// n'est plus coupé à `turn_timeout`. Aucun `idea_reply` ne peut résoudre ce tour.
|
||||
let drain = drain_with_readiness(session, &task, None, input.as_ref(), agent_id);
|
||||
|
||||
// L'attente du rendez-vous (drain OU reply), rendue comme `Result<String, AppError>`
|
||||
// pour être enveloppée par le watchdog (au lieu de `return` directs dans le `select!`).
|
||||
let wait = async {
|
||||
tokio::select! {
|
||||
biased;
|
||||
drained = drain => match drained {
|
||||
Ok(content) => Ok(content),
|
||||
Err(err) => Err(AppError::from(err)),
|
||||
},
|
||||
replied = pending => match replied {
|
||||
Ok(TurnResolution::Replied(content)) => {
|
||||
// La session draine encore en arrière-plan ; on fait avancer la FIFO.
|
||||
input.mark_idle(agent_id);
|
||||
Ok(content)
|
||||
}
|
||||
Ok(TurnResolution::ReturnedToPromptNoReply) => {
|
||||
input.mark_idle(agent_id);
|
||||
Err(AppError::TargetReturnedNoReply(target.to_owned()))
|
||||
}
|
||||
Err(_cancelled) => Err(AppError::Process(format!(
|
||||
"agent {target} : canal de réponse fermé avant un résultat"
|
||||
))),
|
||||
},
|
||||
}
|
||||
};
|
||||
// L'attente du rendez-vous structured, rendue comme `Result<String, AppError>`
|
||||
// pour être enveloppée par le watchdog.
|
||||
let wait = async { drain.await.map_err(AppError::from) };
|
||||
|
||||
// Borne par la fenêtre d'inactivité (réarmée sur signe de vie) sous plafond absolu.
|
||||
let result = match self
|
||||
@ -1733,8 +1539,10 @@ impl OrchestratorService {
|
||||
}
|
||||
};
|
||||
|
||||
// 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`.
|
||||
// Succès : le `Final` a rendu la réponse. On retire explicitement le ticket de
|
||||
// comptabilité (aucun `idea_reply` ne le fera), puis on désarme le garde RAII.
|
||||
mailbox.cancel_head(agent_id, ticket_id);
|
||||
input.mark_idle(agent_id);
|
||||
busy_guard.disarm();
|
||||
|
||||
// Checkpoint Response (best-effort), AVANT de déplacer `result`.
|
||||
@ -2107,15 +1915,14 @@ impl OrchestratorService {
|
||||
/// - `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`.
|
||||
/// ou profil introuvable) ⇒ `AskAgent` refuse la conversation inter-agent.
|
||||
///
|
||||
/// 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.
|
||||
/// L'auto-lancement réutilise le **même** [`LaunchAgent`] que les lancements UI :
|
||||
/// 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 runtime MCP peut encore être matérialisé
|
||||
/// pour les outils non conversationnels, mais il ne participe plus à la résolution
|
||||
/// de la réponse inter-agent.
|
||||
async fn ensure_structured_session(
|
||||
&self,
|
||||
project: &Project,
|
||||
@ -2129,7 +1936,8 @@ impl OrchestratorService {
|
||||
}
|
||||
|
||||
// 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.
|
||||
// structuré. Sinon (agent PTY/TUI legacy, ou profil introuvable) ⇒ refus par
|
||||
// l'appelant.
|
||||
let is_structured = self
|
||||
.profiles
|
||||
.list()
|
||||
@ -2141,9 +1949,26 @@ impl OrchestratorService {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
// Une cible peut déjà être vivante dans le registre PTY parce qu'elle a été
|
||||
// ouverte depuis la surface humaine historique (cellule/menu), alors que le
|
||||
// chemin inter-agent actuel exige une session `AgentSession` headless. Si on
|
||||
// laisse ce PTY en place, `LaunchAgent` applique correctement l'invariant
|
||||
// « 1 session vivante/agent » et rend le PTY existant, donc aucune session
|
||||
// structurée n'est insérée et l'ask échoue avec « aucune session structurée ».
|
||||
//
|
||||
// Le rendez-vous inter-agent est propriétaire du canal headless : on retire
|
||||
// d'abord l'éventuelle session PTY de la cible, puis on relance via le launcher
|
||||
// structuré partagé. Le node hôte est conservé best-effort pour que la surface
|
||||
// puisse se rattacher au même emplacement si elle observe l'événement de relance.
|
||||
let previous_node = self.sessions.node_for_agent(&agent_id);
|
||||
if let Some(session_id) = self.sessions.session_for_agent(&agent_id) {
|
||||
self.close_terminal
|
||||
.execute(CloseTerminalInput { session_id })
|
||||
.await?;
|
||||
}
|
||||
|
||||
// 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
|
||||
// insère dans CE registre). `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
|
||||
@ -2152,7 +1977,7 @@ impl OrchestratorService {
|
||||
agent_id,
|
||||
rows: DEFAULT_ROWS,
|
||||
cols: DEFAULT_COLS,
|
||||
node_id: None,
|
||||
node_id: previous_node,
|
||||
conversation_id: None,
|
||||
mcp_runtime: self
|
||||
.mcp_runtime_provider
|
||||
@ -2386,55 +2211,6 @@ impl OrchestratorService {
|
||||
.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 **submit config**
|
||||
/// (`submit_sequence`/`submit_delay_ms`, ARCHITECTURE §20.3) **and** whether it
|
||||
/// declares an **MCP bridge**, in a single profile lookup. The submit config is
|
||||
@ -2623,6 +2399,7 @@ impl OrchestratorService {
|
||||
/// service outlives the guard.
|
||||
struct WaitEdgeGuard<'a> {
|
||||
graph: &'a StdMutex<WaitForGraph>,
|
||||
active_waits: &'a StdMutex<Vec<(AgentId, AgentId)>>,
|
||||
from: AgentId,
|
||||
to: AgentId,
|
||||
}
|
||||
@ -2636,8 +2413,18 @@ impl<'a> WaitEdgeGuard<'a> {
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
g.add_edge(from, to);
|
||||
}
|
||||
{
|
||||
let mut waits = service
|
||||
.active_waits
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
if !waits.contains(&(from, to)) {
|
||||
waits.push((from, to));
|
||||
}
|
||||
}
|
||||
Self {
|
||||
graph: &service.wait_for,
|
||||
active_waits: &service.active_waits,
|
||||
from,
|
||||
to,
|
||||
}
|
||||
@ -2651,6 +2438,11 @@ impl Drop for WaitEdgeGuard<'_> {
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
g.remove_edge(self.from, self.to);
|
||||
let mut waits = self
|
||||
.active_waits
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
waits.retain(|edge| *edge != (self.from, self.to));
|
||||
}
|
||||
}
|
||||
|
||||
@ -2737,6 +2529,20 @@ mod tests {
|
||||
assert_eq!(submit.delay_ms, Some(900));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn active_wait_dependencies_are_transitive_and_deduplicated() {
|
||||
let edges = vec![
|
||||
(aid(1), aid(2)),
|
||||
(aid(2), aid(3)),
|
||||
(aid(1), aid(3)),
|
||||
(aid(9), aid(10)),
|
||||
];
|
||||
|
||||
let deps = OrchestratorService::active_wait_dependencies_from_edges(&edges, aid(1));
|
||||
|
||||
assert_eq!(deps, vec![aid(3), aid(2)]);
|
||||
}
|
||||
|
||||
// --- BusyTurnGuard (RAII de fin de tour) -------------------------------
|
||||
//
|
||||
// Fakes minimaux pour observer ce que le garde appelle à son Drop : un médiateur
|
||||
|
||||
Reference in New Issue
Block a user