//! [`TokioScheduler`] — adapter du port [`domain::ports::Scheduler`] (ARCHITECTURE //! §21.4), minuterie **one-shot annulable** en mémoire. //! //! # Rôle //! //! Armer un réveil à une **échéance absolue** (époche-ms) qui, à expiration, **remet //! une [`ScheduledTask`] (donnée pure) au drain applicatif** via un canal `mpsc` //! fourni à la construction — exactement le patron du watcher d'orchestrateur //! (§14.3) : l'adapter dispatche une *donnée*, il **n'exécute jamais** la reprise //! lui-même (aucune dépendance vers l'application). `cancel` désarme un réveil non //! encore tiré (c'est ce qui rend la « reprise auto annulable », §21.1-4). //! //! # En mémoire uniquement (§21.1-3) //! //! Aucune persistance : les réveils armés ne survivent pas à un redémarrage d'IdeA. //! La table `id → JoinHandle` vit derrière un `Mutex` (thread-safe). L'attente réelle //! se fait dans une tâche tokio de fond ; `cancel` l'`abort()`. use std::collections::HashMap; use std::sync::{Arc, Mutex}; use std::time::Duration; use tokio::sync::mpsc::UnboundedSender; use tokio::task::JoinHandle; use domain::ids::ScheduleId; use domain::ports::{Clock, Scheduler, ScheduledTask}; /// Adapter tokio du port [`Scheduler`] (§21.4). Construire avec [`TokioScheduler::new`]. /// /// Thread-safe et clonable-en-`Arc` : injecté au composition root comme /// `Arc` (LS7, hors périmètre de ce lot). pub struct TokioScheduler { /// Canal de **remise** des tâches échues : à l'échéance, la tâche de fond y pousse /// la [`ScheduledTask`]. Le drain (application, LS4) en est l'autre bout. Non borné /// pour ne **jamais** perdre/retarder une reprise (événements rares, basse /// fréquence) ni bloquer la tâche de fond. tx: UnboundedSender, /// Horloge injectée (déterminisme/testabilité) pour traduire l'échéance absolue en /// délai relatif (`deadline_ms - now`). clock: Arc, /// Table `id → handle` des réveils armés non encore tirés. `cancel` y retire + /// `abort()` le handle. handles: Arc>>>, } impl TokioScheduler { /// Construit l'adapter. `tx` est le bout **émetteur** du canal de remise (le /// drain applicatif détient le récepteur) ; `clock` fournit « maintenant » pour /// calculer le délai d'attente. #[must_use] pub fn new(tx: UnboundedSender, clock: Arc) -> Self { Self { tx, clock, handles: Arc::new(Mutex::new(HashMap::new())), } } } impl Scheduler for TokioScheduler { fn arm(&self, deadline_ms: i64, task: ScheduledTask) -> ScheduleId { let id = ScheduleId::new_random(); // Délai relatif : une échéance déjà passée (≤ now) ⇒ délai nul ⇒ déclenchement // au plus tôt (contrat du port). `saturating_sub` + `max(0)` évitent tout // débordement/négatif avant le cast non signé. let now = self.clock.now_millis(); let delay = Duration::from_millis(deadline_ms.saturating_sub(now).max(0) as u64); let tx = self.tx.clone(); let handle = tokio::spawn(async move { tokio::time::sleep(delay).await; // Remet la DONNÉE au drain applicatif (jamais d'exécution ici). Une erreur // d'envoi signifie que le récepteur a été lâché (IdeA s'arrête) ⇒ ignorée. let _ = tx.send(task); }); let mut map = self.handles.lock().expect("scheduler mutex sain"); // Élague au passage les réveils déjà tirés (handles terminés) pour borner la // table sans course : sous le même verrou, `is_finished()` est un fait stable. map.retain(|_, h| !h.is_finished()); map.insert(id, handle); id } fn cancel(&self, id: ScheduleId) -> bool { let mut map = self.handles.lock().expect("scheduler mutex sain"); match map.remove(&id) { // Inconnu (jamais armé, ou déjà élagué après tir) ⇒ rien à annuler. None => false, // Déjà tiré mais encore présent ⇒ pas une annulation (la tâche est partie). Some(handle) if handle.is_finished() => false, // Armé et non tiré ⇒ on l'abort : la reprise n'aura pas lieu. Some(handle) => { handle.abort(); true } } } } #[cfg(test)] mod tests { use std::sync::Arc; use std::time::Duration; use tokio::sync::mpsc::error::TryRecvError; use tokio::sync::mpsc::{unbounded_channel, UnboundedReceiver}; use tokio::time::timeout; use uuid::Uuid; use domain::ids::{AgentId, NodeId, ScheduleId}; use domain::ports::{Clock, Scheduler, ScheduledTask}; use super::TokioScheduler; use crate::SystemClock; /// Délai d'ARMEMENT court (les réveils doivent tomber vite pour ne pas ralentir /// la suite). const SHORT_MS: i64 = 50; /// Timeout de RÉCEPTION généreux (anti-flaky) : on attend bien plus longtemps que /// le délai d'armement pour absorber l'aléa d'ordonnancement sans jamais expirer /// à tort. const RECV_TIMEOUT: Duration = Duration::from_secs(2); /// Construit un scheduler + le bout récepteur du canal de remise + l'horloge réelle /// (partagée, pour calculer des échéances cohérentes avec celle qu'`arm` lira). fn make() -> (TokioScheduler, UnboundedReceiver, Arc) { let (tx, rx) = unbounded_channel(); let clock: Arc = Arc::new(SystemClock::new()); (TokioScheduler::new(tx, clock.clone()), rx, clock) } /// Une `ScheduledTask` identifiable par son `conversation_id` (pour distinguer les /// tâches concurrentes à l'arrivée). fn task_with(conv: &str) -> ScheduledTask { ScheduledTask::ResumeAgent { agent_id: AgentId::from_uuid(Uuid::from_u128(1)), node_id: NodeId::from_uuid(Uuid::from_u128(2)), conversation_id: Some(conv.to_owned()), } } /// `arm` à `now + SHORT_MS` ⇒ la tâche EXACTE arrive (sous timeout de sécurité) ; /// et AVANT l'échéance, le canal est vide (try_recv juste après arm). #[tokio::test] async fn arm_fires_after_deadline_with_exact_task() { let (sched, mut rx, clock) = make(); let t = task_with("c-fire"); sched.arm(clock.now_millis() + SHORT_MS, t.clone()); // Juste après arm : rien n'est encore tiré (échéance dans le futur). assert!( matches!(rx.try_recv(), Err(TryRecvError::Empty)), "aucune tâche ne doit arriver AVANT l'échéance" ); let got = timeout(RECV_TIMEOUT, rx.recv()) .await .expect("le réveil doit tirer sous le timeout de sécurité") .expect("le canal de remise doit rester ouvert"); assert_eq!(got, t, "la tâche reçue doit être exactement celle armée"); } /// `cancel` AVANT l'échéance empêche le tir : on arme un délai court, on annule /// aussitôt (large marge avant les 50 ms), puis on attend 4× l'échéance — la tâche /// NE doit JAMAIS arriver. Prouve une absence qui, SANS l'annulation, se serait /// produite (le délai court aurait tiré bien avant la fin de l'attente). #[tokio::test] async fn cancel_before_deadline_prevents_fire() { let (sched, mut rx, clock) = make(); let id = sched.arm(clock.now_millis() + SHORT_MS, task_with("c-cancel")); assert!(sched.cancel(id), "cancel d'un réveil armé non tiré ⇒ true"); // Attend bien au-delà de l'échéance (4×) : si l'annulation avait échoué, la // tâche serait déjà là. tokio::time::sleep(Duration::from_millis((SHORT_MS as u64) * 4)).await; assert!( matches!(rx.try_recv(), Err(TryRecvError::Empty)), "après cancel, aucune tâche ne doit être remise" ); // L'id a été retiré de la table : un second cancel ne trouve plus rien. assert!(!sched.cancel(id), "un id déjà annulé/retiré ⇒ false"); } /// `cancel` d'un id JAMAIS armé ⇒ false (rien à désarmer). #[tokio::test] async fn cancel_unknown_id_is_false() { let (sched, _rx, _clock) = make(); assert!(!sched.cancel(ScheduleId::new_random())); } /// `cancel` APRÈS le tir ⇒ false : on arme court, on attend la réception (donc le /// tir a eu lieu), puis on annule — la tâche est déjà partie, rien à désarmer. #[tokio::test] async fn cancel_after_fire_is_false() { let (sched, mut rx, clock) = make(); let id = sched.arm(clock.now_millis() + SHORT_MS, task_with("c-after")); let _ = timeout(RECV_TIMEOUT, rx.recv()) .await .expect("doit tirer") .expect("canal ouvert"); // Laisse l'exécuteur finaliser l'état du JoinHandle (is_finished) après l'envoi. for _ in 0..8 { tokio::task::yield_now().await; } assert!( !sched.cancel(id), "annuler un réveil déjà tiré ne désarme rien ⇒ false" ); } /// Échéance DÉJÀ PASSÉE (`now - 1000`) ⇒ tir quasi-immédiat (délai nul) : la tâche /// arrive sous un timeout court. #[tokio::test] async fn past_deadline_fires_immediately() { let (sched, mut rx, clock) = make(); let t = task_with("c-past"); sched.arm(clock.now_millis() - 1000, t.clone()); let got = timeout(Duration::from_millis(500), rx.recv()) .await .expect("une échéance passée doit tirer au plus tôt") .expect("canal ouvert"); assert_eq!(got, t); } /// Plusieurs `arm` concurrents (même échéance) : chacun tire indépendamment, on /// reçoit les TROIS tâches (ensemble, sans dépendre de l'ordre d'arrivée). #[tokio::test] async fn multiple_concurrent_arms_all_fire() { let (sched, mut rx, clock) = make(); let deadline = clock.now_millis() + SHORT_MS; let ids: Vec<_> = ["a", "b", "c"] .iter() .map(|c| sched.arm(deadline, task_with(c))) .collect(); // Les ids armés sont distincts (pas de collision). assert_ne!(ids[0], ids[1]); assert_ne!(ids[1], ids[2]); assert_ne!(ids[0], ids[2]); let mut got = Vec::new(); for _ in 0..3 { let task = timeout(RECV_TIMEOUT, rx.recv()) .await .expect("chaque réveil doit tirer") .expect("canal ouvert"); // `ScheduledTask` est mono-variante : déstructuration directe (pas de `if let`). let ScheduledTask::ResumeAgent { conversation_id, .. } = task; got.push(conversation_id.expect("conv présent")); } got.sort(); assert_eq!(got, vec!["a".to_owned(), "b".to_owned(), "c".to_owned()]); } /// `cancel` ciblé parmi plusieurs réveils concurrents : seul l'annulé ne tire pas ; /// les autres arrivent normalement. Prouve l'indépendance des handles (pas de fuite, /// pas d'annulation collatérale). #[tokio::test] async fn cancel_one_among_many_leaves_others_firing() { let (sched, mut rx, clock) = make(); let deadline = clock.now_millis() + SHORT_MS; let keep1 = sched.arm(deadline, task_with("keep-1")); let drop_id = sched.arm(deadline, task_with("dropped")); let keep2 = sched.arm(deadline, task_with("keep-2")); let _ = (keep1, keep2); assert!(sched.cancel(drop_id), "le réveil ciblé est annulé ⇒ true"); // Les deux survivants arrivent ; l'annulé jamais. let mut got = Vec::new(); for _ in 0..2 { let task = timeout(RECV_TIMEOUT, rx.recv()) .await .expect("les survivants doivent tirer") .expect("canal ouvert"); let ScheduledTask::ResumeAgent { conversation_id, .. } = task; got.push(conversation_id.expect("conv présent")); } got.sort(); assert_eq!(got, vec!["keep-1".to_owned(), "keep-2".to_owned()]); // Plus rien derrière (l'annulé n'a pas tiré). assert!( matches!(rx.try_recv(), Err(TryRecvError::Empty)), "le réveil annulé ne doit jamais être remis" ); } }