fix(orchestrator): garde RAII contre une cible coincée Busy après une délégation interrompue
Une délégation `idea_ask_agent` interrompue/annulée côté demandeur (futur `ask_agent` *dropped*) laissait la cible en `Busy` à vie : sur un drop, aucune branche du `select!` ne s'exécute, donc `mark_idle` n'était jamais appelé. Toute délégation suivante restait en file derrière ce tour fantôme et n'était jamais livrée au PTY (symptôme : « DevFrontend/QA ne répondent plus », pont MCP pourtant ESTAB). Ajoute un garde RAII `BusyTurnGuard` posé juste après l'enqueue sur les deux chemins (`ask` PTY et `ask_structured`) : son `Drop` fait `cancel_head` + `mark_idle` quoi qu'il arrive (erreur, timeout, drop), sauf `disarm()` sur la branche succès. Indispensable d'être un garde et non un `mark_idle` dans les branches : le cas réel est un futur droppé. Retire les `cancel_head` redondants des branches erreur/timeout (positionnel/idempotent). `sweep_stalled` reste advisory (n'appelle jamais `mark_idle`) — non concerné. Tests: dropped_ask_future_frees_busy_target, second_delegation_delivered_after_dropped_ask, cancelled_ask_marks_target_idle + 2 tests du garde. cargo test --workspace vert (80 suites). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@ -94,6 +94,68 @@ 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)))
|
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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Fournit les faits OS/runtime (exe + endpoint projet) pour écrire la déclaration MCP
|
/// 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
|
/// 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).
|
/// app-tauri (seul détenteur de current_exe/$APPIMAGE/mcp_endpoint).
|
||||||
@ -827,12 +889,23 @@ impl OrchestratorService {
|
|||||||
|
|
||||||
// 1. Garantir la cible vivante en PTY pour CE fil ; lier sa session à la
|
// 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).
|
// conversation, et brancher son handle d'entrée sur le médiateur (livraison).
|
||||||
let handle = self
|
let (handle, cold_launch) = self
|
||||||
.ensure_live_pty(project, agent_id, conversation_id, &target)
|
.ensure_live_pty(project, agent_id, conversation_id, &target)
|
||||||
.await?;
|
.await?;
|
||||||
// Arm prompt-ready detection (C5) with the target profile's literal marker, so a
|
// 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`).
|
// return-to-prompt frees the turn (the other OR signal being `idea_reply`).
|
||||||
let (prompt_pattern, submit) = self.prompt_and_submit_for_agent(project, agent_id).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.clone(), prompt_pattern, submit);
|
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
|
// 2. Enregistrer le ticket (slot de réponse) + livrer le tour via le médiateur
|
||||||
@ -866,12 +939,22 @@ impl OrchestratorService {
|
|||||||
// le médiateur AVANT l'enqueue (consommé au start_turn).
|
// le médiateur AVANT l'enqueue (consommé au start_turn).
|
||||||
let turn_timeout = self.turn_timeout_for(project, agent_id).await;
|
let turn_timeout = self.turn_timeout_for(project, agent_id).await;
|
||||||
let pending = input.enqueue(agent_id, ticket);
|
let pending = input.enqueue(agent_id, ticket);
|
||||||
|
// 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
|
// Delivery is the mediator's responsibility (`InputMediator::enqueue` writes the
|
||||||
// turn into the bound handle). The service no longer writes the PTY directly —
|
// 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).
|
// 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é ⇒ retirer le ticket
|
// 3. Attendre la réponse, bornée. Timeout/canal fermé ⇒ le garde retire le ticket
|
||||||
// (cible laissée vivante) et renvoyer une erreur typée.
|
// (cible laissée vivante) et la ramène `Idle` au Drop ; renvoie une erreur typée.
|
||||||
match tokio::time::timeout(turn_timeout, pending).await {
|
match tokio::time::timeout(turn_timeout, pending).await {
|
||||||
Ok(Ok(result)) => {
|
Ok(Ok(result)) => {
|
||||||
// Checkpoint Response (P6b, best-effort) : persister la réponse AVANT de
|
// Checkpoint Response (P6b, best-effort) : persister la réponse AVANT de
|
||||||
@ -885,18 +968,19 @@ impl OrchestratorService {
|
|||||||
result.clone(),
|
result.clone(),
|
||||||
)
|
)
|
||||||
.await;
|
.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();
|
||||||
Ok(self.reply_outcome(agent_id, &target, result))
|
Ok(self.reply_outcome(agent_id, &target, result))
|
||||||
}
|
}
|
||||||
Ok(Err(_cancelled)) => {
|
// Erreur / timeout : on laisse le garde faire `cancel_head` + `mark_idle` au
|
||||||
mailbox.cancel_head(agent_id, ticket_id);
|
// Drop (retrait des `cancel_head` redondants — `cancel_head` reste idempotent).
|
||||||
Err(AppError::Process(format!(
|
Ok(Err(_cancelled)) => Err(AppError::Process(format!(
|
||||||
"agent {target} : canal de réponse fermé avant un résultat"
|
"agent {target} : canal de réponse fermé avant un résultat"
|
||||||
)))
|
))),
|
||||||
}
|
Err(_elapsed) => Err(AppError::from(domain::ports::AgentSessionError::Timeout)),
|
||||||
Err(_elapsed) => {
|
|
||||||
mailbox.cancel_head(agent_id, ticket_id);
|
|
||||||
Err(AppError::from(domain::ports::AgentSessionError::Timeout))
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -925,11 +1009,13 @@ impl OrchestratorService {
|
|||||||
) -> Result<OrchestratorOutcome, AppError> {
|
) -> Result<OrchestratorOutcome, AppError> {
|
||||||
let (input, mailbox) = match (&self.input, &self.mailbox) {
|
let (input, mailbox) = match (&self.input, &self.mailbox) {
|
||||||
(Some(i), Some(m)) => (i, m),
|
(Some(i), Some(m)) => (i, m),
|
||||||
_ => return Err(AppError::Invalid(
|
_ => {
|
||||||
|
return Err(AppError::Invalid(
|
||||||
"la messagerie inter-agents (idea_ask_agent) n'est pas disponible : \
|
"la messagerie inter-agents (idea_ask_agent) n'est pas disponible : \
|
||||||
médiateur d'entrée non câblé"
|
médiateur d'entrée non câblé"
|
||||||
.to_owned(),
|
.to_owned(),
|
||||||
)),
|
))
|
||||||
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
// Checkpoint Prompt (best-effort), AVANT de déplacer `task` dans le ticket.
|
// Checkpoint Prompt (best-effort), AVANT de déplacer `task` dans le ticket.
|
||||||
@ -963,27 +1049,30 @@ impl OrchestratorService {
|
|||||||
// l'enqueue qui démarre le tour (le médiateur arme alors sa fenêtre de vivacité).
|
// 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 turn_timeout = self.turn_timeout_for(project, agent_id).await;
|
||||||
let pending = input.enqueue(agent_id, ticket);
|
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,
|
||||||
|
);
|
||||||
|
|
||||||
// Drainer le tour structuré (le `Final` ⇒ contenu + `mark_idle`), borné par le
|
// 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
|
// **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.
|
// **première** issue : le tour structuré OU un `idea_reply` explicite.
|
||||||
let drain = drain_with_readiness(
|
let drain =
|
||||||
session,
|
drain_with_readiness(session, &task, Some(turn_timeout), input.as_ref(), agent_id);
|
||||||
&task,
|
|
||||||
Some(turn_timeout),
|
|
||||||
input.as_ref(),
|
|
||||||
agent_id,
|
|
||||||
);
|
|
||||||
|
|
||||||
let result = tokio::select! {
|
let result = tokio::select! {
|
||||||
biased;
|
biased;
|
||||||
// Issue déterministe : la session a rendu son `Final`.
|
// Issue déterministe : la session a rendu son `Final`.
|
||||||
drained = drain => match drained {
|
drained = drain => match drained {
|
||||||
Ok(content) => content,
|
Ok(content) => content,
|
||||||
Err(err) => {
|
// Erreur de drain : le garde fait `cancel_head` + `mark_idle` au Drop.
|
||||||
mailbox.cancel_head(agent_id, ticket_id);
|
Err(err) => return Err(AppError::from(err)),
|
||||||
return Err(AppError::from(err));
|
|
||||||
}
|
|
||||||
},
|
},
|
||||||
// Issue alternative : un `idea_reply` explicite a résolu le ticket d'abord.
|
// Issue alternative : un `idea_reply` explicite a résolu le ticket d'abord.
|
||||||
replied = pending => match replied {
|
replied = pending => match replied {
|
||||||
@ -993,8 +1082,8 @@ impl OrchestratorService {
|
|||||||
input.mark_idle(agent_id);
|
input.mark_idle(agent_id);
|
||||||
content
|
content
|
||||||
}
|
}
|
||||||
|
// Canal fermé : le garde fait `cancel_head` + `mark_idle` au Drop.
|
||||||
Err(_cancelled) => {
|
Err(_cancelled) => {
|
||||||
mailbox.cancel_head(agent_id, ticket_id);
|
|
||||||
return Err(AppError::Process(format!(
|
return Err(AppError::Process(format!(
|
||||||
"agent {target} : canal de réponse fermé avant un résultat"
|
"agent {target} : canal de réponse fermé avant un résultat"
|
||||||
)));
|
)));
|
||||||
@ -1002,6 +1091,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`.
|
||||||
|
busy_guard.disarm();
|
||||||
|
|
||||||
// Checkpoint Response (best-effort), AVANT de déplacer `result`.
|
// Checkpoint Response (best-effort), AVANT de déplacer `result`.
|
||||||
self.record_turn_best_effort(
|
self.record_turn_best_effort(
|
||||||
&project.root,
|
&project.root,
|
||||||
@ -1058,10 +1151,21 @@ impl OrchestratorService {
|
|||||||
|
|
||||||
// Ensure the target is live for this thread and bind its input handle on the
|
// Ensure the target is live for this thread and bind its input handle on the
|
||||||
// mediator (delivery path). Same call the ask path uses.
|
// mediator (delivery path). Same call the ask path uses.
|
||||||
let handle = self
|
let (handle, cold_launch) = self
|
||||||
.ensure_live_pty(project, agent_id, conversation_id, target)
|
.ensure_live_pty(project, agent_id, conversation_id, target)
|
||||||
.await?;
|
.await?;
|
||||||
let (prompt_pattern, submit) = self.prompt_and_submit_for_agent(project, agent_id).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);
|
input.bind_handle_with_prompt(agent_id, handle, prompt_pattern, submit);
|
||||||
|
|
||||||
// Enqueue a human-sourced ticket in the SAME FIFO as delegations. Fire-and-
|
// Enqueue a human-sourced ticket in the SAME FIFO as delegations. Fire-and-
|
||||||
@ -1145,6 +1249,26 @@ impl OrchestratorService {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 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) {
|
||||||
|
if let Some(input) = &self.input {
|
||||||
|
input.set_front_attached(agent, attached);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Resolves the conversation thread id for an ask: `A↔B` when an agent requests,
|
/// 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
|
/// 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).
|
/// stable per-agent id derived from the target (legacy routing — never panics).
|
||||||
@ -1219,13 +1343,19 @@ impl OrchestratorService {
|
|||||||
/// a launch the handle is resolved from the registry; a missing handle is a
|
/// 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
|
/// [`AppError::Process`] (the launch did not register a PTY session, e.g. a profile
|
||||||
/// IdeA cannot drive as a terminal).
|
/// 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(
|
async fn ensure_live_pty(
|
||||||
&self,
|
&self,
|
||||||
project: &Project,
|
project: &Project,
|
||||||
agent_id: AgentId,
|
agent_id: AgentId,
|
||||||
conversation_id: domain::conversation::ConversationId,
|
conversation_id: domain::conversation::ConversationId,
|
||||||
target: &str,
|
target: &str,
|
||||||
) -> Result<PtyHandle, AppError> {
|
) -> Result<(PtyHandle, bool), AppError> {
|
||||||
// «1 session vivante / conversation» (cadrage C3 §5.2) : on cherche d'abord la
|
// «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
|
// 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).
|
// agent mono-fil dont la session n'a pas encore été liée à sa conversation).
|
||||||
@ -1235,9 +1365,10 @@ impl OrchestratorService {
|
|||||||
.or_else(|| self.sessions.session_for_agent(&agent_id));
|
.or_else(|| self.sessions.session_for_agent(&agent_id));
|
||||||
if let Some(session_id) = existing {
|
if let Some(session_id) = existing {
|
||||||
if let Some(handle) = self.sessions.handle(&session_id) {
|
if let Some(handle) = self.sessions.handle(&session_id) {
|
||||||
// (Re)lier le fil à cette session vivante (idempotent).
|
// (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);
|
self.bind_conversation_session(conversation_id, session_id);
|
||||||
return Ok(handle);
|
return Ok((handle, false));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -1269,11 +1400,13 @@ impl OrchestratorService {
|
|||||||
// Lier la session fraîchement lancée à CE fil (registre terminal + registre de
|
// 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.
|
// conversations) ⇒ un prochain ask sur le même fil la réutilise.
|
||||||
self.bind_conversation_session(conversation_id, session_id);
|
self.bind_conversation_session(conversation_id, session_id);
|
||||||
self.sessions.handle(&session_id).ok_or_else(|| {
|
let handle = self.sessions.handle(&session_id).ok_or_else(|| {
|
||||||
AppError::Process(format!(
|
AppError::Process(format!(
|
||||||
"handle PTY de l'agent {target} introuvable après lancement"
|
"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`.
|
/// Lot 1b — garantit une **session structurée vivante** pour la cible d'un `ask`.
|
||||||
@ -1337,7 +1470,10 @@ impl OrchestratorService {
|
|||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
// Le launcher a inséré la session dans le registre partagé : la relire.
|
// Le launcher a inséré la session dans le registre partagé : la relire.
|
||||||
structured.session_for_agent(&agent_id).map(Some).ok_or_else(|| {
|
structured
|
||||||
|
.session_for_agent(&agent_id)
|
||||||
|
.map(Some)
|
||||||
|
.ok_or_else(|| {
|
||||||
AppError::Process(format!(
|
AppError::Process(format!(
|
||||||
"agent {agent_id} : aucune session structurée vivante après lancement"
|
"agent {agent_id} : aucune session structurée vivante après lancement"
|
||||||
))
|
))
|
||||||
@ -1616,7 +1752,7 @@ impl OrchestratorService {
|
|||||||
&self,
|
&self,
|
||||||
project: &Project,
|
project: &Project,
|
||||||
agent_id: AgentId,
|
agent_id: AgentId,
|
||||||
) -> (Option<String>, SubmitConfig) {
|
) -> (Option<String>, SubmitConfig, bool) {
|
||||||
let Some(agent) = self
|
let Some(agent) = self
|
||||||
.list_agents
|
.list_agents
|
||||||
.execute(ListAgentsInput {
|
.execute(ListAgentsInput {
|
||||||
@ -1626,7 +1762,7 @@ impl OrchestratorService {
|
|||||||
.ok()
|
.ok()
|
||||||
.and_then(|out| out.agents.into_iter().find(|a| a.id == agent_id))
|
.and_then(|out| out.agents.into_iter().find(|a| a.id == agent_id))
|
||||||
else {
|
else {
|
||||||
return (None, SubmitConfig::default());
|
return (None, SubmitConfig::default(), false);
|
||||||
};
|
};
|
||||||
let Some(profile) = self
|
let Some(profile) = self
|
||||||
.profiles
|
.profiles
|
||||||
@ -1635,10 +1771,13 @@ impl OrchestratorService {
|
|||||||
.ok()
|
.ok()
|
||||||
.and_then(|ps| ps.into_iter().find(|p| p.id == agent.profile_id))
|
.and_then(|ps| ps.into_iter().find(|p| p.id == agent.profile_id))
|
||||||
else {
|
else {
|
||||||
return (None, SubmitConfig::default());
|
return (None, SubmitConfig::default(), false);
|
||||||
};
|
};
|
||||||
let submit = SubmitConfig::new(profile.submit_sequence, profile.submit_delay_ms);
|
let submit = SubmitConfig::new(profile.submit_sequence, profile.submit_delay_ms);
|
||||||
(profile.prompt_ready_pattern, submit)
|
// 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) :
|
/// Résout la [`domain::profile::LivenessStrategy`] du profil de la cible (lot 2) :
|
||||||
@ -1809,4 +1948,109 @@ mod tests {
|
|||||||
assert!(by_name);
|
assert!(by_name);
|
||||||
assert_eq!(p.id.to_string(), p.id.to_string());
|
assert_eq!(p.id.to_string(), p.id.to_string());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --- 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"
|
||||||
|
);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@ -1255,6 +1255,124 @@ async fn idea_reply_marks_emitter_idle() {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Régression (garde RAII de fin de tour) — LE test décisif : un `ask_agent` dont le
|
||||||
|
/// **futur est abandonné (drop)** alors qu'il attendait la réponse doit laisser la cible
|
||||||
|
/// `Idle`, **pas** `Busy` à vie. Avant le fix, aucune branche du `match`/`select!` ne
|
||||||
|
/// s'exécutait sur un drop ⇒ l'agent restait `Busy` et toute délégation suivante était
|
||||||
|
/// mise en file derrière ce tour fantôme sans jamais être livrée.
|
||||||
|
#[tokio::test]
|
||||||
|
async fn dropped_ask_future_frees_busy_target() {
|
||||||
|
let agent = scratch_agent(aid(1), "architect", "agents/architect.md");
|
||||||
|
let fx = ask_fixture(FakeContexts::with_agent(&agent, "# persona"));
|
||||||
|
seed_live_pty(&fx.sessions, aid(1), sid(800));
|
||||||
|
|
||||||
|
let svc = Arc::clone(&fx.service);
|
||||||
|
let ask = tokio::spawn(async move { svc.dispatch(&project(), cmd(ASK_JSON)).await });
|
||||||
|
// Le tour a démarré : la cible est Busy, un ticket est en file.
|
||||||
|
await_until(|| fx.mailbox.pending(&aid(1)) == 1).await;
|
||||||
|
assert!(
|
||||||
|
fx.mediator.busy_state(aid(1)).is_busy(),
|
||||||
|
"cible Busy pendant le tour délégué"
|
||||||
|
);
|
||||||
|
|
||||||
|
// Le demandeur interrompt son appel ⇒ le futur `ask_agent` est dropped.
|
||||||
|
ask.abort();
|
||||||
|
let _ = ask.await; // récolte la JoinError(Cancelled), ignorée.
|
||||||
|
|
||||||
|
// Le garde RAII a ramené la cible Idle ET retiré le ticket fantôme de la FIFO.
|
||||||
|
await_until(|| !fx.mediator.busy_state(aid(1)).is_busy()).await;
|
||||||
|
assert_eq!(
|
||||||
|
fx.mediator.busy_state(aid(1)),
|
||||||
|
AgentBusyState::Idle,
|
||||||
|
"futur dropped ⇒ la cible est ramenée Idle par le garde (fix cause racine)"
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
fx.mailbox.pending(&aid(1)),
|
||||||
|
0,
|
||||||
|
"ticket fantôme retiré de la FIFO au Drop du garde"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Régression (garde RAII) — une **2e délégation** vers une cible dont un 1er `ask` a
|
||||||
|
/// été abandonné est bien livrée : l'agent n'est pas coincé `Busy`, son tour suivant
|
||||||
|
/// démarre et peut être résolu par `idea_reply`.
|
||||||
|
#[tokio::test]
|
||||||
|
async fn second_delegation_delivered_after_dropped_ask() {
|
||||||
|
let agent = scratch_agent(aid(1), "architect", "agents/architect.md");
|
||||||
|
let fx = ask_fixture(FakeContexts::with_agent(&agent, "# persona"));
|
||||||
|
seed_live_pty(&fx.sessions, aid(1), sid(800));
|
||||||
|
|
||||||
|
// 1er ask : démarre puis est abandonné (drop).
|
||||||
|
let svc = Arc::clone(&fx.service);
|
||||||
|
let ask1 = tokio::spawn(async move { svc.dispatch(&project(), cmd(ASK_JSON)).await });
|
||||||
|
await_until(|| fx.mailbox.pending(&aid(1)) == 1).await;
|
||||||
|
ask1.abort();
|
||||||
|
let _ = ask1.await;
|
||||||
|
await_until(|| fx.mailbox.pending(&aid(1)) == 0).await;
|
||||||
|
|
||||||
|
// 2e ask vers la MÊME cible : doit démarrer un nouveau tour (pas coincé derrière un
|
||||||
|
// tour fantôme) et se laisser résoudre.
|
||||||
|
let svc2 = Arc::clone(&fx.service);
|
||||||
|
let ask2 = tokio::spawn(async move { svc2.dispatch(&project(), cmd(ASK_JSON)).await });
|
||||||
|
await_until(|| fx.mailbox.pending(&aid(1)) == 1).await;
|
||||||
|
assert!(
|
||||||
|
fx.mediator.busy_state(aid(1)).is_busy(),
|
||||||
|
"le 2e tour démarre bien (cible Busy) — preuve qu'elle n'était pas coincée"
|
||||||
|
);
|
||||||
|
|
||||||
|
fx.service
|
||||||
|
.dispatch(&project(), reply_cmd(aid(1), "réponse au 2e tour"))
|
||||||
|
.await
|
||||||
|
.expect("reply ok");
|
||||||
|
let out = timeout(TEST_GUARD, ask2)
|
||||||
|
.await
|
||||||
|
.expect("le 2e ask se termine")
|
||||||
|
.expect("join ok")
|
||||||
|
.expect("ask ok");
|
||||||
|
assert_eq!(out.reply.as_deref(), Some("réponse au 2e tour"));
|
||||||
|
assert_eq!(
|
||||||
|
fx.mediator.busy_state(aid(1)),
|
||||||
|
AgentBusyState::Idle,
|
||||||
|
"cible Idle après résolution du 2e tour"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Régression (garde RAII) — la **branche erreur/timeout** (canal fermé) ramène aussi la
|
||||||
|
/// cible `Idle`. Avant le fix, ces branches faisaient `cancel_head` mais jamais
|
||||||
|
/// `mark_idle` ⇒ l'agent restait `Busy`. On déclenche la fermeture du canal en retirant
|
||||||
|
/// le ticket de tête (même nettoyage que le timeout de tour).
|
||||||
|
#[tokio::test]
|
||||||
|
async fn cancelled_ask_marks_target_idle() {
|
||||||
|
let agent = scratch_agent(aid(1), "architect", "agents/architect.md");
|
||||||
|
let fx = ask_fixture(FakeContexts::with_agent(&agent, "# persona"));
|
||||||
|
seed_live_pty(&fx.sessions, aid(1), sid(800));
|
||||||
|
|
||||||
|
let svc = Arc::clone(&fx.service);
|
||||||
|
let ask = tokio::spawn(async move { svc.dispatch(&project(), cmd(ASK_JSON)).await });
|
||||||
|
await_until(|| fx.mailbox.pending(&aid(1)) == 1).await;
|
||||||
|
assert!(fx.mediator.busy_state(aid(1)).is_busy());
|
||||||
|
|
||||||
|
// Retire le ticket de tête (drop du sender) ⇒ l'ask voit un canal fermé et part en
|
||||||
|
// erreur typée (PROCESS), le MÊME nettoyage que le timeout de tour.
|
||||||
|
let t = fx.mailbox.ticket_ids(&aid(1))[0];
|
||||||
|
fx.mailbox.cancel_head(aid(1), t);
|
||||||
|
let err = timeout(TEST_GUARD, ask)
|
||||||
|
.await
|
||||||
|
.expect("ask retourne vite sur canal fermé")
|
||||||
|
.expect("join ok")
|
||||||
|
.unwrap_err();
|
||||||
|
assert!(
|
||||||
|
matches!(err.code(), "PROCESS" | "INVALID"),
|
||||||
|
"ask sur canal fermé ⇒ erreur typée : {err:?}"
|
||||||
|
);
|
||||||
|
// Le garde a ramené la cible Idle (la FIFO peut avancer).
|
||||||
|
assert_eq!(
|
||||||
|
fx.mediator.busy_state(aid(1)),
|
||||||
|
AgentBusyState::Idle,
|
||||||
|
"branche erreur ⇒ cible Idle (garde RAII)"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
/// A dead target is launched in the background (PTY) before the task is written.
|
/// A dead target is launched in the background (PTY) before the task is written.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn ask_dead_target_launches_pty_then_writes_and_replies() {
|
async fn ask_dead_target_launches_pty_then_writes_and_replies() {
|
||||||
|
|||||||
Reference in New Issue
Block a user