merge(input): intègre le latch released anti-blocage cold-start dans develop
Race cold-start en ordre inverse (libération avant acquisition) du portail d'écriture : latch released garantissant exactly-once. Validé QA (fmt/check verts, 42 tests input passés). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@ -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]
|
||||||
|
|||||||
Reference in New Issue
Block a user