Introduit le modèle AgentManifest { version, entries, orchestrator } et la
garde d'écriture directe may_write_directly(..., &OrchestratorDesignation) :
seul l'orchestrateur désigné peut écrire directement, les autres passent par
le rendez-vous médié. Câble la désignation à travers domain → application →
infrastructure → app-tauri (context_guard, service, lifecycle, ports).
Ajoute crates/application/src/diag.rs : sink de diagnostic best-effort, sans
dépendance, qui miroite les traces du rendez-vous inter-agents de
l'orchestrateur vers un fichier de log persistant (utile au lancement via
AppImage où stderr est jeté), avec la même discipline « zéro dépendance,
ne casse jamais le rendez-vous ».
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
295 lines
12 KiB
Rust
295 lines
12 KiB
Rust
//! [`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<String>,
|
|
resume_prompt: &str,
|
|
) -> Result<(), AppError>;
|
|
}
|
|
|
|
/// Service d'orchestration des limites de session (§21.5).
|
|
pub struct SessionLimitService {
|
|
clock: Arc<dyn Clock>,
|
|
scheduler: Arc<dyn Scheduler>,
|
|
events: Arc<dyn EventBus>,
|
|
resumer: Arc<dyn AgentResumer>,
|
|
/// Reprises **armées** non encore tirées : `agent_id → ScheduleId` (en mémoire).
|
|
armed: Mutex<HashMap<AgentId, ScheduleId>>,
|
|
}
|
|
|
|
impl SessionLimitService {
|
|
/// Construit le service depuis ses ports injectés (composition root).
|
|
#[must_use]
|
|
pub fn new(
|
|
clock: Arc<dyn Clock>,
|
|
scheduler: Arc<dyn Scheduler>,
|
|
events: Arc<dyn EventBus>,
|
|
resumer: Arc<dyn AgentResumer>,
|
|
) -> 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<String>,
|
|
resets_at_ms: Option<i64>,
|
|
) {
|
|
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<String>,
|
|
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<String>,
|
|
resets_at_ms: Option<i64>,
|
|
) {
|
|
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);
|
|
}
|
|
}
|
|
}
|