fix(runtime): isolation d'usage multi-projet #107

Implémentation du garde ModelServerProjectUseGuard pour séparer
les retours des agents entre projets concurrents.

- Garde RAII par LocalModelServerId partagé entre projets
- Refus inter-projets concurrent via model_server_in_use
- Partage intra-projet conservé avec refcount
- Libération automatique à la fermeture/erreur de session

Tests de validation: rejet projet distinct, compteur refs, libération retrait
This commit is contained in:
2026-07-27 09:47:23 +02:00
parent c807a70fea
commit cf074a7d61
6 changed files with 265 additions and 15 deletions

View File

@ -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<ModelServerProjectUseGuard>,
) -> Result<LaunchAgentOutput, AppError> {
// 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<Option<ModelServerProjectUseGuard>, 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);

View File

@ -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<Mutex<HashMap<LocalModelServerId, ProjectServerUse>>>,
}
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<dyn EventBus>,
active: Mutex<HashMap<LocalModelServerId, ActiveServer>>,
inflight: AsyncMutex<HashMap<LocalModelServerId, Arc<InflightEnsure>>>,
project_usages: Arc<Mutex<HashMap<LocalModelServerId, ProjectServerUse>>>,
download_cancels: Mutex<HashMap<LocalModelServerId, ModelArtifactCancel>>,
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<ModelServerProjectUseGuard, AppError> {
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

View File

@ -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<ModelServerProjectUseGuard>,
}
/// 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<ModelServerProjectUseGuard>,
) {
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<ModelServerProjectUseGuard>,
}
#[derive(Clone)]
@ -446,6 +464,21 @@ impl StructuredSessions {
session: Arc<dyn AgentSession>,
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<dyn AgentSession>,
agent_id: AgentId,
node_id: NodeId,
model_server_guard: Option<ModelServerProjectUseGuard>,
) {
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,
},
);
}

View File

@ -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());