//! [`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, resume_prompt: &str, ) -> Result<(), AppError>; } /// Service d'orchestration des limites de session (§21.5). pub struct SessionLimitService { clock: Arc, scheduler: Arc, events: Arc, resumer: Arc, /// Reprises **armées** non encore tirées : `agent_id → ScheduleId` (en mémoire). armed: Mutex>, } impl SessionLimitService { /// Construit le service depuis ses ports injectés (composition root). #[must_use] pub fn new( clock: Arc, scheduler: Arc, events: Arc, resumer: Arc, ) -> 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, resets_at_ms: Option, ) { 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.arm_scheduled(agent_id, fire_at_ms, node_id, conversation_id, resets_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, }); } } } /// **(d) Filet humain (§21.1 niveau 3).** L'utilisateur a saisi l'heure de reset /// pour un agent en limite **suspectée** (rien n'a matché automatiquement). On /// construit une [`SessionLimit`] de source [`RateLimitSource::Human`], on calcule /// le plan via [`plan_resume`] et on **arme exactement la même reprise** que la /// branche auto : mêmes événements (`AgentRateLimited{Some}` + `AgentResumeScheduled`), /// même dédoublonnage, même annulabilité via [`Self::cancel_resume`]. /// /// L'heure saisie est traitée par le domaine sans privilège particulier : un reset /// déjà passé est clampé à `now` par [`plan_resume`] ⇒ reprise immédiate. Le cas /// [`ResumePlan::HumanFallback`] est ici inatteignable (`resets_at_ms` est toujours /// `Some`) ; on le traite en no-op défensif pour rester total. pub fn confirm_human_resume( &self, agent_id: AgentId, node_id: NodeId, conversation_id: Option, resets_at_ms: i64, ) { let now = self.clock.now_millis(); let limit = SessionLimit::new(Some(resets_at_ms), now, RateLimitSource::Human); if let ResumePlan::Scheduled { fire_at_ms, conversation_id, } = plan_resume(now, &limit, conversation_id) { self.arm_scheduled(agent_id, fire_at_ms, node_id, conversation_id, Some(resets_at_ms)); } } /// Arme (ou ré-arme) une reprise **programmée** pour `agent_id`, fabrique commune aux /// deux entrées (auto §21.1 niveaux 1/2 et filet humain niveau 3). Séquence stricte, /// identique à l'origine — d'où **zéro régression** : publie `AgentRateLimited` /// (avec l'heure de reset connue), **dédoublonne** l'armement précédent via /// [`Self::disarm`] (interne, sans événement), arme le réveil via [`Scheduler::arm`], /// mémorise le [`ScheduleId`], puis publie `AgentResumeScheduled`. /// /// `resets_at_ms` est l'heure de reset **annoncée à l'UI** (countdown) ; `fire_at_ms` /// est l'échéance effective (déjà clampée anti-passé par le domaine). Les deux ne /// coïncident que si le reset est futur — on conserve donc la sémantique d'origine en /// publiant l'heure de reset brute, pas l'échéance clampée. fn arm_scheduled( &self, agent_id: AgentId, fire_at_ms: i64, node_id: NodeId, conversation_id: Option, resets_at_ms: Option, ) { 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 }); } /// **(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); } } }