diff --git a/docs/development/OpenNexus生产化实施进度-2026-09-08.md b/docs/development/OpenNexus生产化实施进度-2026-09-08.md index 6495133..45261ee 100644 --- a/docs/development/OpenNexus生产化实施进度-2026-09-08.md +++ b/docs/development/OpenNexus生产化实施进度-2026-09-08.md @@ -519,3 +519,15 @@ Core 的独立数据目录目前不等于已授权 Vault。Python 旧笔记写 - 单元测试覆盖预览对象被改写后实际参数仍冻结、重复确认、跨实例、目录代次、确认前/后过期、条数/字节上限及清理后复用。真实 AppContainer 测试核对预览来源/包/Vault,修改展示副本不会改变实际请求,旧目录凭据被拒绝;新增确认后锁定保险库导致 CREDENTIALS_LOCKED,重新解锁也不能恢复已失效会话。 - 64 项扩展回归通过,5 项 ignored 为三个父测试驱动 Job 辅助入口及两个显式长时验收,本轮未重跑长时项。desktop 全目标 Clippy -D warnings 通过。日志 `.build/extension-call-binding-tests.log`、`.build/extension-call-binding-clippy.log`。 - 当前 confirm_call 仍只是 Host 内部组件接口,未登记为 renderer/Core 命令,也没有自动批准策略或真实用户确认 UI。完整产品入口必须通过经过认证的 Host 路由确认用户批准,再与活动安装/在线信任、实例注册和 Core 路由组合;全部沙箱资源与其他生产化验收仍未完成,第三方执行能力保持禁用。 + + +## 增量:后台实例注册、串行派发与停止清理 + +- 新增 extension_instance::Registry。启动先用当前时间核验执行许可,再按执行类型、来源/命名空间/包和 Vault 限制同键单实例;最多 16 个注册项。每个工作线程在自身栈中创建并持有 Profile、PinnedPackage、BoundEntry、Running 与 Session,没有把带裸 SID 的 Profile 强行标为 Send。 +- 工作线程完成包校验、AppContainer/Job 挂起创建后调用 before_resume 复核钩子,再执行带租约恢复和 MCP 初始化/目录发现。该 start API 为 unsafe 内部集成边界,调用者仍须建立完整资源、用户与活动安装授权;本轮没有将“钩子存在”当作在线信任或完整沙箱策略已接通,也未新增 renderer/Core 启动命令。 +- 每实例四槽请求队列,预览参数入队前即执行 256 KiB/结构限制;预览、确认并调用、目录刷新由实例线程串行处理。确认和一次性凭据消费合在同一线程,ApprovedCall 不通过 IPC 返回或另行长期排队。回执等待限于后台,取消/丢弃 Ticket 会通知排队或进行中的工作;超时不证明远端工具尚未产生副作用,后续 Core 路由仍需正确呈现未知完成状态。 +- stop 使用独立原子信号并取消当前请求,不会被满队列阻塞;后台循环持续处理空闲 MCP 通知和租约撤销。终态在原生栈所有者退出后发布,Profile 显式删除失败有稳定错误;Registry 只回收确认已完成的线程,清理失败或 worker panic 的注册项保留隔离,避免直接允许同键替换。Registry 析构会停止并 join 工作线程,因此其生命周期管理不能置于 UI 同步调用中。 +- 真实后台 AppContainer 集成测试覆盖启动/目录就绪、预览与确认调用、重复确认拒绝、停止、线程回收后以新实例 ID 重建、活动工具取消、启动前策略拒绝及保险库锁定自动退出。测试同时确认旧 Endpoint 不能再提交请求,并通过保留的仅测试 Job 观察句柄确认最终 ActiveProcesses=0。 +- 集成测试最初以固定 ActiveProcesses=2 判断工具子进程启动,本机在 Ready 时已计入 2 个进程,导致实际取消的是排队请求。已改为同时确认当前 Ticket 处于执行槽、且进程数超过该实例就绪基线,再触发取消;修正后得到进行中工具的 MCP_CANCELLED 并清空 Job。此前仅凭固定计数的证据不能单独用于识别某个具体后代进程。 +- 新增队列满/输入超限/停止后队列不执行,以及已结束但清理失败的注册项不回收测试。67 项扩展回归通过,5 项 ignored 为三个父测试驱动 Job 辅助入口及两个显式长时验收,本轮未重跑长时项。desktop 全目标 Clippy -D warnings 通过,首次检查的回调类型复杂度提示已通过类型别名修正。日志 `.build/extension-instance-tests.log`、`.build/extension-instance-regression.log`、`.build/extension-instance-clippy.log`。 +- 仍需将注册表纳入实际 Host 依赖与 UI/Core 的认证路由,完成活动安装/在线信任钩子、包对象权限回收/残留检查、全部 Host 创建路径继承审计、scratch/CPU 和完整攻击矩阵;其他同步/部署/安装包验收也继续推进。第三方执行能力保持禁用,完整生产化目标尚未完成。 diff --git a/frontend/src-tauri/src/extension_instance.rs b/frontend/src-tauri/src/extension_instance.rs new file mode 100644 index 0000000..ecddea6 --- /dev/null +++ b/frontend/src-tauri/src/extension_instance.rs @@ -0,0 +1,815 @@ +//! Native instance workers. Construct/drop this manager off the UI thread. +//! Production callers must still satisfy the complete launch policy contract. +use crate::{ + credentials::CredentialBroker, + extension_call_authorization::{Identity, Review}, + extension_container::Profile, + extension_launch_authorization::Context, + extension_mcp::Session, + extension_mcp_tools::Description, + extension_package::Inventory, + extension_permit::{Authority, Claims, Permit}, + extension_pinned::PinnedPackage, + workspace::{HostError, Result}, +}; +use serde::Serialize; +use serde_json::Value; +use sha2::{Digest, Sha256}; +use std::{ + collections::BTreeMap, + path::PathBuf, + sync::{ + atomic::{AtomicBool, AtomicU8, Ordering}, + mpsc::{self, Receiver, RecvTimeoutError, SyncSender}, + Arc, Mutex, + }, + thread::JoinHandle, + time::{Duration, SystemTime, UNIX_EPOCH}, +}; +fn now_ms() -> Result { + u64::try_from( + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_err(|_| HostError::new("EXTENSION_CLOCK_INVALID"))? + .as_millis(), + ) + .map_err(|_| HostError::new("EXTENSION_CLOCK_INVALID")) +} +pub type ResumeCheck = Box Result<()> + Send>; +pub struct LaunchSpec { + pub package: cap_std::fs::Dir, + pub inventory: Inventory, + pub claims: Claims, + pub permit: Permit, + pub authority: Arc, + pub credentials: Arc>, + pub vault_id: String, + pub policy_version: String, + pub system_root: PathBuf, + /// Revalidate active install/current trust and all external policy immediately + /// before resume, after expensive package checks. Errors prohibit execution. + pub before_resume: ResumeCheck, +} +#[derive(Clone, Copy, Serialize, PartialEq, Eq, Debug)] +#[serde(rename_all = "snake_case")] +pub enum Status { + Starting, + Ready, + Stopping, + Stopped, + Failed, +} +#[derive(Serialize)] +pub struct Snapshot { + pub status: Status, + pub identity: Option, + pub tool_count: usize, + pub error: Option, +} +struct Control { + stop: AtomicBool, + #[cfg(test)] + job: Mutex>, + active: Mutex>>, + status: AtomicU8, + identity: Mutex>, + tools: Mutex>, + error: Mutex>, +} +impl Control { + fn status(&self) -> Status { + match self.status.load(Ordering::Acquire) { + 0 => Status::Starting, + 1 => Status::Ready, + 2 => Status::Stopping, + 3 => Status::Stopped, + _ => Status::Failed, + } + } + fn stop(&self) { + self.stop.store(true, Ordering::Release); + if let Some(active) = &*self.active.lock().unwrap_or_else(|e| e.into_inner()) { + active.store(true, Ordering::Release); + } + let _ = self + .status + .fetch_update(Ordering::AcqRel, Ordering::Acquire, |value| { + (value < 2).then_some(2) + }); + } +} +struct Active<'a>(&'a Control); +impl Drop for Active<'_> { + fn drop(&mut self) { + self.0 + .active + .lock() + .unwrap_or_else(|e| e.into_inner()) + .take(); + } +} +struct Request { + cancel: Arc, + reply: SyncSender>, +} +impl Request { + fn execute( + self, + control: &Control, + action: impl FnOnce(&AtomicBool) -> Result, + ) -> Option { + let mut active = control.active.lock().unwrap_or_else(|e| e.into_inner()); + if control.stop.load(Ordering::Acquire) || self.cancel.load(Ordering::Acquire) { + let _ = self + .reply + .try_send(Err(HostError::new("EXTENSION_INSTANCE_CANCELLED"))); + return None; + } + *active = Some(Arc::clone(&self.cancel)); + drop(active); + let _active = Active(control); + let result = action(&self.cancel); + let error = result.as_ref().err().map(|e| e.code.clone()); + let _ = self.reply.try_send(result); + error + } +} +enum Command { + Review { + name: String, + arguments: Value, + request: Request, + }, + Invoke { + review_id: String, + request: Request, + }, + Refresh { + request: Request>, + }, +} +pub struct Ticket { + receiver: Receiver>, + cancel: Arc, +} +impl Ticket { + /// Background wait only; dropping a ticket cancels its queued/in-flight work. + pub fn wait(self, timeout: Duration) -> Result { + if timeout > Duration::from_secs(65) { + return Err(HostError::new("EXTENSION_INSTANCE_WAIT_INVALID")); + } + self.receiver.recv_timeout(timeout).map_err(|error| { + HostError::new(match error { + RecvTimeoutError::Timeout => "EXTENSION_INSTANCE_TIMEOUT", + RecvTimeoutError::Disconnected => "EXTENSION_INSTANCE_CLOSED", + }) + })? + } + pub fn cancel(&self) { + self.cancel.store(true, Ordering::Release); + } +} +impl Drop for Ticket { + fn drop(&mut self) { + self.cancel.store(true, Ordering::Release); + } +} +#[derive(Clone)] +pub struct Endpoint { + control: Arc, + commands: SyncSender, +} +impl Endpoint { + fn submit(&self, build: impl FnOnce(Request) -> Command) -> Result> { + if self.control.status() != Status::Ready || self.control.stop.load(Ordering::Acquire) { + return Err(HostError::new("EXTENSION_INSTANCE_NOT_READY")); + } + let (reply, receiver) = mpsc::sync_channel(1); + let cancel = Arc::new(AtomicBool::new(false)); + self.commands + .try_send(build(Request { + cancel: Arc::clone(&cancel), + reply, + })) + .map_err(|_| HostError::new("EXTENSION_INSTANCE_BACKPRESSURE"))?; + Ok(Ticket { receiver, cancel }) + } + pub fn review(&self, name: String, arguments: Value) -> Result> { + if name.is_empty() || name.len() > 128 { + return Err(HostError::new("EXTENSION_MCP_ARGUMENTS_INVALID")); + } + crate::extension_mcp_tools::argument_bounds(&arguments)?; + self.submit(|request| Command::Review { + name, + arguments, + request, + }) + } + /// Host route only: the user must have approved the exact saved review. + /// Confirmation and consumption happen together on the instance thread. + pub fn invoke_confirmed(&self, review_id: String) -> Result> { + if uuid::Uuid::parse_str(&review_id).is_err() || review_id.len() != 36 { + return Err(HostError::new("EXTENSION_CALL_REVIEW_UNKNOWN")); + } + self.submit(|request| Command::Invoke { review_id, request }) + } + pub fn refresh(&self) -> Result>> { + self.submit(|request| Command::Refresh { request }) + } + pub fn stop(&self) { + self.control.stop(); + } + pub fn snapshot(&self) -> Snapshot { + Snapshot { + status: self.control.status(), + identity: self + .control + .identity + .lock() + .unwrap_or_else(|e| e.into_inner()) + .clone(), + tool_count: self + .control + .tools + .lock() + .unwrap_or_else(|e| e.into_inner()) + .len(), + error: self + .control + .error + .lock() + .unwrap_or_else(|e| e.into_inner()) + .clone(), + } + } +} +struct Entry { + endpoint: Endpoint, + worker: JoinHandle<()>, +} +#[derive(Default)] +pub struct Registry { + entries: BTreeMap, +} +impl Registry { + /// # Safety + /// The caller must establish all sandbox limits and current install/user + /// authorization. before_resume must recheck live trust/active installation. + /// This API is not exposed to renderer/Core and does not enable extensions. + pub unsafe fn start(&mut self, spec: LaunchSpec) -> Result { + self.reap(); + spec.authority + .verify(&spec.permit, &spec.claims, now_ms()?)?; + let key = format!( + "{:x}", + Sha256::digest( + serde_json::to_vec(&( + spec.claims.kind, + &spec.claims.source, + &spec.claims.namespace, + &spec.claims.package_id, + &spec.claims.vault_id + )) + .map_err(|_| HostError::new("EXTENSION_INSTANCE_INVALID"))? + ) + ); + if self.entries.contains_key(&key) { + return Err(HostError::new("EXTENSION_INSTANCE_ALREADY_RUNNING")); + } + if self.entries.len() >= 16 { + return Err(HostError::new("EXTENSION_INSTANCE_LIMIT")); + } + let control = Arc::new(Control { + stop: AtomicBool::new(false), + #[cfg(test)] + job: Mutex::new(None), + active: Mutex::new(None), + status: AtomicU8::new(0), + identity: Mutex::new(None), + tools: Mutex::new(Vec::new()), + error: Mutex::new(None), + }); + let (commands, receiver) = mpsc::sync_channel(4); + let endpoint = Endpoint { + control: Arc::clone(&control), + commands, + }; + let worker = std::thread::Builder::new() + .name("opennexus-instance".into()) + .spawn(move || { + let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + run(spec, &control, receiver) + })); + // All native stack owners have dropped before publishing terminal state. + let error = match result { + Ok(Ok(())) => None, + Ok(Err(error)) + if control.stop.load(Ordering::Acquire) + && matches!( + error.code.as_str(), + "EXTENSION_MCP_CANCELLED" | "EXTENSION_INSTANCE_CANCELLED" + ) => + { + None + } + Ok(Err(error)) => Some(error.code), + Err(_) => Some("EXTENSION_INSTANCE_WORKER_FAILED".into()), + }; + let failed = error.is_some(); + *control.error.lock().unwrap_or_else(|e| e.into_inner()) = error; + control + .tools + .lock() + .unwrap_or_else(|e| e.into_inner()) + .clear(); + control + .status + .store(if failed { 4 } else { 3 }, Ordering::Release); + }) + .map_err(|_| HostError::new("EXTENSION_INSTANCE_WORKER_UNAVAILABLE"))?; + self.entries.insert( + key, + Entry { + endpoint: endpoint.clone(), + worker, + }, + ); + Ok(endpoint) + } + /// Reap only threads confirmed finished, so an old generation cannot overlap + /// a replacement merely because stop was requested or status was changed. + pub fn reap(&mut self) { + let done: Vec<_> = self + .entries + .iter() + .filter(|(_, entry)| { + entry.worker.is_finished() + && !entry + .endpoint + .control + .error + .lock() + .unwrap_or_else(|e| e.into_inner()) + .as_deref() + .is_some_and(|code| { + matches!( + code, + "EXTENSION_CONTAINER_CLEANUP_FAILED" + | "EXTENSION_RESOURCE_TERMINATE_FAILED" + | "EXTENSION_INSTANCE_WORKER_FAILED" + ) + }) + }) + .map(|(key, _)| key.clone()) + .collect(); + for key in done { + if let Some(entry) = self.entries.remove(&key) { + let _ = entry.worker.join(); + } + } + } + pub fn stop_all(&self) { + for entry in self.entries.values() { + entry.endpoint.stop(); + } + } +} +impl Drop for Registry { + fn drop(&mut self) { + self.stop_all(); + for (_, entry) in std::mem::take(&mut self.entries) { + let _ = entry.worker.join(); + } + } +} +fn run(spec: LaunchSpec, control: &Control, receiver: Receiver) -> Result<()> { + if control.stop.load(Ordering::Acquire) { + return Ok(()); + } + let profile = Profile::create()?; + let result = run_in_profile(spec, control, receiver, &profile); + profile.remove()?; + result +} +fn run_in_profile( + spec: LaunchSpec, + control: &Control, + receiver: Receiver, + profile: &Profile, +) -> Result<()> { + let pinned = PinnedPackage::open(&spec.package, &spec.inventory, &spec.claims.tree_sha256)?; + pinned.grant_read_execute(profile)?; + let entry = pinned.bind_entry(&spec.claims.entry)?; + let folder = profile.folder()?; + let scratch = folder.join("Temp"); + let context = Context { + vault_id: &spec.vault_id, + policy_version: &spec.policy_version, + system_root: &spec.system_root, + container_data: &folder, + scratch: &scratch, + }; + let prepared = { + let credentials = spec + .credentials + .lock() + .map_err(|_| HostError::new("CREDENTIALS_LOCKED"))?; + context.prepare( + &spec.authority, + &spec.permit, + &spec.claims, + &entry, + &credentials, + now_ms()?, + )? + }; + let (suspended, io) = prepared.create_suspended_with_stdio(profile, &entry)?; + (spec.before_resume)(&spec.claims)?; + if control.stop.load(Ordering::Acquire) { + return Ok(()); + } + // Safety obligation belongs to Registry::start's caller, rechecked above. + let running = unsafe { suspended.resume()? }; + #[cfg(test)] + { + *control.job.lock().unwrap() = Some(running.test_job()?); + } + let mut session = Session::new(&running, io)?; + session.initialize(&control.stop)?; + let tools = session.refresh_tools(&control.stop)?; + *control.identity.lock().unwrap_or_else(|e| e.into_inner()) = Some(running.call_identity()?); + *control.tools.lock().unwrap_or_else(|e| e.into_inner()) = tools; + if control.stop.load(Ordering::Acquire) { + return session.shutdown(); + } + let _ = control + .status + .compare_exchange(0, 1, Ordering::AcqRel, Ordering::Acquire); + while !control.stop.load(Ordering::Acquire) { + session.drain_pending()?; + if session.take_tools_changed() { + control + .tools + .lock() + .unwrap_or_else(|e| e.into_inner()) + .clear(); + } + let error = match receiver.recv_timeout(Duration::from_millis(20)) { + Ok(Command::Review { + name, + arguments, + request, + }) => request.execute(control, |_| session.review_call(&name, arguments)), + Ok(Command::Invoke { review_id, request }) => request.execute(control, |cancel| { + let approved = session.confirm_call(&review_id)?; + session.call_tool(approved, cancel) + }), + Ok(Command::Refresh { request }) => request.execute(control, |cancel| { + let tools = session.refresh_tools(cancel)?; + *control.tools.lock().unwrap_or_else(|e| e.into_inner()) = tools.clone(); + Ok(tools) + }), + Err(RecvTimeoutError::Timeout) => None, + Err(RecvTimeoutError::Disconnected) => break, + }; + if session.is_failed() { + if error.as_deref() == Some("EXTENSION_MCP_CANCELLED") { + return session.shutdown(); + } + return Err(HostError::new( + error.as_deref().unwrap_or("EXTENSION_INSTANCE_FAILED"), + )); + } + } + session.shutdown() +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + use std::time::Instant; + use zeroize::Zeroizing; + fn wait_for(mut condition: impl FnMut() -> bool) { + let started = Instant::now(); + while !condition() { + assert!(started.elapsed() < Duration::from_secs(10)); + std::thread::sleep(Duration::from_millis(5)); + } + } + #[test] + fn native_worker_routes_reviews_cancels_calls_and_reaps_generations() { + let temp = tempfile::tempdir().unwrap(); + let package = temp.path().join("package"); + std::fs::create_dir(&package).unwrap(); + let executable = package.join("entry.exe"); + let fixture = PathBuf::from(env!("CARGO_MANIFEST_DIR")) + .join("tests/fixtures/sandbox_network_probe.rs"); + let compiled = std::process::Command::new("rustc") + .arg("--edition=2021") + .arg(fixture) + .arg("-o") + .arg(&executable) + .output() + .unwrap(); + assert!( + compiled.status.success(), + "{}", + String::from_utf8_lossy(&compiled.stderr) + ); + let bytes = std::fs::read(&executable).unwrap(); + let files: BTreeMap = + [("entry.exe".into(), format!("{:x}", Sha256::digest(&bytes)))] + .into_iter() + .collect(); + let dir = + cap_std::fs::Dir::open_ambient_dir(&package, cap_std::ambient_authority()).unwrap(); + let inventory = || Inventory { + files: files.clone(), + expanded_size: bytes.len() as u64, + manifest: "entry.exe".into(), + }; + let tree = crate::extension_unpack::verify_tree(&dir, &inventory()).unwrap(); + let authority = Arc::new(Authority::default()); + let credentials = Arc::new(Mutex::new(CredentialBroker::new( + temp.path().join("credentials.v1"), + ))); + credentials + .lock() + .unwrap() + .unlock(Zeroizing::new(b"instance fixture password".to_vec())) + .unwrap(); + let vault_id = uuid::Uuid::new_v4().to_string(); + let make = |mode: &str| { + let claims = Claims { + kind: crate::extension_permit::ExecutionKind::Mcp, + source: "https://catalog.example/".into(), + namespace: "examples".into(), + package_id: "native-worker".into(), + version: "1.0.0".into(), + archive_sha256: "a".repeat(64), + tree_sha256: tree.clone(), + signer_sha256: "b".repeat(64), + entry: "entry.exe".into(), + arguments: vec![mode.into()], + environment: BTreeMap::new(), + permissions: Default::default(), + vault_id: vault_id.clone(), + platform: "windows".into(), + policy_version: "1".into(), + expires_at_ms: now_ms().unwrap() + 120_000, + }; + LaunchSpec { + package: dir.try_clone().unwrap(), + inventory: inventory(), + permit: authority.issue(&claims, now_ms().unwrap()).unwrap(), + claims, + authority: Arc::clone(&authority), + credentials: Arc::clone(&credentials), + vault_id: vault_id.clone(), + policy_version: "1".into(), + system_root: PathBuf::from(std::env::var_os("SystemRoot").unwrap()), + before_resume: Box::new(|_| Ok(())), + } + }; + let mut registry = Registry::default(); + let mut invalid = make("mcp"); + invalid.claims.arguments = vec!["changed".into()]; + assert!(unsafe { registry.start(invalid) }.is_err()); + assert!(registry.entries.is_empty()); + let endpoint = unsafe { registry.start(make("mcp")) }.unwrap(); + assert_eq!( + unsafe { registry.start(make("mcp")) }.err().unwrap().code, + "EXTENSION_INSTANCE_ALREADY_RUNNING" + ); + wait_for(|| endpoint.snapshot().status != Status::Starting); + assert_eq!( + endpoint.snapshot().status, + Status::Ready, + "{:?}", + endpoint.snapshot().error + ); + assert_eq!(endpoint.snapshot().tool_count, 1); + let first_identity = serde_json::to_value(endpoint.snapshot().identity.unwrap()).unwrap(); + let review = endpoint + .review("echo".into(), json!({})) + .unwrap() + .wait(Duration::from_secs(5)) + .unwrap(); + let result = endpoint + .invoke_confirmed(review.review_id.clone()) + .unwrap() + .wait(Duration::from_secs(5)) + .unwrap(); + assert_eq!(result["content"][0]["text"], "native MCP success"); + assert_eq!( + endpoint + .invoke_confirmed(review.review_id) + .unwrap() + .wait(Duration::from_secs(5)) + .unwrap_err() + .code, + "EXTENSION_CALL_REVIEW_UNKNOWN" + ); + endpoint.stop(); + wait_for(|| { + registry.reap(); + registry.entries.is_empty() + }); + assert_eq!( + endpoint.snapshot().status, + Status::Stopped, + "{:?}", + endpoint.snapshot().error + ); + assert!(endpoint.review("echo".into(), json!({})).is_err()); + assert_eq!( + endpoint + .control + .job + .lock() + .unwrap() + .as_ref() + .unwrap() + .active_processes() + .unwrap(), + 0 + ); + let second = unsafe { registry.start(make("mcp_cancel")) }.unwrap(); + wait_for(|| second.snapshot().status != Status::Starting); + assert_eq!( + second.snapshot().status, + Status::Ready, + "{:?}", + second.snapshot().error + ); + assert_ne!( + serde_json::to_value(second.snapshot().identity.unwrap()).unwrap()["instance_id"], + first_identity["instance_id"] + ); + let review = second + .review("echo".into(), json!({})) + .unwrap() + .wait(Duration::from_secs(5)) + .unwrap(); + let baseline = second + .control + .job + .lock() + .unwrap() + .as_ref() + .unwrap() + .active_processes() + .unwrap(); + let ticket = second.invoke_confirmed(review.review_id).unwrap(); + wait_for(|| { + let active = second + .control + .active + .lock() + .unwrap() + .as_ref() + .is_some_and(|active| Arc::ptr_eq(active, &ticket.cancel)); + let processes = second + .control + .job + .lock() + .unwrap() + .as_ref() + .unwrap() + .active_processes() + .unwrap(); + active && processes > baseline + }); + ticket.cancel(); + assert_eq!( + ticket.wait(Duration::from_secs(5)).unwrap_err().code, + "EXTENSION_MCP_CANCELLED" + ); + wait_for(|| { + registry.reap(); + registry.entries.is_empty() + }); + assert_eq!( + second.snapshot().status, + Status::Stopped, + "{:?}", + second.snapshot().error + ); + assert_eq!( + second + .control + .job + .lock() + .unwrap() + .as_ref() + .unwrap() + .active_processes() + .unwrap(), + 0 + ); + let called = Arc::new(AtomicBool::new(false)); + let marker = Arc::clone(&called); + let mut denied = make("mcp"); + denied.before_resume = Box::new(move |_| { + marker.store(true, Ordering::Release); + Err(HostError::new("FIXTURE_POLICY_DENIED")) + }); + let denied = unsafe { registry.start(denied) }.unwrap(); + wait_for(|| { + registry.reap(); + registry.entries.is_empty() + }); + assert!(called.load(Ordering::Acquire)); + assert_eq!(denied.snapshot().status, Status::Failed); + assert_eq!( + denied.snapshot().error.as_deref(), + Some("FIXTURE_POLICY_DENIED") + ); + assert!(denied.control.job.lock().unwrap().is_none()); + let locked = unsafe { registry.start(make("mcp")) }.unwrap(); + wait_for(|| locked.snapshot().status != Status::Starting); + assert_eq!(locked.snapshot().status, Status::Ready); + credentials.lock().unwrap().lock(); + wait_for(|| { + registry.reap(); + registry.entries.is_empty() + }); + assert_eq!(locked.snapshot().status, Status::Failed); + assert_eq!( + locked.snapshot().error.as_deref(), + Some("CREDENTIALS_LOCKED") + ); + assert_eq!( + locked + .control + .job + .lock() + .unwrap() + .as_ref() + .unwrap() + .active_processes() + .unwrap(), + 0 + ); + } + fn channel() -> (Endpoint, Receiver) { + let (commands, receiver) = mpsc::sync_channel(4); + let control = Arc::new(Control { + stop: AtomicBool::new(false), + job: Mutex::new(None), + active: Mutex::new(None), + status: AtomicU8::new(1), + identity: Mutex::new(None), + tools: Mutex::new(Vec::new()), + error: Mutex::new(None), + }); + (Endpoint { control, commands }, receiver) + } + #[test] + fn bounded_queue_and_stop_prevent_queued_work_from_executing() { + let (endpoint, receiver) = channel(); + let mut tickets = Vec::new(); + for _ in 0..4 { + tickets.push(endpoint.review("echo".into(), json!({})).unwrap()); + } + assert_eq!( + endpoint + .review("echo".into(), json!({})) + .err() + .unwrap() + .code, + "EXTENSION_INSTANCE_BACKPRESSURE" + ); + assert!(endpoint + .review("echo".into(), json!({"large":"x".repeat(256 * 1024)})) + .is_err()); + endpoint.stop(); + assert_eq!(endpoint.snapshot().status, Status::Stopping); + for ticket in tickets { + let Command::Review { request, .. } = receiver.try_recv().unwrap() else { + panic!("wrong request") + }; + request.execute(&endpoint.control, |_| panic!("stopped queue executed")); + assert_eq!( + ticket.wait(Duration::from_secs(1)).err().unwrap().code, + "EXTENSION_INSTANCE_CANCELLED" + ); + } + assert!(endpoint.refresh().is_err()); + } + #[test] + fn cleanup_failures_remain_quarantined_even_after_worker_exit() { + let (endpoint, _receiver) = channel(); + endpoint.control.status.store(4, Ordering::Release); + *endpoint.control.error.lock().unwrap() = Some("EXTENSION_CONTAINER_CLEANUP_FAILED".into()); + let worker = std::thread::spawn(|| {}); + wait_for(|| worker.is_finished()); + let mut registry = Registry::default(); + registry + .entries + .insert("quarantined".into(), Entry { endpoint, worker }); + registry.reap(); + assert_eq!(registry.entries.len(), 1); + } +} diff --git a/frontend/src-tauri/src/extension_mcp.rs b/frontend/src-tauri/src/extension_mcp.rs index 9318d9e..f32eb5f 100644 --- a/frontend/src-tauri/src/extension_mcp.rs +++ b/frontend/src-tauri/src/extension_mcp.rs @@ -420,6 +420,9 @@ impl<'a, 'p> Session<'a, 'p> { } } } + pub(crate) fn is_failed(&self) -> bool { + self.failed + } pub fn take_tools_changed(&mut self) -> bool { std::mem::take(&mut self.tools_changed) } diff --git a/frontend/src-tauri/src/extension_mcp_tools.rs b/frontend/src-tauri/src/extension_mcp_tools.rs index b8abd70..ee47a52 100644 --- a/frontend/src-tauri/src/extension_mcp_tools.rs +++ b/frontend/src-tauri/src/extension_mcp_tools.rs @@ -44,6 +44,13 @@ fn bounded(value: &Value, bytes: usize, schema: bool) -> Result<()> { } Ok(()) } +pub(crate) fn argument_bounds(value: &Value) -> Result<()> { + bounded(value, 256 * 1024, false)?; + if !value.is_object() { + return Err(HostError::new("EXTENSION_MCP_ARGUMENTS_INVALID")); + } + Ok(()) +} fn compile(schema: &Value) -> Result { bounded(schema, 64 * 1024, true)?; if !schema.is_object() diff --git a/frontend/src-tauri/src/extension_process.rs b/frontend/src-tauri/src/extension_process.rs index 0762c20..e9d85c5 100644 --- a/frontend/src-tauri/src/extension_process.rs +++ b/frontend/src-tauri/src/extension_process.rs @@ -298,6 +298,10 @@ impl Running<'_> { Ok(()) } #[cfg(test)] + pub(crate) fn test_job(&self) -> Result { + self.process.job.clone_for_deadline() + } + #[cfg(test)] pub(crate) fn active_test_processes(&self) -> Result { self.process.job.active_processes() } diff --git a/frontend/src-tauri/src/lib.rs b/frontend/src-tauri/src/lib.rs index d604c0c..3f5fb26 100644 --- a/frontend/src-tauri/src/lib.rs +++ b/frontend/src-tauri/src/lib.rs @@ -93,3 +93,6 @@ pub mod extension_mcp_tools; #[cfg(all(windows, feature = "desktop"))] pub mod extension_call_authorization; + +#[cfg(all(windows, feature = "desktop"))] +pub mod extension_instance;