diff --git a/apps/rocm/src/main.rs b/apps/rocm/src/main.rs index 09d49be4f..a4e0354bd 100644 --- a/apps/rocm/src/main.rs +++ b/apps/rocm/src/main.rs @@ -6941,22 +6941,7 @@ fn start_managed_service( #[cfg(windows)] thread::sleep(Duration::from_millis(200)); - let readiness = wait_for_service_http_ready_with_progress( - engine, - host, - port, - &resolve.canonical_model_id, - endpoint_api_key, - Duration::from_secs(45), - on_wait_tick, - ); - let launch_status = status_for_readiness(readiness); - record.status = launch_status.to_owned(); - if readiness == EndpointReadiness::Serving { - // Latch the verification the wait just performed, so the readiness checks - // behind `services list` and chat read it instead of re-probing. - record.inference_verified_at_unix_ms = Some(rocm_core::unix_time_millis() as u64); - } + await_managed_readiness(&mut record, endpoint_api_key, on_wait_tick); record.write()?; let endpoint_url = format!("{}/v1", format_http_base_url(host, port)); record_cli_audit_event( @@ -6966,14 +6951,14 @@ fn start_managed_service( "info", format!( "launched managed service engine={} model={} endpoint={} readiness={}", - engine, resolve.canonical_model_id, endpoint_url, launch_status + engine, resolve.canonical_model_id, endpoint_url, record.status ), Some(service_id), ); Ok(ManagedLaunchReport { service_id: service_id.to_owned(), endpoint_url, - status: launch_status.to_owned(), + status: record.status, already_running: false, child_pid: Some(child_pid), log_path: Some(record.log_path), @@ -18306,6 +18291,10 @@ fn restart_internal_managed_service( fs::create_dir_all(parent) .with_context(|| format!("failed to create {}", parent.display()))?; } + // The readiness wait below ends as soon as the engine state file records a + // terminal status. The previous run's file — often `failed`, which is why it + // is being restarted — would end it before the new child writes its own. + let _ = fs::remove_file(&record.engine_state_path); let current_exe = managed_service_launcher_path() .context("failed to resolve current rocm executable path")?; let recipe = parse_engine_recipe_json_arg(record.engine_recipe_json.clone())?; @@ -18385,18 +18374,11 @@ fn restart_internal_managed_service( // the new child has an unloaded model, so the old verdict says nothing about // it. record.reset_for_restart(); - let readiness = wait_for_service_http_ready( - &record.engine, - &record.host, - record.port, - &record.canonical_model_id, - endpoint_api_key.as_deref(), - Duration::from_secs(45), - ); - record.status = status_for_readiness(readiness).to_owned(); - if readiness == EndpointReadiness::Serving { - record.inference_verified_at_unix_ms = Some(rocm_core::unix_time_millis() as u64); - } + // Persist the new child before the wait, which can run for minutes: a caller + // that kills this command mid-wait must not leave the record naming the old + // supervisor while the new one runs. + record.write()?; + await_managed_readiness(&mut record, endpoint_api_key.as_deref(), &mut |_elapsed| {}); record.write()?; Ok(record) } @@ -21922,6 +21904,61 @@ fn apply_app_path_env(command: &mut ProcessCommand, paths: &AppPaths) { } } +/// Readiness budget for a managed launch, taken from the engine's own startup +/// timeout so the CLI never contradicts it. vLLM's is 5 minutes by default and +/// tunable via `ROCM_CLI_VLLM_READY_TIMEOUT_SECS`; other engines start fast +/// enough that the historical 45 s is still generous. +fn managed_ready_timeout(engine: &str) -> Duration { + match engine { + "vllm" => rocm_engine_vllm::ready_timeout(), + _ => Duration::from_secs(45), + } +} + +/// Wait for a freshly (re)spawned managed service and settle `record.status`. +/// +/// Waits on the engine's own budget ([`managed_ready_timeout`]). Each tick +/// re-reads the engine state file so a crash ends the wait instead of spinning +/// out that whole budget: the supervisor is our child, so `process_is_running` +/// keeps reporting it alive while it sits unreaped as a zombie, but the +/// supervisor writes a terminal status when its engine exits. Only if the +/// supervisor itself is killed is nothing written; the wait then runs out the +/// budget and reports how far the endpoint got. +fn await_managed_readiness( + record: &mut ManagedServiceRecord, + endpoint_api_key: Option<&str>, + on_wait_tick: &mut dyn FnMut(Duration), +) { + let engine = record.engine.clone(); + let host = record.host.clone(); + let model = record.canonical_model_id.clone(); + let readiness = wait_for_service_http_ready_with_progress( + &engine, + &host, + record.port, + &model, + endpoint_api_key, + managed_ready_timeout(&engine), + &mut |elapsed| { + on_wait_tick(elapsed); + let _ = record.refresh_from_engine_state(); + managed_service_running_state(&record.status) != "not_running" + }, + ); + // `readiness` is the furthest the endpoint got during the wait, not where it + // is now. A terminal status the engine recorded since is newer, so it stands: + // a model that listed and then died in warmup has failed, not "still loading". + if managed_service_running_state(&record.status) != "not_running" { + record.status = status_for_readiness(readiness).to_owned(); + } + if readiness == EndpointReadiness::Serving { + // Latch the verification the wait just performed, so the readiness checks + // behind `services list` and chat read it instead of re-probing. + record.inference_verified_at_unix_ms = Some(rocm_core::unix_time_millis() as u64); + } +} + +#[cfg(test)] fn wait_for_service_http_ready( engine: &str, host: &str, @@ -21937,7 +21974,7 @@ fn wait_for_service_http_ready( canonical_model_id, endpoint_api_key, timeout, - &mut |_elapsed| {}, + &mut |_elapsed| true, ) } @@ -21952,10 +21989,12 @@ fn wait_for_service_http_ready( /// result seen before `timeout` is what gets returned. /// /// Polls until `timeout` elapses, invoking `on_tick(elapsed)` once per iteration -/// so a caller can animate a spinner. Engine-neutral: `service_http_readiness_paths` -/// maps each engine to the right listing path. `endpoint_api_key` is sent as a -/// bearer token so both the listing check and the inference probe still succeed -/// against a public endpoint that requires authentication. +/// so a caller can animate a spinner and, by returning `false`, abandon a wait it +/// knows is hopeless (e.g. the engine died) rather than burn the whole budget. +/// Engine-neutral: `service_http_readiness_paths` maps each engine to the right +/// listing path. `endpoint_api_key` is sent as a bearer token so both the listing +/// check and the inference probe still succeed against a public endpoint that +/// requires authentication. fn wait_for_service_http_ready_with_progress( engine: &str, host: &str, @@ -21963,10 +22002,13 @@ fn wait_for_service_http_ready_with_progress( canonical_model_id: &str, endpoint_api_key: Option<&str>, timeout: Duration, - on_tick: &mut dyn FnMut(Duration), + on_tick: &mut dyn FnMut(Duration) -> bool, ) -> EndpointReadiness { let start = std::time::Instant::now(); let endpoint = format_http_base_url(host, port); + // High-water mark: once the model has listed this never drops back, even if + // the endpoint later stops answering. Callers that need "where it is now" + // must look elsewhere (see `await_managed_readiness`). let mut best = EndpointReadiness::Unreachable; while start.elapsed() < timeout { let listed = service_http_readiness_paths(engine).iter().any(|path| { @@ -22000,7 +22042,9 @@ fn wait_for_service_http_ready_with_progress( return EndpointReadiness::Serving; } } - on_tick(start.elapsed()); + if !on_tick(start.elapsed()) { + return best; + } thread::sleep(Duration::from_millis(250)); } best @@ -22867,6 +22911,88 @@ mod tests { ); } + /// Regression: the managed launch used a hardcoded 45 s wait while vLLM's own + /// startup budget is minutes, so every cold vLLM start was reported as "not + /// ready" seconds before the server came up. + #[test] + fn managed_ready_timeout_follows_the_engine_budget() { + assert_eq!( + managed_ready_timeout("vllm"), + rocm_engine_vllm::ready_timeout() + ); + assert!(managed_ready_timeout("vllm") > Duration::from_secs(45)); + assert_eq!(managed_ready_timeout("lemonade"), Duration::from_secs(45)); + } + + /// A launch the engine gave up on is `failed`, even after its endpoint listed + /// the model, and the wait ends on the engine's terminal status instead of + /// running out the (minutes-long) vLLM budget. The readiness the wait returns + /// is a high-water mark, so `Listing` alone must not read as "still loading". + #[test] + fn managed_readiness_reports_an_engine_that_died_after_listing_as_failed() -> Result<()> { + use std::io::{Read, Write}; + use std::net::TcpListener; + + // Lists the model but cannot serve it: what vLLM looks like when it dies + // during KV-cache warmup, after the OpenAI server is already up. + let listener = TcpListener::bind(("127.0.0.1", 0))?; + let port = listener.local_addr()?.port(); + let server = thread::spawn(move || { + while let Ok((mut stream, _)) = listener.accept() { + stream.set_read_timeout(Some(Duration::from_secs(2))).ok(); + let mut buffer = [0_u8; 1024]; + let Ok(read) = stream.read(&mut buffer) else { + continue; + }; + let request = String::from_utf8_lossy(&buffer[..read]).into_owned(); + let (status_line, body) = if request.starts_with("POST /v1/chat/completions ") { + ("HTTP/1.1 503 Service Unavailable", r#"{"error":"loading"}"#) + } else { + ("HTTP/1.1 200 OK", r#"{"data":[{"id":"Qwen3-0.6B"}]}"#) + }; + let _ = write!( + stream, + "{status_line}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), + body + ); + } + }); + + let (root, paths) = test_paths("managed-readiness-engine-died"); + paths.ensure()?; + let mut record = ManagedServiceRecord::new( + &paths, + "svc-died", + "vllm", + "Qwen3-0.6B", + "Qwen3-0.6B", + "127.0.0.1", + port, + "managed", + std::process::id(), + None, + None, + None, + ); + record.status = "running".to_owned(); + fs::create_dir_all(record.engine_state_path.parent().expect("state dir"))?; + fs::write(&record.engine_state_path, r#"{"status":"failed"}"#)?; + + let started = std::time::Instant::now(); + await_managed_readiness(&mut record, None, &mut |_elapsed| {}); + + assert_eq!(record.status, "failed"); + assert!( + started.elapsed() < Duration::from_secs(30), + "a dead engine must end the wait, not run out the budget: {:?}", + started.elapsed() + ); + drop(server); + fs::remove_dir_all(root).ok(); + Ok(()) + } + /// Regression test for the engine child-stdin write under the process-wide /// `SIG_DFL` that `main` installs via [`reset_sigpipe`]. If an engine child /// exits before reading its stdin, the parent's write to that pipe must diff --git a/apps/rocm/src/serve_summary.rs b/apps/rocm/src/serve_summary.rs index c4bfa0213..c73d48490 100644 --- a/apps/rocm/src/serve_summary.rs +++ b/apps/rocm/src/serve_summary.rs @@ -53,7 +53,9 @@ pub(crate) struct DeploymentSummary { /// Full chat-completions endpoint, e.g. `http://127.0.0.1:1337/v1/chat/completions`. pub chat_endpoint: String, pub service_id: String, - /// `"ready"`, `"starting"`, or an existing service's status. + /// `"ready"`, `"running"` (lists the model, cannot serve yet), `"starting"`, + /// `"failed"` (the engine exited during startup), or an existing service's + /// status. pub status: String, /// True when an equivalent server was already running and nothing was spawned. pub already_running: bool, @@ -94,14 +96,17 @@ pub(crate) fn format_tps(tps: Option) -> String { /// printed so it is unit-testable and identical across engines). pub(crate) fn render_summary(summary: &DeploymentSummary) -> String { // The server answered its health check within the startup window. A launch - // that timed out lands here with `status == "starting"`; the heading and a - // note must make that visibly different from a healthy deployment so the - // summary is never mistaken for success. + // that timed out lands here with `status == "starting"`, and one whose engine + // exited with `status == "failed"`; the heading and a note must make both + // visibly different from a healthy deployment so the summary is never + // mistaken for success, and a failed launch is never read as "wait longer". let ready = summary.status == "ready"; let heading = if summary.already_running { "Deployment summary (already running)" } else if ready { "Deployment summary" + } else if summary.status == "failed" { + "Deployment summary (failed)" } else { "Deployment summary (not ready yet)" }; @@ -160,11 +165,15 @@ pub(crate) fn render_summary(summary: &DeploymentSummary) -> String { ); } if !ready && !summary.already_running { - // Two different ways to miss "ready", and conflating them misleads: a + // Three different ways to miss "ready", and conflating them misleads: a + // `failed` service is over and needs its engine error read, while a // `running` service answered and advertised the model but could not serve // a request yet, which is also why the metrics rows are empty — the smoke // test only runs against a service that can actually serve. - let explanation = if summary.status == "running" { + let explanation = if summary.status == "failed" { + "the server exited during startup, so the launch is over; the engine error is \ + at the end of its log" + } else if summary.status == "running" { "the server is up and lists the model, but it could not serve a request yet, \ so no smoke test was run; it is most likely still loading" } else { @@ -279,6 +288,25 @@ mod tests { assert_eq!(format_ttft(None), "n/a"); } + /// A launch whose engine died must not be described as "may still be + /// loading" or "not ready yet" — the wait ended because the server is gone, + /// not because it is slow, and the next action is reading the engine error. + #[test] + fn failed_launch_points_at_the_engine_error() { + let summary = DeploymentSummary { + status: "failed".to_owned(), + ..base_summary() + }; + let rendered = render_summary(&summary); + assert!( + rendered.starts_with("Deployment summary (failed)"), + "{rendered}" + ); + assert!(rendered.contains("exited during startup"), "{rendered}"); + assert!(!rendered.contains("not ready yet"), "{rendered}"); + assert!(!rendered.contains("may still be loading"), "{rendered}"); + } + #[test] fn format_tps_renders_rate_or_na() { assert_eq!(format_tps(Some(42.15)), "42.1 tok/s"); diff --git a/engines/vllm/src/lib.rs b/engines/vllm/src/lib.rs index a806d3101..fc3981bc7 100644 --- a/engines/vllm/src/lib.rs +++ b/engines/vllm/src/lib.rs @@ -22,6 +22,8 @@ mod process; mod runtime; mod state; +pub use process::ready_timeout; + pub(crate) const ENGINE_NAME: &str = "vllm"; pub(crate) const DEFAULT_HOST: &str = "127.0.0.1"; diff --git a/engines/vllm/src/process.rs b/engines/vllm/src/process.rs index fef713a6c..f62ac9629 100644 --- a/engines/vllm/src/process.rs +++ b/engines/vllm/src/process.rs @@ -152,7 +152,7 @@ pub(crate) fn serve_http(mut request: ServeHttpRequest) -> Result<()> { &request.host, request.port, &request.model_ref, - vllm_ready_timeout(), + ready_timeout(), request.log_path.as_deref(), ) { // Terminate the whole vLLM process tree so the EngineCore worker (which @@ -582,8 +582,9 @@ fn vllm_enforce_eager_enabled() -> bool { /// Defaults to [`DEFAULT_VLLM_READY_TIMEOUT`]. A valid-but-slow cold start /// (large weight download, first-decode compile) can exceed the default, so the /// timeout is configurable via `ROCM_CLI_VLLM_READY_TIMEOUT_SECS` (a positive -/// integer number of seconds). -fn vllm_ready_timeout() -> Duration { +/// integer number of seconds). Public so the CLI's managed-launch path waits on +/// the same budget the engine itself honors instead of a shorter one of its own. +pub fn ready_timeout() -> Duration { resolve_vllm_ready_timeout(std::env::var("ROCM_CLI_VLLM_READY_TIMEOUT_SECS").ok()) } diff --git a/tests/e2e-cucumber/features/model_serving.feature b/tests/e2e-cucumber/features/model_serving.feature index 0cecc8c25..7f9bd089d 100644 --- a/tests/e2e-cucumber/features/model_serving.feature +++ b/tests/e2e-cucumber/features/model_serving.feature @@ -268,3 +268,16 @@ Feature: Model serving Given a local server attempt has failed When the user lists running services Then the list reports the attempt and how to look at it + + # The vLLM adapter and GPU preflight are real, but the server executable is a + # short-lived fixture. Deliberately no `@requires-engine:vllm`: that gate asks + # whether real vLLM can start on this GPU family, while this scenario overrides + # the runtime and would otherwise skip valid adapter coverage on Strix. + @id:serve-vllm-engine-exit-reported-failed @requires-gpu @requires-os:linux + Scenario: serve-23 - A vLLM server that exits during startup is reported as failed promptly + Given a managed runtime is active + And a vLLM server will exit during startup + When the user launches it as a managed model server + Then the managed launch is reported as failed + And the failed launch returns promptly + And the service state reports the launch as failed diff --git a/tests/e2e-cucumber/tests/e2e/serving_steps.rs b/tests/e2e-cucumber/tests/e2e/serving_steps.rs index a36489808..8528b0bc0 100644 --- a/tests/e2e-cucumber/tests/e2e/serving_steps.rs +++ b/tests/e2e-cucumber/tests/e2e/serving_steps.rs @@ -167,6 +167,9 @@ async fn serve_and_wait(world: &mut E2eWorld, args: &[&str], model: &str, ready_ /// the job times out. const SERVE_PORT: u16 = 11435; +const FAILING_VLLM_MODEL: &str = "e2e/vllm-startup-exit"; +const FAKE_VLLM_ERROR: &str = "e2e fake vLLM: startup failed"; + /// The port the CLI's built-in local assistant (lemonade Qwen3-4B) listens on. /// The CLI auto-starts this assistant independently of any scenario; on Instinct /// it falls back to a Vulkan llama-server that pins a GPU core (EAI-7052), @@ -706,6 +709,38 @@ async fn setup_failed_local_server(world: &mut E2eWorld) { }); } +#[given("a vLLM server will exit during startup")] +async fn setup_vllm_startup_exit(world: &mut E2eWorld) { + let root = world.isolated_root.as_ref().expect("no isolated root"); + let script = root.path().join("fake-vllm-exit.sh"); + std::fs::write( + &script, + format!( + "#!/bin/sh\nsleep 3\nprintf '%s\\n' '{FAKE_VLLM_ERROR}' >&2\nexit 1\n" + ), + ) + .unwrap_or_else(|e| panic!("failed to write {}: {e}", script.display())); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt as _; + std::fs::set_permissions(&script, std::fs::Permissions::from_mode(0o755)) + .unwrap_or_else(|e| panic!("failed to mark {} executable: {e}", script.display())); + } + assert!( + script.is_absolute(), + "ROCM_CLI_VLLM_COMMAND requires an absolute path: {}", + script.display() + ); + world + .command_env + .push(("ROCM_CLI_VLLM_COMMAND", script.into_os_string())); + // Keep the engine budget longer than the assertion below: only the terminal + // state written after the fake exits should end this launch promptly. + world + .command_env + .push(("ROCM_CLI_VLLM_READY_TIMEOUT_SECS", "60".into())); +} + #[given("the served model has been detected")] async fn setup_model_detected(world: &mut E2eWorld) { let (stdout, _, _) = crate::run_rocm(world, &["services", "list"]); @@ -1004,6 +1039,37 @@ async fn user_serves_with_rocr_hiding_hip_named_gpus(world: &mut E2eWorld) { world.cli_rc = Some(rc); } +#[when("the user launches it as a managed model server")] +async fn user_launches_failing_vllm_server(world: &mut E2eWorld) { + let listener = std::net::TcpListener::bind(("127.0.0.1", 0)) + .expect("failed to reserve an OS-assigned serve port"); + let port = listener + .local_addr() + .expect("reserved listener has no local address") + .port() + .to_string(); + drop(listener); + + let started = Instant::now(); + let (stdout, stderr, rc) = crate::run_rocm_with_scenario_env( + world, + &[ + "serve", + FAILING_VLLM_MODEL, + "--engine", + "vllm", + "--managed", + "--port", + &port, + ], + ); + world.cli_elapsed = Some(started.elapsed()); + world.model_name = Some(FAILING_VLLM_MODEL.to_owned()); + world.cli_output = Some(stdout); + world.cli_stderr = Some(stderr); + world.cli_rc = Some(rc); +} + #[then("the user is told to allow public binding first")] async fn assert_public_bind_message(world: &mut E2eWorld) { let output = serve_output(world); @@ -1103,6 +1169,48 @@ fn serve_output(world: &E2eWorld) -> String { ) } +#[then("the managed launch is reported as failed")] +async fn assert_managed_launch_failed(world: &mut E2eWorld) { + let output = serve_output(world); + assert!( + output + .lines() + .any(|line| line.trim() == "readiness: failed"), + "expected a failed readiness verdict, got:\n{output}" + ); + assert!( + !output.contains("readiness: starting") + && !output.contains("not ready") + && !output.contains("may still be loading"), + "a terminal launch must not be described as still loading:\n{output}" + ); +} + +#[then("the failed launch returns promptly")] +async fn assert_failed_launch_returns_promptly(world: &mut E2eWorld) { + let elapsed = world.cli_elapsed.expect("managed launch was not timed"); + assert!( + elapsed < Duration::from_secs(30), + "failed launch took {:.1}s, so it did not stop on the terminal engine state:\n{}", + elapsed.as_secs_f64(), + serve_output(world) + ); +} + +#[then("the service state reports the launch as failed")] +async fn assert_failed_launch_service_state(world: &mut E2eWorld) { + let listed = crate::run_rocm_ok(world, &["services", "list", "--all"]); + let expected_model = format!(" model: {FAILING_VLLM_MODEL}"); + let failed_record = listed.split("\n- ").any(|record| { + record.lines().any(|line| line == expected_model) + && record.lines().any(|line| line == " status: failed") + }); + assert!( + failed_record, + "expected the failed launch and service record to agree:\n{listed}" + ); +} + #[then("serving is refused before any engine starts")] async fn assert_serve_refused(world: &mut E2eWorld) { // A non-zero exit is the observable "refused". The GPU-required pre-flight