Redémarre le ticket #4 sur la base develop. Émission d'annonces live d'un agent vers l'UI, distinctes du Final de délégation : - domain: ReplyEvent::Announcement/Final et DomainEvent::AgentAnnouncement (events.rs), port d'émission (ports.rs), gating de readiness (readiness.rs). - application: mapping des événements structurés en annonces (agent/structured.rs, agent/mod.rs, lib.rs) et relais côté orchestrateur (orchestrator/service.rs). - infrastructure/session: parse des annonces + fix du Final pour Claude et Codex, propagé aux adaptateurs et à la conformance (claude.rs, codex.rs, conformance.rs, mod.rs, process.rs, sandbox_e2e.rs). - app-tauri: relais Tauri des annonces vers le front (events.rs, chat.rs). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
553 lines
22 KiB
Rust
553 lines
22 KiB
Rust
//! Helper applicatif `send_blocking` — le rendez-vous synchrone inter-agents
|
|
//! au-dessus du port [`AgentSession`] (ARCHITECTURE §17.1 / §17.4).
|
|
//!
|
|
//! `AgentSession::send` retourne un **flux** d'événements (`ReplyStream`), à la
|
|
//! manière de `PtyPort::subscribe_output`, mais **typé** : deltas de texte →
|
|
//! activités d'outil → **un** événement terminal déterministe
|
|
//! [`ReplyEvent::Final`]. Le rendez-vous synchrone dont l'orchestrateur a besoin
|
|
//! (§17.4) s'obtient en **drainant ce flux jusqu'au `Final`** : c'est *la* primitive
|
|
//! de la messagerie inter-agents, **sans outbox, sans corrélation fichier** (le
|
|
//! `Final` *est* la fin de tour).
|
|
//!
|
|
//! DRY : un seul chemin de lecture (le flux). `send_blocking` n'est qu'un *consom-
|
|
//! mateur* du même flux que la cellule chat utilise pour le rendu incrémental.
|
|
|
|
use std::sync::Arc;
|
|
use std::time::{Duration, SystemTime, UNIX_EPOCH};
|
|
|
|
use domain::conversation::ConversationParty;
|
|
use domain::events::DomainEvent;
|
|
use domain::ids::AgentId;
|
|
use domain::input::InputMediator;
|
|
use domain::mailbox::TicketId;
|
|
use domain::ports::{AgentSession, AgentSessionError, EventBus, ReplyEvent, ReplyStream};
|
|
use domain::readiness::{ReadinessPolicy, ReadinessSignal};
|
|
use domain::ProjectId;
|
|
|
|
/// 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>,
|
|
},
|
|
}
|
|
|
|
/// Attribution applicative des annonces live d'un tour inter-agent structuré.
|
|
#[derive(Clone)]
|
|
pub struct AnnouncementPublisher {
|
|
/// Bus d'événements domaine.
|
|
pub bus: Arc<dyn EventBus>,
|
|
/// Projet hôte du rendez-vous.
|
|
pub project_id: ProjectId,
|
|
/// Partie qui a demandé le tour.
|
|
pub requester: ConversationParty,
|
|
/// Agent cible qui produit les annonces.
|
|
pub target: AgentId,
|
|
/// Ticket FIFO corrélant le tour.
|
|
pub ticket: TicketId,
|
|
}
|
|
|
|
impl AnnouncementPublisher {
|
|
fn publish_event(&self, event: &ReplyEvent) {
|
|
if let ReplyEvent::Announcement { text } = event {
|
|
self.bus.publish(DomainEvent::AgentAnnouncement {
|
|
project_id: self.project_id,
|
|
requester: self.requester,
|
|
target: self.target,
|
|
ticket: self.ticket,
|
|
text: text.clone(),
|
|
at_ms: now_epoch_ms(),
|
|
});
|
|
}
|
|
}
|
|
}
|
|
|
|
fn now_epoch_ms() -> u64 {
|
|
SystemTime::now()
|
|
.duration_since(UNIX_EPOCH)
|
|
.unwrap_or_default()
|
|
.as_millis()
|
|
.try_into()
|
|
.unwrap_or(u64::MAX)
|
|
}
|
|
|
|
/// Envoie `prompt` à la session vivante puis **draine le flux de réponse jusqu'au
|
|
/// [`ReplyEvent::Final`]**, et retourne son contenu agrégé.
|
|
///
|
|
/// C'est le rendez-vous synchrone (§17.4) : on attend que le tour soit
|
|
/// déterministiquement terminé (`Final`) avant de rendre la main. Les deltas de
|
|
/// texte et les activités d'outil traversés en chemin sont **ignorés** ici (ils
|
|
/// servent le rendu incrémental côté UI, pas l'appelant synchrone).
|
|
///
|
|
/// `timeout`, lorsqu'il est fourni, borne l'attente : si aucun `Final` n'est observé
|
|
/// dans le délai, on retourne [`AgentSessionError::Timeout`] **sans tuer la
|
|
/// session** (elle reste vivante dans le registre ; l'appelant décide de la suite).
|
|
/// `None` ⇒ pas de borne temporelle (on attend la fin du tour).
|
|
///
|
|
/// # Errors
|
|
/// - [`AgentSessionError::Io`]/[`AgentSessionError::Decode`] remontées par `send`
|
|
/// (échec de communication / décodage de la sortie structurée) ;
|
|
/// - [`AgentSessionError::Io`] si le flux se termine **sans** `Final` (tour
|
|
/// interrompu) — y compris un tour clos sur une **limite de session** (le
|
|
/// rendez-vous synchrone n'a pas de contenu à rendre ; cf. [`TurnOutcome`] et la
|
|
/// variante riche [`drain_with_readiness_outcome`] pour exploiter la limite) ;
|
|
/// - [`AgentSessionError::Timeout`] si `timeout` expire avant le `Final`.
|
|
pub async fn send_blocking(
|
|
session: &dyn AgentSession,
|
|
prompt: &str,
|
|
timeout: Option<Duration>,
|
|
) -> Result<String, AgentSessionError> {
|
|
match drain_bounded_events(session, prompt, timeout, |_event| {}, |_signal| {}).await? {
|
|
TurnOutcome::Completed(content) => Ok(content),
|
|
// Le rendez-vous synchrone (ask) attend un contenu : un tour limité n'en a pas
|
|
// ⇒ on conserve le comportement historique (erreur), sans casser le contrat.
|
|
TurnOutcome::RateLimited { .. } => Err(AgentSessionError::Io(
|
|
"le tour s'est clos en limite de session, sans contenu Final".to_string(),
|
|
)),
|
|
}
|
|
}
|
|
|
|
/// Comme [`send_blocking`], mais **branche la readiness** : à chaque événement du
|
|
/// tour, [`ReadinessPolicy::classify`] est consulté et, dès qu'il renvoie
|
|
/// [`ReadinessSignal::TurnEnded`] (le `Final`), le médiateur d'entrée est notifié
|
|
/// (`mark_idle(agent)`) pour que la FIFO de l'agent avance — **sans** dépendre d'un
|
|
/// `idea_reply` explicite ni du sniff de prompt PTY (chantier readiness/heartbeat,
|
|
/// lot 1, fix de la cause racine du blocage `Busy`).
|
|
///
|
|
/// DRY : **un seul** chemin de lecture du flux (la boucle de [`drain_bounded`]) ;
|
|
/// cette fonction n'est que `send_blocking` muni d'un *sink* de readiness. Le `Final`
|
|
/// réveille donc à la fois le `pending` (via la valeur de retour) **et** la FIFO (via
|
|
/// `mark_idle`). `idea_reply` reste un signal alternatif (premier arrivé gagne) côté
|
|
/// orchestrateur.
|
|
///
|
|
/// # Errors
|
|
/// Identiques à [`send_blocking`] (échec `send`/décodage, flux clos sans `Final`,
|
|
/// timeout). Un tour clos sur une **limite de session** ⇒ [`AgentSessionError::Io`]
|
|
/// **ici** (signature historique `Result<String>`, zéro régression pour l'appelant
|
|
/// orchestrateur) ; utilise [`drain_with_readiness_outcome`] pour exploiter la limite.
|
|
pub async fn drain_with_readiness(
|
|
session: &dyn AgentSession,
|
|
prompt: &str,
|
|
timeout: Option<Duration>,
|
|
mediator: &dyn InputMediator,
|
|
agent: AgentId,
|
|
) -> 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(),
|
|
)),
|
|
}
|
|
}
|
|
|
|
/// Comme [`drain_with_readiness`], mais publie les [`ReplyEvent::Announcement`]
|
|
/// produits live par la session avec l'attribution applicative complète.
|
|
pub async fn drain_with_readiness_and_announcements(
|
|
session: &dyn AgentSession,
|
|
prompt: &str,
|
|
timeout: Option<Duration>,
|
|
mediator: &dyn InputMediator,
|
|
agent: AgentId,
|
|
announcements: Option<AnnouncementPublisher>,
|
|
) -> Result<String, AgentSessionError> {
|
|
let Some(publisher) = announcements else {
|
|
return drain_with_readiness(session, prompt, timeout, mediator, agent).await;
|
|
};
|
|
|
|
let (tap_tx, tap_rx) = std::sync::mpsc::channel();
|
|
let live_publisher = publisher.clone();
|
|
let live_pump = std::thread::spawn(move || {
|
|
for event in tap_rx {
|
|
live_publisher.publish_event(&event);
|
|
}
|
|
});
|
|
|
|
let stream_result = match timeout {
|
|
Some(dur) => tokio::time::timeout(dur, session.send_with_tap(prompt, tap_tx))
|
|
.await
|
|
.map_err(|_elapsed| AgentSessionError::Timeout)?,
|
|
None => session.send_with_tap(prompt, tap_tx).await,
|
|
};
|
|
let _ = live_pump.join();
|
|
let stream = stream_result?;
|
|
|
|
match drain_stream_to_final(
|
|
stream,
|
|
|event| {
|
|
if !matches!(event, ReplyEvent::Final { .. }) {
|
|
mediator.mark_alive(agent);
|
|
}
|
|
publisher.publish_event(event);
|
|
},
|
|
|signal| {
|
|
if signal == ReadinessSignal::TurnEnded {
|
|
mediator.mark_idle(agent);
|
|
}
|
|
},
|
|
)? {
|
|
TurnOutcome::Completed(content) => Ok(content),
|
|
TurnOutcome::RateLimited { .. } => Err(AgentSessionError::Io(
|
|
"le tour s'est clos en limite de session, sans contenu Final".to_string(),
|
|
)),
|
|
}
|
|
}
|
|
|
|
/// Draine un flux de réponse déjà ouvert, avec la même politique de readiness que
|
|
/// [`drain_with_readiness`].
|
|
///
|
|
/// Ce point d'entrée sert aux appelants qui doivent effectuer une action durable
|
|
/// après `AgentSession::send` réussi mais avant le drain complet du tour.
|
|
pub async fn drain_reply_stream_with_readiness(
|
|
stream: ReplyStream,
|
|
mediator: &dyn InputMediator,
|
|
agent: AgentId,
|
|
) -> Result<String, AgentSessionError> {
|
|
match drain_stream_to_final(
|
|
stream,
|
|
|event| {
|
|
if !matches!(event, ReplyEvent::Final { .. }) {
|
|
mediator.mark_alive(agent);
|
|
}
|
|
},
|
|
|signal| {
|
|
if signal == ReadinessSignal::TurnEnded {
|
|
mediator.mark_idle(agent);
|
|
}
|
|
},
|
|
)? {
|
|
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`) :
|
|
// la readiness ne classe pas les non-terminaux. Pour le **battement** de vivacité
|
|
// (lot 2) on a besoin de notifier le médiateur à CHAQUE événement non terminal
|
|
// (delta / activité / heartbeat) ⇒ on passe un sink d'événement bruts `on_event`.
|
|
drain_bounded_events(
|
|
session,
|
|
prompt,
|
|
timeout,
|
|
|event| {
|
|
// Tout événement **non terminal** prouve la vivacité ⇒ un battement.
|
|
if !matches!(event, ReplyEvent::Final { .. }) {
|
|
mediator.mark_alive(agent);
|
|
}
|
|
},
|
|
|signal| {
|
|
// Seul le `Final` marque `Idle` : un `RateLimited` ne fait PAS avancer la
|
|
// FIFO (la reprise est gérée par le service de limite, §21.5).
|
|
if signal == ReadinessSignal::TurnEnded {
|
|
mediator.mark_idle(agent);
|
|
}
|
|
},
|
|
)
|
|
.await
|
|
}
|
|
|
|
/// Ouvre le flux du tour (`send`) et le **draine jusqu'au `Final`**, en appliquant
|
|
/// la borne temporelle `timeout`, en notifiant `on_event` à **chaque** événement brut
|
|
/// (pour le battement de vivacité, lot 2) et `on_signal` à chaque [`ReadinessSignal`]
|
|
/// dérivé par [`ReadinessPolicy`] (le `Final` ⇒ `TurnEnded`).
|
|
///
|
|
/// **Chemin de lecture unique** (DRY) : `send_blocking` et `drain_with_readiness`
|
|
/// passent tous deux par ici, en différant seulement par leurs *sinks*. La session
|
|
/// **reste vivante** sur timeout (on ne `shutdown` rien ici, §17.1).
|
|
async fn drain_bounded_events(
|
|
session: &dyn AgentSession,
|
|
prompt: &str,
|
|
timeout: Option<Duration>,
|
|
on_event: impl FnMut(&ReplyEvent),
|
|
on_signal: impl FnMut(ReadinessSignal),
|
|
) -> Result<TurnOutcome, AgentSessionError> {
|
|
match timeout {
|
|
Some(dur) => {
|
|
match tokio::time::timeout(dur, drain_to_final(session, prompt, on_event, on_signal))
|
|
.await
|
|
{
|
|
Ok(result) => result,
|
|
// La session **reste vivante** : on ne `shutdown` rien ici (§17.1).
|
|
Err(_elapsed) => Err(AgentSessionError::Timeout),
|
|
}
|
|
}
|
|
None => drain_to_final(session, prompt, on_event, on_signal).await,
|
|
}
|
|
}
|
|
|
|
/// Ouvre le flux du tour (`send`) et le **draine jusqu'au `Final`**.
|
|
///
|
|
/// Le flux ([`domain::ports::ReplyStream`]) est un itérateur synchrone et borné :
|
|
/// après le `Final` il ne produit plus rien. On le parcourt donc simplement
|
|
/// jusqu'à rencontrer le `Final` (et on retourne son contenu) ; si le flux s'épuise
|
|
/// avant, l'issue dépend de ce qu'on a vu (réconciliation §21.2-T4) :
|
|
/// - un [`ReplyEvent::RateLimited`] a été observé ⇒ fin **gracieuse**
|
|
/// [`TurnOutcome::RateLimited`] (l'agent est limité, pas en erreur) ;
|
|
/// - sinon ⇒ tour réellement interrompu → erreur [`AgentSessionError::Io`]
|
|
/// (comportement **inchangé**).
|
|
///
|
|
/// Chaque événement est classé par [`ReadinessPolicy`] et le signal éventuel est
|
|
/// remonté à `on_signal` (le `Final` ⇒ [`ReadinessSignal::TurnEnded`] ; un
|
|
/// `RateLimited` ⇒ [`ReadinessSignal::RateLimited`]). Deltas, activités et heartbeats
|
|
/// sont non terminaux ⇒ ignorés par le rendez-vous synchrone.
|
|
async fn drain_to_final(
|
|
session: &dyn AgentSession,
|
|
prompt: &str,
|
|
on_event: impl FnMut(&ReplyEvent),
|
|
on_signal: impl FnMut(ReadinessSignal),
|
|
) -> Result<TurnOutcome, AgentSessionError> {
|
|
let stream = session.send(prompt).await?;
|
|
drain_stream_to_final(stream, on_event, on_signal)
|
|
}
|
|
|
|
fn drain_stream_to_final(
|
|
stream: ReplyStream,
|
|
mut on_event: impl FnMut(&ReplyEvent),
|
|
mut on_signal: impl FnMut(ReadinessSignal),
|
|
) -> 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;
|
|
for event in stream {
|
|
// Battement de vivacité (lot 2) : notifié pour CHAQUE événement brut, avant le
|
|
// classement readiness. Le sink décide (les non-terminaux prouvent la vivacité).
|
|
on_event(&event);
|
|
if let Some(signal) = ReadinessPolicy::classify(&event) {
|
|
on_signal(signal);
|
|
}
|
|
match event {
|
|
ReplyEvent::Final { content } => return Ok(TurnOutcome::Completed(content)),
|
|
ReplyEvent::RateLimited { resets_at_ms } => last_rate_limit = Some(resets_at_ms),
|
|
// TextDelta / Announcement / ToolActivity / Heartbeat : non terminaux.
|
|
ReplyEvent::TextDelta { .. }
|
|
| ReplyEvent::Announcement { .. }
|
|
| ReplyEvent::ToolActivity { .. }
|
|
| ReplyEvent::Heartbeat => {}
|
|
}
|
|
}
|
|
// Flux clos sans `Final` : fin gracieuse SI une limite a été vue (§21.2-T4),
|
|
// sinon erreur comme avant.
|
|
match last_rate_limit {
|
|
Some(resets_at_ms) => Ok(TurnOutcome::RateLimited { resets_at_ms }),
|
|
None => Err(AgentSessionError::Io(
|
|
"le flux de réponse s'est terminé sans événement Final".to_string(),
|
|
)),
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use std::sync::Mutex;
|
|
|
|
use domain::ids::SessionId;
|
|
use domain::input::{AgentBusyState, SubmitConfig};
|
|
use domain::mailbox::{PendingReply, Ticket};
|
|
use domain::ports::PtyHandle;
|
|
|
|
fn agent(n: u128) -> AgentId {
|
|
AgentId::from_uuid(uuid::Uuid::from_u128(n))
|
|
}
|
|
|
|
/// Session factice : `send` rejoue une liste fixe d'événements (terminée par un
|
|
/// `Final`).
|
|
struct FakeSession {
|
|
events: Vec<ReplyEvent>,
|
|
}
|
|
#[async_trait::async_trait]
|
|
impl AgentSession for FakeSession {
|
|
fn id(&self) -> SessionId {
|
|
SessionId::from_uuid(uuid::Uuid::from_u128(1))
|
|
}
|
|
fn conversation_id(&self) -> Option<String> {
|
|
None
|
|
}
|
|
async fn send(
|
|
&self,
|
|
_prompt: &str,
|
|
) -> Result<domain::ports::ReplyStream, AgentSessionError> {
|
|
Ok(Box::new(self.events.clone().into_iter()))
|
|
}
|
|
async fn shutdown(&self) -> Result<(), AgentSessionError> {
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
/// Médiateur factice qui enregistre l'ordre des `mark_alive` / `mark_idle`.
|
|
#[derive(Default)]
|
|
struct RecordingMediator {
|
|
calls: Mutex<Vec<&'static str>>,
|
|
}
|
|
impl InputMediator for RecordingMediator {
|
|
fn enqueue(&self, _agent: AgentId, _ticket: Ticket) -> PendingReply {
|
|
unreachable!("non utilisé par drain_with_readiness")
|
|
}
|
|
fn preempt(&self, _agent: AgentId) {}
|
|
fn mark_idle(&self, _agent: AgentId) {
|
|
self.calls.lock().unwrap().push("idle");
|
|
}
|
|
fn mark_alive(&self, _agent: AgentId) {
|
|
self.calls.lock().unwrap().push("alive");
|
|
}
|
|
fn busy_state(&self, _agent: AgentId) -> AgentBusyState {
|
|
AgentBusyState::Idle
|
|
}
|
|
fn bind_handle(&self, _agent: AgentId, _handle: PtyHandle) {}
|
|
fn bind_handle_with_submit(
|
|
&self,
|
|
_agent: AgentId,
|
|
_handle: PtyHandle,
|
|
_submit: SubmitConfig,
|
|
) {
|
|
}
|
|
}
|
|
|
|
#[derive(Default)]
|
|
struct RecordingBus {
|
|
events: Mutex<Vec<DomainEvent>>,
|
|
}
|
|
impl EventBus for RecordingBus {
|
|
fn publish(&self, event: DomainEvent) {
|
|
self.events.lock().unwrap().push(event);
|
|
}
|
|
|
|
fn subscribe(&self) -> domain::ports::EventStream {
|
|
Box::new(std::iter::empty())
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn drain_marks_alive_on_each_non_terminal_then_idle_on_final() {
|
|
let session = FakeSession {
|
|
events: vec![
|
|
ReplyEvent::TextDelta { text: "a".into() },
|
|
ReplyEvent::ToolActivity {
|
|
label: "lit".into(),
|
|
},
|
|
ReplyEvent::Heartbeat,
|
|
ReplyEvent::Final {
|
|
content: "fini".into(),
|
|
},
|
|
],
|
|
};
|
|
let mediator = RecordingMediator::default();
|
|
let out = drain_with_readiness(&session, "go", None, &mediator, agent(1))
|
|
.await
|
|
.expect("drain ok");
|
|
assert_eq!(out, "fini");
|
|
// Trois battements (delta, activité, heartbeat) PUIS l'idle sur le Final.
|
|
assert_eq!(
|
|
*mediator.calls.lock().unwrap(),
|
|
vec!["alive", "alive", "alive", "idle"],
|
|
"un battement par événement non terminal, idle au Final (pas de battement sur le Final)"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn drain_publishes_announcement_but_only_final_marks_idle_and_resolves() {
|
|
let session = FakeSession {
|
|
events: vec![
|
|
ReplyEvent::Announcement {
|
|
text: "je travaille".into(),
|
|
},
|
|
ReplyEvent::Final {
|
|
content: "fini".into(),
|
|
},
|
|
],
|
|
};
|
|
let mediator = RecordingMediator::default();
|
|
let bus = Arc::new(RecordingBus::default());
|
|
let project_id = ProjectId::from_uuid(uuid::Uuid::from_u128(10));
|
|
let requester = ConversationParty::User;
|
|
let target = agent(2);
|
|
let ticket = TicketId::from_uuid(uuid::Uuid::from_u128(11));
|
|
|
|
let out = drain_with_readiness_and_announcements(
|
|
&session,
|
|
"go",
|
|
None,
|
|
&mediator,
|
|
target,
|
|
Some(AnnouncementPublisher {
|
|
bus: bus.clone(),
|
|
project_id,
|
|
requester,
|
|
target,
|
|
ticket,
|
|
}),
|
|
)
|
|
.await
|
|
.expect("drain ok");
|
|
|
|
assert_eq!(out, "fini");
|
|
assert_eq!(
|
|
*mediator.calls.lock().unwrap(),
|
|
vec!["alive", "idle"],
|
|
"Announcement est non terminale ; seul Final marque Idle"
|
|
);
|
|
let events = bus.events.lock().unwrap();
|
|
assert_eq!(events.len(), 1);
|
|
match &events[0] {
|
|
DomainEvent::AgentAnnouncement {
|
|
project_id: p,
|
|
requester: r,
|
|
target: t,
|
|
ticket: tk,
|
|
text,
|
|
at_ms,
|
|
} => {
|
|
assert_eq!(*p, project_id);
|
|
assert_eq!(*r, requester);
|
|
assert_eq!(*t, target);
|
|
assert_eq!(*tk, ticket);
|
|
assert_eq!(text, "je travaille");
|
|
assert!(*at_ms > 0);
|
|
}
|
|
other => panic!("événement inattendu: {other:?}"),
|
|
}
|
|
}
|
|
}
|