| 1 | //! Real committed host bundle, SDK, Rust ticket redemption and contained |
| 2 | //! fixture child. The child is a local MCP peer, never a provider or real app. |
| 3 | use super::*; |
| 4 | use crate::extension_host::tests::{FixturePlugins, node_for_tests}; |
| 5 | use crate::extension_host::{ExtensionHostOptions, TestManagerGuard}; |
| 6 | use crate::mcp::{McpBackend, McpConnection, McpTimeouts, McpTransport}; |
| 7 | use crate::plugins::activation::TestPolicyGuard; |
| 8 | use std::sync::atomic::Ordering; |
| 9 | |
| 10 | const PEER: &str = r#" |
| 11 | import readline from 'node:readline'; import fs from 'node:fs'; |
| 12 | const record = process.argv[2]; |
| 13 | const rl = readline.createInterface({input:process.stdin}); |
| 14 | rl.on('line', line => { |
| 15 | const m=JSON.parse(line); fs.appendFileSync(record, JSON.stringify({method:m.method,wire_id:m.id,params:m.params,ping_reply:m.id==='fixture-ping'&&m.result&&Object.keys(m.result).length===0})+'\n'); |
| 16 | if(m.method==='notifications/initialized') { process.stdout.write(JSON.stringify({jsonrpc:'2.0',id:'fixture-ping',method:'ping',params:{}})+'\n'); return; } |
| 17 | if (!m.method || !Object.hasOwn(m,'id')) return; |
| 18 | if(typeof m.id !== 'string') throw new Error('SDK numeric ID escaped the Rust wire adapter'); |
| 19 | const result=m.method==='initialize'?{protocolVersion:'2025-06-18',serverInfo:{name:'broker-fixture',version:'1'},capabilities:{tools:{}}} |
| 20 | :m.method==='tools/list'?{tools:[{name:'echo',inputSchema:{type:'object'}}]} |
| 21 | :{content:[{type:'text',text:JSON.stringify(m.params.arguments)}],isError:false}; |
| 22 | process.stdout.write(JSON.stringify({jsonrpc:'2.0',id:m.id,result})+'\n'); |
| 23 | }); |
| 24 | "#; |
| 25 | |
| 26 | #[cfg(test)] |
| 27 | async fn fixture( |
| 28 | test: &str, |
| 29 | ) -> Option<( |
| 30 | FixturePlugins, |
| 31 | Arc<ExtensionHostManager>, |
| 32 | std::path::PathBuf, |
| 33 | McpServerConfig, |
| 34 | )> { |
| 35 | fixture_with_plugins(test, &[]).await |
| 36 | } |
| 37 | #[cfg(test)] |
| 38 | async fn fixture_with_plugins( |
| 39 | test: &str, |
| 40 | names: &[&str], |
| 41 | ) -> Option<( |
| 42 | FixturePlugins, |
| 43 | Arc<ExtensionHostManager>, |
| 44 | std::path::PathBuf, |
| 45 | McpServerConfig, |
| 46 | )> { |
| 47 | let node = node_for_tests(test)?; |
| 48 | let fixture = FixturePlugins::new(names).await; |
| 49 | std::fs::create_dir_all(&fixture.root).unwrap(); |
| 50 | let script = fixture.root.join("mcp-peer.mjs"); |
| 51 | let record = fixture.root.join("mcp-frames.jsonl"); |
| 52 | std::fs::write(&script, PEER).unwrap(); |
| 53 | std::fs::write(&record, "").unwrap(); |
| 54 | let config: McpServerConfig = serde_json::from_value(json!({ |
| 55 | "command": node.to_string_lossy(), |
| 56 | "args": [script.to_string_lossy(), record.to_string_lossy()], |
| 57 | "connect_timeout": 5, |
| 58 | })) |
| 59 | .unwrap(); |
| 60 | let manager = fixture.manager(node); |
| 61 | Some((fixture, manager, record, config)) |
| 62 | } |
| 63 | async fn transport(config: &McpServerConfig) -> SdkTransport { |
| 64 | let mut transport = SdkTransport::connect( |
| 65 | "fixture", |
| 66 | config, |
| 67 | CancellationToken::new(), |
| 68 | Duration::from_secs(5), |
| 69 | ) |
| 70 | .await |
| 71 | .unwrap(); |
| 72 | transport.send(serde_json::to_vec(&json!({"jsonrpc":"2.0","id":"init","method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{},"clientInfo":{"name":"codewhale-tui","version":env!("CARGO_PKG_VERSION")}}})).unwrap()).await.unwrap(); |
| 73 | let init: Value = serde_json::from_slice(&transport.recv().await.unwrap()).unwrap(); |
| 74 | assert_eq!(init["id"], "init"); |
| 75 | transport |
| 76 | } |
| 77 | fn request( |
| 78 | transport: &SdkTransport, |
| 79 | grant: &McpOperationGrant, |
| 80 | params: Value, |
| 81 | id: u64, |
| 82 | ) -> HostRequest { |
| 83 | let params = ProcWriteParams { |
| 84 | owner: transport.session.owner.clone(), |
| 85 | session_id: transport.session_id.clone(), |
| 86 | ticket: Some(grant.ticket.clone()), |
| 87 | operation_id: Some(grant.operation_id.clone()), |
| 88 | frame: json!({"jsonrpc":"2.0","id":grant.wire_id,"method":grant.method,"params":params}), |
| 89 | }; |
| 90 | let frame = json!({"jsonrpc":"2.0","id":id,"method":"proc/write","params":params}); |
| 91 | match parse_host_message(frame, HostTier::Builtin).unwrap() { |
| 92 | HostMessage::Request { request, .. } => request, |
| 93 | _ => panic!("decoded process write must be a request"), |
| 94 | } |
| 95 | } |
| 96 | |
| 97 | #[tokio::test(flavor = "current_thread")] |
| 98 | async fn host_sdk_real_connection_keeps_rust_catalog_and_exact_tool_call() { |
| 99 | let _policy = TestPolicyGuard::extension_host(true); |
| 100 | let Some((_fixture, manager, record, config)) = fixture("host_sdk_real_connection").await |
| 101 | else { |
| 102 | return; |
| 103 | }; |
| 104 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 105 | let mut connection = McpConnection::connect_with_backend( |
| 106 | "fixture".into(), |
| 107 | config, |
| 108 | &McpTimeouts::default(), |
| 109 | None, |
| 110 | McpBackend::Host, |
| 111 | ) |
| 112 | .await |
| 113 | .unwrap(); |
| 114 | assert!(connection.is_ready()); |
| 115 | assert_eq!(connection.tools().len(), 1); |
| 116 | assert_eq!(connection.tools()[0].name, "echo"); |
| 117 | let result = connection |
| 118 | .call_tool("echo", json!({"payload":"exact"}), 5) |
| 119 | .await |
| 120 | .unwrap(); |
| 121 | assert_eq!(result["content"][0]["text"], "{\"payload\":\"exact\"}"); |
| 122 | tokio::time::timeout(Duration::from_secs(2), async { |
| 123 | while !std::fs::read_to_string(&record) |
| 124 | .unwrap() |
| 125 | .contains("\"ping_reply\":true") |
| 126 | { |
| 127 | tokio::task::yield_now().await; |
| 128 | } |
| 129 | }) |
| 130 | .await |
| 131 | .unwrap(); |
| 132 | let frames = std::fs::read_to_string(record).unwrap(); |
| 133 | assert_eq!( |
| 134 | frames |
| 135 | .lines() |
| 136 | .filter(|line| line.contains("\"method\":\"initialize\"")) |
| 137 | .count(), |
| 138 | 1 |
| 139 | ); |
| 140 | assert_eq!( |
| 141 | frames |
| 142 | .lines() |
| 143 | .filter(|line| line.contains("\"method\":\"tools/call\"")) |
| 144 | .count(), |
| 145 | 1 |
| 146 | ); |
| 147 | assert!(!frames.contains("server/discover")); |
| 148 | assert!(frames.contains("\"ping_reply\":true")); |
| 149 | let frames: Vec<Value> = frames |
| 150 | .lines() |
| 151 | .map(|line| serde_json::from_str(line).unwrap()) |
| 152 | .collect(); |
| 153 | for (method, expected) in [ |
| 154 | ("initialize", "1"), |
| 155 | ("tools/list", "2"), |
| 156 | ("tools/call", "3"), |
| 157 | ] { |
| 158 | assert_eq!( |
| 159 | frames |
| 160 | .iter() |
| 161 | .find(|frame| frame["method"] == method) |
| 162 | .unwrap()["wire_id"], |
| 163 | expected, |
| 164 | "the peer must see the original Rust facade ID" |
| 165 | ); |
| 166 | } |
| 167 | } |
| 168 | |
| 169 | #[tokio::test(flavor = "current_thread")] |
| 170 | async fn decoded_operation_mismatch_refuses_before_pipe_and_retires_session() { |
| 171 | let _policy = TestPolicyGuard::extension_host(true); |
| 172 | let Some((_fixture, manager, record, config)) = fixture("decoded_operation_mismatch").await |
| 173 | else { |
| 174 | return; |
| 175 | }; |
| 176 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 177 | let transport = transport(&config).await; |
| 178 | let grant = transport |
| 179 | .grant( |
| 180 | "tools/list", |
| 181 | json!({"cursor":"admitted"}), |
| 182 | Duration::from_secs(5), |
| 183 | None, |
| 184 | Some("raw-accepted-id"), |
| 185 | ) |
| 186 | .unwrap(); |
| 187 | let (cx, _, _) = HostRequestContext::for_test(7); |
| 188 | let result = manager |
| 189 | .shared |
| 190 | .mcp_broker |
| 191 | .serve( |
| 192 | &manager.shared, |
| 193 | transport.session.host_generation, |
| 194 | request(&transport, &grant, json!({"cursor":"rewritten"}), 7), |
| 195 | cx, |
| 196 | ) |
| 197 | .await; |
| 198 | assert!(result.is_err()); |
| 199 | assert!(transport.session.cancel.is_cancelled()); |
| 200 | assert!( |
| 201 | !std::fs::read_to_string(record) |
| 202 | .unwrap() |
| 203 | .contains("rewritten") |
| 204 | ); |
| 205 | } |
| 206 | |
| 207 | #[tokio::test(flavor = "current_thread")] |
| 208 | async fn exact_grant_replay_and_stale_owner_host_are_refused() { |
| 209 | let _policy = TestPolicyGuard::extension_host(true); |
| 210 | let Some((_fixture, manager, record, config)) = fixture("exact_grant_replay").await else { |
| 211 | return; |
| 212 | }; |
| 213 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 214 | let transport = transport(&config).await; |
| 215 | let grant = transport |
| 216 | .grant( |
| 217 | "tools/list", |
| 218 | json!({}), |
| 219 | Duration::from_secs(5), |
| 220 | None, |
| 221 | Some("raw-accepted-id"), |
| 222 | ) |
| 223 | .unwrap(); |
| 224 | let mut foreign = request(&transport, &grant, json!({}), 8); |
| 225 | if let HostRequest::ProcWrite(params) = &mut foreign { |
| 226 | params.owner.generation += 1; |
| 227 | } |
| 228 | let (cx, _, _) = HostRequestContext::for_test(8); |
| 229 | assert!( |
| 230 | manager |
| 231 | .shared |
| 232 | .mcp_broker |
| 233 | .serve( |
| 234 | &manager.shared, |
| 235 | transport.session.host_generation, |
| 236 | foreign, |
| 237 | cx |
| 238 | ) |
| 239 | .await |
| 240 | .is_err() |
| 241 | ); |
| 242 | let (cx, _, _) = HostRequestContext::for_test(9); |
| 243 | assert!( |
| 244 | manager |
| 245 | .shared |
| 246 | .mcp_broker |
| 247 | .serve( |
| 248 | &manager.shared, |
| 249 | transport.session.host_generation + 1, |
| 250 | request(&transport, &grant, json!({}), 9), |
| 251 | cx |
| 252 | ) |
| 253 | .await |
| 254 | .is_err() |
| 255 | ); |
| 256 | let (cx, _, _) = HostRequestContext::for_test(10); |
| 257 | manager |
| 258 | .shared |
| 259 | .mcp_broker |
| 260 | .serve( |
| 261 | &manager.shared, |
| 262 | transport.session.host_generation, |
| 263 | request(&transport, &grant, json!({}), 10), |
| 264 | cx, |
| 265 | ) |
| 266 | .await |
| 267 | .unwrap(); |
| 268 | tokio::time::timeout(Duration::from_secs(2), async { |
| 269 | loop { |
| 270 | if std::fs::read_to_string(&record) |
| 271 | .unwrap() |
| 272 | .contains("tools/list") |
| 273 | { |
| 274 | break; |
| 275 | } |
| 276 | tokio::task::yield_now().await; |
| 277 | } |
| 278 | }) |
| 279 | .await |
| 280 | .unwrap(); |
| 281 | let (cx, _, _) = HostRequestContext::for_test(11); |
| 282 | assert!( |
| 283 | manager |
| 284 | .shared |
| 285 | .mcp_broker |
| 286 | .serve( |
| 287 | &manager.shared, |
| 288 | transport.session.host_generation, |
| 289 | request(&transport, &grant, json!({}), 11), |
| 290 | cx |
| 291 | ) |
| 292 | .await |
| 293 | .is_err() |
| 294 | ); |
| 295 | assert_eq!( |
| 296 | std::fs::read_to_string(record) |
| 297 | .unwrap() |
| 298 | .lines() |
| 299 | .filter(|line| line.contains("tools/list")) |
| 300 | .count(), |
| 301 | 1 |
| 302 | ); |
| 303 | assert!(transport.session.cancel.is_cancelled()); |
| 304 | } |
| 305 | |
| 306 | #[tokio::test(flavor = "current_thread")] |
| 307 | async fn host_withdrawal_revokes_unwritten_grant_and_broker_lifetime() { |
| 308 | let _policy = TestPolicyGuard::extension_host(true); |
| 309 | let Some((_fixture, manager, record, config)) = fixture("host_withdrawal").await else { |
| 310 | return; |
| 311 | }; |
| 312 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 313 | let transport = transport(&config).await; |
| 314 | let grant = transport |
| 315 | .grant( |
| 316 | "tools/list", |
| 317 | json!({"cursor":"queued"}), |
| 318 | Duration::from_secs(5), |
| 319 | None, |
| 320 | Some("raw-accepted-id"), |
| 321 | ) |
| 322 | .unwrap(); |
| 323 | manager |
| 324 | .shared |
| 325 | .core_calls |
| 326 | .revoke_host(HostTier::Builtin, transport.session.host_generation); |
| 327 | manager.shared.mcp_users.fetch_sub( |
| 328 | manager |
| 329 | .shared |
| 330 | .mcp_broker |
| 331 | .revoke_host(HostTier::Builtin, transport.session.host_generation), |
| 332 | std::sync::atomic::Ordering::SeqCst, |
| 333 | ); |
| 334 | let (cx, _, _) = HostRequestContext::for_test(12); |
| 335 | assert!( |
| 336 | manager |
| 337 | .shared |
| 338 | .mcp_broker |
| 339 | .serve( |
| 340 | &manager.shared, |
| 341 | transport.session.host_generation, |
| 342 | request(&transport, &grant, json!({"cursor":"queued"}), 12), |
| 343 | cx |
| 344 | ) |
| 345 | .await |
| 346 | .is_err() |
| 347 | ); |
| 348 | assert!(transport.probe_dead()); |
| 349 | assert!(!std::fs::read_to_string(record).unwrap().contains("queued")); |
| 350 | } |
| 351 | |
| 352 | #[tokio::test(flavor = "current_thread")] |
| 353 | async fn host_http_requires_rust_prepared_client_without_fallback_or_spawn() { |
| 354 | let _policy = TestPolicyGuard::extension_host(true); |
| 355 | let Some((_fixture, manager, _, _)) = fixture("host_backend_refuses_http").await else { |
| 356 | return; |
| 357 | }; |
| 358 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 359 | let config: McpServerConfig = |
| 360 | serde_json::from_value(json!({"url": "https://example.invalid/mcp"})).unwrap(); |
| 361 | let error = SdkTransport::connect( |
| 362 | "http", |
| 363 | &config, |
| 364 | CancellationToken::new(), |
| 365 | Duration::from_secs(5), |
| 366 | ) |
| 367 | .await |
| 368 | .err() |
| 369 | .unwrap(); |
| 370 | assert!(error.to_string().contains("Rust-prepared HTTP client")); |
| 371 | assert!( |
| 372 | manager |
| 373 | .shared |
| 374 | .mcp_broker |
| 375 | .sessions |
| 376 | .lock() |
| 377 | .unwrap() |
| 378 | .is_empty() |
| 379 | ); |
| 380 | assert_eq!( |
| 381 | manager |
| 382 | .shared |
| 383 | .mcp_users |
| 384 | .load(std::sync::atomic::Ordering::SeqCst), |
| 385 | 0 |
| 386 | ); |
| 387 | } |
| 388 | |
| 389 | #[tokio::test(flavor = "current_thread")] |
| 390 | async fn real_computer_use_bootstrap_and_human_decision_stay_in_rust() { |
| 391 | if !cfg!(target_os = "macos") { |
| 392 | return; |
| 393 | } |
| 394 | let _env = crate::test_support::lock_test_env(); |
| 395 | let _policy = TestPolicyGuard::extension_host(true); |
| 396 | let Some(node) = node_for_tests("real_computer_use_host_sdk") else { |
| 397 | return; |
| 398 | }; |
| 399 | let (root, _registry, pool, server) = crate::mcp::computer_use_test_fixture(); |
| 400 | let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", root.path()); |
| 401 | let _backend = crate::test_support::EnvVarGuard::set("CODEWHALE_SECRET_BACKEND", "file"); |
| 402 | let manager = Arc::new(ExtensionHostManager::new( |
| 403 | super::super::ExtensionHostOptions { |
| 404 | node_override: Some(node), |
| 405 | // Keep the bundle under the Codewhale home's readable |
| 406 | // extension-host entry; a nested `host` entry is sandbox-denied. |
| 407 | root: Some(root.path().to_path_buf()), |
| 408 | ..Default::default() |
| 409 | }, |
| 410 | )); |
| 411 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 412 | let mut pool = pool.with_backend(McpBackend::Host); |
| 413 | pool.get_or_connect(&server).await.unwrap(); |
| 414 | let session = manager |
| 415 | .shared |
| 416 | .mcp_broker |
| 417 | .sessions |
| 418 | .lock() |
| 419 | .unwrap() |
| 420 | .values() |
| 421 | .next() |
| 422 | .cloned() |
| 423 | .unwrap(); |
| 424 | assert!(session.decision_key.lock().unwrap().is_some()); |
| 425 | let input = json!({"action":"allow","scope":"foreground","remember":false}); |
| 426 | let tool = McpPool::mcp_model_tool_name(&server, "consent"); |
| 427 | let decode = |result: Value| { |
| 428 | serde_json::from_str::<Value>(result["content"][0]["text"].as_str().unwrap()).unwrap() |
| 429 | }; |
| 430 | let unsigned = decode( |
| 431 | pool.call_tool_with_decision(&tool, input.clone(), None) |
| 432 | .await |
| 433 | .unwrap(), |
| 434 | ); |
| 435 | assert_eq!(unsigned["error"]["code"], "consent_needs_user"); |
| 436 | let decision = HumanDecision::for_test(&tool, &input); |
| 437 | let approved = decode( |
| 438 | pool.call_tool_with_decision(&tool, input.clone(), Some(&decision)) |
| 439 | .await |
| 440 | .unwrap(), |
| 441 | ); |
| 442 | assert_eq!(approved["ok"], true); |
| 443 | assert!( |
| 444 | pool.call_tool_with_decision( |
| 445 | &tool, |
| 446 | json!({"action":"allow","scope":"foreground","remember":true}), |
| 447 | Some(&decision) |
| 448 | ) |
| 449 | .await |
| 450 | .is_err() |
| 451 | ); |
| 452 | // Ticket-visible params do not contain the attestation added at the pipe. |
| 453 | let transport = SdkTransport { |
| 454 | manager: Arc::clone(&manager), |
| 455 | host: manager.shared.ready_host(HostTier::Builtin).unwrap(), |
| 456 | session: Arc::clone(&session), |
| 457 | session_id: session.session_id.clone(), |
| 458 | replies: Default::default(), |
| 459 | discovery_timeout: Duration::from_secs(5), |
| 460 | }; |
| 461 | let grant = transport |
| 462 | .grant( |
| 463 | "tools/call", |
| 464 | json!({"name":"consent","arguments":input}), |
| 465 | Duration::from_secs(5), |
| 466 | Some(&decision), |
| 467 | Some("raw-decision-id"), |
| 468 | ) |
| 469 | .unwrap(); |
| 470 | assert!(grant.params.get("_meta").is_none()); |
| 471 | pool.shutdown_all().await; |
| 472 | manager.shutdown().await; |
| 473 | } |
| 474 | |
| 475 | #[tokio::test(flavor = "current_thread")] |
| 476 | async fn cancelled_or_aborted_partial_pipe_write_retires_contained_child() { |
| 477 | let _policy = TestPolicyGuard::extension_host(true); |
| 478 | for abort_handler in [false, true] { |
| 479 | let Some((_fixture, manager, record, mut config)) = fixture("partial_pipe_write").await |
| 480 | else { |
| 481 | return; |
| 482 | }; |
| 483 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 484 | let paused = r#" |
| 485 | import fs from 'node:fs'; const record=process.argv[2]; let pending=''; |
| 486 | function parse(chunk) { pending+=chunk; while(pending.includes('\n')) { |
| 487 | const at=pending.indexOf('\n'), m=JSON.parse(pending.slice(0,at)); pending=pending.slice(at+1); |
| 488 | if(m.method==='initialize') process.stdout.write(JSON.stringify({jsonrpc:'2.0',id:m.id,result:{protocolVersion:'2025-06-18',serverInfo:{name:'paused',version:'1'},capabilities:{tools:{}}}})+'\n'); |
| 489 | if(m.method==='notifications/initialized') {process.stdin.off('data',parse); process.stdin.once('data',()=>{fs.appendFileSync(record,'partial-byte-observed\n'); process.stdin.pause()}); return;} |
| 490 | }} |
| 491 | process.stdin.on('data',parse); |
| 492 | "#; |
| 493 | let path = std::path::PathBuf::from(&config.args[0]); |
| 494 | std::fs::write(path, paused).unwrap(); |
| 495 | config.connect_timeout = Some(5); |
| 496 | let transport = transport(&config).await; |
| 497 | let body = json!({"name":"echo","arguments":{"large":"x".repeat(2*1024*1024)}}); |
| 498 | let grant = transport |
| 499 | .grant( |
| 500 | "tools/call", |
| 501 | body.clone(), |
| 502 | Duration::from_secs(5), |
| 503 | None, |
| 504 | Some("raw-accepted-id"), |
| 505 | ) |
| 506 | .unwrap(); |
| 507 | let request = request(&transport, &grant, body, 91); |
| 508 | let (cx, _, cancel) = HostRequestContext::for_test(91); |
| 509 | let task_manager = Arc::clone(&manager); |
| 510 | let generation = transport.session.host_generation; |
| 511 | let task = tokio::spawn(async move { |
| 512 | task_manager |
| 513 | .shared |
| 514 | .mcp_broker |
| 515 | .serve(&task_manager.shared, generation, request, cx) |
| 516 | .await |
| 517 | }); |
| 518 | tokio::time::timeout(Duration::from_secs(3), async { |
| 519 | loop { |
| 520 | if std::fs::read_to_string(&record) |
| 521 | .unwrap() |
| 522 | .contains("partial-byte-observed") |
| 523 | { |
| 524 | break; |
| 525 | } |
| 526 | tokio::time::sleep(Duration::from_millis(5)).await; |
| 527 | } |
| 528 | }) |
| 529 | .await |
| 530 | .unwrap(); |
| 531 | if abort_handler { |
| 532 | task.abort(); |
| 533 | } else { |
| 534 | cancel.cancel(); |
| 535 | } |
| 536 | let result = tokio::time::timeout(Duration::from_secs(3), task) |
| 537 | .await |
| 538 | .unwrap(); |
| 539 | assert!(result.is_err() || result.unwrap().is_err()); |
| 540 | assert!(transport.session.cancel.is_cancelled()); |
| 541 | tokio::time::timeout(Duration::from_secs(3), async { |
| 542 | loop { |
| 543 | if transport.session.broker.lock().await.is_none() { |
| 544 | break; |
| 545 | } |
| 546 | tokio::time::sleep(Duration::from_millis(5)).await; |
| 547 | } |
| 548 | }) |
| 549 | .await |
| 550 | .unwrap(); |
| 551 | assert_eq!( |
| 552 | std::fs::read_to_string(record) |
| 553 | .unwrap() |
| 554 | .lines() |
| 555 | .filter(|line| *line == "partial-byte-observed") |
| 556 | .count(), |
| 557 | 1 |
| 558 | ); |
| 559 | } |
| 560 | } |
| 561 | |
| 562 | #[test] |
| 563 | fn production_process_surface_is_builtin_only_and_strict() { |
| 564 | let owner = OwnerRef { |
| 565 | plugin_id: "host:mcp".into(), |
| 566 | generation: 1, |
| 567 | owner_token: "token".into(), |
| 568 | }; |
| 569 | let frame = json!({"jsonrpc":"2.0","id":1,"method":"proc/write","params":{"owner":owner,"session_id":"session","ticket":"ticket","operation_id":"operation","frame":{"jsonrpc":"2.0","id":1,"method":"tools/list","params":{}}}}); |
| 570 | assert!(parse_host_message(frame.clone(), HostTier::Builtin).is_ok()); |
| 571 | assert!(parse_host_message(frame.clone(), HostTier::Plugin).is_err()); |
| 572 | let mut foreign = frame; |
| 573 | foreign["params"]["argv"] = json!(["unreviewed"]); |
| 574 | assert!(parse_host_message(foreign, HostTier::Builtin).is_err()); |
| 575 | } |
| 576 | |
| 577 | #[tokio::test(flavor = "current_thread")] |
| 578 | async fn host_restart_mints_fresh_owner_and_never_revives_old_grant() { |
| 579 | let _policy = TestPolicyGuard::extension_host(true); |
| 580 | let Some((_fixture, manager, _, config)) = fixture("host_restart_generation").await else { |
| 581 | return; |
| 582 | }; |
| 583 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 584 | let first = transport(&config).await; |
| 585 | let grant = first |
| 586 | .grant( |
| 587 | "tools/list", |
| 588 | json!({}), |
| 589 | Duration::from_secs(5), |
| 590 | None, |
| 591 | Some("raw-accepted-id"), |
| 592 | ) |
| 593 | .unwrap(); |
| 594 | // Hold the existing reconcile lock so an automatic scheduled restart |
| 595 | // cannot race this test's explicit retry. Status alone is insufficient: |
| 596 | // a dead Ready slot projects Restarting before the exit callback runs. |
| 597 | let reconcile = manager.shared.sync_lock.lock().await; |
| 598 | first.host.shutdown().await; |
| 599 | tokio::time::timeout(Duration::from_secs(5), async { |
| 600 | while matches!( |
| 601 | &*manager |
| 602 | .shared |
| 603 | .tier_runtime(HostTier::Builtin) |
| 604 | .host |
| 605 | .lock() |
| 606 | .expect("builtin slot lock"), |
| 607 | super::super::HostSlot::Ready(_) | super::super::HostSlot::Unresponsive(_) |
| 608 | ) { |
| 609 | tokio::time::sleep(Duration::from_millis(5)).await; |
| 610 | } |
| 611 | }) |
| 612 | .await |
| 613 | .expect("builtin exit callback withdrew the actual slot before retry"); |
| 614 | assert!(first.session.cancel.is_cancelled()); |
| 615 | manager.retry(); |
| 616 | drop(reconcile); |
| 617 | let second = transport(&config).await; |
| 618 | assert_ne!(first.session.owner, second.session.owner); |
| 619 | assert_ne!( |
| 620 | first.session.host_generation, |
| 621 | second.session.host_generation |
| 622 | ); |
| 623 | let (cx, _, _) = HostRequestContext::for_test(51); |
| 624 | assert!( |
| 625 | manager |
| 626 | .shared |
| 627 | .mcp_broker |
| 628 | .serve( |
| 629 | &manager.shared, |
| 630 | second.session.host_generation, |
| 631 | request(&first, &grant, json!({}), 51), |
| 632 | cx |
| 633 | ) |
| 634 | .await |
| 635 | .is_err() |
| 636 | ); |
| 637 | } |
| 638 | |
| 639 | #[tokio::test(flavor = "current_thread")] |
| 640 | async fn host_mcp_without_native_keeps_restart_policy_and_harness_closed() { |
| 641 | let _review_policy = TestPolicyGuard::extension_host(true); |
| 642 | let Some((fixture, manager, _, config)) = |
| 643 | fixture_with_plugins("host_mcp_without_native_restart", &["dsh-workspace-deps"]).await |
| 644 | else { |
| 645 | return; |
| 646 | }; |
| 647 | let plugins = fixture.registry(); |
| 648 | let _native_off = TestPolicyGuard::extension_host(false); |
| 649 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 650 | let attachment = manager.attach(plugins); |
| 651 | attachment.sync().await.unwrap(); |
| 652 | assert_eq!(manager.spawn_attempts(), 0, "config/attachment are lazy"); |
| 653 | |
| 654 | // Even a retained harness demand cannot admit a new harness with Native off. |
| 655 | manager.shared.harness_users.store(1, Ordering::SeqCst); |
| 656 | let first = transport(&config).await; |
| 657 | assert_eq!(manager.tier_spawn_attempts(HostTier::Builtin), 1); |
| 658 | assert_eq!(manager.tier_spawn_attempts(HostTier::Plugin), 0); |
| 659 | let replay_policy = manager.shared.builtin.supervision.lock().unwrap().policy; |
| 660 | assert!(!replay_policy, "MCP demand must not become Native policy"); |
| 661 | assert!( |
| 662 | manager |
| 663 | .shared |
| 664 | .registry |
| 665 | .lock() |
| 666 | .unwrap() |
| 667 | .owner("host:harness") |
| 668 | .is_none() |
| 669 | ); |
| 670 | assert!(manager.ensure_harness_builtin().await.is_err()); |
| 671 | |
| 672 | // Use the real host exit/owner withdrawal before the same replay method |
| 673 | // that schedule_restart invokes with its captured supervision policy. |
| 674 | let reconcile = manager.shared.sync_lock.lock().await; |
| 675 | first.host.shutdown().await; |
| 676 | tokio::time::timeout(Duration::from_secs(5), async { |
| 677 | while matches!( |
| 678 | &*manager.shared.builtin.host.lock().unwrap(), |
| 679 | super::super::HostSlot::Ready(_) | super::super::HostSlot::Unresponsive(_) |
| 680 | ) { |
| 681 | tokio::time::sleep(Duration::from_millis(5)).await; |
| 682 | } |
| 683 | }) |
| 684 | .await |
| 685 | .expect("actual Builtin exit callback withdraws the slot"); |
| 686 | assert!(first.session.cancel.is_cancelled()); |
| 687 | manager.retry(); |
| 688 | drop(reconcile); |
| 689 | manager.reconcile_with_policy(replay_policy).await.unwrap(); |
| 690 | assert_eq!(manager.tier_spawn_attempts(HostTier::Plugin), 0); |
| 691 | assert!( |
| 692 | !manager |
| 693 | .shared |
| 694 | .registry |
| 695 | .lock() |
| 696 | .unwrap() |
| 697 | .owners() |
| 698 | .any(|owner| owner.tier == HostTier::Plugin) |
| 699 | ); |
| 700 | let second = transport(&config).await; |
| 701 | assert_ne!(first.session.owner, second.session.owner); |
| 702 | assert!(!manager.shared.builtin.supervision.lock().unwrap().policy); |
| 703 | assert!( |
| 704 | manager |
| 705 | .shared |
| 706 | .registry |
| 707 | .lock() |
| 708 | .unwrap() |
| 709 | .owner("host:harness") |
| 710 | .is_none() |
| 711 | ); |
| 712 | manager.shared.harness_users.store(0, Ordering::SeqCst); |
| 713 | } |
| 714 | |
| 715 | #[tokio::test(flavor = "current_thread")] |
| 716 | async fn selected_host_without_node_uses_runtime_diagnostic_not_native_gate() { |
| 717 | let _policy = TestPolicyGuard::extension_host(false); |
| 718 | let root = tempfile::tempdir().unwrap(); |
| 719 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 720 | node_override: Some(root.path().join("missing-node")), |
| 721 | root: Some(root.path().join("home")), |
| 722 | ..Default::default() |
| 723 | })); |
| 724 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 725 | let config: McpServerConfig = |
| 726 | serde_json::from_value(json!({"command":"missing-peer"})).unwrap(); |
| 727 | let error = McpConnection::connect_with_backend( |
| 728 | "missing-runtime".into(), |
| 729 | config, |
| 730 | &McpTimeouts::default(), |
| 731 | None, |
| 732 | McpBackend::Host, |
| 733 | ) |
| 734 | .await |
| 735 | .err() |
| 736 | .expect("selected Host must fail without its runtime"); |
| 737 | let diagnostic = format!("{error:#}"); |
| 738 | assert!( |
| 739 | diagnostic.contains("Node.js ^22.19 || >=24"), |
| 740 | "{diagnostic}" |
| 741 | ); |
| 742 | assert!(diagnostic.contains("[extension_host] node"), "{diagnostic}"); |
| 743 | assert!(!diagnostic.contains("requires features.extension_host")); |
| 744 | assert_eq!(manager.tier_spawn_attempts(HostTier::Plugin), 0); |
| 745 | assert!( |
| 746 | manager |
| 747 | .shared |
| 748 | .mcp_broker |
| 749 | .sessions |
| 750 | .lock() |
| 751 | .unwrap() |
| 752 | .is_empty() |
| 753 | ); |
| 754 | } |
| 755 | |
| 756 | #[tokio::test(flavor = "current_thread")] |
| 757 | async fn explicit_rust_backend_without_native_never_starts_builtin_runtime() { |
| 758 | let _setup_policy = TestPolicyGuard::extension_host(true); |
| 759 | let Some((fixture, _, _, config)) = fixture("rust_backend_no_builtin").await else { |
| 760 | return; |
| 761 | }; |
| 762 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 763 | node_override: Some(fixture.root.join("missing-host-node")), |
| 764 | root: Some(fixture.root.clone()), |
| 765 | ..Default::default() |
| 766 | })); |
| 767 | let _native_off = TestPolicyGuard::extension_host(false); |
| 768 | let _manager = TestManagerGuard::install(Arc::clone(&manager)); |
| 769 | let connection = McpConnection::connect_with_backend( |
| 770 | "fixture".into(), |
| 771 | config, |
| 772 | &McpTimeouts::default(), |
| 773 | None, |
| 774 | McpBackend::Rust, |
| 775 | ) |
| 776 | .await |
| 777 | .unwrap(); |
| 778 | assert!(connection.is_ready()); |
| 779 | assert_eq!(connection.tools().len(), 1); |
| 780 | assert_eq!(manager.spawn_attempts(), 0); |
| 781 | } |
| 782 |