diff --git a/apps/rocm/src/main.rs b/apps/rocm/src/main.rs index e600faf20..4c79b2a28 100644 --- a/apps/rocm/src/main.rs +++ b/apps/rocm/src/main.rs @@ -66,7 +66,6 @@ use std::io::{self, BufRead, Read, Write}; use std::net::{TcpStream, ToSocketAddrs}; use std::path::{Path, PathBuf}; use std::process::ExitCode; -#[cfg(not(windows))] use std::process::ExitStatus; use std::process::{Command as ProcessCommand, Stdio}; use std::sync::{Mutex, OnceLock}; @@ -6064,8 +6063,12 @@ fn serve(args: ServeArgs) -> Result<()> { // test reaches Lemonade's backend boundary without real GPU hardware. let scripted_backend_failure = cfg!(feature = "e2e-test-hooks") && std::env::var_os("ROCM_E2E_LEMONADE_BACKEND_INSTALL_FAILURE").is_some(); + // The scripted startup-death scenario bypasses the same host precondition, for + // the same reason: it asserts what happens *after* a spawn, so it must reach + // one without real GPU hardware. if !cpu_only && !scripted_backend_failure + && !scripted_managed_engine_startup_failure() && let Some(usable) = visible_gpu_indices.as_deref() && usable.is_empty() { @@ -6124,8 +6127,12 @@ fn serve(args: ServeArgs) -> Result<()> { device_policy_name(&device_policy) ); } + // Preparation is skipped for the scripted startup-death scenario too: it + // asserts the launch, not the install, and a real runtime download is exactly + // what the ungated lane cannot do. if !matches!(device_policy, DevicePolicy::CpuOnly) && engine_manages_own_runtime(&selected_engine) + && !scripted_managed_engine_startup_failure() { ensure_self_managed_engine_ready(&paths, &mut config, &selected_engine)?; } @@ -6591,12 +6598,166 @@ fn attach_background_stdio(command: &mut ProcessCommand, log_path: Option<&Path> Ok(()) } -#[cfg(not(windows))] +/// How long a freshly spawned managed engine is watched before it is treated as +/// started. Long enough to catch an engine that dies on the way up (bad runtime, +/// missing model, port already taken), short enough to be invisible next to the +/// readiness wait that follows. Both platforms use the same budget so a launch +/// that fails at startup fails the same way on each. +const MANAGED_ENGINE_STARTUP_SETTLE: Duration = Duration::from_millis(200); + +/// E2E-only switch that makes the managed engine die during startup, so the +/// black-box suite can assert what a user sees when it does. +/// +/// Read in the *parent*, which is what spawns the engine, so nothing test-only +/// leaks into the child's own code path. The child is the real `rocm` binary, +/// really spawned and really dead — only its arguments are swapped for ones the +/// CLI rejects outright, so the launch takes exactly the path a broken engine +/// takes. +#[cfg(feature = "e2e-test-hooks")] +const MANAGED_ENGINE_STARTUP_FAILURE_TEST_ENV: &str = "ROCM_E2E_MANAGED_ENGINE_STARTUP_FAILURE"; + +/// Whether the E2E startup-failure switch is armed. +#[cfg(feature = "e2e-test-hooks")] +fn scripted_managed_engine_startup_failure() -> bool { + std::env::var_os(MANAGED_ENGINE_STARTUP_FAILURE_TEST_ENV).is_some() +} + +/// Release builds do not enable `e2e-test-hooks`, so the switch does not exist +/// there and no environment variable can arm it. +#[cfg(not(feature = "e2e-test-hooks"))] +const fn scripted_managed_engine_startup_failure() -> bool { + false +} + +/// Arguments that make the spawned `rocm` exit immediately with a non-zero +/// status. Used only under [`scripted_managed_engine_startup_failure`]. +/// +/// An unrecognised *flag*, deliberately: clap rejects it with exit code 2 before +/// any work. An unrecognised subcommand would not do — the CLI falls back to +/// natural-language request parsing for those and exits 0, which is a different +/// (and much slower) thing than an engine dying. +#[cfg(feature = "e2e-test-hooks")] +fn managed_engine_startup_failure_args() -> Vec { + vec!["--e2e-managed-engine-startup-failure".to_owned()] +} + +/// Retire the service record of a launch that did not get an engine running. +/// +/// Callers must reach this on *every* path that abandons a record, because a +/// record left in a non-terminal state keeps claiming its engine + model: +/// `managed_service_is_live` counts `"starting"` as alive, so the idempotency +/// guard would report the corpse as `AlreadyRunning` and refuse every later +/// `rocm serve` for the same pair. `"failed"` is outside the live set, so the +/// next launch proceeds. +/// +/// On the launch path the record is persisted *before* the spawn, to claim the +/// GPU for concurrent auto-selection, so it still carries the constructor's +/// `"starting"` status and a `supervisor_pid` of 0 — and pid 0 is skipped by +/// `recorded_service_pids`, so the liveness refresh cannot demote it either and +/// the wedge is permanent. On the restart path the record is pre-existing and +/// its pid is merely stale, which the liveness refresh would eventually clean +/// up; retiring it here is what makes the outcome immediate and identical on +/// both paths. +fn mark_managed_launch_failed(record: &mut ManagedServiceRecord) -> Result<()> { + record.status = "failed".to_owned(); + record.write() +} + +/// Retire `record` when `outcome` failed, so a bail-out before the child exists +/// stops it claiming its engine + model. +/// +/// Every fallible step between `record.write()` and the spawn has to go through +/// here. Each one that does not is another way to strand the pid-0 `"starting"` +/// corpse — the failures differ, the wedge is identical. +/// +/// **The span ends at the spawn, deliberately.** Once a child exists, retiring the +/// record is no longer the safe default: it frees this engine + model for the next +/// `rocm serve`, which would start a rival engine on the same port while the child +/// this launch spawned may still be alive. Past that point the caller has to decide +/// per failure whether the child is known dead — which is what +/// [`fail_managed_launch_if_engine_died`] does, and why a failed *liveness query* +/// is not routed here. +/// +/// Takes an already-evaluated `Result` rather than a closure so it can wrap that +/// prefix without holding a second borrow of `record` across the call. +/// +/// The retirement error is deliberately dropped: what the user needs to see is +/// the failure that aborted the launch, not a bookkeeping failure behind it. +fn retire_record_on_error(record: &mut ManagedServiceRecord, outcome: Result) -> Result { + outcome.inspect_err(|_| { + let _ = mark_managed_launch_failed(record); + }) +} + +/// Fail the launch if the freshly spawned engine had already exited, retiring the +/// record on the way out so the dead launch stops claiming its engine + model. +/// +/// Both platforms funnel their startup check through here — Unix from +/// `Child::try_wait`, Windows from the watched detached spawn — so a child that +/// is *observed* to have died produces the same error, carrying the same tail of +/// the child's own log, whichever platform the user is on. +/// +/// A liveness query that *fails* is a different case, and the platforms do not +/// agree on it. Windows degrades to "assume it is alive": `observe_early_exit` +/// reports `None` when the wait times out or the exit code cannot be read, so the +/// launch proceeds. Unix propagates the `try_wait` error and fails the launch. +/// Neither retires the record, because at that point the child may be running and +/// releasing the engine + model would invite a rival engine onto the same port. +fn fail_managed_launch_if_engine_died( + record: &mut ManagedServiceRecord, + startup_exit: Option, +) -> Result<()> { + let Some(status) = startup_exit else { + return Ok(()); + }; + // Dropped, not propagated, for the same reason as `retire_record_on_error`: + // the engine's own exit status and log tail are what diagnose this launch, and + // a failed record write must not displace them. + let _ = mark_managed_launch_failed(record); + bail!( + "{}", + managed_engine_startup_failure_detail(status, &record.log_path) + ); +} + +/// Read a watched detached spawn's early exit as an `ExitStatus`. +/// +/// `GetExitCodeProcess` yields a bare `u32`, which is exactly what +/// `ExitStatusExt::from_raw` takes on Windows; the conversion is isolated here so +/// the shared startup check above stays platform-neutral. +#[cfg(windows)] +fn managed_startup_exit(spawn: &rocm_core::DetachedSpawn) -> Option { + use std::os::windows::process::ExitStatusExt as _; + + spawn.early_exit_code.map(ExitStatus::from_raw) +} + +/// Render a managed engine's immediate exit for the user: the exit status plus a +/// tail of the child's own service log, which is the only record of why it died — +/// the child is detached, so nothing else reaches the terminal. fn managed_engine_startup_failure_detail(status: ExitStatus, log_path: &Path) -> String { + managed_engine_startup_failure_detail_polling( + status, + log_path, + MANAGED_ENGINE_STARTUP_LOG_POLL_INTERVAL, + ) +} + +/// How long to wait between re-reads of an empty service log before giving up on +/// a tail. A child that died at startup may still be flushing its last lines. +const MANAGED_ENGINE_STARTUP_LOG_POLL_INTERVAL: Duration = Duration::from_millis(120); + +/// [`managed_engine_startup_failure_detail`] with the re-read interval injected, +/// so a test of the no-log branch need not sit through the real waits. +fn managed_engine_startup_failure_detail_polling( + status: ExitStatus, + log_path: &Path, + poll_interval: Duration, +) -> String { let mut recent_lines = read_optional_tail_lines(log_path, 80, "service log"); if recent_lines.is_empty() { for _ in 0..5 { - thread::sleep(Duration::from_millis(120)); + thread::sleep(poll_interval); recent_lines = read_optional_tail_lines(log_path, 80, "service log"); if !recent_lines.is_empty() { break; @@ -6780,15 +6941,25 @@ fn spawn_managed_engine_child( record.requires_api_key = require_api_key; record.write()?; - if let Some(parent) = record.engine_state_path.parent() { - fs::create_dir_all(parent) - .with_context(|| format!("failed to create {}", parent.display()))?; - } - fs::File::create(&record.log_path) - .with_context(|| format!("failed to create {}", record.log_path.display()))?; - let current_exe = managed_service_launcher_path() - .context("failed to resolve current rocm executable path")?; - let serve_args = builtin_engine_serve_http_args( + // From here up to and including the spawn, every fallible step goes through + // `retire_record_on_error`: the record is on disk claiming this engine + model, + // so any bail-out that skips the retirement wedges the service just as an + // unretired startup death would. The span stops at the spawn — see that + // helper's doc for why retiring is no longer the safe default once a child + // exists. + let engine_state_parent = record.engine_state_path.parent().map(Path::to_path_buf); + if let Some(parent) = engine_state_parent { + let created = fs::create_dir_all(&parent) + .with_context(|| format!("failed to create {}", parent.display())); + retire_record_on_error(&mut record, created)?; + } + let log_created = fs::File::create(&record.log_path) + .with_context(|| format!("failed to create {}", record.log_path.display())); + retire_record_on_error(&mut record, log_created)?; + let launcher = + managed_service_launcher_path().context("failed to resolve current rocm executable path"); + let current_exe = retire_record_on_error(&mut record, launcher)?; + let built_args = builtin_engine_serve_http_args( engine, service_id, &resolve.canonical_model_id, @@ -6801,8 +6972,18 @@ fn spawn_managed_engine_child( engine_recipe, &record.engine_state_path, Some(&record.log_path), - )?; - let engine_envs_root = env_root_for_service(paths, engine, runtime_id, env_id)?; + ); + let serve_args = retire_record_on_error(&mut record, built_args)?; + // E2E-only: swap in arguments the child rejects immediately, so the launch + // below observes a real engine that really died during startup. + #[cfg(feature = "e2e-test-hooks")] + let serve_args = if scripted_managed_engine_startup_failure() { + managed_engine_startup_failure_args() + } else { + serve_args + }; + let envs_root = env_root_for_service(paths, engine, runtime_id, env_id); + let engine_envs_root = retire_record_on_error(&mut record, envs_root)?; // Hand the child the *path* to the endpoint key file (public bind only) via the // environment. A path — not the secret value — is what the detached-spawn // primitives accept as an env override, and it keeps the key off both the argv @@ -6826,11 +7007,15 @@ fn spawn_managed_engine_child( // // The field is threaded, never derived from the key file. Deriving it marked // every public bind as having demanded auth — see the assignment above. - ensure_public_service_has_endpoint_key( + // + // Evaluated before the guard, not inside the call, so reading the record here + // does not overlap the `&mut` borrow the retirement needs. + let public_key_guard = ensure_public_service_has_endpoint_key( host, endpoint_key_file.is_some(), record.requires_api_key, - )?; + ); + retire_record_on_error(&mut record, public_key_guard)?; #[cfg(windows)] let child_pid = { let env_values = app_path_env_var_values(paths, engine_envs_root.as_deref()); @@ -6838,14 +7023,30 @@ fn spawn_managed_engine_child( if let Some(key_file) = endpoint_key_file.as_deref() { env_refs.push((rocm_engine_protocol::ENDPOINT_API_KEY_FILE_ENV, key_file)); } - rocm_core::spawn_detached_no_inherit(¤t_exe, &serve_args, &env_refs) - .context("failed to launch managed engine process")? + // The detached spawn hands back a bare PID, not a `Child`, so there is no + // `try_wait()` to lean on. Watch the process while its handle is still + // open instead — the Windows helper does that internally, which is what + // makes the check free of the PID-reuse race a later `OpenProcess` would + // have. + let spawned = rocm_core::spawn_detached_no_inherit_watching_startup( + ¤t_exe, + &serve_args, + &env_refs, + MANAGED_ENGINE_STARTUP_SETTLE, + ) + .context("failed to launch managed engine process"); + let spawn = retire_record_on_error(&mut record, spawned)?; + fail_managed_launch_if_engine_died(&mut record, managed_startup_exit(&spawn))?; + spawn.pid }; #[cfg(not(windows))] let child_pid = { let mut command = managed_service_process_command(¤t_exe, &serve_args); command.stdin(Stdio::null()); - attach_background_stdio(&mut command, Some(&record.log_path))?; + // Evaluated first so the borrow of `record.log_path` ends before the + // `&mut record` the retirement needs. + let attached = attach_background_stdio(&mut command, Some(&record.log_path)); + retire_record_on_error(&mut record, attached)?; detach_background_command(&mut command); apply_app_path_env(&mut command, paths); if let Some(engine_envs_root) = engine_envs_root.as_deref() { @@ -6854,20 +7055,20 @@ fn spawn_managed_engine_child( if let Some(key_file) = endpoint_key_file.as_deref() { command.env(rocm_engine_protocol::ENDPOINT_API_KEY_FILE_ENV, key_file); } - let mut child = command + let spawned = command .spawn() - .context("failed to launch managed engine process")?; + .context("failed to launch managed engine process"); + let mut child = retire_record_on_error(&mut record, spawned)?; let child_pid = child.id(); - thread::sleep(Duration::from_millis(200)); - if let Some(status) = child + thread::sleep(MANAGED_ENGINE_STARTUP_SETTLE); + // NOT routed through `retire_record_on_error`: the child is already + // spawned, and a failed query means its state is unknown, not that it + // died. Retiring here would release the engine + model to the next serve + // while this child may still be listening on the port. + let startup_exit = child .try_wait() - .context("failed to check managed engine startup state")? - { - bail!( - "{}", - managed_engine_startup_failure_detail(status, &record.log_path) - ); - } + .context("failed to check managed engine startup state")?; + fail_managed_launch_if_engine_died(&mut record, startup_exit)?; child_pid }; record.supervisor_pid = child_pid; @@ -6925,11 +7126,14 @@ fn start_managed_service( // visible to any concurrent auto-selection. Release the launch lock before // the readiness wait below, which can block for many seconds — holding it // that long would needlessly serialize unrelated serves. + // + // What the lock *does* still cover is the startup check above: both platforms + // spend `MANAGED_ENGINE_STARTUP_SETTLE` inside `spawn_managed_engine_child`, + // so a concurrent serve waits that much longer for the lock. That is + // deliberate — the check has to see the child before the record is promoted to + // `"running"` — and bounded, unlike the readiness wait. drop(launch_lock); - #[cfg(windows)] - thread::sleep(Duration::from_millis(200)); - let readiness = wait_for_service_http_ready_with_progress( engine, host, @@ -18124,9 +18328,19 @@ fn restart_internal_managed_service( } // The stop above may have recorded an unconfirmed-stop marker; a successful // restart supersedes it. Leaving it set would let the next liveness refresh - // delete the key of the service we are bringing back up. Reaches disk with - // the record writes below; if the restart bails before one of those, the - // marker stays set on disk — correct, since then the stop is what stands. + // delete the key of the service we are bringing back up. + // + // A restart that fails at or after the spawn persists the cleared marker too, + // because those bails retire the record and retiring writes it. So the key + // outlives a restart that got an engine up and lost it, rather than being + // reclaimed by the next refresh. That is the safe direction — the record is + // left `"failed"`, which is not live, so nothing reuses the service while the + // key waits for a retry. + // + // An earlier bail — any of the `?`s between here and the spawn — leaves the + // record unwritten, so the marker stays set on disk and the next refresh + // reclaims the key. That is correct too: nothing was restarted, so the stop + // stands. record.stop_requested_unix_ms = None; let policy = parse_device_policy(record.device_policy.as_deref())?; fs::OpenOptions::new() @@ -18173,14 +18387,28 @@ fn restart_internal_managed_service( if let Some(key_file) = endpoint_key_file.as_deref() { env_refs.push((rocm_engine_protocol::ENDPOINT_API_KEY_FILE_ENV, key_file)); } - rocm_core::spawn_detached_no_inherit(¤t_exe, &serve_args, &env_refs) - .context("failed to restart managed engine process")? + // Same reasoning as the launch path: without a `Child` to `try_wait()` on, + // the only race-free liveness check is the one the Windows helper performs + // while the process handle is still open. + let spawned = rocm_core::spawn_detached_no_inherit_watching_startup( + ¤t_exe, + &serve_args, + &env_refs, + MANAGED_ENGINE_STARTUP_SETTLE, + ) + .context("failed to restart managed engine process"); + let spawn = retire_record_on_error(&mut record, spawned)?; + fail_managed_launch_if_engine_died(&mut record, managed_startup_exit(&spawn))?; + spawn.pid }; #[cfg(not(windows))] let child_pid = { let mut command = managed_service_process_command(¤t_exe, &serve_args); command.stdin(Stdio::null()); - attach_background_stdio(&mut command, Some(&record.log_path))?; + // Evaluated first so the borrow of `record.log_path` ends before the + // `&mut record` the retirement needs. + let attached = attach_background_stdio(&mut command, Some(&record.log_path)); + retire_record_on_error(&mut record, attached)?; detach_background_command(&mut command); apply_app_path_env(&mut command, paths); if let Some(engine_envs_root) = engine_envs_root.as_deref() { @@ -18189,25 +18417,19 @@ fn restart_internal_managed_service( if let Some(key_file) = endpoint_key_file.as_deref() { command.env(rocm_engine_protocol::ENDPOINT_API_KEY_FILE_ENV, key_file); } - let mut child = command + let spawned = command .spawn() - .context("failed to restart managed engine process")?; - thread::sleep(Duration::from_millis(200)); - if let Some(status) = child + .context("failed to restart managed engine process"); + let mut child = retire_record_on_error(&mut record, spawned)?; + thread::sleep(MANAGED_ENGINE_STARTUP_SETTLE); + // Not retired on a failed query, for the reason given at the launch site: + // an unknown child state is not a dead child. + let startup_exit = child .try_wait() - .context("failed to check restarted engine startup state")? - { - record.status = "failed".to_owned(); - record.write()?; - bail!( - "{}", - managed_engine_startup_failure_detail(status, &record.log_path) - ); - } + .context("failed to check restarted engine startup state")?; + fail_managed_launch_if_engine_died(&mut record, startup_exit)?; child.id() }; - #[cfg(windows)] - thread::sleep(Duration::from_millis(200)); record.status = "running".to_owned(); record.supervisor_pid = child_pid; record.engine_pid = Some(child_pid); @@ -29148,6 +29370,253 @@ install therock"; assert!(found.is_none()); } + /// An engine that dies during startup must reach the user with the child's + /// own log tail, not just an exit status — the log is the only place the + /// reason is recorded, since the child is detached and writes nowhere else. + /// Both platforms feed this from the same startup check, so both produce this + /// message. + #[test] + fn managed_engine_startup_failure_detail_carries_the_child_log_tail() { + let (root, _paths) = test_paths("managed-startup-failure-detail"); + fs::create_dir_all(&root).expect("test dir"); + let log_path = root.join("service.log"); + fs::write(&log_path, "loading runtime\nfatal: no usable device\n").expect("write log"); + + let detail = managed_engine_startup_failure_detail(exit_status_from_code(3), &log_path); + + let _ = fs::remove_dir_all(&root); + assert!( + detail.contains("managed engine exited immediately"), + "unexpected detail: {detail}" + ); + assert!( + detail.contains("fatal: no usable device"), + "log tail missing from detail: {detail}" + ); + assert!( + detail.contains(&log_path.display().to_string()), + "log path missing from detail: {detail}" + ); + } + + /// The same failure with no log written yet still has to name the log path so + /// the user knows where to look once the child flushes. + #[test] + fn managed_engine_startup_failure_detail_without_a_log_still_points_at_it() { + let (root, _paths) = test_paths("managed-startup-failure-no-log"); + fs::create_dir_all(&root).expect("test dir"); + let log_path = root.join("absent.log"); + + // Zero interval: the re-read loop still runs all of its iterations, so the + // branch it guards is exercised, without ~600 ms of real sleeps. + let detail = managed_engine_startup_failure_detail_polling( + exit_status_from_code(1), + &log_path, + Duration::ZERO, + ); + + let _ = fs::remove_dir_all(&root); + assert!( + detail.contains(&log_path.display().to_string()), + "log path missing from detail: {detail}" + ); + assert!( + !detail.contains("recent startup log output"), + "empty log should not advertise a tail: {detail}" + ); + } + + /// Build an `ExitStatus` from a plain exit code the way the Windows startup + /// check does, i.e. straight from `GetExitCodeProcess`. + #[cfg(windows)] + fn exit_status_from_code(code: i32) -> ExitStatus { + use std::os::windows::process::ExitStatusExt as _; + ExitStatus::from_raw(code.cast_unsigned()) + } + + /// Build an `ExitStatus` from a plain exit code the way the Unix startup + /// check does, i.e. from the raw wait status `try_wait` reports. + #[cfg(not(windows))] + fn exit_status_from_code(code: i32) -> ExitStatus { + use std::os::unix::process::ExitStatusExt as _; + // Unix wait status: exit code in the high byte, no signal. + ExitStatus::from_raw(code << 8) + } + + /// A launch that dies at startup must retire its own service record. + /// + /// The record is written before the spawn, so it still says `"starting"` with + /// a `supervisor_pid` of 0. Pid 0 is filtered out of the liveness refresh, so + /// nothing ever demotes such a record, while `"starting"` counts as live — + /// leaving it would make the idempotency guard report the corpse as already + /// running and refuse every later `rocm serve` for the same engine + model. + #[test] + fn a_failed_launch_record_stops_blocking_the_next_serve() -> Result<()> { + let (root, paths) = test_paths("managed-failed-launch-unblocks"); + paths.ensure()?; + let mut record = ManagedServiceRecord::new( + &paths, + "lemonade-qwen-3000", + "lemonade", + "qwen", + "qwen-canonical", + "127.0.0.1", + 11512, + "managed", + // The pre-spawn state: no child pid recorded yet. + 0, + None, + None, + None, + ); + record.write()?; + + // Precondition: abandoned as written, the record blocks the next launch. + assert!( + existing_live_managed_service(&paths, "lemonade", "qwen-canonical").is_some(), + "a pid-0 `starting` record should look live — that is the trap being closed" + ); + + mark_managed_launch_failed(&mut record)?; + + let still_blocking = existing_live_managed_service(&paths, "lemonade", "qwen-canonical"); + let _ = fs::remove_dir_all(root); + assert!( + still_blocking.is_none(), + "a failed launch must not keep claiming the engine+model" + ); + Ok(()) + } + + /// Build the pre-spawn record every managed launch persists to claim its + /// engine + model: `"starting"`, with no child pid recorded yet. + fn pre_spawn_record(paths: &AppPaths, port: u16) -> ManagedServiceRecord { + ManagedServiceRecord::new( + paths, + "lemonade-qwen-3000", + "lemonade", + "qwen", + "qwen-canonical", + "127.0.0.1", + port, + "managed", + 0, + None, + None, + None, + ) + } + + /// The whole user-visible fix in one assertion: an engine that exited during + /// startup must fail the launch, name its log, *and* let go of the engine + + /// model so the next `rocm serve` can try again. + /// + /// Both platforms route their startup check through this function — Unix from + /// `Child::try_wait`, Windows from the watched detached spawn — so this runs + /// everywhere and covers the wiring between detecting the death and acting on + /// it, which the two halves tested in isolation do not. + #[test] + fn a_dead_engine_fails_the_launch_and_frees_the_engine_for_the_next_serve() -> Result<()> { + let (root, paths) = test_paths("managed-dead-engine-fails-launch"); + paths.ensure()?; + let mut record = pre_spawn_record(&paths, 11513); + record.write()?; + fs::write(&record.log_path, "fatal: no usable device\n")?; + + let outcome = + fail_managed_launch_if_engine_died(&mut record, Some(exit_status_from_code(3))); + + let still_blocking = existing_live_managed_service(&paths, "lemonade", "qwen-canonical"); + let _ = fs::remove_dir_all(root); + let error = outcome + .expect_err("an engine that already exited must fail the launch") + .to_string(); + assert!( + error.contains("fatal: no usable device"), + "the child's log tail is the only account of why it died: {error}" + ); + assert!( + still_blocking.is_none(), + "a launch that failed at startup must not keep claiming the engine+model" + ); + Ok(()) + } + + /// The converse, and the reason the check can be trusted: a child that is + /// still running must leave the launch untouched. If this ever failed, every + /// healthy serve would be rejected and its record retired. + #[test] + fn a_live_engine_leaves_the_launch_running() -> Result<()> { + let (root, paths) = test_paths("managed-live-engine-proceeds"); + paths.ensure()?; + let mut record = pre_spawn_record(&paths, 11514); + record.write()?; + + let outcome = fail_managed_launch_if_engine_died(&mut record, None); + + let _ = fs::remove_dir_all(root); + assert!(outcome.is_ok(), "a running engine must not fail the launch"); + assert_eq!( + record.status, "starting", + "a running engine must not have its record retired" + ); + Ok(()) + } + + /// A spawn that fails outright strands the same record an unretired early + /// exit would: it was written before the spawn, so returning past it leaves a + /// pid-0 `"starting"` corpse that no liveness refresh can demote. The + /// original spawn error must still be what reaches the user. + #[test] + fn a_spawn_that_never_started_frees_the_engine_for_the_next_serve() -> Result<()> { + let (root, paths) = test_paths("managed-spawn-failure-unblocks"); + paths.ensure()?; + let mut record = pre_spawn_record(&paths, 11515); + record.write()?; + + let outcome: Result = retire_record_on_error( + &mut record, + Err(anyhow::anyhow!("failed to launch managed engine process")), + ); + + let still_blocking = existing_live_managed_service(&paths, "lemonade", "qwen-canonical"); + let _ = fs::remove_dir_all(root); + let error = outcome + .expect_err("a failed spawn must fail the launch") + .to_string(); + assert!( + error.contains("failed to launch managed engine process"), + "the spawn failure must survive the retirement: {error}" + ); + assert!( + still_blocking.is_none(), + "a spawn that never started must not keep claiming the engine+model" + ); + Ok(()) + } + + /// The Windows adapter between the watched spawn and the shared startup + /// check: an observed exit code has to arrive as that exit code, and a child + /// still running has to arrive as "no exit". + #[cfg(windows)] + #[test] + fn managed_startup_exit_reads_the_watched_exit_code() { + let exited = rocm_core::DetachedSpawn { + pid: 4242, + early_exit_code: Some(7), + }; + assert_eq!( + managed_startup_exit(&exited).and_then(|status| status.code()), + Some(7) + ); + + let running = rocm_core::DetachedSpawn { + pid: 4242, + early_exit_code: None, + }; + assert!(managed_startup_exit(&running).is_none()); + } + #[test] fn spawn_managed_engine_child_blocks_reuse_with_mismatched_recipe() -> Result<()> { // A live service recorded with one recipe (e.g. a tool-call parser flag) diff --git a/crates/rocm-core/src/lib.rs b/crates/rocm-core/src/lib.rs index 387e22859..789714a7a 100644 --- a/crates/rocm-core/src/lib.rs +++ b/crates/rocm-core/src/lib.rs @@ -1164,6 +1164,25 @@ fn http_response_is_complete(response: &[u8]) -> bool { false } +/// Outcome of a detached Windows spawn that was briefly watched for an early exit. +/// +/// The observation runs while the process handle returned by `CreateProcessW` is +/// still open, so a reported exit code is always the exit code of the process that +/// was just spawned. Re-opening the process by PID after the fact could not offer +/// that guarantee — Windows recycles PIDs, so by then the PID may name an +/// unrelated process. +#[cfg(windows)] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct DetachedSpawn { + /// PID of the spawned process. + pub pid: u32, + /// `Some(code)` if the process had already exited when the observation window + /// elapsed, `None` if it was still running. `None` is also reported if the + /// exit code could not be read, which degrades to the unwatched behaviour + /// rather than inventing a failure. + pub early_exit_code: Option, +} + #[cfg(windows)] pub fn spawn_detached_no_inherit( program: &Path, @@ -1177,6 +1196,38 @@ pub fn spawn_detached_no_inherit( DETACHED_PROCESS | CREATE_NEW_PROCESS_GROUP | CREATE_NO_WINDOW | CREATE_UNICODE_ENVIRONMENT, false, None, + None, + ) + .map(|spawn| spawn.pid) +} + +/// As [`spawn_detached_no_inherit`], but watches the new process for up to +/// `settle` before returning, so the caller can tell "started" apart from +/// "started and died immediately". +/// +/// Creation flags and handle inheritance are identical to +/// [`spawn_detached_no_inherit`]: the child stays fully detached and outlives this +/// process. The wait is bounded and only delays the caller by `settle`; it is not +/// a join. +/// +/// This exists because the detached spawn primitives hand back a bare PID rather +/// than a `std::process::Child`, so a caller has no equivalent of `try_wait()` to +/// notice a child that failed during startup. +#[cfg(windows)] +pub fn spawn_detached_no_inherit_watching_startup( + program: &Path, + args: &[String], + env_overrides: &[(&str, &Path)], + settle: Duration, +) -> Result { + spawn_windows_no_inherit( + program, + args, + env_overrides, + DETACHED_PROCESS | CREATE_NEW_PROCESS_GROUP | CREATE_NO_WINDOW | CREATE_UNICODE_ENVIRONMENT, + false, + None, + Some(settle), ) } @@ -1193,7 +1244,9 @@ pub fn spawn_hidden_console_no_inherit( CREATE_NEW_CONSOLE | CREATE_NEW_PROCESS_GROUP | CREATE_UNICODE_ENVIRONMENT, true, None, + None, ) + .map(|spawn| spawn.pid) } #[cfg(windows)] @@ -1262,12 +1315,13 @@ pub fn spawn_hidden_console_with_log( CREATE_NEW_CONSOLE | CREATE_NEW_PROCESS_GROUP | CREATE_UNICODE_ENVIRONMENT, true, Some((stdout_handle, stderr_handle)), + None, ); unsafe { CloseHandle(stdout_handle); CloseHandle(stderr_handle); } - result + result.map(|spawn| spawn.pid) } #[cfg(windows)] @@ -1584,7 +1638,8 @@ fn spawn_windows_no_inherit( windows_sys::Win32::Foundation::HANDLE, windows_sys::Win32::Foundation::HANDLE, )>, -) -> Result { + settle: Option, +) -> Result { use std::ptr::{null, null_mut}; use windows_sys::Win32::Foundation::CloseHandle; @@ -1628,11 +1683,59 @@ fn spawn_windows_no_inherit( std::io::Error::last_os_error() ); } + // Observe the child *before* the handle is closed. Waiting on this handle is + // race-free; re-opening the process later by PID would not be, since Windows + // recycles PIDs and the PID could by then belong to something else. + let early_exit_code = + settle.and_then(|settle| unsafe { observe_early_exit(process_info.hProcess, settle) }); unsafe { CloseHandle(process_info.hThread); CloseHandle(process_info.hProcess); } - Ok(process_info.dwProcessId) + Ok(DetachedSpawn { + pid: process_info.dwProcessId, + early_exit_code, + }) +} + +/// Wait up to `settle` for `process` to exit, reporting its exit code if it did. +/// +/// Returns `None` both when the process is still running and when the exit code +/// could not be read, so a failed query degrades to "assume it is alive" rather +/// than reporting a startup failure that may not have happened. +/// +/// # Safety +/// +/// `process` must be a live process handle granting `SYNCHRONIZE` and +/// `PROCESS_QUERY_INFORMATION` access. The handle is not closed here. +#[cfg(windows)] +#[allow(unsafe_code)] // Win32 FFI +unsafe fn observe_early_exit( + process: windows_sys::Win32::Foundation::HANDLE, + settle: Duration, +) -> Option { + use windows_sys::Win32::Foundation::WAIT_OBJECT_0; + + if unsafe { WaitForSingleObject(process, wait_timeout_millis(settle)) } != WAIT_OBJECT_0 { + return None; + } + let mut exit_code: u32 = 0; + if unsafe { GetExitCodeProcess(process, &mut exit_code) } == 0 { + return None; + } + Some(exit_code) +} + +/// Clamp a wait budget to the `u32` milliseconds `WaitForSingleObject` takes, +/// never yielding `INFINITE`. +/// +/// `INFINITE` is `u32::MAX`, so a saturating conversion of a long duration would +/// silently turn a bounded startup check into one that blocks until the child +/// exits. Cap one millisecond below it instead. +#[cfg(any(windows, test))] +fn wait_timeout_millis(budget: Duration) -> u32 { + const LONGEST_FINITE_WAIT_MS: u128 = (u32::MAX - 1) as u128; + budget.as_millis().min(LONGEST_FINITE_WAIT_MS) as u32 } #[cfg(windows)] @@ -8163,6 +8266,73 @@ mod tests { use std::path::Path; use std::path::PathBuf; + #[test] + fn wait_timeout_millis_keeps_the_startup_wait_bounded() { + assert_eq!(wait_timeout_millis(Duration::from_millis(200)), 200); + assert_eq!(wait_timeout_millis(Duration::ZERO), 0); + // Sub-millisecond budgets truncate to a poll rather than rounding up. + assert_eq!(wait_timeout_millis(Duration::from_micros(900)), 0); + // A budget past the `u32` millisecond range must never land on + // `INFINITE` (`u32::MAX`), which would block until the child exits. + let huge = wait_timeout_millis(Duration::from_secs(u64::from(u32::MAX))); + assert_eq!(huge, u32::MAX - 1); + assert_ne!(huge, u32::MAX); + } + + /// Watching a detached spawn must surface a child that died during startup. + /// Only meaningful on Windows: the detached spawn primitives return a bare + /// PID there, so this wait is the only chance to notice the exit. + #[cfg(windows)] + #[test] + fn watched_detached_spawn_reports_a_child_that_exits_immediately() { + // Each token is a separate argument so the assembled command line needs no + // quoting, keeping `cmd.exe`'s quote-stripping rules out of the test. + let spawn = spawn_detached_no_inherit_watching_startup( + &test_system32_tool("cmd.exe"), + &["/C".to_owned(), "exit".to_owned(), "7".to_owned()], + &[], + Duration::from_secs(10), + ) + .expect("spawn should succeed"); + assert_eq!(spawn.early_exit_code, Some(7)); + assert_ne!(spawn.pid, 0); + } + + /// The converse: a child that is still running must not be reported as a + /// startup failure, or every successful launch would be rejected. + #[cfg(windows)] + #[test] + fn watched_detached_spawn_does_not_report_a_child_that_keeps_running() { + // `ping` against localhost is the dependency-free Windows sleep, and ~9s + // outlasts the 200 ms observation window by well over an order of + // magnitude. Spawned directly rather than through `cmd /C` so the + // returned PID is the process the test has to clean up: on Windows + // `terminate_process_tree` is plain `terminate_process`, so an + // intermediate shell would leave `ping` orphaned. + let spawn = spawn_detached_no_inherit_watching_startup( + &test_system32_tool("ping.exe"), + &["-n".to_owned(), "10".to_owned(), "127.0.0.1".to_owned()], + &[], + Duration::from_millis(200), + ) + .expect("spawn should succeed"); + // Clean up before asserting, so a failing assertion cannot leak the child. + let early_exit_code = spawn.early_exit_code; + let was_running = process_is_running(spawn.pid); + let _ = terminate_process(spawn.pid); + assert_eq!(early_exit_code, None); + assert!(was_running); + } + + /// `CreateProcessW` is called with an explicit application name, so it does + /// no `PATH` search — the program has to be a full path. + #[cfg(windows)] + fn test_system32_tool(exe: &str) -> PathBuf { + PathBuf::from(std::env::var_os("SystemRoot").unwrap_or_else(|| r"C:\Windows".into())) + .join("System32") + .join(exe) + } + #[test] fn file_lock_creates_missing_parent_dirs_and_lock_file() { let dir = diff --git a/tests/e2e-cucumber/features/model_serving.feature b/tests/e2e-cucumber/features/model_serving.feature index 0cecc8c25..52afd5ed5 100644 --- a/tests/e2e-cucumber/features/model_serving.feature +++ b/tests/e2e-cucumber/features/model_serving.feature @@ -268,3 +268,25 @@ 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 + + # A managed engine is spawned detached, so if it dies on the way up there is no + # terminal for it to report to: the child's own service log is the only account + # of why. Without a startup check the CLI reported such a launch as a success + # and left the user waiting on a server that would never come up. The engine is + # scripted to die at startup in test builds (`rocm/e2e-test-hooks`), which also + # waives the no-GPU pre-flight and engine preparation so this reaches a real + # spawn without GPU hardware or a runtime download — the death itself is real, + # only its trigger is scripted. + # + # Carries no host tag on purpose, so it runs on every lane. `@requires-no-gpu` + # would have been wrong twice over: the premise is a dead engine, not an absent + # GPU, and that tag is skipped on any host that HAS a GPU — which would have + # excluded the Windows lane, the one place the Windows half of the startup check + # can actually execute. The waivers above are what make the premise hold + # regardless of the host's hardware. + @id:serve-managed-engine-dies-at-startup + Scenario: serve-23 - An engine that dies at startup fails the serve and names its log + Given the managed engine dies during startup + When the user serves a model with Lemonade + Then serving fails and names the engine's own log + And the failed launch does not block the next serve diff --git a/tests/e2e-cucumber/tests/e2e/serving_steps.rs b/tests/e2e-cucumber/tests/e2e/serving_steps.rs index a36489808..1189ba60f 100644 --- a/tests/e2e-cucumber/tests/e2e/serving_steps.rs +++ b/tests/e2e-cucumber/tests/e2e/serving_steps.rs @@ -832,7 +832,7 @@ async fn lemonade_preparation_cannot_complete(world: &mut E2eWorld) { } #[when("the user serves a model with Lemonade")] -async fn user_serves_with_failing_lemonade_preparation(world: &mut E2eWorld) { +async fn user_serves_with_lemonade(world: &mut E2eWorld) { let (stdout, stderr, rc) = crate::run_rocm_with_scenario_env( world, &[ @@ -1445,3 +1445,105 @@ async fn assert_response_model_correct(world: &mut E2eWorld) { fn model_ids_match(resp_model: &str, expected: &str) -> bool { e2e_cucumber::model_id::model_ids_match(resp_model, expected) } + +/// Arms the scripted startup death. Also waives serve's no-GPU pre-flight and +/// engine preparation, so the black-box test reaches a real spawn on a host with +/// neither GPU hardware nor an installed runtime. +#[given("the managed engine dies during startup")] +async fn managed_engine_dies_during_startup(world: &mut E2eWorld) { + world + .command_env + .push((MANAGED_ENGINE_STARTUP_FAILURE_ENV, "1".into())); +} + +const MANAGED_ENGINE_STARTUP_FAILURE_ENV: &str = "ROCM_E2E_MANAGED_ENGINE_STARTUP_FAILURE"; + +#[then("serving fails and names the engine's own log")] +async fn assert_startup_death_names_log(world: &mut E2eWorld) { + let output = serve_output(world); + assert_ne!( + world.cli_rc, + Some(0), + "an engine that died at startup must fail the serve:\n{output}" + ); + // See the note in assert_lemonade_preparation_retry_is_bounded: the seam that + // scripts the death also waives the no-GPU pre-flight, so this refusal can + // only appear when the binary under test was built without the feature. + assert!( + !output.contains("no usable AMD GPU detected"), + "serve stopped at the no-GPU pre-flight, so the binary under test was \ + built without the `rocm/e2e-test-hooks` feature and never reached the \ + engine launch:\n{output}" + ); + assert!( + output.contains("managed engine exited immediately"), + "expected the launch to report the engine's immediate exit:\n{output}" + ); + // The child is detached, so its own log is the only account of why it died — + // the user has to be told where it is. The path is read off the record the + // launch left behind, not matched as a `.log` suffix that any incidental + // mention of a log would satisfy; and it is the path `rocm` itself recorded, + // so a Windows 8.3 short form cannot make the comparison disagree with itself. + // + // The two platforms reach this through different branches of the message. + // On Unix the parent redirects the child's stdio into the log, so the dead + // child's output is there and the message carries a tail. On Windows the + // detached spawn does not redirect, and the scripted child never receives + // `--log`, so the log stays empty and only the no-tail branch runs. Naming + // the path is the claim both branches share, and it is the one pinned here. + let services = world + .isolated_root + .as_ref() + .expect("scenario has no isolated root") + .path() + .join("data") + .join("services"); + let log_paths: Vec = std::fs::read_dir(&services) + .unwrap_or_else(|e| panic!("no services dir at {}: {e}", services.display())) + .filter_map(|entry| std::fs::read_to_string(entry.ok()?.path()).ok()) + .filter_map(|json| serde_json::from_str::(&json).ok()) + .filter(|record| record["engine"] == "lemonade") + .filter_map(|record| record["log_path"].as_str().map(str::to_owned)) + .collect(); + assert!( + !log_paths.is_empty(), + "the failed launch left no lemonade service record in {}", + services.display() + ); + assert!( + log_paths.iter().any(|path| output.contains(path.as_str())), + "expected the failure to name the engine's own service log {log_paths:?}:\n{output}" + ); +} + +#[then("the failed launch does not block the next serve")] +async fn assert_failed_launch_does_not_block_retry(world: &mut E2eWorld) { + // The record is written before the engine is spawned, so a launch that dies + // leaves one behind. If it is not retired, it reads as a live service and the + // idempotency guard refuses every later serve of the same engine + model — + // the user could never retry. Re-arm the seam: the runner consumes the + // scenario env on each invocation. + world + .command_env + .push((MANAGED_ENGINE_STARTUP_FAILURE_ENV, "1".into())); + let (stdout, stderr, rc) = crate::run_rocm_with_scenario_env( + world, + &[ + "serve", + "Qwen3-0.6B-GGUF", + "--engine", + "lemonade", + "--managed", + ], + ); + let output = format!("{stdout}\n{stderr}"); + assert_ne!( + rc, 0, + "the retry should fail the same way, not succeed:\n{output}" + ); + assert!( + output.contains("managed engine exited immediately"), + "the retry must reach the engine launch again rather than being turned \ + away by the previous failed launch:\n{output}" + ); +}