//! Helper applicatif `send_blocking` — le rendez-vous synchrone inter-agents //! au-dessus du port [`AgentSession`] (ARCHITECTURE §17.1 / §17.4). //! //! `AgentSession::send` retourne un **flux** d'événements (`ReplyStream`), à la //! manière de `PtyPort::subscribe_output`, mais **typé** : deltas de texte → //! activités d'outil → **un** événement terminal déterministe //! [`ReplyEvent::Final`]. Le rendez-vous synchrone dont l'orchestrateur a besoin //! (§17.4) s'obtient en **drainant ce flux jusqu'au `Final`** : c'est *la* primitive //! de la messagerie inter-agents, **sans outbox, sans corrélation fichier** (le //! `Final` *est* la fin de tour). //! //! DRY : un seul chemin de lecture (le flux). `send_blocking` n'est qu'un *consom- //! mateur* du même flux que la cellule chat utilise pour le rendu incrémental. use std::time::Duration; use domain::ids::AgentId; use domain::input::InputMediator; use domain::ports::{AgentSession, AgentSessionError, ReplyEvent}; use domain::readiness::{ReadinessPolicy, ReadinessSignal}; /// Issue d'un tour drainé (ARCHITECTURE §21.2-T4, réconciliation de la limite de /// session avec le contrat « seul `Final` est terminal »). /// /// Un tour se termine de deux façons **gracieuses** : /// - [`TurnOutcome::Completed`] : le flux a rendu son [`ReplyEvent::Final`] — fin /// déterministe normale (cas historique, contenu agrégé) ; /// - [`TurnOutcome::RateLimited`] : le flux s'est **clos sans `Final`** *parce que* /// l'agent est entré en **limite de session** (un [`ReplyEvent::RateLimited`] a été /// observé dans le tour). Ce n'est **pas** une erreur (§21.2-T4) : l'agent reste /// vivant, le service de limite (lot LS4) arme la reprise. /// /// Un flux clos **sans `Final` ET sans `RateLimited`** reste une **erreur** /// [`AgentSessionError::Io`] (tour réellement interrompu) — comportement inchangé. #[derive(Debug, Clone, PartialEq, Eq)] pub enum TurnOutcome { /// Tour terminé normalement par un `Final` ; porte le contenu agrégé. Completed(String), /// Tour clos sur une **limite de session** (sans `Final`) ; porte l'heure de /// reset éventuelle (époche-ms) telle que vue dans le dernier `RateLimited`. RateLimited { /// Instant de reset en époche-ms (`None` ⇒ heure inconnue ⇒ filet humain). resets_at_ms: Option, }, } /// Envoie `prompt` à la session vivante puis **draine le flux de réponse jusqu'au /// [`ReplyEvent::Final`]**, et retourne son contenu agrégé. /// /// C'est le rendez-vous synchrone (§17.4) : on attend que le tour soit /// déterministiquement terminé (`Final`) avant de rendre la main. Les deltas de /// texte et les activités d'outil traversés en chemin sont **ignorés** ici (ils /// servent le rendu incrémental côté UI, pas l'appelant synchrone). /// /// `timeout`, lorsqu'il est fourni, borne l'attente : si aucun `Final` n'est observé /// dans le délai, on retourne [`AgentSessionError::Timeout`] **sans tuer la /// session** (elle reste vivante dans le registre ; l'appelant décide de la suite). /// `None` ⇒ pas de borne temporelle (on attend la fin du tour). /// /// # Errors /// - [`AgentSessionError::Io`]/[`AgentSessionError::Decode`] remontées par `send` /// (échec de communication / décodage de la sortie structurée) ; /// - [`AgentSessionError::Io`] si le flux se termine **sans** `Final` (tour /// interrompu) — y compris un tour clos sur une **limite de session** (le /// rendez-vous synchrone n'a pas de contenu à rendre ; cf. [`TurnOutcome`] et la /// variante riche [`drain_with_readiness_outcome`] pour exploiter la limite) ; /// - [`AgentSessionError::Timeout`] si `timeout` expire avant le `Final`. pub async fn send_blocking( session: &dyn AgentSession, prompt: &str, timeout: Option, ) -> Result { match drain_bounded_events(session, prompt, timeout, |_event| {}, |_signal| {}).await? { TurnOutcome::Completed(content) => Ok(content), // Le rendez-vous synchrone (ask) attend un contenu : un tour limité n'en a pas // ⇒ on conserve le comportement historique (erreur), sans casser le contrat. TurnOutcome::RateLimited { .. } => Err(AgentSessionError::Io( "le tour s'est clos en limite de session, sans contenu Final".to_string(), )), } } /// Comme [`send_blocking`], mais **branche la readiness** : à chaque événement du /// tour, [`ReadinessPolicy::classify`] est consulté et, dès qu'il renvoie /// [`ReadinessSignal::TurnEnded`] (le `Final`), le médiateur d'entrée est notifié /// (`mark_idle(agent)`) pour que la FIFO de l'agent avance — **sans** dépendre d'un /// `idea_reply` explicite ni du sniff de prompt PTY (chantier readiness/heartbeat, /// lot 1, fix de la cause racine du blocage `Busy`). /// /// DRY : **un seul** chemin de lecture du flux (la boucle de [`drain_bounded`]) ; /// cette fonction n'est que `send_blocking` muni d'un *sink* de readiness. Le `Final` /// réveille donc à la fois le `pending` (via la valeur de retour) **et** la FIFO (via /// `mark_idle`). `idea_reply` reste un signal alternatif (premier arrivé gagne) côté /// orchestrateur. /// /// # Errors /// Identiques à [`send_blocking`] (échec `send`/décodage, flux clos sans `Final`, /// timeout). Un tour clos sur une **limite de session** ⇒ [`AgentSessionError::Io`] /// **ici** (signature historique `Result`, zéro régression pour l'appelant /// orchestrateur) ; utilise [`drain_with_readiness_outcome`] pour exploiter la limite. pub async fn drain_with_readiness( session: &dyn AgentSession, prompt: &str, timeout: Option, mediator: &dyn InputMediator, agent: AgentId, ) -> Result { match drain_with_readiness_outcome(session, prompt, timeout, mediator, agent).await? { TurnOutcome::Completed(content) => Ok(content), TurnOutcome::RateLimited { .. } => Err(AgentSessionError::Io( "le tour s'est clos en limite de session, sans contenu Final".to_string(), )), } } /// Variante **riche** de [`drain_with_readiness`] : même branchement readiness, mais /// retourne le [`TurnOutcome`] complet au lieu de réduire la limite à une erreur. /// /// C'est le point d'entrée du chemin **conscient de la limite** (réconciliation /// §21.2-T4) : sur un tour clos sans `Final` mais ayant vu un /// [`ReplyEvent::RateLimited`], il renvoie `Ok(`[`TurnOutcome::RateLimited`]`)` (fin /// gracieuse) plutôt qu'une `Io`. Le drain applicatif / le `SessionLimitService` /// (LS4) consomment cette issue pour armer la reprise. Le câblage de ce chemin sur /// le tour délégué de l'orchestrateur est du ressort de LS7. /// /// **Readiness inchangée** : `mark_idle` reste piloté **uniquement** par le `Final` /// (`TurnEnded`) — un `RateLimited` ne marque **pas** l'agent `Idle` (§21.5 : « le /// `mark_idle`/FIFO reste piloté par `Final`/timeout »). /// /// # Errors /// Comme [`drain_with_readiness`], **sauf** qu'un tour limité n'est plus une erreur /// (il devient [`TurnOutcome::RateLimited`]). pub async fn drain_with_readiness_outcome( session: &dyn AgentSession, prompt: &str, timeout: Option, mediator: &dyn InputMediator, agent: AgentId, ) -> Result { // `on_signal` ne reçoit QUE les événements terminaux (le `Final` ⇒ `TurnEnded`) : // la readiness ne classe pas les non-terminaux. Pour le **battement** de vivacité // (lot 2) on a besoin de notifier le médiateur à CHAQUE événement non terminal // (delta / activité / heartbeat) ⇒ on passe un sink d'événement bruts `on_event`. drain_bounded_events( session, prompt, timeout, |event| { // Tout événement **non terminal** prouve la vivacité ⇒ un battement. if !matches!(event, ReplyEvent::Final { .. }) { mediator.mark_alive(agent); } }, |signal| { // Seul le `Final` marque `Idle` : un `RateLimited` ne fait PAS avancer la // FIFO (la reprise est gérée par le service de limite, §21.5). if signal == ReadinessSignal::TurnEnded { mediator.mark_idle(agent); } }, ) .await } /// Ouvre le flux du tour (`send`) et le **draine jusqu'au `Final`**, en appliquant /// la borne temporelle `timeout`, en notifiant `on_event` à **chaque** événement brut /// (pour le battement de vivacité, lot 2) et `on_signal` à chaque [`ReadinessSignal`] /// dérivé par [`ReadinessPolicy`] (le `Final` ⇒ `TurnEnded`). /// /// **Chemin de lecture unique** (DRY) : `send_blocking` et `drain_with_readiness` /// passent tous deux par ici, en différant seulement par leurs *sinks*. La session /// **reste vivante** sur timeout (on ne `shutdown` rien ici, §17.1). async fn drain_bounded_events( session: &dyn AgentSession, prompt: &str, timeout: Option, on_event: impl FnMut(&ReplyEvent), on_signal: impl FnMut(ReadinessSignal), ) -> Result { match timeout { Some(dur) => { match tokio::time::timeout(dur, drain_to_final(session, prompt, on_event, on_signal)) .await { Ok(result) => result, // La session **reste vivante** : on ne `shutdown` rien ici (§17.1). Err(_elapsed) => Err(AgentSessionError::Timeout), } } None => drain_to_final(session, prompt, on_event, on_signal).await, } } /// Ouvre le flux du tour (`send`) et le **draine jusqu'au `Final`**. /// /// Le flux ([`domain::ports::ReplyStream`]) est un itérateur synchrone et borné : /// après le `Final` il ne produit plus rien. On le parcourt donc simplement /// jusqu'à rencontrer le `Final` (et on retourne son contenu) ; si le flux s'épuise /// avant, l'issue dépend de ce qu'on a vu (réconciliation §21.2-T4) : /// - un [`ReplyEvent::RateLimited`] a été observé ⇒ fin **gracieuse** /// [`TurnOutcome::RateLimited`] (l'agent est limité, pas en erreur) ; /// - sinon ⇒ tour réellement interrompu → erreur [`AgentSessionError::Io`] /// (comportement **inchangé**). /// /// Chaque événement est classé par [`ReadinessPolicy`] et le signal éventuel est /// remonté à `on_signal` (le `Final` ⇒ [`ReadinessSignal::TurnEnded`] ; un /// `RateLimited` ⇒ [`ReadinessSignal::RateLimited`]). Deltas, activités et heartbeats /// sont non terminaux ⇒ ignorés par le rendez-vous synchrone. async fn drain_to_final( session: &dyn AgentSession, prompt: &str, mut on_event: impl FnMut(&ReplyEvent), mut on_signal: impl FnMut(ReadinessSignal), ) -> Result { // Mémorise la dernière limite vue (§21.2-T4) : `Some(resets_at_ms)` dès qu'un // `RateLimited` traverse le flux. Sert UNIQUEMENT au cas « clos sans Final » — // un `Final` ultérieur l'emporte toujours (le tour a réellement abouti). let mut last_rate_limit: Option> = None; let stream = session.send(prompt).await?; for event in stream { // Battement de vivacité (lot 2) : notifié pour CHAQUE événement brut, avant le // classement readiness. Le sink décide (les non-terminaux prouvent la vivacité). on_event(&event); if let Some(signal) = ReadinessPolicy::classify(&event) { on_signal(signal); } match event { ReplyEvent::Final { content } => return Ok(TurnOutcome::Completed(content)), ReplyEvent::RateLimited { resets_at_ms } => last_rate_limit = Some(resets_at_ms), // TextDelta / ToolActivity / Heartbeat : non terminaux, ignorés ici. ReplyEvent::TextDelta { .. } | ReplyEvent::ToolActivity { .. } | ReplyEvent::Heartbeat => {} } } // Flux clos sans `Final` : fin gracieuse SI une limite a été vue (§21.2-T4), // sinon erreur comme avant. match last_rate_limit { Some(resets_at_ms) => Ok(TurnOutcome::RateLimited { resets_at_ms }), None => Err(AgentSessionError::Io( "le flux de réponse s'est terminé sans événement Final".to_string(), )), } } #[cfg(test)] mod tests { use super::*; use std::sync::Mutex; use domain::ids::SessionId; use domain::input::{AgentBusyState, SubmitConfig}; use domain::mailbox::{PendingReply, Ticket}; use domain::ports::PtyHandle; fn agent(n: u128) -> AgentId { AgentId::from_uuid(uuid::Uuid::from_u128(n)) } /// Session factice : `send` rejoue une liste fixe d'événements (terminée par un /// `Final`). struct FakeSession { events: Vec, } #[async_trait::async_trait] impl AgentSession for FakeSession { fn id(&self) -> SessionId { SessionId::from_uuid(uuid::Uuid::from_u128(1)) } fn conversation_id(&self) -> Option { None } async fn send( &self, _prompt: &str, ) -> Result { Ok(Box::new(self.events.clone().into_iter())) } async fn shutdown(&self) -> Result<(), AgentSessionError> { Ok(()) } } /// Médiateur factice qui enregistre l'ordre des `mark_alive` / `mark_idle`. #[derive(Default)] struct RecordingMediator { calls: Mutex>, } impl InputMediator for RecordingMediator { fn enqueue(&self, _agent: AgentId, _ticket: Ticket) -> PendingReply { unreachable!("non utilisé par drain_with_readiness") } fn preempt(&self, _agent: AgentId) {} fn mark_idle(&self, _agent: AgentId) { self.calls.lock().unwrap().push("idle"); } fn mark_alive(&self, _agent: AgentId) { self.calls.lock().unwrap().push("alive"); } fn busy_state(&self, _agent: AgentId) -> AgentBusyState { AgentBusyState::Idle } fn bind_handle(&self, _agent: AgentId, _handle: PtyHandle) {} fn bind_handle_with_submit( &self, _agent: AgentId, _handle: PtyHandle, _submit: SubmitConfig, ) { } } #[tokio::test] async fn drain_marks_alive_on_each_non_terminal_then_idle_on_final() { let session = FakeSession { events: vec![ ReplyEvent::TextDelta { text: "a".into() }, ReplyEvent::ToolActivity { label: "lit".into(), }, ReplyEvent::Heartbeat, ReplyEvent::Final { content: "fini".into(), }, ], }; let mediator = RecordingMediator::default(); let out = drain_with_readiness(&session, "go", None, &mediator, agent(1)) .await .expect("drain ok"); assert_eq!(out, "fini"); // Trois battements (delta, activité, heartbeat) PUIS l'idle sur le Final. assert_eq!( *mediator.calls.lock().unwrap(), vec!["alive", "alive", "alive", "idle"], "un battement par événement non terminal, idle au Final (pas de battement sur le Final)" ); } }