Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
200 changes: 163 additions & 37 deletions apps/rocm/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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),
Expand Down Expand Up @@ -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())?;
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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,
Expand All @@ -21937,7 +21974,7 @@ fn wait_for_service_http_ready(
canonical_model_id,
endpoint_api_key,
timeout,
&mut |_elapsed| {},
&mut |_elapsed| true,
)
}

Expand All @@ -21952,21 +21989,26 @@ 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,
port: u16,
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| {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
40 changes: 34 additions & 6 deletions apps/rocm/src/serve_summary.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -94,14 +96,17 @@ pub(crate) fn format_tps(tps: Option<f64>) -> 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)"
};
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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");
Expand Down
2 changes: 2 additions & 0 deletions engines/vllm/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
7 changes: 4 additions & 3 deletions engines/vllm/src/process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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())
}

Expand Down
Loading