feat(session-limits): LS4 — service application + réconciliation T4
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 <noreply@anthropic.com>
This commit is contained in:
@ -10,13 +10,17 @@ mod catalogue;
|
|||||||
mod inspect;
|
mod inspect;
|
||||||
mod lifecycle;
|
mod lifecycle;
|
||||||
mod resume;
|
mod resume;
|
||||||
|
mod session_limit;
|
||||||
mod structured;
|
mod structured;
|
||||||
mod usecases;
|
mod usecases;
|
||||||
|
|
||||||
pub(crate) use lifecycle::unique_md_path;
|
pub(crate) use lifecycle::unique_md_path;
|
||||||
pub(crate) use lifecycle::ReattachDecision;
|
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 catalogue::{reference_profile_id, reference_profiles, selectable_reference_profiles};
|
||||||
pub use inspect::{InspectConversation, InspectConversationInput, InspectConversationOutput};
|
pub use inspect::{InspectConversation, InspectConversationInput, InspectConversationOutput};
|
||||||
|
|||||||
229
crates/application/src/agent/session_limit.rs
Normal file
229
crates/application/src/agent/session_limit.rs
Normal file
@ -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<String>,
|
||||||
|
resume_prompt: &str,
|
||||||
|
) -> Result<(), AppError>;
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Service d'orchestration des limites de session (§21.5).
|
||||||
|
pub struct SessionLimitService {
|
||||||
|
clock: Arc<dyn Clock>,
|
||||||
|
scheduler: Arc<dyn Scheduler>,
|
||||||
|
events: Arc<dyn EventBus>,
|
||||||
|
resumer: Arc<dyn AgentResumer>,
|
||||||
|
/// Reprises **armées** non encore tirées : `agent_id → ScheduleId` (en mémoire).
|
||||||
|
armed: Mutex<HashMap<AgentId, ScheduleId>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl SessionLimitService {
|
||||||
|
/// Construit le service depuis ses ports injectés (composition root).
|
||||||
|
#[must_use]
|
||||||
|
pub fn new(
|
||||||
|
clock: Arc<dyn Clock>,
|
||||||
|
scheduler: Arc<dyn Scheduler>,
|
||||||
|
events: Arc<dyn EventBus>,
|
||||||
|
resumer: Arc<dyn AgentResumer>,
|
||||||
|
) -> 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<String>,
|
||||||
|
resets_at_ms: Option<i64>,
|
||||||
|
) {
|
||||||
|
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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -19,6 +19,31 @@ use domain::ids::AgentId;
|
|||||||
use domain::ports::{AgentSession, AgentSessionError, ReplyEvent};
|
use domain::ports::{AgentSession, AgentSessionError, ReplyEvent};
|
||||||
use domain::readiness::{ReadinessPolicy, ReadinessSignal};
|
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<i64>,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
/// Envoie `prompt` à la session vivante puis **draine le flux de réponse jusqu'au
|
/// Envoie `prompt` à la session vivante puis **draine le flux de réponse jusqu'au
|
||||||
/// [`ReplyEvent::Final`]**, et retourne son contenu agrégé.
|
/// [`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`
|
/// - [`AgentSessionError::Io`]/[`AgentSessionError::Decode`] remontées par `send`
|
||||||
/// (échec de communication / décodage de la sortie structurée) ;
|
/// (échec de communication / décodage de la sortie structurée) ;
|
||||||
/// - [`AgentSessionError::Io`] si le flux se termine **sans** `Final` (tour
|
/// - [`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`.
|
/// - [`AgentSessionError::Timeout`] si `timeout` expire avant le `Final`.
|
||||||
pub async fn send_blocking(
|
pub async fn send_blocking(
|
||||||
session: &dyn AgentSession,
|
session: &dyn AgentSession,
|
||||||
prompt: &str,
|
prompt: &str,
|
||||||
timeout: Option<Duration>,
|
timeout: Option<Duration>,
|
||||||
) -> Result<String, AgentSessionError> {
|
) -> Result<String, AgentSessionError> {
|
||||||
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
|
/// Comme [`send_blocking`], mais **branche la readiness** : à chaque événement du
|
||||||
@ -61,7 +95,9 @@ pub async fn send_blocking(
|
|||||||
///
|
///
|
||||||
/// # Errors
|
/// # Errors
|
||||||
/// Identiques à [`send_blocking`] (échec `send`/décodage, flux clos sans `Final`,
|
/// 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<String>`, zéro régression pour l'appelant
|
||||||
|
/// orchestrateur) ; utilise [`drain_with_readiness_outcome`] pour exploiter la limite.
|
||||||
pub async fn drain_with_readiness(
|
pub async fn drain_with_readiness(
|
||||||
session: &dyn AgentSession,
|
session: &dyn AgentSession,
|
||||||
prompt: &str,
|
prompt: &str,
|
||||||
@ -69,6 +105,38 @@ pub async fn drain_with_readiness(
|
|||||||
mediator: &dyn InputMediator,
|
mediator: &dyn InputMediator,
|
||||||
agent: AgentId,
|
agent: AgentId,
|
||||||
) -> Result<String, AgentSessionError> {
|
) -> Result<String, AgentSessionError> {
|
||||||
|
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<Duration>,
|
||||||
|
mediator: &dyn InputMediator,
|
||||||
|
agent: AgentId,
|
||||||
|
) -> Result<TurnOutcome, AgentSessionError> {
|
||||||
// `on_signal` ne reçoit QUE les événements terminaux (le `Final` ⇒ `TurnEnded`) :
|
// `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é
|
// 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
|
// (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| {
|
|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 {
|
if signal == ReadinessSignal::TurnEnded {
|
||||||
mediator.mark_idle(agent);
|
mediator.mark_idle(agent);
|
||||||
}
|
}
|
||||||
@ -106,7 +176,7 @@ async fn drain_bounded_events(
|
|||||||
timeout: Option<Duration>,
|
timeout: Option<Duration>,
|
||||||
on_event: impl FnMut(&ReplyEvent),
|
on_event: impl FnMut(&ReplyEvent),
|
||||||
on_signal: impl FnMut(ReadinessSignal),
|
on_signal: impl FnMut(ReadinessSignal),
|
||||||
) -> Result<String, AgentSessionError> {
|
) -> Result<TurnOutcome, AgentSessionError> {
|
||||||
match timeout {
|
match timeout {
|
||||||
Some(dur) => match tokio::time::timeout(
|
Some(dur) => match tokio::time::timeout(
|
||||||
dur,
|
dur,
|
||||||
@ -126,18 +196,27 @@ async fn drain_bounded_events(
|
|||||||
///
|
///
|
||||||
/// Le flux ([`domain::ports::ReplyStream`]) est un itérateur synchrone et borné :
|
/// 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
|
/// 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
|
/// jusqu'à rencontrer le `Final` (et on retourne son contenu) ; si le flux s'épuise
|
||||||
/// s'épuise avant, c'est un tour interrompu → erreur [`AgentSessionError::Io`].
|
/// 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
|
/// Chaque événement est classé par [`ReadinessPolicy`] et le signal éventuel est
|
||||||
/// remonté à `on_signal` (le `Final` ⇒ [`ReadinessSignal::TurnEnded`]). Deltas,
|
/// remonté à `on_signal` (le `Final` ⇒ [`ReadinessSignal::TurnEnded`] ; un
|
||||||
/// activités et heartbeats sont non terminaux ⇒ ignorés par le rendez-vous synchrone.
|
/// `RateLimited` ⇒ [`ReadinessSignal::RateLimited`]). Deltas, activités et heartbeats
|
||||||
|
/// sont non terminaux ⇒ ignorés par le rendez-vous synchrone.
|
||||||
async fn drain_to_final(
|
async fn drain_to_final(
|
||||||
session: &dyn AgentSession,
|
session: &dyn AgentSession,
|
||||||
prompt: &str,
|
prompt: &str,
|
||||||
mut on_event: impl FnMut(&ReplyEvent),
|
mut on_event: impl FnMut(&ReplyEvent),
|
||||||
mut on_signal: impl FnMut(ReadinessSignal),
|
mut on_signal: impl FnMut(ReadinessSignal),
|
||||||
) -> Result<String, AgentSessionError> {
|
) -> Result<TurnOutcome, AgentSessionError> {
|
||||||
|
// 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<Option<i64>> = None;
|
||||||
let stream = session.send(prompt).await?;
|
let stream = session.send(prompt).await?;
|
||||||
for event in stream {
|
for event in stream {
|
||||||
// Battement de vivacité (lot 2) : notifié pour CHAQUE événement brut, avant le
|
// 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) {
|
if let Some(signal) = ReadinessPolicy::classify(&event) {
|
||||||
on_signal(signal);
|
on_signal(signal);
|
||||||
}
|
}
|
||||||
if let ReplyEvent::Final { content } = event {
|
match event {
|
||||||
return Ok(content);
|
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.
|
// TextDelta / ToolActivity / Heartbeat : non terminaux, ignorés ici.
|
||||||
|
ReplyEvent::TextDelta { .. }
|
||||||
|
| ReplyEvent::ToolActivity { .. }
|
||||||
|
| ReplyEvent::Heartbeat => {}
|
||||||
}
|
}
|
||||||
Err(AgentSessionError::Io(
|
}
|
||||||
|
// 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(),
|
"le flux de réponse s'est terminé sans événement Final".to_string(),
|
||||||
))
|
)),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
|
|||||||
@ -29,8 +29,9 @@ pub mod terminal;
|
|||||||
pub mod window;
|
pub mod window;
|
||||||
|
|
||||||
pub use agent::{
|
pub use agent::{
|
||||||
drain_with_readiness, reference_profile_id, reference_profiles, selectable_reference_profiles,
|
drain_with_readiness, drain_with_readiness_outcome, reference_profile_id, reference_profiles,
|
||||||
send_blocking, ChangeAgentProfile, ChangeAgentProfileInput, ChangeAgentProfileOutput,
|
selectable_reference_profiles, send_blocking, AgentResumer, ChangeAgentProfile,
|
||||||
|
ChangeAgentProfileInput, ChangeAgentProfileOutput,
|
||||||
ConfigureProfiles, ConfigureProfilesInput, ConfigureProfilesOutput, CreateAgentFromScratch,
|
ConfigureProfiles, ConfigureProfilesInput, ConfigureProfilesOutput, CreateAgentFromScratch,
|
||||||
CreateAgentInput, CreateAgentOutput, DeleteAgent, DeleteAgentInput, DeleteProfile,
|
CreateAgentInput, CreateAgentOutput, DeleteAgent, DeleteAgentInput, DeleteProfile,
|
||||||
DeleteProfileInput, DetectProfiles, DetectProfilesInput, DetectProfilesOutput, FirstRunState,
|
DeleteProfileInput, DetectProfiles, DetectProfilesInput, DetectProfilesOutput, FirstRunState,
|
||||||
@ -41,8 +42,8 @@ pub use agent::{
|
|||||||
ProfileAvailability,
|
ProfileAvailability,
|
||||||
ProviderSessionProvider, ReadAgentContext, ReadAgentContextInput, ReadAgentContextOutput,
|
ProviderSessionProvider, ReadAgentContext, ReadAgentContextInput, ReadAgentContextOutput,
|
||||||
ReferenceProfiles, ReferenceProfilesOutput, ResumableAgent, SaveProfile, SaveProfileInput,
|
ReferenceProfiles, ReferenceProfilesOutput, ResumableAgent, SaveProfile, SaveProfileInput,
|
||||||
SaveProfileOutput, StructuredSessionDescriptor, UpdateAgentContext, UpdateAgentContextInput,
|
SaveProfileOutput, SessionLimitService, StructuredSessionDescriptor, TurnOutcome,
|
||||||
AGENT_MEMORY_RECALL_BUDGET,
|
UpdateAgentContext, UpdateAgentContextInput, AGENT_MEMORY_RECALL_BUDGET, RESUME_PROMPT,
|
||||||
};
|
};
|
||||||
pub use conversation::RecordTurn;
|
pub use conversation::RecordTurn;
|
||||||
pub use embedder::{
|
pub use embedder::{
|
||||||
|
|||||||
440
crates/application/tests/session_limit_service.rs
Normal file
440
crates/application/tests/session_limit_service.rs
Normal file
@ -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<Mutex<Vec<DomainEvent>>>);
|
||||||
|
impl SpyBus {
|
||||||
|
fn events(&self) -> Vec<DomainEvent> {
|
||||||
|
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<Mutex<Vec<(i64, ScheduledTask)>>>,
|
||||||
|
issued: Arc<Mutex<Vec<ScheduleId>>>,
|
||||||
|
cancels: Arc<Mutex<Vec<ScheduleId>>>,
|
||||||
|
next_id: Arc<Mutex<u128>>,
|
||||||
|
cancel_result: Arc<AtomicBool>,
|
||||||
|
}
|
||||||
|
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<ScheduleId> {
|
||||||
|
self.issued.lock().unwrap().clone()
|
||||||
|
}
|
||||||
|
fn cancels(&self) -> Vec<ScheduleId> {
|
||||||
|
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<Mutex<Vec<(AgentId, NodeId, Option<String>, String)>>>,
|
||||||
|
fail: Arc<AtomicBool>,
|
||||||
|
}
|
||||||
|
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>, 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<String>,
|
||||||
|
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é"
|
||||||
|
);
|
||||||
|
}
|
||||||
205
crates/application/tests/session_limit_t4.rs
Normal file
205
crates/application/tests/session_limit_t4.rs
Normal file
@ -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<String>` :
|
||||||
|
//! - `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<ReplyEvent>,
|
||||||
|
}
|
||||||
|
#[async_trait]
|
||||||
|
impl AgentSession for FakeSession {
|
||||||
|
fn id(&self) -> SessionId {
|
||||||
|
SessionId::from_uuid(Uuid::from_u128(1))
|
||||||
|
}
|
||||||
|
fn conversation_id(&self) -> Option<String> {
|
||||||
|
None
|
||||||
|
}
|
||||||
|
async fn send(&self, _prompt: &str) -> Result<ReplyStream, AgentSessionError> {
|
||||||
|
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<Vec<&'static str>>,
|
||||||
|
}
|
||||||
|
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<String>
|
||||||
|
// ===========================================================================
|
||||||
|
|
||||||
|
/// `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"));
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user