feat(session-limits): LS3 — port Scheduler + adapter TokioScheduler
Introduit l'abstraction de planification pour la reprise différée à la levée d'une limite de session : - domaine : trait Scheduler + enum ScheduledTask (ports.rs), ScheduleId via typed_id! (ids.rs), re-exports (lib.rs). - infra : TokioScheduler (scheduler/mod.rs, nouveau) + pub mod scheduler et re-export (lib.rs). Tests QA inline (#[cfg(test)]) : 7 tests scheduler (3× sans flaky). `cargo test -p infrastructure` = 195 passed / 0 failed ; domaine + infra builds 0 warning ; LS1/LS2 toujours verts. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
300
crates/infrastructure/src/scheduler/mod.rs
Normal file
300
crates/infrastructure/src/scheduler/mod.rs
Normal file
@ -0,0 +1,300 @@
|
||||
//! [`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<dyn Scheduler>` (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<ScheduledTask>,
|
||||
/// Horloge injectée (déterminisme/testabilité) pour traduire l'échéance absolue en
|
||||
/// délai relatif (`deadline_ms - now`).
|
||||
clock: Arc<dyn Clock>,
|
||||
/// Table `id → handle` des réveils armés non encore tirés. `cancel` y retire +
|
||||
/// `abort()` le handle.
|
||||
handles: Arc<Mutex<HashMap<ScheduleId, JoinHandle<()>>>>,
|
||||
}
|
||||
|
||||
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<ScheduledTask>, clock: Arc<dyn Clock>) -> 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<ScheduledTask>, Arc<dyn Clock>) {
|
||||
let (tx, rx) = unbounded_channel();
|
||||
let clock: Arc<dyn Clock> = 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");
|
||||
if 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");
|
||||
if 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"
|
||||
);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user