From 0a701f7b65edc19873cb971e6d03e1cf76f1a887 Mon Sep 17 00:00:00 2001 From: ankssjain <308424510+ankssjain@users.noreply.github.com> Date: Wed, 30 Sep 2026 01:28:38 +0530 Subject: [PATCH 1/3] fix(native-sidecar): service Python RPCs beneath WASM parents --- .../src/execution/child_process.rs | 23 ++- crates/native-sidecar/tests/service.rs | 173 ++++++++++++++++++ 2 files changed, 194 insertions(+), 2 deletions(-) diff --git a/crates/native-sidecar/src/execution/child_process.rs b/crates/native-sidecar/src/execution/child_process.rs index 53a652c345..ddf3d4b649 100644 --- a/crates/native-sidecar/src/execution/child_process.rs +++ b/crates/native-sidecar/src/execution/child_process.rs @@ -2769,7 +2769,26 @@ where }) .unwrap_or(false); if parent_is_pull_driven_wasm { - continue; + // The parent's poller requeues Python bridge events for + // the supervisor. Service those requests even though the + // parent still owns delivery of stdout, stderr, and exit. + let queued_python_event = 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)) + .and_then(|child| child.pending_execution_events.front()) + .is_some_and(|event| { + matches!( + event, + ActiveExecutionEvent::PythonVfsRpcRequest(_) + | ActiveExecutionEvent::PythonSocketConnectCompletion(_) + ) + }) + }); + if !queued_python_event { + continue; + } } self.expire_child_process_sync_if_needed( vm_id, @@ -2783,7 +2802,7 @@ where process_id, &parent_path, &child_process_id, - false, + parent_is_pull_driven_wasm, javascript_services, python_services, python_socket_completions, diff --git a/crates/native-sidecar/tests/service.rs b/crates/native-sidecar/tests/service.rs index 03ef99a087..5b7f5ac866 100644 --- a/crates/native-sidecar/tests/service.rs +++ b/crates/native-sidecar/tests/service.rs @@ -25585,6 +25585,179 @@ console.log(JSON.stringify({ } } + #[test] + fn wasm_parent_services_queued_python_requests_without_consuming_child_output() { + let mut sidecar = create_test_sidecar(); + let (connection_id, session_id) = + authenticate_and_open_session(&mut sidecar).expect("authenticate sidecar"); + let vm_id = create_vm( + &mut sidecar, + &connection_id, + &session_id, + PermissionsPolicy::allow_all(), + ) + .expect("create vm"); + let context = create_python_context_for_vm_test( + &sidecar, + &vm_id, + PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../execution/assets/pyodide"), + ); + let limits = { + let vm = sidecar.vms.get(&vm_id).expect("Python test VM"); + agentos_execution::PythonExecutionLimits { + reactor_work_quantum: Some(vm.limits.reactor.work_quantum), + bridge_call_timeout_ms: Some( + vm.limits + .reactor + .operation_deadline_ms + .saturating_add(1_000), + ), + max_open_fds: vm.kernel.resource_limits().max_open_fds, + ..Default::default() + } + }; + let execution = start_python_execution_for_vm_test( + &sidecar, + &vm_id, + StartPythonExecutionRequest { + guest_runtime: Default::default(), + limits, + vm_id: vm_id.clone(), + context_id: context.context_id, + code: String::from("print(42)"), + file_path: None, + env: BTreeMap::new(), + cwd: temp_dir("agentos-wasm-parent-python-rpc"), + }, + ) + .expect("start Python execution"); + let root_handle = create_kernel_process_handle_for_tests(); + let mut root = active_process_for_tests( + root_handle.pid(), + root_handle, + GuestRuntimeKind::WebAssembly, + ActiveExecution::HostFunction(HostFunctionExecution::default()), + ); + let child_handle = create_kernel_process_handle_for_tests(); + let mut child = active_process_for_tests( + child_handle.pid(), + child_handle, + GuestRuntimeKind::Python, + ActiveExecution::Python(execution), + ); + // The WASM parent's poller requeues these events for the supervisor. + // Neither request may be stranded behind the WASM-parent guard. + child + .queue_pending_execution_event(ActiveExecutionEvent::PythonVfsRpcRequest(Box::new( + PythonVfsRpcRequest { + id: 1, + method: PythonVfsRpcMethod::ReadDir, + path: String::from("/"), + destination: None, + target: None, + mode: None, + uid: None, + gid: None, + atime_ms: None, + mtime_ms: None, + content_base64: None, + recursive: false, + url: None, + http_method: None, + headers: BTreeMap::new(), + body_base64: None, + hostname: None, + family: None, + port: None, + socket_id: None, + command: None, + args: Vec::new(), + argv0: None, + cwd: None, + env: BTreeMap::new(), + shell: false, + max_buffer: None, + timeout_ms: None, + }, + ))) + .expect("queue Python filesystem request"); + child + .queue_pending_execution_event(ActiveExecutionEvent::PythonSocketConnectCompletion( + Box::new(crate::state::PythonSocketConnectCompletion { + request_id: 2, + result: Err(crate::state::DeferredRpcError { + code: String::from("ECONNREFUSED"), + message: String::from("test connection refused"), + }), + }), + )) + .expect("queue Python socket completion"); + for event in [ + ActiveExecutionEvent::Stdout(b"42\n".to_vec()), + ActiveExecutionEvent::Stderr(b"diagnostic\n".to_vec()), + ActiveExecutionEvent::Exited(0), + ] { + child + .queue_pending_execution_event(event) + .expect("queue shell-owned event"); + } + root.child_processes + .insert(String::from("python-child"), child); + sidecar + .vms + .get_mut(&vm_id) + .expect("test vm") + .active_processes + .insert(String::from("wasm-root"), root); + + let mut javascript_services = Vec::new(); + let mut python_services = Vec::new(); + let mut socket_completions = Vec::new(); + let mut child_bridge_services = Vec::new(); + // One claim per pump also checks that the second queued request + // makes progress on the next supervisor turn. + for _ in 0..2 { + sidecar + .pump_child_process_events_nowait( + &vm_id, + &mut javascript_services, + &mut python_services, + &mut socket_completions, + &mut child_bridge_services, + 2, + ) + .expect("pump Python child requests"); + } + assert_eq!( + python_services.len(), + 1, + "supervisor must claim Python VFS request" + ); + assert_eq!( + python_services[0].request.method, + PythonVfsRpcMethod::ReadDir + ); + assert_eq!(python_services[0].child_path, ["python-child"]); + assert_eq!( + socket_completions.len(), + 1, + "supervisor must claim Python socket completion" + ); + assert_eq!(socket_completions[0].completion.request_id, 2); + assert!(javascript_services.is_empty()); + assert!(child_bridge_services.is_empty()); + assert!(sidecar.pending_process_events.is_empty()); + let vm = sidecar.vms.get(&vm_id).expect("test vm"); + let child = &vm.active_processes["wasm-root"].child_processes["python-child"]; + let events = &child.pending_execution_events; + assert_eq!(events.len(), 3, "shell retains stdout, stderr, and exit"); + assert!(matches!(&events[0], ActiveExecutionEvent::Stdout(bytes) if bytes == b"42\n")); + assert!( + matches!(&events[1], ActiveExecutionEvent::Stderr(bytes) if bytes == b"diagnostic\n") + ); + assert!(matches!(&events[2], ActiveExecutionEvent::Exited(0))); + } + #[test] fn wasm_parent_child_write_deadline_wakes_after_parent_stops_polling() { assert_node_available(); From 57165577c1f2579457a1004fbf258d15638d7f53 Mon Sep 17 00:00:00 2001 From: Anks Metaforms Date: Wed, 30 Sep 2026 15:59:10 +0530 Subject: [PATCH 2/3] fix(native-sidecar): account for all child service claims --- .../src/execution/child_process.rs | 220 ++++++++++++++++-- .../src/execution/process_events.rs | 2 +- crates/native-sidecar/tests/service.rs | 69 ++++-- 3 files changed, 254 insertions(+), 37 deletions(-) diff --git a/crates/native-sidecar/src/execution/child_process.rs b/crates/native-sidecar/src/execution/child_process.rs index ddf3d4b649..5862fbc165 100644 --- a/crates/native-sidecar/src/execution/child_process.rs +++ b/crates/native-sidecar/src/execution/child_process.rs @@ -289,6 +289,9 @@ impl Drop for ChildBridgeRelayLease { struct ClaimedDescendantBridgeEvent { event: Value, reservation: Option, + // Claiming an owned service is progress even when there is no event to + // deliver to the parent. Both attached and detached pumps must budget it. + service_claims: usize, } impl OwnedChildBridgeEventService { @@ -2702,7 +2705,7 @@ where let mut yielded = false; loop { - let mut emitted_this_round = false; + let mut progressed_this_round = false; for (candidate_index, (process_id, child_path)) in child_candidates.iter().enumerate() { if javascript_services .len() @@ -2769,24 +2772,18 @@ where }) .unwrap_or(false); if parent_is_pull_driven_wasm { - // The parent's poller requeues Python bridge events for - // the supervisor. Service those requests even though the - // parent still owns delivery of stdout, stderr, and exit. - let queued_python_event = self.vms.get(vm_id).is_some_and(|vm| { + // The parent's poller leaves internal requests and + // completions for the supervisor, regardless of language. + // Stream and exit events remain owned by the parent. + let queued_internal_event = 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)) .and_then(|child| child.pending_execution_events.front()) - .is_some_and(|event| { - matches!( - event, - ActiveExecutionEvent::PythonVfsRpcRequest(_) - | ActiveExecutionEvent::PythonSocketConnectCompletion(_) - ) - }) + .is_some_and(Self::internal_execution_event) }); - if !queued_python_event { + if !queued_internal_event { continue; } } @@ -2812,7 +2809,16 @@ where Err(error) if is_javascript_child_process_gone_error(&error) => continue, Err(error) => return Err(error), }; - let ClaimedDescendantBridgeEvent { event, reservation } = claimed; + let ClaimedDescendantBridgeEvent { + event, + reservation, + service_claims, + } = claimed; + if service_claims > 0 { + progressed_this_round = true; + work += service_claims; + child_work[candidate_index] += service_claims; + } if event.is_null() { continue; } @@ -2829,11 +2835,11 @@ where break; } emitted_any = true; - emitted_this_round = true; + progressed_this_round = true; work += 1; child_work[candidate_index] += 1; } - if yielded || !emitted_this_round { + if yielded || !progressed_this_round { break; } } @@ -3607,9 +3613,20 @@ where } Err(error) => return Err(error), }; - let ClaimedDescendantBridgeEvent { event, reservation } = claimed; + let ClaimedDescendantBridgeEvent { + event, + reservation, + service_claims, + } = claimed; + work += service_claims; + detached_work += service_claims; let Some(event_type) = event.get("type").and_then(Value::as_str) else { + if service_claims > 0 { + // Continue draining owned services, or rearm at the + // capacity/fairness checks at the top of the loop. + continue; + } break; }; let public_event_queued = event @@ -8534,6 +8551,9 @@ where owned_python_socket_completions: &mut Vec, public_process_id: Option<&str>, ) -> Result { + let service_count_before = owned_javascript_services.len() + + owned_python_services.len() + + owned_python_socket_completions.len(); let mut claimed_reservation = None; let mut future = Box::pin(self.poll_descendant_javascript_child_process( vm_id, @@ -8542,19 +8562,24 @@ where child_process_id, 0, preserve_pull_owned_events, - Some(owned_javascript_services), - Some(owned_python_services), - Some(owned_python_socket_completions), + Some(&mut *owned_javascript_services), + Some(&mut *owned_python_services), + Some(&mut *owned_python_socket_completions), Some(&mut claimed_reservation), public_process_id, )); let mut context = Context::from_waker(Waker::noop()); let poll = future.as_mut().poll(&mut context); drop(future); + let service_claims = owned_javascript_services.len() + + owned_python_services.len() + + owned_python_socket_completions.len() + - service_count_before; match poll { Poll::Ready(result) => result.map(|event| ClaimedDescendantBridgeEvent { event, reservation: claimed_reservation, + service_claims, }), Poll::Pending => Err(SidecarError::InvalidState(String::from( "ERR_AGENTOS_CHILD_EVENT_TURN_SUSPENDED: bounded child event claim unexpectedly suspended", @@ -11141,6 +11166,161 @@ mod child_event_claim_tests { .await; } + async fn assert_child_service_claims_make_bounded_progress(detached: bool) { + // Test ordinary parents as well as pull-driven shells. Every producer + // runs before the first turn, so further turns require pump rearming. + for wasm_parent in [false, true] { + for (capacity, vm_quantum, child_quantum, batch_size) in + [(8, 8, 8, 3), (1, 8, 8, 1), (8, 1, 8, 1), (8, 8, 1, 1)] + { + let (mut sidecar, vm_id) = + sidecar_with_test_vm(agentos_runtime::DEFAULT_PROTOCOL_MAX_PROCESS_EVENTS) + .await; + sidecar.config.runtime.fairness.vm_quantum_operations = vm_quantum; + sidecar + .config + .runtime + .fairness + .capability_quantum_operations = child_quantum; + let notify = Arc::clone(&sidecar.process_event_notify); + let root_id = String::from("service-progress-root"); + let child_id = String::from("service-progress-child"); + { + let mut vm = sidecar.vms.get_mut(&vm_id).expect("service progress VM"); + let mut root = host_function_process(&mut vm, "service parent", None); + if wasm_parent { + root.runtime = GuestRuntimeKind::WebAssembly; + } + let mut child = + host_function_process(&mut vm, "service child", Some(root.kernel_pid)) + .with_event_notify(Arc::clone(¬ify)); + for id in 1..=3 { + child + .queue_pending_execution_event(rpc(id)) + .expect("queue child service request"); + } + if wasm_parent && !detached { + for event in [ + ActiveExecutionEvent::Stdout(b"shell stdout".to_vec()), + ActiveExecutionEvent::Stderr(b"shell stderr".to_vec()), + ActiveExecutionEvent::Exited(0), + ] { + child + .queue_pending_execution_event(event) + .expect("queue shell-owned output"); + } + } + root.child_processes.insert(child_id.clone(), child); + vm.active_processes.insert(root_id.clone(), root); + if detached { + vm.detached_child_processes + .insert(format!("{root_id}/{child_id}")); + } + } + let ownership = + OwnershipScope::vm("child-claim-connection", "child-claim-session", &vm_id); + let mut claimed_ids = Vec::new(); + while claimed_ids.len() < 3 { + tokio::time::timeout(Duration::from_millis(50), notify.notified()) + .await + .expect("queued child services must retain a continuation edge"); + let services = if detached { + let mut javascript = Vec::new(); + sidecar + .pump_detached_child_process_events_nowait( + &vm_id, + &mut javascript, + &mut Vec::new(), + &mut Vec::new(), + &mut [], + capacity, + ) + .expect("pump detached descendant services"); + javascript + } else { + sidecar + .pump_process_events_nowait(&ownership, capacity) + .expect("pump attached descendant services") + .javascript_services + }; + assert_eq!( + services.len(), + batch_size, + "detached={detached}, wasm_parent={wasm_parent}, capacity={capacity}, \ + vm_quantum={vm_quantum}, child_quantum={child_quantum}" + ); + for service in services { + assert_eq!(service.process_id, root_id); + assert_eq!(service.child_path, [child_id.clone()]); + claimed_ids.push(service.request.id); + } + } + assert_eq!( + claimed_ids, + [1, 2, 3], + "claims are ordered and exactly once" + ); + + // A bounded last turn may leave one conservative continuation. + // An empty follow-up must quiesce rather than repeatedly wake. + consume_notify_permit(¬ify).await; + let empty = if detached { + let mut javascript = Vec::new(); + sidecar + .pump_detached_child_process_events_nowait( + &vm_id, + &mut javascript, + &mut Vec::new(), + &mut Vec::new(), + &mut [], + capacity, + ) + .expect("empty detached follow-up"); + javascript + } else { + sidecar + .pump_process_events_nowait(&ownership, capacity) + .expect("empty attached follow-up") + .javascript_services + }; + assert!(empty.is_empty()); + assert!( + tokio::time::timeout(Duration::from_millis(10), notify.notified()) + .await + .is_err(), + "an empty child turn must not hot-rearm itself" + ); + if wasm_parent && !detached { + let vm = sidecar.vms.get(&vm_id).expect("shell ownership VM"); + let events = &vm.active_processes[&root_id].child_processes[&child_id] + .pending_execution_events; + assert_eq!(events.len(), 3); + assert!( + matches!(&events[0], ActiveExecutionEvent::Stdout(bytes) if bytes == b"shell stdout") + ); + assert!( + matches!(&events[1], ActiveExecutionEvent::Stderr(bytes) if bytes == b"shell stderr") + ); + assert!(matches!(&events[2], ActiveExecutionEvent::Exited(0))); + } + } + } + } + + #[tokio::test(flavor = "current_thread")] + async fn attached_service_claims_make_bounded_progress_after_coalesced_notify() { + tokio::task::LocalSet::new() + .run_until(assert_child_service_claims_make_bounded_progress(false)) + .await; + } + + #[tokio::test(flavor = "current_thread")] + async fn detached_service_claims_make_bounded_progress_after_coalesced_notify() { + tokio::task::LocalSet::new() + .run_until(assert_child_service_claims_make_bounded_progress(true)) + .await; + } + #[tokio::test(flavor = "current_thread")] async fn root_public_event_rearm_reaches_terminal_after_exact_256_update_quantum() { tokio::task::LocalSet::new() diff --git a/crates/native-sidecar/src/execution/process_events.rs b/crates/native-sidecar/src/execution/process_events.rs index 912da20eba..3eee0059ef 100644 --- a/crates/native-sidecar/src/execution/process_events.rs +++ b/crates/native-sidecar/src/execution/process_events.rs @@ -1273,7 +1273,7 @@ where Ok(()) } - fn internal_execution_event(event: &ActiveExecutionEvent) -> bool { + pub(super) fn internal_execution_event(event: &ActiveExecutionEvent) -> bool { matches!( event, ActiveExecutionEvent::JavascriptSyncRpcRequest(_) diff --git a/crates/native-sidecar/tests/service.rs b/crates/native-sidecar/tests/service.rs index 5b7f5ac866..5e37f2511b 100644 --- a/crates/native-sidecar/tests/service.rs +++ b/crates/native-sidecar/tests/service.rs @@ -25587,7 +25587,21 @@ console.log(JSON.stringify({ #[test] fn wasm_parent_services_queued_python_requests_without_consuming_child_output() { + assert_wasm_parent_python_service_progress(8); + } + + #[test] + fn wasm_parent_python_services_rearm_after_coalesced_notification() { + assert_wasm_parent_python_service_progress(1); + } + + fn assert_wasm_parent_python_service_progress(child_quantum: usize) { let mut sidecar = create_test_sidecar(); + sidecar + .config + .runtime + .fairness + .capability_quantum_operations = child_quantum; let (connection_id, session_id) = authenticate_and_open_session(&mut sidecar).expect("authenticate sidecar"); let vm_id = create_vm( @@ -25645,6 +25659,8 @@ console.log(JSON.stringify({ GuestRuntimeKind::Python, ActiveExecution::Python(execution), ); + let notify = Arc::clone(&sidecar.process_event_notify); + child = child.with_event_notify(Arc::clone(¬ify)); // The WASM parent's poller requeues these events for the supervisor. // Neither request may be stranded behind the WASM-parent guard. child @@ -25710,23 +25726,32 @@ console.log(JSON.stringify({ .active_processes .insert(String::from("wasm-root"), root); - let mut javascript_services = Vec::new(); let mut python_services = Vec::new(); let mut socket_completions = Vec::new(); - let mut child_bridge_services = Vec::new(); - // One claim per pump also checks that the second queued request - // makes progress on the next supervisor turn. - for _ in 0..2 { - sidecar - .pump_child_process_events_nowait( - &vm_id, - &mut javascript_services, - &mut python_services, - &mut socket_completions, - &mut child_bridge_services, - 2, - ) + let ownership = OwnershipScope::vm(&connection_id, &session_id, &vm_id); + let mut cx = std::task::Context::from_waker(std::task::Waker::noop()); + // All producer notifications have coalesced. Each subsequent turn + // must be justified by a continuation, not a manual second call. + let turns = if child_quantum == 1 { 2 } else { 1 }; + for _ in 0..turns { + assert!( + std::future::Future::poll(Box::pin(notify.notified()).as_mut(), &mut cx) + .is_ready(), + "queued Python requests need a wake-up before each turn" + ); + let turn = sidecar + .pump_process_events_nowait(&ownership, 8) .expect("pump Python child requests"); + assert_eq!( + turn.python_services.len() + turn.python_socket_completions.len(), + if child_quantum == 1 { 1 } else { 2 }, + "claim queued services up to the fairness limit" + ); + assert!(!turn.emitted_any, "internal claims are not public output"); + assert!(turn.javascript_services.is_empty()); + assert!(turn.child_bridge_services.is_empty()); + python_services.extend(turn.python_services); + socket_completions.extend(turn.python_socket_completions); } assert_eq!( python_services.len(), @@ -25744,8 +25769,20 @@ console.log(JSON.stringify({ "supervisor must claim Python socket completion" ); assert_eq!(socket_completions[0].completion.request_id, 2); - assert!(javascript_services.is_empty()); - assert!(child_bridge_services.is_empty()); + // Drain a possible final fairness continuation. Shell-owned output + // must neither be consumed nor cause the supervisor to hot-spin. + let _ = std::future::Future::poll(Box::pin(notify.notified()).as_mut(), &mut cx); + let empty = sidecar + .pump_process_events_nowait(&ownership, 8) + .expect("shell-output-only follow-up"); + assert!(!empty.emitted_any); + assert!(empty.python_services.is_empty()); + assert!(empty.python_socket_completions.is_empty()); + assert!( + std::future::Future::poll(Box::pin(notify.notified()).as_mut(), &mut cx) + .is_pending(), + "shell-owned output must not rearm the supervisor" + ); assert!(sidecar.pending_process_events.is_empty()); let vm = sidecar.vms.get(&vm_id).expect("test vm"); let child = &vm.active_processes["wasm-root"].child_processes["python-child"]; From 6cb3b41c299b2df2f017c56ca16063a0452e1fe0 Mon Sep 17 00:00:00 2001 From: Anks Metaforms Date: Wed, 30 Sep 2026 16:29:29 +0530 Subject: [PATCH 3/3] test(native-sidecar): cover mixed child service wakeups --- crates/native-sidecar/tests/service.rs | 65 +++++++++++++++++++++++--- 1 file changed, 58 insertions(+), 7 deletions(-) diff --git a/crates/native-sidecar/tests/service.rs b/crates/native-sidecar/tests/service.rs index 5e37f2511b..5b6c4397f4 100644 --- a/crates/native-sidecar/tests/service.rs +++ b/crates/native-sidecar/tests/service.rs @@ -25587,15 +25587,24 @@ console.log(JSON.stringify({ #[test] fn wasm_parent_services_queued_python_requests_without_consuming_child_output() { - assert_wasm_parent_python_service_progress(8); + assert_wasm_parent_python_service_progress(8, false); } #[test] fn wasm_parent_python_services_rearm_after_coalesced_notification() { - assert_wasm_parent_python_service_progress(1); + assert_wasm_parent_python_service_progress(1, false); } - fn assert_wasm_parent_python_service_progress(child_quantum: usize) { + #[test] + fn wasm_parent_mixed_javascript_python_services_make_progress() { + assert_wasm_parent_python_service_progress(8, true); + assert_wasm_parent_python_service_progress(1, true); + } + + fn assert_wasm_parent_python_service_progress( + child_quantum: usize, + include_javascript: bool, + ) { let mut sidecar = create_test_sidecar(); sidecar .config @@ -25661,6 +25670,26 @@ console.log(JSON.stringify({ ); let notify = Arc::clone(&sidecar.process_event_notify); child = child.with_event_notify(Arc::clone(¬ify)); + if include_javascript { + // This method has an inline handler in the compatibility poll + // path. The supervisor must instead claim it as an owned + // service, account for it, and continue to the Python request. + child + .queue_pending_execution_event(ActiveExecutionEvent::JavascriptSyncRpcRequest( + JavascriptSyncRpcRequest { + id: 3, + method: String::from("process.signal_state"), + args: vec![ + Value::from(libc::SIGTERM), + Value::from("ignore"), + Value::from("[]"), + Value::from(0), + ], + raw_bytes_args: Default::default(), + }, + )) + .expect("queue JavaScript request before Python requests"); + } // The WASM parent's poller requeues these events for the supervisor. // Neither request may be stranded behind the WASM-parent guard. child @@ -25726,13 +25755,19 @@ console.log(JSON.stringify({ .active_processes .insert(String::from("wasm-root"), root); + let mut javascript_services = Vec::new(); let mut python_services = Vec::new(); let mut socket_completions = Vec::new(); let ownership = OwnershipScope::vm(&connection_id, &session_id, &vm_id); let mut cx = std::task::Context::from_waker(std::task::Waker::noop()); // All producer notifications have coalesced. Each subsequent turn // must be justified by a continuation, not a manual second call. - let turns = if child_quantum == 1 { 2 } else { 1 }; + let total_services = 2 + usize::from(include_javascript); + let turns = if child_quantum == 1 { + total_services + } else { + 1 + }; for _ in 0..turns { assert!( std::future::Future::poll(Box::pin(notify.notified()).as_mut(), &mut cx) @@ -25743,16 +25778,31 @@ console.log(JSON.stringify({ .pump_process_events_nowait(&ownership, 8) .expect("pump Python child requests"); assert_eq!( - turn.python_services.len() + turn.python_socket_completions.len(), - if child_quantum == 1 { 1 } else { 2 }, + turn.javascript_services.len() + + turn.python_services.len() + + turn.python_socket_completions.len(), + if child_quantum == 1 { + 1 + } else { + total_services + }, "claim queued services up to the fairness limit" ); assert!(!turn.emitted_any, "internal claims are not public output"); - assert!(turn.javascript_services.is_empty()); assert!(turn.child_bridge_services.is_empty()); + javascript_services.extend(turn.javascript_services); python_services.extend(turn.python_services); socket_completions.extend(turn.python_socket_completions); } + assert_eq!(javascript_services.len(), usize::from(include_javascript)); + if include_javascript { + assert_eq!(javascript_services[0].request.id, 3); + assert_eq!( + javascript_services[0].request.method, + "process.signal_state" + ); + assert_eq!(javascript_services[0].child_path, ["python-child"]); + } assert_eq!( python_services.len(), 1, @@ -25776,6 +25826,7 @@ console.log(JSON.stringify({ .pump_process_events_nowait(&ownership, 8) .expect("shell-output-only follow-up"); assert!(!empty.emitted_any); + assert!(empty.javascript_services.is_empty()); assert!(empty.python_services.is_empty()); assert!(empty.python_socket_completions.is_empty()); assert!(