diff --git a/crates/client/tests/process_e2e.rs b/crates/client/tests/process_e2e.rs index fed8bd693d..c80343e54b 100644 --- a/crates/client/tests/process_e2e.rs +++ b/crates/client/tests/process_e2e.rs @@ -154,6 +154,60 @@ async fn output_read_right_after_a_failed_spawn_returns_the_launch_error() { .expect("output read returns the launch error"); } +#[tokio::test] +async fn concurrent_launches_on_a_fresh_vm_all_run() { + if !common::require_sidecar("concurrent_launches_on_a_fresh_vm_all_run") { + return; + } + let Some(os) = common::new_vm_with_commands().await else { + if common::allow_local_e2e_skips() { + eprintln!( + "skipping concurrent_launches_on_a_fresh_vm_all_run: coreutils package absent" + ); + return; + } + panic!("concurrent_launches_on_a_fresh_vm_all_run: coreutils package absent"); + }; + // A fresh VM is cold: the first WebAssembly and Python launches await runtime setup while + // the others arrive, so every launch here overlaps another one of the same runtime. + let launches: Vec<(&str, Vec)> = (0..3) + .map(|index| ("echo", vec![format!("wasm-{index}")])) + .chain((0..3).map(|index| { + ( + "python3", + vec!["-c".to_owned(), format!("print('python-{index}')")], + ) + })) + .collect(); + let result = tokio::time::timeout(std::time::Duration::from_secs(120), async { + let mut handles = Vec::new(); + for (command, args) in &launches { + handles.push(os.spawn_process(command, args.clone(), SpawnOptions::default())?); + } + let mut outcomes = Vec::new(); + for ((command, args), handle) in launches.iter().zip(&handles) { + outcomes.push((command, args, os.wait_process(handle.pid).await)); + } + let failed: Vec<_> = outcomes + .iter() + .filter(|(_, _, outcome)| !matches!(outcome, Ok(0))) + .collect(); + anyhow::ensure!( + failed.is_empty(), + "every concurrent launch must run to exit 0; failed: {failed:?}" + ); + Ok::<_, anyhow::Error>(()) + }) + .await; + tokio::time::timeout(std::time::Duration::from_secs(10), os.shutdown()) + .await + .expect("bounded VM cleanup") + .expect("shutdown VM"); + result + .expect("concurrent launches finish within the timeout") + .expect("concurrent launches all run"); +} + #[tokio::test] async fn exec_argv_timeout_confirms_node_guest_exit() { if !common::require_sidecar("exec_argv_timeout_confirms_node_guest_exit") { diff --git a/crates/native-sidecar/src/execution/child_process.rs b/crates/native-sidecar/src/execution/child_process.rs index 2535927155..aa72715856 100644 --- a/crates/native-sidecar/src/execution/child_process.rs +++ b/crates/native-sidecar/src/execution/child_process.rs @@ -4907,6 +4907,8 @@ where ActiveExecution::Javascript(execution) } GuestRuntimeKind::WebAssembly => { + // Take turns with other launches on this VM: a cold start awaits while holding the engine. + let launch_turn = execution_engines.wasm_launch_turn().await; let mut wasm_engine = execution_engines.wasm("start WebAssembly child process")?; // These values configure the trusted WASM runner, not @@ -4947,6 +4949,8 @@ where ) .await; wasm_engine.dispose_context(&context_id); + drop(wasm_engine); + drop(launch_turn); let execution = execution_result.map_err(wasm_error)?; ActiveExecution::Wasm(Box::new(execution)) } @@ -6558,6 +6562,8 @@ where let runtime_context = vm.runtime_context.clone(); let execution_engines = execution_engines.clone(); Box::pin(async move { + // Take turns with other launches on this VM: a cold start awaits while holding the engine. + let _launch_turn = execution_engines.wasm_launch_turn().await; let mut wasm_engine = execution_engines.wasm("start nested WebAssembly child process")?; let context = wasm_engine.create_context(CreateWasmContextRequest { @@ -6601,6 +6607,8 @@ where let code = resolved.entrypoint.clone(); let cwd = resolved.host_cwd.clone(); Box::pin(async move { + // Take turns with other launches on this VM: a cold start awaits while holding the engine. + let _launch_turn = execution_engines.python_launch_turn().await; let mut python_engine = execution_engines.python("start nested Python child process")?; let pyodide_dist_path = python_engine diff --git a/crates/native-sidecar/src/execution/launch.rs b/crates/native-sidecar/src/execution/launch.rs index 10afe5817a..e276eaffcc 100644 --- a/crates/native-sidecar/src/execution/launch.rs +++ b/crates/native-sidecar/src/execution/launch.rs @@ -5115,6 +5115,9 @@ where // coordinator, so filesystem/kernel commands can still enter this // VM and operations in other VMs remain independent. drop(vm); + // A cold start awaits Pyodide warmup while holding the engine; take turns with + // other launches on this VM instead of failing with an execution conflict. + let launch_turn = execution_engines.python_launch_turn().await; let mut python_engine = execution_engines.python("start Python execution")?; let pyodide_dist_path = python_engine .bundled_pyodide_dist_path_for_vm_async(&vm_id, &runtime_context) @@ -5170,6 +5173,7 @@ where .await .map_err(python_error)?; drop(python_engine); + drop(launch_turn); vm = input .vm .try_borrow_mut("register started Python execution")?; @@ -5191,6 +5195,10 @@ where // WASM import-cache materialization and prewarm are external // waits. Release mutable VM state before entering them. drop(vm); + // A cold start awaits import-cache materialization and prewarm while holding the + // engine; take turns with other launches on this VM instead of failing with an + // execution conflict. + let launch_turn = execution_engines.wasm_launch_turn().await; let mut wasm_engine = execution_engines.wasm("start WebAssembly execution")?; let context = wasm_engine.create_context(CreateWasmContextRequest { vm_id: vm_id.clone(), @@ -5213,6 +5221,7 @@ where .await .map_err(wasm_error)?; drop(wasm_engine); + drop(launch_turn); vm = input .vm .try_borrow_mut("register started WebAssembly execution")?; diff --git a/crates/native-sidecar/src/state.rs b/crates/native-sidecar/src/state.rs index 2a7e304030..f635db0e0c 100644 --- a/crates/native-sidecar/src/state.rs +++ b/crates/native-sidecar/src/state.rs @@ -1100,6 +1100,10 @@ struct VmExecutionEnginesInner { javascript: RefCell, python: RefCell, wasm: RefCell, + /// Launches that hold an engine across an await take turns here, so a second launch + /// waits instead of failing with an execution conflict. + python_launch: tokio::sync::Mutex<()>, + wasm_launch: tokio::sync::Mutex<()>, } impl VmExecutionEngines { @@ -1120,6 +1124,8 @@ impl VmExecutionEngines { javascript: RefCell::new(javascript), python: RefCell::new(python), wasm: RefCell::new(wasm), + python_launch: tokio::sync::Mutex::new(()), + wasm_launch: tokio::sync::Mutex::new(()), }), } } @@ -1143,6 +1149,14 @@ impl VmExecutionEngines { .map_err(|_| execution_engine_conflict_error(&self.inner.vm_id, "Python", operation)) } + pub(crate) async fn python_launch_turn(&self) -> tokio::sync::MutexGuard<'_, ()> { + self.inner.python_launch.lock().await + } + + pub(crate) async fn wasm_launch_turn(&self) -> tokio::sync::MutexGuard<'_, ()> { + self.inner.wasm_launch.lock().await + } + pub(crate) fn wasm( &self, operation: &'static str,