diff --git a/docs/development/OpenNexus生产化实施进度-2026-09-08.md b/docs/development/OpenNexus生产化实施进度-2026-09-08.md index 3b00a4f..6df8731 100644 --- a/docs/development/OpenNexus生产化实施进度-2026-09-08.md +++ b/docs/development/OpenNexus生产化实施进度-2026-09-08.md @@ -564,3 +564,13 @@ Core 的独立数据目录目前不等于已授权 Vault。Python 旧笔记写 - 这些证据补齐了 CPU 故障通过原生会话和后台注册表传播的路径,不能将 9.77/12.80 秒的结果改称严格连续 10 秒的计时证明;窗口语义差异仍保留。实际产品启用、资源/网络/文件攻击全矩阵和完整 C-04 验收仍未完成。 - 默认扩展回归 68 通过、9 ignored,新增两项显式长时测试已单独执行通过,其余 ignored 仍有其他长时验收及辅助进程入口;desktop 全目标 Clippy -D warnings 通过。日志 `.build/extension-cpu-integration-regression.log`、`.build/extension-cpu-integration-clippy.log`。 - 同轮只读复查测试 Sync:HTTP /health=200,/ready=503,响应错误 DEPENDENCY_UNAVAILABLE;SSH 22 可以建立 TCP,但在认证前读到空 banner,连接被对端关闭。未尝试写入部署,也未将 /health 成功当作服务就绪。证据 `.build/sync-cpu-integration-readiness.json`;该远端障碍不阻止继续本地生产化工作。 + + +## 增量:内存与进程数超限的错误传播和整组清理 + +- 将 CPU 专用监视器扩展为 ResourceMonitor,继续保留 512 MiB Job 内存、16 个活动进程及单核等效 CPU 硬限额。在独立 completion port 上处理 JOB_OBJECT_MSG_JOB_MEMORY_LIMIT 与 JOB_OBJECT_MSG_ACTIVE_PROCESS_LIMIT,分别发布 EXTENSION_RESOURCE_MEMORY_EXCEEDED / EXTENSION_RESOURCE_PROCESSES_EXCEEDED,随后终止整个 Job,沿用原有监视器失败/终止失败错误链。 +- 原先第二个挂起进程超过测试限额时只拒绝分配,现收到超限通知后也撤销已有 Job;相应原生测试改为核验明确错误与已存在进程全部退出。内存测试仍在受限子进程尝试正常 32 MiB 分配及超预算 600 MiB 分配,父进程改验资源错误与 Job 清空;新增子进程连续创建后代耗尽生产 16 进程预算的实际测试。该组 5 项通过,6 ignored 为长时项/父测试辅助进程,日志 `.build/extension-memory-process-limits.log`。 +- 依据 [Windows completion port 文档](https://learn.microsoft.com/en-us/windows/win32/api/winnt/ns-winnt-jobobject_associate_completion_port),这两类硬限额通知不保证投递;内核分配限额仍然执行,但当前事件驱动清理不能被宣称为丢通知时也有严格清理时限。没有将通知当作保证,也未据此启用第三方执行或关闭 C-04 缺口。 +- 原生 MCP fixture 新增 mcp_memory(尝试预留 600 MiB)与 mcp_processes(连续启动同包受管后代)模式。原 AppContainer 和后台注册表显式测试扩展为 CPU/内存/进程三类:核对调用错误、失败会话、资源状态、Job 清空、Host 保存;后台额外核对旧 Endpoint、工具目录清理、包 ACL 基线和每次失败后同包新实例实际调用成功。内存/进程回执限定 10 秒,CPU 保留既有 20 秒观察上限且仍不作为连续 10 秒计时证明。 +- 初次后台三类测试通过,CPU/内存/进程调用到错误分别约 12.339 秒、2.509 毫秒、395.929 毫秒,整项 16.34 秒;AppContainer 协议组三类也通过,整项 30.59 秒。日志 `.build/extension-resource-registry.log`、`.build/extension-resource-container.log`。后续以收紧后的内存/进程 10 秒测试边界重新执行最终版本,结果另记。 +- 最终版本两项显式原生/后台资源集成测试同时执行通过,合计 26.42 秒,日志 `.build/extension-resource-explicit-final.log`;内存错误约 1.0–1.3 毫秒,进程数错误约 378–445 毫秒。默认扩展回归 69 通过、10 ignored,日志 `.build/extension-resource-regression.log`;desktop 全目标 Clippy -D warnings 通过,日志 `.build/extension-resource-clippy.log`。本轮没有重跑全部前端、后端或远端部署测试,也未声称完整生产化完成。 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));