diff --git a/crates/infrastructure/src/input/mod.rs b/crates/infrastructure/src/input/mod.rs index 8a9d40c..eaea177 100644 --- a/crates/infrastructure/src/input/mod.rs +++ b/crates/infrastructure/src/input/mod.rs @@ -57,6 +57,17 @@ struct BusyTracker { /// pour un démarrage à froid, publiée par [`BusyTracker::prompt_ready`] à l'apparition /// du prompt (jamais avant). Absent ⇒ aucun tour en attente de gate. deferred: Mutex>, + /// **Latch « déjà libéré »** (fix race cold-start, ordre inverse) : ensemble des + /// agents pour lesquels un signal de readiness (`release_cold_start` / + /// `prompt_ready`) est arrivé **avant** que l'`enqueue` n'ait parqué son tour dans + /// `deferred`. Sans ce latch, ce signal trouverait `deferred` vide, ne ferait rien + /// d'utile (l'agent encore `starting` ⇒ pas de `mark_idle`), puis l'`enqueue` + /// parquerait le tour dans `deferred` — où il resterait à jamais (Claude Code n'a pas + /// de `prompt_ready_pattern`, donc aucun second signal ne viendrait le drainer). Le + /// latch enregistre « cet agent en démarrage est déjà prêt » : l'`enqueue` qui suit + /// livre alors **immédiatement** au lieu de parquer. Consommé (retiré) à la livraison + /// ⇒ exactement-une-fois, quel que soit l'ordre. Vide ⇒ comportement inchangé. + released: Mutex>, events: Option>, /// Sink de livraison headless (cf. [`HeadlessSink`]). Câblé par /// [`MediatedInbox::with_pty`]/[`MediatedInbox::with_events`] quand un PTY est @@ -139,6 +150,7 @@ impl BusyTracker { liveness: Mutex::new(HashMap::new()), starting: Mutex::new(HashSet::new()), deferred: Mutex::new(HashMap::new()), + released: Mutex::new(HashSet::new()), events, headless_sink: None, mailbox: None, @@ -158,6 +170,12 @@ impl BusyTracker { .unwrap_or_else(std::sync::PoisonError::into_inner) } + fn lock_released(&self) -> std::sync::MutexGuard<'_, HashSet> { + self.released + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } + /// Marque un agent **en démarrage à froid** : son tout premier tour devra être /// *gaté* sur le prompt-ready watcher (la `DelegationReady` sera différée à /// l'apparition du prompt). À n'appeler **que** lorsqu'un watcher est effectivement @@ -343,8 +361,10 @@ impl BusyTracker { /// - Sinon : comportement historique ⇒ `mark_idle` (fait avancer la FIFO). fn prompt_ready(&self, agent: AgentId) { // On ne retire l'agent de `starting` qu'ici : tant que son premier tour n'a pas - // été livré, un re-arming éventuel doit rester gaté. - self.lock_starting().remove(&agent); + // été livré, un re-arming éventuel doit rester gaté. On capture `was_starting` + // pour distinguer le retour-au-prompt d'un tour chaud (fin de tour ⇒ `mark_idle` + + // grâce) du prompt-ready arrivé AVANT l'enqueue différé d'un démarrage à froid. + let was_starting = self.lock_starting().remove(&agent); let deferred = self.lock_deferred().remove(&agent); match deferred { Some(d) => { @@ -354,6 +374,17 @@ impl BusyTracker { ); self.publish_deferred(agent, d); } + // **Ordre inverse de la race** : l'agent est encore en démarrage à froid mais + // aucun tour n'est parqué ⇒ le prompt-ready est arrivé AVANT l'enqueue. On pose + // le latch « déjà libéré » et SURTOUT pas de `mark_idle`/grâce : le tour n'a pas + // encore couru, il n'y a aucun appelant à réveiller. L'`enqueue` qui suit verra + // le latch et livrera immédiatement (exactement-une-fois). + None if was_starting => { + application::diag!( + "[input-mediator] prompt_ready before deferred enqueue (cold-start) -> released latch agent={agent}" + ); + self.lock_released().insert(agent); + } None => { // Pas de tour différé : la cible est revenue à son prompt sans // `idea_reply`. On capture le ticket actif AVANT `mark_idle` (qui efface @@ -393,13 +424,25 @@ impl BusyTracker { /// pas de fin de tour). Idempotent et OR-safe avec `prompt_ready` (le `remove` ne rend /// `Some` qu'une fois). fn release_cold_start(&self, agent: AgentId) { - self.lock_starting().remove(&agent); + let was_starting = self.lock_starting().remove(&agent); if let Some(d) = self.lock_deferred().remove(&agent) { + // Cas nominal : l'`enqueue` a déjà parqué le tour ⇒ on le draine. application::diag!( "[input-mediator] release_cold_start (mcp-ready) released deferred delegation agent={agent} ticket={}", d.ticket ); self.publish_deferred(agent, d); + } else if was_starting { + // **Ordre inverse de la race** : la readiness MCP est arrivée AVANT que + // l'`enqueue` ait parqué son tour. On pose le latch « déjà libéré » : le + // prochain `enqueue` (agent encore vu comme démarrant) livrera immédiatement au + // lieu de parquer dans `deferred` — où le tour resterait à jamais (Claude Code + // n'ayant pas de prompt_ready_pattern, aucun second signal ne le drainerait). + // Hors démarrage à froid (`was_starting == false`), aucun latch : no-op. + application::diag!( + "[input-mediator] release_cold_start before deferred enqueue (cold-start) -> released latch agent={agent}" + ); + self.lock_released().insert(agent); } } @@ -939,10 +982,26 @@ impl InputMediator for MediatedInbox { submit_delay_ms: submit.delay_ms, }; if self.tracker.lock_starting().contains(&agent) { - application::diag!( - "[input-mediator] enqueue deferred (cold-start gate) agent={agent} ticket={ticket_id}" - ); - self.tracker.lock_deferred().insert(agent, deferred); + // Agent en démarrage à froid. Si un signal de readiness est DÉJÀ arrivé + // (latch « déjà libéré » posé par `release_cold_start`/`prompt_ready` + // AVANT cet enqueue — l'ordre inverse de la race), on consomme le latch, + // on retire l'agent du gate et on livre TOUT DE SUITE : sinon ce tour + // serait parqué dans `deferred` sans plus aucun signal pour le drainer + // (blocage éternel). Sinon (readiness pas encore arrivée), on parque + // normalement : le prochain signal le drainera. Exactement-une-fois dans + // les deux ordres (release-avant-enqueue ET enqueue-avant-release). + if self.tracker.lock_released().remove(&agent) { + application::diag!( + "[input-mediator] enqueue sees released latch (cold-start, reverse order) -> deliver now agent={agent} ticket={ticket_id}" + ); + self.tracker.lock_starting().remove(&agent); + self.tracker.publish_deferred(agent, deferred); + } else { + application::diag!( + "[input-mediator] enqueue deferred (cold-start gate) agent={agent} ticket={ticket_id}" + ); + self.tracker.lock_deferred().insert(agent, deferred); + } } else { self.tracker.publish_deferred(agent, deferred); } @@ -1983,6 +2042,92 @@ mod tests { ); } + /// **Régression de la race cold-start (ordre INVERSE)** : `release_cold_start` + /// (signal de connexion du pont MCP) arrive AVANT que l'`enqueue` n'ait parqué son + /// tour différé. Avant le fix, le signal trouvait `deferred` vide (no-op), puis + /// l'`enqueue` parquait le tour dans `deferred` — où il restait à JAMAIS (aucun second + /// signal ne venait le drainer ⇒ `idea_ask_agent` bloquait jusqu'au timeout, Claude + /// Code n'ayant pas de `prompt_ready_pattern`). Après le fix : exactement UNE + /// `DelegationReady` est livrée (le latch « déjà libéré » garantit la livraison). + #[test] + fn release_cold_start_before_deferred_enqueue_still_delivers_once() { + let bus = Arc::new(RecordingBus::default()); + let inbox = MediatedInbox::new(Arc::new(InMemoryMailbox::new()), Arc::new(FixedClock(1))) + .with_events(Arc::clone(&bus) as Arc); + let a = agent(1); + + // Démarrage à froid : gate armé AVANT l'enqueue (ordre de l'orchestrateur). + inbox.mark_starting(a); + // **La race** : la readiness MCP arrive AVANT l'enqueue qui parquera le tour. + inbox.release_cold_start(a); + assert!( + bus.delegation_ready().is_empty(), + "rien à livrer tant que l'enqueue n'a pas fourni le tour" + ); + + // Puis l'enqueue fournit le tour : grâce au latch « déjà libéré », il livre + // immédiatement au lieu de le perdre dans `deferred`. + inbox.enqueue(a, ticket(10, "cold task")); + let ready = bus.delegation_ready(); + assert_eq!( + ready.len(), + 1, + "release AVANT enqueue ⇒ exactement une DelegationReady (jamais zéro, jamais deux)" + ); + assert!( + ready[0].1.ends_with("\ncold task"), + "le texte livré porte la tâche brute: {:?}", + ready[0].1 + ); + assert!( + inbox.busy_state(a).is_busy(), + "le tour froid livré ⇒ l'agent reste Busy (idle viendra d'idea_reply)" + ); + } + + /// Variante prompt-ready de la race ordre INVERSE : le watcher signale `prompt_ready` + /// AVANT l'enqueue ⇒ pas de `mark_idle` prématuré, et le tour est livré une seule fois + /// par l'enqueue qui suit (via le latch « déjà libéré »). + #[test] + fn prompt_ready_before_deferred_enqueue_still_delivers_once() { + let pty = Arc::new(FakePty::new()); + let bus = Arc::new(RecordingBus::default()); + let a = agent(1); + let h = handle(23); + // Le prompt "\n> " est présent d'emblée ⇒ le watcher déclenche tôt. + pty.seed(&h, vec![b"ready\n> ".to_vec()]); + let inbox = MediatedInbox::with_pty( + Arc::new(InMemoryMailbox::new()), + Arc::new(FixedClock(1)), + Arc::clone(&pty) as Arc, + ) + .with_events(Arc::clone(&bus) as Arc); + + // Cellule front montée ⇒ livraison observable via l'événement. + inbox.set_front_attached(a, true); + inbox.mark_starting(a); + // Le watcher s'arme et déclenche `prompt_ready` AVANT l'enqueue : il consomme son + // stream (prompt présent d'emblée), pose le latch « déjà libéré » et n'émet + // surtout PAS de `mark_idle` (le tour n'a pas encore couru). + inbox.bind_handle_with_prompt(a, h, Some("\n> ".to_owned()), SubmitConfig::default()); + // Laisse le watcher (thread détaché) consommer son flux fini et poser le latch. + // Le latch n'est pas observable publiquement et `prompt_ready` n'est déclenchable + // que par le watcher ⇒ on attend brièvement (stream in-memory ⇒ déclenche très + // vite ; la marge couvre l'ordonnancement du thread). + std::thread::sleep(std::time::Duration::from_millis(100)); + assert!( + bus.delegation_ready().is_empty(), + "prompt-ready avant enqueue ⇒ rien encore livré" + ); + + inbox.enqueue(a, ticket(10, "cold task")); + assert_eq!( + bus.delegation_ready().len(), + 1, + "prompt-ready AVANT enqueue ⇒ exactement une DelegationReady" + ); + } + /// `release_cold_start` sans tour différé ⇒ no-op : aucune `DelegationReady` et, /// surtout, PAS de transition Idle (signal de DÉMARRAGE, pas de fin de tour). #[test]