fix(input): latch released anti-blocage de la race cold-start en ordre inverse

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 <noreply@anthropic.com>
This commit is contained in:
2026-06-21 18:09:27 +02:00
parent a66881d73d
commit 8bb832c3b9

View File

@ -57,6 +57,17 @@ struct BusyTracker {
/// pour un démarrage à froid, publiée par [`BusyTracker::prompt_ready`] à l'apparition /// pour un démarrage à froid, publiée par [`BusyTracker::prompt_ready`] à l'apparition
/// du prompt (jamais avant). Absent ⇒ aucun tour en attente de gate. /// du prompt (jamais avant). Absent ⇒ aucun tour en attente de gate.
deferred: Mutex<HashMap<AgentId, DeferredDelegation>>, deferred: Mutex<HashMap<AgentId, DeferredDelegation>>,
/// **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<HashSet<AgentId>>,
events: Option<Arc<dyn EventBus>>, events: Option<Arc<dyn EventBus>>,
/// Sink de livraison headless (cf. [`HeadlessSink`]). Câblé par /// Sink de livraison headless (cf. [`HeadlessSink`]). Câblé par
/// [`MediatedInbox::with_pty`]/[`MediatedInbox::with_events`] quand un PTY est /// [`MediatedInbox::with_pty`]/[`MediatedInbox::with_events`] quand un PTY est
@ -139,6 +150,7 @@ impl BusyTracker {
liveness: Mutex::new(HashMap::new()), liveness: Mutex::new(HashMap::new()),
starting: Mutex::new(HashSet::new()), starting: Mutex::new(HashSet::new()),
deferred: Mutex::new(HashMap::new()), deferred: Mutex::new(HashMap::new()),
released: Mutex::new(HashSet::new()),
events, events,
headless_sink: None, headless_sink: None,
mailbox: None, mailbox: None,
@ -158,6 +170,12 @@ impl BusyTracker {
.unwrap_or_else(std::sync::PoisonError::into_inner) .unwrap_or_else(std::sync::PoisonError::into_inner)
} }
fn lock_released(&self) -> std::sync::MutexGuard<'_, HashSet<AgentId>> {
self.released
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
/// Marque un agent **en démarrage à froid** : son tout premier tour devra être /// 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 à /// *gaté* sur le prompt-ready watcher (la `DelegationReady` sera différée à
/// l'apparition du prompt). À n'appeler **que** lorsqu'un watcher est effectivement /// 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). /// - Sinon : comportement historique ⇒ `mark_idle` (fait avancer la FIFO).
fn prompt_ready(&self, agent: AgentId) { fn prompt_ready(&self, agent: AgentId) {
// On ne retire l'agent de `starting` qu'ici : tant que son premier tour n'a pas // 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é. // été livré, un re-arming éventuel doit rester gaté. On capture `was_starting`
self.lock_starting().remove(&agent); // 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); let deferred = self.lock_deferred().remove(&agent);
match deferred { match deferred {
Some(d) => { Some(d) => {
@ -354,6 +374,17 @@ impl BusyTracker {
); );
self.publish_deferred(agent, d); 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 => { None => {
// Pas de tour différé : la cible est revenue à son prompt sans // Pas de tour différé : la cible est revenue à son prompt sans
// `idea_reply`. On capture le ticket actif AVANT `mark_idle` (qui efface // `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 /// pas de fin de tour). Idempotent et OR-safe avec `prompt_ready` (le `remove` ne rend
/// `Some` qu'une fois). /// `Some` qu'une fois).
fn release_cold_start(&self, agent: AgentId) { 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) { if let Some(d) = self.lock_deferred().remove(&agent) {
// Cas nominal : l'`enqueue` a déjà parqué le tour ⇒ on le draine.
application::diag!( application::diag!(
"[input-mediator] release_cold_start (mcp-ready) released deferred delegation agent={agent} ticket={}", "[input-mediator] release_cold_start (mcp-ready) released deferred delegation agent={agent} ticket={}",
d.ticket d.ticket
); );
self.publish_deferred(agent, d); 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, submit_delay_ms: submit.delay_ms,
}; };
if self.tracker.lock_starting().contains(&agent) { if self.tracker.lock_starting().contains(&agent) {
// 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!( application::diag!(
"[input-mediator] enqueue deferred (cold-start gate) agent={agent} ticket={ticket_id}" "[input-mediator] enqueue deferred (cold-start gate) agent={agent} ticket={ticket_id}"
); );
self.tracker.lock_deferred().insert(agent, deferred); self.tracker.lock_deferred().insert(agent, deferred);
}
} else { } else {
self.tracker.publish_deferred(agent, deferred); 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<dyn EventBus>);
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<dyn PtyPort>,
)
.with_events(Arc::clone(&bus) as Arc<dyn EventBus>);
// 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, /// `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). /// surtout, PAS de transition Idle (signal de DÉMARRAGE, pas de fin de tour).
#[test] #[test]