From 9000b4d09f7d803d0be466de17a3918822522ad8 Mon Sep 17 00:00:00 2001 From: Blomios Date: Tue, 16 Jun 2026 18:55:26 +0200 Subject: [PATCH] =?UTF-8?q?feat(session-limits):=20LS4=20=E2=80=94=20servi?= =?UTF-8?q?ce=20application=20+=20r=C3=A9conciliation=20T4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Orchestre la détection et la reprise au niveau application : - session_limit.rs (nouveau) : SessionLimitService + port AgentResumer + const RESUME_PROMPT. - structured.rs : réconciliation T4 — enum TurnOutcome + drain_with_readiness_outcome ; signatures historiques préservées. - agent/mod.rs + lib.rs : modules et re-exports. Tests QA (nouveaux) : tests/session_limit_service.rs (9) + tests/session_limit_t4.rs (7). `cargo test -p application` = tous binaires verts / 0 failed (16 nouveaux), zéro régression (drain_with_readiness_lot1 7/7, send_blocking_d1 9/9) ; builds domaine+infra+application 0 warning. Co-Authored-By: Claude Opus 4.8 --- crates/application/src/agent/mod.rs | 6 +- crates/application/src/agent/session_limit.rs | 229 +++++++++ crates/application/src/agent/structured.rs | 118 ++++- crates/application/src/lib.rs | 9 +- .../tests/session_limit_service.rs | 440 ++++++++++++++++++ crates/application/tests/session_limit_t4.rs | 205 ++++++++ 6 files changed, 987 insertions(+), 20 deletions(-) create mode 100644 crates/application/src/agent/session_limit.rs create mode 100644 crates/application/tests/session_limit_service.rs create mode 100644 crates/application/tests/session_limit_t4.rs diff --git a/crates/application/src/agent/mod.rs b/crates/application/src/agent/mod.rs index e93f6e9..8f2389e 100644 --- a/crates/application/src/agent/mod.rs +++ b/crates/application/src/agent/mod.rs @@ -10,13 +10,17 @@ mod catalogue; mod inspect; mod lifecycle; mod resume; +mod session_limit; mod structured; mod usecases; pub(crate) use lifecycle::unique_md_path; pub(crate) use lifecycle::ReattachDecision; -pub use structured::{drain_with_readiness, send_blocking}; +pub use session_limit::{AgentResumer, SessionLimitService, RESUME_PROMPT}; +pub use structured::{ + drain_with_readiness, drain_with_readiness_outcome, send_blocking, TurnOutcome, +}; pub use catalogue::{reference_profile_id, reference_profiles, selectable_reference_profiles}; pub use inspect::{InspectConversation, InspectConversationInput, InspectConversationOutput}; diff --git a/crates/application/src/agent/session_limit.rs b/crates/application/src/agent/session_limit.rs new file mode 100644 index 0000000..0cc756f --- /dev/null +++ b/crates/application/src/agent/session_limit.rs @@ -0,0 +1,229 @@ +//! [`SessionLimitService`] — orchestration applicative des **limites de session** +//! des agents (ARCHITECTURE §21.5) : **détecter → planifier → reprendre**, et +//! **annuler**. +//! +//! Service **pur-ports** (SOLID/hexagonal) : il ne dépend que de traits du domaine +//! ([`Clock`], [`Scheduler`], [`EventBus`]) et d'un port applicatif de reprise +//! ([`AgentResumer`], implémenté au composition root en LS7 par-dessus `LaunchAgent`). +//! Aucune dépendance vers un adapter concret ⇒ entièrement testable avec des fakes. +//! +//! # État en mémoire uniquement (§21.1-3) +//! +//! La seule mémoire du service est une table `agent_id → ScheduleId` des reprises +//! **armées** (pour pouvoir les annuler). Aucune persistance : à un redémarrage +//! d'IdeA le chemin `ListResumableAgents` existant prend le relais. +//! +//! # Dédoublonnage par agent (§21.10-4) +//! +//! Un agent n'a qu'**une** reprise armée à la fois : un second signal de limite +//! **rafraîchit** l'armement (annule l'ancien, arme le nouveau) au lieu d'empiler. + +use std::collections::HashMap; +use std::sync::{Arc, Mutex}; + +use async_trait::async_trait; + +use domain::ids::{AgentId, NodeId, ScheduleId}; +use domain::ports::{Clock, EventBus, ScheduledTask, Scheduler}; +use domain::session_limit::{plan_resume, RateLimitSource, ResumePlan, SessionLimit}; +use domain::DomainEvent; + +use crate::error::AppError; + +/// Prompt **court** envoyé à l'agent au moment de la reprise automatique (§21.5). +/// +/// `--resume` (via [`domain::ports::SessionPlan::Resume`]) porte déjà tout +/// l'historique : ce prompt n'a qu'à **réamorcer** le tour, pas reconstruire le +/// contexte. Volontairement neutre et model-agnostique. +pub const RESUME_PROMPT: &str = + "La limite de session est levée. Reprends là où tu t'étais arrêté."; + +/// Port applicatif de **reprise d'un agent** (frontière implémentée au composition +/// root, LS7). Calqué sur les autres traits-passerelles de l'application +/// ([`crate::agent::HandoffProvider`], [`crate::agent::ProviderSessionProvider`]) : +/// l'app-tauri le branche par-dessus le mécanisme de lancement existant +/// (`LaunchAgent` + `AgentSessionFactory`) avec [`domain::ports::SessionPlan::Resume`]. +/// +/// Le service ne sait **pas** relancer un agent lui-même (cela exige le `Project`, le +/// profil, le contexte préparé, le PTY… que seul `LaunchAgent` résout) ; il délègue +/// donc à ce port, en restant testable avec un fake. +#[async_trait] +pub trait AgentResumer: Send + Sync { + /// Relance/réattache l'agent `agent_id` dans sa cellule `node_id`, en reprenant la + /// conversation moteur `conversation_id` (`SessionPlan::Resume` côté lancement) et + /// en lui transmettant `resume_prompt` comme premier tour. + /// + /// # Errors + /// [`AppError`] si la relance échoue (profil/contexte introuvable, échec de + /// démarrage de session…). Le service propage l'erreur sans publier `AgentResumed`. + async fn resume( + &self, + agent_id: AgentId, + node_id: NodeId, + conversation_id: Option, + resume_prompt: &str, + ) -> Result<(), AppError>; +} + +/// Service d'orchestration des limites de session (§21.5). +pub struct SessionLimitService { + clock: Arc, + scheduler: Arc, + events: Arc, + resumer: Arc, + /// Reprises **armées** non encore tirées : `agent_id → ScheduleId` (en mémoire). + armed: Mutex>, +} + +impl SessionLimitService { + /// Construit le service depuis ses ports injectés (composition root). + #[must_use] + pub fn new( + clock: Arc, + scheduler: Arc, + events: Arc, + resumer: Arc, + ) -> Self { + Self { + clock, + scheduler, + events, + resumer, + armed: Mutex::new(HashMap::new()), + } + } + + /// **(a) Détection → planification.** À partir d'un signal de limite + /// (`ReadinessSignal::RateLimited`/`TurnOutcome::RateLimited`) pour `agent_id` dans + /// la cellule `node_id`, construit un [`SessionLimit`] (source `Structured`), + /// calcule le plan via [`plan_resume`] et agit : + /// + /// - [`ResumePlan::Scheduled`] ⇒ publie `AgentRateLimited`, **arme** la reprise via + /// [`Scheduler::arm`] (dédoublonnée : un éventuel armement antérieur du même agent + /// est annulé d'abord, §21.10-4), mémorise le [`ScheduleId`], puis publie + /// `AgentResumeScheduled`. + /// - [`ResumePlan::HumanFallback`] (heure de reset inconnue) ⇒ publie + /// `AgentRateLimited{None}` puis `AgentRateLimitSuspected{None}` (filet humain ; + /// la confirmation UI est LS6/LS8 — ici on émet seulement l'événement). + pub fn on_rate_limited( + &self, + agent_id: AgentId, + node_id: NodeId, + conversation_id: Option, + resets_at_ms: Option, + ) { + let now = self.clock.now_millis(); + let limit = SessionLimit::new(resets_at_ms, now, RateLimitSource::Structured); + + match plan_resume(now, &limit, conversation_id) { + ResumePlan::Scheduled { + fire_at_ms, + conversation_id, + } => { + self.events.publish(DomainEvent::AgentRateLimited { + agent_id, + resets_at_ms, + }); + // Dédoublonnage (§21.10-4) : un signal de rafraîchissement annule + // l'armement précédent (sans événement d'annulation : c'est interne). + self.disarm(agent_id); + let id = self.scheduler.arm( + fire_at_ms, + ScheduledTask::ResumeAgent { + agent_id, + node_id, + conversation_id, + }, + ); + self.armed.lock().expect("session-limit mutex sain").insert(agent_id, id); + self.events + .publish(DomainEvent::AgentResumeScheduled { agent_id, fire_at_ms }); + } + ResumePlan::HumanFallback => { + self.events.publish(DomainEvent::AgentRateLimited { + agent_id, + resets_at_ms: None, + }); + self.events.publish(DomainEvent::AgentRateLimitSuspected { + agent_id, + resets_at_ms: None, + }); + } + } + } + + /// **(b) Exécution de la reprise.** Consomme une [`ScheduledTask::ResumeAgent`] + /// échue (celle que `TokioScheduler` pousse dans le canal de remise ; le câblage du + /// récepteur dans le runtime est LS7). Retire l'entrée armée (le réveil a tiré), + /// relance l'agent via [`AgentResumer`] avec [`RESUME_PROMPT`], puis publie + /// `AgentResumed`. + /// + /// # Errors + /// [`AppError`] propagée par [`AgentResumer::resume`] (la relance a échoué) ; dans + /// ce cas `AgentResumed` n'est **pas** publié. + pub async fn execute_resume(&self, task: ScheduledTask) -> Result<(), AppError> { + let ScheduledTask::ResumeAgent { + agent_id, + node_id, + conversation_id, + } = task; + + // Le réveil a tiré : l'entrée armée n'a plus lieu d'être (qu'on réussisse ou non). + self.disarm(agent_id); + + self.resumer + .resume(agent_id, node_id, conversation_id, RESUME_PROMPT) + .await?; + + self.events.publish(DomainEvent::AgentResumed { agent_id }); + Ok(()) + } + + /// **(c) Annulation.** Désarme la reprise auto de `agent_id` (socle du « annulable »). + /// Retrouve le [`ScheduleId`], appelle [`Scheduler::cancel`] et, **seulement si** + /// l'annulation a réussi, retire l'entrée et publie `AgentResumeCancelled`. + /// + /// Renvoie `true` ssi la reprise a effectivement été annulée. + /// + /// # Course « cancel pile au tir » (vigilance QA, LS3) + /// Sous runtime multi-thread, le réveil peut tirer **pile** au moment de l'annulation : + /// [`Scheduler::cancel`] renvoie alors `false` (déjà tiré). Dans ce cas on **ne + /// publie pas** `AgentResumeCancelled` et on **laisse l'entrée** (l'`execute_resume` + /// en cours la retirera) : la reprise **suit son cours**, cohérent et sans + /// événement trompeur. + pub fn cancel_resume(&self, agent_id: AgentId) -> bool { + let id = self + .armed + .lock() + .expect("session-limit mutex sain") + .get(&agent_id) + .copied(); + let Some(id) = id else { + return false; // aucune reprise armée pour cet agent. + }; + + if self.scheduler.cancel(id) { + self.armed.lock().expect("session-limit mutex sain").remove(&agent_id); + self.events + .publish(DomainEvent::AgentResumeCancelled { agent_id }); + true + } else { + // Déjà tiré : la reprise suivra son cours, pas d'événement d'annulation. + false + } + } + + /// Retire (best-effort) l'armement de `agent_id` et annule le réveil sous-jacent + /// s'il existe. Usage interne (rafraîchissement / nettoyage post-tir) — **ne publie + /// aucun événement** (contrairement à [`Self::cancel_resume`]). + fn disarm(&self, agent_id: AgentId) { + let previous = self + .armed + .lock() + .expect("session-limit mutex sain") + .remove(&agent_id); + if let Some(id) = previous { + self.scheduler.cancel(id); + } + } +} diff --git a/crates/application/src/agent/structured.rs b/crates/application/src/agent/structured.rs index b2a452a..8884118 100644 --- a/crates/application/src/agent/structured.rs +++ b/crates/application/src/agent/structured.rs @@ -19,6 +19,31 @@ use domain::ids::AgentId; 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é. /// @@ -36,14 +61,23 @@ use domain::readiness::{ReadinessPolicy, ReadinessSignal}; /// - [`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) ; +/// 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 { - drain_bounded_events(session, prompt, timeout, |_event| {}, |_signal| {}).await + 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 @@ -61,7 +95,9 @@ pub async fn send_blocking( /// /// # Errors /// Identiques à [`send_blocking`] (échec `send`/décodage, flux clos sans `Final`, -/// timeout). +/// 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, @@ -69,6 +105,38 @@ pub async fn drain_with_readiness( 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 @@ -84,6 +152,8 @@ pub async fn drain_with_readiness( } }, |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); } @@ -106,7 +176,7 @@ async fn drain_bounded_events( timeout: Option, on_event: impl FnMut(&ReplyEvent), on_signal: impl FnMut(ReadinessSignal), -) -> Result { +) -> Result { match timeout { Some(dur) => match tokio::time::timeout( dur, @@ -126,18 +196,27 @@ async fn drain_bounded_events( /// /// 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, c'est un tour interrompu → erreur [`AgentSessionError::Io`]. +/// 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`]). Deltas, -/// activités et heartbeats sont non terminaux ⇒ ignorés par le rendez-vous synchrone. +/// 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 { +) -> 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 @@ -146,14 +225,23 @@ async fn drain_to_final( if let Some(signal) = ReadinessPolicy::classify(&event) { on_signal(signal); } - if let ReplyEvent::Final { content } = event { - return Ok(content); + 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 => {} } - // TextDelta / ToolActivity / Heartbeat : non terminaux, ignorés ici. } - Err(AgentSessionError::Io( - "le flux de réponse s'est terminé sans événement Final".to_string(), - )) + // 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)] diff --git a/crates/application/src/lib.rs b/crates/application/src/lib.rs index 744699d..4a0a42c 100644 --- a/crates/application/src/lib.rs +++ b/crates/application/src/lib.rs @@ -29,8 +29,9 @@ pub mod terminal; pub mod window; pub use agent::{ - drain_with_readiness, reference_profile_id, reference_profiles, selectable_reference_profiles, - send_blocking, ChangeAgentProfile, ChangeAgentProfileInput, ChangeAgentProfileOutput, + drain_with_readiness, drain_with_readiness_outcome, reference_profile_id, reference_profiles, + selectable_reference_profiles, send_blocking, AgentResumer, ChangeAgentProfile, + ChangeAgentProfileInput, ChangeAgentProfileOutput, ConfigureProfiles, ConfigureProfilesInput, ConfigureProfilesOutput, CreateAgentFromScratch, CreateAgentInput, CreateAgentOutput, DeleteAgent, DeleteAgentInput, DeleteProfile, DeleteProfileInput, DetectProfiles, DetectProfilesInput, DetectProfilesOutput, FirstRunState, @@ -41,8 +42,8 @@ pub use agent::{ ProfileAvailability, ProviderSessionProvider, ReadAgentContext, ReadAgentContextInput, ReadAgentContextOutput, ReferenceProfiles, ReferenceProfilesOutput, ResumableAgent, SaveProfile, SaveProfileInput, - SaveProfileOutput, StructuredSessionDescriptor, UpdateAgentContext, UpdateAgentContextInput, - AGENT_MEMORY_RECALL_BUDGET, + SaveProfileOutput, SessionLimitService, StructuredSessionDescriptor, TurnOutcome, + UpdateAgentContext, UpdateAgentContextInput, AGENT_MEMORY_RECALL_BUDGET, RESUME_PROMPT, }; pub use conversation::RecordTurn; pub use embedder::{ diff --git a/crates/application/tests/session_limit_service.rs b/crates/application/tests/session_limit_service.rs new file mode 100644 index 0000000..bb3564d --- /dev/null +++ b/crates/application/tests/session_limit_service.rs @@ -0,0 +1,440 @@ +//! LS4 — tests unitaires QA du `SessionLimitService` (ARCHITECTURE §21.5), +//! **100 % fakes** des ports : `Clock` fixe, `Scheduler` enregistreur/contrôlable, +//! `EventBus` espion, `AgentResumer` espion/contrôlable. +//! +//! Couvre les trois responsabilités du service : +//! - (a) détection → planification (`on_rate_limited`) : armement + ordre des events ; +//! - (b) exécution de la reprise (`execute_resume`) : prompt, event, propagation d'erreur ; +//! - (c) annulation (`cancel_resume`) : contrat anti-course (cancel `false` ⇒ pas d'event). + +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex}; + +use async_trait::async_trait; + +use application::{AgentResumer, AppError, SessionLimitService, RESUME_PROMPT}; +use domain::ids::{AgentId, NodeId, ScheduleId}; +use domain::ports::{Clock, EventBus, EventStream, ScheduledTask, Scheduler}; +use domain::DomainEvent; +use uuid::Uuid; + +fn aid(n: u128) -> AgentId { + AgentId::from_uuid(Uuid::from_u128(n)) +} +fn nid(n: u128) -> NodeId { + NodeId::from_uuid(Uuid::from_u128(n)) +} + +// --------------------------------------------------------------------------- +// Fakes des ports +// --------------------------------------------------------------------------- + +/// Horloge fixe (déterminisme) : `now` injecté. +struct FixedClock(i64); +impl Clock for FixedClock { + fn now_millis(&self) -> i64 { + self.0 + } +} + +/// EventBus espion : journalise les events publiés, dans l'ordre. +#[derive(Default, Clone)] +struct SpyBus(Arc>>); +impl SpyBus { + fn events(&self) -> Vec { + self.0.lock().unwrap().clone() + } +} +impl EventBus for SpyBus { + fn publish(&self, event: DomainEvent) { + self.0.lock().unwrap().push(event); + } + fn subscribe(&self) -> EventStream { + Box::new(std::iter::empty()) + } +} + +/// Scheduler fake : enregistre chaque `arm` (deadline + task) et chaque `cancel`, +/// distribue des `ScheduleId` déterministes, et rend un résultat de `cancel` +/// contrôlable (pour simuler le cas « déjà tiré »). +#[derive(Clone)] +struct FakeScheduler { + armed: Arc>>, + issued: Arc>>, + cancels: Arc>>, + next_id: Arc>, + cancel_result: Arc, +} +impl FakeScheduler { + fn new() -> Self { + Self { + armed: Arc::new(Mutex::new(Vec::new())), + issued: Arc::new(Mutex::new(Vec::new())), + cancels: Arc::new(Mutex::new(Vec::new())), + next_id: Arc::new(Mutex::new(1)), + cancel_result: Arc::new(AtomicBool::new(true)), + } + } + fn set_cancel_result(&self, v: bool) { + self.cancel_result.store(v, Ordering::SeqCst); + } + fn armed(&self) -> Vec<(i64, ScheduledTask)> { + self.armed.lock().unwrap().clone() + } + fn issued(&self) -> Vec { + self.issued.lock().unwrap().clone() + } + fn cancels(&self) -> Vec { + self.cancels.lock().unwrap().clone() + } +} +impl Scheduler for FakeScheduler { + fn arm(&self, deadline_ms: i64, task: ScheduledTask) -> ScheduleId { + self.armed.lock().unwrap().push((deadline_ms, task)); + let mut n = self.next_id.lock().unwrap(); + let id = ScheduleId::from_uuid(Uuid::from_u128(*n)); + *n += 1; + self.issued.lock().unwrap().push(id); + id + } + fn cancel(&self, id: ScheduleId) -> bool { + self.cancels.lock().unwrap().push(id); + self.cancel_result.load(Ordering::SeqCst) + } +} + +/// AgentResumer fake : enregistre l'appel `resume` (et le prompt reçu), et peut +/// être configuré pour échouer (propagation d'erreur). +#[derive(Clone)] +struct FakeResumer { + calls: Arc, String)>>>, + fail: Arc, +} +impl FakeResumer { + fn new() -> Self { + Self { + calls: Arc::new(Mutex::new(Vec::new())), + fail: Arc::new(AtomicBool::new(false)), + } + } + fn set_fail(&self, v: bool) { + self.fail.store(v, Ordering::SeqCst); + } + fn calls(&self) -> Vec<(AgentId, NodeId, Option, String)> { + self.calls.lock().unwrap().clone() + } +} +#[async_trait] +impl AgentResumer for FakeResumer { + async fn resume( + &self, + agent_id: AgentId, + node_id: NodeId, + conversation_id: Option, + resume_prompt: &str, + ) -> Result<(), AppError> { + self.calls + .lock() + .unwrap() + .push((agent_id, node_id, conversation_id, resume_prompt.to_owned())); + if self.fail.load(Ordering::SeqCst) { + Err(AppError::Internal("reprise échouée (fake)".to_owned())) + } else { + Ok(()) + } + } +} + +/// Banc d'essai : assemble le service avec ses fakes (et garde les poignées). +struct Env { + service: SessionLimitService, + scheduler: FakeScheduler, + bus: SpyBus, + resumer: FakeResumer, +} +fn env_at(now_ms: i64) -> Env { + let scheduler = FakeScheduler::new(); + let bus = SpyBus::default(); + let resumer = FakeResumer::new(); + let service = SessionLimitService::new( + Arc::new(FixedClock(now_ms)), + Arc::new(scheduler.clone()), + Arc::new(bus.clone()), + Arc::new(resumer.clone()), + ); + Env { + service, + scheduler, + bus, + resumer, + } +} + +// =========================================================================== +// (a) Détection → planification +// =========================================================================== + +const NOW: i64 = 1_700_000_000_000; + +/// `on_rate_limited(Some(reset futur))` ⇒ EXACTEMENT un `arm(fire_at_ms, ResumeAgent{..})` +/// avec `fire_at_ms == reset` (futur) ; events `AgentRateLimited` PUIS `AgentResumeScheduled` +/// dans CET ordre. +#[test] +fn on_rate_limited_future_arms_and_emits_in_order() { + let env = env_at(NOW); + let reset = NOW + 60_000; + env.service + .on_rate_limited(aid(1), nid(2), Some("conv-1".to_owned()), Some(reset)); + + // Exactement un arm, avec la bonne échéance et la bonne tâche. + let armed = env.scheduler.armed(); + assert_eq!(armed.len(), 1, "exactement un arm"); + assert_eq!(armed[0].0, reset, "fire_at_ms == reset (futur)"); + assert_eq!( + armed[0].1, + ScheduledTask::ResumeAgent { + agent_id: aid(1), + node_id: nid(2), + conversation_id: Some("conv-1".to_owned()), + } + ); + + // Ordre des events : RateLimited puis ResumeScheduled. + assert_eq!( + env.bus.events(), + vec![ + DomainEvent::AgentRateLimited { + agent_id: aid(1), + resets_at_ms: Some(reset), + }, + DomainEvent::AgentResumeScheduled { + agent_id: aid(1), + fire_at_ms: reset, + }, + ] + ); +} + +/// Reset DÉJÀ PASSÉ ⇒ `plan_resume` clampe à `now` ⇒ `fire_at_ms == now` (jamais dans +/// le passé) ; le `AgentRateLimited` garde l'heure brute (passée), le `ResumeScheduled` +/// porte le `now` clampé. +#[test] +fn on_rate_limited_past_reset_clamps_fire_at_to_now() { + let env = env_at(NOW); + let past = NOW - 60_000; + env.service + .on_rate_limited(aid(1), nid(2), None, Some(past)); + + let armed = env.scheduler.armed(); + assert_eq!(armed.len(), 1); + assert_eq!(armed[0].0, NOW, "fire_at_ms clampé à now (jamais le passé)"); + + assert_eq!( + env.bus.events(), + vec![ + DomainEvent::AgentRateLimited { + agent_id: aid(1), + resets_at_ms: Some(past), // l'heure brute (passée) est conservée dans l'event + }, + DomainEvent::AgentResumeScheduled { + agent_id: aid(1), + fire_at_ms: NOW, // clampé + }, + ] + ); +} + +/// `on_rate_limited(None)` ⇒ AUCUN arm ; events `AgentRateLimited{None}` puis +/// `AgentRateLimitSuspected{None}` (filet humain). +#[test] +fn on_rate_limited_without_reset_is_human_fallback_no_arm() { + let env = env_at(NOW); + env.service + .on_rate_limited(aid(1), nid(2), Some("conv-1".to_owned()), None); + + assert!(env.scheduler.armed().is_empty(), "aucun arm sans heure de reset"); + assert!(env.scheduler.cancels().is_empty(), "aucun cancel non plus"); + assert_eq!( + env.bus.events(), + vec![ + DomainEvent::AgentRateLimited { + agent_id: aid(1), + resets_at_ms: None, + }, + DomainEvent::AgentRateLimitSuspected { + agent_id: aid(1), + resets_at_ms: None, + }, + ] + ); +} + +/// Dédoublonnage (§21.10-4) : deux `on_rate_limited` successifs pour le MÊME agent ⇒ +/// l'ancien `ScheduleId` est annulé sur le Scheduler (cancel interne), un nouveau réveil +/// est armé, et AUCUN `AgentResumeCancelled` n'est émis (le dédoublonnage est interne). +#[test] +fn on_rate_limited_twice_same_agent_dedups_cancelling_previous() { + let env = env_at(NOW); + let reset1 = NOW + 60_000; + let reset2 = NOW + 120_000; + env.service + .on_rate_limited(aid(1), nid(2), Some("conv-1".to_owned()), Some(reset1)); + env.service + .on_rate_limited(aid(1), nid(2), Some("conv-1".to_owned()), Some(reset2)); + + // Deux arms (un par signal), ids distincts. + let issued = env.scheduler.issued(); + assert_eq!(issued.len(), 2, "deux arms (rafraîchissement, pas empilement)"); + // Le premier id émis a été annulé par le dédoublonnage du 2ᵉ signal. + assert_eq!( + env.scheduler.cancels(), + vec![issued[0]], + "l'ancien ScheduleId est cancel-é avant de réarmer" + ); + + // Aucun AgentResumeCancelled : le dédoublonnage est silencieux. + let cancelled: Vec<_> = env + .bus + .events() + .into_iter() + .filter(|e| matches!(e, DomainEvent::AgentResumeCancelled { .. })) + .collect(); + assert!( + cancelled.is_empty(), + "le dédoublonnage interne n'émet PAS AgentResumeCancelled" + ); +} + +// =========================================================================== +// (b) Exécution de la reprise +// =========================================================================== + +/// `execute_resume(ResumeAgent{..})` ⇒ `AgentResumer::resume` appelé avec +/// (agent, node, conv, RESUME_PROMPT) ; `AgentResumed` publié ; l'entrée armée est +/// retirée (un `cancel_resume` ultérieur ⇒ false). +#[tokio::test] +async fn execute_resume_calls_resumer_with_prompt_and_emits_resumed() { + let env = env_at(NOW); + // Arme d'abord (pour prouver que l'entrée est ensuite retirée). + env.service + .on_rate_limited(aid(1), nid(2), Some("conv-1".to_owned()), Some(NOW + 60_000)); + + let task = ScheduledTask::ResumeAgent { + agent_id: aid(1), + node_id: nid(2), + conversation_id: Some("conv-1".to_owned()), + }; + env.service.execute_resume(task).await.expect("resume ok"); + + // Le resumer a vu exactement l'appel attendu, avec le prompt constant. + assert_eq!( + env.resumer.calls(), + vec![(aid(1), nid(2), Some("conv-1".to_owned()), RESUME_PROMPT.to_owned())] + ); + + // AgentResumed publié. + assert!( + env.bus + .events() + .iter() + .any(|e| *e == DomainEvent::AgentResumed { agent_id: aid(1) }), + "AgentResumed doit être publié" + ); + + // L'entrée armée a été retirée ⇒ cancel_resume ultérieur = false. + assert!( + !env.service.cancel_resume(aid(1)), + "après execute_resume, plus rien à annuler" + ); +} + +/// `AgentResumer` qui échoue ⇒ l'erreur est propagée ET `AgentResumed` n'est PAS publié. +#[tokio::test] +async fn execute_resume_propagates_error_without_emitting_resumed() { + let env = env_at(NOW); + env.resumer.set_fail(true); + + let task = ScheduledTask::ResumeAgent { + agent_id: aid(1), + node_id: nid(2), + conversation_id: None, + }; + let err = env + .service + .execute_resume(task) + .await + .expect_err("la reprise doit échouer"); + assert!(matches!(err, AppError::Internal(_)), "erreur propagée: {err:?}"); + + assert!( + !env + .bus + .events() + .iter() + .any(|e| matches!(e, DomainEvent::AgentResumed { .. })), + "AgentResumed ne doit PAS être publié si la reprise a échoué" + ); +} + +// =========================================================================== +// (c) Annulation +// =========================================================================== + +/// `cancel_resume` après un armement, Scheduler renvoyant `true` ⇒ renvoie `true` et +/// publie `AgentResumeCancelled`. +#[test] +fn cancel_resume_after_arm_returns_true_and_emits_cancelled() { + let env = env_at(NOW); + env.service + .on_rate_limited(aid(1), nid(2), Some("conv-1".to_owned()), Some(NOW + 60_000)); + let issued = env.scheduler.issued(); + + assert!(env.service.cancel_resume(aid(1)), "cancel d'un réveil armé ⇒ true"); + // Le bon ScheduleId a été passé au Scheduler. + assert_eq!(env.scheduler.cancels(), vec![issued[0]]); + // AgentResumeCancelled publié. + assert!( + env.bus + .events() + .iter() + .any(|e| *e == DomainEvent::AgentResumeCancelled { agent_id: aid(1) }), + "AgentResumeCancelled doit être publié" + ); +} + +/// `cancel_resume` sans armement préalable ⇒ `false`, aucun event, et le Scheduler +/// n'est même pas sollicité. +#[test] +fn cancel_resume_without_arm_is_false_no_event() { + let env = env_at(NOW); + assert!(!env.service.cancel_resume(aid(1))); + assert!(env.scheduler.cancels().is_empty(), "Scheduler non sollicité"); + assert!(env.bus.events().is_empty(), "aucun event"); +} + +/// Contrat ANTI-COURSE : Scheduler renvoyant `false` (« déjà tiré ») ⇒ `cancel_resume` +/// renvoie `false` et n'émet PAS `AgentResumeCancelled` (la reprise suit son cours). +#[test] +fn cancel_resume_when_scheduler_already_fired_is_false_no_event() { + let env = env_at(NOW); + env.service + .on_rate_limited(aid(1), nid(2), Some("conv-1".to_owned()), Some(NOW + 60_000)); + // Simule un réveil déjà tiré : cancel renvoie false. + env.scheduler.set_cancel_result(false); + + assert!( + !env.service.cancel_resume(aid(1)), + "cancel pile au tir ⇒ false (la reprise suit son cours)" + ); + // Le Scheduler a bien été sollicité (mais a répondu false). + assert_eq!(env.scheduler.cancels().len(), 1); + // AUCUN AgentResumeCancelled (pas d'event trompeur). + assert!( + !env + .bus + .events() + .iter() + .any(|e| matches!(e, DomainEvent::AgentResumeCancelled { .. })), + "pas d'AgentResumeCancelled quand le réveil a déjà tiré" + ); +} diff --git a/crates/application/tests/session_limit_t4.rs b/crates/application/tests/session_limit_t4.rs new file mode 100644 index 0000000..2a225b6 --- /dev/null +++ b/crates/application/tests/session_limit_t4.rs @@ -0,0 +1,205 @@ +//! LS4 — réconciliation §21.2-T4 au niveau applicatif : `drain_with_readiness_outcome` +//! traduit un tour clos par une **limite de session** (un `RateLimited` vu, pas de +//! `Final`) en `Ok(TurnOutcome::RateLimited{..})` — une **fin gracieuse**, pas une +//! erreur. **100 % fakes** (fake `AgentSession` scriptable + fake `InputMediator`). +//! +//! On vérifie AUSSI la NON-RÉGRESSION des signatures historiques `Result` : +//! - `drain_with_readiness` (et `send_blocking`) ⇒ un tour limité reste une `Io` +//! (zéro régression pour l'appelant orchestrateur) ; +//! - un flux clos SANS `Final` ET SANS `RateLimited` reste une `Io` (tour tronqué). + +use std::sync::Mutex; + +use async_trait::async_trait; + +use application::{drain_with_readiness, drain_with_readiness_outcome, send_blocking, TurnOutcome}; +use domain::ids::AgentId; +use domain::input::{AgentBusyState, InputMediator}; +use domain::mailbox::{PendingReply, Ticket}; +use domain::ports::{AgentSession, AgentSessionError, ReplyEvent, ReplyStream}; +use domain::SessionId; +use uuid::Uuid; + +fn aid(n: u128) -> AgentId { + AgentId::from_uuid(Uuid::from_u128(n)) +} + +// --------------------------------------------------------------------------- +// Fake AgentSession : `send` rejoue une liste fixe d'événements. +// --------------------------------------------------------------------------- + +struct FakeSession { + events: Vec, +} +#[async_trait] +impl AgentSession for FakeSession { + fn id(&self) -> SessionId { + SessionId::from_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(()) + } +} + +// --------------------------------------------------------------------------- +// Fake InputMediator : journalise mark_idle / mark_alive (pour prouver que la +// readiness reste pilotée par le Final, pas par un RateLimited). +// --------------------------------------------------------------------------- + +#[derive(Default)] +struct RecordingMediator { + calls: Mutex>, +} +impl InputMediator for RecordingMediator { + fn enqueue(&self, _agent: AgentId, _ticket: Ticket) -> PendingReply { + PendingReply::new(Box::pin(std::future::pending())) + } + 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 + } +} + +// =========================================================================== +// drain_with_readiness_outcome — issue riche T4 +// =========================================================================== + +/// `[RateLimited{Some(t)}]` sans Final ⇒ `Ok(TurnOutcome::RateLimited{Some(t)})` +/// (fin gracieuse, PAS d'Err). +#[tokio::test] +async fn outcome_rate_limited_some_without_final_is_graceful() { + let session = FakeSession { + events: vec![ReplyEvent::RateLimited { + resets_at_ms: Some(1_700_000_000_000), + }], + }; + let mediator = RecordingMediator::default(); + let out = drain_with_readiness_outcome(&session, "go", None, &mediator, aid(1)) + .await + .expect("un tour limité est une fin gracieuse, pas une erreur"); + assert_eq!( + out, + TurnOutcome::RateLimited { + resets_at_ms: Some(1_700_000_000_000) + } + ); + // Readiness : un RateLimited ne marque PAS Idle (piloté par Final uniquement). + assert!( + !mediator.calls.lock().unwrap().contains(&"idle"), + "un RateLimited ne doit pas marquer l'agent Idle" + ); +} + +/// `[RateLimited{None}]` sans Final ⇒ `Ok(TurnOutcome::RateLimited{None})`. +#[tokio::test] +async fn outcome_rate_limited_none_without_final_is_graceful() { + let session = FakeSession { + events: vec![ReplyEvent::RateLimited { resets_at_ms: None }], + }; + let mediator = RecordingMediator::default(); + let out = drain_with_readiness_outcome(&session, "go", None, &mediator, aid(1)) + .await + .expect("fin gracieuse"); + assert_eq!(out, TurnOutcome::RateLimited { resets_at_ms: None }); +} + +/// `[.., RateLimited, Final]` ⇒ `Ok(TurnOutcome::Completed(contenu))` : le Final +/// l'emporte toujours (le tour a réellement abouti). +#[tokio::test] +async fn outcome_rate_limited_then_final_is_completed() { + let session = FakeSession { + events: vec![ + ReplyEvent::TextDelta { text: "a".into() }, + ReplyEvent::RateLimited { + resets_at_ms: Some(42), + }, + ReplyEvent::Final { + content: "fini".into(), + }, + ], + }; + let mediator = RecordingMediator::default(); + let out = drain_with_readiness_outcome(&session, "go", None, &mediator, aid(1)) + .await + .expect("ok"); + assert_eq!(out, TurnOutcome::Completed("fini".to_owned())); + // Le Final marque bien Idle. + assert!(mediator.calls.lock().unwrap().contains(&"idle")); +} + +/// `[TextDelta]` seul (ni Final ni RateLimited) ⇒ `Err(Io)` INCHANGÉ : un vrai flux +/// tronqué reste une erreur (non-régression critique). +#[tokio::test] +async fn outcome_truncated_stream_without_final_or_ratelimit_is_io_error() { + let session = FakeSession { + events: vec![ReplyEvent::TextDelta { text: "a".into() }], + }; + let mediator = RecordingMediator::default(); + let err = drain_with_readiness_outcome(&session, "go", None, &mediator, aid(1)) + .await + .expect_err("flux tronqué sans limite ⇒ erreur"); + assert!(matches!(err, AgentSessionError::Io(_)), "vu: {err:?}"); +} + +// =========================================================================== +// Non-régression des signatures historiques Result +// =========================================================================== + +/// `drain_with_readiness` (signature historique) : un tour limité reste `Err(Io)` +/// (zéro régression pour l'appelant orchestrateur). +#[tokio::test] +async fn drain_with_readiness_rate_limited_is_io_error() { + let session = FakeSession { + events: vec![ReplyEvent::RateLimited { + resets_at_ms: Some(1), + }], + }; + let mediator = RecordingMediator::default(); + let err = drain_with_readiness(&session, "go", None, &mediator, aid(1)) + .await + .expect_err("limite ⇒ Io sur la signature historique"); + assert!(matches!(err, AgentSessionError::Io(_)), "vu: {err:?}"); +} + +/// `send_blocking` (signature historique) : un tour limité reste `Err(Io)`. +#[tokio::test] +async fn send_blocking_rate_limited_is_io_error() { + let session = FakeSession { + events: vec![ReplyEvent::RateLimited { resets_at_ms: None }], + }; + let err = send_blocking(&session, "go", None) + .await + .expect_err("limite ⇒ Io"); + assert!(matches!(err, AgentSessionError::Io(_)), "vu: {err:?}"); +} + +/// `drain_with_readiness` nominal : `[.., Final]` ⇒ le contenu, et `mark_idle` au Final. +#[tokio::test] +async fn drain_with_readiness_nominal_still_completes() { + let session = FakeSession { + events: vec![ + ReplyEvent::TextDelta { text: "a".into() }, + ReplyEvent::Final { + content: "fini".into(), + }, + ], + }; + let mediator = RecordingMediator::default(); + let content = drain_with_readiness(&session, "go", None, &mediator, aid(1)) + .await + .expect("ok"); + assert_eq!(content, "fini"); + assert!(mediator.calls.lock().unwrap().contains(&"idle")); +}