feat(model-server): progression fine du téléchargement du modèle llamacpp (#54)

Stretch B2/F2 de #54, par-dessus le MVP déjà mergé (B1/F1).

Backend : le port de téléchargement HF publie une progression débouncée
(bytes reçus / total, pourcentage) via le stream de statut du serveur
modèle, avec gestion du total inconnu (pas de faux %), du cache hit,
de l'annulation et du timeout.

Frontend : l'overlay plein-cellule de préparation du serveur affiche la
progression réelle (barre, %, octets, source) en mappant le fil de
statut, avec la règle « pas de faux % » quand le total est inconnu.

Tests : application + infrastructure (téléchargement débouncé, cancel,
timeout, cache hit, total inconnu) et vitest (overlay + formatage pur).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-07-14 10:37:14 +02:00
parent bd335c3a1c
commit 2183dfd291
15 changed files with 1368 additions and 87 deletions

View File

@ -18,6 +18,7 @@ use domain::model_server::{
};
use domain::ports::{
DirEntry, EventBus, EventStream, FileSystem, FsError, ManagedProcess, ManagedProcessHandle,
ModelArtifactCancel, ModelArtifactDownloader, ModelArtifactProgress, ModelArtifactResolution,
ModelServerArgv, ModelServerError, ModelServerProbe, ModelServerRegistry, ModelServerRuntime,
ProcessStatus, ProfileStore, RemotePath, SpawnSpec, StoreError,
};
@ -280,6 +281,74 @@ impl ModelServerRuntime for FakeRuntime {
}
}
#[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>>,
@ -345,6 +414,26 @@ fn ensure(
.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,
}
}
#[tokio::test]
async fn reachable_server_is_reused_without_spawn() {
let registry = Arc::new(FakeRegistry::default());
@ -738,6 +827,270 @@ async fn hf_download_phase_fails_when_process_exits() {
}));
}
#[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(&registry),
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(&registry),
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(&registry),
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(&registry),
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(&registry),
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(