diff --git a/crates/domain/src/ids.rs b/crates/domain/src/ids.rs index 00fb82c..0e67a71 100644 --- a/crates/domain/src/ids.rs +++ b/crates/domain/src/ids.rs @@ -92,3 +92,10 @@ typed_id!( /// Identifies a node in a [`crate::layout::LayoutTree`]. NodeId ); +typed_id!( + /// Identifies one armed one-shot wake-up of a [`crate::ports::Scheduler`] + /// (ARCHITECTURE §21.4). Opaque, cancellable handle returned by + /// [`crate::ports::Scheduler::arm`] and consumed by + /// [`crate::ports::Scheduler::cancel`]. + ScheduleId +); diff --git a/crates/domain/src/lib.rs b/crates/domain/src/lib.rs index a1c619c..8fca7d6 100644 --- a/crates/domain/src/lib.rs +++ b/crates/domain/src/lib.rs @@ -65,8 +65,8 @@ mod validation; pub use error::DomainError; pub use ids::{ - AgentId, LayoutId, NodeId, ProfileId, ProjectId, SessionId, SkillId, TabId, TemplateId, - WindowId, + AgentId, LayoutId, NodeId, ProfileId, ProjectId, ScheduleId, SessionId, SkillId, TabId, + TemplateId, WindowId, }; pub use project::{Project, ProjectPath}; @@ -145,6 +145,6 @@ pub use ports::{ FsError, GitCommitInfo, GitError, GitFileStatus, GitPort, GraphCommit, IdGenerator, MemoryError, MemoryQuery, MemoryRecall, MemoryStore, Output, OutputStream, PermissionStore, PreparedContext, ProcessError, ProcessSpawner, ProfileStore, ProjectStore, PtyError, PtyHandle, - PtyPort, RemoteError, RemoteHost, RemotePath, RuntimeError, SpawnSpec, StoreError, - TemplateStore, + PtyPort, RemoteError, RemoteHost, RemotePath, RuntimeError, ScheduledTask, Scheduler, SpawnSpec, + StoreError, TemplateStore, }; diff --git a/crates/domain/src/ports.rs b/crates/domain/src/ports.rs index de71833..6a8fd47 100644 --- a/crates/domain/src/ports.rs +++ b/crates/domain/src/ports.rs @@ -29,7 +29,7 @@ use thiserror::Error; use crate::agent::AgentManifest; use crate::events::DomainEvent; -use crate::ids::{AgentId, SessionId}; +use crate::ids::{AgentId, NodeId, ScheduleId, SessionId}; use crate::markdown::MarkdownDoc; use crate::memory::{Memory, MemoryIndexEntry, MemoryLink, MemorySlug}; use crate::permission::ProjectPermissions; @@ -262,6 +262,34 @@ pub enum ReplyEvent { /// clôt le flux ; deltas, activités et heartbeats sont tous non terminaux. pub type ReplyStream = Box + Send>; +/// Intention **model-agnostique** exécutée à l'échéance d'un réveil de +/// [`Scheduler`] (ARCHITECTURE §21.4). +/// +/// **Donnée pure, pas de closure** : à l'image du dispatch de l'orchestrateur +/// (§14.3, où une requête est une *donnée* validée puis dispatchée), une tâche +/// programmée est une **valeur** qui franchit la frontière, jamais un effet ni une +/// fonction. C'est ce qui garde le port [`Scheduler`] sans aucune dépendance vers +/// l'application : l'adapter ne fait que **remettre** cette donnée à un drain +/// applicatif, qui seul l'exécute. +/// +/// Enum **extensible** : d'autres intentions programmées pourront s'ajouter sans +/// toucher au port (Open/Closed). +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum ScheduledTask { + /// Reprendre un agent à l'heure de reset de sa limite de session (§21) : relancer + /// via [`SessionPlan::Resume`] avec un prompt de reprise court (logique côté + /// application, lot LS4). Porte exactement le pivot de reprise model-agnostique. + ResumeAgent { + /// L'agent à reprendre. + agent_id: AgentId, + /// La cellule (nœud du layout) qui héberge sa session. + node_id: NodeId, + /// Id de conversation du moteur à reprendre (`None` ⇒ reprise dégradée sans + /// id, comme le reste du chemin de reprise §15). + conversation_id: Option, + }, +} + // --------------------------------------------------------------------------- // Per-port error types // --------------------------------------------------------------------------- @@ -1266,3 +1294,50 @@ pub trait IdGenerator: Send + Sync { /// Returns a fresh UUID. fn new_uuid(&self) -> uuid::Uuid; } + +/// Minuterie **one-shot annulable** (ARCHITECTURE §21.4) — le **seul** port neuf de +/// la feature « limites de session ». Là où [`Clock`] dit *quelle heure il est*, ce +/// port *réveille* à une échéance absolue donnée. +/// +/// # Interface Segregation (§21.8-I) +/// Volontairement minimal — `arm` / `cancel`, rien d'autre — et **distinct** de +/// [`Clock`] : un consommateur qui n'a besoin que de « réveille-moi à T » ne dépend +/// pas d'un store ni d'une horloge. +/// +/// # Synchrone (pas d'`async_trait`) +/// Comme [`Clock`], [`IdGenerator`] et [`EventBus`], ce port est **non bloquant** : +/// `arm` ne fait qu'enregistrer une minuterie (l'attente réelle se passe en tâche de +/// fond dans l'adapter) et `cancel` ne fait qu'annuler — aucun `await` au point +/// d'appel. On le garde donc en `fn` simple, sans payer le boxing d'`async_trait`. +/// +/// # Remise de la tâche échue — dispatch par DONNÉE (§14.3, §21.4) +/// Le port **n'exécute jamais** la reprise (aucune dépendance vers l'application). +/// À l'échéance, l'adapter **pousse la [`ScheduledTask`] (une valeur)** dans un canal +/// fourni à sa construction, que l'application **draine** et exécute — exactement le +/// patron du watcher d'orchestrateur, qui dispatche une requête-donnée vers +/// l'`OrchestratorService`. Aucune closure ne franchit la frontière domaine. +/// +/// # État en mémoire (§21.1-3) +/// Aucune persistance : un réveil armé ne survit pas à un redémarrage d'IdeA (le +/// chemin `ListResumableAgents` existant prend alors le relais). +/// +/// # Substituabilité (Liskov) +/// Tout adapter (le `TokioScheduler` réel comme un fake de test qui tire à la +/// demande) respecte le même contrat : `arm` rend un id annulable ; `cancel` renvoie +/// `true` ssi il a effectivement désarmé un réveil non encore tiré. +pub trait Scheduler: Send + Sync { + /// Arme une minuterie one-shot à `deadline_ms` (**époche-millisecondes absolues**, + /// cohérent avec [`Clock::now_millis`] et `ResumePlan::Scheduled.fire_at_ms`). À + /// l'échéance, l'adapter remet `task` au drain applicatif (cf. doc du trait). + /// Renvoie un [`ScheduleId`] annulable. + /// + /// `deadline_ms` **≤ now** ⇒ déclenchement **au plus tôt** (immédiat). Le clamp + /// anti-passé est déjà assuré en amont par `plan_resume` (domaine) ; ce contrat ne + /// fait que garantir qu'une échéance passée ne « se perd » pas. + fn arm(&self, deadline_ms: i64, task: ScheduledTask) -> ScheduleId; + + /// Annule un réveil armé **non encore tiré**. Idempotent et **sans erreur** : + /// renvoie `true` s'il a effectivement été désarmé, `false` s'il était inconnu ou + /// déjà tiré. C'est ce qui sous-tend la « reprise auto **annulable** » (§21.1-4). + fn cancel(&self, id: ScheduleId) -> bool; +} diff --git a/crates/infrastructure/src/lib.rs b/crates/infrastructure/src/lib.rs index df9e914..fc8d01a 100644 --- a/crates/infrastructure/src/lib.rs +++ b/crates/infrastructure/src/lib.rs @@ -30,6 +30,7 @@ pub mod pty; pub mod remote; pub mod runtime; pub mod sandbox; +pub mod scheduler; pub mod session; pub mod store; @@ -59,6 +60,7 @@ pub use runtime::CliAgentRuntime; #[cfg(target_os = "linux")] pub use sandbox::LandlockSandbox; pub use sandbox::{default_enforcer, NoopSandbox}; +pub use scheduler::TokioScheduler; pub use session::{ClaudeSdkSession, CodexExecSession, FakeCli, StructuredSessionFactory}; #[cfg(feature = "vector-onnx")] pub use store::OnnxEmbedder; diff --git a/crates/infrastructure/src/scheduler/mod.rs b/crates/infrastructure/src/scheduler/mod.rs new file mode 100644 index 0000000..6a2100c --- /dev/null +++ b/crates/infrastructure/src/scheduler/mod.rs @@ -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` (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"); + 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" + ); + } +}