Skip to content
Merged
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
54 changes: 54 additions & 0 deletions crates/client/tests/process_e2e.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String>)> = (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") {
Expand Down
8 changes: 8 additions & 0 deletions crates/native-sidecar/src/execution/child_process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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))
}
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
Expand Down
9 changes: 9 additions & 0 deletions crates/native-sidecar/src/execution/launch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Comment on lines +5118 to +5120

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔴 High · Queue before creating an unregistered process

This await happens after the duplicate-ID check and after spawn_trusted_root_process. If a queued request is cancelled, dropping KernelProcessHandle does not finish or reap that process, so it remains in the VM kernel without an active_processes entry. Two same-runtime requests with the same process_id can also both pass the check; the queue then lets both start, and the later active_processes.insert silently replaces the first. Acquire the runtime turn before creating/reserving the process and hold it through registration, or add a cancellation-safe pending reservation that claims the ID and rolls back the kernel process.

Comment thread
eersnington marked this conversation as resolved.
Comment thread
eersnington marked this conversation as resolved.
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)
Expand Down Expand Up @@ -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")?;
Expand All @@ -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(),
Expand All @@ -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")?;
Expand Down
14 changes: 14 additions & 0 deletions crates/native-sidecar/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1100,6 +1100,10 @@ struct VmExecutionEnginesInner {
javascript: RefCell<JavascriptExecutionEngine>,
python: RefCell<PythonExecutionEngine>,
wasm: RefCell<WasmExecutionEngine>,
/// 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 {
Expand All @@ -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(()),
}),
}
}
Expand All @@ -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<'_, ()> {
Comment on lines 1149 to +1152

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 Medium · Engine admission remains bypassable

The queue is independent of python()/wasm(), so existing launch paths can still take the raw RefCell and return the conflict this change is intended to remove. In production, exec_javascript_process_image_owned directly borrows the WASM/Python engines at child_process.rs:5586 and :5620 while another cold start may hold them; the legacy spawn path also queues WASM but still directly borrows Python at :4953. Encapsulate queue acquisition with the engine borrow and route every Python/WASM launch or replacement through that admission path.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

not introduced by this pr. but need to look into this deeper

Comment thread
eersnington marked this conversation as resolved.
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,
Expand Down
Loading