diff --git a/crates/client/tests/process_e2e.rs b/crates/client/tests/process_e2e.rs index efc14f0c85..04cbb63025 100644 --- a/crates/client/tests/process_e2e.rs +++ b/crates/client/tests/process_e2e.rs @@ -406,3 +406,59 @@ async fn process_surface_exec_spawn_and_snapshot() { os.shutdown().await.expect("shutdown"); } + +#[tokio::test] +async fn shell_children_of_every_runtime_run_to_completion() { + if !common::require_sidecar("shell_children_of_every_runtime_run_to_completion") { + return; + } + let Some(os) = common::new_vm_with_commands().await else { + if common::allow_local_e2e_skips() { + eprintln!( + "skipping shell_children_of_every_runtime_run_to_completion: coreutils package absent" + ); + return; + } + panic!("shell_children_of_every_runtime_run_to_completion: coreutils package absent"); + }; + // The shell waits on each child before it runs the next line, so every line of output + // proves the previous child ran, reported its exit status, and let the shell continue. + let script = [ + r#"python3 -c "print(6*7)""#, + r#"echo 'print(1+1)' | python3 -"#, + r#"python3 -c "raise ValueError" 2>/dev/null; echo "python-exit-$?""#, + "echo nested > /tmp/nested.txt", + r#"sh -c "cat /tmp/nested.txt""#, + ] + .join("\n"); + let result = tokio::time::timeout(std::time::Duration::from_secs(60), async { + let handle = os.spawn_process( + "sh", + vec!["-c".to_owned(), script], + SpawnOptions { + retain_output: true, + ..Default::default() + }, + )?; + let exit_code = os.wait_process(handle.pid).await?; + let stdout = os + .read_process_output(handle.pid, None, None, None) + .await? + .events + .into_iter() + .filter(|event| event.stream == agentos_client::ProcessStream::Stdout) + .flat_map(|event| event.data) + .collect::>(); + Ok::<_, anyhow::Error>((exit_code, String::from_utf8_lossy(&stdout).into_owned())) + }) + .await; + tokio::time::timeout(std::time::Duration::from_secs(10), os.shutdown()) + .await + .expect("bounded VM cleanup") + .expect("shutdown VM"); + let (exit_code, stdout) = result + .expect("the shell and its children finish within the timeout") + .expect("the shell runs"); + assert_eq!(stdout, "42\n2\npython-exit-1\nnested\n"); + assert_eq!(exit_code, 0); +} diff --git a/crates/native-sidecar/src/execution/child_process.rs b/crates/native-sidecar/src/execution/child_process.rs index a4a279d8ef..2535927155 100644 --- a/crates/native-sidecar/src/execution/child_process.rs +++ b/crates/native-sidecar/src/execution/child_process.rs @@ -2773,75 +2773,81 @@ where // The standalone WASM runner pulls descendant output through // child_process.poll while implementing waitpid. Keep stream // and exit delivery single-owner, but still claim internal - // runtime RPCs from JavaScript children. Otherwise a command - // such as `npm test` deadlocks when npm asks the sidecar to - // spawn its script shell while the WASM parent is waiting for - // npm to exit. - let (parent_is_pull_driven_wasm, child_is_javascript, pending_needs_supervisor) = self + // runtime RPCs from every child. Otherwise `npm test` deadlocks + // when npm asks the sidecar to spawn its script shell while the + // WASM parent is waiting for npm to exit, and a Python child or + // a nested shell never gets its filesystem or spawn RPCs answered. + let (parent_is_pull_driven_wasm, needs_supervisor) = self .vms .get(vm_id) .map(|vm| { - let parent = vm.active_processes + let parent = vm + .active_processes .get(process_id) .and_then(|root| Self::active_process_by_path(root, &parent_path)); - let child = parent.and_then(|parent| parent.child_processes.get(&child_process_id)); - ( - parent.is_some_and(|parent| parent.runtime == GuestRuntimeKind::WebAssembly), - child.is_some_and(|child| child.runtime == GuestRuntimeKind::JavaScript), - child - .and_then(|child| child.pending_execution_events.front()) - .is_some_and(|event| { - matches!(event, ActiveExecutionEvent::JavascriptSyncRpcRequest(request) if javascript_rpc_requires_owned_supervisor(&request.method)) - }), - ) + let child = + parent.and_then(|parent| parent.child_processes.get(&child_process_id)); + ( + parent.is_some_and(|parent| { + parent.runtime == GuestRuntimeKind::WebAssembly + }), + child.is_some_and(|child| { + match child.pending_execution_events.front() { + Some(ActiveExecutionEvent::JavascriptSyncRpcRequest( + request, + )) => javascript_rpc_requires_owned_supervisor(&request.method), + Some( + ActiveExecutionEvent::JavascriptSyncRpcCompletion(_) + | ActiveExecutionEvent::PythonVfsRpcRequest(_) + | ActiveExecutionEvent::PythonSocketConnectCompletion(_) + | ActiveExecutionEvent::SignalState { .. }, + ) => true, + Some( + ActiveExecutionEvent::Stdout(_) + | ActiveExecutionEvent::Stderr(_) + | ActiveExecutionEvent::Exited(_), + ) => false, + // The parent's poll queues every supervisor-owned event, but a + // JavaScript child can also expose one before the parent polls. + None => child.runtime == GuestRuntimeKind::JavaScript, + } + }), + ) }) - .unwrap_or((false, false, false)); + .unwrap_or((false, false)); if parent_is_pull_driven_wasm { - if child_is_javascript - && self.vms.get(vm_id).is_some_and(|vm| { - vm.active_processes - .get(process_id) - .and_then(|root| Self::active_process_by_path(root, &parent_path)) - .and_then(|parent| parent.child_processes.get(&child_process_id)) - .is_some_and(|child| !child.pending_execution_events.is_empty()) - }) - && !pending_needs_supervisor - { + if !needs_supervisor { continue; } - if child_is_javascript { - let services_before = javascript_services - .len() - .saturating_add(python_services.len()) - .saturating_add(python_socket_completions.len()); - let claimed = match self.poll_descendant_javascript_child_process_nowait( - vm_id, - process_id, - &parent_path, - &child_process_id, - true, - javascript_services, - python_services, - python_socket_completions, - None, - ) { - Ok(event) => event, - Err(error) if is_javascript_child_process_gone_error(&error) => { - continue - } - Err(error) => return Err(error), - }; - drop(claimed.reservation); - let services_after = javascript_services - .len() - .saturating_add(python_services.len()) - .saturating_add(python_socket_completions.len()); - if services_after > services_before { - emitted_any = true; - emitted_this_round = true; - work += 1; - child_work[candidate_index] += 1; - } + let services_before = javascript_services + .len() + .saturating_add(python_services.len()) + .saturating_add(python_socket_completions.len()); + let claimed = match self.poll_descendant_javascript_child_process_nowait( + vm_id, + process_id, + &parent_path, + &child_process_id, + true, + javascript_services, + python_services, + python_socket_completions, + None, + ) { + Ok(event) => event, + Err(error) if is_javascript_child_process_gone_error(&error) => continue, + Err(error) => return Err(error), + }; + drop(claimed.reservation); + let services_after = javascript_services + .len() + .saturating_add(python_services.len()) + .saturating_add(python_socket_completions.len()); + if services_after > services_before { + emitted_any = true; + emitted_this_round = true; + work += 1; + child_work[candidate_index] += 1; } continue; }