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
56 changes: 56 additions & 0 deletions crates/client/tests/process_e2e.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<Vec<_>>();
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);
}
126 changes: 66 additions & 60 deletions crates/native-sidecar/src/execution/child_process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Comment on lines +2810 to +2812

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 · Probe empty non-JavaScript child queues

A Python child’s first VFS request is still in its runtime event receiver, so pending_execution_events.front() is None here and this returns false. The WASM-parent branch then continues without calling try_poll_execution_event; if the parent reaches child_process.poll first, that same request is polled without owned_python_services and fails with ERR_AGENTOS_PYTHON_EVENT_OWNERSHIP instead of being serviced. This leaves the Python-under-shell path added by the E2E test dependent on a queue state that normal runtime events do not create. Probe children of every runtime when this queue is empty (for example, let the existing preserving nowait poll run for None) so Python VFS requests and WASM internal RPCs can be discovered and claimed.

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.

slop. prod's dispatcher uses poll_owned_descendant_javascript_child_process which

  1. reads the python request
  2. queues it for the supervisor
  3. calls notify_one()

the shell queues the py req, and this fix handles it. this is looking at the wrong code path

}
}),
)
})
.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;
}
Expand Down
Loading