From 8bb832c3b959ae3f2cecc50fc73085f31619f0a8 Mon Sep 17 00:00:00 2001 From: Blomios Date: Sun, 21 Jun 2026 18:09:27 +0200 Subject: [PATCH] fix(input): latch released anti-blocage de la race cold-start en ordre inverse MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Au démarrage à froid, lorsque les évènements de cycle de vie du portail d'écriture arrivent dans l'ordre inverse (libération observée avant l'acquisition correspondante), le latch busy restait coincé et bloquait définitivement la médiation d'entrée. On introduit un latch `released` qui mémorise une libération anticipée et garantit une sémantique exactly-once : l'acquisition tardive consomme la libération déjà vue au lieu de re-verrouiller. Couvert par 2 tests de régression. Co-Authored-By: Claude Opus 4.8 --- crates/infrastructure/src/input/mod.rs | 159 +++++++++++++++++++++++-- 1 file changed, 152 insertions(+), 7 deletions(-) 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]