//! Unit tests for the local model-server ensure use case. use std::collections::{HashMap, VecDeque}; use std::sync::{Arc, Mutex}; use std::time::Duration; use async_trait::async_trait; use application::{ DeleteModelServer, DeleteModelServerInput, EnsureLocalModelServer, EnsureLocalModelServerInput, ModelServerReadinessPolicy, }; use domain::events::DomainEvent; use domain::model_server::{ ExecutablePath, HfModelRef, LlamaCppOptions, LocalModelRef, LocalModelServerConfig, LocalModelServerKind, ModelPath, ModelServerEndpoint, ModelServerLifecycleStatus, ModelServerStatus, ModelSource, StopPolicy, }; use domain::ports::{ DirEntry, EventBus, EventStream, FileSystem, FsError, ManagedProcess, ManagedProcessHandle, ModelArtifactCancel, ModelArtifactDownloader, ModelArtifactProgress, ModelArtifactResolution, ModelServerArgv, ModelServerError, ModelServerProbe, ModelServerRegistry, ModelServerRuntime, ProcessStatus, ProfileStore, RemotePath, SpawnSpec, StoreError, }; use domain::profile::{AgentProfile, ContextInjection, OpenCodeConfig, StructuredAdapter}; use domain::{LocalModelServerId, ProfileId, ProjectPath}; fn sid(n: u128) -> LocalModelServerId { LocalModelServerId::from_uuid(uuid::Uuid::from_u128(n)) } fn config( id: LocalModelServerId, port: u16, path: &str, auto_start: bool, ) -> LocalModelServerConfig { LocalModelServerConfig::new( id, LocalModelServerKind::LlamaCpp, "llama.cpp", ModelServerEndpoint::new(format!("http://localhost:{port}"), port).unwrap(), LocalModelRef::new( "qwen", "Qwen", Some(ModelSource::LocalPath { path: ModelPath::new(path).unwrap(), }), "qwen3-coder-30b", ) .unwrap(), Some(ExecutablePath::new("llama-server").unwrap()), LlamaCppOptions::default(), Vec::new(), auto_start, StopPolicy::StopOnAppExit, ) .unwrap() } fn hf_config(id: LocalModelServerId, port: u16, repo: &str) -> LocalModelServerConfig { LocalModelServerConfig::new( id, LocalModelServerKind::LlamaCpp, "llama.cpp", ModelServerEndpoint::new(format!("http://localhost:{port}"), port).unwrap(), LocalModelRef::new( "qwen", "Qwen", Some(ModelSource::HuggingFace { repo: HfModelRef::new(repo).unwrap(), }), "qwen3-coder-30b", ) .unwrap(), Some(ExecutablePath::new("llama-server").unwrap()), LlamaCppOptions::default(), Vec::new(), true, StopPolicy::StopOnAppExit, ) .unwrap() } #[derive(Default)] struct FakeRegistry(Mutex>); #[async_trait] impl ModelServerRegistry for FakeRegistry { async fn get( &self, id: &LocalModelServerId, ) -> Result, ModelServerError> { Ok(self.0.lock().unwrap().get(id).cloned()) } async fn list(&self) -> Result, ModelServerError> { Ok(self.0.lock().unwrap().values().cloned().collect()) } async fn save(&self, config: LocalModelServerConfig) -> Result<(), ModelServerError> { self.0.lock().unwrap().insert(config.id, config); Ok(()) } async fn delete(&self, id: LocalModelServerId) -> Result<(), ModelServerError> { self.0.lock().unwrap().remove(&id); Ok(()) } } #[derive(Default)] struct FakeProfiles(Mutex>); #[async_trait] impl ProfileStore for FakeProfiles { async fn list(&self) -> Result, StoreError> { Ok(self.0.lock().unwrap().clone()) } async fn save(&self, profile: &AgentProfile) -> Result<(), StoreError> { self.0.lock().unwrap().push(profile.clone()); Ok(()) } async fn delete(&self, _id: ProfileId) -> Result<(), StoreError> { Ok(()) } async fn is_configured(&self) -> Result { Ok(true) } async fn mark_configured(&self) -> Result<(), StoreError> { Ok(()) } } fn opencode_profile(id: u128, server_id: LocalModelServerId) -> AgentProfile { AgentProfile::new( ProfileId::from_uuid(uuid::Uuid::from_u128(id)), "Local OpenCode", "opencode", Vec::new(), ContextInjection::stdin(), None, "{projectRoot}", None, ) .unwrap() .with_structured_adapter(StructuredAdapter::OpenCode) .with_opencode( OpenCodeConfig::new("http://localhost:8080/v1", None, "qwen", None, None) .unwrap() .with_local_model_server_id(server_id), ) } struct FakeProbe(Mutex>); impl FakeProbe { fn new(statuses: Vec) -> Self { Self(Mutex::new(statuses.into())) } } #[async_trait] impl ModelServerProbe for FakeProbe { async fn probe( &self, _endpoint: &ModelServerEndpoint, ) -> Result { Ok(self .0 .lock() .unwrap() .pop_front() .unwrap_or(ModelServerStatus::Unreachable)) } } #[derive(Default)] struct FakeProcess { spawns: Mutex>, kills: Mutex>, statuses: Mutex>, status_sequence: Mutex>, spawn_delay: Duration, } #[async_trait] impl ManagedProcess for FakeProcess { async fn spawn(&self, spec: SpawnSpec) -> Result { if !self.spawn_delay.is_zero() { tokio::time::sleep(self.spawn_delay).await; } self.spawns.lock().unwrap().push(spec); let handle = ManagedProcessHandle { id: format!("h{}", self.spawns.lock().unwrap().len()), }; self.statuses .lock() .unwrap() .insert(handle.id.clone(), ProcessStatus::Running); Ok(handle) } async fn kill(&self, handle: &ManagedProcessHandle) -> Result<(), ModelServerError> { self.kills.lock().unwrap().push(handle.id.clone()); Ok(()) } async fn status( &self, handle: &ManagedProcessHandle, ) -> Result { if let Some(status) = self.status_sequence.lock().unwrap().pop_front() { return Ok(status); } Ok(*self .statuses .lock() .unwrap() .get(&handle.id) .unwrap_or(&ProcessStatus::Unknown)) } } struct FakeRuntime; impl ModelServerRuntime for FakeRuntime { fn build_argv( &self, config: &LocalModelServerConfig, ) -> Result { let mut args = Vec::new(); match config .model .source .as_ref() .ok_or_else(|| ModelServerError::PathNotAccessible("model.source missing".to_owned()))? { ModelSource::LocalPath { path } => { args.push("--model".to_owned()); args.push(path.as_str().to_owned()); } ModelSource::HuggingFace { repo } => { args.push("-hf".to_owned()); args.push(repo.as_str().to_owned()); } } args.extend([ "--port".to_owned(), config.endpoint.port.to_string(), "--host".to_owned(), config.options.host.clone(), ]); Ok(ModelServerArgv { command: config .binary .as_ref() .map(|binary| binary.as_str().to_owned()) .unwrap_or_else(|| "llama-server".to_owned()), args, }) } fn build_spawn_spec( &self, config: &LocalModelServerConfig, ) -> Result { let argv = self.build_argv(config)?; Ok(SpawnSpec { command: argv.command, args: argv.args, cwd: ProjectPath::new("/").unwrap(), env: Vec::new(), context_plan: None, sandbox: None, }) } } #[derive(Clone)] enum FakeDownloadOutcome { Resolve { progress: Vec, path: &'static str, cache_hit: bool, }, WaitForCancel, Sleep(Duration), } struct FakeModelArtifactDownloader { outcome: Mutex, } impl FakeModelArtifactDownloader { fn new(outcome: FakeDownloadOutcome) -> Self { Self { outcome: Mutex::new(outcome), } } } #[async_trait] impl ModelArtifactDownloader for FakeModelArtifactDownloader { async fn resolve_hf_model( &self, repo: &HfModelRef, progress: Arc, cancel: ModelArtifactCancel, ) -> Result { let outcome = self.outcome.lock().unwrap().clone(); match outcome { FakeDownloadOutcome::Resolve { progress: events, path, cache_hit, } => { if cancel.is_cancelled() { return Err(ModelServerError::Cancelled); } if !cache_hit { for event in events { progress(ModelArtifactProgress { source: event.source.or_else(|| Some(repo.as_str().to_owned())), ..event }); } } Ok(ModelArtifactResolution { path: ModelPath::new(path).unwrap(), cache_hit, }) } FakeDownloadOutcome::WaitForCancel => loop { if cancel.is_cancelled() { return Err(ModelServerError::Cancelled); } tokio::time::sleep(Duration::from_millis(1)).await; }, FakeDownloadOutcome::Sleep(duration) => { tokio::time::sleep(duration).await; Err(ModelServerError::Probe("unexpected wake".to_owned())) } } } } #[derive(Default)] struct FakeFs { existing: Mutex>, } #[async_trait] impl FileSystem for FakeFs { async fn read(&self, _path: &RemotePath) -> Result, FsError> { Err(FsError::NotFound("unused".to_owned())) } async fn write(&self, _path: &RemotePath, _data: &[u8]) -> Result<(), FsError> { Ok(()) } async fn exists(&self, path: &RemotePath) -> Result { Ok(self .existing .lock() .unwrap() .iter() .any(|p| p == path.as_str())) } async fn create_dir_all(&self, _path: &RemotePath) -> Result<(), FsError> { Ok(()) } async fn list(&self, _path: &RemotePath) -> Result, FsError> { Ok(Vec::new()) } async fn symlink(&self, _src: &RemotePath, _dst: &RemotePath) -> Result<(), FsError> { Ok(()) } } #[derive(Default)] struct FakeEvents(Mutex>); impl EventBus for FakeEvents { fn publish(&self, event: DomainEvent) { self.0.lock().unwrap().push(event); } fn subscribe(&self) -> EventStream { Box::new(std::iter::empty()) } } fn ensure( registry: Arc, probe: Arc, process: Arc, fs: Arc, events: Arc, ) -> EnsureLocalModelServer { EnsureLocalModelServer::new(registry, probe, process, Arc::new(FakeRuntime), fs, events) .with_readiness_policy(ModelServerReadinessPolicy { attempts: 2, backoff: Duration::ZERO, }) .with_hf_download_deadline(Duration::from_secs(5)) } fn ensure_with_downloader( registry: Arc, probe: Arc, process: Arc, fs: Arc, events: Arc, downloader: Arc, ) -> EnsureLocalModelServer { ensure(registry, probe, process, fs, events) .with_model_artifact_downloader(downloader as Arc) } fn progress(downloaded: Option, total: Option) -> ModelArtifactProgress { ModelArtifactProgress { downloaded_bytes: downloaded, total_bytes: total, source: None, } } #[tokio::test] async fn reachable_server_is_reused_without_spawn() { let registry = Arc::new(FakeRegistry::default()); registry .save(config(sid(1), 8080, "/models/qwen.gguf", true)) .await .unwrap(); let process = Arc::new(FakeProcess::default()); let usecase = ensure( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ModelServerStatus::ReadyReused])), Arc::clone(&process), Arc::new(FakeFs::default()), Arc::new(FakeEvents::default()), ); let out = usecase .execute(EnsureLocalModelServerInput { server_id: sid(1) }) .await .unwrap(); assert_eq!(out.ready.base_url, "http://localhost:8080/v1"); assert_eq!(out.ready.model, "qwen3-coder-30b"); assert!(process.spawns.lock().unwrap().is_empty()); } #[tokio::test] async fn absent_auto_start_spawns_and_waits_until_ready() { let registry = Arc::new(FakeRegistry::default()); registry .save(config(sid(2), 8081, "/models/qwen.gguf", true)) .await .unwrap(); let fs = Arc::new(FakeFs::default()); fs.existing .lock() .unwrap() .push("/models/qwen.gguf".to_owned()); let process = Arc::new(FakeProcess::default()); let usecase = ensure( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ ModelServerStatus::Unreachable, ModelServerStatus::ReadyReused, ])), Arc::clone(&process), fs, Arc::new(FakeEvents::default()), ); let out = usecase .execute(EnsureLocalModelServerInput { server_id: sid(2) }) .await .unwrap(); assert_eq!(out.ready.status, ModelServerStatus::ReadyStarted); let spawns = process.spawns.lock().unwrap(); assert_eq!(spawns.len(), 1); assert_eq!(spawns[0].command, "llama-server"); assert_eq!( spawns[0].args, vec![ "--model", "/models/qwen.gguf", "--port", "8081", "--host", "127.0.0.1" ] ); } #[tokio::test] async fn concurrent_ensure_same_server_shares_one_start_attempt() { let registry = Arc::new(FakeRegistry::default()); registry .save(config(sid(10), 8090, "/models/qwen.gguf", true)) .await .unwrap(); let fs = Arc::new(FakeFs::default()); fs.existing .lock() .unwrap() .push("/models/qwen.gguf".to_owned()); let process = Arc::new(FakeProcess { spawn_delay: Duration::from_millis(25), ..FakeProcess::default() }); let events = Arc::new(FakeEvents::default()); let usecase = Arc::new(ensure( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ ModelServerStatus::Unreachable, ModelServerStatus::ReadyReused, ])), Arc::clone(&process), fs, Arc::clone(&events), )); let mut tasks = Vec::new(); for _ in 0..5 { let usecase = Arc::clone(&usecase); tasks.push(tokio::spawn(async move { usecase .execute(EnsureLocalModelServerInput { server_id: sid(10) }) .await .unwrap() })); } let mut outputs = Vec::new(); for task in tasks { outputs.push(task.await.unwrap()); } assert!(outputs .iter() .all(|out| out.ready.status == ModelServerStatus::ReadyStarted)); assert_eq!(process.spawns.lock().unwrap().len(), 1); let starting_events = events .0 .lock() .unwrap() .iter() .filter(|event| { matches!( event, DomainEvent::ModelServerStatusChanged { status: domain::model_server::ModelServerLifecycleStatus::Starting, .. } ) }) .count(); assert_eq!(starting_events, 1); } #[tokio::test] async fn missing_model_path_is_path_not_accessible() { let registry = Arc::new(FakeRegistry::default()); registry .save(config(sid(3), 8082, "/models/missing.gguf", true)) .await .unwrap(); let usecase = ensure( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ModelServerStatus::Unreachable])), Arc::new(FakeProcess::default()), Arc::new(FakeFs::default()), Arc::new(FakeEvents::default()), ); let err = usecase .execute(EnsureLocalModelServerInput { server_id: sid(3) }) .await .unwrap_err(); assert_eq!(err.code(), "MODEL_SERVER"); assert!(err.to_string().contains("path_not_accessible")); } #[tokio::test] async fn active_managed_port_collision_is_explicit_error() { let registry = Arc::new(FakeRegistry::default()); registry .save(config(sid(4), 8083, "/models/a.gguf", true)) .await .unwrap(); registry .save(config(sid(5), 8083, "/models/b.gguf", true)) .await .unwrap(); let fs = Arc::new(FakeFs::default()); fs.existing .lock() .unwrap() .extend(["/models/a.gguf".to_owned(), "/models/b.gguf".to_owned()]); let process = Arc::new(FakeProcess::default()); let events = Arc::new(FakeEvents::default()); let usecase = ensure( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ ModelServerStatus::Unreachable, ModelServerStatus::ReadyReused, ModelServerStatus::Unreachable, ])), Arc::clone(&process), fs, events, ); usecase .execute(EnsureLocalModelServerInput { server_id: sid(4) }) .await .unwrap(); let err = usecase .execute(EnsureLocalModelServerInput { server_id: sid(5) }) .await .unwrap_err(); assert!(err.to_string().contains("port_occupied")); assert_eq!(process.spawns.lock().unwrap().len(), 1); } #[tokio::test] async fn readiness_timeout_kills_started_process() { let registry = Arc::new(FakeRegistry::default()); registry .save(config(sid(6), 8084, "/models/qwen.gguf", true)) .await .unwrap(); let fs = Arc::new(FakeFs::default()); fs.existing .lock() .unwrap() .push("/models/qwen.gguf".to_owned()); let process = Arc::new(FakeProcess::default()); let usecase = ensure( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ ModelServerStatus::Unreachable, ModelServerStatus::Unreachable, ModelServerStatus::Unreachable, ])), Arc::clone(&process), fs, Arc::new(FakeEvents::default()), ); let err = usecase .execute(EnsureLocalModelServerInput { server_id: sid(6) }) .await .unwrap_err(); assert!(err.to_string().contains("timeout")); assert_eq!(process.spawns.lock().unwrap().len(), 1); assert_eq!(process.kills.lock().unwrap().as_slice(), ["h1"]); } #[tokio::test] async fn hf_unreachable_alive_publishes_downloading_without_short_timeout_kill() { let registry = Arc::new(FakeRegistry::default()); registry .save(hf_config( sid(11), 8091, "Qwen/Qwen3-Coder-30B-A3B-Instruct-GGUF", )) .await .unwrap(); let process = Arc::new(FakeProcess::default()); let events = Arc::new(FakeEvents::default()); let usecase = ensure( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ ModelServerStatus::Unreachable, ModelServerStatus::Unreachable, ModelServerStatus::Unreachable, ModelServerStatus::ReadyReused, ])), Arc::clone(&process), Arc::new(FakeFs::default()), Arc::clone(&events), ); let out = usecase .execute(EnsureLocalModelServerInput { server_id: sid(11) }) .await .unwrap(); assert_eq!(out.ready.status, ModelServerStatus::ReadyStarted); assert!(process.kills.lock().unwrap().is_empty()); assert!(events.0.lock().unwrap().iter().any(|event| { matches!( event, DomainEvent::ModelServerStatusChanged { server_id, status: ModelServerLifecycleStatus::Downloading { downloaded_bytes: None, total_bytes: None, percent: None, source: Some(source), }, } if *server_id == sid(11) && source == "Qwen/Qwen3-Coder-30B-A3B-Instruct-GGUF" ) })); } #[tokio::test] async fn hf_download_phase_reports_ready_when_probe_becomes_ok() { let registry = Arc::new(FakeRegistry::default()); registry .save(hf_config( sid(12), 8092, "Qwen/Qwen3-Coder-30B-A3B-Instruct-GGUF", )) .await .unwrap(); let process = Arc::new(FakeProcess::default()); let events = Arc::new(FakeEvents::default()); let usecase = ensure( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ ModelServerStatus::Unreachable, ModelServerStatus::Unreachable, ModelServerStatus::Unreachable, ModelServerStatus::ReadyStarted, ])), Arc::clone(&process), Arc::new(FakeFs::default()), Arc::clone(&events), ); let out = usecase .execute(EnsureLocalModelServerInput { server_id: sid(12) }) .await .unwrap(); assert_eq!(out.ready.status, ModelServerStatus::ReadyStarted); let statuses: Vec = events .0 .lock() .unwrap() .iter() .filter_map(|event| match event { DomainEvent::ModelServerStatusChanged { status, .. } => Some(status.clone()), _ => None, }) .collect(); assert!(matches!( statuses.as_slice(), [ ModelServerLifecycleStatus::Probing, ModelServerLifecycleStatus::Starting, ModelServerLifecycleStatus::Downloading { .. }, ModelServerLifecycleStatus::Ready { reused: false }, ] )); } #[tokio::test] async fn hf_download_phase_fails_when_process_exits() { let registry = Arc::new(FakeRegistry::default()); registry .save(hf_config( sid(13), 8093, "Qwen/Qwen3-Coder-30B-A3B-Instruct-GGUF", )) .await .unwrap(); let process = Arc::new(FakeProcess::default()); process .status_sequence .lock() .unwrap() .push_back(ProcessStatus::Exited { code: Some(42) }); let events = Arc::new(FakeEvents::default()); let usecase = ensure( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ ModelServerStatus::Unreachable, ModelServerStatus::Unreachable, ModelServerStatus::Unreachable, ])), Arc::clone(&process), Arc::new(FakeFs::default()), Arc::clone(&events), ); let err = usecase .execute(EnsureLocalModelServerInput { server_id: sid(13) }) .await .unwrap_err(); assert!(err.to_string().contains("process")); assert!(err.to_string().contains("42")); assert!(process.kills.lock().unwrap().is_empty()); assert!(events.0.lock().unwrap().iter().any(|event| { matches!( event, DomainEvent::ModelServerStatusChanged { status: ModelServerLifecycleStatus::Failed { code, .. }, .. } if code == "process" ) })); } #[tokio::test] async fn hf_downloader_publishes_debounced_progress_and_spawns_resolved_model_path() { let registry = Arc::new(FakeRegistry::default()); registry .save(hf_config( sid(14), 8094, "Qwen/Qwen3-Coder-30B-A3B-Instruct-GGUF", )) .await .unwrap(); let process = Arc::new(FakeProcess::default()); let events = Arc::new(FakeEvents::default()); let downloader = Arc::new(FakeModelArtifactDownloader::new( FakeDownloadOutcome::Resolve { progress: vec![ progress(Some(10), Some(100)), progress(Some(50), Some(100)), progress(Some(100), Some(100)), ], path: "/cache/model.gguf", cache_hit: false, }, )); let usecase = ensure_with_downloader( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ ModelServerStatus::Unreachable, ModelServerStatus::ReadyStarted, ])), Arc::clone(&process), Arc::new(FakeFs::default()), Arc::clone(&events), downloader, ); let out = usecase .execute(EnsureLocalModelServerInput { server_id: sid(14) }) .await .unwrap(); assert_eq!(out.ready.status, ModelServerStatus::ReadyStarted); let statuses: Vec = events .0 .lock() .unwrap() .iter() .filter_map(|event| match event { DomainEvent::ModelServerStatusChanged { status, .. } => Some(status.clone()), _ => None, }) .collect(); assert!(matches!( statuses.as_slice(), [ ModelServerLifecycleStatus::Probing, ModelServerLifecycleStatus::Downloading { downloaded_bytes: Some(10), total_bytes: Some(100), percent: Some(10.0), .. }, ModelServerLifecycleStatus::Downloading { downloaded_bytes: Some(100), total_bytes: Some(100), percent: Some(100.0), .. }, ModelServerLifecycleStatus::Starting, ModelServerLifecycleStatus::Ready { reused: false }, ] )); let spawns = process.spawns.lock().unwrap(); assert_eq!(spawns.len(), 1); assert_eq!(spawns[0].args[0..2], ["--model", "/cache/model.gguf"]); assert!(!spawns[0].args.iter().any(|arg| arg == "-hf")); } #[tokio::test] async fn hf_downloader_cache_hit_skips_downloading_events() { let registry = Arc::new(FakeRegistry::default()); registry .save(hf_config(sid(15), 8095, "Qwen/Qwen3-Coder:Q4_K_M")) .await .unwrap(); let process = Arc::new(FakeProcess::default()); let events = Arc::new(FakeEvents::default()); let downloader = Arc::new(FakeModelArtifactDownloader::new( FakeDownloadOutcome::Resolve { progress: Vec::new(), path: "/cache/q4.gguf", cache_hit: true, }, )); let usecase = ensure_with_downloader( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ ModelServerStatus::Unreachable, ModelServerStatus::ReadyStarted, ])), Arc::clone(&process), Arc::new(FakeFs::default()), Arc::clone(&events), downloader, ); usecase .execute(EnsureLocalModelServerInput { server_id: sid(15) }) .await .unwrap(); let statuses: Vec = events .0 .lock() .unwrap() .iter() .filter_map(|event| match event { DomainEvent::ModelServerStatusChanged { status, .. } => Some(status.clone()), _ => None, }) .collect(); assert!(matches!( statuses.as_slice(), [ ModelServerLifecycleStatus::Probing, ModelServerLifecycleStatus::Starting, ModelServerLifecycleStatus::Ready { reused: false }, ] )); } #[tokio::test] async fn hf_downloader_unknown_total_publishes_no_percent() { let registry = Arc::new(FakeRegistry::default()); registry .save(hf_config(sid(16), 8096, "Qwen/Qwen3-Coder:Q4_K_M")) .await .unwrap(); let process = Arc::new(FakeProcess::default()); let events = Arc::new(FakeEvents::default()); let downloader = Arc::new(FakeModelArtifactDownloader::new( FakeDownloadOutcome::Resolve { progress: vec![progress(Some(10), None)], path: "/cache/q4.gguf", cache_hit: false, }, )); let usecase = ensure_with_downloader( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ ModelServerStatus::Unreachable, ModelServerStatus::ReadyStarted, ])), Arc::clone(&process), Arc::new(FakeFs::default()), Arc::clone(&events), downloader, ); usecase .execute(EnsureLocalModelServerInput { server_id: sid(16) }) .await .unwrap(); assert!(events.0.lock().unwrap().iter().any(|event| { matches!( event, DomainEvent::ModelServerStatusChanged { status: ModelServerLifecycleStatus::Downloading { downloaded_bytes: Some(10), total_bytes: None, percent: None, .. }, .. } ) })); } #[tokio::test] async fn hf_downloader_cancel_fails_without_process_spawn() { let registry = Arc::new(FakeRegistry::default()); registry .save(hf_config(sid(17), 8097, "Qwen/Qwen3-Coder:Q4_K_M")) .await .unwrap(); let process = Arc::new(FakeProcess::default()); let events = Arc::new(FakeEvents::default()); let downloader = Arc::new(FakeModelArtifactDownloader::new( FakeDownloadOutcome::WaitForCancel, )); let usecase = Arc::new(ensure_with_downloader( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ModelServerStatus::Unreachable])), Arc::clone(&process), Arc::new(FakeFs::default()), Arc::clone(&events), downloader, )); let task_usecase = Arc::clone(&usecase); let task = tokio::spawn(async move { task_usecase .execute(EnsureLocalModelServerInput { server_id: sid(17) }) .await }); tokio::time::sleep(Duration::from_millis(5)).await; usecase.stop_on_app_exit().await.unwrap(); let err = task.await.unwrap().unwrap_err(); assert!(err.to_string().contains("cancelled")); assert!(process.spawns.lock().unwrap().is_empty()); assert!(events.0.lock().unwrap().iter().any(|event| { matches!( event, DomainEvent::ModelServerStatusChanged { status: ModelServerLifecycleStatus::Failed { code, .. }, .. } if code == "cancelled" ) })); } #[tokio::test] async fn hf_downloader_timeout_fails_without_process_spawn() { let registry = Arc::new(FakeRegistry::default()); registry .save(hf_config(sid(18), 8098, "Qwen/Qwen3-Coder:Q4_K_M")) .await .unwrap(); let process = Arc::new(FakeProcess::default()); let events = Arc::new(FakeEvents::default()); let downloader = Arc::new(FakeModelArtifactDownloader::new( FakeDownloadOutcome::Sleep(Duration::from_millis(50)), )); let usecase = ensure_with_downloader( Arc::clone(®istry), Arc::new(FakeProbe::new(vec![ModelServerStatus::Unreachable])), Arc::clone(&process), Arc::new(FakeFs::default()), Arc::clone(&events), downloader, ) .with_hf_download_deadline(Duration::from_millis(1)); let err = usecase .execute(EnsureLocalModelServerInput { server_id: sid(18) }) .await .unwrap_err(); assert!(err.to_string().contains("timeout")); assert!(process.spawns.lock().unwrap().is_empty()); assert!(events.0.lock().unwrap().iter().any(|event| { matches!( event, DomainEvent::ModelServerStatusChanged { status: ModelServerLifecycleStatus::Failed { code, .. }, .. } if code == "timeout" ) })); } #[tokio::test] async fn missing_registry_entry_is_model_server_not_configured() { 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 err = usecase .execute(EnsureLocalModelServerInput { server_id: sid(7) }) .await .unwrap_err(); assert_eq!(err.code(), "MODEL_SERVER"); assert!(err.to_string().contains("not_configured")); } #[tokio::test] async fn delete_model_server_refuses_when_profile_references_it() { let registry = Arc::new(FakeRegistry::default()); registry .save(config(sid(8), 8085, "/models/qwen.gguf", false)) .await .unwrap(); let profiles = Arc::new(FakeProfiles(Mutex::new(vec![opencode_profile(18, sid(8))]))); let usecase = DeleteModelServer::new(registry, profiles); let err = usecase .execute(DeleteModelServerInput { server_id: sid(8) }) .await .unwrap_err(); assert_eq!(err.code(), "MODEL_SERVER"); match err { application::AppError::ModelServer { code, .. } => { assert_eq!(code, "model_server_in_use"); } other => panic!("unexpected error: {other}"), } } #[tokio::test] async fn delete_model_server_removes_unused_config() { let registry = Arc::new(FakeRegistry::default()); registry .save(config(sid(9), 8086, "/models/qwen.gguf", false)) .await .unwrap(); let profiles = Arc::new(FakeProfiles::default()); let usecase = DeleteModelServer::new( Arc::clone(®istry) as Arc, profiles, ); usecase .execute(DeleteModelServerInput { server_id: sid(9) }) .await .unwrap(); assert!(registry.get(&sid(9)).await.unwrap().is_none()); }