feat(session-limits): LS7 — câblage backend app-tauri (taps niveaux 1&2 + reprise annulable)
Branche le SessionLimitService dans l'application Tauri et expose la surface de reprise/annulation au front. - application/agent/lifecycle.rs : LaunchAgentOutput.profile exposé (None sur réattache/idempotent, Some sur lancement effectif). - application/terminal/registry.rs : StructuredSessions::meta_for_session() (lookup agent/node par SessionId pour le tap niveau 1). - app-tauri/state.rs : ResumeContext(s), AppAgentResumer (impl du port AgentResumer au-dessus de LaunchAgent), instanciation + câblage du SessionLimitService (TokioScheduler + drain des réveils) dans AppState::build. - app-tauri/commands.rs : taps niveau 1 (agent_send) et niveau 2 (launch_agent, parser regex confiné), alimentation de resume_contexts, commande cancel_resume. - app-tauri/lib.rs : enregistrement de cancel_resume dans le handler. - app-tauri/Cargo.toml : dépendance async-trait. Tests : session_limit_wiring.rs (2 tests de composition) + meta_for_session dans structured_registry_d1.rs ; fixtures dto_agents/dto_chat ajustées (profile: None). Tout vert. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@ -12,9 +12,11 @@ use std::path::PathBuf;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use application::{
|
||||
AssignSkillToAgent, ChangeAgentProfile, CheckEmbedderSuggestion, CloseProject, CloseTab,
|
||||
AgentResumer, AppError, AssignSkillToAgent, ChangeAgentProfile, CheckEmbedderSuggestion,
|
||||
CloseProject, CloseTab,
|
||||
CloseTerminal, ConfigureProfiles, CreateAgentFromScratch, CreateAgentFromTemplate,
|
||||
CreateLayout, CreateMemory, CreateProject, CreateSkill, CreateTemplate, DeleteAgent,
|
||||
LaunchAgentInput, SessionLimitService,
|
||||
DeleteEmbedderProfile, DeleteLayout, DeleteMemory, DeleteProfile, DeleteSkill, DeleteTemplate,
|
||||
DescribeEmbedderEngines, DetectAgentDrift, DetectProfiles, DismissEmbedderSuggestion,
|
||||
FirstRunState, GetMemory, GetProjectPermissions, GitBranches, GitCheckout, GitCommit, GitGraph,
|
||||
@ -35,7 +37,7 @@ use domain::ports::{
|
||||
AgentContextStore, AgentRuntime, AgentSessionFactory, Clock, Embedder, EmbedderEnvInspector,
|
||||
EmbedderProfileStore, EmbedderPromptStore, EventBus, FileSystem, GitPort, IdGenerator,
|
||||
MemoryRecall, MemoryStore, PermissionStore, ProcessSpawner, ProfileStore, ProjectStore,
|
||||
PtyPort, SkillStore, TemplateStore,
|
||||
PtyPort, ScheduledTask, Scheduler, SkillStore, TemplateStore,
|
||||
};
|
||||
use domain::profile::{
|
||||
AgentProfile, ContextInjection, McpConfigStrategy, McpTransport, StructuredAdapter,
|
||||
@ -54,7 +56,7 @@ use infrastructure::{
|
||||
HeuristicHandoffSummarizer, IdeaiContextStore, InMemoryConversationRegistry, InMemoryMailbox,
|
||||
LocalFileSystem, LocalProcessSpawner, McpServer, MediatedInbox, NaiveMemoryRecall,
|
||||
OrchestratorWatchHandle, PortablePtyAdapter, StructuredSessionFactory, SystemClock,
|
||||
SystemMillisClock, TokioBroadcastEventBus, UuidGenerator, VectorMemoryRecall,
|
||||
SystemMillisClock, TokioBroadcastEventBus, TokioScheduler, UuidGenerator, VectorMemoryRecall,
|
||||
DEFAULT_OLLAMA_BASE_URL, ONNX_CACHE_SUBDIR, RECOMMENDED_ONNX_MODELS, VECTOR_HTTP_ENABLED,
|
||||
VECTOR_ONNX_ENABLED,
|
||||
};
|
||||
@ -127,6 +129,108 @@ impl application::ProviderSessionProvider for AppProviderSessionProvider {
|
||||
}
|
||||
}
|
||||
|
||||
/// Contexte minimal de **relance** d'un agent (LS7, ARCHITECTURE §21.5).
|
||||
///
|
||||
/// [`AgentResumer::resume`] et [`ScheduledTask::ResumeAgent`] ne portent **pas** le
|
||||
/// `Project` ni la taille de la cellule, alors que [`LaunchAgentInput`] les exige.
|
||||
/// La commande `launch_agent` (seul endroit où ces faits sont en main) alimente ce
|
||||
/// contexte par `agent_id` ; [`AppAgentResumer`] le relit à l'échéance pour
|
||||
/// recomposer un lancement complet.
|
||||
#[derive(Clone)]
|
||||
pub struct ResumeContext {
|
||||
/// Le projet hôte de l'agent (pour recomposer `LaunchAgentInput`).
|
||||
pub project: Project,
|
||||
/// Hauteur de la cellule au dernier lancement (lignes PTY).
|
||||
pub rows: u16,
|
||||
/// Largeur de la cellule au dernier lancement (colonnes PTY).
|
||||
pub cols: u16,
|
||||
}
|
||||
|
||||
/// Registre partagé `agent_id → ResumeContext` (composition root ↔ commande
|
||||
/// `launch_agent`). Le **même** `Arc` est injecté dans [`AppAgentResumer`] et conservé
|
||||
/// sur [`AppState`] pour que la commande l'alimente à chaque lancement.
|
||||
pub type ResumeContexts = Arc<Mutex<HashMap<AgentId, ResumeContext>>>;
|
||||
|
||||
/// Implémente le port applicatif [`AgentResumer`] (LS7) **par-dessus** le mécanisme de
|
||||
/// lancement existant ([`LaunchAgent`]).
|
||||
///
|
||||
/// À l'échéance d'une limite de session, [`SessionLimitService::execute_resume`]
|
||||
/// délègue ici : on relit le [`ResumeContext`] alimenté par `launch_agent`, on
|
||||
/// recompose un [`LaunchAgentInput`] (avec le `conversation_id` de la cellule ⇒
|
||||
/// `LaunchAgent` applique [`domain::ports::SessionPlan::Resume`]), puis on transmet le
|
||||
/// `resume_prompt` comme **premier tour** via le portail d'entrée unique
|
||||
/// ([`InputMediator`](domain::input::InputMediator)) — jamais un write PTY brut
|
||||
/// (ARCHITECTURE §20). Sans contexte connu (resume « à l'aveugle »), on échoue
|
||||
/// proprement (l'erreur empêche `AgentResumed` d'être publié), jamais de panique.
|
||||
struct AppAgentResumer {
|
||||
/// Le **même** `Arc<LaunchAgent>` que la commande `launch_agent` (relance/réattache).
|
||||
launch_agent: Arc<LaunchAgent>,
|
||||
/// Portail d'entrée unique : injecte le prompt de reprise comme premier tour.
|
||||
input_mediator: Arc<dyn domain::input::InputMediator>,
|
||||
/// Contexte de relance alimenté par la commande `launch_agent`.
|
||||
contexts: ResumeContexts,
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl AgentResumer for AppAgentResumer {
|
||||
async fn resume(
|
||||
&self,
|
||||
agent_id: AgentId,
|
||||
node_id: domain::NodeId,
|
||||
conversation_id: Option<String>,
|
||||
resume_prompt: &str,
|
||||
) -> Result<(), AppError> {
|
||||
// Repli propre (jamais de panique) : sans contexte de relance connu, on ne
|
||||
// reprend pas à l'aveugle. L'erreur remonte ⇒ `AgentResumed` n'est pas publié.
|
||||
let ctx = self
|
||||
.contexts
|
||||
.lock()
|
||||
.ok()
|
||||
.and_then(|m| m.get(&agent_id).cloned());
|
||||
let Some(ctx) = ctx else {
|
||||
return Err(AppError::NotFound(format!(
|
||||
"resume context for agent {agent_id}"
|
||||
)));
|
||||
};
|
||||
|
||||
// Recompose la déclaration MCP réelle (même recette que la commande
|
||||
// `launch_agent`) pour que l'agent repris retrouve ses outils `idea_*`.
|
||||
let mcp_runtime = crate::mcp_endpoint::idea_exe_path().map(|exe| McpRuntime {
|
||||
exe,
|
||||
endpoint: crate::mcp_endpoint::mcp_endpoint(&ctx.project.id)
|
||||
.as_cli_arg()
|
||||
.to_owned(),
|
||||
project_id: ctx.project.id.as_uuid().simple().to_string(),
|
||||
requester: agent_id.to_string(),
|
||||
});
|
||||
|
||||
self.launch_agent
|
||||
.execute(LaunchAgentInput {
|
||||
project: ctx.project,
|
||||
agent_id,
|
||||
rows: ctx.rows,
|
||||
cols: ctx.cols,
|
||||
node_id: Some(node_id),
|
||||
conversation_id,
|
||||
mcp_runtime,
|
||||
})
|
||||
.await?;
|
||||
|
||||
// Premier tour de reprise par le **portail d'entrée** (pas de write brut, §20) :
|
||||
// l'enqueue publie `DelegationReady`, la cellule l'écrit quand le prompt est prêt.
|
||||
// On ne corrèle aucune réponse (reprise, pas une délégation) ⇒ on lâche le
|
||||
// `PendingReply`.
|
||||
let ticket = domain::mailbox::Ticket::new(
|
||||
domain::mailbox::TicketId::new_random(),
|
||||
"IdeA",
|
||||
resume_prompt,
|
||||
);
|
||||
let _ = self.input_mediator.enqueue(agent_id, ticket);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
/// Everything the IPC layer needs at runtime, managed by Tauri.
|
||||
///
|
||||
/// Use cases are stored behind `Arc` so handlers clone cheaply. The concrete
|
||||
@ -339,6 +443,15 @@ pub struct AppState {
|
||||
/// stopped on close, exactly like the watcher, via the same `Mutex`-guarded
|
||||
/// per-project registry.
|
||||
pub mcp_servers: Mutex<HashMap<ProjectId, McpServerHandle>>,
|
||||
/// Service de gestion des **limites de session** des agents (ARCHITECTURE §21) :
|
||||
/// détecte → planifie → reprend, et annule. Alimenté par les taps niveau 1
|
||||
/// (structuré, `agent_send`) et niveau 2 (PTY, `launch_agent`) ; sa reprise auto est
|
||||
/// annulable via la commande `cancel_resume`.
|
||||
pub session_limit_service: Arc<SessionLimitService>,
|
||||
/// Registre `agent_id → ResumeContext` (LS7) partagé avec [`AppAgentResumer`] :
|
||||
/// la commande `launch_agent` y dépose le `Project`/taille du dernier lancement pour
|
||||
/// que la reprise auto puisse recomposer un `LaunchAgentInput` complet.
|
||||
pub resume_contexts: ResumeContexts,
|
||||
}
|
||||
|
||||
impl AppState {
|
||||
@ -910,6 +1023,45 @@ impl AppState {
|
||||
});
|
||||
}
|
||||
let input_mediator = Arc::clone(&mediated_inbox) as Arc<dyn domain::input::InputMediator>;
|
||||
|
||||
// --- Limites de session des agents (ARCHITECTURE §21, LS7) ---
|
||||
// Service pur-ports « détecter → planifier → reprendre » câblé sur l'existant :
|
||||
// l'horloge système, le bus partagé, un `TokioScheduler` (minuterie one-shot
|
||||
// annulable) dont les tâches échues sont drainées plus bas, et un `AppAgentResumer`
|
||||
// qui relance via le *même* `LaunchAgent` que la commande `launch_agent`.
|
||||
let resume_contexts: ResumeContexts = Arc::new(Mutex::new(HashMap::new()));
|
||||
let (resume_tx, mut resume_rx) = tokio::sync::mpsc::unbounded_channel::<ScheduledTask>();
|
||||
let scheduler = Arc::new(TokioScheduler::new(
|
||||
resume_tx,
|
||||
Arc::clone(&clock) as Arc<dyn Clock>,
|
||||
)) as Arc<dyn Scheduler>;
|
||||
let resumer = Arc::new(AppAgentResumer {
|
||||
launch_agent: Arc::clone(&launch_agent),
|
||||
input_mediator: Arc::clone(&input_mediator),
|
||||
contexts: Arc::clone(&resume_contexts),
|
||||
}) as Arc<dyn AgentResumer>;
|
||||
let session_limit_service = Arc::new(SessionLimitService::new(
|
||||
Arc::clone(&clock) as Arc<dyn Clock>,
|
||||
scheduler,
|
||||
Arc::clone(&events_port),
|
||||
resumer,
|
||||
));
|
||||
// Drain du scheduler (§21.5-b) : à chaque réveil tiré, `TokioScheduler` pousse une
|
||||
// `ScheduledTask::ResumeAgent` ; on l'exécute via le service (relance + AgentResumed).
|
||||
// Patron `sweep_stalled` : `tauri::async_runtime::spawn` (pas `tokio::spawn` — `build`
|
||||
// tourne dans le hook `setup` sans runtime ambiant). Détaché ⇒ vit autant que l'app.
|
||||
{
|
||||
let service = Arc::clone(&session_limit_service);
|
||||
tauri::async_runtime::spawn(async move {
|
||||
while let Some(task) = resume_rx.recv().await {
|
||||
if let Err(e) = service.execute_resume(task).await {
|
||||
// Best-effort : une relance qui échoue ne fige pas le drain.
|
||||
eprintln!("[session-limit] reprise auto échouée : {e}");
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
// Registre des conversations par paire (cadrage C3) : un fil par paire, session
|
||||
// vivante keyée par conversation (lève l'ambiguïté session/agent).
|
||||
let conversation_registry = Arc::new(InMemoryConversationRegistry::new())
|
||||
@ -1051,6 +1203,8 @@ impl AppState {
|
||||
orchestrator_service,
|
||||
orchestrator_watchers: Mutex::new(HashMap::new()),
|
||||
mcp_servers: Mutex::new(HashMap::new()),
|
||||
session_limit_service,
|
||||
resume_contexts,
|
||||
move_tab,
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user