diff --git a/crates/native-sidecar/src/execution/child_process.rs b/crates/native-sidecar/src/execution/child_process.rs index a4a279d8ef..4590840393 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, + // Includes both staged services and completions settled in place. Neither + // needs a parent-facing event, but both consume the caller's work budget. + internal_work: usize, } impl OwnedChildBridgeEventService { @@ -2718,7 +2721,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() @@ -2772,77 +2775,43 @@ 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 + // and exit delivery single-owner. JavaScript process-control + // RPCs and Python bridge work still need the supervisor; + // ordinary JavaScript RPCs remain owned by the pull lane. + 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(_), + ) => true, + Some(_) => false, + // JavaScript can expose a request directly from + // its runtime before the parent has polled it. + None => child.runtime == GuestRuntimeKind::JavaScript, + } + }), + ) }) - .unwrap_or((false, 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 - { - 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; - } - } + .unwrap_or((false, false)); + if parent_is_pull_driven_wasm && !needs_supervisor { continue; } self.expire_child_process_sync_if_needed( @@ -2857,7 +2826,7 @@ where process_id, &parent_path, &child_process_id, - false, + parent_is_pull_driven_wasm, javascript_services, python_services, python_socket_completions, @@ -2867,7 +2836,17 @@ 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, + internal_work, + } = claimed; + if internal_work > 0 { + emitted_any = true; + progressed_this_round = true; + work += internal_work; + child_work[candidate_index] += internal_work; + } if event.is_null() { continue; } @@ -2884,11 +2863,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; } } @@ -3662,9 +3641,20 @@ where } Err(error) => return Err(error), }; - let ClaimedDescendantBridgeEvent { event, reservation } = claimed; + let ClaimedDescendantBridgeEvent { + event, + reservation, + internal_work, + } = claimed; + work += internal_work; + detached_work += internal_work; let Some(event_type) = event.get("type").and_then(Value::as_str) else { + if internal_work > 0 { + // Continue draining internal work, or rearm at the + // capacity/fairness checks at the top of the loop. + continue; + } break; }; let public_event_queued = event @@ -8078,6 +8068,7 @@ where None, None, None, + None, )) .await .map(Into::into) @@ -8589,6 +8580,7 @@ where owned_python_socket_completions: &mut Vec, public_process_id: Option<&str>, ) -> Result { + let mut internal_work = 0; let mut claimed_reservation = None; let mut future = Box::pin(self.poll_descendant_javascript_child_process( vm_id, @@ -8597,6 +8589,7 @@ where child_process_id, 0, preserve_pull_owned_events, + Some(&mut internal_work), Some(owned_javascript_services), Some(owned_python_services), Some(owned_python_socket_completions), @@ -8610,6 +8603,7 @@ where Poll::Ready(result) => result.map(|event| ClaimedDescendantBridgeEvent { event, reservation: claimed_reservation, + internal_work, }), Poll::Pending => Err(SidecarError::InvalidState(String::from( "ERR_AGENTOS_CHILD_EVENT_TURN_SUSPENDED: bounded child event claim unexpectedly suspended", @@ -8628,6 +8622,7 @@ where child_process_id: &str, wait_ms: u64, preserve_pull_owned_events: bool, + mut claimed_internal_work: Option<&mut usize>, mut owned_javascript_services: Option<&mut Vec>, mut owned_python_services: Option<&mut Vec>, mut owned_python_socket_completions: Option<&mut Vec>, @@ -8643,6 +8638,15 @@ where let _ = wait_ms; loop { + // A supervised poll may consume one internal event before yielding + // to the caller's VM/child budget checks. In-place completions must + // not drain an unbounded queue behind a single counted service. + if claimed_internal_work + .as_deref() + .is_some_and(|work| *work > 0) + { + return Ok(Value::Null); + } self.drain_queued_descendant_javascript_child_process_events( vm_id, process_id, @@ -8734,6 +8738,7 @@ where ActiveExecutionEvent::Stdout(_) | ActiveExecutionEvent::Stderr(_) | ActiveExecutionEvent::Exited(_) + | ActiveExecutionEvent::SignalState { .. } ) { let Some(mut vm) = self.vms.get_mut(vm_id) else { @@ -8747,10 +8752,8 @@ where let Some(child) = parent.child_processes.get_mut(child_process_id) else { return Ok(Value::Null); }; - child.queue_pending_polled_execution_event(PolledExecutionEvent { - event, - reservation, - })?; + child + .requeue_pending_execution_event(PolledExecutionEvent { event, reservation })?; return Ok(Value::Null); } match event { @@ -9112,6 +9115,9 @@ where })?; return Ok(Value::Null); } + if let Some(work) = claimed_internal_work.as_deref_mut() { + *work += 1; + } if let Some(services) = owned_javascript_services.as_deref_mut() { let Some((connection_id, session_id, vm)) = self.vms.get(vm_id).and_then(|state| { @@ -9510,6 +9516,9 @@ where } } ActiveExecutionEvent::PythonVfsRpcRequest(request) => { + if let Some(work) = claimed_internal_work.as_deref_mut() { + *work += 1; + } let Some(services) = owned_python_services.as_deref_mut() else { return Err(SidecarError::InvalidState(String::from( "ERR_AGENTOS_PYTHON_EVENT_OWNERSHIP: nested Python VFS RPC must be claimed by the owned process-event supervisor", @@ -9558,6 +9567,9 @@ where return Ok(Value::Null); } ActiveExecutionEvent::PythonSocketConnectCompletion(completion) => { + if let Some(work) = claimed_internal_work.as_deref_mut() { + *work += 1; + } let Some(services) = owned_python_socket_completions.as_deref_mut() else { return Err(SidecarError::InvalidState(String::from( "ERR_AGENTOS_PYTHON_SOCKET_COMPLETION_OWNERSHIP: nested completion must be claimed by the owned process-event supervisor", @@ -9606,6 +9618,9 @@ where return Ok(Value::Null); } ActiveExecutionEvent::JavascriptSyncRpcCompletion(completion) => { + if let Some(work) = claimed_internal_work.as_deref_mut() { + *work += 1; + } drop(reservation); let Some(mut vm) = self.vms.get_mut(vm_id) else { return Ok(Value::Null); @@ -10104,6 +10119,7 @@ where None, None, None, + None, ) .await } @@ -11363,6 +11379,442 @@ mod child_event_claim_tests { .await; } + #[tokio::test(flavor = "current_thread")] + async fn supervised_child_poll_preserves_pull_owned_events_without_counting_work() { + tokio::task::LocalSet::new() + .run_until(async { + let (mut sidecar, vm_id) = + sidecar_with_test_vm(agentos_runtime::DEFAULT_PROTOCOL_MAX_PROCESS_EVENTS) + .await; + let notify = Arc::clone(&sidecar.process_event_notify); + let root_id = String::from("pull-owned-root"); + let child_id = String::from("pull-owned-child"); + { + let mut vm = sidecar.vms.get_mut(&vm_id).expect("pull-owned VM"); + let mut root = host_function_process(&mut vm, "pull-owned parent", None); + root.runtime = GuestRuntimeKind::WebAssembly; + let mut child = + host_function_process(&mut vm, "pull-owned child", Some(root.kernel_pid)) + .with_event_notify(Arc::clone(¬ify)); + let (signal, registration) = parse_process_signal_state_request( + &signal_state_rpc(1, libc::SIGTERM as u32, "ignore").args, + ) + .expect("signal-state fixture"); + for event in [ + rpc(1), + ActiveExecutionEvent::SignalState { + signal, + registration, + }, + ActiveExecutionEvent::Stdout(b"stdout".to_vec()), + ActiveExecutionEvent::Stderr(b"stderr".to_vec()), + ActiveExecutionEvent::Exited(0), + ] { + child + .queue_pending_execution_event(event) + .expect("queue pull-owned event"); + } + root.child_processes.insert(child_id.clone(), child); + vm.active_processes.insert(root_id.clone(), root); + } + consume_notify_permit(¬ify).await; + for remaining in (1..=5).rev() { + let mut javascript = Vec::new(); + let claimed = sidecar + .poll_descendant_javascript_child_process_nowait( + &vm_id, + &root_id, + &[], + &child_id, + true, + &mut javascript, + &mut Vec::new(), + &mut Vec::new(), + None, + ) + .expect("probe pull-owned event"); + assert!(claimed.event.is_null()); + assert_eq!( + claimed.internal_work, 0, + "requeued events are not consumed work" + ); + assert!(javascript.is_empty()); + // Requeueing restores one producer notification. The normal + // pump must skip this parent-owned queue on that wake and + // quiesce, rather than repeatedly probe and requeue it. + consume_notify_permit(¬ify).await; + let turn = sidecar + .pump_process_events_nowait( + &OwnershipScope::vm( + "child-claim-connection", + "child-claim-session", + &vm_id, + ), + 8, + ) + .expect("parent-owned queue follow-up"); + assert!(!turn.emitted_any); + assert!(turn.javascript_services.is_empty()); + assert!(turn.child_bridge_services.is_empty()); + assert!( + tokio::time::timeout(Duration::from_millis(10), notify.notified()) + .await + .is_err(), + "parent-owned queues must not hot-rearm the supervisor" + ); + let mut vm = sidecar.vms.get_mut(&vm_id).expect("pull-owned VM"); + let child = vm + .active_processes + .get_mut(&root_id) + .unwrap() + .child_processes + .get_mut(&child_id) + .unwrap(); + assert_eq!(child.pending_execution_events.len(), remaining); + // Simulate the parent consuming its next event. The + // supervisor's probe must not rotate it behind later output. + let event = child + .pop_pending_execution_event() + .expect("parent-owned event"); + assert!(match remaining { + 5 => matches!(event, ActiveExecutionEvent::JavascriptSyncRpcRequest(_)), + 4 => matches!(event, ActiveExecutionEvent::SignalState { .. }), + 3 => matches!(event, ActiveExecutionEvent::Stdout(_)), + 2 => matches!(event, ActiveExecutionEvent::Stderr(_)), + 1 => matches!(event, ActiveExecutionEvent::Exited(0)), + _ => unreachable!(), + }); + } + }) + .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( + ActiveExecutionEvent::JavascriptSyncRpcRequest(signal_state_rpc( + id, + libc::SIGTERM as u32, + "ignore", + )), + ) + .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; + } + + async fn assert_child_completions_make_bounded_progress(detached: bool) { + for wasm_parent in [false, true] { + for (vm_quantum, child_quantum, batch_size) in [(1, 8, 1), (8, 1, 1), (8, 8, 3)] { + 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("completion-root"); + let child_id = String::from("completion-child"); + let (mut root, mut child) = { + let mut vm = sidecar.vms.get_mut(&vm_id).expect("completion VM"); + let mut root = host_function_process(&mut vm, "completion parent", None); + if wasm_parent { + root.runtime = GuestRuntimeKind::WebAssembly; + } + let child = javascript_process( + &vm_id, + &mut vm, + "completion child", + r#" +const fs = require("node:fs/promises"); +Promise.allSettled(["/one", "/two", "/three"].map(path => fs.stat(path))); +setInterval(() => {}, 60_000); +"#, + Arc::clone(¬ify), + ); + (root, child) + }; + // Use real outstanding bridge calls so settling each completion + // exercises the production response path, not a mocked handler. + let mut request_ids = Vec::new(); + for _ in 0..3 { + let event = tokio::time::timeout( + Duration::from_secs(2), + child.execution.poll_event(Duration::from_secs(2)), + ) + .await + .expect("async filesystem request deadline") + .expect("poll async filesystem request") + .expect("async filesystem request"); + let ActiveExecutionEvent::JavascriptSyncRpcRequest(request) = event else { + panic!("expected async filesystem request, got {event:?}"); + }; + assert_eq!(request.method, "fs.promises.stat"); + request_ids.push(request.id); + } + for &request_id in &request_ids { + child + .queue_pending_execution_event( + ActiveExecutionEvent::JavascriptSyncRpcCompletion( + crate::state::JavascriptSyncRpcCompletion { + request_id, + result: Err(crate::state::DeferredRpcError { + code: String::from("ENOENT"), + message: String::from( + "completion fixture: file does not exist", + ), + }), + }, + ), + ) + .expect("queue completion"); + } + root.child_processes.insert(child_id.clone(), child); + { + let mut vm = sidecar.vms.get_mut(&vm_id).expect("completion VM"); + 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 remaining = request_ids; + while !remaining.is_empty() { + tokio::time::timeout(Duration::from_millis(50), notify.notified()) + .await + .expect("completion work must retain a continuation edge"); + let mut javascript = Vec::new(); + if detached { + sidecar + .pump_detached_child_process_events_nowait( + &vm_id, + &mut javascript, + &mut Vec::new(), + &mut Vec::new(), + &mut [], + 8, + ) + .expect("pump detached completions"); + } else { + let turn = sidecar + .pump_process_events_nowait(&ownership, 8) + .expect("pump attached completions"); + javascript = turn.javascript_services; + assert!(turn.child_bridge_services.is_empty()); + } + assert!( + javascript.is_empty(), + "completions do not stage owned services" + ); + remaining.drain(..batch_size.min(remaining.len())); + let vm = sidecar.vms.get(&vm_id).expect("completion VM"); + let child = &vm.active_processes[&root_id].child_processes[&child_id]; + let queued = child + .pending_execution_events + .iter() + .filter_map(|event| match event { + ActiveExecutionEvent::JavascriptSyncRpcCompletion(completion) => { + Some(completion.request_id) + } + _ => None, + }) + .collect::>(); + assert_eq!(queued, remaining, + "detached={detached}, wasm_parent={wasm_parent}, vm_quantum={vm_quantum}, child_quantum={child_quantum}"); + } + consume_notify_permit(¬ify).await; + if detached { + sidecar + .pump_detached_child_process_events_nowait( + &vm_id, + &mut Vec::new(), + &mut Vec::new(), + &mut Vec::new(), + &mut [], + 8, + ) + .expect("empty detached completion turn"); + } else { + sidecar + .pump_process_events_nowait(&ownership, 8) + .expect("empty attached completion turn"); + } + assert!( + tokio::time::timeout(Duration::from_millis(10), notify.notified()) + .await + .is_err(), + "empty completion turns must not hot-rearm" + ); + } + } + } + + #[tokio::test(flavor = "current_thread")] + async fn attached_completions_make_bounded_progress_after_coalesced_notify() { + tokio::task::LocalSet::new() + .run_until(assert_child_completions_make_bounded_progress(false)) + .await; + } + + #[tokio::test(flavor = "current_thread")] + async fn detached_completions_make_bounded_progress_after_coalesced_notify() { + tokio::task::LocalSet::new() + .run_until(assert_child_completions_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/tests/service.rs b/crates/native-sidecar/tests/service.rs index c1a64b22e9..db30c8aa4d 100644 --- a/crates/native-sidecar/tests/service.rs +++ b/crates/native-sidecar/tests/service.rs @@ -25594,6 +25594,266 @@ console.log(JSON.stringify({ } } + #[test] + fn wasm_parent_services_queued_python_requests_without_consuming_child_output() { + 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, false); + } + + #[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 + .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( + &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), + ); + 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 + .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 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 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) + .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.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.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, + "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); + // 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.javascript_services.is_empty()); + 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"]; + 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();