diff --git a/frontend/src-tauri/src/extension_container.rs b/frontend/src-tauri/src/extension_container.rs index 3024e02..4e4473d 100644 --- a/frontend/src-tauri/src/extension_container.rs +++ b/frontend/src-tauri/src/extension_container.rs @@ -690,11 +690,11 @@ mod tests { real_native_protocol_probes(true, false); } #[test] - #[ignore = "real AppContainer MCP descendant CPU pressure; run explicitly"] - fn real_mcp_cpu_pressure_reports_resource_error_and_reaps_tree() { + #[ignore = "real AppContainer MCP CPU/memory/process exhaustion; run explicitly"] + fn real_mcp_resource_pressure_reports_error_and_reaps_tree() { real_native_protocol_probes(false, true); } - fn real_native_protocol_probes(_mcp_deadline: bool, cpu: bool) { + fn real_native_protocol_probes(_mcp_deadline: bool, resources: bool) { use std::{ net::{TcpListener, UdpSocket}, os::windows::fs::OpenOptionsExt, @@ -971,8 +971,8 @@ mod tests { if _mcp_deadline { mcp_modes.push("mcp_deadline"); } - if cpu { - mcp_modes.push("mcp_cpu"); + if resources { + mcp_modes.extend(["mcp_cpu", "mcp_memory", "mcp_processes"]); } for mode in mcp_modes { use std::sync::{ @@ -1092,27 +1092,36 @@ mod tests { ); assert!(session.take_tools_changed()); } - "mcp_cpu" => { - assert_eq!(result.unwrap_err().code, "EXTENSION_RESOURCE_CPU_EXCEEDED"); + "mcp_cpu" | "mcp_memory" | "mcp_processes" => { + let expected = match mode { + "mcp_memory" => "EXTENSION_RESOURCE_MEMORY_EXCEEDED", + "mcp_processes" => "EXTENSION_RESOURCE_PROCESSES_EXCEEDED", + _ => "EXTENSION_RESOURCE_CPU_EXCEEDED", + }; + assert_eq!(result.unwrap_err().code, expected); let elapsed = tool_started.elapsed(); - assert!(elapsed < std::time::Duration::from_secs(20)); - eprintln!("AppContainer MCP descendant CPU error: {elapsed:?}"); + assert!( + elapsed + < std::time::Duration::from_secs(if mode == "mcp_cpu" { + 20 + } else { + 10 + }) + ); + eprintln!("AppContainer {mode} resource error: {elapsed:?}"); assert_eq!( session.refresh_tools(&cancel).err().unwrap().code, "EXTENSION_MCP_SESSION_FAILED" ); - assert_eq!( - running.check_authorization().unwrap_err().code, - "EXTENSION_RESOURCE_CPU_EXCEEDED" - ); + assert_eq!(running.check_authorization().unwrap_err().code, expected); let vault = tempfile::tempdir().unwrap(); let mut workspace = crate::workspace::Workspace::open(vault.path()).unwrap(); workspace .write( - "cpu-recovery.md", + "resource-recovery.md", "", - b"Host saved after descendant CPU termination", + b"Host saved after resource termination", "local", ) .unwrap(); @@ -1120,10 +1129,10 @@ mod tests { assert_eq!( crate::workspace::Workspace::open(vault.path()) .unwrap() - .read("cpu-recovery.md") + .read("resource-recovery.md") .unwrap() .content, - "Host saved after descendant CPU termination" + "Host saved after resource termination" ); } "mcp_deadline" => { diff --git a/frontend/src-tauri/src/extension_instance.rs b/frontend/src-tauri/src/extension_instance.rs index 2aacdde..548b4e0 100644 --- a/frontend/src-tauri/src/extension_instance.rs +++ b/frontend/src-tauri/src/extension_instance.rs @@ -514,11 +514,11 @@ mod tests { native_worker_lifecycle(false); } #[test] - #[ignore = "real background MCP descendant CPU exhaustion and restart; run explicitly"] - fn native_cpu_failure_is_reported_reaped_and_replacement_can_start() { + #[ignore = "real background MCP CPU/memory/process exhaustion and restart; run explicitly"] + fn native_resource_failures_are_reaped_and_replacements_can_start() { native_worker_lifecycle(true); } - fn native_worker_lifecycle(cpu: bool) { + fn native_worker_lifecycle(resources: bool) { let temp = tempfile::tempdir().unwrap(); let package = temp.path().join("package"); std::fs::create_dir(&package).unwrap(); @@ -732,8 +732,16 @@ mod tests { .unwrap(), 0 ); - if cpu { - let exhausted = unsafe { registry.start(make("mcp_cpu")) }.unwrap(); + for (mode, expected) in if resources { + vec![ + ("mcp_cpu", "EXTENSION_RESOURCE_CPU_EXCEEDED"), + ("mcp_memory", "EXTENSION_RESOURCE_MEMORY_EXCEEDED"), + ("mcp_processes", "EXTENSION_RESOURCE_PROCESSES_EXCEEDED"), + ] + } else { + Vec::new() + } { + let exhausted = unsafe { registry.start(make(mode)) }.unwrap(); wait_for(|| exhausted.snapshot().status != Status::Starting); assert_eq!(exhausted.snapshot().status, Status::Ready); let old_identity = @@ -748,21 +756,18 @@ mod tests { exhausted .invoke_confirmed(review.review_id) .unwrap() - .wait(Duration::from_secs(20)) + .wait(Duration::from_secs(if mode == "mcp_cpu" { 20 } else { 10 })) .unwrap_err() .code, - "EXTENSION_RESOURCE_CPU_EXCEEDED" + expected ); - eprintln!("background CPU call error after {:?}", started.elapsed()); + eprintln!("background {mode} call error after {:?}", started.elapsed()); wait_for(|| { registry.reap(); registry.entries.is_empty() }); assert_eq!(exhausted.snapshot().status, Status::Failed); - assert_eq!( - exhausted.snapshot().error.as_deref(), - Some("EXTENSION_RESOURCE_CPU_EXCEEDED") - ); + assert_eq!(exhausted.snapshot().error.as_deref(), Some(expected)); assert_eq!(exhausted.snapshot().tool_count, 0); assert!(exhausted.review("echo".into(), json!({})).is_err()); assert_eq!( diff --git a/frontend/src-tauri/src/extension_job.rs b/frontend/src-tauri/src/extension_job.rs index bca2936..7f4b56d 100644 --- a/frontend/src-tauri/src/extension_job.rs +++ b/frontend/src-tauri/src/extension_job.rs @@ -12,13 +12,13 @@ use windows_sys::Win32::System::{ pub struct Job { handle: OwnedHandle, state: std::sync::Arc, - _cpu: Option, + _monitor: Option, } impl Job { pub(crate) fn clone_for_deadline(&self) -> Result { Ok(Self { state: std::sync::Arc::clone(&self.state), - _cpu: None, + _monitor: None, handle: self .handle .try_clone() @@ -36,7 +36,7 @@ impl Job { let mut job = Self { handle: unsafe { OwnedHandle::from_raw_handle(raw) }, state: std::sync::Arc::new(std::sync::atomic::AtomicU8::new(0)), - _cpu: None, + _monitor: None, }; let limits = JOBOBJECT_EXTENDED_LIMIT_INFORMATION { BasicLimitInformation: JOBOBJECT_BASIC_LIMIT_INFORMATION { @@ -63,7 +63,7 @@ impl Job { }, }; job.set(JobObjectCpuRateControlInformation, &cpu)?; - job._cpu = Some(CpuMonitor::arm(&job)?); + job._monitor = Some(ResourceMonitor::arm(&job)?); Ok(job) } pub fn check_resources(&self) -> Result<()> { @@ -72,6 +72,8 @@ impl Job { 0 => Ok(()), 1 => Err(HostError::new("EXTENSION_RESOURCE_CPU_EXCEEDED")), 3 => Err(HostError::new("EXTENSION_RESOURCE_TERMINATE_FAILED")), + 4 => Err(HostError::new("EXTENSION_RESOURCE_MEMORY_EXCEEDED")), + 5 => Err(HostError::new("EXTENSION_RESOURCE_PROCESSES_EXCEEDED")), _ => Err(HostError::new("EXTENSION_RESOURCE_MONITOR_FAILED")), } } @@ -130,12 +132,14 @@ impl Job { /// The Windows notification uses a ten-second window and ToleranceHigh (60% /// over budget). This is not a measurement of ten uninterrupted busy seconds. /// Only the original Job owns this monitor; observation/deadline clones do not. +const JOB_MEMORY_LIMIT: u32 = 10; // JOB_OBJECT_MSG_JOB_MEMORY_LIMIT +const JOB_PROCESS_LIMIT: u32 = 3; // JOB_OBJECT_MSG_ACTIVE_PROCESS_LIMIT const JOB_NOTIFICATION_LIMIT: u32 = 11; // JOB_OBJECT_MSG_NOTIFICATION_LIMIT (Windows SDK) -struct CpuMonitor { +struct ResourceMonitor { stop: std::sync::Arc, worker: Option>, } -impl CpuMonitor { +impl ResourceMonitor { fn arm(job: &Job) -> Result { use std::sync::{ atomic::{AtomicBool, Ordering}, @@ -172,7 +176,7 @@ impl CpuMonitor { let stop = Arc::new(AtomicBool::new(false)); let thread_stop = Arc::clone(&stop); let worker = std::thread::Builder::new() - .name("extension-cpu".into()) + .name("extension-resources".into()) .spawn(move || { // All exits, including a caught panic or completion-port failure, // terminate the tree while this worker still owns a Job handle. @@ -199,6 +203,14 @@ impl CpuMonitor { if key != 1 { return 2; } + // These hard-limit notifications are best effort on Windows; + // the kernel still enforces the configured allocation caps. + if code == JOB_MEMORY_LIMIT { + return 4; + } + if code == JOB_PROCESS_LIMIT { + return 5; + } if code != JOB_NOTIFICATION_LIMIT { continue; } @@ -234,7 +246,7 @@ impl CpuMonitor { }) } } -impl Drop for CpuMonitor { +impl Drop for ResourceMonitor { fn drop(&mut self) { self.stop.store(true, std::sync::atomic::Ordering::Release); if let Some(worker) = self.worker.take() { @@ -449,16 +461,48 @@ mod tests { job.assign_suspended(child.process.as_handle()).unwrap(); assert_ne!(ResumeThread(child.thread.as_raw_handle()), u32::MAX); } - assert_eq!( - unsafe { WaitForSingleObject(child.process.as_raw_handle(), 10000) }, - WAIT_OBJECT_0 - ); - let mut code = 1; - assert_ne!( - unsafe { GetExitCodeProcess(child.process.as_raw_handle(), &mut code) }, - 0 - ); - assert_eq!(code, 0); + assert_resource_cleanup(&job, "EXTENSION_RESOURCE_MEMORY_EXCEEDED"); + } + fn assert_resource_cleanup(job: &Job, expected: &str) { + let started = std::time::Instant::now(); + while job.check_resources().is_ok() && started.elapsed() < Duration::from_secs(5) { + std::thread::sleep(Duration::from_millis(5)); + } + assert_eq!(job.check_resources().unwrap_err().code, expected); + while job.active_processes().unwrap() != 0 && started.elapsed() < Duration::from_secs(5) { + std::thread::sleep(Duration::from_millis(5)); + } + assert_eq!(job.active_processes().unwrap(), 0); + } + #[test] + #[ignore = "helper attempts to exceed production child process budget"] + fn worker_processes() { + let mut children = Vec::new(); + for _ in 0..24 { + match std::process::Command::new(std::env::current_exe().unwrap()) + .creation_flags(CREATE_NO_WINDOW) + .args(["--ignored", "--exact", "extension_job::tests::worker_wait"]) + .spawn() + { + Ok(child) => children.push(child), + Err(_) => break, + } + } + std::thread::sleep(Duration::from_secs(30)); + for mut child in children { + let _ = child.kill(); + let _ = child.wait(); + } + } + #[test] + fn descendant_process_exhaustion_terminates_production_job() { + let job = Job::new().unwrap(); + let child = worker_named("worker_processes"); + unsafe { + job.assign_suspended(child.process.as_handle()).unwrap(); + assert_ne!(ResumeThread(child.thread.as_raw_handle()), u32::MAX); + } + assert_resource_cleanup(&job, "EXTENSION_RESOURCE_PROCESSES_EXCEEDED"); } #[test] fn managed_descendants_terminate_with_their_job() { @@ -494,18 +538,10 @@ mod tests { assert_eq!(job.active_processes().unwrap(), 1); let second = worker(); assert!(unsafe { job.assign_suspended(second.process.as_handle()) }.is_err()); - // The second process has never been resumed, even after assignment failure. + // Both processes were still suspended. The attempted limit violation + // now revokes the entire first job, rather than leaving it runnable. drop(second); - assert_ne!( - unsafe { ResumeThread(first.thread.as_raw_handle()) }, - u32::MAX - ); - let deadline = std::time::Instant::now() + Duration::from_secs(5); - while job.active_processes().unwrap() != 1 && std::time::Instant::now() < deadline { - std::thread::sleep(Duration::from_millis(10)); - } - assert_eq!(job.active_processes().unwrap(), 1); - drop(job); + assert_resource_cleanup(&job, "EXTENSION_RESOURCE_PROCESSES_EXCEEDED"); assert_eq!( unsafe { WaitForSingleObject(first.process.as_raw_handle(), 5000) }, WAIT_OBJECT_0 diff --git a/frontend/src-tauri/tests/fixtures/sandbox_network_probe.rs b/frontend/src-tauri/tests/fixtures/sandbox_network_probe.rs index 24ac507..6ac10e1 100644 --- a/frontend/src-tauri/tests/fixtures/sandbox_network_probe.rs +++ b/frontend/src-tauri/tests/fixtures/sandbox_network_probe.rs @@ -41,6 +41,24 @@ fn main() { } let mut request = read(&mut input); assert!(request.contains("tools/call")); + if args[1] == "mcp_memory" { + let mut excessive = Vec::::new(); + let _ = excessive.try_reserve_exact(600 * 1024 * 1024); + std::thread::sleep(Duration::from_secs(120)); + return; + } + if args[1] == "mcp_processes" { + let mut children = Vec::new(); + for _ in 0..24 { + match std::process::Command::new(std::env::current_exe().unwrap()).arg("wait").spawn() { + Ok(child) => children.push(child), + Err(_) => break, + } + } + std::thread::sleep(Duration::from_secs(120)); + for mut child in children { let _ = child.kill(); let _ = child.wait(); } + return; + } if args[1] == "mcp_cancel" || args[1] == "mcp_deadline" || args[1] == "mcp_cpu" { let _child = std::process::Command::new(std::env::current_exe().unwrap()).arg(if args[1] == "mcp_cpu" { "cpu_burn" } else { "wait" }).spawn().unwrap(); std::thread::sleep(Duration::from_secs(120));