diff --git a/.ideai/tickets/107/carnet.md b/.ideai/tickets/107/carnet.md new file mode 100644 index 0000000..47158a8 --- /dev/null +++ b/.ideai/tickets/107/carnet.md @@ -0,0 +1,6 @@ +--- +issueRef: "#107" +version: 3 +updatedBy: {"kind":"user"} +updatedAt: 1785136945441 +--- diff --git a/.ideai/tickets/107/issue.md b/.ideai/tickets/107/issue.md new file mode 100644 index 0000000..ab84b6e --- /dev/null +++ b/.ideai/tickets/107/issue.md @@ -0,0 +1,17 @@ +--- +id: "6a79006b-0201-4176-ae54-39a05cc3baa6" +number: 107 +title: "[Bug] croisement entre les projet des retours des agents" +status: "open" +priority: "critical" +sprint: null +links: [] +agentRefs: [{"agentId":"a6ced819-b893-4213-b003-9e9dc79b9641","role":"assigned"}] +createdBy: {"kind":"user"} +updatedBy: {"kind":"user"} +createdAt: 1785136727824 +updatedAt: 1785136945441 +version: 3 +--- +Il y a un soucis très important que j'ai constatés. J'ai actuellement 2 projets ouverts: IdeA et GameTime. Les deux projets travaillaient en même temps et j'ai vu GameTime qui semblait récupérer une requete du projet IdeA. Pour plus de précision, sur mes deux projets, j'ai un agent Git qui utilise OpenCode et un modelle local llamacpp, et j'ia eu l'impression qu'ils ont tous les deux appelé leur agent Git mais GameTime à recus la réponse de IdeA car la réponse parlait d'une branche du projet IdeA. +Je ne suis pas totalement sur de ce que j'avance, la seule chose dont je suis sur, c'est que GameTime s'est vu adressé une réponse qui était déstinée au projet IdeA. \ No newline at end of file diff --git a/crates/application/src/agent/lifecycle.rs b/crates/application/src/agent/lifecycle.rs index 6f35b94..9bf23b1 100644 --- a/crates/application/src/agent/lifecycle.rs +++ b/crates/application/src/agent/lifecycle.rs @@ -35,7 +35,9 @@ use domain::live_state::WorkStatus; use crate::error::AppError; use crate::layout::{persist_doc, resolve_doc}; -use crate::model_server::{EnsureLocalModelServer, EnsureLocalModelServerInput}; +use crate::model_server::{ + EnsureLocalModelServer, EnsureLocalModelServerInput, ModelServerProjectUseGuard, +}; use crate::project::project_context_path; use crate::terminal::{StructuredSessions, TerminalSessions}; use crate::workstate::GetLiveStateLean; @@ -1782,7 +1784,8 @@ impl LaunchAgent { ) .await?; - self.ensure_local_model_server_for_opencode(&agent, &mut profile) + let local_model_server_guard = self + .ensure_local_model_server_for_opencode(&input.project, &agent, &mut profile) .await?; // 5a. ── INJECTION DE LA CONF MCP (cadrage v3, Décision 3) ── @@ -1886,6 +1889,7 @@ impl LaunchAgent { &spec.env, spec.sandbox.as_ref(), structured_policy.as_ref(), + local_model_server_guard, ) .await; } @@ -1905,7 +1909,10 @@ impl LaunchAgent { // 7. For the Stdin strategy, pipe the context once the PTY is live. if matches!(spec.context_plan, Some(ContextInjectionPlan::Stdin)) { - self.pty.write(&handle, content.as_str().as_bytes())?; + if let Err(err) = self.pty.write(&handle, content.as_str().as_bytes()) { + let _ = self.pty.kill(&handle).await; + return Err(err.into()); + } } let node_id = input.node_id.unwrap_or_else(NodeId::new_random); @@ -1917,8 +1924,12 @@ impl LaunchAgent { size, ); session.status = SessionStatus::Running; - self.sessions - .insert_in_project(input.project.id, handle, session.clone()); + self.sessions.insert_in_project_with_model_server_guard( + input.project.id, + handle, + session.clone(), + local_model_server_guard, + ); self.events.publish(DomainEvent::AgentLaunched { agent_id: agent.id, @@ -1964,6 +1975,7 @@ impl LaunchAgent { env: &[(String, String)], sandbox: Option<&SandboxPlan>, structured_policy: Option<&StructuredProviderLaunchPolicy>, + model_server_guard: Option, ) -> Result { // Relaie le plan de sandbox OS (lot LP4-4) à la fabrique : `spec.sandbox`, // déjà compilé (pur, domaine) en step 5d. `None` ⇒ exécution native inchangée. @@ -1986,7 +1998,13 @@ impl LaunchAgent { // Enregistre la session vivante (invariant « 1 session/agent » : déjà gardé en // amont sur les deux registres). - structured.insert_in_project(project_id, Arc::clone(&session), agent.id, node_id); + structured.insert_in_project_with_model_server_guard( + project_id, + Arc::clone(&session), + agent.id, + node_id, + model_server_guard, + ); // ── SÉPARATION DES DEUX CLÉS (ARCHITECTURE §19.7, lot P8a) ── // - **id de paire** (`pair_conversation_id`) : clé **logique** persistée sur @@ -2570,17 +2588,18 @@ impl LaunchAgent { async fn ensure_local_model_server_for_opencode( &self, + project: &Project, agent: &Agent, profile: &mut AgentProfile, - ) -> Result<(), AppError> { + ) -> Result, AppError> { if profile.structured_adapter != Some(StructuredAdapter::OpenCode) { - return Ok(()); + return Ok(None); } let Some(opencode) = profile.opencode.as_mut() else { - return Ok(()); + return Ok(None); }; let Some(server_id) = opencode.local_model_server_id else { - return Ok(()); + return Ok(None); }; let Some(ensure) = self.local_model_server.as_ref() else { let err = AppError::ModelServer { @@ -2592,6 +2611,13 @@ impl LaunchAgent { self.publish_agent_launch_failed(agent.id, &err); return Err(err); }; + let guard = match ensure.acquire_project_use(server_id, project.id) { + Ok(guard) => guard, + Err(err) => { + self.publish_agent_launch_failed(agent.id, &err); + return Err(err); + } + }; match ensure .execute(EnsureLocalModelServerInput { server_id }) @@ -2600,7 +2626,7 @@ impl LaunchAgent { Ok(output) => { opencode.base_url = output.ready.base_url; opencode.model = output.ready.model; - Ok(()) + Ok(Some(guard)) } Err(err) => { self.publish_agent_launch_failed(agent.id, &err); diff --git a/crates/application/src/model_server.rs b/crates/application/src/model_server.rs index 1aa9eba..4f9163e 100644 --- a/crates/application/src/model_server.rs +++ b/crates/application/src/model_server.rs @@ -15,7 +15,7 @@ use domain::ports::{ ModelArtifactDownloader, ModelArtifactProgress, ModelServerError, ModelServerProbe, ModelServerRegistry, ModelServerRuntime, ProcessStatus, ProfileStore, RemotePath, }; -use domain::{LocalModelServerId, StopPolicy}; +use domain::{LocalModelServerId, ProjectId, StopPolicy}; use tokio::sync::{Mutex as AsyncMutex, Notify}; use tokio::time::Instant; @@ -146,6 +146,42 @@ pub struct EnsureLocalModelServerOutput { pub ready: ModelServerReady, } +#[derive(Debug)] +struct ProjectServerUse { + project_id: ProjectId, + refs: usize, +} + +/// RAII guard held by live OpenCode sessions while they use a local model server. +/// +/// A local llama.cpp server may be shared by several agents of the same project, +/// but concurrent use by distinct projects is refused to avoid context cross-talk +/// through the shared OpenAI-compatible endpoint. +#[derive(Debug)] +pub struct ModelServerProjectUseGuard { + server_id: LocalModelServerId, + project_id: ProjectId, + usages: Arc>>, +} + +impl Drop for ModelServerProjectUseGuard { + fn drop(&mut self) { + let Ok(mut usages) = self.usages.lock() else { + return; + }; + let Some(active) = usages.get_mut(&self.server_id) else { + return; + }; + if active.project_id != self.project_id { + return; + } + active.refs = active.refs.saturating_sub(1); + if active.refs == 0 { + usages.remove(&self.server_id); + } + } +} + /// Readiness retry policy. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub struct ReadinessPolicy { @@ -220,6 +256,7 @@ pub struct EnsureLocalModelServer { events: Arc, active: Mutex>, inflight: AsyncMutex>>, + project_usages: Arc>>, download_cancels: Mutex>, readiness: ReadinessPolicy, hf_download_deadline: Duration, @@ -247,6 +284,7 @@ impl EnsureLocalModelServer { events, active: Mutex::new(HashMap::new()), inflight: AsyncMutex::new(HashMap::new()), + project_usages: Arc::new(Mutex::new(HashMap::new())), download_cancels: Mutex::new(HashMap::new()), readiness: ReadinessPolicy::default(), hf_download_deadline: DEFAULT_HF_DOWNLOAD_DEADLINE, @@ -287,6 +325,51 @@ impl EnsureLocalModelServer { self } + /// Acquires exclusive cross-project use of `server_id` for `project_id`. + /// + /// Multiple agents from the same project can hold the guard concurrently. A + /// different project receives the existing `model_server_in_use` error channel. + /// + /// # Errors + /// [`AppError::ModelServer`] with `code=model_server_in_use` when another + /// project currently owns the server usage guard. + pub fn acquire_project_use( + &self, + server_id: LocalModelServerId, + project_id: ProjectId, + ) -> Result { + let mut usages = self + .project_usages + .lock() + .map_err(|_| ModelServerError::InUse(server_id.to_string()))?; + match usages.get_mut(&server_id) { + Some(active) if active.project_id == project_id => { + active.refs = active.refs.saturating_add(1); + } + Some(active) => { + return Err(ModelServerError::InUse(format!( + "local model server {server_id} is already in use by project {}", + active.project_id + )) + .into()); + } + None => { + usages.insert( + server_id, + ProjectServerUse { + project_id, + refs: 1, + }, + ); + } + } + Ok(ModelServerProjectUseGuard { + server_id, + project_id, + usages: Arc::clone(&self.project_usages), + }) + } + /// Ensures the server is reachable. /// /// # Errors diff --git a/crates/application/src/terminal/registry.rs b/crates/application/src/terminal/registry.rs index 7ffff67..0d4e657 100644 --- a/crates/application/src/terminal/registry.rs +++ b/crates/application/src/terminal/registry.rs @@ -13,6 +13,8 @@ use domain::conversation::ConversationId; use domain::ports::{AgentSession, PtyHandle}; use domain::{AgentId, IssueRef, NodeId, ProjectId, SessionId, SessionKind, TerminalSession}; +use crate::model_server::ModelServerProjectUseGuard; + /// Runtime family of a live agent session. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum LiveSessionKind { @@ -38,11 +40,12 @@ pub struct LiveSessionSnapshot { } /// A registered, live terminal: its PTY handle plus the domain snapshot. -#[derive(Debug, Clone)] +#[derive(Debug)] struct Entry { project_id: ProjectId, handle: PtyHandle, session: TerminalSession, + _model_server_guard: Option, } /// Read-only liveness query over the agents that currently own a live PTY. @@ -164,6 +167,18 @@ impl TerminalSessions { project_id: ProjectId, handle: PtyHandle, session: TerminalSession, + ) { + self.insert_in_project_with_model_server_guard(project_id, handle, session, None); + } + + /// Inserts a freshly-opened session and retains a model-server usage guard + /// until the session is removed. + pub fn insert_in_project_with_model_server_guard( + &self, + project_id: ProjectId, + handle: PtyHandle, + session: TerminalSession, + model_server_guard: Option, ) { if let Ok(mut map) = self.entries.lock() { map.insert( @@ -172,6 +187,7 @@ impl TerminalSessions { project_id, handle, session, + _model_server_guard: model_server_guard, }, ); } @@ -388,6 +404,8 @@ struct StructuredEntry { agent_id: AgentId, /// La cellule (feuille de layout) qui héberge actuellement la vue. node_id: NodeId, + /// Optional local-model-server usage guard retained for the live session. + _model_server_guard: Option, } #[derive(Clone)] @@ -446,6 +464,21 @@ impl StructuredSessions { session: Arc, agent_id: AgentId, node_id: NodeId, + ) { + self.insert_in_project_with_model_server_guard( + project_id, session, agent_id, node_id, None, + ); + } + + /// Enregistre une session structurée et retient un garde d'usage de serveur + /// local jusqu'au retrait de la session. + pub fn insert_in_project_with_model_server_guard( + &self, + project_id: ProjectId, + session: Arc, + agent_id: AgentId, + node_id: NodeId, + model_server_guard: Option, ) { if let Ok(mut map) = self.entries.lock() { let id = session.id(); @@ -456,6 +489,7 @@ impl StructuredSessions { session, agent_id, node_id, + _model_server_guard: model_server_guard, }, ); } diff --git a/crates/application/tests/model_server.rs b/crates/application/tests/model_server.rs index 4dcffbd..9cd50f4 100644 --- a/crates/application/tests/model_server.rs +++ b/crates/application/tests/model_server.rs @@ -8,7 +8,7 @@ use async_trait::async_trait; use application::{ DeleteModelServer, DeleteModelServerInput, EnsureLocalModelServer, EnsureLocalModelServerInput, - ModelServerReadinessPolicy, + ModelServerReadinessPolicy, TerminalSessions, }; use domain::events::DomainEvent; use domain::model_server::{ @@ -23,12 +23,19 @@ use domain::ports::{ ProcessStatus, ProfileStore, RemotePath, SpawnSpec, StoreError, }; use domain::profile::{AgentProfile, ContextInjection, OpenCodeConfig, StructuredAdapter}; -use domain::{LocalModelServerId, ProfileId, ProjectPath}; +use domain::{ + AgentId, LocalModelServerId, NodeId, ProfileId, ProjectId, ProjectPath, PtySize, SessionId, + SessionKind, SessionStatus, TerminalSession, +}; fn sid(n: u128) -> LocalModelServerId { LocalModelServerId::from_uuid(uuid::Uuid::from_u128(n)) } +fn pid(n: u128) -> ProjectId { + ProjectId::from_uuid(uuid::Uuid::from_u128(n)) +} + fn config( id: LocalModelServerId, port: u16, @@ -611,6 +618,83 @@ async fn concurrent_ensure_same_server_shares_one_start_attempt() { assert_eq!(starting_events, 1); } +#[tokio::test] +async fn local_model_server_use_rejects_distinct_project_until_guard_is_released() { + let usecase = ensure( + Arc::new(FakeRegistry::default()), + Arc::new(FakeProbe::new(Vec::new())), + Arc::new(FakeProcess::default()), + Arc::new(FakeFs::default()), + Arc::new(FakeEvents::default()), + ); + + let guard = usecase.acquire_project_use(sid(30), pid(1)).unwrap(); + let err = usecase.acquire_project_use(sid(30), pid(2)).unwrap_err(); + + assert_eq!(err.code(), "MODEL_SERVER"); + assert!(err.to_string().contains("model_server_in_use")); + + drop(guard); + assert!(usecase.acquire_project_use(sid(30), pid(2)).is_ok()); +} + +#[tokio::test] +async fn local_model_server_use_allows_same_project_with_refcount() { + let usecase = ensure( + Arc::new(FakeRegistry::default()), + Arc::new(FakeProbe::new(Vec::new())), + Arc::new(FakeProcess::default()), + Arc::new(FakeFs::default()), + Arc::new(FakeEvents::default()), + ); + + let first = usecase.acquire_project_use(sid(31), pid(1)).unwrap(); + let second = usecase.acquire_project_use(sid(31), pid(1)).unwrap(); + + assert!(usecase.acquire_project_use(sid(31), pid(2)).is_err()); + drop(first); + assert!(usecase.acquire_project_use(sid(31), pid(2)).is_err()); + drop(second); + assert!(usecase.acquire_project_use(sid(31), pid(2)).is_ok()); +} + +#[tokio::test] +async fn terminal_session_removal_releases_local_model_server_use_guard() { + let usecase = ensure( + Arc::new(FakeRegistry::default()), + Arc::new(FakeProbe::new(Vec::new())), + Arc::new(FakeProcess::default()), + Arc::new(FakeFs::default()), + Arc::new(FakeEvents::default()), + ); + let project = pid(1); + let server_id = sid(32); + let session_id = SessionId::from_uuid(uuid::Uuid::from_u128(33)); + let guard = usecase.acquire_project_use(server_id, project).unwrap(); + let sessions = TerminalSessions::new(); + let mut session = TerminalSession::starting( + session_id, + NodeId::from_uuid(uuid::Uuid::from_u128(34)), + ProjectPath::new("/tmp/project-a").unwrap(), + SessionKind::Agent { + agent_id: AgentId::from_uuid(uuid::Uuid::from_u128(35)), + }, + PtySize { rows: 24, cols: 80 }, + ); + session.status = SessionStatus::Running; + + sessions.insert_in_project_with_model_server_guard( + project, + domain::ports::PtyHandle { session_id }, + session, + Some(guard), + ); + assert!(usecase.acquire_project_use(server_id, pid(2)).is_err()); + + sessions.remove(&session_id); + assert!(usecase.acquire_project_use(server_id, pid(2)).is_ok()); +} + #[tokio::test] async fn missing_model_path_is_path_not_accessible() { let registry = Arc::new(FakeRegistry::default());