Corrige la régression high « Error on loading local model » : le premier chargement du serveur modèle local (cold-start llama.cpp) dépassait la fenêtre de readiness et échouait. Le warmup dispose désormais d'un deadline par défaut de 600 s, surchargable par config optionnelle `warmup_deadline_secs` (validée dans [30, 1800]). - domain: champ `warmup_deadline_secs: Option<u64>` + validation de borne - application: policy de readiness effective (défaut 600 s, override par config) - app-tauri: DTO `warmupDeadlineSecs` - infrastructure: application de la deadline effective au warmup Verdict QA (vert) : domain 252, application 81 + 22 model_server, app-tauri dto_model_servers 6, infrastructure model_server 2, build OK. Contrat readiness couvert sur ports mockés. Caveat : le cold-start end-to-end réel (llama-server) sort du sandbox de test -> vérification manuelle utilisateur restante, non couverte par les tests unitaires. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1328 lines
39 KiB
Rust
1328 lines
39 KiB
Rust
//! 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<HashMap<LocalModelServerId, LocalModelServerConfig>>);
|
|
|
|
#[async_trait]
|
|
impl ModelServerRegistry for FakeRegistry {
|
|
async fn get(
|
|
&self,
|
|
id: &LocalModelServerId,
|
|
) -> Result<Option<LocalModelServerConfig>, ModelServerError> {
|
|
Ok(self.0.lock().unwrap().get(id).cloned())
|
|
}
|
|
|
|
async fn list(&self) -> Result<Vec<LocalModelServerConfig>, 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<Vec<AgentProfile>>);
|
|
|
|
#[async_trait]
|
|
impl ProfileStore for FakeProfiles {
|
|
async fn list(&self) -> Result<Vec<AgentProfile>, 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<bool, StoreError> {
|
|
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<VecDeque<ModelServerStatus>>);
|
|
|
|
impl FakeProbe {
|
|
fn new(statuses: Vec<ModelServerStatus>) -> Self {
|
|
Self(Mutex::new(statuses.into()))
|
|
}
|
|
}
|
|
|
|
#[async_trait]
|
|
impl ModelServerProbe for FakeProbe {
|
|
async fn probe(
|
|
&self,
|
|
_endpoint: &ModelServerEndpoint,
|
|
) -> Result<ModelServerStatus, ModelServerError> {
|
|
Ok(self
|
|
.0
|
|
.lock()
|
|
.unwrap()
|
|
.pop_front()
|
|
.unwrap_or(ModelServerStatus::Unreachable))
|
|
}
|
|
}
|
|
|
|
#[derive(Default)]
|
|
struct FakeProcess {
|
|
spawns: Mutex<Vec<SpawnSpec>>,
|
|
kills: Mutex<Vec<String>>,
|
|
statuses: Mutex<HashMap<String, ProcessStatus>>,
|
|
status_sequence: Mutex<VecDeque<ProcessStatus>>,
|
|
spawn_delay: Duration,
|
|
}
|
|
|
|
#[async_trait]
|
|
impl ManagedProcess for FakeProcess {
|
|
async fn spawn(&self, spec: SpawnSpec) -> Result<ManagedProcessHandle, ModelServerError> {
|
|
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<ProcessStatus, ModelServerError> {
|
|
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<ModelServerArgv, ModelServerError> {
|
|
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<SpawnSpec, ModelServerError> {
|
|
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<ModelArtifactProgress>,
|
|
path: &'static str,
|
|
cache_hit: bool,
|
|
},
|
|
WaitForCancel,
|
|
Sleep(Duration),
|
|
}
|
|
|
|
struct FakeModelArtifactDownloader {
|
|
outcome: Mutex<FakeDownloadOutcome>,
|
|
}
|
|
|
|
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<dyn Fn(ModelArtifactProgress) + Send + Sync>,
|
|
cancel: ModelArtifactCancel,
|
|
) -> Result<ModelArtifactResolution, ModelServerError> {
|
|
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<Vec<String>>,
|
|
}
|
|
|
|
#[async_trait]
|
|
impl FileSystem for FakeFs {
|
|
async fn read(&self, _path: &RemotePath) -> Result<Vec<u8>, 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<bool, FsError> {
|
|
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<Vec<DirEntry>, FsError> {
|
|
Ok(Vec::new())
|
|
}
|
|
|
|
async fn symlink(&self, _src: &RemotePath, _dst: &RemotePath) -> Result<(), FsError> {
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
#[derive(Default)]
|
|
struct FakeEvents(Mutex<Vec<DomainEvent>>);
|
|
|
|
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<FakeRegistry>,
|
|
probe: Arc<FakeProbe>,
|
|
process: Arc<FakeProcess>,
|
|
fs: Arc<FakeFs>,
|
|
events: Arc<FakeEvents>,
|
|
) -> EnsureLocalModelServer {
|
|
EnsureLocalModelServer::new(registry, probe, process, Arc::new(FakeRuntime), fs, events)
|
|
.with_readiness_policy(ModelServerReadinessPolicy {
|
|
attempts: 2,
|
|
backoff: Duration::from_millis(1),
|
|
warmup_deadline: Duration::from_millis(25),
|
|
})
|
|
.with_hf_download_deadline(Duration::from_secs(5))
|
|
}
|
|
|
|
fn ensure_with_downloader(
|
|
registry: Arc<FakeRegistry>,
|
|
probe: Arc<FakeProbe>,
|
|
process: Arc<FakeProcess>,
|
|
fs: Arc<FakeFs>,
|
|
events: Arc<FakeEvents>,
|
|
downloader: Arc<FakeModelArtifactDownloader>,
|
|
) -> EnsureLocalModelServer {
|
|
ensure(registry, probe, process, fs, events)
|
|
.with_model_artifact_downloader(downloader as Arc<dyn ModelArtifactDownloader>)
|
|
}
|
|
|
|
fn progress(downloaded: Option<u64>, total: Option<u64>) -> ModelArtifactProgress {
|
|
ModelArtifactProgress {
|
|
downloaded_bytes: downloaded,
|
|
total_bytes: total,
|
|
source: None,
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn readiness_policy_default_warmup_deadline_is_ten_minutes() {
|
|
assert_eq!(
|
|
ModelServerReadinessPolicy::default().warmup_deadline,
|
|
Duration::from_secs(600)
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn effective_readiness_policy_uses_configured_warmup_deadline_when_present() {
|
|
let configured = config(sid(23), 8103, "/models/qwen.gguf", true)
|
|
.with_warmup_deadline_secs(Some(900))
|
|
.unwrap();
|
|
let unconfigured = config(sid(24), 8104, "/models/qwen.gguf", true);
|
|
let usecase = EnsureLocalModelServer::new(
|
|
Arc::new(FakeRegistry::default()),
|
|
Arc::new(FakeProbe::new(Vec::new())),
|
|
Arc::new(FakeProcess::default()),
|
|
Arc::new(FakeRuntime),
|
|
Arc::new(FakeFs::default()),
|
|
Arc::new(FakeEvents::default()),
|
|
);
|
|
|
|
assert_eq!(
|
|
usecase
|
|
.effective_readiness_policy(&configured)
|
|
.warmup_deadline,
|
|
Duration::from_secs(900)
|
|
);
|
|
assert_eq!(
|
|
usecase
|
|
.effective_readiness_policy(&unconfigured)
|
|
.warmup_deadline,
|
|
Duration::from_secs(600)
|
|
);
|
|
}
|
|
|
|
#[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 slow_local_warmup_exceeding_short_readiness_window_succeeds_without_stop() {
|
|
let registry = Arc::new(FakeRegistry::default());
|
|
registry
|
|
.save(config(sid(19), 8099, "/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 mut statuses = vec![ModelServerStatus::Unreachable];
|
|
statuses.extend(std::iter::repeat(ModelServerStatus::Unreachable).take(22));
|
|
statuses.push(ModelServerStatus::ReadyStarted);
|
|
let registry_port: Arc<dyn ModelServerRegistry> = registry.clone();
|
|
let process_port: Arc<dyn ManagedProcess> = process.clone();
|
|
let usecase = EnsureLocalModelServer::new(
|
|
registry_port,
|
|
Arc::new(FakeProbe::new(statuses)),
|
|
process_port,
|
|
Arc::new(FakeRuntime),
|
|
fs,
|
|
Arc::new(FakeEvents::default()),
|
|
)
|
|
.with_readiness_policy(ModelServerReadinessPolicy {
|
|
attempts: 20,
|
|
backoff: Duration::from_millis(1),
|
|
warmup_deadline: Duration::from_millis(100),
|
|
});
|
|
|
|
let out = usecase
|
|
.execute(EnsureLocalModelServerInput { server_id: sid(19) })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert_eq!(out.ready.status, ModelServerStatus::ReadyStarted);
|
|
assert_eq!(process.spawns.lock().unwrap().len(), 1);
|
|
assert!(process.kills.lock().unwrap().is_empty());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn process_exit_during_local_warmup_fails_fast_and_cleans_active_registry() {
|
|
let registry = Arc::new(FakeRegistry::default());
|
|
registry
|
|
.save(config(sid(20), 8100, "/models/a.gguf", true))
|
|
.await
|
|
.unwrap();
|
|
registry
|
|
.save(config(sid(21), 8100, "/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());
|
|
process
|
|
.status_sequence
|
|
.lock()
|
|
.unwrap()
|
|
.push_back(ProcessStatus::Exited { code: Some(42) });
|
|
let usecase = ensure(
|
|
Arc::clone(®istry),
|
|
Arc::new(FakeProbe::new(vec![
|
|
ModelServerStatus::Unreachable,
|
|
ModelServerStatus::Unreachable,
|
|
ModelServerStatus::Unreachable,
|
|
ModelServerStatus::ReadyStarted,
|
|
])),
|
|
Arc::clone(&process),
|
|
fs,
|
|
Arc::new(FakeEvents::default()),
|
|
);
|
|
|
|
let err = usecase
|
|
.execute(EnsureLocalModelServerInput { server_id: sid(20) })
|
|
.await
|
|
.unwrap_err();
|
|
|
|
assert!(err.to_string().contains("process"));
|
|
assert!(err.to_string().contains("42"));
|
|
assert!(process.kills.lock().unwrap().is_empty());
|
|
|
|
let out = usecase
|
|
.execute(EnsureLocalModelServerInput { server_id: sid(21) })
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(out.ready.status, ModelServerStatus::ReadyStarted);
|
|
assert_eq!(process.spawns.lock().unwrap().len(), 2);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn warmup_deadline_reached_returns_timeout_and_stops_started_process() {
|
|
let registry = Arc::new(FakeRegistry::default());
|
|
registry
|
|
.save(config(sid(22), 8101, "/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 registry_port: Arc<dyn ModelServerRegistry> = registry.clone();
|
|
let process_port: Arc<dyn ManagedProcess> = process.clone();
|
|
let usecase = EnsureLocalModelServer::new(
|
|
registry_port,
|
|
Arc::new(FakeProbe::new(vec![ModelServerStatus::Unreachable])),
|
|
process_port,
|
|
Arc::new(FakeRuntime),
|
|
fs,
|
|
Arc::new(FakeEvents::default()),
|
|
)
|
|
.with_readiness_policy(ModelServerReadinessPolicy {
|
|
attempts: 20,
|
|
backoff: Duration::from_millis(1),
|
|
warmup_deadline: Duration::from_millis(3),
|
|
});
|
|
|
|
let err = usecase
|
|
.execute(EnsureLocalModelServerInput { server_id: sid(22) })
|
|
.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<ModelServerLifecycleStatus> = 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<ModelServerLifecycleStatus> = 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<ModelServerLifecycleStatus> = 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<dyn ModelServerRegistry>,
|
|
profiles,
|
|
);
|
|
|
|
usecase
|
|
.execute(DeleteModelServerInput { server_id: sid(9) })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert!(registry.get(&sid(9)).await.unwrap().is_none());
|
|
}
|