| 1 | //! Extension host tests. |
| 2 | //! |
| 3 | //! Unit tests (protocol corpus, registry rules, tool gating) need no Node. |
| 4 | //! Integration tests spawn the *committed* bundle under a real Node ≥22.19: |
| 5 | //! they skip with a printed reason when none is found, unless |
| 6 | //! `CODEWHALE_EXT_HOST_TESTS=1` is set (CI), where a missing Node fails. |
| 7 | |
| 8 | use std::collections::BTreeSet; |
| 9 | use std::path::{Path, PathBuf}; |
| 10 | use std::sync::Arc; |
| 11 | use std::time::{Duration, Instant}; |
| 12 | |
| 13 | use serde_json::{Value, json}; |
| 14 | |
| 15 | use super::command::{BuiltinCommandCatalog, BuiltinCommandsGuard}; |
| 16 | use super::protocol::{ |
| 17 | self, OwnerRef, RegisterKind, RegisterParams, RegisterSpecWire, parse_core_message, |
| 18 | parse_host_message, |
| 19 | }; |
| 20 | use super::registry::{OwnerRegistry, OwnerState}; |
| 21 | use super::tier::HostTier; |
| 22 | use super::{ExtensionHostManager, ExtensionHostOptions, HostAttachment, HostStatus}; |
| 23 | use crate::plugins::PluginRegistry; |
| 24 | use crate::plugins::activation::TestPolicyGuard; |
| 25 | use crate::plugins::discovery::{DiscoveryConfig, discover_with_config}; |
| 26 | use crate::tools::spec::{ApprovalRequirement, ToolContext, ToolError, ToolSpec}; |
| 27 | |
| 28 | /// The integration tests below run the host on Node unless they say |
| 29 | /// otherwise; the Bun ones pin Bun (`bun_for_tests`). |
| 30 | const NODE: crate::config::ExtensionHostRuntime = crate::config::ExtensionHostRuntime::Node; |
| 31 | |
| 32 | /// What the registry is told is a built-in command in the tests that are not |
| 33 | /// about the command table. Those that are (`commands::extension_host_tests`) |
| 34 | /// install the real one: this module may not depend on `crate::commands`. |
| 35 | #[derive(Debug)] |
| 36 | struct StubBuiltinCommands; |
| 37 | |
| 38 | impl BuiltinCommandCatalog for StubBuiltinCommands { |
| 39 | fn answers_to(&self, name: &str) -> bool { |
| 40 | matches!( |
| 41 | name, |
| 42 | "help" | "trust" | "model" | "jihua" | "zidong" | "stub-alias" |
| 43 | ) |
| 44 | } |
| 45 | } |
| 46 | |
| 47 | /// The stub catalog, for this thread until the guard drops. |
| 48 | fn stub_builtin_commands() -> BuiltinCommandsGuard { |
| 49 | BuiltinCommandsGuard::install(Arc::new(StubBuiltinCommands)) |
| 50 | } |
| 51 | |
| 52 | fn fixtures_dir() -> PathBuf { |
| 53 | Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/extension_host") |
| 54 | } |
| 55 | |
| 56 | // --------------------------------------------------------------------------- |
| 57 | // Protocol |
| 58 | // --------------------------------------------------------------------------- |
| 59 | |
| 60 | #[test] |
| 61 | fn protocol_corpus_parses_and_round_trips_in_both_directions() { |
| 62 | let dir = fixtures_dir().join("protocol"); |
| 63 | let mut entries: Vec<_> = std::fs::read_dir(&dir) |
| 64 | .expect("corpus dir") |
| 65 | .map(|entry| entry.expect("entry").path()) |
| 66 | .filter(|path| path.extension().is_some_and(|ext| ext == "json")) |
| 67 | .collect(); |
| 68 | entries.sort(); |
| 69 | let (mut valid, mut invalid) = (0, 0); |
| 70 | for path in entries { |
| 71 | let case: Value = serde_json::from_slice(&std::fs::read(&path).unwrap()).unwrap(); |
| 72 | let frame = case["frame"].clone(); |
| 73 | let direction = case["direction"].as_str().unwrap(); |
| 74 | let expect_valid = case["valid"].as_bool().unwrap(); |
| 75 | let name = path.file_name().unwrap().to_string_lossy(); |
| 76 | // Every method allows both tiers today, so a case reads the same under |
| 77 | // either (the tier rule itself is `protocol::tests`). |
| 78 | for tier in HostTier::ALL { |
| 79 | let reencoded = match direction { |
| 80 | "host_to_core" => parse_host_message(frame.clone(), tier).map(|m| m.to_value()), |
| 81 | "core_to_host" => parse_core_message(frame.clone(), tier).map(|m| m.to_value()), |
| 82 | other => panic!("unknown direction {other}"), |
| 83 | }; |
| 84 | if expect_valid { |
| 85 | let reencoded = reencoded.unwrap_or_else(|e| panic!("{name} ({tier:?}): {e}")); |
| 86 | assert_eq!( |
| 87 | reencoded, frame, |
| 88 | "{name} ({tier:?}) must round-trip exactly" |
| 89 | ); |
| 90 | } else { |
| 91 | assert!(reencoded.is_err(), "{name} ({tier:?}) must be rejected"); |
| 92 | } |
| 93 | } |
| 94 | if expect_valid { |
| 95 | let bytes = protocol::encode_frame(&frame).unwrap(); |
| 96 | assert_eq!(&bytes[..4], b"CWX1"); |
| 97 | valid += 1; |
| 98 | } else { |
| 99 | invalid += 1; |
| 100 | } |
| 101 | } |
| 102 | assert!( |
| 103 | valid >= 15 && invalid >= 8, |
| 104 | "corpus too small: {valid} valid, {invalid} invalid" |
| 105 | ); |
| 106 | } |
| 107 | |
| 108 | #[tokio::test] |
| 109 | async fn frames_decode_in_order_and_violations_are_typed() { |
| 110 | let mut bytes = protocol::encode_frame(&json!({"a": 1})).unwrap(); |
| 111 | bytes.extend(protocol::encode_frame(&json!({"b": "ü"})).unwrap()); |
| 112 | let mut reader = bytes.as_slice(); |
| 113 | assert_eq!( |
| 114 | protocol::read_frame(&mut reader).await.unwrap(), |
| 115 | Some(json!({"a": 1})) |
| 116 | ); |
| 117 | assert_eq!( |
| 118 | protocol::read_frame(&mut reader).await.unwrap(), |
| 119 | Some(json!({"b": "ü"})) |
| 120 | ); |
| 121 | assert_eq!(protocol::read_frame(&mut reader).await.unwrap(), None); |
| 122 | |
| 123 | let mut bad = b"NOPE\0\0\0\0".as_slice(); |
| 124 | assert!(matches!( |
| 125 | protocol::read_frame(&mut bad).await, |
| 126 | Err(protocol::FrameError::BadMagic) |
| 127 | )); |
| 128 | let mut huge = Vec::from(*b"CWX1"); |
| 129 | huge.extend(((protocol::MAX_FRAME + 1) as u32).to_le_bytes()); |
| 130 | assert!(matches!( |
| 131 | protocol::read_frame(&mut huge.as_slice()).await, |
| 132 | Err(protocol::FrameError::TooLarge(_)) |
| 133 | )); |
| 134 | } |
| 135 | |
| 136 | // --------------------------------------------------------------------------- |
| 137 | // Owner registry |
| 138 | // --------------------------------------------------------------------------- |
| 139 | |
| 140 | pub(crate) fn fake_authority(plugin_id: &str) -> crate::plugins::types::PluginAuthority { |
| 141 | crate::plugins::types::PluginAuthority { |
| 142 | plugin_id: crate::plugins::types::PluginId(plugin_id.to_string()), |
| 143 | plugin_name: plugin_id.to_string(), |
| 144 | workspace: PathBuf::from("/w"), |
| 145 | state_path: PathBuf::from("/s"), |
| 146 | source_manifest: PathBuf::from("/m"), |
| 147 | staged_manifest: PathBuf::from("/sm"), |
| 148 | content_hash: format!("hash-{plugin_id}"), |
| 149 | capability_hash: "cap".to_string(), |
| 150 | state_generation: 1, |
| 151 | } |
| 152 | } |
| 153 | |
| 154 | fn register(registry: &mut OwnerRegistry, owner: &OwnerRef, name: &str) -> Result<u64, String> { |
| 155 | registry.register_tool(&RegisterParams { |
| 156 | scope: None, |
| 157 | owner: owner.clone(), |
| 158 | kind: RegisterKind::Tool, |
| 159 | spec: RegisterSpecWire { |
| 160 | name: name.to_string(), |
| 161 | description: "d".to_string(), |
| 162 | input_schema: json!({"type": "object", "properties": {}}) |
| 163 | .as_object() |
| 164 | .cloned(), |
| 165 | argument_hint: None, |
| 166 | }, |
| 167 | }) |
| 168 | } |
| 169 | |
| 170 | #[test] |
| 171 | fn registry_refuses_shadowing_and_foreign_names_and_undoes_exactly_one_entry() { |
| 172 | let mut registry = OwnerRegistry::new(); |
| 173 | registry.add_native_names(["grep_files"]); |
| 174 | let a = registry |
| 175 | .begin_owner( |
| 176 | HostTier::Plugin, |
| 177 | "a", |
| 178 | "a", |
| 179 | Some(fake_authority("a")), |
| 180 | "hash-a", |
| 181 | ) |
| 182 | .unwrap(); |
| 183 | let b = registry |
| 184 | .begin_owner( |
| 185 | HostTier::Plugin, |
| 186 | "b", |
| 187 | "b", |
| 188 | Some(fake_authority("b")), |
| 189 | "hash-b", |
| 190 | ) |
| 191 | .unwrap(); |
| 192 | |
| 193 | // Built-ins (static, snapshot, and case-folded) and reserved prefixes. |
| 194 | for name in [ |
| 195 | "read_file", |
| 196 | "grep_files", |
| 197 | "READ", |
| 198 | "tool_search", |
| 199 | "mcp_x_y", |
| 200 | "ext_z", |
| 201 | ] { |
| 202 | assert!( |
| 203 | register(&mut registry, &a, name).is_err(), |
| 204 | "{name} must be refused" |
| 205 | ); |
| 206 | } |
| 207 | assert!(register(&mut registry, &a, "bad name").is_err()); |
| 208 | |
| 209 | let first = register(&mut registry, &a, "shared_name").unwrap(); |
| 210 | // Another owner cannot take it, in any case. |
| 211 | assert!(register(&mut registry, &b, "Shared_Name").is_err()); |
| 212 | // Same owner re-registering retires the old handle. |
| 213 | let second = register(&mut registry, &a, "shared_name").unwrap(); |
| 214 | assert_ne!(first, second); |
| 215 | registry.mark_active(&a); |
| 216 | registry.unregister(&a, first); // stale: must not remove the newer entry |
| 217 | assert_eq!(registry.live_tools().len(), 1); |
| 218 | assert!(registry.is_live(second, &a)); |
| 219 | // A foreign owner cannot unregister it either. |
| 220 | registry.unregister(&b, second); |
| 221 | assert!(registry.is_live(second, &a)); |
| 222 | |
| 223 | // A stale token is refused. |
| 224 | let mut stale = a.clone(); |
| 225 | stale.owner_token = "not-the-token".to_string(); |
| 226 | assert!(register(&mut registry, &stale, "other").is_err()); |
| 227 | |
| 228 | // Revocation is synchronous and total. |
| 229 | assert_eq!(registry.revoke_owner("a"), Some(a.clone())); |
| 230 | assert!(!registry.is_live(second, &a)); |
| 231 | assert!(registry.live_tools().is_empty()); |
| 232 | assert!(register(&mut registry, &a, "after_revoke").is_err()); |
| 233 | } |
| 234 | |
| 235 | /// Names that the approval tables key by name must never reach an extension: |
| 236 | /// a `fetch_url` session grant for github.com is `net:github.com`, and a |
| 237 | /// plugin tool called `web_fetch` would otherwise get that same key. |
| 238 | #[test] |
| 239 | fn registry_refuses_names_the_approval_tables_special_case() { |
| 240 | let mut registry = OwnerRegistry::new(); |
| 241 | let a = registry |
| 242 | .begin_owner( |
| 243 | HostTier::Plugin, |
| 244 | "a", |
| 245 | "a", |
| 246 | Some(fake_authority("a")), |
| 247 | "hash-a", |
| 248 | ) |
| 249 | .unwrap(); |
| 250 | // Special-cased by name somewhere in the approval path; some are also |
| 251 | // natives in some modes. |
| 252 | for name in [ |
| 253 | "web_fetch", |
| 254 | "exec_wait", |
| 255 | "exec_interact", |
| 256 | "task_shell_start", |
| 257 | "web_search", |
| 258 | "run_tests", |
| 259 | "run_verifiers", |
| 260 | "fim_edit", |
| 261 | "Bash", |
| 262 | "read_workspace_deps", |
| 263 | "list_things", |
| 264 | "get_secret", |
| 265 | "start_mcp_server", |
| 266 | ] { |
| 267 | let refused = |
| 268 | register(&mut registry, &a, name).expect_err(&format!("{name} must be refused")); |
| 269 | assert!( |
| 270 | refused.contains("reserved") || refused.contains("collides with a built-in"), |
| 271 | "{name}: {refused}" |
| 272 | ); |
| 273 | } |
| 274 | // Not natives in any mode: only the classifier probe refuses these. |
| 275 | for name in [ |
| 276 | "web_fetch", |
| 277 | "exec_wait", |
| 278 | "exec_interact", |
| 279 | "read_workspace_deps", |
| 280 | ] { |
| 281 | let refused = register(&mut registry, &a, name).unwrap_err(); |
| 282 | assert!(refused.contains("reserved"), "{name}: {refused}"); |
| 283 | } |
| 284 | // The fetch-family key really is shared by name: this is what the refusal |
| 285 | // protects. |
| 286 | let input = json!({"url": "https://github.com/x"}); |
| 287 | assert_eq!( |
| 288 | crate::tools::approval_cache::build_approval_grouping_key("web_fetch", &input), |
| 289 | crate::tools::approval_cache::build_approval_grouping_key("fetch_url", &input), |
| 290 | ); |
| 291 | // Opaque names are admitted and keyed as themselves. |
| 292 | for name in ["load_workspace_dependencies", "slow_wait", "probe_read"] { |
| 293 | register(&mut registry, &a, name).unwrap_or_else(|e| panic!("{name}: {e}")); |
| 294 | } |
| 295 | } |
| 296 | |
| 297 | #[test] |
| 298 | fn registry_enforces_schema_and_count_caps() { |
| 299 | let mut registry = OwnerRegistry::new(); |
| 300 | let a = registry |
| 301 | .begin_owner( |
| 302 | HostTier::Plugin, |
| 303 | "a", |
| 304 | "a", |
| 305 | Some(fake_authority("a")), |
| 306 | "hash-a", |
| 307 | ) |
| 308 | .unwrap(); |
| 309 | let mut params = RegisterParams { |
| 310 | scope: None, |
| 311 | owner: a.clone(), |
| 312 | kind: RegisterKind::Tool, |
| 313 | spec: RegisterSpecWire { |
| 314 | name: "big".to_string(), |
| 315 | description: "x".repeat(super::registry::MAX_DESCRIPTION_BYTES + 1), |
| 316 | input_schema: json!({"type": "object"}).as_object().cloned(), |
| 317 | argument_hint: None, |
| 318 | }, |
| 319 | }; |
| 320 | assert!( |
| 321 | registry |
| 322 | .register_tool(¶ms) |
| 323 | .unwrap_err() |
| 324 | .contains("description") |
| 325 | ); |
| 326 | params.spec.description = "ok".to_string(); |
| 327 | params.spec.input_schema = |
| 328 | json!({"type": "object", "description": "y".repeat(super::registry::MAX_SCHEMA_BYTES)}) |
| 329 | .as_object() |
| 330 | .cloned(); |
| 331 | assert!( |
| 332 | registry |
| 333 | .register_tool(¶ms) |
| 334 | .unwrap_err() |
| 335 | .contains("schema") |
| 336 | ); |
| 337 | params.spec.input_schema = json!({"type": "string"}).as_object().cloned(); |
| 338 | assert!( |
| 339 | registry |
| 340 | .register_tool(¶ms) |
| 341 | .unwrap_err() |
| 342 | .contains("object") |
| 343 | ); |
| 344 | for index in 0..super::registry::MAX_TOOLS_PER_OWNER { |
| 345 | register(&mut registry, &a, &format!("t{index}")).unwrap(); |
| 346 | } |
| 347 | assert!( |
| 348 | register(&mut registry, &a, "one_too_many") |
| 349 | .unwrap_err() |
| 350 | .contains("at most") |
| 351 | ); |
| 352 | } |
| 353 | |
| 354 | /// A tool registered through the real admission path, as the registry hands it |
| 355 | /// to `HostToolSpec`, with `schema` as its input schema. |
| 356 | fn admitted_tool(schema: Value) -> Result<super::registry::ToolRegistration, String> { |
| 357 | let mut registry = OwnerRegistry::new(); |
| 358 | let owner = registry |
| 359 | .begin_owner( |
| 360 | HostTier::Plugin, |
| 361 | "probe", |
| 362 | "probe", |
| 363 | Some(fake_authority("probe")), |
| 364 | "hash-probe", |
| 365 | ) |
| 366 | .unwrap(); |
| 367 | registry.register_tool(&RegisterParams { |
| 368 | scope: None, |
| 369 | owner: owner.clone(), |
| 370 | kind: RegisterKind::Tool, |
| 371 | spec: RegisterSpecWire { |
| 372 | name: "probe_tool".to_string(), |
| 373 | description: "d".to_string(), |
| 374 | input_schema: schema.as_object().cloned(), |
| 375 | argument_hint: None, |
| 376 | }, |
| 377 | })?; |
| 378 | assert!(registry.mark_active(&owner)); |
| 379 | Ok(registry.live_tools().remove(0)) |
| 380 | } |
| 381 | |
| 382 | #[test] |
| 383 | fn an_uncompilable_tool_schema_is_refused_at_registration_with_a_reason() { |
| 384 | for (label, schema) in [ |
| 385 | ( |
| 386 | "an unknown type", |
| 387 | json!({"type": "object", "properties": {"a": {"type": "nonsense"}}}), |
| 388 | ), |
| 389 | ( |
| 390 | "an external $ref the core will not fetch", |
| 391 | json!({"type": "object", "properties": {"a": {"$ref": "https://example.invalid/schema.json"}}}), |
| 392 | ), |
| 393 | ( |
| 394 | "a required list that is not a list", |
| 395 | json!({"type": "object", "required": "a"}), |
| 396 | ), |
| 397 | ] { |
| 398 | let reason = admitted_tool(schema).expect_err(label); |
| 399 | assert!( |
| 400 | reason.contains("probe_tool") && reason.contains("not a valid JSON Schema"), |
| 401 | "{label}: {reason}" |
| 402 | ); |
| 403 | } |
| 404 | admitted_tool(json!({"type": "object", "properties": {"a": {"type": "string"}}})) |
| 405 | .expect("a valid schema is admitted"); |
| 406 | } |
| 407 | |
| 408 | /// What the model gets back, and what the host never sees: the tool checks the |
| 409 | /// input against its registered schema in `prepare` (before any approval |
| 410 | /// card) and again in `execute`. |
| 411 | #[tokio::test] |
| 412 | async fn tool_input_is_checked_against_the_registered_schema_before_approval_and_execution() { |
| 413 | let registration = admitted_tool(json!({ |
| 414 | "type": "object", |
| 415 | "properties": { |
| 416 | "name": {"type": "string"}, |
| 417 | "count": {"type": "integer", "minimum": 0} |
| 418 | }, |
| 419 | "required": ["name"], |
| 420 | "additionalProperties": false |
| 421 | })) |
| 422 | .unwrap(); |
| 423 | // No host is running: a call that reached it would say so (`NotAvailable`). |
| 424 | let manager = ExtensionHostManager::new(ExtensionHostOptions::default()); |
| 425 | let tool = super::tool::HostToolSpec::new(registration, Arc::clone(&manager.shared)); |
| 426 | let context = ToolContext::new(Path::new("/w")); |
| 427 | |
| 428 | let rejected = |input: Value, expect: &[&str]| { |
| 429 | let error = tool.prepare(input.clone(), &context).unwrap_err(); |
| 430 | let ToolError::InvalidInput { message } = error else { |
| 431 | panic!("{input}: not an invalid-input error: {error:?}"); |
| 432 | }; |
| 433 | assert!( |
| 434 | message.contains("probe_tool") && expect.iter().all(|part| message.contains(part)), |
| 435 | "{input}: {message}" |
| 436 | ); |
| 437 | input |
| 438 | }; |
| 439 | // Valid input passes `prepare`, with the Rust-composed approval card. |
| 440 | let prepared = tool |
| 441 | .prepare(json!({"name": "x", "count": 2}), &context) |
| 442 | .expect("valid input is admitted"); |
| 443 | assert_eq!(prepared.approval, ApprovalRequirement::Required); |
| 444 | // An extra property. |
| 445 | let extra = rejected(json!({"name": "x", "surprise": true}), &["surprise"]); |
| 446 | // A wrong type, and the path of the field. |
| 447 | let wrong_type = rejected(json!({"name": 7}), &["name"]); |
| 448 | rejected(json!({"name": "x", "count": -1}), &["count"]); |
| 449 | // A missing required field, and a non-object. |
| 450 | rejected(json!({}), &["name"]); |
| 451 | rejected(json!("name"), &[]); |
| 452 | |
| 453 | // `execute` refuses the same inputs without touching the host (which is |
| 454 | // not running, so reaching it would be `NotAvailable`)... |
| 455 | for input in [extra, wrong_type] { |
| 456 | let error = tool.execute(input.clone(), &context).await.unwrap_err(); |
| 457 | assert!( |
| 458 | matches!(error, ToolError::InvalidInput { .. }), |
| 459 | "{input}: {error:?}" |
| 460 | ); |
| 461 | } |
| 462 | assert_eq!(manager.spawn_attempts(), 0); |
| 463 | // ...while valid input gets past the check and meets the dead host. |
| 464 | let error = tool |
| 465 | .execute(json!({"name": "x"}), &context) |
| 466 | .await |
| 467 | .unwrap_err(); |
| 468 | assert!(matches!(error, ToolError::NotAvailable { .. }), "{error:?}"); |
| 469 | } |
| 470 | |
| 471 | // --------------------------------------------------------------------------- |
| 472 | // Licence notices for the embedded bundle |
| 473 | // --------------------------------------------------------------------------- |
| 474 | |
| 475 | /// Every `node_modules/<package>` the bundle was built from, read from the |
| 476 | /// bundler's own `// node_modules/...` markers in the embedded bundle (not |
| 477 | /// from the generator that writes the notices). |
| 478 | fn bundled_packages() -> BTreeSet<String> { |
| 479 | let mut packages = BTreeSet::new(); |
| 480 | for bytes in [ |
| 481 | super::BUNDLE, |
| 482 | include_bytes!("../../extension-host/dist/builtin/mcp.mjs").as_slice(), |
| 483 | ] { |
| 484 | let text = std::str::from_utf8(bytes).expect("the bundle is UTF-8"); |
| 485 | for line in text.lines() { |
| 486 | let Some(path) = line.strip_prefix("// ") else { |
| 487 | continue; |
| 488 | }; |
| 489 | let Some((_, after)) = path.rsplit_once("node_modules/") else { |
| 490 | continue; |
| 491 | }; |
| 492 | let mut parts = after.split('/'); |
| 493 | let first = parts.next().unwrap_or_default(); |
| 494 | let name = if first.starts_with('@') { |
| 495 | format!("{first}/{}", parts.next().unwrap_or_default()) |
| 496 | } else { |
| 497 | first.to_string() |
| 498 | }; |
| 499 | packages.insert(name); |
| 500 | } |
| 501 | } |
| 502 | packages |
| 503 | } |
| 504 | |
| 505 | /// `name@version` of every package the embedded notices list. |
| 506 | fn noticed_packages() -> Vec<(String, String)> { |
| 507 | let text = std::str::from_utf8(super::NOTICES).expect("the notices are UTF-8"); |
| 508 | let listed = text |
| 509 | .split_once("\nPackages:\n") |
| 510 | .expect("the notices list their packages") |
| 511 | .1; |
| 512 | listed |
| 513 | .lines() |
| 514 | .take_while(|line| line.starts_with(" ")) |
| 515 | .map(|line| { |
| 516 | let entry = line.trim().split(" (").next().unwrap(); |
| 517 | let (name, version) = entry.rsplit_once('@').expect("name@version"); |
| 518 | (name.to_string(), version.to_string()) |
| 519 | }) |
| 520 | .collect() |
| 521 | } |
| 522 | |
| 523 | #[test] |
| 524 | fn materialized_bundle_directory_carries_its_licence_notices() { |
| 525 | let home = tempfile::tempdir().unwrap(); |
| 526 | let bundle = super::materialize_bundle(home.path()).unwrap(); |
| 527 | let dir = bundle.parent().unwrap().to_path_buf(); |
| 528 | assert_eq!( |
| 529 | dir, |
| 530 | super::supervisor::bundle_dir(home.path(), super::bundle_sha256()), |
| 531 | "the notices live in the bundle's own digest-named directory" |
| 532 | ); |
| 533 | let notices = dir.join("LICENSES.txt"); |
| 534 | assert_eq!(std::fs::read(¬ices).unwrap(), super::NOTICES); |
| 535 | assert_eq!(std::fs::read(&bundle).unwrap(), super::BUNDLE); |
| 536 | let builtin = dir.join("builtin/mcp.mjs"); |
| 537 | let builtin_bytes = include_bytes!("../../extension-host/dist/builtin/mcp.mjs"); |
| 538 | assert_eq!(std::fs::read(&builtin).unwrap(), builtin_bytes); |
| 539 | let pinned = super::tier::BUILTIN_MODULES |
| 540 | .iter() |
| 541 | .find(|module| module.id == "mcp") |
| 542 | .unwrap(); |
| 543 | assert_eq!( |
| 544 | super::hex(Sha256::digest(builtin_bytes)), |
| 545 | pinned.source_sha256 |
| 546 | ); |
| 547 | assert!( |
| 548 | std::str::from_utf8(super::NOTICES) |
| 549 | .unwrap() |
| 550 | .contains("Copyright (c) 2021-present Shigma"), |
| 551 | "the notices carry the licence text, not only package names" |
| 552 | ); |
| 553 | // Written like the bundle: read-only, no staging files left behind. |
| 554 | #[cfg(unix)] |
| 555 | for path in [&bundle, ¬ices, &builtin] { |
| 556 | use std::os::unix::fs::PermissionsExt; |
| 557 | assert_eq!( |
| 558 | std::fs::metadata(path).unwrap().permissions().mode() & 0o777, |
| 559 | 0o400, |
| 560 | "{}", |
| 561 | path.display() |
| 562 | ); |
| 563 | } |
| 564 | let mut names: Vec<_> = std::fs::read_dir(&dir) |
| 565 | .unwrap() |
| 566 | .map(|entry| entry.unwrap().file_name().to_string_lossy().into_owned()) |
| 567 | .collect(); |
| 568 | names.sort(); |
| 569 | assert_eq!( |
| 570 | names, |
| 571 | ["LICENSES.txt", "builtin", "codewhale-extension-host.mjs"] |
| 572 | ); |
| 573 | assert_eq!( |
| 574 | std::fs::read_dir(dir.join("builtin")).unwrap().count(), |
| 575 | super::tier::BUILTIN_MODULES.len() |
| 576 | ); |
| 577 | for module in super::tier::BUILTIN_MODULES { |
| 578 | let path = dir.join("builtin").join(format!("{}.mjs", module.id)); |
| 579 | let bytes = std::fs::read(&path).unwrap(); |
| 580 | assert_eq!(super::hex(Sha256::digest(&bytes)), module.source_sha256); |
| 581 | #[cfg(unix)] |
| 582 | { |
| 583 | use std::os::unix::fs::PermissionsExt; |
| 584 | assert_eq!( |
| 585 | std::fs::metadata(path).unwrap().permissions().mode() & 0o777, |
| 586 | 0o400 |
| 587 | ); |
| 588 | } |
| 589 | } |
| 590 | |
| 591 | // A directory from a build that wrote only the bundle gets its notices; a |
| 592 | // tampered or replaced notices file is rewritten, never trusted. |
| 593 | std::fs::remove_file(¬ices).unwrap(); |
| 594 | super::materialize_bundle(home.path()).unwrap(); |
| 595 | assert_eq!(std::fs::read(¬ices).unwrap(), super::NOTICES); |
| 596 | std::fs::remove_file(¬ices).unwrap(); |
| 597 | std::fs::write(¬ices, "tampered").unwrap(); |
| 598 | super::materialize_bundle(home.path()).unwrap(); |
| 599 | assert_eq!(std::fs::read(¬ices).unwrap(), super::NOTICES); |
| 600 | #[cfg(unix)] |
| 601 | { |
| 602 | let elsewhere = home.path().join("elsewhere.txt"); |
| 603 | std::fs::write(&elsewhere, "not the notices").unwrap(); |
| 604 | std::fs::remove_file(¬ices).unwrap(); |
| 605 | std::os::unix::fs::symlink(&elsewhere, ¬ices).unwrap(); |
| 606 | super::materialize_bundle(home.path()).unwrap(); |
| 607 | assert!( |
| 608 | !std::fs::symlink_metadata(¬ices) |
| 609 | .unwrap() |
| 610 | .file_type() |
| 611 | .is_symlink(), |
| 612 | "a symlink at the notices name is replaced, not followed" |
| 613 | ); |
| 614 | assert_eq!(std::fs::read(¬ices).unwrap(), super::NOTICES); |
| 615 | assert_eq!( |
| 616 | std::fs::read_to_string(&elsewhere).unwrap(), |
| 617 | "not the notices", |
| 618 | "the symlink's target is never written through" |
| 619 | ); |
| 620 | } |
| 621 | } |
| 622 | |
| 623 | #[test] |
| 624 | fn every_package_in_the_bundle_has_a_licence_notice_and_a_third_party_entry() { |
| 625 | let bundled = bundled_packages(); |
| 626 | assert!( |
| 627 | bundled.contains("@deepseek-ai/cordis"), |
| 628 | "the scan of the bundle found no packages: {bundled:?}" |
| 629 | ); |
| 630 | let noticed = noticed_packages(); |
| 631 | let third_party = std::fs::read_to_string( |
| 632 | Path::new(env!("CARGO_MANIFEST_DIR")).join("../../THIRD_PARTY_NOTICES.md"), |
| 633 | ) |
| 634 | .expect("THIRD_PARTY_NOTICES.md"); |
| 635 | let notice_text = std::str::from_utf8(super::NOTICES).unwrap(); |
| 636 | for package in &bundled { |
| 637 | let Some((_, version)) = noticed.iter().find(|(name, _)| name == package) else { |
| 638 | panic!("`{package}` is in the host bundle but not in dist/LICENSES.txt: {noticed:?}"); |
| 639 | }; |
| 640 | assert!( |
| 641 | third_party.contains(&format!("`{package}` {version}")), |
| 642 | "`{package}` {version} is bundled but THIRD_PARTY_NOTICES.md does not list it" |
| 643 | ); |
| 644 | let heading = format!("{package}@{version} ("); |
| 645 | assert!( |
| 646 | notice_text.matches(&heading).count() >= 2, |
| 647 | "dist/LICENSES.txt lists `{package}` but carries no section with its licence text" |
| 648 | ); |
| 649 | } |
| 650 | // The verbatim excerpts are not node_modules inputs; the notices still |
| 651 | // name them, and THIRD_PARTY_NOTICES.md must too. |
| 652 | for (package, version) in ¬iced { |
| 653 | assert!( |
| 654 | third_party.contains(&format!("`{package}` {version}")), |
| 655 | "`{package}` {version} is in dist/LICENSES.txt but not in THIRD_PARTY_NOTICES.md" |
| 656 | ); |
| 657 | } |
| 658 | } |
| 659 | |
| 660 | // --------------------------------------------------------------------------- |
| 661 | // Integration: real Node, real bundle |
| 662 | // --------------------------------------------------------------------------- |
| 663 | |
| 664 | /// A Node for the integration tests, or `None` (skip) when there is none and |
| 665 | /// the tests were not explicitly required. |
| 666 | pub(crate) fn node_for_tests(test: &str) -> Option<PathBuf> { |
| 667 | let resolution = crate::dependencies::resolve_extension_host_runtime( |
| 668 | crate::config::ExtensionHostRuntime::Node, |
| 669 | None, |
| 670 | None, |
| 671 | ); |
| 672 | match resolution.selected.as_ref() { |
| 673 | Some(runtime) => Some(runtime.path.clone()), |
| 674 | None if std::env::var_os("CODEWHALE_EXT_HOST_TESTS").is_some() => panic!( |
| 675 | "{test}: CODEWHALE_EXT_HOST_TESTS is set but no Node ^22.19 || >=24 was found: {}", |
| 676 | resolution.failure() |
| 677 | ), |
| 678 | None => { |
| 679 | eprintln!( |
| 680 | "skipping {test}: no Node ^22.19 || >=24 ({})", |
| 681 | resolution.failure() |
| 682 | ); |
| 683 | None |
| 684 | } |
| 685 | } |
| 686 | } |
| 687 | |
| 688 | /// Fixture plugins installed into a private user plugin dir through the |
| 689 | /// reviewed installer (`plugins::install`, local path), then reviewed |
| 690 | /// (trusted) and enabled through the real registry. Callers must hold a |
| 691 | /// `TestPolicyGuard::extension_host(true)` on this thread. |
| 692 | pub(crate) struct FixturePlugins { |
| 693 | _temp: tempfile::TempDir, |
| 694 | pub config: DiscoveryConfig, |
| 695 | pub root: PathBuf, |
| 696 | /// The stub built-in command catalog, so a fixture plugin's commands can |
| 697 | /// register; a test about the real table installs its own after this. |
| 698 | _commands: BuiltinCommandsGuard, |
| 699 | } |
| 700 | |
| 701 | impl FixturePlugins { |
| 702 | pub(crate) async fn new(names: &[&str]) -> Self { |
| 703 | use crate::plugins::install::{ |
| 704 | DEFAULT_MAX_SIZE_BYTES, PluginInstallOutcome, PluginInstallSource, install, |
| 705 | }; |
| 706 | let temp = tempfile::tempdir().unwrap(); |
| 707 | let workspace = temp.path().join("project"); |
| 708 | let user = temp.path().join("user"); |
| 709 | std::fs::create_dir_all(&workspace).unwrap(); |
| 710 | for name in names { |
| 711 | let outcome = install( |
| 712 | PluginInstallSource::LocalPath(fixtures_dir().join(name)), |
| 713 | &user, |
| 714 | DEFAULT_MAX_SIZE_BYTES, |
| 715 | &crate::network_policy::NetworkPolicy::default(), |
| 716 | false, |
| 717 | &|_| None, |
| 718 | ) |
| 719 | .await |
| 720 | .unwrap_or_else(|error| panic!("install {name}: {error:#}")); |
| 721 | assert!( |
| 722 | matches!(outcome, PluginInstallOutcome::Installed(ref installed) if installed.name == *name), |
| 723 | "install {name}: {outcome:?}" |
| 724 | ); |
| 725 | } |
| 726 | let config = DiscoveryConfig { |
| 727 | workspace: workspace.clone(), |
| 728 | user_plugins_dir: user, |
| 729 | workspace_plugins_dir: workspace.join(".codewhale/plugins"), |
| 730 | builtin_plugin_dirs: Vec::new(), |
| 731 | state_path: temp.path().join("state/plugin-state.json"), |
| 732 | }; |
| 733 | let mut registry = discover_with_config(&config); |
| 734 | for name in names { |
| 735 | registry |
| 736 | .trust(name) |
| 737 | .unwrap_or_else(|e| panic!("trust {name}: {e}")); |
| 738 | registry |
| 739 | .enable(name) |
| 740 | .unwrap_or_else(|e| panic!("enable {name}: {e}")); |
| 741 | } |
| 742 | let root = temp.path().join("home"); |
| 743 | let fixture = Self { |
| 744 | _temp: temp, |
| 745 | config, |
| 746 | root, |
| 747 | _commands: stub_builtin_commands(), |
| 748 | }; |
| 749 | let registry = fixture.registry(); |
| 750 | for name in names { |
| 751 | assert!( |
| 752 | registry.is_active(name), |
| 753 | "{name} must be active under policy v4" |
| 754 | ); |
| 755 | } |
| 756 | fixture |
| 757 | } |
| 758 | |
| 759 | pub(crate) fn registry(&self) -> Arc<PluginRegistry> { |
| 760 | Arc::new(discover_with_config(&self.config)) |
| 761 | } |
| 762 | |
| 763 | pub(crate) fn disable(&self, name: &str) -> Arc<PluginRegistry> { |
| 764 | let mut registry = discover_with_config(&self.config); |
| 765 | registry.disable(name).unwrap(); |
| 766 | self.registry() |
| 767 | } |
| 768 | |
| 769 | pub(crate) fn workspace(&self) -> &Path { |
| 770 | &self.config.workspace |
| 771 | } |
| 772 | |
| 773 | pub(crate) fn manager(&self, node: PathBuf) -> Arc<ExtensionHostManager> { |
| 774 | Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 775 | runtime: NODE, |
| 776 | node_override: Some(node), |
| 777 | root: Some(self.root.clone()), |
| 778 | ..Default::default() |
| 779 | })) |
| 780 | } |
| 781 | } |
| 782 | |
| 783 | pub(crate) fn host_tool( |
| 784 | engine: &HostAttachment, |
| 785 | workspace: &Path, |
| 786 | name: &str, |
| 787 | ) -> Arc<dyn ToolSpec> { |
| 788 | let mut registry = crate::tools::registry::ToolRegistryBuilder::new() |
| 789 | .build(ToolContext::new(workspace).with_plugin_registry(engine.plugin_view())); |
| 790 | let installed = engine.install_tools(&mut registry); |
| 791 | assert!( |
| 792 | installed.contains(&name.to_string()), |
| 793 | "{name} not installed: {installed:?}" |
| 794 | ); |
| 795 | registry.get(name).unwrap() |
| 796 | } |
| 797 | |
| 798 | /// Start budget for `dsh_plugin_runs_end_to_end_behind_the_approval_gate`: |
| 799 | /// bundle materialization, the runtime probe, the launch plan (on Linux, the |
| 800 | /// bwrap probe), spawn, handshake and one plugin's activation, under |
| 801 | /// nextest's full-core load on every CI OS. The host alone is ready in |
| 802 | /// ~50 ms (the JS suite's `READY_BUDGET_MS` gates that), but Windows CI |
| 803 | /// under load has taken over 5 s to hand-shake (`HANDSHAKE_DEADLINE`). |
| 804 | /// 20 s fails a start drifting toward that 30 s deadline, whose miss |
| 805 | /// disables every extension for the session, without failing on CI load. |
| 806 | const HOST_START_BUDGET: Duration = Duration::from_secs(20); |
| 807 | |
| 808 | /// Resident-size budget for a host with one plugin active, summed over its |
| 809 | /// process tree (bwrap's two processes included on Linux). 160 MiB is 2.4x |
| 810 | /// the largest idle host measured (67 MB, Node in a Linux container: |
| 811 | /// `supervisor::HOST_MEMORY_CAP`; 57 MiB Node 26 and 41 MiB Bun 1.4 with this |
| 812 | /// plugin on macOS arm64, 2026-09-30) and far below the 1 GiB cap. Not |
| 813 | /// checked where there is no `ps` (Windows). |
| 814 | const HOST_RSS_BUDGET_MIB: u64 = 160; |
| 815 | |
| 816 | /// Resident size in KiB of `pid` and every process under it — under bwrap, |
| 817 | /// `pid` is bwrap's and the runtime two levels down — from `ps`; `None` where |
| 818 | /// `ps` is missing or does not list `pid`. |
| 819 | fn tree_rss_kib(pid: u32) -> Option<u64> { |
| 820 | let output = std::process::Command::new("ps") |
| 821 | .args(["-A", "-o", "pid=,ppid=,rss="]) |
| 822 | .output() |
| 823 | .ok()?; |
| 824 | let rows: Vec<[u64; 3]> = String::from_utf8_lossy(&output.stdout) |
| 825 | .lines() |
| 826 | .filter_map(|line| { |
| 827 | let fields: Vec<u64> = line |
| 828 | .split_whitespace() |
| 829 | .filter_map(|field| field.parse().ok()) |
| 830 | .collect(); |
| 831 | <[u64; 3]>::try_from(fields).ok() |
| 832 | }) |
| 833 | .collect(); |
| 834 | let root = u64::from(pid); |
| 835 | if !rows.iter().any(|row| row[0] == root) { |
| 836 | return None; |
| 837 | } |
| 838 | let mut members = vec![root]; |
| 839 | let mut next = 0; |
| 840 | while let Some(&parent) = members.get(next) { |
| 841 | members.extend( |
| 842 | rows.iter() |
| 843 | .filter(|row| row[1] == parent && row[0] != parent) |
| 844 | .map(|row| row[0]), |
| 845 | ); |
| 846 | next += 1; |
| 847 | } |
| 848 | Some( |
| 849 | rows.iter() |
| 850 | .filter(|row| members.contains(&row[0])) |
| 851 | .map(|row| row[2]) |
| 852 | .sum(), |
| 853 | ) |
| 854 | } |
| 855 | |
| 856 | #[tokio::test] |
| 857 | async fn dsh_plugin_runs_end_to_end_behind_the_approval_gate() { |
| 858 | let Some(node) = node_for_tests("dsh_plugin_runs_end_to_end_behind_the_approval_gate") else { |
| 859 | return; |
| 860 | }; |
| 861 | let _policy = TestPolicyGuard::extension_host(true); |
| 862 | let fixture = FixturePlugins::new(&["dsh-workspace-deps"]).await; |
| 863 | let manager = fixture.manager(node); |
| 864 | assert_eq!( |
| 865 | manager.status(), |
| 866 | HostStatus::Idle, |
| 867 | "nothing starts before sync" |
| 868 | ); |
| 869 | |
| 870 | let started = Instant::now(); |
| 871 | let engine = manager.attach(fixture.registry()); |
| 872 | engine.sync().await.unwrap(); |
| 873 | let elapsed = started.elapsed(); |
| 874 | let pid = manager.host_pid().expect("host running"); |
| 875 | let rss = tree_rss_kib(pid); |
| 876 | eprintln!( |
| 877 | "extension host: spawn + handshake + activation {:.1} ms (budget {HOST_START_BUDGET:?}); RSS {} KiB (budget {HOST_RSS_BUDGET_MIB} MiB)", |
| 878 | elapsed.as_secs_f64() * 1000.0, |
| 879 | rss.map_or_else(|| "?".to_string(), |kib| kib.to_string()) |
| 880 | ); |
| 881 | assert!( |
| 882 | elapsed <= HOST_START_BUDGET, |
| 883 | "the extension host took {elapsed:?} to start, over its {HOST_START_BUDGET:?} budget" |
| 884 | ); |
| 885 | match rss { |
| 886 | Some(kib) => assert!( |
| 887 | kib <= HOST_RSS_BUDGET_MIB * 1024, |
| 888 | "the extension host is {kib} KiB resident, over its {HOST_RSS_BUDGET_MIB} MiB budget" |
| 889 | ), |
| 890 | None => eprintln!("extension host RSS budget not checked: no `ps` listing here"), |
| 891 | } |
| 892 | assert_eq!( |
| 893 | manager.live_tool_names(), |
| 894 | vec!["load_workspace_dependencies"] |
| 895 | ); |
| 896 | assert_eq!( |
| 897 | manager.owner_state( |
| 898 | fixture |
| 899 | .registry() |
| 900 | .get("dsh-workspace-deps") |
| 901 | .unwrap() |
| 902 | .id |
| 903 | .as_str() |
| 904 | ), |
| 905 | Some(OwnerState::Active) |
| 906 | ); |
| 907 | |
| 908 | let tool = host_tool(&engine, fixture.workspace(), "load_workspace_dependencies"); |
| 909 | assert_eq!(tool.registration_origin(), "extension:dsh-workspace-deps"); |
| 910 | // The plugin declares `presentCall: kind 'read'`; approval stays Required. |
| 911 | assert_eq!( |
| 912 | tool.approval_requirement_for(&json!({})), |
| 913 | ApprovalRequirement::Required |
| 914 | ); |
| 915 | assert!(!tool.is_read_only_for(&json!({}))); |
| 916 | assert!(tool.defer_loading()); |
| 917 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()); |
| 918 | let prepared = tool.prepare(json!({}), &context).unwrap(); |
| 919 | assert_eq!(prepared.approval, ApprovalRequirement::Required); |
| 920 | assert!( |
| 921 | prepared |
| 922 | .description |
| 923 | .contains("extension:dsh-workspace-deps") |
| 924 | ); |
| 925 | |
| 926 | let result = tool.execute(json!({}), &context).await.unwrap(); |
| 927 | assert!(result.success); |
| 928 | let payload: Value = serde_json::from_str(&result.content).unwrap(); |
| 929 | assert_eq!(payload["pythonDistributions"]["numpy"], "2.1.0"); |
| 930 | assert!(payload["python"].as_str().unwrap().contains("dependencies")); |
| 931 | manager.shutdown().await; |
| 932 | } |
| 933 | |
| 934 | /// A manifest may declare several `native` entries. They activate under one |
| 935 | /// owner, in order. Entry-scoped failure retires that entry without taking |
| 936 | /// away a sibling selected by another caller. |
| 937 | #[tokio::test] |
| 938 | async fn a_plugin_with_two_native_entries_activates_both_under_one_owner() { |
| 939 | let Some(node) = node_for_tests("a_plugin_with_two_native_entries") else { |
| 940 | return; |
| 941 | }; |
| 942 | let _policy = TestPolicyGuard::extension_host(true); |
| 943 | let fixture = FixturePlugins::new(&["two-entries", "two-entries-failing"]).await; |
| 944 | let manager = fixture.manager(node); |
| 945 | let engine = manager.attach(fixture.registry()); |
| 946 | engine.sync().await.unwrap(); |
| 947 | let registry = fixture.registry(); |
| 948 | let id = |name: &str| registry.get(name).unwrap().id.as_str().to_string(); |
| 949 | |
| 950 | // Both entries are live under the one owner. |
| 951 | assert_eq!( |
| 952 | manager.owner_state(&id("two-entries")), |
| 953 | Some(OwnerState::Active) |
| 954 | ); |
| 955 | let mut tools = manager.live_tool_names(); |
| 956 | tools.sort(); |
| 957 | assert_eq!(tools, ["tef_first", "two_first", "two_second"]); |
| 958 | assert_eq!(manager.live_command_names(), ["two-hello"]); |
| 959 | let report = manager.owner_report(&id("two-entries")).unwrap(); |
| 960 | assert!( |
| 961 | report.diagnostics.iter().any( |
| 962 | |line| line.contains("tools: two_first, two_second") && line.contains("/two-hello") |
| 963 | ), |
| 964 | "{:?}", |
| 965 | report.diagnostics |
| 966 | ); |
| 967 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()); |
| 968 | for (name, answer) in [("two_first", "first"), ("two_second", "second")] { |
| 969 | let result = host_tool(&engine, fixture.workspace(), name) |
| 970 | .execute(json!({}), &context) |
| 971 | .await |
| 972 | .unwrap(); |
| 973 | assert_eq!(result.content, answer); |
| 974 | } |
| 975 | |
| 976 | // The failing entry is retired while its already-active sibling survives. |
| 977 | let failing = id("two-entries-failing"); |
| 978 | assert_eq!(manager.owner_state(&failing), Some(OwnerState::Active)); |
| 979 | assert!(manager.live_tool_names().contains(&"tef_first".to_string())); |
| 980 | assert!( |
| 981 | manager |
| 982 | .shared |
| 983 | .registry |
| 984 | .lock() |
| 985 | .unwrap() |
| 986 | .owner(&failing) |
| 987 | .unwrap() |
| 988 | .scopes |
| 989 | .values() |
| 990 | .any(|state| matches!(state, OwnerState::Failed(_))) |
| 991 | ); |
| 992 | // The healthy plugin sharing the host is untouched, and a later reconcile |
| 993 | // does not retry the failed bytes. |
| 994 | engine.sync().await.unwrap(); |
| 995 | let mut tools = manager.live_tool_names(); |
| 996 | tools.sort(); |
| 997 | assert_eq!(tools, ["tef_first", "two_first", "two_second"]); |
| 998 | manager.shutdown().await; |
| 999 | } |
| 1000 | |
| 1001 | fn plugin_settings( |
| 1002 | name: &str, |
| 1003 | config: &str, |
| 1004 | ) -> std::collections::BTreeMap<String, crate::config::PluginSettings> { |
| 1005 | std::collections::BTreeMap::from([( |
| 1006 | name.to_string(), |
| 1007 | crate::config::PluginSettings { |
| 1008 | config: Some(toml::from_str(config).expect("config TOML")), |
| 1009 | }, |
| 1010 | )]) |
| 1011 | } |
| 1012 | |
| 1013 | /// Run an extension tool and parse its JSON answer. |
| 1014 | async fn call_json( |
| 1015 | tool: &Arc<dyn ToolSpec>, |
| 1016 | input: Value, |
| 1017 | workspace: &Path, |
| 1018 | engine: &HostAttachment, |
| 1019 | ) -> Value { |
| 1020 | let result = tool |
| 1021 | .execute( |
| 1022 | input, |
| 1023 | &ToolContext::new(workspace).with_plugin_registry(engine.plugin_view()), |
| 1024 | ) |
| 1025 | .await |
| 1026 | .unwrap_or_else(|error| panic!("{error:?}")); |
| 1027 | serde_json::from_str(&result.content).expect("a JSON answer") |
| 1028 | } |
| 1029 | |
| 1030 | /// The plugin's settings come from the user's config file, are bounded, and |
| 1031 | /// the keys (never the values) are what `/plugin show` can list. |
| 1032 | #[test] |
| 1033 | fn plugin_settings_are_read_from_the_user_config_and_bounded() { |
| 1034 | use super::plugin_config::{MAX_PLUGIN_CONFIG_BYTES, PluginConfigs}; |
| 1035 | let dir = tempfile::tempdir().unwrap(); |
| 1036 | let path = dir.path().join("config.toml"); |
| 1037 | // Other keys are none of this reader's business; `[plugins]` is. |
| 1038 | std::fs::write( |
| 1039 | &path, |
| 1040 | "model = \"x\"\n[plugins.\"greeter\".config]\ngreeting = \"Hi\"\nlimit = 3\n[plugins.\"greeter\".config.nested]\ntags = [\"a\", \"b\"]\n[plugins.\"quiet\"]\n", |
| 1041 | ) |
| 1042 | .unwrap(); |
| 1043 | let settings = crate::config::read_plugin_settings(&path).unwrap(); |
| 1044 | let mut configs = PluginConfigs::default(); |
| 1045 | configs.replace(&settings, Some(path.clone())); |
| 1046 | assert_eq!(configs.source(), Some(path.clone())); |
| 1047 | let greeter = configs.select("greeter").unwrap(); |
| 1048 | assert_eq!( |
| 1049 | greeter.value, |
| 1050 | json!({"greeting": "Hi", "limit": 3, "nested": {"tags": ["a", "b"]}}) |
| 1051 | ); |
| 1052 | // Keys only, in order; a plugin with a table but no config has nothing to list. |
| 1053 | assert_eq!( |
| 1054 | configs.summary("greeter").unwrap().unwrap(), |
| 1055 | ["greeting", "limit", "nested"] |
| 1056 | ); |
| 1057 | assert!(configs.summary("quiet").is_none()); |
| 1058 | assert!(configs.summary("nobody").is_none()); |
| 1059 | // No settings is an empty object, with a digest of its own. |
| 1060 | let none = configs.select("nobody").unwrap(); |
| 1061 | assert_eq!(none.value, json!({})); |
| 1062 | assert_ne!(none.hash, greeter.hash); |
| 1063 | // The digest follows the value: equal for the same settings, new for a change. |
| 1064 | let mut again = PluginConfigs::default(); |
| 1065 | again.replace(&settings, None); |
| 1066 | assert_eq!(again.select("greeter").unwrap().hash, greeter.hash); |
| 1067 | assert_eq!(again.source(), None, "a reload names no new source"); |
| 1068 | configs.replace_reloaded(&plugin_settings("greeter", "greeting = \"Yo\"")); |
| 1069 | assert_ne!(configs.select("greeter").unwrap().hash, greeter.hash); |
| 1070 | assert_eq!( |
| 1071 | configs.source(), |
| 1072 | Some(path.clone()), |
| 1073 | "the source survives a reload" |
| 1074 | ); |
| 1075 | |
| 1076 | // A missing file has no settings; an unreadable shape or a stray key fails loudly. |
| 1077 | assert!( |
| 1078 | crate::config::read_plugin_settings(&dir.path().join("absent.toml")) |
| 1079 | .unwrap() |
| 1080 | .is_empty() |
| 1081 | ); |
| 1082 | std::fs::write(&path, "[plugins.\"greeter\"]\nenabled = true\n").unwrap(); |
| 1083 | assert!( |
| 1084 | crate::config::read_plugin_settings(&path) |
| 1085 | .unwrap_err() |
| 1086 | .contains("cannot parse") |
| 1087 | ); |
| 1088 | |
| 1089 | // Refusals name the plugin and the rule; the config is never half-delivered. |
| 1090 | let refuse = |config: &str| -> String { |
| 1091 | let mut configs = PluginConfigs::default(); |
| 1092 | configs.replace(&plugin_settings("greeter", config), None); |
| 1093 | configs.select("greeter").unwrap_err() |
| 1094 | }; |
| 1095 | let big = refuse(&format!( |
| 1096 | "blob = \"{}\"", |
| 1097 | "x".repeat(MAX_PLUGIN_CONFIG_BYTES) |
| 1098 | )); |
| 1099 | assert!( |
| 1100 | big.contains("greeter") && big.contains("byte limit"), |
| 1101 | "{big}" |
| 1102 | ); |
| 1103 | let at_limit = format!("blob = \"{}\"", "x".repeat(MAX_PLUGIN_CONFIG_BYTES - 20)); |
| 1104 | let mut ok = PluginConfigs::default(); |
| 1105 | ok.replace(&plugin_settings("greeter", &at_limit), None); |
| 1106 | assert!( |
| 1107 | ok.select("greeter").is_ok(), |
| 1108 | "just under the cap is accepted" |
| 1109 | ); |
| 1110 | let when = refuse("when = 1979-05-27T07:32:00Z"); |
| 1111 | assert!(when.contains("date-time"), "{when}"); |
| 1112 | let deep = refuse(&format!("x = {}1{}", "[".repeat(20), "]".repeat(20))); |
| 1113 | assert!(deep.contains("nests deeper"), "{deep}"); |
| 1114 | // A refused config has a digest too, so it is not retried every turn but is once the file changes. |
| 1115 | let refused: Result<super::plugin_config::PluginConfig, String> = Err(big); |
| 1116 | assert!(super::plugin_config::activation_hash(&refused).starts_with("refused:")); |
| 1117 | } |
| 1118 | |
| 1119 | /// The plugin's context in a real host: its settings (checked by its own |
| 1120 | /// `Config` schema), the workspace of each call, and its own data directory; |
| 1121 | /// a change of settings re-activates it under a new generation; a refused one |
| 1122 | /// fails it with the reason. |
| 1123 | #[tokio::test] |
| 1124 | async fn plugin_context_reaches_the_plugin_and_changed_settings_reactivate_it() { |
| 1125 | let Some(node) = node_for_tests("plugin_context_reaches_the_plugin") else { |
| 1126 | return; |
| 1127 | }; |
| 1128 | let _policy = TestPolicyGuard::extension_host(true); |
| 1129 | let fixture = FixturePlugins::new(&["plugin-context"]).await; |
| 1130 | let manager = fixture.manager(node); |
| 1131 | let _manager = super::TestManagerGuard::install(Arc::clone(&manager)); |
| 1132 | let id = fixture |
| 1133 | .registry() |
| 1134 | .get("plugin-context") |
| 1135 | .unwrap() |
| 1136 | .id |
| 1137 | .as_str() |
| 1138 | .to_string(); |
| 1139 | let generation = |manager: &ExtensionHostManager| { |
| 1140 | manager |
| 1141 | .shared |
| 1142 | .registry |
| 1143 | .lock() |
| 1144 | .unwrap() |
| 1145 | .owner(&id) |
| 1146 | .map(|entry| entry.owner.generation) |
| 1147 | }; |
| 1148 | manager.set_plugin_settings( |
| 1149 | &plugin_settings("plugin-context", "greeting = \"Hi\""), |
| 1150 | None, |
| 1151 | ); |
| 1152 | let engine = manager.attach(fixture.registry()); |
| 1153 | engine.sync().await.unwrap(); |
| 1154 | assert_eq!(manager.owner_state(&id), Some(OwnerState::Active)); |
| 1155 | let first_generation = generation(&manager).unwrap(); |
| 1156 | |
| 1157 | let data_dir = super::supervisor::plugin_data_dir(&fixture.root, &id, "plugin-context"); |
| 1158 | assert!( |
| 1159 | data_dir.is_dir(), |
| 1160 | "the core made the plugin's directory before activation" |
| 1161 | ); |
| 1162 | assert!( |
| 1163 | data_dir.starts_with(fixture.root.join("extension-host/data/plugins")), |
| 1164 | "{}", |
| 1165 | data_dir.display() |
| 1166 | ); |
| 1167 | #[cfg(unix)] |
| 1168 | { |
| 1169 | use std::os::unix::fs::PermissionsExt; |
| 1170 | assert_eq!( |
| 1171 | std::fs::metadata(&data_dir).unwrap().permissions().mode() & 0o777, |
| 1172 | 0o700 |
| 1173 | ); |
| 1174 | } |
| 1175 | |
| 1176 | let probe = host_tool(&engine, fixture.workspace(), "ctx_probe"); |
| 1177 | let seen = call_json(&probe, json!({}), fixture.workspace(), &engine).await; |
| 1178 | // The plugin's own schema supplied the default for `limit`. |
| 1179 | assert_eq!(seen["config"], json!({"greeting": "Hi", "limit": 3})); |
| 1180 | assert_eq!(seen["workspace"], fixture.workspace().to_str().unwrap()); |
| 1181 | assert_eq!(seen["dataDir"], data_dir.to_str().unwrap()); |
| 1182 | // The workspace is the call's own, not the process's or the plugin's. |
| 1183 | let elsewhere = fixture.workspace().join("elsewhere"); |
| 1184 | let seen = call_json(&probe, json!({}), &elsewhere, &engine).await; |
| 1185 | assert_eq!(seen["workspace"], elsewhere.to_str().unwrap()); |
| 1186 | // Nothing else about the machine is in the context. |
| 1187 | assert_eq!( |
| 1188 | seen["keys"], |
| 1189 | json!(["args", "callId", "dataDir", "signal", "workspace"]) |
| 1190 | ); |
| 1191 | // The schema refuses an unknown field before the host is asked. |
| 1192 | let error = probe |
| 1193 | .execute( |
| 1194 | json!({"home": true}), |
| 1195 | &ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()), |
| 1196 | ) |
| 1197 | .await |
| 1198 | .unwrap_err(); |
| 1199 | assert!(matches!(error, ToolError::InvalidInput { .. }), "{error:?}"); |
| 1200 | |
| 1201 | // The plugin can write in its own directory, inside the host's sandbox. |
| 1202 | let note = host_tool(&engine, fixture.workspace(), "ctx_note"); |
| 1203 | let written = call_json( |
| 1204 | ¬e, |
| 1205 | json!({"text": "remember"}), |
| 1206 | fixture.workspace(), |
| 1207 | &engine, |
| 1208 | ) |
| 1209 | .await; |
| 1210 | assert_eq!(written["file"], data_dir.join("note.txt").to_str().unwrap()); |
| 1211 | assert_eq!( |
| 1212 | std::fs::read_to_string(data_dir.join("note.txt")).unwrap(), |
| 1213 | "remember" |
| 1214 | ); |
| 1215 | |
| 1216 | // A command is told the workspace it was loaded for, and the same directory. |
| 1217 | let entries = manager.commands_for_plugins(engine.plugin_view().as_ref()); |
| 1218 | let command = entries |
| 1219 | .first() |
| 1220 | .expect("the plugin's command is live") |
| 1221 | .reference(); |
| 1222 | let super::command::CommandOutcome::Show { text } = |
| 1223 | super::command::run(&manager.shared, &command, "", None) |
| 1224 | .await |
| 1225 | .unwrap() |
| 1226 | else { |
| 1227 | panic!("a show answer"); |
| 1228 | }; |
| 1229 | let said: Value = serde_json::from_str(&text).unwrap(); |
| 1230 | assert_eq!(said["workspace"], fixture.workspace().to_str().unwrap()); |
| 1231 | assert_eq!(said["dataDir"], data_dir.to_str().unwrap()); |
| 1232 | assert_eq!(said["config"]["greeting"], "Hi"); |
| 1233 | |
| 1234 | // What `/plugin show` can list: keys, not values. |
| 1235 | assert_eq!( |
| 1236 | manager |
| 1237 | .plugin_config_summary("plugin-context") |
| 1238 | .unwrap() |
| 1239 | .unwrap(), |
| 1240 | ["greeting"] |
| 1241 | ); |
| 1242 | |
| 1243 | // Unchanged settings never churn the owner... |
| 1244 | manager.set_plugin_settings( |
| 1245 | &plugin_settings("plugin-context", "greeting = \"Hi\""), |
| 1246 | None, |
| 1247 | ); |
| 1248 | engine.sync().await.unwrap(); |
| 1249 | assert_eq!(generation(&manager), Some(first_generation)); |
| 1250 | // ...a changed value is a new generation with the new value... |
| 1251 | manager.set_plugin_settings( |
| 1252 | &plugin_settings("plugin-context", "greeting = \"Yo\""), |
| 1253 | None, |
| 1254 | ); |
| 1255 | engine.sync().await.unwrap(); |
| 1256 | assert_eq!(manager.owner_state(&id), Some(OwnerState::Active)); |
| 1257 | assert!(generation(&manager).unwrap() > first_generation); |
| 1258 | let probe = host_tool(&engine, fixture.workspace(), "ctx_probe"); |
| 1259 | assert_eq!( |
| 1260 | call_json(&probe, json!({}), fixture.workspace(), &engine).await["config"]["greeting"], |
| 1261 | "Yo" |
| 1262 | ); |
| 1263 | // The same directory serves every generation: the note survived. |
| 1264 | assert_eq!( |
| 1265 | std::fs::read_to_string(data_dir.join("note.txt")).unwrap(), |
| 1266 | "remember" |
| 1267 | ); |
| 1268 | |
| 1269 | // ...a config the plugin's own schema refuses fails it with the field named |
| 1270 | // and leaves nothing live... |
| 1271 | manager.set_plugin_settings(&plugin_settings("plugin-context", "limit = 99"), None); |
| 1272 | engine.sync().await.unwrap(); |
| 1273 | assert!( |
| 1274 | matches!(manager.owner_state(&id), Some(OwnerState::Failed(ref reason)) if reason.contains("limit")), |
| 1275 | "{:?}", |
| 1276 | manager.owner_state(&id) |
| 1277 | ); |
| 1278 | assert!(manager.live_tool_names().is_empty()); |
| 1279 | assert!( |
| 1280 | manager |
| 1281 | .commands_for_plugins(engine.plugin_view().as_ref()) |
| 1282 | .is_empty() |
| 1283 | ); |
| 1284 | // ...so does one over the size cap (refused before the host is asked), and |
| 1285 | // `/plugin show` says why... |
| 1286 | manager.set_plugin_settings( |
| 1287 | &plugin_settings( |
| 1288 | "plugin-context", |
| 1289 | &format!( |
| 1290 | "blob = \"{}\"", |
| 1291 | "x".repeat(super::plugin_config::MAX_PLUGIN_CONFIG_BYTES) |
| 1292 | ), |
| 1293 | ), |
| 1294 | None, |
| 1295 | ); |
| 1296 | engine.sync().await.unwrap(); |
| 1297 | assert!( |
| 1298 | matches!(manager.owner_state(&id), Some(OwnerState::Failed(ref reason)) if reason.contains("byte limit")), |
| 1299 | "{:?}", |
| 1300 | manager.owner_state(&id) |
| 1301 | ); |
| 1302 | assert!( |
| 1303 | manager |
| 1304 | .plugin_config_summary("plugin-context") |
| 1305 | .unwrap() |
| 1306 | .is_err() |
| 1307 | ); |
| 1308 | // ...and fixing the settings brings it back, without a restart. |
| 1309 | manager.set_plugin_settings( |
| 1310 | &plugin_settings("plugin-context", "greeting = \"Hello again\""), |
| 1311 | None, |
| 1312 | ); |
| 1313 | engine.sync().await.unwrap(); |
| 1314 | assert_eq!(manager.owner_state(&id), Some(OwnerState::Active)); |
| 1315 | manager.shutdown().await; |
| 1316 | } |
| 1317 | |
| 1318 | #[tokio::test] |
| 1319 | async fn execute_tools_refuses_extension_tools_before_any_host_call() { |
| 1320 | let Some(node) = node_for_tests("execute_tools_refuses_extension_tools_before_any_host_call") |
| 1321 | else { |
| 1322 | return; |
| 1323 | }; |
| 1324 | let _policy = TestPolicyGuard::extension_host(true); |
| 1325 | let fixture = FixturePlugins::new(&["slow-tool"]).await; |
| 1326 | let manager = fixture.manager(node); |
| 1327 | let engine = manager.attach(fixture.registry()); |
| 1328 | engine.sync().await.unwrap(); |
| 1329 | let mut registry = crate::tools::registry::ToolRegistryBuilder::new() |
| 1330 | .build(ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view())); |
| 1331 | engine.install_tools(&mut registry); |
| 1332 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()); |
| 1333 | let started = Instant::now(); |
| 1334 | let result = crate::tools::codemode::execute_tools_tool( |
| 1335 | &json!({"code": "return await tools.call('slow_wait', { ms: 5000 })"}), |
| 1336 | ®istry, |
| 1337 | &context, |
| 1338 | ) |
| 1339 | .await |
| 1340 | .unwrap(); |
| 1341 | // Refused at the gate: had the call reached the host it would take 5 s. |
| 1342 | assert!(started.elapsed() < Duration::from_secs(4)); |
| 1343 | assert!(!result.success, "{}", result.content); |
| 1344 | assert!( |
| 1345 | result.content.contains("can mutate") || result.content.contains("needs approval"), |
| 1346 | "{}", |
| 1347 | result.content |
| 1348 | ); |
| 1349 | manager.shutdown().await; |
| 1350 | } |
| 1351 | |
| 1352 | #[tokio::test] |
| 1353 | async fn disabling_mid_call_revokes_at_once_and_teardown_waits_for_async_disposers() { |
| 1354 | let Some(node) = node_for_tests("disabling_mid_call") else { |
| 1355 | return; |
| 1356 | }; |
| 1357 | let _policy = TestPolicyGuard::extension_host(true); |
| 1358 | let fixture = FixturePlugins::new(&["slow-tool"]).await; |
| 1359 | let manager = fixture.manager(node); |
| 1360 | let engine = manager.attach(fixture.registry()); |
| 1361 | engine.sync().await.unwrap(); |
| 1362 | let tool = host_tool(&engine, fixture.workspace(), "slow_wait"); |
| 1363 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()); |
| 1364 | let call = tokio::spawn(async move { tool.execute(json!({}), &context).await }); |
| 1365 | tokio::time::sleep(Duration::from_millis(150)).await; |
| 1366 | |
| 1367 | let disabled = fixture.disable("slow-tool"); |
| 1368 | // The revocation scan (a full re-hash of the staged tree) runs first; |
| 1369 | // the clock for the 500 ms bound starts when the registry drops the |
| 1370 | // handle, which is the moment revocation takes effect. A plain thread |
| 1371 | // watches for it, because this runtime is single-threaded (the policy |
| 1372 | // override is thread-local) and would only look when `sync` yields. |
| 1373 | let watcher = { |
| 1374 | let manager = Arc::clone(&manager); |
| 1375 | std::thread::spawn(move || { |
| 1376 | let started = Instant::now(); |
| 1377 | while !manager.live_tool_names().is_empty() { |
| 1378 | assert!( |
| 1379 | started.elapsed() < Duration::from_secs(10), |
| 1380 | "revocation never happened" |
| 1381 | ); |
| 1382 | std::thread::yield_now(); |
| 1383 | } |
| 1384 | Instant::now() |
| 1385 | }) |
| 1386 | }; |
| 1387 | let sync_started = Instant::now(); |
| 1388 | engine.set_plugins(disabled); |
| 1389 | let sync = { |
| 1390 | let manager = Arc::clone(&manager); |
| 1391 | tokio::spawn(async move { manager.reconcile().await }) |
| 1392 | }; |
| 1393 | let outcome = tokio::time::timeout(Duration::from_secs(2), call) |
| 1394 | .await |
| 1395 | .expect("call resolves") |
| 1396 | .unwrap(); |
| 1397 | let resolved_at = Instant::now(); |
| 1398 | let revoked_at = watcher.join().unwrap(); |
| 1399 | let call_resolved = resolved_at.saturating_duration_since(revoked_at); |
| 1400 | assert!( |
| 1401 | matches!(outcome, Err(ToolError::Cancelled { .. })), |
| 1402 | "{outcome:?}" |
| 1403 | ); |
| 1404 | eprintln!( |
| 1405 | "extension host: in-flight call resolved as cancelled {:.1} ms after revocation", |
| 1406 | call_resolved.as_secs_f64() * 1000.0 |
| 1407 | ); |
| 1408 | assert!( |
| 1409 | call_resolved < Duration::from_millis(500), |
| 1410 | "{call_resolved:?}" |
| 1411 | ); |
| 1412 | sync.await.unwrap().unwrap(); |
| 1413 | let teardown = sync_started.elapsed(); |
| 1414 | assert!( |
| 1415 | teardown >= Duration::from_millis(300), |
| 1416 | "ack must wait for the 300 ms async disposer (got {teardown:?})" |
| 1417 | ); |
| 1418 | let diagnostics = manager.diagnostics(); |
| 1419 | assert!( |
| 1420 | !diagnostics.iter().any(|d| d.contains("teardown")), |
| 1421 | "disposed with nothing leaked: {diagnostics:?}" |
| 1422 | ); |
| 1423 | manager.shutdown().await; |
| 1424 | } |
| 1425 | |
| 1426 | #[tokio::test] |
| 1427 | async fn killed_host_fails_calls_once_and_replays_with_fresh_owners() { |
| 1428 | let Some(node) = node_for_tests("killed_host") else { |
| 1429 | return; |
| 1430 | }; |
| 1431 | let _policy = TestPolicyGuard::extension_host(true); |
| 1432 | let fixture = FixturePlugins::new(&["slow-tool"]).await; |
| 1433 | let manager = fixture.manager(node); |
| 1434 | let engine = manager.attach(fixture.registry()); |
| 1435 | engine.sync().await.unwrap(); |
| 1436 | assert_eq!(manager.spawn_attempts(), 1); |
| 1437 | let pid = manager.host_pid().unwrap(); |
| 1438 | let tool = host_tool(&engine, fixture.workspace(), "slow_wait"); |
| 1439 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()); |
| 1440 | let call = tokio::spawn(async move { tool.execute(json!({}), &context).await }); |
| 1441 | tokio::time::sleep(Duration::from_millis(150)).await; |
| 1442 | #[cfg(unix)] |
| 1443 | let status = std::process::Command::new("kill") |
| 1444 | .args(["-9", &pid.to_string()]) |
| 1445 | .status() |
| 1446 | .unwrap(); |
| 1447 | #[cfg(windows)] |
| 1448 | let status = std::process::Command::new("taskkill") |
| 1449 | .args(["/F", "/PID", &pid.to_string()]) |
| 1450 | .status() |
| 1451 | .unwrap(); |
| 1452 | assert!(status.success()); |
| 1453 | let outcome = tokio::time::timeout(Duration::from_secs(2), call) |
| 1454 | .await |
| 1455 | .expect("call resolves") |
| 1456 | .unwrap(); |
| 1457 | match outcome { |
| 1458 | Err(ToolError::NotAvailable { message }) => { |
| 1459 | assert!(message.contains("extension host exited"), "{message}") |
| 1460 | } |
| 1461 | other => panic!("expected a typed not-available error, got {other:?}"), |
| 1462 | } |
| 1463 | wait_host(&manager, || { |
| 1464 | manager.spawn_attempts() == 2 && manager.live_tool_names().contains(&"slow_wait".into()) |
| 1465 | }) |
| 1466 | .await; |
| 1467 | assert!(matches!(manager.status(), HostStatus::Ready { .. })); |
| 1468 | assert_ne!(manager.host_pid(), Some(pid)); |
| 1469 | assert_eq!( |
| 1470 | manager |
| 1471 | .shared |
| 1472 | .plugin |
| 1473 | .supervision |
| 1474 | .lock() |
| 1475 | .unwrap() |
| 1476 | .crashes |
| 1477 | .len(), |
| 1478 | 1 |
| 1479 | ); |
| 1480 | manager.shutdown().await; |
| 1481 | } |
| 1482 | |
| 1483 | #[tokio::test] |
| 1484 | async fn approval_providing_plugin_fails_activation_and_leaves_nothing_registered() { |
| 1485 | let Some(node) = node_for_tests("approval_providing_plugin") else { |
| 1486 | return; |
| 1487 | }; |
| 1488 | let _policy = TestPolicyGuard::extension_host(true); |
| 1489 | let fixture = FixturePlugins::new(&["refuses-approval", "clash-native"]).await; |
| 1490 | let manager = fixture.manager(node); |
| 1491 | let engine = manager.attach(fixture.registry()); |
| 1492 | engine.sync().await.unwrap(); |
| 1493 | let registry = fixture.registry(); |
| 1494 | for (name, needle) in [ |
| 1495 | ("refuses-approval", "approval"), |
| 1496 | ("clash-native", "read_file"), |
| 1497 | ] { |
| 1498 | let id = registry.get(name).unwrap().id.as_str().to_string(); |
| 1499 | match manager.owner_state(&id) { |
| 1500 | Some(OwnerState::Failed(reason)) => { |
| 1501 | assert!(reason.contains(needle), "{name}: {reason}") |
| 1502 | } |
| 1503 | other => panic!("{name}: expected failed activation, got {other:?}"), |
| 1504 | } |
| 1505 | } |
| 1506 | assert!(manager.live_tool_names().is_empty()); |
| 1507 | // A failed activation of the same bytes is not retried every turn. |
| 1508 | let attempts = manager.spawn_attempts(); |
| 1509 | engine.sync().await.unwrap(); |
| 1510 | assert_eq!(manager.spawn_attempts(), attempts); |
| 1511 | manager.shutdown().await; |
| 1512 | } |
| 1513 | |
| 1514 | /// Stands in for a `~/.codewhale/tools` script tool. |
| 1515 | struct FakeScriptTool; |
| 1516 | |
| 1517 | #[async_trait::async_trait] |
| 1518 | impl ToolSpec for FakeScriptTool { |
| 1519 | fn name(&self) -> &str { |
| 1520 | "fixture_script_tool" |
| 1521 | } |
| 1522 | fn registration_origin(&self) -> std::borrow::Cow<'_, str> { |
| 1523 | "plugin script fixture_script_tool".into() |
| 1524 | } |
| 1525 | fn description(&self) -> &str { |
| 1526 | "script" |
| 1527 | } |
| 1528 | fn input_schema(&self) -> Value { |
| 1529 | json!({"type": "object"}) |
| 1530 | } |
| 1531 | fn capabilities(&self) -> Vec<crate::tools::spec::ToolCapability> { |
| 1532 | Vec::new() |
| 1533 | } |
| 1534 | async fn execute( |
| 1535 | &self, |
| 1536 | _input: Value, |
| 1537 | _context: &ToolContext, |
| 1538 | ) -> Result<crate::tools::spec::ToolResult, ToolError> { |
| 1539 | Ok(crate::tools::spec::ToolResult::success("from the script")) |
| 1540 | } |
| 1541 | } |
| 1542 | |
| 1543 | #[tokio::test] |
| 1544 | async fn an_extension_named_like_a_script_tool_is_skipped_at_turn_build() { |
| 1545 | let Some(node) = node_for_tests("script_name_clash") else { |
| 1546 | return; |
| 1547 | }; |
| 1548 | let _policy = TestPolicyGuard::extension_host(true); |
| 1549 | let fixture = FixturePlugins::new(&["clash-script"]).await; |
| 1550 | let manager = fixture.manager(node); |
| 1551 | let engine = manager.attach(fixture.registry()); |
| 1552 | engine.sync().await.unwrap(); |
| 1553 | assert_eq!(manager.live_tool_names(), vec!["fixture_script_tool"]); |
| 1554 | let mut registry = crate::tools::registry::ToolRegistryBuilder::new() |
| 1555 | .build(ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view())); |
| 1556 | registry.register(Arc::new(FakeScriptTool)); |
| 1557 | let installed = engine.install_tools(&mut registry); |
| 1558 | assert!(installed.is_empty()); |
| 1559 | assert_eq!( |
| 1560 | registry |
| 1561 | .get("fixture_script_tool") |
| 1562 | .unwrap() |
| 1563 | .registration_origin(), |
| 1564 | "plugin script fixture_script_tool", |
| 1565 | "the script tool is unaffected" |
| 1566 | ); |
| 1567 | assert!( |
| 1568 | manager |
| 1569 | .diagnostics() |
| 1570 | .iter() |
| 1571 | .any(|d| d.contains("fixture_script_tool") && d.contains("skipped")), |
| 1572 | "{:?}", |
| 1573 | manager.diagnostics() |
| 1574 | ); |
| 1575 | manager.shutdown().await; |
| 1576 | } |
| 1577 | |
| 1578 | #[tokio::test] |
| 1579 | async fn with_no_native_plugin_the_host_is_never_spawned() { |
| 1580 | let _policy = TestPolicyGuard::extension_host(true); |
| 1581 | let temp = tempfile::tempdir().unwrap(); |
| 1582 | let registry = Arc::new(PluginRegistry::empty(temp.path())); |
| 1583 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 1584 | node_override: None, |
| 1585 | root: Some(temp.path().join("home")), |
| 1586 | ..Default::default() |
| 1587 | })); |
| 1588 | let engine = manager.attach(registry); |
| 1589 | engine.sync().await.unwrap(); |
| 1590 | assert_eq!(manager.spawn_attempts(), 0); |
| 1591 | assert_eq!(manager.status(), HostStatus::Idle); |
| 1592 | assert!(!temp.path().join("home").exists(), "nothing materialized"); |
| 1593 | } |
| 1594 | |
| 1595 | async fn probe(tool: &Arc<dyn ToolSpec>, path: &Path, context: &ToolContext) -> Value { |
| 1596 | let result = tool |
| 1597 | .execute(json!({"path": path.to_string_lossy()}), context) |
| 1598 | .await |
| 1599 | .unwrap(); |
| 1600 | serde_json::from_str(&result.content).unwrap() |
| 1601 | } |
| 1602 | |
| 1603 | #[tokio::test] |
| 1604 | async fn sandboxed_host_cannot_read_codewhale_secrets_or_write_outside_its_data_dir() { |
| 1605 | let Some(node) = node_for_tests("sandboxed_host") else { |
| 1606 | return; |
| 1607 | }; |
| 1608 | sandboxed_host_boundary(NODE, node).await; |
| 1609 | } |
| 1610 | |
| 1611 | /// Same actual Rust manager, review/activation, tool and filesystem scenario, |
| 1612 | /// using the exact compiled image selected by the existing Bun resolver. |
| 1613 | #[cfg(any(target_os = "linux", target_os = "macos", windows))] |
| 1614 | #[tokio::test] |
| 1615 | async fn compiled_native_host_cannot_read_secrets_or_write_outside_its_data_dir() { |
| 1616 | let Some(binary) = compiled_image_for_tests() else { |
| 1617 | return; |
| 1618 | }; |
| 1619 | sandboxed_host_boundary(crate::config::ExtensionHostRuntime::Bun, binary).await; |
| 1620 | // Emitted only after the complete actual Rust admission/tool scenario. |
| 1621 | // CI extracts this same full-run success output; no fake-Core promotion. |
| 1622 | eprintln!( |
| 1623 | "compiled-native-containment=passed platform={} arch={}", |
| 1624 | std::env::consts::OS, |
| 1625 | std::env::consts::ARCH |
| 1626 | ); |
| 1627 | } |
| 1628 | |
| 1629 | /// Resolve the exact requested image through the same production runtime |
| 1630 | /// admission for every compiled Native scenario. Required inputs cannot skip. |
| 1631 | #[cfg(any(target_os = "linux", target_os = "macos", windows))] |
| 1632 | fn compiled_image_for_tests() -> Option<PathBuf> { |
| 1633 | let Some(binary) = std::env::var_os("CODEWHALE_COMPILED_HOST_TEST_BINARY") else { |
| 1634 | assert!( |
| 1635 | std::env::var_os("CODEWHALE_EXT_HOST_TESTS").is_none(), |
| 1636 | "required compiled Native image input is missing" |
| 1637 | ); |
| 1638 | eprintln!("compiled Native receipt unavailable: name a matching canonical compiled image"); |
| 1639 | return None; |
| 1640 | }; |
| 1641 | let binary = PathBuf::from(binary); |
| 1642 | let resolution = crate::dependencies::resolve_extension_host_runtime( |
| 1643 | crate::config::ExtensionHostRuntime::Bun, |
| 1644 | None, |
| 1645 | Some(&binary), |
| 1646 | ); |
| 1647 | let runtime = resolution |
| 1648 | .selected |
| 1649 | .as_ref() |
| 1650 | .unwrap_or_else(|| panic!("compiled Native runtime: {}", resolution.failure())); |
| 1651 | assert!( |
| 1652 | runtime.compiled, |
| 1653 | "receipt requires an actual canonical compiled image, not system Bun" |
| 1654 | ); |
| 1655 | Some(binary) |
| 1656 | } |
| 1657 | |
| 1658 | async fn sandboxed_host_boundary( |
| 1659 | choice: crate::config::ExtensionHostRuntime, |
| 1660 | runtime_path: PathBuf, |
| 1661 | ) { |
| 1662 | let _policy = TestPolicyGuard::extension_host(true); |
| 1663 | let fixture = FixturePlugins::new(&["secret-probe"]).await; |
| 1664 | // Created before launch: the deny-list records the canonical spelling of |
| 1665 | // paths that exist (macOS `/var` → `/private/var`). |
| 1666 | let secrets = fixture.root.join("secrets"); |
| 1667 | std::fs::create_dir_all(&secrets).unwrap(); |
| 1668 | let token = secrets.join("token"); |
| 1669 | std::fs::write(&token, "s3cret-value").unwrap(); |
| 1670 | // Any other entry of the Codewhale home is denied too (config backups, |
| 1671 | // OAuth tokens, state), not only the named stores. |
| 1672 | let backup = fixture.root.join("config.toml.bak-20260925"); |
| 1673 | std::fs::write(&backup, "api_key = \"s3cret-backup\"").unwrap(); |
| 1674 | let tokens = fixture.root.join("tokens"); |
| 1675 | std::fs::create_dir_all(&tokens).unwrap(); |
| 1676 | std::fs::write(tokens.join("codex.json"), "s3cret-oauth").unwrap(); |
| 1677 | // Outside Core home, POSIX ordinary reads remain allowed; Windows LPAC |
| 1678 | // refuses ungranted workspace reads. Neither path grants Core secrets. |
| 1679 | let readable = fixture.workspace().join("readable.txt"); |
| 1680 | std::fs::write(&readable, "plain").unwrap(); |
| 1681 | |
| 1682 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 1683 | runtime: choice, |
| 1684 | node_override: (choice == NODE).then_some(runtime_path.clone()), |
| 1685 | bun_override: (choice == crate::config::ExtensionHostRuntime::Bun) |
| 1686 | .then_some(runtime_path.clone()), |
| 1687 | root: Some(fixture.root.clone()), |
| 1688 | ..Default::default() |
| 1689 | })); |
| 1690 | let engine = manager.attach(fixture.registry()); |
| 1691 | engine.sync().await.unwrap(); |
| 1692 | let HostStatus::Ready { sandbox, .. } = manager.status() else { |
| 1693 | panic!("host not ready: {:?}", manager.status()); |
| 1694 | }; |
| 1695 | let sandbox = match sandbox { |
| 1696 | super::supervisor::HostSandbox::Wrapped(name) => name, |
| 1697 | super::supervisor::HostSandbox::Unsandboxed(reason) => { |
| 1698 | panic!("Native containment requires a verified sandbox: {reason}"); |
| 1699 | } |
| 1700 | }; |
| 1701 | assert!(super::render_status(&manager).contains(&format!("{sandbox} sandbox"))); |
| 1702 | |
| 1703 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()); |
| 1704 | let read = host_tool(&engine, fixture.workspace(), "probe_read"); |
| 1705 | let write = host_tool(&engine, fixture.workspace(), "probe_write"); |
| 1706 | |
| 1707 | let plain = probe(&read, &readable, &context).await; |
| 1708 | #[cfg(not(windows))] |
| 1709 | assert_eq!( |
| 1710 | plain, |
| 1711 | json!({"ok": true, "text": "plain"}), |
| 1712 | "ordinary reads work" |
| 1713 | ); |
| 1714 | #[cfg(windows)] |
| 1715 | assert_eq!(plain["ok"], false, "LPAC refuses ungranted workspace reads"); |
| 1716 | for denied in [token, backup, tokens.join("codex.json")] { |
| 1717 | let secret = probe(&read, &denied, &context).await; |
| 1718 | assert_eq!(secret["ok"], false, "{} was readable", denied.display()); |
| 1719 | assert!( |
| 1720 | !secret.to_string().contains("s3cret"), |
| 1721 | "{} leaked its contents", |
| 1722 | denied.display() |
| 1723 | ); |
| 1724 | } |
| 1725 | // A store created after the host started is denied by name. |
| 1726 | let state = fixture.root.join("state"); |
| 1727 | std::fs::create_dir_all(&state).unwrap(); |
| 1728 | std::fs::write(state.join("late.json"), "s3cret-late").unwrap(); |
| 1729 | let late = probe(&read, &state.join("late.json"), &context).await; |
| 1730 | assert_eq!( |
| 1731 | late["ok"], false, |
| 1732 | "a store created after start was readable" |
| 1733 | ); |
| 1734 | // The migrated history can first appear after host launch too. Its name |
| 1735 | // must be denied before enumeration can observe the file. |
| 1736 | let history = fixture.root.join("composer_history.jsonl"); |
| 1737 | std::fs::write(&history, "\"private synthetic prompt\"\n").unwrap(); |
| 1738 | let late_history = probe(&read, &history, &context).await; |
| 1739 | assert_eq!(late_history["ok"], false, "new history was readable"); |
| 1740 | // The Codex credential file Codewhale itself reads, when this machine has |
| 1741 | // one. Only `ok` is reported, never the content. |
| 1742 | let codex_auth = crate::oauth::auth_file_path(); |
| 1743 | if codex_auth.is_file() { |
| 1744 | let codex = probe(&read, &codex_auth, &context).await; |
| 1745 | assert_eq!(codex["ok"], false, "code: {}", codex["code"]); |
| 1746 | } |
| 1747 | |
| 1748 | let data = fixture.root.join("extension-host/data/probe.txt"); |
| 1749 | assert_eq!(probe(&write, &data, &context).await["ok"], true); |
| 1750 | // Outside the data dir and the temp dirs (the fixture itself lives under |
| 1751 | // TMPDIR, which the profile leaves writable): the crate's source dir, |
| 1752 | // unless the checkout itself sits in a temp dir. |
| 1753 | let crate_dir = std::fs::canonicalize(env!("CARGO_MANIFEST_DIR")).unwrap(); |
| 1754 | let in_temp = [std::env::temp_dir(), PathBuf::from("/tmp")] |
| 1755 | .iter() |
| 1756 | .filter_map(|dir| std::fs::canonicalize(dir).ok()) |
| 1757 | .any(|dir| crate_dir.starts_with(dir)); |
| 1758 | if !in_temp { |
| 1759 | let escape = crate_dir.join(format!(".ext-host-probe-{}", uuid::Uuid::new_v4().simple())); |
| 1760 | let escaped = probe(&write, &escape, &context).await; |
| 1761 | let leaked = escape.exists(); |
| 1762 | let _ = std::fs::remove_file(&escape); |
| 1763 | assert_eq!(escaped["ok"], false, "{escaped}"); |
| 1764 | assert!(!leaked); |
| 1765 | } |
| 1766 | // Bun's TCP error adapter collapses Darwin's EPERM into ECONNREFUSED. |
| 1767 | // Its UDP adapter retains the OS errno, so use a real bound datagram |
| 1768 | // receiver there; every other runtime/platform keeps the TCP control. |
| 1769 | let (denied, accepted_count, protocol) = |
| 1770 | if cfg!(target_os = "macos") && choice == crate::config::ExtensionHostRuntime::Bun { |
| 1771 | let listener = tokio::net::UdpSocket::bind(("127.0.0.1", 0)).await.unwrap(); |
| 1772 | let address = listener.local_addr().unwrap(); |
| 1773 | let positive = tokio::net::UdpSocket::bind(("127.0.0.1", 0)).await.unwrap(); |
| 1774 | let marker = b"controller-network-probe"; |
| 1775 | assert_eq!( |
| 1776 | positive.send_to(marker, address).await.unwrap(), |
| 1777 | marker.len() |
| 1778 | ); |
| 1779 | let mut packet = [0_u8; 64]; |
| 1780 | let (bytes, sender) = |
| 1781 | tokio::time::timeout(Duration::from_secs(2), listener.recv_from(&mut packet)) |
| 1782 | .await |
| 1783 | .unwrap() |
| 1784 | .unwrap(); |
| 1785 | assert_eq!(&packet[..bytes], marker); |
| 1786 | assert_eq!(sender, positive.local_addr().unwrap()); |
| 1787 | let mut accepted_count = 1; |
| 1788 | let send = host_tool(&engine, fixture.workspace(), "probe_send"); |
| 1789 | let result = tokio::time::timeout( |
| 1790 | Duration::from_secs(5), |
| 1791 | send.execute(json!({"port": address.port()}), &context), |
| 1792 | ) |
| 1793 | .await |
| 1794 | .unwrap() |
| 1795 | .unwrap(); |
| 1796 | let denied: Value = serde_json::from_str(&result.content).unwrap(); |
| 1797 | if tokio::time::timeout(Duration::from_millis(100), listener.recv_from(&mut packet)) |
| 1798 | .await |
| 1799 | .is_ok() |
| 1800 | { |
| 1801 | accepted_count += 1; |
| 1802 | } |
| 1803 | (denied, accepted_count, "udp") |
| 1804 | } else { |
| 1805 | // The unsandboxed controller reaches this actual listener. The Native |
| 1806 | // host must fail the same connection at its own OS sandbox boundary. |
| 1807 | let listener = tokio::net::TcpListener::bind(("127.0.0.1", 0)) |
| 1808 | .await |
| 1809 | .unwrap(); |
| 1810 | let address = listener.local_addr().unwrap(); |
| 1811 | let positive = tokio::net::TcpStream::connect(address).await.unwrap(); |
| 1812 | let (accepted, _) = tokio::time::timeout(Duration::from_secs(2), listener.accept()) |
| 1813 | .await |
| 1814 | .unwrap() |
| 1815 | .unwrap(); |
| 1816 | let mut accepted_count = 1; |
| 1817 | drop(positive); |
| 1818 | drop(accepted); |
| 1819 | let connect = host_tool(&engine, fixture.workspace(), "probe_connect"); |
| 1820 | let result = tokio::time::timeout( |
| 1821 | Duration::from_secs(5), |
| 1822 | connect.execute(json!({"port": address.port()}), &context), |
| 1823 | ) |
| 1824 | .await |
| 1825 | .unwrap() |
| 1826 | .unwrap(); |
| 1827 | let denied: Value = serde_json::from_str(&result.content).unwrap(); |
| 1828 | if tokio::time::timeout(Duration::from_millis(100), listener.accept()) |
| 1829 | .await |
| 1830 | .is_ok() |
| 1831 | { |
| 1832 | accepted_count += 1; |
| 1833 | } |
| 1834 | (denied, accepted_count, "tcp") |
| 1835 | }; |
| 1836 | assert_eq!( |
| 1837 | denied["ok"], false, |
| 1838 | "Native host reached {protocol} controller: {denied}" |
| 1839 | ); |
| 1840 | #[cfg(target_os = "linux")] |
| 1841 | { |
| 1842 | let controller = std::fs::read_link("/proc/self/ns/net").unwrap(); |
| 1843 | let controller = controller.to_str().unwrap(); |
| 1844 | let host = denied["network_namespace"].as_str().unwrap(); |
| 1845 | let namespace_id = |name: &str| { |
| 1846 | name.strip_prefix("net:[") |
| 1847 | .and_then(|name| name.strip_suffix(']')) |
| 1848 | .and_then(|id| id.parse::<u64>().ok()) |
| 1849 | .filter(|id| *id != 0) |
| 1850 | .expect("an actual bounded Linux kernel network namespace") |
| 1851 | }; |
| 1852 | assert_ne!( |
| 1853 | namespace_id(host), |
| 1854 | namespace_id(controller), |
| 1855 | "Native host must remain in its isolated kernel network namespace" |
| 1856 | ); |
| 1857 | assert_eq!( |
| 1858 | denied["code"], "ECONNREFUSED", |
| 1859 | "the isolated loopback cannot reach the controller: {denied}" |
| 1860 | ); |
| 1861 | } |
| 1862 | #[cfg(not(target_os = "linux"))] |
| 1863 | assert!( |
| 1864 | matches!(denied["code"].as_str(), Some("EPERM" | "EACCES")), |
| 1865 | "connection must be refused by the OS sandbox: {denied}" |
| 1866 | ); |
| 1867 | assert_eq!( |
| 1868 | accepted_count, 1, |
| 1869 | "only the controller reached the listener" |
| 1870 | ); |
| 1871 | if choice == crate::config::ExtensionHostRuntime::Bun { |
| 1872 | eprintln!( |
| 1873 | "compiled-native-network=passed controller_accepts=1 host_errno={} platform={} arch={} protocol={protocol}", |
| 1874 | denied["code"].as_str().unwrap(), |
| 1875 | std::env::consts::OS, |
| 1876 | std::env::consts::ARCH, |
| 1877 | ); |
| 1878 | } |
| 1879 | manager.shutdown().await; |
| 1880 | } |
| 1881 | |
| 1882 | // --------------------------------------------------------------------------- |
| 1883 | // Receipt-bound approval keys (design §4.3) |
| 1884 | // --------------------------------------------------------------------------- |
| 1885 | |
| 1886 | fn keys_for( |
| 1887 | manager: &ExtensionHostManager, |
| 1888 | registration: super::registry::ToolRegistration, |
| 1889 | input: &Value, |
| 1890 | ) -> (String, String) { |
| 1891 | let name = registration.name.clone(); |
| 1892 | let mut registry = crate::tools::ToolRegistry::new(ToolContext::new(Path::new("/w"))); |
| 1893 | registry.register(Arc::new(super::tool::HostToolSpec::new( |
| 1894 | registration, |
| 1895 | Arc::clone(&manager.shared), |
| 1896 | ))); |
| 1897 | let (exact, grouping) = |
| 1898 | crate::tools::approval_cache::approval_keys_for_call(Some(®istry), &name, input); |
| 1899 | (exact.0, grouping.0) |
| 1900 | } |
| 1901 | |
| 1902 | /// A session grant for an extension tool covers one reviewed plugin build: |
| 1903 | /// an update of the plugin, or another plugin that later registers the same |
| 1904 | /// tool name, gets a different key and is asked again. |
| 1905 | #[test] |
| 1906 | fn extension_approval_keys_are_bound_to_the_plugin_receipt() { |
| 1907 | let manager = ExtensionHostManager::new(ExtensionHostOptions::default()); |
| 1908 | let input = json!({"path": "x"}); |
| 1909 | let mut owners = OwnerRegistry::new(); |
| 1910 | let live = |owners: &mut OwnerRegistry, owner: &OwnerRef| { |
| 1911 | register(owners, owner, "shared_tool").unwrap(); |
| 1912 | owners.mark_active(owner); |
| 1913 | owners.live_tools().pop().unwrap() |
| 1914 | }; |
| 1915 | |
| 1916 | let first = owners |
| 1917 | .begin_owner( |
| 1918 | HostTier::Plugin, |
| 1919 | "a", |
| 1920 | "a", |
| 1921 | Some(fake_authority("a")), |
| 1922 | "hash-a1", |
| 1923 | ) |
| 1924 | .unwrap(); |
| 1925 | let first = keys_for(&manager, live(&mut owners, &first), &input); |
| 1926 | assert!( |
| 1927 | first.0.starts_with("ext:a@hash-a1:") && first.0.contains(":shared_tool:"), |
| 1928 | "{first:?}" |
| 1929 | ); |
| 1930 | assert_eq!(first.0, first.1, "a grant covers the exact call only"); |
| 1931 | let generic = crate::tools::approval_cache::build_approval_grouping_key("shared_tool", &input); |
| 1932 | assert_ne!(first.1, generic.0, "never the name-derived family key"); |
| 1933 | |
| 1934 | // Same plugin, same input, updated bytes: a different grant. |
| 1935 | let updated = owners |
| 1936 | .begin_owner( |
| 1937 | HostTier::Plugin, |
| 1938 | "a", |
| 1939 | "a", |
| 1940 | Some(fake_authority("a")), |
| 1941 | "hash-a2", |
| 1942 | ) |
| 1943 | .unwrap(); |
| 1944 | let updated = keys_for(&manager, live(&mut owners, &updated), &input); |
| 1945 | assert_ne!(first.1, updated.1); |
| 1946 | |
| 1947 | // Another plugin takes the name once the first is gone. |
| 1948 | owners.revoke_owner("a"); |
| 1949 | let other = owners |
| 1950 | .begin_owner( |
| 1951 | HostTier::Plugin, |
| 1952 | "b", |
| 1953 | "b", |
| 1954 | Some(fake_authority("b")), |
| 1955 | "hash-a1", |
| 1956 | ) |
| 1957 | .unwrap(); |
| 1958 | let other = keys_for(&manager, live(&mut owners, &other), &input); |
| 1959 | assert_ne!(first.1, other.1); |
| 1960 | assert_ne!(updated.1, other.1); |
| 1961 | |
| 1962 | // Tools without a scope keep their existing keys. |
| 1963 | let shell = json!({"command": "cargo build --release"}); |
| 1964 | let (exact, grouping) = |
| 1965 | crate::tools::approval_cache::approval_keys_for_call(None, "exec_shell", &shell); |
| 1966 | assert_eq!( |
| 1967 | exact, |
| 1968 | crate::tools::approval_cache::build_approval_key("exec_shell", &shell) |
| 1969 | ); |
| 1970 | assert_eq!( |
| 1971 | grouping, |
| 1972 | crate::tools::approval_cache::build_approval_grouping_key("exec_shell", &shell) |
| 1973 | ); |
| 1974 | } |
| 1975 | |
| 1976 | // --------------------------------------------------------------------------- |
| 1977 | // The native-entry rule at validate / review time |
| 1978 | // --------------------------------------------------------------------------- |
| 1979 | |
| 1980 | fn native_bundle(user: &Path, name: &str, native_path: &str, files: &[&str]) { |
| 1981 | let root = user.join(name); |
| 1982 | for file in files { |
| 1983 | let path = root.join(file); |
| 1984 | std::fs::create_dir_all(path.parent().unwrap()).unwrap(); |
| 1985 | std::fs::write( |
| 1986 | &path, |
| 1987 | "export const name = 'x'\nexport function apply() {}\n", |
| 1988 | ) |
| 1989 | .unwrap(); |
| 1990 | } |
| 1991 | std::fs::write( |
| 1992 | root.join("plugin.json"), |
| 1993 | serde_json::to_vec_pretty(&json!({ |
| 1994 | "$schema": "https://agent-plugins.org/schemas/plugin.json", |
| 1995 | "name": name, |
| 1996 | "version": "0.1.0", |
| 1997 | "description": "native entry rule fixture", |
| 1998 | "license": "MIT", |
| 1999 | "extensions": {"net.codewhale": {"native": {"path": native_path}}} |
| 2000 | })) |
| 2001 | .unwrap(), |
| 2002 | ) |
| 2003 | .unwrap(); |
| 2004 | } |
| 2005 | |
| 2006 | /// `/plugin validate` and the review screen read plugin diagnostics, so an |
| 2007 | /// entry that activation would refuse must fail there too, with the flag on; |
| 2008 | /// with it off `native` is inventory-only and any path stays valid. |
| 2009 | #[test] |
| 2010 | fn native_entry_rule_fails_validation_when_the_host_is_enabled() { |
| 2011 | let temp = tempfile::tempdir().unwrap(); |
| 2012 | let user = temp.path().join("user"); |
| 2013 | native_bundle(&user, "dir-entry", "lib", &["lib/index.mjs"]); |
| 2014 | native_bundle(&user, "ts-entry", "index.ts", &["index.ts"]); |
| 2015 | native_bundle(&user, "good-entry", "index.mjs", &["index.mjs"]); |
| 2016 | native_bundle(&user, "typed-entry", "index.mts", &["index.mts"]); |
| 2017 | let config = DiscoveryConfig { |
| 2018 | workspace: temp.path().join("project"), |
| 2019 | user_plugins_dir: user, |
| 2020 | workspace_plugins_dir: temp.path().join("project/.codewhale/plugins"), |
| 2021 | builtin_plugin_dirs: Vec::new(), |
| 2022 | state_path: temp.path().join("state/plugin-state.json"), |
| 2023 | }; |
| 2024 | let native_errors = |registry: &PluginRegistry, name: &str| -> Vec<String> { |
| 2025 | registry |
| 2026 | .get(name) |
| 2027 | .unwrap_or_else(|| panic!("{name} not discovered: {:?}", registry.diagnostics())) |
| 2028 | .diagnostics |
| 2029 | .iter() |
| 2030 | .filter(|diagnostic| diagnostic.code == "native-entry-invalid") |
| 2031 | .map(|diagnostic| { |
| 2032 | assert_eq!( |
| 2033 | diagnostic.level, |
| 2034 | crate::plugins::types::PluginDiagnosticLevel::Error |
| 2035 | ); |
| 2036 | diagnostic.message.clone() |
| 2037 | }) |
| 2038 | .collect() |
| 2039 | }; |
| 2040 | |
| 2041 | { |
| 2042 | let _policy = TestPolicyGuard::extension_host(true); |
| 2043 | let registry = discover_with_config(&config); |
| 2044 | for name in ["dir-entry", "ts-entry"] { |
| 2045 | let errors = native_errors(®istry, name); |
| 2046 | assert_eq!(errors.len(), 1, "{name}: {errors:?}"); |
| 2047 | assert!( |
| 2048 | errors[0].contains(".mjs, .js or .mts"), |
| 2049 | "{name}: {errors:?}" |
| 2050 | ); |
| 2051 | } |
| 2052 | assert!(native_errors(®istry, "good-entry").is_empty()); |
| 2053 | assert!(native_errors(®istry, "typed-entry").is_empty()); |
| 2054 | assert!(!registry.validation_is_clean()); |
| 2055 | } |
| 2056 | let _policy = TestPolicyGuard::extension_host(false); |
| 2057 | let registry = discover_with_config(&config); |
| 2058 | for name in ["dir-entry", "ts-entry", "good-entry", "typed-entry"] { |
| 2059 | assert!(native_errors(®istry, name).is_empty(), "{name}"); |
| 2060 | } |
| 2061 | } |
| 2062 | |
| 2063 | #[test] |
| 2064 | fn owner_reports_keep_bounded_attributed_logs_and_ignore_stale_hosts() { |
| 2065 | use super::supervisor::HostEvents; |
| 2066 | |
| 2067 | let manager = ExtensionHostManager::new(ExtensionHostOptions::default()); |
| 2068 | manager |
| 2069 | .shared |
| 2070 | .plugin |
| 2071 | .host_generation |
| 2072 | .store(2, std::sync::atomic::Ordering::SeqCst); |
| 2073 | { |
| 2074 | let mut registry = manager.shared.registry.lock().unwrap(); |
| 2075 | for id in ["alpha", "beta"] { |
| 2076 | let owner = registry |
| 2077 | .begin_owner(HostTier::Plugin, id, id, Some(fake_authority(id)), "hash") |
| 2078 | .unwrap(); |
| 2079 | registry.mark_active(&owner); |
| 2080 | register(&mut registry, &owner, &format!("{id}_probe")).unwrap(); |
| 2081 | } |
| 2082 | } |
| 2083 | let current = super::Events { |
| 2084 | shared: Arc::downgrade(&manager.shared), |
| 2085 | tier: HostTier::Plugin, |
| 2086 | generation: 2, |
| 2087 | }; |
| 2088 | let stale = super::Events { |
| 2089 | shared: Arc::downgrade(&manager.shared), |
| 2090 | tier: HostTier::Plugin, |
| 2091 | generation: 1, |
| 2092 | }; |
| 2093 | for index in 0..25 { |
| 2094 | current.log(&protocol::LogParams { |
| 2095 | level: "warn".into(), |
| 2096 | msg: format!("alpha message {index}"), |
| 2097 | plugin_id: Some("alpha".into()), |
| 2098 | }); |
| 2099 | } |
| 2100 | let mut log = protocol::LogParams { |
| 2101 | level: "error".into(), |
| 2102 | msg: "beta only".into(), |
| 2103 | plugin_id: Some("beta".into()), |
| 2104 | }; |
| 2105 | current.log(&log); |
| 2106 | log.msg = "stale message".into(); |
| 2107 | stale.log(&log); |
| 2108 | log.plugin_id = Some("unknown".into()); |
| 2109 | current.log(&log); |
| 2110 | log.plugin_id = Some("alpha".into()); |
| 2111 | log.level = "debug".into(); |
| 2112 | current.log(&log); |
| 2113 | let alpha = manager.owner_report("alpha").unwrap(); |
| 2114 | assert!(matches!(alpha.state, Some(OwnerState::Active))); |
| 2115 | assert_eq!(alpha.tools, ["alpha_probe"]); |
| 2116 | assert_eq!(alpha.diagnostics.len(), 20); |
| 2117 | assert_eq!(alpha.diagnostics[0], "warn: alpha message 5"); |
| 2118 | assert_eq!( |
| 2119 | manager.owner_report("beta").unwrap().diagnostics, |
| 2120 | ["error: beta only"] |
| 2121 | ); |
| 2122 | assert!(manager.owner_report("unknown").is_none()); |
| 2123 | manager |
| 2124 | .shared |
| 2125 | .plugin_diagnostic(&"x".repeat(10_000), "oversized id".into()); |
| 2126 | assert!( |
| 2127 | manager |
| 2128 | .shared |
| 2129 | .diagnostics |
| 2130 | .lock() |
| 2131 | .unwrap() |
| 2132 | .back() |
| 2133 | .unwrap() |
| 2134 | .plugin_id |
| 2135 | .is_none() |
| 2136 | ); |
| 2137 | for _ in 0..80 { |
| 2138 | manager.shared.plugin_diagnostic("alpha", "🦀".repeat(3000)); |
| 2139 | } |
| 2140 | assert_eq!(manager.diagnostics().len(), 64); |
| 2141 | assert!( |
| 2142 | manager |
| 2143 | .diagnostics() |
| 2144 | .iter() |
| 2145 | .all(|line| line.len() <= super::MAX_DIAGNOSTIC_BYTES + '…'.len_utf8()) |
| 2146 | ); |
| 2147 | assert!( |
| 2148 | manager |
| 2149 | .owner_report("alpha") |
| 2150 | .unwrap() |
| 2151 | .diagnostics |
| 2152 | .iter() |
| 2153 | .all(|line| line.ends_with('…')) |
| 2154 | ); |
| 2155 | } |
| 2156 | |
| 2157 | #[tokio::test] |
| 2158 | async fn typed_author_example_is_reviewed_before_its_tool_can_execute() { |
| 2159 | use crate::plugins::install::{DEFAULT_MAX_SIZE_BYTES, PluginInstallSource, install}; |
| 2160 | |
| 2161 | let Some(node) = node_for_tests("typed author example") else { |
| 2162 | return; |
| 2163 | }; |
| 2164 | let _policy = TestPolicyGuard::extension_host(true); |
| 2165 | let fixture = FixturePlugins::new(&[]).await; |
| 2166 | let example = |
| 2167 | Path::new(env!("CARGO_MANIFEST_DIR")).join("../../docs/examples/plugins/hello-extension"); |
| 2168 | install( |
| 2169 | PluginInstallSource::LocalPath(example), |
| 2170 | &fixture.config.user_plugins_dir, |
| 2171 | DEFAULT_MAX_SIZE_BYTES, |
| 2172 | &crate::network_policy::NetworkPolicy::default(), |
| 2173 | false, |
| 2174 | &|_| None, |
| 2175 | ) |
| 2176 | .await |
| 2177 | .unwrap(); |
| 2178 | let mut plugins = discover_with_config(&fixture.config); |
| 2179 | assert!(!plugins.is_active("hello-extension")); |
| 2180 | plugins.trust("hello-extension").unwrap(); |
| 2181 | assert!( |
| 2182 | !plugins.is_active("hello-extension"), |
| 2183 | "trust alone does not enable code" |
| 2184 | ); |
| 2185 | plugins.enable("hello-extension").unwrap(); |
| 2186 | // This tests plugin review and typed loading, not hang detection. The |
| 2187 | // 600 ms watchdog used by supervision fault tests can kill a healthy |
| 2188 | // typed-plugin load on a busy runner before registration completes. |
| 2189 | let manager = fixture.manager(node); |
| 2190 | let engine = manager.attach(Arc::new(plugins)); |
| 2191 | engine.sync().await.unwrap(); |
| 2192 | let tool = host_tool(&engine, fixture.workspace(), "hello_greet"); |
| 2193 | assert_eq!(tool.approval_requirement(), ApprovalRequirement::Required); |
| 2194 | let result = tool |
| 2195 | .execute( |
| 2196 | json!({"name": "Codewhale"}), |
| 2197 | &ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()), |
| 2198 | ) |
| 2199 | .await |
| 2200 | .unwrap(); |
| 2201 | let payload: Value = serde_json::from_str(&result.content).unwrap(); |
| 2202 | assert_eq!(payload["greeting"], "Hello, Codewhale!"); |
| 2203 | assert!(payload["callId"].as_str().is_some_and(|id| !id.is_empty())); |
| 2204 | |
| 2205 | // The example registers a plain object whose schema says |
| 2206 | // `additionalProperties: false` and `name: string`. The core enforces it; |
| 2207 | // the example's own `execute` would have accepted either input. Neither |
| 2208 | // call reaches the host. |
| 2209 | let sent = manager.host_requests_started(); |
| 2210 | for input in [ |
| 2211 | json!({"name": "Codewhale", "surprise": true}), |
| 2212 | json!({"name": 7}), |
| 2213 | ] { |
| 2214 | let error = tool |
| 2215 | .execute( |
| 2216 | input.clone(), |
| 2217 | &ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()), |
| 2218 | ) |
| 2219 | .await |
| 2220 | .unwrap_err(); |
| 2221 | assert!( |
| 2222 | matches!(error, ToolError::InvalidInput { ref message } if message.contains("hello_greet")), |
| 2223 | "{input}: {error:?}" |
| 2224 | ); |
| 2225 | } |
| 2226 | assert_eq!( |
| 2227 | manager.host_requests_started(), |
| 2228 | sent, |
| 2229 | "a call the schema refuses sends the host nothing" |
| 2230 | ); |
| 2231 | manager.shutdown().await; |
| 2232 | } |
| 2233 | |
| 2234 | // --------------------------------------------------------------------------- |
| 2235 | // Engines sharing one process-wide host |
| 2236 | // --------------------------------------------------------------------------- |
| 2237 | |
| 2238 | #[test] |
| 2239 | fn attachment_changes_discard_the_complete_scan_before_owner_side_effects() { |
| 2240 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions::default())); |
| 2241 | let old = Arc::new(PluginRegistry::empty(Path::new("/old"))); |
| 2242 | let current = Arc::new(PluginRegistry::empty(Path::new("/current"))); |
| 2243 | let attachment = manager.attach(Arc::clone(&old)); |
| 2244 | let scan = |plugins: &Arc<PluginRegistry>| super::DesiredScan { |
| 2245 | attachments: vec![( |
| 2246 | attachment.id, |
| 2247 | Arc::clone(plugins), |
| 2248 | [("plugin".into(), "hash-plugin".into())].into(), |
| 2249 | )], |
| 2250 | owners: [( |
| 2251 | "plugin".into(), |
| 2252 | super::DesiredOwner { |
| 2253 | plugin_name: "plugin".into(), |
| 2254 | authority: fake_authority("plugin"), |
| 2255 | entries: Vec::new(), |
| 2256 | }, |
| 2257 | )] |
| 2258 | .into(), |
| 2259 | errors: Vec::new(), |
| 2260 | }; |
| 2261 | |
| 2262 | // An old scan finishes after the engine has already changed workspace. |
| 2263 | attachment.set_plugins(Arc::clone(¤t)); |
| 2264 | let current_view = attachment.plugin_view(); |
| 2265 | let mut attachments = manager.shared.attachments.lock().unwrap(); |
| 2266 | assert!(scan(&old).publish(&mut attachments).is_none()); |
| 2267 | assert!(attachments[&attachment.id].desired.is_empty()); |
| 2268 | // A current scan publishes both the engine view and owner union. |
| 2269 | let (owners, _) = scan(¤t_view).publish(&mut attachments).unwrap(); |
| 2270 | assert!(owners.contains_key("plugin")); |
| 2271 | assert_eq!(attachments[&attachment.id].desired["plugin"], "hash-plugin"); |
| 2272 | drop(attachments); |
| 2273 | |
| 2274 | // A newly attached engine also invalidates the complete scan, even when |
| 2275 | // it uses the same snapshot: otherwise its owners could be revoked. |
| 2276 | let other = manager.attach(Arc::clone(¤t)); |
| 2277 | assert!( |
| 2278 | scan(¤t_view) |
| 2279 | .publish(&mut manager.shared.attachments.lock().unwrap()) |
| 2280 | .is_none() |
| 2281 | ); |
| 2282 | drop(other); |
| 2283 | let stale = scan(¤t_view); |
| 2284 | drop(attachment); |
| 2285 | assert!( |
| 2286 | stale |
| 2287 | .publish(&mut manager.shared.attachments.lock().unwrap()) |
| 2288 | .is_none() |
| 2289 | ); |
| 2290 | } |
| 2291 | |
| 2292 | pub(crate) fn installed(engine: &HostAttachment, workspace: &Path) -> Vec<String> { |
| 2293 | let mut registry = crate::tools::registry::ToolRegistryBuilder::new() |
| 2294 | .build(ToolContext::new(workspace).with_plugin_registry(engine.plugin_view())); |
| 2295 | engine.install_tools(&mut registry) |
| 2296 | } |
| 2297 | |
| 2298 | pub(crate) fn plugin_id(fixture: &FixturePlugins, name: &str) -> String { |
| 2299 | fixture |
| 2300 | .registry() |
| 2301 | .get(name) |
| 2302 | .unwrap() |
| 2303 | .id |
| 2304 | .as_str() |
| 2305 | .to_string() |
| 2306 | } |
| 2307 | |
| 2308 | /// Two engines for different workspaces in one process: syncing either |
| 2309 | /// keeps the other's plugin active and its in-flight call running, neither |
| 2310 | /// receives the other's tools, and detaching one revokes only its plugin. |
| 2311 | #[tokio::test] |
| 2312 | async fn engines_in_one_process_never_revoke_each_others_plugins() { |
| 2313 | let Some(node) = node_for_tests("engines_in_one_process") else { |
| 2314 | return; |
| 2315 | }; |
| 2316 | let _policy = TestPolicyGuard::extension_host(true); |
| 2317 | let slow = FixturePlugins::new(&["slow-tool"]).await; |
| 2318 | let deps = FixturePlugins::new(&["dsh-workspace-deps"]).await; |
| 2319 | let (slow_id, deps_id) = ( |
| 2320 | plugin_id(&slow, "slow-tool"), |
| 2321 | plugin_id(&deps, "dsh-workspace-deps"), |
| 2322 | ); |
| 2323 | let manager = slow.manager(node); |
| 2324 | let first = manager.attach(slow.registry()); |
| 2325 | first.sync().await.unwrap(); |
| 2326 | let tool = host_tool(&first, slow.workspace(), "slow_wait"); |
| 2327 | let context = ToolContext::new(slow.workspace()).with_plugin_registry(first.plugin_view()); |
| 2328 | let call = tokio::spawn(async move { tool.execute(json!({"ms": 600}), &context).await }); |
| 2329 | tokio::time::sleep(Duration::from_millis(100)).await; |
| 2330 | |
| 2331 | // A second workspace's engine attaches and syncs mid-call. |
| 2332 | let second = manager.attach(deps.registry()); |
| 2333 | second.sync().await.unwrap(); |
| 2334 | assert_eq!(manager.owner_state(&slow_id), Some(OwnerState::Active)); |
| 2335 | assert_eq!(manager.owner_state(&deps_id), Some(OwnerState::Active)); |
| 2336 | let result = tokio::time::timeout(Duration::from_secs(5), call) |
| 2337 | .await |
| 2338 | .expect("call resolves") |
| 2339 | .unwrap() |
| 2340 | .expect("the first engine's in-flight call completes"); |
| 2341 | assert!(result.success, "{}", result.content); |
| 2342 | assert!(result.content.contains("600"), "{}", result.content); |
| 2343 | |
| 2344 | assert_eq!(installed(&first, slow.workspace()), vec!["slow_wait"]); |
| 2345 | assert_eq!( |
| 2346 | installed(&second, deps.workspace()), |
| 2347 | vec!["load_workspace_dependencies"] |
| 2348 | ); |
| 2349 | first.sync().await.unwrap(); |
| 2350 | assert_eq!(manager.owner_state(&deps_id), Some(OwnerState::Active)); |
| 2351 | assert!( |
| 2352 | !manager.diagnostics().iter().any(|d| d.contains("revoked")), |
| 2353 | "{:?}", |
| 2354 | manager.diagnostics() |
| 2355 | ); |
| 2356 | |
| 2357 | // Detaching does not revoke by itself; the next reconcile revokes only |
| 2358 | // what no remaining engine desires. |
| 2359 | drop(second); |
| 2360 | assert_eq!(manager.owner_state(&deps_id), Some(OwnerState::Active)); |
| 2361 | first.sync().await.unwrap(); |
| 2362 | assert_eq!(manager.owner_state(&deps_id), None); |
| 2363 | assert_eq!(manager.owner_state(&slow_id), Some(OwnerState::Active)); |
| 2364 | assert_eq!(manager.spawn_attempts(), 1); |
| 2365 | manager.shutdown().await; |
| 2366 | } |
| 2367 | |
| 2368 | /// Each snapshot is re-verified against persisted plugin state, so a disable |
| 2369 | /// made through one engine's registry revokes the plugin for an engine still |
| 2370 | /// holding the older snapshot. |
| 2371 | #[tokio::test] |
| 2372 | async fn a_disable_through_either_registry_revokes_for_every_engine() { |
| 2373 | let Some(node) = node_for_tests("a_disable_through_either_registry") else { |
| 2374 | return; |
| 2375 | }; |
| 2376 | let _policy = TestPolicyGuard::extension_host(true); |
| 2377 | let fixture = FixturePlugins::new(&["slow-tool"]).await; |
| 2378 | let id = plugin_id(&fixture, "slow-tool"); |
| 2379 | let manager = fixture.manager(node); |
| 2380 | let first = manager.attach(fixture.registry()); |
| 2381 | let second = manager.attach(fixture.registry()); |
| 2382 | first.sync().await.unwrap(); |
| 2383 | assert_eq!(installed(&first, fixture.workspace()), vec!["slow_wait"]); |
| 2384 | assert_eq!(installed(&second, fixture.workspace()), vec!["slow_wait"]); |
| 2385 | |
| 2386 | second.set_plugins(fixture.disable("slow-tool")); |
| 2387 | second.sync().await.unwrap(); |
| 2388 | assert_eq!(manager.owner_state(&id), None, "revoked and forgotten"); |
| 2389 | assert!(installed(&first, fixture.workspace()).is_empty()); |
| 2390 | assert!(installed(&second, fixture.workspace()).is_empty()); |
| 2391 | manager.shutdown().await; |
| 2392 | } |
| 2393 | |
| 2394 | fn fast_supervision() -> super::SupervisionOptions { |
| 2395 | super::SupervisionOptions { |
| 2396 | heartbeat_interval: Duration::from_millis(50), |
| 2397 | ping_timeout: Duration::from_millis(150), |
| 2398 | hang_timeout: Duration::from_millis(600), |
| 2399 | restart_backoff: Duration::from_millis(25), |
| 2400 | ..Default::default() |
| 2401 | } |
| 2402 | } |
| 2403 | |
| 2404 | fn supervised_manager(fixture: &FixturePlugins, node: PathBuf) -> Arc<ExtensionHostManager> { |
| 2405 | Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 2406 | runtime: NODE, |
| 2407 | node_override: Some(node), |
| 2408 | bun_override: None, |
| 2409 | root: Some(fixture.root.clone()), |
| 2410 | supervision: fast_supervision(), |
| 2411 | })) |
| 2412 | } |
| 2413 | |
| 2414 | async fn wait_host(manager: &ExtensionHostManager, predicate: impl Fn() -> bool) { |
| 2415 | tokio::time::timeout(Duration::from_secs(10), async { |
| 2416 | while !predicate() { |
| 2417 | tokio::time::sleep(Duration::from_millis(10)).await; |
| 2418 | } |
| 2419 | }) |
| 2420 | .await |
| 2421 | .unwrap_or_else(|_| { |
| 2422 | panic!( |
| 2423 | "host wait timed out: {:?}; {:?}", |
| 2424 | manager.status(), |
| 2425 | manager.diagnostics() |
| 2426 | ) |
| 2427 | }); |
| 2428 | } |
| 2429 | |
| 2430 | #[test] |
| 2431 | fn crash_budget_is_bounded_and_expires_only_with_the_window() { |
| 2432 | let options = super::SupervisionOptions::default(); |
| 2433 | let mut state = super::SupervisionState::default(); |
| 2434 | let now = Instant::now(); |
| 2435 | assert!(state.record_crash(now, &options)); |
| 2436 | assert!(state.record_crash(now + Duration::from_secs(1), &options)); |
| 2437 | assert!(!state.record_crash(now + Duration::from_secs(2), &options)); |
| 2438 | assert_eq!(state.crashes.len(), 3); |
| 2439 | assert!(state.record_crash( |
| 2440 | now + options.crash_window + Duration::from_secs(3), |
| 2441 | &options |
| 2442 | )); |
| 2443 | assert_eq!(state.crashes.len(), 1); |
| 2444 | } |
| 2445 | |
| 2446 | #[test] |
| 2447 | fn dirty_teardown_window_is_bounded_and_a_requested_restart_waits_for_idle() { |
| 2448 | let options = super::SupervisionOptions::default(); |
| 2449 | let mut state = super::SupervisionState::default(); |
| 2450 | let now = Instant::now(); |
| 2451 | state.record_dirty_teardown(now, &options); |
| 2452 | assert!(!state.dirty_restart_pending); |
| 2453 | state.record_dirty_teardown(now + options.dirty_window, &options); |
| 2454 | assert!(!state.dirty_restart_pending, "the first event expired"); |
| 2455 | state.record_dirty_teardown( |
| 2456 | now + options.dirty_window + Duration::from_secs(1), |
| 2457 | &options, |
| 2458 | ); |
| 2459 | assert!(state.dirty_restart_pending); |
| 2460 | for second in 2..100 { |
| 2461 | state.record_dirty_teardown( |
| 2462 | now + options.dirty_window + Duration::from_secs(second), |
| 2463 | &options, |
| 2464 | ); |
| 2465 | } |
| 2466 | assert_eq!(state.dirty_teardowns.len(), 2); |
| 2467 | state.record_dirty_teardown(now + options.dirty_window * 3, &options); |
| 2468 | assert!( |
| 2469 | state.dirty_restart_pending, |
| 2470 | "an idle request does not expire" |
| 2471 | ); |
| 2472 | assert!(state.crashes.is_empty()); |
| 2473 | } |
| 2474 | |
| 2475 | #[tokio::test] |
| 2476 | async fn a_call_past_its_method_deadline_is_cancelled_and_the_host_stays_usable() { |
| 2477 | let Some(node) = node_for_tests("method deadline") else { |
| 2478 | return; |
| 2479 | }; |
| 2480 | let _policy = TestPolicyGuard::extension_host(true); |
| 2481 | let fixture = FixturePlugins::new(&["slow-tool", "clash-script"]).await; |
| 2482 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 2483 | runtime: NODE, |
| 2484 | node_override: Some(node), |
| 2485 | bun_override: None, |
| 2486 | root: Some(fixture.root.clone()), |
| 2487 | supervision: super::SupervisionOptions { |
| 2488 | tool_call_deadline: Duration::from_millis(300), |
| 2489 | ..Default::default() |
| 2490 | }, |
| 2491 | })); |
| 2492 | let engine = manager.attach(fixture.registry()); |
| 2493 | engine.sync().await.unwrap(); |
| 2494 | let pid = manager.host_pid(); |
| 2495 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()); |
| 2496 | // `slow_wait` answers after 30 s unless it is cancelled. |
| 2497 | let slow = host_tool(&engine, fixture.workspace(), "slow_wait"); |
| 2498 | let started = Instant::now(); |
| 2499 | let outcome = tokio::time::timeout(Duration::from_secs(5), slow.execute(json!({}), &context)) |
| 2500 | .await |
| 2501 | .expect("the method deadline bounds the call"); |
| 2502 | let elapsed = started.elapsed(); |
| 2503 | assert!( |
| 2504 | matches!(outcome, Err(ToolError::Timeout { .. })), |
| 2505 | "{outcome:?}" |
| 2506 | ); |
| 2507 | assert!(elapsed < Duration::from_secs(2), "{elapsed:?}"); |
| 2508 | // The same host takes the next call. |
| 2509 | let quick = host_tool(&engine, fixture.workspace(), "fixture_script_tool"); |
| 2510 | let result = quick.execute(json!({}), &context).await.unwrap(); |
| 2511 | assert!(result.success, "{}", result.content); |
| 2512 | assert_eq!(manager.host_pid(), pid); |
| 2513 | assert_eq!(manager.spawn_attempts(), 1); |
| 2514 | // `slow-tool`'s disposer awaits its in-flight call, so a clean teardown |
| 2515 | // proves the host received `$/cancel` and aborted the timed-out call |
| 2516 | // (uncancelled, it outlives the host's 2 s dispose deadline). |
| 2517 | engine.set_plugins(fixture.disable("slow-tool")); |
| 2518 | engine.sync().await.unwrap(); |
| 2519 | let diagnostics = manager.diagnostics(); |
| 2520 | assert!( |
| 2521 | !diagnostics.iter().any(|line| line.contains("teardown")), |
| 2522 | "{diagnostics:?}" |
| 2523 | ); |
| 2524 | manager.shutdown().await; |
| 2525 | } |
| 2526 | |
| 2527 | #[tokio::test] |
| 2528 | async fn ordinary_exit_rejects_requests_from_a_drained_calls_waker() { |
| 2529 | use std::future::Future; |
| 2530 | use std::pin::Pin; |
| 2531 | use std::sync::Mutex; |
| 2532 | use std::task::{Context, Wake, Waker}; |
| 2533 | |
| 2534 | use super::supervisor::{HostCallError, HostProcess}; |
| 2535 | |
| 2536 | struct AdmissionProbe { |
| 2537 | host: Arc<HostProcess>, |
| 2538 | result: Mutex<Option<Result<(), HostCallError>>>, |
| 2539 | woke: tokio::sync::Notify, |
| 2540 | } |
| 2541 | |
| 2542 | impl Wake for AdmissionProbe { |
| 2543 | fn wake(self: Arc<Self>) { |
| 2544 | self.wake_by_ref(); |
| 2545 | } |
| 2546 | |
| 2547 | fn wake_by_ref(self: &Arc<Self>) { |
| 2548 | let mut result = self.result.lock().unwrap(); |
| 2549 | if result.is_none() { |
| 2550 | // oneshot send wakes synchronously: probe the exact gap after |
| 2551 | // the exit watcher drains pending and starts failing its calls. |
| 2552 | *result = Some( |
| 2553 | self.host |
| 2554 | .start_request(protocol::CoreRequest::Ping, None) |
| 2555 | .map(|(id, _)| self.host.forget(id)), |
| 2556 | ); |
| 2557 | self.woke.notify_one(); |
| 2558 | } |
| 2559 | } |
| 2560 | } |
| 2561 | |
| 2562 | let Some(node) = node_for_tests("ordinary exit admission") else { |
| 2563 | return; |
| 2564 | }; |
| 2565 | let _policy = TestPolicyGuard::extension_host(true); |
| 2566 | let fixture = FixturePlugins::new(&["slow-tool"]).await; |
| 2567 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 2568 | runtime: NODE, |
| 2569 | node_override: Some(node), |
| 2570 | bun_override: None, |
| 2571 | root: Some(fixture.root.clone()), |
| 2572 | supervision: super::SupervisionOptions { |
| 2573 | heartbeat_interval: Duration::from_secs(60), |
| 2574 | ..Default::default() |
| 2575 | }, |
| 2576 | })); |
| 2577 | let engine = manager.attach(fixture.registry()); |
| 2578 | engine.sync().await.unwrap(); |
| 2579 | let host = manager.shared.ready_host(HostTier::Plugin).unwrap(); |
| 2580 | let registration = manager.shared.registry.lock().unwrap().live_tools()[0].clone(); |
| 2581 | let (_, mut call) = host |
| 2582 | .start_request( |
| 2583 | protocol::CoreRequest::ToolCall(protocol::ToolCallParams { |
| 2584 | handle: registration.handle, |
| 2585 | call_id: "exit-admission".into(), |
| 2586 | input: json!({"ms": 30_000}), |
| 2587 | deadline_ms: 60_000, |
| 2588 | workspace: None, |
| 2589 | ticket: None, |
| 2590 | session_id: None, |
| 2591 | agent_id: None, |
| 2592 | origin_turn_id: None, |
| 2593 | }), |
| 2594 | Some(registration.owner.plugin_id), |
| 2595 | ) |
| 2596 | .unwrap(); |
| 2597 | // A following response proves the writer flushed the slow call and is |
| 2598 | // waiting for another frame, so a closed outbound queue cannot mask the bug. |
| 2599 | host.call(protocol::CoreRequest::Ping, None).await.unwrap(); |
| 2600 | let probe = Arc::new(AdmissionProbe { |
| 2601 | host: Arc::clone(&host), |
| 2602 | result: Mutex::new(None), |
| 2603 | woke: tokio::sync::Notify::new(), |
| 2604 | }); |
| 2605 | let waker = Waker::from(Arc::clone(&probe)); |
| 2606 | assert!( |
| 2607 | Pin::new(&mut call) |
| 2608 | .poll(&mut Context::from_waker(&waker)) |
| 2609 | .is_pending() |
| 2610 | ); |
| 2611 | assert!( |
| 2612 | !host.is_retiring(), |
| 2613 | "this exercises exit, not maintenance sealing" |
| 2614 | ); |
| 2615 | host.terminate("ordinary exit admission regression".into()); |
| 2616 | tokio::time::timeout(Duration::from_secs(5), probe.woke.notified()) |
| 2617 | .await |
| 2618 | .expect("exit must fail the pending call"); |
| 2619 | assert!(matches!( |
| 2620 | probe.result.lock().unwrap().take().unwrap(), |
| 2621 | Err(HostCallError::Exited(_)) |
| 2622 | )); |
| 2623 | assert!(matches!(call.await.unwrap(), Err(HostCallError::Exited(_)))); |
| 2624 | manager.shutdown().await; |
| 2625 | } |
| 2626 | |
| 2627 | #[tokio::test] |
| 2628 | async fn idle_retirement_seals_admission_and_does_not_wait_for_heartbeat() { |
| 2629 | use futures_util::FutureExt; |
| 2630 | |
| 2631 | let Some(node) = node_for_tests("idle admission") else { |
| 2632 | return; |
| 2633 | }; |
| 2634 | let _policy = TestPolicyGuard::extension_host(true); |
| 2635 | let fixture = FixturePlugins::new(&["slow-tool"]).await; |
| 2636 | let manager = supervised_manager(&fixture, node); |
| 2637 | let engine = manager.attach(fixture.registry()); |
| 2638 | engine.sync().await.unwrap(); |
| 2639 | let host = manager.shared.ready_host(HostTier::Plugin).unwrap(); |
| 2640 | let registration = manager.shared.registry.lock().unwrap().live_tools()[0].clone(); |
| 2641 | let (_, call) = host |
| 2642 | .start_request( |
| 2643 | protocol::CoreRequest::ToolCall(protocol::ToolCallParams { |
| 2644 | handle: registration.handle, |
| 2645 | call_id: "idle-admission".into(), |
| 2646 | input: json!({"ms": 100}), |
| 2647 | deadline_ms: 5000, |
| 2648 | workspace: None, |
| 2649 | ticket: None, |
| 2650 | session_id: None, |
| 2651 | agent_id: None, |
| 2652 | origin_turn_id: None, |
| 2653 | }), |
| 2654 | Some(registration.owner.plugin_id), |
| 2655 | ) |
| 2656 | .unwrap(); |
| 2657 | assert!(!host.terminate_if_idle(super::DIRTY_RESTART_REASON)); |
| 2658 | assert!(call.await.unwrap().is_ok()); |
| 2659 | // No await between admission and retirement: this heartbeat is still in |
| 2660 | // the pending map, but it does not make the process busy. |
| 2661 | let (_, _heartbeat) = host |
| 2662 | .start_request(protocol::CoreRequest::Ping, None) |
| 2663 | .unwrap(); |
| 2664 | assert!(host.terminate_if_idle(super::DIRTY_RESTART_REASON)); |
| 2665 | assert!(matches!( |
| 2666 | manager.shared.ready_host(HostTier::Plugin), |
| 2667 | Err(HostStatus::Restarting { .. }) |
| 2668 | )); |
| 2669 | assert!(matches!(manager.status(), HostStatus::Restarting { .. })); |
| 2670 | // Poll without yielding to the exit watcher. A reconcile in this exact |
| 2671 | // gap must not start activation on the sealed process and falsely fail |
| 2672 | // a valid receipt before replay. |
| 2673 | assert!( |
| 2674 | manager |
| 2675 | .ensure_host(HostTier::Plugin, true) |
| 2676 | .now_or_never() |
| 2677 | .unwrap() |
| 2678 | .is_err() |
| 2679 | ); |
| 2680 | assert_eq!( |
| 2681 | manager.owner_state(&plugin_id(&fixture, "slow-tool")), |
| 2682 | Some(OwnerState::Active) |
| 2683 | ); |
| 2684 | assert!(matches!( |
| 2685 | host.start_request(protocol::CoreRequest::Ping, None), |
| 2686 | Err(super::supervisor::HostCallError::Exited(_)) |
| 2687 | )); |
| 2688 | assert!(!host.terminate_if_idle(super::DIRTY_RESTART_REASON)); |
| 2689 | manager.shutdown().await; |
| 2690 | } |
| 2691 | |
| 2692 | #[tokio::test] |
| 2693 | async fn two_dirty_teardowns_wait_for_a_live_call_then_replay_without_spending_crash_budget() { |
| 2694 | let Some(node) = node_for_tests("dirty teardown") else { |
| 2695 | return; |
| 2696 | }; |
| 2697 | let _policy = TestPolicyGuard::extension_host(true); |
| 2698 | let fixture = FixturePlugins::new(&["dirty-dispose", "slow-tool"]).await; |
| 2699 | let manager = supervised_manager(&fixture, node); |
| 2700 | let engine = manager.attach(fixture.registry()); |
| 2701 | engine.sync().await.unwrap(); |
| 2702 | // An existing unexpected crash must survive planned maintenance. |
| 2703 | manager |
| 2704 | .shared |
| 2705 | .plugin |
| 2706 | .supervision |
| 2707 | .lock() |
| 2708 | .unwrap() |
| 2709 | .record_crash(Instant::now(), &manager.shared.options.supervision); |
| 2710 | engine.set_plugins(fixture.disable("dirty-dispose")); |
| 2711 | engine.sync().await.unwrap(); |
| 2712 | assert_eq!( |
| 2713 | manager |
| 2714 | .shared |
| 2715 | .plugin |
| 2716 | .supervision |
| 2717 | .lock() |
| 2718 | .unwrap() |
| 2719 | .dirty_teardowns |
| 2720 | .len(), |
| 2721 | 1 |
| 2722 | ); |
| 2723 | tokio::time::sleep(Duration::from_millis(150)).await; |
| 2724 | assert_eq!( |
| 2725 | manager.spawn_attempts(), |
| 2726 | 1, |
| 2727 | "one dirty event is insufficient" |
| 2728 | ); |
| 2729 | |
| 2730 | let mut enabled = discover_with_config(&fixture.config); |
| 2731 | enabled.enable("dirty-dispose").unwrap(); |
| 2732 | engine.set_plugins(Arc::new(enabled)); |
| 2733 | engine.sync().await.unwrap(); |
| 2734 | let old = host_tool(&engine, fixture.workspace(), "slow_wait"); |
| 2735 | let host = manager.shared.ready_host(HostTier::Plugin).unwrap(); |
| 2736 | let registration = manager.shared.registry.lock().unwrap().live_tools()[0].clone(); |
| 2737 | let (_, call) = host |
| 2738 | .start_request( |
| 2739 | protocol::CoreRequest::ToolCall(protocol::ToolCallParams { |
| 2740 | handle: registration.handle, |
| 2741 | call_id: "survives-dirty-teardown".into(), |
| 2742 | input: json!({"ms": 4000}), |
| 2743 | deadline_ms: 10000, |
| 2744 | workspace: None, |
| 2745 | ticket: None, |
| 2746 | session_id: None, |
| 2747 | agent_id: None, |
| 2748 | origin_turn_id: None, |
| 2749 | }), |
| 2750 | Some(registration.owner.plugin_id.clone()), |
| 2751 | ) |
| 2752 | .unwrap(); |
| 2753 | engine.set_plugins(fixture.disable("dirty-dispose")); |
| 2754 | engine.sync().await.unwrap(); |
| 2755 | assert!( |
| 2756 | manager |
| 2757 | .shared |
| 2758 | .plugin |
| 2759 | .supervision |
| 2760 | .lock() |
| 2761 | .unwrap() |
| 2762 | .dirty_restart_pending |
| 2763 | ); |
| 2764 | tokio::time::sleep(Duration::from_millis(150)).await; |
| 2765 | assert_eq!(manager.spawn_attempts(), 1, "a live call defers retirement"); |
| 2766 | let completed = tokio::time::timeout(Duration::from_secs(10), call) |
| 2767 | .await |
| 2768 | .unwrap() |
| 2769 | .unwrap() |
| 2770 | .unwrap(); |
| 2771 | assert_eq!(completed["structured"]["waited"], 4000); |
| 2772 | wait_host(&manager, || { |
| 2773 | manager.spawn_attempts() == 2 && manager.live_tool_names().contains(&"slow_wait".into()) |
| 2774 | }) |
| 2775 | .await; |
| 2776 | assert_eq!( |
| 2777 | manager |
| 2778 | .shared |
| 2779 | .plugin |
| 2780 | .supervision |
| 2781 | .lock() |
| 2782 | .unwrap() |
| 2783 | .crashes |
| 2784 | .len(), |
| 2785 | 1 |
| 2786 | ); |
| 2787 | assert!( |
| 2788 | !manager |
| 2789 | .shared |
| 2790 | .plugin |
| 2791 | .supervision |
| 2792 | .lock() |
| 2793 | .unwrap() |
| 2794 | .dirty_restart_pending |
| 2795 | ); |
| 2796 | let replayed = manager.shared.registry.lock().unwrap().live_tools()[0].clone(); |
| 2797 | assert_ne!(registration.owner, replayed.owner); |
| 2798 | assert_ne!(registration.handle, replayed.handle); |
| 2799 | assert!(matches!( |
| 2800 | old.execute( |
| 2801 | json!({"ms": 1}), |
| 2802 | &ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()) |
| 2803 | ) |
| 2804 | .await, |
| 2805 | Err(ToolError::NotAvailable { .. }) |
| 2806 | )); |
| 2807 | // A delayed outcome from the retired process cannot dirty its replacement. |
| 2808 | manager.shared.record_dirty_teardown(&host); |
| 2809 | assert!( |
| 2810 | manager |
| 2811 | .shared |
| 2812 | .plugin |
| 2813 | .supervision |
| 2814 | .lock() |
| 2815 | .unwrap() |
| 2816 | .dirty_teardowns |
| 2817 | .is_empty() |
| 2818 | ); |
| 2819 | manager.shutdown().await; |
| 2820 | } |
| 2821 | |
| 2822 | #[tokio::test] |
| 2823 | async fn failed_activation_cleanup_also_records_a_dirty_teardown() { |
| 2824 | let Some(node) = node_for_tests("failed activation teardown") else { |
| 2825 | return; |
| 2826 | }; |
| 2827 | let _policy = TestPolicyGuard::extension_host(true); |
| 2828 | let fixture = FixturePlugins::new(&["clash-script"]).await; |
| 2829 | native_bundle( |
| 2830 | &fixture.config.user_plugins_dir, |
| 2831 | "failed-disposal", |
| 2832 | "index.mjs", |
| 2833 | &["index.mjs"], |
| 2834 | ); |
| 2835 | std::fs::write( |
| 2836 | fixture.config.user_plugins_dir.join("failed-disposal/index.mjs"), |
| 2837 | "export const name = 'failed-disposal';\nexport function apply(ctx) {\n ctx.effect(() => () => new Promise(() => {}), 'unfinished activation cleanup');\n throw new Error('fixture activation failure');\n}\n", |
| 2838 | ).unwrap(); |
| 2839 | let mut plugins = discover_with_config(&fixture.config); |
| 2840 | plugins.trust("failed-disposal").unwrap(); |
| 2841 | plugins.enable("failed-disposal").unwrap(); |
| 2842 | let manager = supervised_manager(&fixture, node); |
| 2843 | let engine = manager.attach(Arc::new(plugins)); |
| 2844 | tokio::time::timeout(Duration::from_secs(15), engine.sync()) |
| 2845 | .await |
| 2846 | .unwrap() |
| 2847 | .unwrap(); |
| 2848 | assert!(matches!( |
| 2849 | manager.owner_state(&plugin_id(&fixture, "failed-disposal")), |
| 2850 | Some(OwnerState::Failed(_)) |
| 2851 | )); |
| 2852 | assert_eq!( |
| 2853 | manager |
| 2854 | .shared |
| 2855 | .plugin |
| 2856 | .supervision |
| 2857 | .lock() |
| 2858 | .unwrap() |
| 2859 | .dirty_teardowns |
| 2860 | .len(), |
| 2861 | 1 |
| 2862 | ); |
| 2863 | assert!( |
| 2864 | manager |
| 2865 | .shared |
| 2866 | .plugin |
| 2867 | .supervision |
| 2868 | .lock() |
| 2869 | .unwrap() |
| 2870 | .crashes |
| 2871 | .is_empty() |
| 2872 | ); |
| 2873 | assert_eq!(manager.spawn_attempts(), 1); |
| 2874 | assert!( |
| 2875 | manager |
| 2876 | .live_tool_names() |
| 2877 | .contains(&"fixture_script_tool".into()) |
| 2878 | ); |
| 2879 | manager.shutdown().await; |
| 2880 | } |
| 2881 | |
| 2882 | #[test] |
| 2883 | fn host_exit_preserves_failed_receipts_and_blames_only_the_activating_owner() { |
| 2884 | let mut registry = OwnerRegistry::new(); |
| 2885 | let active = registry |
| 2886 | .begin_owner( |
| 2887 | HostTier::Plugin, |
| 2888 | "healthy", |
| 2889 | "healthy", |
| 2890 | Some(fake_authority("healthy")), |
| 2891 | "hash-healthy", |
| 2892 | ) |
| 2893 | .unwrap(); |
| 2894 | registry.mark_active(&active); |
| 2895 | register(&mut registry, &active, "healthy_probe").unwrap(); |
| 2896 | let failed = registry |
| 2897 | .begin_owner( |
| 2898 | HostTier::Plugin, |
| 2899 | "failed", |
| 2900 | "failed", |
| 2901 | Some(fake_authority("failed")), |
| 2902 | "hash-failed", |
| 2903 | ) |
| 2904 | .unwrap(); |
| 2905 | registry.mark_failed(&failed, OwnerState::Faulted("existing fault".into())); |
| 2906 | registry |
| 2907 | .begin_owner( |
| 2908 | HostTier::Plugin, |
| 2909 | "activating", |
| 2910 | "activating", |
| 2911 | Some(fake_authority("activating")), |
| 2912 | "hash-activating", |
| 2913 | ) |
| 2914 | .unwrap(); |
| 2915 | registry.host_exited(HostTier::Plugin, "fixture crash"); |
| 2916 | assert!(registry.owner("healthy").is_none()); |
| 2917 | assert!(registry.live_tools().is_empty()); |
| 2918 | assert!(matches!( |
| 2919 | registry.owner("failed").unwrap().state, |
| 2920 | OwnerState::Faulted(_) |
| 2921 | )); |
| 2922 | assert!(matches!( |
| 2923 | registry.owner("activating").unwrap().state, |
| 2924 | OwnerState::Failed(_) |
| 2925 | )); |
| 2926 | let replay = registry |
| 2927 | .begin_owner( |
| 2928 | HostTier::Plugin, |
| 2929 | "healthy", |
| 2930 | "healthy", |
| 2931 | Some(fake_authority("healthy")), |
| 2932 | "hash-healthy", |
| 2933 | ) |
| 2934 | .unwrap(); |
| 2935 | assert_ne!(replay.generation, active.generation); |
| 2936 | assert_ne!(replay.owner_token, active.owner_token); |
| 2937 | } |
| 2938 | |
| 2939 | #[test] |
| 2940 | fn opening_an_engine_never_resets_a_crash_budget() { |
| 2941 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions::default())); |
| 2942 | { |
| 2943 | *manager.shared.plugin.host.lock().unwrap() = super::HostSlot::Failed { |
| 2944 | reason: "budget".into(), |
| 2945 | stderr_tail: String::new(), |
| 2946 | }; |
| 2947 | let mut state = manager.shared.plugin.supervision.lock().unwrap(); |
| 2948 | for _ in 0..3 { |
| 2949 | state.record_crash(Instant::now(), &manager.shared.options.supervision); |
| 2950 | } |
| 2951 | } |
| 2952 | let _engine = manager.attach(Arc::new(PluginRegistry::empty(Path::new("/fixture")))); |
| 2953 | assert_eq!( |
| 2954 | manager |
| 2955 | .shared |
| 2956 | .plugin |
| 2957 | .supervision |
| 2958 | .lock() |
| 2959 | .unwrap() |
| 2960 | .crashes |
| 2961 | .len(), |
| 2962 | 3 |
| 2963 | ); |
| 2964 | assert!(matches!(manager.status(), HostStatus::Failed { .. })); |
| 2965 | manager.retry(); |
| 2966 | assert!( |
| 2967 | manager |
| 2968 | .shared |
| 2969 | .plugin |
| 2970 | .supervision |
| 2971 | .lock() |
| 2972 | .unwrap() |
| 2973 | .crashes |
| 2974 | .is_empty() |
| 2975 | ); |
| 2976 | assert_eq!(manager.status(), HostStatus::Idle); |
| 2977 | } |
| 2978 | |
| 2979 | /// One typed state, `HostStatus`, names why the host is down in `/plugin` |
| 2980 | /// and in the error a tool call routed to it gets. |
| 2981 | #[tokio::test] |
| 2982 | async fn a_host_that_is_down_names_why_in_plugin_status_and_tool_errors() { |
| 2983 | let policy = TestPolicyGuard::extension_host(true); |
| 2984 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions::default())); |
| 2985 | let registration = { |
| 2986 | let mut registry = manager.shared.registry.lock().unwrap(); |
| 2987 | let owner = registry |
| 2988 | .begin_owner( |
| 2989 | HostTier::Plugin, |
| 2990 | "probe", |
| 2991 | "probe", |
| 2992 | Some(fake_authority("probe")), |
| 2993 | "hash", |
| 2994 | ) |
| 2995 | .unwrap(); |
| 2996 | registry.mark_active(&owner); |
| 2997 | register(&mut registry, &owner, "probe_tool").unwrap(); |
| 2998 | registry.live_tools()[0].clone() |
| 2999 | }; |
| 3000 | let tool = super::tool::HostToolSpec::new(registration, Arc::clone(&manager.shared)); |
| 3001 | let context = ToolContext::new(Path::new("/fixture")); |
| 3002 | let refused = "start failed: host did not apply the requested 1024 MiB kernel memory limit; initialization refused"; |
| 3003 | for (slot, why) in [ |
| 3004 | ( |
| 3005 | super::HostSlot::Restarting { |
| 3006 | reason: "exited with signal: 9 (SIGKILL)".into(), |
| 3007 | }, |
| 3008 | "restarting after: exited with signal: 9 (SIGKILL)".to_string(), |
| 3009 | ), |
| 3010 | ( |
| 3011 | super::HostSlot::Failed { |
| 3012 | reason: refused.into(), |
| 3013 | stderr_tail: String::new(), |
| 3014 | }, |
| 3015 | format!("{refused} (change or reload a plugin to retry)"), |
| 3016 | ), |
| 3017 | (super::HostSlot::Idle, "not started".to_string()), |
| 3018 | ] { |
| 3019 | *manager.shared.plugin.host.lock().unwrap() = slot; |
| 3020 | let report = super::render_status(&manager); |
| 3021 | assert!( |
| 3022 | report.starts_with(&format!("Extension host (experimental): {why}")), |
| 3023 | "{report}" |
| 3024 | ); |
| 3025 | match tool.execute(json!({}), &context).await { |
| 3026 | Err(ToolError::NotAvailable { message }) => assert!( |
| 3027 | message.starts_with(&format!("extension host is down: {why}")), |
| 3028 | "{message}" |
| 3029 | ), |
| 3030 | other => panic!("expected a typed not-available error, got {other:?}"), |
| 3031 | } |
| 3032 | } |
| 3033 | drop(policy); |
| 3034 | let _off = TestPolicyGuard::extension_host(false); |
| 3035 | let disabled = "disabled by config ([features] extension_host is off)"; |
| 3036 | assert_eq!( |
| 3037 | super::status_report(), |
| 3038 | format!("Extension host (experimental): {disabled}") |
| 3039 | ); |
| 3040 | match tool.execute(json!({}), &context).await { |
| 3041 | Err(ToolError::NotAvailable { message }) => { |
| 3042 | assert_eq!(message, format!("extension host is down: {disabled}")) |
| 3043 | } |
| 3044 | other => panic!("expected a typed not-available error, got {other:?}"), |
| 3045 | } |
| 3046 | } |
| 3047 | |
| 3048 | #[tokio::test] |
| 3049 | async fn three_crashes_stop_replay_until_explicit_retry() { |
| 3050 | let Some(node) = node_for_tests("crash budget") else { |
| 3051 | return; |
| 3052 | }; |
| 3053 | let _policy = TestPolicyGuard::extension_host(true); |
| 3054 | let fixture = FixturePlugins::new(&["crash-tool", "refuses-approval"]).await; |
| 3055 | let manager = supervised_manager(&fixture, node); |
| 3056 | let engine = manager.attach(fixture.registry()); |
| 3057 | engine.sync().await.unwrap(); |
| 3058 | let failed_id = plugin_id(&fixture, "refuses-approval"); |
| 3059 | let mut previous = None; |
| 3060 | for crash in 1..=3 { |
| 3061 | let tool = host_tool(&engine, fixture.workspace(), "crash_probe"); |
| 3062 | let registration = manager |
| 3063 | .shared |
| 3064 | .registry |
| 3065 | .lock() |
| 3066 | .unwrap() |
| 3067 | .live_tools() |
| 3068 | .into_iter() |
| 3069 | .find(|t| t.name == "crash_probe") |
| 3070 | .unwrap(); |
| 3071 | if let Some(old) = previous { |
| 3072 | assert_ne!(registration.owner, old); |
| 3073 | } |
| 3074 | previous = Some(registration.owner); |
| 3075 | let outcome = tool |
| 3076 | .execute( |
| 3077 | json!({}), |
| 3078 | &ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()), |
| 3079 | ) |
| 3080 | .await; |
| 3081 | assert!(matches!(outcome, Err(ToolError::NotAvailable { .. }))); |
| 3082 | if crash < 3 { |
| 3083 | wait_host(&manager, || { |
| 3084 | manager.spawn_attempts() == crash + 1 |
| 3085 | && manager.live_tool_names().contains(&"crash_probe".into()) |
| 3086 | }) |
| 3087 | .await; |
| 3088 | assert!(matches!( |
| 3089 | manager.owner_state(&failed_id), |
| 3090 | Some(OwnerState::Failed(_)) |
| 3091 | )); |
| 3092 | } else { |
| 3093 | wait_host(&manager, || { |
| 3094 | matches!(manager.status(), HostStatus::Failed { .. }) |
| 3095 | }) |
| 3096 | .await; |
| 3097 | } |
| 3098 | } |
| 3099 | assert_eq!(manager.spawn_attempts(), 3); |
| 3100 | let _another = manager.attach(fixture.registry()); |
| 3101 | engine.sync().await.ok(); |
| 3102 | assert_eq!(manager.spawn_attempts(), 3); |
| 3103 | manager.retry(); |
| 3104 | engine.sync().await.unwrap(); |
| 3105 | assert_eq!(manager.spawn_attempts(), 4); |
| 3106 | assert!( |
| 3107 | manager |
| 3108 | .shared |
| 3109 | .plugin |
| 3110 | .supervision |
| 3111 | .lock() |
| 3112 | .unwrap() |
| 3113 | .crashes |
| 3114 | .is_empty() |
| 3115 | ); |
| 3116 | manager.shutdown().await; |
| 3117 | } |
| 3118 | |
| 3119 | #[tokio::test] |
| 3120 | async fn an_activation_crash_does_not_prevent_other_receipts_replaying() { |
| 3121 | let Some(node) = node_for_tests("activation crash") else { |
| 3122 | return; |
| 3123 | }; |
| 3124 | let _policy = TestPolicyGuard::extension_host(true); |
| 3125 | let fixture = FixturePlugins::new(&["crash-activation", "clash-script"]).await; |
| 3126 | let manager = supervised_manager(&fixture, node); |
| 3127 | let engine = manager.attach(fixture.registry()); |
| 3128 | engine.sync().await.unwrap(); |
| 3129 | wait_host(&manager, || { |
| 3130 | manager.spawn_attempts() == 2 |
| 3131 | && manager |
| 3132 | .live_tool_names() |
| 3133 | .contains(&"fixture_script_tool".into()) |
| 3134 | }) |
| 3135 | .await; |
| 3136 | assert!(matches!( |
| 3137 | manager.owner_state(&plugin_id(&fixture, "crash-activation")), |
| 3138 | Some(OwnerState::Failed(_)) |
| 3139 | )); |
| 3140 | assert_eq!( |
| 3141 | manager |
| 3142 | .shared |
| 3143 | .plugin |
| 3144 | .supervision |
| 3145 | .lock() |
| 3146 | .unwrap() |
| 3147 | .crashes |
| 3148 | .len(), |
| 3149 | 1 |
| 3150 | ); |
| 3151 | manager.shutdown().await; |
| 3152 | } |
| 3153 | |
| 3154 | #[tokio::test] |
| 3155 | async fn heartbeat_recovers_a_delayed_pong_then_kills_a_hung_host() { |
| 3156 | let Some(node) = node_for_tests("heartbeat") else { |
| 3157 | return; |
| 3158 | }; |
| 3159 | let _policy = TestPolicyGuard::extension_host(true); |
| 3160 | let fixture = FixturePlugins::new(&["hang-tool"]).await; |
| 3161 | let manager = supervised_manager(&fixture, node); |
| 3162 | let engine = manager.attach(fixture.registry()); |
| 3163 | engine.sync().await.unwrap(); |
| 3164 | let tool = host_tool(&engine, fixture.workspace(), "hang_probe"); |
| 3165 | let workspace = fixture.workspace().to_path_buf(); |
| 3166 | let plugins = engine.plugin_view(); |
| 3167 | let task = tokio::spawn(async move { |
| 3168 | tool.execute( |
| 3169 | json!({"ms": 400}), |
| 3170 | &ToolContext::new(&workspace).with_plugin_registry(plugins), |
| 3171 | ) |
| 3172 | .await |
| 3173 | }); |
| 3174 | wait_host(&manager, || { |
| 3175 | matches!(manager.status(), HostStatus::Unresponsive { .. }) |
| 3176 | }) |
| 3177 | .await; |
| 3178 | let result = task.await.unwrap(); |
| 3179 | assert!(result.is_ok(), "delayed live call failed: {result:?}"); |
| 3180 | wait_host(&manager, || { |
| 3181 | matches!(manager.status(), HostStatus::Ready { .. }) |
| 3182 | }) |
| 3183 | .await; |
| 3184 | assert_eq!(manager.spawn_attempts(), 1); |
| 3185 | let tool = host_tool(&engine, fixture.workspace(), "hang_probe"); |
| 3186 | let result = tokio::time::timeout( |
| 3187 | Duration::from_secs(5), |
| 3188 | tool.execute( |
| 3189 | json!({}), |
| 3190 | &ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()), |
| 3191 | ), |
| 3192 | ) |
| 3193 | .await |
| 3194 | .unwrap(); |
| 3195 | assert!(matches!(result, Err(ToolError::NotAvailable { .. }))); |
| 3196 | wait_host(&manager, || { |
| 3197 | manager.spawn_attempts() == 2 && manager.live_tool_names().contains(&"hang_probe".into()) |
| 3198 | }) |
| 3199 | .await; |
| 3200 | assert_eq!( |
| 3201 | manager |
| 3202 | .shared |
| 3203 | .plugin |
| 3204 | .supervision |
| 3205 | .lock() |
| 3206 | .unwrap() |
| 3207 | .crashes |
| 3208 | .len(), |
| 3209 | 1 |
| 3210 | ); |
| 3211 | manager.shutdown().await; |
| 3212 | } |
| 3213 | |
| 3214 | #[test] |
| 3215 | fn old_host_callbacks_cannot_fault_or_remove_a_new_owner() { |
| 3216 | use super::supervisor::HostEvents; |
| 3217 | let manager = ExtensionHostManager::new(ExtensionHostOptions::default()); |
| 3218 | manager |
| 3219 | .shared |
| 3220 | .plugin |
| 3221 | .host_generation |
| 3222 | .store(2, std::sync::atomic::Ordering::SeqCst); |
| 3223 | let owner = { |
| 3224 | let mut registry = manager.shared.registry.lock().unwrap(); |
| 3225 | let owner = registry |
| 3226 | .begin_owner( |
| 3227 | HostTier::Plugin, |
| 3228 | "fixture", |
| 3229 | "fixture", |
| 3230 | Some(fake_authority("fixture")), |
| 3231 | "hash-fixture", |
| 3232 | ) |
| 3233 | .unwrap(); |
| 3234 | registry.mark_active(&owner); |
| 3235 | register(&mut registry, &owner, "fixture_probe").unwrap(); |
| 3236 | owner |
| 3237 | }; |
| 3238 | let old = super::Events { |
| 3239 | shared: Arc::downgrade(&manager.shared), |
| 3240 | tier: HostTier::Plugin, |
| 3241 | generation: 1, |
| 3242 | }; |
| 3243 | old.faulted(&protocol::FaultedParams { |
| 3244 | owner, |
| 3245 | error: "stale fault".into(), |
| 3246 | }); |
| 3247 | old.exited(1, "stale exit".into(), String::new()); |
| 3248 | assert_eq!(manager.owner_state("fixture"), Some(OwnerState::Active)); |
| 3249 | assert_eq!(manager.live_tool_names(), ["fixture_probe"]); |
| 3250 | assert!( |
| 3251 | manager |
| 3252 | .shared |
| 3253 | .plugin |
| 3254 | .supervision |
| 3255 | .lock() |
| 3256 | .unwrap() |
| 3257 | .crashes |
| 3258 | .is_empty() |
| 3259 | ); |
| 3260 | } |
| 3261 | |
| 3262 | #[test] |
| 3263 | fn a_new_attachment_retries_only_cooled_down_launch_failures() { |
| 3264 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions::default())); |
| 3265 | *manager.shared.plugin.host.lock().unwrap() = super::HostSlot::Failed { |
| 3266 | reason: "missing Node".into(), |
| 3267 | stderr_tail: String::new(), |
| 3268 | }; |
| 3269 | { |
| 3270 | let mut state = manager.shared.plugin.supervision.lock().unwrap(); |
| 3271 | state.launch_failed = true; |
| 3272 | state.last_start = Some(Instant::now()); |
| 3273 | } |
| 3274 | let plugins = Arc::new(PluginRegistry::empty(Path::new("/fixture"))); |
| 3275 | let _first = manager.attach(Arc::clone(&plugins)); |
| 3276 | assert!(matches!(manager.status(), HostStatus::Failed { .. })); |
| 3277 | manager.shared.plugin.supervision.lock().unwrap().last_start = |
| 3278 | Some(Instant::now() - Duration::from_secs(61)); |
| 3279 | let _later = manager.attach(plugins); |
| 3280 | assert_eq!(manager.status(), HostStatus::Idle); |
| 3281 | } |
| 3282 | |
| 3283 | #[tokio::test] |
| 3284 | async fn replay_rechecks_persisted_disable_and_keeps_workspace_tools_separate() { |
| 3285 | let Some(node) = node_for_tests("replay authority") else { |
| 3286 | return; |
| 3287 | }; |
| 3288 | let _policy = TestPolicyGuard::extension_host(true); |
| 3289 | let a = FixturePlugins::new(&["crash-tool"]).await; |
| 3290 | let b = FixturePlugins::new(&["clash-script"]).await; |
| 3291 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 3292 | runtime: NODE, |
| 3293 | node_override: Some(node), |
| 3294 | bun_override: None, |
| 3295 | root: Some(a.root.clone()), |
| 3296 | supervision: super::SupervisionOptions { |
| 3297 | restart_backoff: Duration::from_secs(1), |
| 3298 | ..fast_supervision() |
| 3299 | }, |
| 3300 | })); |
| 3301 | let first = manager.attach(a.registry()); |
| 3302 | let second = manager.attach(b.registry()); |
| 3303 | first.sync().await.unwrap(); |
| 3304 | let tool = host_tool(&first, a.workspace(), "crash_probe"); |
| 3305 | assert!(matches!( |
| 3306 | tool.execute( |
| 3307 | json!({}), |
| 3308 | &ToolContext::new(a.workspace()).with_plugin_registry(first.plugin_view()) |
| 3309 | ) |
| 3310 | .await, |
| 3311 | Err(ToolError::NotAvailable { .. }) |
| 3312 | )); |
| 3313 | wait_host(&manager, || { |
| 3314 | matches!(manager.status(), HostStatus::Restarting { .. }) |
| 3315 | }) |
| 3316 | .await; |
| 3317 | // Keep the engine's snapshot stale deliberately. Replay must consult the |
| 3318 | // persisted state rather than restoring the previous owner's authority. |
| 3319 | a.disable("crash-tool"); |
| 3320 | wait_host(&manager, || { |
| 3321 | manager.spawn_attempts() == 2 |
| 3322 | && manager |
| 3323 | .live_tool_names() |
| 3324 | .contains(&"fixture_script_tool".into()) |
| 3325 | }) |
| 3326 | .await; |
| 3327 | assert!(installed(&first, a.workspace()).is_empty()); |
| 3328 | assert_eq!(installed(&second, b.workspace()), ["fixture_script_tool"]); |
| 3329 | assert!(manager.owner_state(&plugin_id(&a, "crash-tool")).is_none()); |
| 3330 | manager.shutdown().await; |
| 3331 | tokio::time::sleep(Duration::from_millis(100)).await; |
| 3332 | assert_eq!(manager.status(), HostStatus::Idle); |
| 3333 | assert_eq!( |
| 3334 | manager.spawn_attempts(), |
| 3335 | 2, |
| 3336 | "planned shutdown never restarts" |
| 3337 | ); |
| 3338 | } |
| 3339 | |
| 3340 | #[tokio::test] |
| 3341 | async fn explicit_retry_refreshes_same_byte_authority_without_inheriting_old_handles() { |
| 3342 | let Some(node) = node_for_tests("same byte retry") else { |
| 3343 | return; |
| 3344 | }; |
| 3345 | let _policy = TestPolicyGuard::extension_host(true); |
| 3346 | let fixture = FixturePlugins::new(&["clash-script"]).await; |
| 3347 | let manager = supervised_manager(&fixture, node); |
| 3348 | let engine = manager.attach(fixture.registry()); |
| 3349 | engine.sync().await.unwrap(); |
| 3350 | let old = host_tool(&engine, fixture.workspace(), "fixture_script_tool"); |
| 3351 | let old_owner = manager.shared.registry.lock().unwrap().live_tools()[0] |
| 3352 | .owner |
| 3353 | .clone(); |
| 3354 | let mut updated = discover_with_config(&fixture.config); |
| 3355 | updated.enable("clash-script").unwrap(); |
| 3356 | manager.refresh_workspace(&Arc::new(updated)); |
| 3357 | manager.retry(); |
| 3358 | engine.sync().await.unwrap(); |
| 3359 | let current_owner = manager.shared.registry.lock().unwrap().live_tools()[0] |
| 3360 | .owner |
| 3361 | .clone(); |
| 3362 | assert_ne!(old_owner, current_owner); |
| 3363 | assert!(matches!( |
| 3364 | old.execute( |
| 3365 | json!({}), |
| 3366 | &ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()) |
| 3367 | ) |
| 3368 | .await, |
| 3369 | Err(ToolError::NotAvailable { .. }) |
| 3370 | )); |
| 3371 | let current = host_tool(&engine, fixture.workspace(), "fixture_script_tool"); |
| 3372 | assert!( |
| 3373 | current |
| 3374 | .execute( |
| 3375 | json!({}), |
| 3376 | &ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()) |
| 3377 | ) |
| 3378 | .await |
| 3379 | .is_ok() |
| 3380 | ); |
| 3381 | assert_eq!( |
| 3382 | manager.spawn_attempts(), |
| 3383 | 1, |
| 3384 | "a healthy process need not restart" |
| 3385 | ); |
| 3386 | manager.shutdown().await; |
| 3387 | } |
| 3388 | |
| 3389 | // --------------------------------------------------------------------------- |
| 3390 | // Runtime selection (Bun / Node) and the memory cap |
| 3391 | // --------------------------------------------------------------------------- |
| 3392 | |
| 3393 | /// A Bun >= 1.4.0 for the Bun integration tests, or `None` (skip) unless |
| 3394 | /// `CODEWHALE_EXT_HOST_BUN_TESTS` requires one. |
| 3395 | fn bun_for_tests(test: &str) -> Option<PathBuf> { |
| 3396 | let resolution = crate::dependencies::resolve_extension_host_runtime( |
| 3397 | crate::config::ExtensionHostRuntime::Bun, |
| 3398 | None, |
| 3399 | None, |
| 3400 | ); |
| 3401 | match resolution.selected { |
| 3402 | Some(runtime) => Some(runtime.path), |
| 3403 | None if std::env::var_os("CODEWHALE_EXT_HOST_BUN_TESTS").is_some() => panic!( |
| 3404 | "{test}: CODEWHALE_EXT_HOST_BUN_TESTS is set but {}", |
| 3405 | resolution.failure() |
| 3406 | ), |
| 3407 | None => { |
| 3408 | eprintln!("skipping {test}: {}", resolution.failure()); |
| 3409 | None |
| 3410 | } |
| 3411 | } |
| 3412 | } |
| 3413 | |
| 3414 | fn bun_manager( |
| 3415 | fixture: &FixturePlugins, |
| 3416 | bun: PathBuf, |
| 3417 | supervision: super::SupervisionOptions, |
| 3418 | ) -> Arc<ExtensionHostManager> { |
| 3419 | Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 3420 | runtime: crate::config::ExtensionHostRuntime::Bun, |
| 3421 | node_override: None, |
| 3422 | bun_override: Some(bun), |
| 3423 | root: Some(fixture.root.clone()), |
| 3424 | supervision, |
| 3425 | })) |
| 3426 | } |
| 3427 | |
| 3428 | #[test] |
| 3429 | fn node_is_the_default_runtime_and_bun_is_an_opt_in() { |
| 3430 | use crate::config::{ExtensionHostConfig, ExtensionHostRuntime as Choice}; |
| 3431 | let parse = |text: &str| toml::from_str::<ExtensionHostConfig>(text).unwrap(); |
| 3432 | // Bun is not the default until it is qualified on every platform. |
| 3433 | assert_eq!(parse("").effective_runtime(), Choice::Node); |
| 3434 | assert_eq!( |
| 3435 | parse("node = \"/opt/node\"").effective_runtime(), |
| 3436 | Choice::Node |
| 3437 | ); |
| 3438 | // A table that names only a Bun asks for Bun. |
| 3439 | assert_eq!(parse("bun = \"/opt/bun\"").effective_runtime(), Choice::Bun); |
| 3440 | assert_eq!( |
| 3441 | parse("node = \"/n\"\nbun = \"/b\"").effective_runtime(), |
| 3442 | Choice::Node |
| 3443 | ); |
| 3444 | assert_eq!( |
| 3445 | parse("runtime = \"bun\"\nnode = \"/n\"").effective_runtime(), |
| 3446 | Choice::Bun |
| 3447 | ); |
| 3448 | assert_eq!( |
| 3449 | parse("runtime = \"auto\"").effective_runtime(), |
| 3450 | Choice::Auto |
| 3451 | ); |
| 3452 | assert!(toml::from_str::<ExtensionHostConfig>("runtime = \"deno\"").is_err()); |
| 3453 | let options = ExtensionHostOptions::from_config(Some(&parse("node = \"~/node\""))); |
| 3454 | assert_eq!(options.runtime, Choice::Node); |
| 3455 | assert!(!options.node_override.unwrap().starts_with("~")); |
| 3456 | assert_eq!( |
| 3457 | ExtensionHostOptions::from_config(None).runtime, |
| 3458 | Choice::Node |
| 3459 | ); |
| 3460 | assert_eq!(ExtensionHostOptions::default().runtime, Choice::Node); |
| 3461 | } |
| 3462 | |
| 3463 | #[test] |
| 3464 | fn launch_plan_gives_each_runtime_its_own_flags() { |
| 3465 | use crate::dependencies::{HostRuntime, HostRuntimeKind}; |
| 3466 | let home = tempfile::tempdir().unwrap(); |
| 3467 | let bundle = home.path().join("codewhale-extension-host.mjs"); |
| 3468 | for (kind, expected) in [ |
| 3469 | ( |
| 3470 | HostRuntimeKind::Bun, |
| 3471 | vec!["--no-install", "--no-env-file", "--config=", "--no-addons"], |
| 3472 | ), |
| 3473 | ( |
| 3474 | HostRuntimeKind::Node, |
| 3475 | [ |
| 3476 | &[ |
| 3477 | "--max-old-space-size=256", |
| 3478 | "--disable-proto=throw", |
| 3479 | "--no-addons", |
| 3480 | ][..], |
| 3481 | if cfg!(windows) { |
| 3482 | &["--preserve-symlinks", "--preserve-symlinks-main"][..] |
| 3483 | } else { |
| 3484 | &[][..] |
| 3485 | }, |
| 3486 | &["--no-experimental-sqlite", "--no-experimental-ffi"][..], |
| 3487 | ] |
| 3488 | .concat(), |
| 3489 | ), |
| 3490 | ] { |
| 3491 | let runtime = HostRuntime { |
| 3492 | kind, |
| 3493 | path: PathBuf::from("/opt/runtime/bin").join(kind.name()), |
| 3494 | version: (1, 4, 0), |
| 3495 | compiled: false, |
| 3496 | native_code_flags: match kind { |
| 3497 | HostRuntimeKind::Bun => Vec::new(), |
| 3498 | HostRuntimeKind::Node => crate::dependencies::NODE_NATIVE_CODE_FLAGS.to_vec(), |
| 3499 | }, |
| 3500 | }; |
| 3501 | let launch = super::supervisor::plan_launch( |
| 3502 | HostTier::Builtin, |
| 3503 | &runtime, |
| 3504 | &bundle, |
| 3505 | home.path(), |
| 3506 | 1 << 30, |
| 3507 | ) |
| 3508 | .unwrap(); |
| 3509 | // Wrapped or not, the runtime's flags come right before the bundle. |
| 3510 | let at = launch |
| 3511 | .args |
| 3512 | .iter() |
| 3513 | .position(|arg| Path::new(arg) == bundle) |
| 3514 | .expect("bundle in argv"); |
| 3515 | let flags = &launch.args[at - expected.len()..at]; |
| 3516 | for (flag, want) in flags.iter().zip(&expected) { |
| 3517 | assert!(flag.starts_with(want), "{kind:?}: {flags:?}"); |
| 3518 | } |
| 3519 | assert_eq!(launch.runtime, runtime); |
| 3520 | assert_eq!(launch.memory_cap, 1 << 30); |
| 3521 | assert_eq!( |
| 3522 | launch.memory, |
| 3523 | super::supervisor::MemoryEnforcement::planned(kind) |
| 3524 | ); |
| 3525 | let shadow_realm_off = launch |
| 3526 | .runtime_env |
| 3527 | .contains(&("BUN_JSC_useShadowRealm".to_string(), "0".to_string())); |
| 3528 | // An inherited NODE_OPTIONS preload must not run before the lockdown. |
| 3529 | assert!( |
| 3530 | launch |
| 3531 | .runtime_env |
| 3532 | .contains(&("NODE_OPTIONS".to_string(), String::new())), |
| 3533 | "{:?}", |
| 3534 | launch.runtime_env |
| 3535 | ); |
| 3536 | assert_eq!(shadow_realm_off, kind == HostRuntimeKind::Bun); |
| 3537 | if kind == HostRuntimeKind::Bun { |
| 3538 | assert!( |
| 3539 | !launch |
| 3540 | .args |
| 3541 | .iter() |
| 3542 | .any(|arg| arg.starts_with("--max-old-space")) |
| 3543 | ); |
| 3544 | } |
| 3545 | } |
| 3546 | // Where the kernel limit comes from on each platform. |
| 3547 | use super::supervisor::MemoryEnforcement; |
| 3548 | let (bun, node) = ( |
| 3549 | MemoryEnforcement::planned(HostRuntimeKind::Bun), |
| 3550 | MemoryEnforcement::planned(HostRuntimeKind::Node), |
| 3551 | ); |
| 3552 | if cfg!(target_os = "macos") { |
| 3553 | assert_eq!( |
| 3554 | (bun, node), |
| 3555 | (MemoryEnforcement::Jetsam, MemoryEnforcement::Heartbeat) |
| 3556 | ); |
| 3557 | } else if cfg!(target_os = "linux") { |
| 3558 | assert_eq!( |
| 3559 | (bun, node), |
| 3560 | (MemoryEnforcement::Rlimit, MemoryEnforcement::Rlimit) |
| 3561 | ); |
| 3562 | } else if cfg!(windows) { |
| 3563 | assert_eq!( |
| 3564 | (bun, node), |
| 3565 | (MemoryEnforcement::JobObject, MemoryEnforcement::JobObject) |
| 3566 | ); |
| 3567 | } |
| 3568 | } |
| 3569 | |
| 3570 | #[tokio::test] |
| 3571 | async fn handshake_refuses_a_runtime_or_version_mismatch_and_an_unapplied_kernel_cap() { |
| 3572 | let Some(node) = node_for_tests("handshake_refuses_a_runtime_or_version_mismatch") else { |
| 3573 | return; |
| 3574 | }; |
| 3575 | use crate::dependencies::HostRuntimeKind; |
| 3576 | let home = tempfile::tempdir().unwrap(); |
| 3577 | let bundle = super::materialize_bundle(home.path()).unwrap(); |
| 3578 | // Resolved, so the Node gets the native-code flags it accepts. |
| 3579 | let runtime = |
| 3580 | crate::dependencies::resolve_extension_host_runtime(NODE, Some(node.as_path()), None) |
| 3581 | .selected |
| 3582 | .expect("the test Node resolves"); |
| 3583 | struct NoEvents; |
| 3584 | impl super::supervisor::HostEvents for NoEvents { |
| 3585 | fn register(&self, _: &protocol::RegisterParams) -> protocol::RegisterResult { |
| 3586 | unreachable!() |
| 3587 | } |
| 3588 | fn unregister(&self, _: &protocol::UnregisterParams) {} |
| 3589 | fn faulted(&self, _: &protocol::FaultedParams) {} |
| 3590 | fn log(&self, _: &protocol::LogParams) {} |
| 3591 | fn exited(&self, _: u64, _: String, _: String) {} |
| 3592 | } |
| 3593 | for case in ["runtime", "version", "cap"] { |
| 3594 | let mut launch = super::supervisor::plan_launch( |
| 3595 | HostTier::Plugin, |
| 3596 | &runtime, |
| 3597 | &bundle, |
| 3598 | home.path(), |
| 3599 | 1 << 30, |
| 3600 | ) |
| 3601 | .unwrap(); |
| 3602 | let expected = match case { |
| 3603 | // The core believes it launched Bun; the host truthfully says Node. |
| 3604 | "runtime" => { |
| 3605 | launch.runtime.kind = HostRuntimeKind::Bun; |
| 3606 | "host reports runtime node but bun was launched" |
| 3607 | } |
| 3608 | // The pinned probe saw another version: the binary at that path |
| 3609 | // was replaced after it was probed. |
| 3610 | "version" => { |
| 3611 | launch.runtime.version = (0, 0, 1); |
| 3612 | "the runtime binary changed mid-session" |
| 3613 | } |
| 3614 | // An actual Node host reports no kernel cap. Even with matching |
| 3615 | // runtime/digest, the requested hard boundary must block admission. |
| 3616 | _ => { |
| 3617 | launch.memory = super::supervisor::MemoryEnforcement::Jetsam; |
| 3618 | "host did not apply the requested 1024 MiB kernel memory limit; initialization refused" |
| 3619 | } |
| 3620 | }; |
| 3621 | let error = match super::supervisor::HostProcess::spawn( |
| 3622 | 1, |
| 3623 | &launch, |
| 3624 | super::bundle_sha256(), |
| 3625 | Arc::new(NoEvents), |
| 3626 | ) |
| 3627 | .await |
| 3628 | { |
| 3629 | Ok(_) => panic!("{expected}"), |
| 3630 | Err(error) => error, |
| 3631 | }; |
| 3632 | assert!(error.contains(expected), "{error}"); |
| 3633 | } |
| 3634 | } |
| 3635 | |
| 3636 | /// A host that reports another tier than the one launched, or built-in module |
| 3637 | /// digests other than the ones the core pins, is refused at the handshake like |
| 3638 | /// a runtime mismatch: before initialization, with the reason. |
| 3639 | #[tokio::test] |
| 3640 | async fn handshake_refuses_a_tier_or_built_in_module_digest_mismatch() { |
| 3641 | let Some(node) = node_for_tests("handshake_refuses_a_tier_or_built_in_module_digest_mismatch") |
| 3642 | else { |
| 3643 | return; |
| 3644 | }; |
| 3645 | // A table that pins a module the host bundle does not embed. |
| 3646 | const PINNED: &[super::tier::BuiltinModule] = &[super::tier::BuiltinModule { |
| 3647 | id: "demo", |
| 3648 | source_sha256: "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef", |
| 3649 | tools: &[], |
| 3650 | }]; |
| 3651 | struct NoEvents; |
| 3652 | impl super::supervisor::HostEvents for NoEvents { |
| 3653 | fn register(&self, _: &protocol::RegisterParams) -> protocol::RegisterResult { |
| 3654 | unreachable!() |
| 3655 | } |
| 3656 | fn unregister(&self, _: &protocol::UnregisterParams) {} |
| 3657 | fn faulted(&self, _: &protocol::FaultedParams) {} |
| 3658 | fn log(&self, _: &protocol::LogParams) {} |
| 3659 | fn exited(&self, _: u64, _: String, _: String) {} |
| 3660 | } |
| 3661 | let home = tempfile::tempdir().unwrap(); |
| 3662 | let bundle = super::materialize_bundle(home.path()).unwrap(); |
| 3663 | let runtime = |
| 3664 | crate::dependencies::resolve_extension_host_runtime(NODE, Some(node.as_path()), None) |
| 3665 | .selected |
| 3666 | .expect("the test Node resolves"); |
| 3667 | let launch_for = |tier: HostTier| { |
| 3668 | super::supervisor::plan_launch(tier, &runtime, &bundle, home.path(), 1 << 30).unwrap() |
| 3669 | }; |
| 3670 | |
| 3671 | // Each tier reports its own tier and the exact embedded module digests. |
| 3672 | for tier in HostTier::ALL { |
| 3673 | let host = super::supervisor::HostProcess::spawn( |
| 3674 | 1, |
| 3675 | &launch_for(tier), |
| 3676 | super::bundle_sha256(), |
| 3677 | Arc::new(NoEvents), |
| 3678 | ) |
| 3679 | .await |
| 3680 | .unwrap_or_else(|error| panic!("{tier:?} host: {error}")); |
| 3681 | assert_eq!(host.tier, tier); |
| 3682 | host.shutdown().await; |
| 3683 | } |
| 3684 | |
| 3685 | let mut embedded: Vec<_> = super::tier::BUILTIN_MODULES |
| 3686 | .iter() |
| 3687 | .map(|module| format!("{}={}", module.id, &module.source_sha256[..12])) |
| 3688 | .collect(); |
| 3689 | embedded.sort(); |
| 3690 | let module_mismatch = format!( |
| 3691 | "host bundle embeds the built-in module digests [{}] but the core pins [demo=0123456789ab]", |
| 3692 | embedded.join(", ") |
| 3693 | ); |
| 3694 | for (case, expected) in [ |
| 3695 | ( |
| 3696 | "tier", |
| 3697 | "host reports the plugin tier but the builtin tier was launched", |
| 3698 | ), |
| 3699 | ("modules", module_mismatch.as_str()), |
| 3700 | ] { |
| 3701 | let mut launch = launch_for(HostTier::Plugin); |
| 3702 | match case { |
| 3703 | // The core believes it launched the builtin tier; the host, started |
| 3704 | // with `--tier=plugin`, truthfully says plugin. |
| 3705 | "tier" => launch.tier = HostTier::Builtin, |
| 3706 | _ => launch.builtin_modules = PINNED, |
| 3707 | } |
| 3708 | let error = match super::supervisor::HostProcess::spawn( |
| 3709 | 1, |
| 3710 | &launch, |
| 3711 | super::bundle_sha256(), |
| 3712 | Arc::new(NoEvents), |
| 3713 | ) |
| 3714 | .await |
| 3715 | { |
| 3716 | Ok(_) => panic!("{case}: {expected}"), |
| 3717 | Err(error) => error, |
| 3718 | }; |
| 3719 | assert!(error.contains(expected), "{case}: {error}"); |
| 3720 | } |
| 3721 | } |
| 3722 | |
| 3723 | #[tokio::test] |
| 3724 | async fn bun_host_runs_the_dsh_plugin_reports_bun_and_restarts_on_bun() { |
| 3725 | let Some(bun) = bun_for_tests("bun_host_runs_the_dsh_plugin") else { |
| 3726 | return; |
| 3727 | }; |
| 3728 | let _policy = TestPolicyGuard::extension_host(true); |
| 3729 | let fixture = FixturePlugins::new(&["dsh-workspace-deps"]).await; |
| 3730 | let manager = bun_manager(&fixture, bun.clone(), fast_supervision()); |
| 3731 | let engine = manager.attach(fixture.registry()); |
| 3732 | engine.sync().await.unwrap(); |
| 3733 | let HostStatus::Ready { |
| 3734 | runtime, |
| 3735 | runtime_version, |
| 3736 | .. |
| 3737 | } = manager.status() |
| 3738 | else { |
| 3739 | panic!("{:?} {:?}", manager.status(), manager.diagnostics()); |
| 3740 | }; |
| 3741 | assert_eq!(runtime, "bun"); |
| 3742 | let banner = std::process::Command::new(&bun) |
| 3743 | .arg("--version") |
| 3744 | .output() |
| 3745 | .unwrap(); |
| 3746 | assert_eq!( |
| 3747 | runtime_version, |
| 3748 | String::from_utf8_lossy(&banner.stdout).trim() |
| 3749 | ); |
| 3750 | let summary = manager.runtime_summary().unwrap(); |
| 3751 | assert!(summary.starts_with("bun "), "{summary}"); |
| 3752 | assert!(summary.contains("(runtime = \"bun\")"), "{summary}"); |
| 3753 | let report = super::render_status(&manager); |
| 3754 | assert!( |
| 3755 | report.contains(&format!("· bun {runtime_version} ·")), |
| 3756 | "{report}" |
| 3757 | ); |
| 3758 | assert!(report.contains("runtime: bun "), "{report}"); |
| 3759 | |
| 3760 | let tool = host_tool(&engine, fixture.workspace(), "load_workspace_dependencies"); |
| 3761 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()); |
| 3762 | let result = tool.execute(json!({}), &context).await.unwrap(); |
| 3763 | assert!(result.success, "{}", result.content); |
| 3764 | let payload: Value = serde_json::from_str(&result.content).unwrap(); |
| 3765 | assert_eq!(payload["pythonDistributions"]["numpy"], "2.1.0"); |
| 3766 | |
| 3767 | // A crash restarts on the pinned runtime; it is never re-resolved. |
| 3768 | let pid = manager.host_pid().unwrap(); |
| 3769 | #[cfg(unix)] |
| 3770 | assert!( |
| 3771 | std::process::Command::new("kill") |
| 3772 | .args(["-9", &pid.to_string()]) |
| 3773 | .status() |
| 3774 | .unwrap() |
| 3775 | .success() |
| 3776 | ); |
| 3777 | #[cfg(windows)] |
| 3778 | assert!( |
| 3779 | std::process::Command::new("taskkill") |
| 3780 | .args(["/F", "/PID", &pid.to_string()]) |
| 3781 | .status() |
| 3782 | .unwrap() |
| 3783 | .success() |
| 3784 | ); |
| 3785 | wait_host(&manager, || { |
| 3786 | manager.spawn_attempts() == 2 && matches!(manager.status(), HostStatus::Ready { .. }) |
| 3787 | }) |
| 3788 | .await; |
| 3789 | assert!(matches!( |
| 3790 | manager.status(), |
| 3791 | HostStatus::Ready { runtime: "bun", .. } |
| 3792 | )); |
| 3793 | assert_eq!(manager.runtime_summary().unwrap(), summary); |
| 3794 | #[cfg(target_os = "macos")] |
| 3795 | { |
| 3796 | // This kill came from the operator, not the memory hog. The observed |
| 3797 | // signal and configured cap must not invent a cause for it. |
| 3798 | let diagnostics = manager.diagnostics(); |
| 3799 | assert!( |
| 3800 | diagnostics |
| 3801 | .iter() |
| 3802 | .any(|line| line.contains("SIGKILL cause unavailable")), |
| 3803 | "{diagnostics:?}" |
| 3804 | ); |
| 3805 | assert!( |
| 3806 | !diagnostics |
| 3807 | .iter() |
| 3808 | .any(|line| line.contains("exceeded its memory cap")), |
| 3809 | "{diagnostics:?}" |
| 3810 | ); |
| 3811 | } |
| 3812 | manager.shutdown().await; |
| 3813 | } |
| 3814 | |
| 3815 | /// `runtime = "auto"`: a Bun that passes the version probe but cannot start |
| 3816 | /// the host is reported once, and Node runs the host for the rest of the |
| 3817 | /// session. `runtime = "bun"` never falls back, and a launch that never |
| 3818 | /// completed a handshake pins nothing. |
| 3819 | #[cfg(unix)] |
| 3820 | #[tokio::test] |
| 3821 | async fn auto_uses_node_for_the_session_when_the_bun_host_fails_to_start() { |
| 3822 | use std::os::unix::fs::PermissionsExt; |
| 3823 | let Some(node) = node_for_tests("auto_uses_node_when_the_bun_host_fails_to_start") else { |
| 3824 | return; |
| 3825 | }; |
| 3826 | let _policy = TestPolicyGuard::extension_host(true); |
| 3827 | let fixture = FixturePlugins::new(&["slow-tool"]).await; |
| 3828 | let bin = tempfile::tempdir().unwrap(); |
| 3829 | let bun = bin.path().join("bun"); |
| 3830 | std::fs::write( |
| 3831 | &bun, |
| 3832 | "#!/bin/sh\ncase \"$1\" in --version) echo 1.4.0;; *) echo 'simulated Bun start failure' >&2; exit 3;; esac\n", |
| 3833 | ) |
| 3834 | .unwrap(); |
| 3835 | std::fs::set_permissions(&bun, std::fs::Permissions::from_mode(0o755)).unwrap(); |
| 3836 | let manager_for = |runtime| { |
| 3837 | Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 3838 | runtime, |
| 3839 | node_override: Some(node.clone()), |
| 3840 | bun_override: Some(bun.clone()), |
| 3841 | root: Some(fixture.root.clone()), |
| 3842 | supervision: fast_supervision(), |
| 3843 | })) |
| 3844 | }; |
| 3845 | |
| 3846 | let manager = manager_for(crate::config::ExtensionHostRuntime::Auto); |
| 3847 | let engine = manager.attach(fixture.registry()); |
| 3848 | engine.sync().await.unwrap(); |
| 3849 | assert!( |
| 3850 | matches!( |
| 3851 | manager.status(), |
| 3852 | HostStatus::Ready { |
| 3853 | runtime: "node", |
| 3854 | .. |
| 3855 | } |
| 3856 | ), |
| 3857 | "{:?} {:?}", |
| 3858 | manager.status(), |
| 3859 | manager.diagnostics() |
| 3860 | ); |
| 3861 | let diagnostics = manager.diagnostics(); |
| 3862 | let reported: Vec<_> = diagnostics |
| 3863 | .iter() |
| 3864 | .filter(|line| line.contains("uses Node for the rest of this session")) |
| 3865 | .collect(); |
| 3866 | assert_eq!(reported.len(), 1, "{diagnostics:?}"); |
| 3867 | assert!( |
| 3868 | reported[0].contains(&bun.display().to_string()), |
| 3869 | "{diagnostics:?}" |
| 3870 | ); |
| 3871 | let summary = manager.runtime_summary().unwrap(); |
| 3872 | assert!(summary.starts_with("node "), "{summary}"); |
| 3873 | assert!(summary.contains("(runtime = \"auto\")"), "{summary}"); |
| 3874 | assert!( |
| 3875 | summary.contains("Bun failed to start this session"), |
| 3876 | "{summary}" |
| 3877 | ); |
| 3878 | assert_eq!(manager.spawn_attempts(), 2); |
| 3879 | manager.shutdown().await; |
| 3880 | |
| 3881 | let explicit = manager_for(crate::config::ExtensionHostRuntime::Bun); |
| 3882 | let engine = explicit.attach(fixture.registry()); |
| 3883 | let _ = engine.sync().await; |
| 3884 | assert!( |
| 3885 | matches!(explicit.status(), HostStatus::Failed { .. }), |
| 3886 | "{:?} {:?}", |
| 3887 | explicit.status(), |
| 3888 | explicit.diagnostics() |
| 3889 | ); |
| 3890 | assert_eq!(explicit.spawn_attempts(), 1); |
| 3891 | assert_eq!(explicit.runtime_summary(), None); |
| 3892 | } |
| 3893 | |
| 3894 | /// The cap holds for memory outside the JS heap too (Buffers), which Node's |
| 3895 | /// `--max-old-space-size` never bounded. Linux and Windows: the kernel fails |
| 3896 | /// the allocation at 1 GiB. macOS: Bun's jetsam limit gets the host killed by |
| 3897 | /// the kernel; a Node host is killed by the heartbeat check. |
| 3898 | async fn memory_hog_is_stopped( |
| 3899 | manager: Arc<ExtensionHostManager>, |
| 3900 | fixture: &FixturePlugins, |
| 3901 | kind: crate::dependencies::HostRuntimeKind, |
| 3902 | ) { |
| 3903 | use super::supervisor::MemoryEnforcement; |
| 3904 | let engine = manager.attach(fixture.registry()); |
| 3905 | engine.sync().await.unwrap(); |
| 3906 | // The host confirmed the enforcement this platform plans for it (for |
| 3907 | // macOS + Bun: the jetsam limit it applied to itself, via `host/hello`). |
| 3908 | let HostStatus::Ready { memory, .. } = manager.status() else { |
| 3909 | panic!("{:?} {:?}", manager.status(), manager.diagnostics()); |
| 3910 | }; |
| 3911 | assert_eq!(memory, MemoryEnforcement::planned(kind)); |
| 3912 | let tool = host_tool(&engine, fixture.workspace(), "memory_hog"); |
| 3913 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()); |
| 3914 | let outcome = tokio::time::timeout( |
| 3915 | Duration::from_secs(60), |
| 3916 | tool.execute(json!({"mib": 2048}), &context), |
| 3917 | ) |
| 3918 | .await |
| 3919 | .expect("the hog is stopped well before 60 s"); |
| 3920 | match &outcome { |
| 3921 | Ok(result) => assert!( |
| 3922 | !result.success, |
| 3923 | "allocating 2 GiB must not succeed under the cap: {}", |
| 3924 | result.content |
| 3925 | ), |
| 3926 | Err(ToolError::NotAvailable { message }) => { |
| 3927 | assert!(message.contains("extension host exited"), "{message}"); |
| 3928 | } |
| 3929 | // Linux (RLIMIT_DATA) and Windows (Job Object): the kernel refuses the |
| 3930 | // allocation, the runtime throws inside the tool, and the host lives on. |
| 3931 | Err(ToolError::ExecutionFailed { message, .. }) |
| 3932 | if matches!( |
| 3933 | memory, |
| 3934 | MemoryEnforcement::Rlimit | MemoryEnforcement::JobObject |
| 3935 | ) => |
| 3936 | { |
| 3937 | assert!( |
| 3938 | message.contains("allocation failed") |
| 3939 | || (kind == crate::dependencies::HostRuntimeKind::Bun |
| 3940 | && message.ends_with("RangeError: Out of memory")), |
| 3941 | "{message}" |
| 3942 | ); |
| 3943 | } |
| 3944 | Err(other) => panic!("unexpected error: {other:?}"), |
| 3945 | } |
| 3946 | let stopped_by = match memory { |
| 3947 | MemoryEnforcement::Jetsam => { |
| 3948 | Some("configured kernel memory limit: 400 MiB; SIGKILL cause unavailable") |
| 3949 | } |
| 3950 | MemoryEnforcement::Heartbeat => Some("exceeded its memory cap"), |
| 3951 | _ => None, |
| 3952 | }; |
| 3953 | if let Some(stopped_by) = stopped_by { |
| 3954 | // Pending calls settle before the existing exit callback publishes |
| 3955 | // its diagnostic. Observe that callback rather than race its delivery. |
| 3956 | wait_host(&manager, || { |
| 3957 | manager |
| 3958 | .diagnostics() |
| 3959 | .iter() |
| 3960 | .any(|line| line.contains(stopped_by)) |
| 3961 | }) |
| 3962 | .await; |
| 3963 | } |
| 3964 | manager.shutdown().await; |
| 3965 | } |
| 3966 | |
| 3967 | fn memory_cap_supervision() -> super::SupervisionOptions { |
| 3968 | super::SupervisionOptions { |
| 3969 | // Linux cannot start either runtime under much less than 1 GiB of |
| 3970 | // RLIMIT_DATA (see `HOST_MEMORY_CAP`); on macOS 400 MiB keeps the |
| 3971 | // test fast for both the jetsam limit and the heartbeat check. |
| 3972 | memory_cap: if cfg!(target_os = "macos") { |
| 3973 | 400 << 20 |
| 3974 | } else { |
| 3975 | super::supervisor::HOST_MEMORY_CAP |
| 3976 | }, |
| 3977 | // Production hang detection. With the fast 600 ms hang timeout, a host |
| 3978 | // stalled in GC near its cap was killed for a missed heartbeat and |
| 3979 | // restarted before the memory limit itself stopped it, so the test |
| 3980 | // never observed the enforcement it exists to prove. |
| 3981 | ping_timeout: Duration::from_secs(3), |
| 3982 | hang_timeout: super::supervisor::PING_DEADLINE, |
| 3983 | ..fast_supervision() |
| 3984 | } |
| 3985 | } |
| 3986 | |
| 3987 | #[cfg(any(target_os = "linux", target_os = "macos", windows))] |
| 3988 | #[tokio::test] |
| 3989 | async fn memory_cap_stops_a_node_host() { |
| 3990 | let Some(node) = node_for_tests("memory_cap_stops_a_node_host") else { |
| 3991 | return; |
| 3992 | }; |
| 3993 | let _policy = TestPolicyGuard::extension_host(true); |
| 3994 | let fixture = FixturePlugins::new(&["memory-hog"]).await; |
| 3995 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 3996 | runtime: NODE, |
| 3997 | node_override: Some(node), |
| 3998 | bun_override: None, |
| 3999 | root: Some(fixture.root.clone()), |
| 4000 | supervision: memory_cap_supervision(), |
| 4001 | })); |
| 4002 | memory_hog_is_stopped( |
| 4003 | manager, |
| 4004 | &fixture, |
| 4005 | crate::dependencies::HostRuntimeKind::Node, |
| 4006 | ) |
| 4007 | .await; |
| 4008 | } |
| 4009 | |
| 4010 | #[cfg(any(target_os = "linux", target_os = "macos", windows))] |
| 4011 | #[tokio::test] |
| 4012 | async fn memory_cap_stops_a_bun_host() { |
| 4013 | let Some(bun) = bun_for_tests("memory_cap_stops_a_bun_host") else { |
| 4014 | return; |
| 4015 | }; |
| 4016 | let _policy = TestPolicyGuard::extension_host(true); |
| 4017 | let fixture = FixturePlugins::new(&["memory-hog"]).await; |
| 4018 | let manager = bun_manager(&fixture, bun, memory_cap_supervision()); |
| 4019 | memory_hog_is_stopped(manager, &fixture, crate::dependencies::HostRuntimeKind::Bun).await; |
| 4020 | } |
| 4021 | |
| 4022 | /// Actual Rust containment/memory authority, using the same exact compiled |
| 4023 | /// image as the containment receipt rather than transferring system-Bun proof. |
| 4024 | #[cfg(any(target_os = "linux", target_os = "macos", windows))] |
| 4025 | #[tokio::test] |
| 4026 | async fn compiled_native_host_memory_cap_is_enforced() { |
| 4027 | let Some(binary) = compiled_image_for_tests() else { |
| 4028 | return; |
| 4029 | }; |
| 4030 | let _policy = TestPolicyGuard::extension_host(true); |
| 4031 | let fixture = FixturePlugins::new(&["memory-hog"]).await; |
| 4032 | let manager = bun_manager(&fixture, binary, memory_cap_supervision()); |
| 4033 | memory_hog_is_stopped(manager, &fixture, crate::dependencies::HostRuntimeKind::Bun).await; |
| 4034 | eprintln!( |
| 4035 | "compiled-native-memory=passed platform={} arch={}", |
| 4036 | std::env::consts::OS, |
| 4037 | std::env::consts::ARCH |
| 4038 | ); |
| 4039 | } |
| 4040 | |
| 4041 | // --------------------------------------------------------------------------- |
| 4042 | // Extension commands |
| 4043 | // --------------------------------------------------------------------------- |
| 4044 | |
| 4045 | fn register_command( |
| 4046 | registry: &mut OwnerRegistry, |
| 4047 | owner: &OwnerRef, |
| 4048 | name: &str, |
| 4049 | hint: Option<&str>, |
| 4050 | ) -> Result<u64, String> { |
| 4051 | registry.register(&RegisterParams { |
| 4052 | scope: None, |
| 4053 | owner: owner.clone(), |
| 4054 | kind: RegisterKind::Command, |
| 4055 | spec: RegisterSpecWire { |
| 4056 | name: name.to_string(), |
| 4057 | description: "d".to_string(), |
| 4058 | input_schema: None, |
| 4059 | argument_hint: hint.map(str::to_string), |
| 4060 | }, |
| 4061 | }) |
| 4062 | } |
| 4063 | |
| 4064 | #[test] |
| 4065 | fn command_registry_refuses_shadowing_and_undoes_exactly_one_entry() { |
| 4066 | let _catalog = stub_builtin_commands(); |
| 4067 | let mut registry = OwnerRegistry::new(); |
| 4068 | let a = registry |
| 4069 | .begin_owner( |
| 4070 | HostTier::Plugin, |
| 4071 | "a", |
| 4072 | "a", |
| 4073 | Some(fake_authority("a")), |
| 4074 | "hash-a", |
| 4075 | ) |
| 4076 | .unwrap(); |
| 4077 | let b = registry |
| 4078 | .begin_owner( |
| 4079 | HostTier::Plugin, |
| 4080 | "b", |
| 4081 | "b", |
| 4082 | Some(fake_authority("b")), |
| 4083 | "hash-b", |
| 4084 | ) |
| 4085 | .unwrap(); |
| 4086 | |
| 4087 | // Built-in names, their aliases, and the fixed mode aliases, as the |
| 4088 | // catalog says (the real table: `commands::extension_host_tests`). |
| 4089 | for name in ["help", "trust", "model", "jihua", "zidong", "stub-alias"] { |
| 4090 | let refused = register_command(&mut registry, &a, name, None) |
| 4091 | .expect_err(&format!("/{name} must be refused")); |
| 4092 | assert!(refused.contains("built-in command"), "{name}: {refused}"); |
| 4093 | } |
| 4094 | // Names outside DSH's grammar (upper case, a leading digit or slash, |
| 4095 | // spaces), and over-long ones. |
| 4096 | for name in [ |
| 4097 | "Hello", |
| 4098 | "9lives", |
| 4099 | "/slash", |
| 4100 | "two words", |
| 4101 | "", |
| 4102 | &"x".repeat(65), |
| 4103 | ] { |
| 4104 | let refused = register_command(&mut registry, &a, name, None) |
| 4105 | .expect_err(&format!("{name:?} must be refused")); |
| 4106 | assert!(refused.contains("invalid"), "{name:?}: {refused}"); |
| 4107 | } |
| 4108 | // A command has no input schema; descriptions and hints are bounded, |
| 4109 | // non-empty, single-line text. |
| 4110 | let mut params = RegisterParams { |
| 4111 | scope: None, |
| 4112 | owner: a.clone(), |
| 4113 | kind: RegisterKind::Command, |
| 4114 | spec: RegisterSpecWire { |
| 4115 | name: "ok-name".to_string(), |
| 4116 | description: "fine".to_string(), |
| 4117 | input_schema: json!({"type": "object"}).as_object().cloned(), |
| 4118 | argument_hint: None, |
| 4119 | }, |
| 4120 | }; |
| 4121 | assert!( |
| 4122 | registry |
| 4123 | .register(¶ms) |
| 4124 | .unwrap_err() |
| 4125 | .contains("input schema") |
| 4126 | ); |
| 4127 | params.spec.input_schema = None; |
| 4128 | for (description, hint, expect) in [ |
| 4129 | (" ", None, "needs a description"), |
| 4130 | ("two\nlines", None, "control characters"), |
| 4131 | ("fine", Some(""), "must not be empty"), |
| 4132 | ("fine", Some("\u{1b}[31m<x>"), "control characters"), |
| 4133 | ] { |
| 4134 | params.spec.description = description.to_string(); |
| 4135 | params.spec.argument_hint = hint.map(str::to_string); |
| 4136 | let refused = registry.register(¶ms).unwrap_err(); |
| 4137 | assert!( |
| 4138 | refused.contains(expect), |
| 4139 | "{description:?}/{hint:?}: {refused}" |
| 4140 | ); |
| 4141 | } |
| 4142 | params.spec.description = "x".repeat(super::registry::MAX_COMMAND_DESCRIPTION_BYTES + 1); |
| 4143 | params.spec.argument_hint = None; |
| 4144 | assert!( |
| 4145 | registry |
| 4146 | .register(¶ms) |
| 4147 | .unwrap_err() |
| 4148 | .contains("description") |
| 4149 | ); |
| 4150 | params.spec.description = "fine".to_string(); |
| 4151 | params.spec.argument_hint = Some("h".repeat(super::registry::MAX_COMMAND_HINT_BYTES + 1)); |
| 4152 | assert!( |
| 4153 | registry |
| 4154 | .register(¶ms) |
| 4155 | .unwrap_err() |
| 4156 | .contains("argument hint") |
| 4157 | ); |
| 4158 | |
| 4159 | let first = register_command(&mut registry, &a, "shared-name", Some("<x>")).unwrap(); |
| 4160 | // Another plugin cannot take it. |
| 4161 | let refused = register_command(&mut registry, &b, "shared-name", None).unwrap_err(); |
| 4162 | assert!( |
| 4163 | refused.contains("already registered by extension"), |
| 4164 | "{refused}" |
| 4165 | ); |
| 4166 | // Commands and tools are separate namespaces: the model calls one, the |
| 4167 | // user the other. |
| 4168 | register(&mut registry, &b, "shared_name_tool").unwrap(); |
| 4169 | register(&mut registry, &a, "shared-name").unwrap(); |
| 4170 | // The same owner re-registering retires the old handle. |
| 4171 | let second = register_command(&mut registry, &a, "shared-name", None).unwrap(); |
| 4172 | assert_ne!(first, second); |
| 4173 | registry.mark_active(&a); |
| 4174 | registry.unregister(&a, first); // stale: must not remove the newer entry |
| 4175 | let live = registry.live_commands(); |
| 4176 | assert_eq!(live.len(), 1); |
| 4177 | assert_eq!(live[0].handle, second); |
| 4178 | assert!(registry.live_command(second, "a", a.generation).is_some()); |
| 4179 | // A foreign owner cannot unregister it, and a stale generation finds nothing. |
| 4180 | registry.unregister(&b, second); |
| 4181 | assert!(registry.live_command(second, "a", a.generation).is_some()); |
| 4182 | assert!( |
| 4183 | registry |
| 4184 | .live_command(second, "a", a.generation + 1) |
| 4185 | .is_none() |
| 4186 | ); |
| 4187 | assert!(registry.live_command(second, "b", a.generation).is_none()); |
| 4188 | // Unregistering a command leaves the same owner's tool alone. |
| 4189 | registry.unregister(&a, second); |
| 4190 | assert!(registry.live_commands().is_empty()); |
| 4191 | let tools = registry.live_tools(); |
| 4192 | assert_eq!(tools.len(), 1, "b is not active yet; a's tool stays"); |
| 4193 | assert_eq!(tools[0].name, "shared-name"); |
| 4194 | |
| 4195 | // Caps. |
| 4196 | for index in 0..super::registry::MAX_COMMANDS_PER_OWNER { |
| 4197 | register_command(&mut registry, &a, &format!("c{index}"), None).unwrap(); |
| 4198 | } |
| 4199 | assert!( |
| 4200 | register_command(&mut registry, &a, "one-too-many", None) |
| 4201 | .unwrap_err() |
| 4202 | .contains("at most") |
| 4203 | ); |
| 4204 | |
| 4205 | // Revocation is synchronous and total, and a revoked owner cannot register. |
| 4206 | registry.mark_active(&a); |
| 4207 | assert_eq!( |
| 4208 | registry.live_commands().len(), |
| 4209 | super::registry::MAX_COMMANDS_PER_OWNER |
| 4210 | ); |
| 4211 | assert_eq!(registry.revoke_owner("a"), Some(a.clone())); |
| 4212 | assert!(registry.live_commands().is_empty()); |
| 4213 | assert!(register_command(&mut registry, &a, "after-revoke", None).is_err()); |
| 4214 | |
| 4215 | // A host crash drops every command, whoever owned it. |
| 4216 | registry.mark_active(&b); |
| 4217 | let handle = register_command(&mut registry, &b, "survivor", None).unwrap(); |
| 4218 | assert!(registry.live_command(handle, "b", b.generation).is_some()); |
| 4219 | registry.host_exited(HostTier::Plugin, "exited"); |
| 4220 | assert!(registry.live_commands().is_empty()); |
| 4221 | assert!(registry.live_command(handle, "b", b.generation).is_none()); |
| 4222 | } |
| 4223 | |
| 4224 | /// The host protocol gains `command/run` without gaining any core |
| 4225 | /// authority: the lint that guards `METHODS` runs in `protocol::tests`; this |
| 4226 | /// pins what the new method's request and its answer look like. |
| 4227 | #[test] |
| 4228 | fn command_run_and_its_answers_have_the_documented_shapes() { |
| 4229 | let request = protocol::CoreRequest::CommandRun(protocol::CommandRunParams { |
| 4230 | handle: 7, |
| 4231 | command_id: "c".to_string(), |
| 4232 | raw_input: "args".to_string(), |
| 4233 | deadline_ms: 30_000, |
| 4234 | workspace: None, |
| 4235 | session_id: None, |
| 4236 | agent_id: None, |
| 4237 | origin_turn_id: None, |
| 4238 | }); |
| 4239 | assert_eq!(request.method(), "command/run"); |
| 4240 | assert_eq!(request.deadline(), Duration::from_secs(30)); |
| 4241 | for (value, expect) in [ |
| 4242 | ( |
| 4243 | json!({"kind": "success", "text": "t"}), |
| 4244 | protocol::CommandResultWire::Success { |
| 4245 | text: Some("t".into()), |
| 4246 | }, |
| 4247 | ), |
| 4248 | ( |
| 4249 | json!({"kind": "success"}), |
| 4250 | protocol::CommandResultWire::Success { text: None }, |
| 4251 | ), |
| 4252 | ( |
| 4253 | json!({"kind": "error", "text": "no"}), |
| 4254 | protocol::CommandResultWire::Error { text: "no".into() }, |
| 4255 | ), |
| 4256 | ( |
| 4257 | json!({"kind": "submit", "prompt": "p"}), |
| 4258 | protocol::CommandResultWire::Submit { |
| 4259 | prompt: "p".into(), |
| 4260 | text: None, |
| 4261 | }, |
| 4262 | ), |
| 4263 | ] { |
| 4264 | let parsed: protocol::CommandResultWire = serde_json::from_value(value.clone()).unwrap(); |
| 4265 | assert_eq!(parsed, expect); |
| 4266 | assert_eq!(serde_json::to_value(&parsed).unwrap(), value); |
| 4267 | } |
| 4268 | assert!( |
| 4269 | serde_json::from_value::<protocol::CommandResultWire>(json!({"kind": "approve"})).is_err() |
| 4270 | ); |
| 4271 | } |
| 4272 | |
| 4273 | /// What a command shows is plugin-controlled text: escape sequences are |
| 4274 | /// stripped, and a prompt too large to submit whole is refused, not cut. |
| 4275 | #[tokio::test] |
| 4276 | async fn command_output_is_stripped_bounded_and_oversized_prompts_are_refused() { |
| 4277 | use super::command::{self, CommandOutcome}; |
| 4278 | // The wire-to-outcome mapping needs a host to answer, so run it against |
| 4279 | // a stub that returns canned results for `command/run`. |
| 4280 | let Some(node) = node_for_tests("command_output") else { |
| 4281 | return; |
| 4282 | }; |
| 4283 | let dir = tempfile::tempdir().unwrap(); |
| 4284 | let plugin = dir.path().join("canned"); |
| 4285 | std::fs::create_dir_all(&plugin).unwrap(); |
| 4286 | std::fs::write( |
| 4287 | plugin.join("plugin.json"), |
| 4288 | r#"{"$schema":"https://agent-plugins.org/schemas/plugin.json","name":"canned","version":"0.1.0","description":"canned results","license":"MIT","extensions":{"net.codewhale":{"native":{"path":"index.mjs"}}}}"#, |
| 4289 | ) |
| 4290 | .unwrap(); |
| 4291 | std::fs::write( |
| 4292 | plugin.join("index.mjs"), |
| 4293 | format!( |
| 4294 | r#"export const name = 'canned' |
| 4295 | export const inject = ['commands'] |
| 4296 | export function apply(ctx) {{ |
| 4297 | ctx.commands.register({{ name: 'big-text', description: 'd', handler: () => 'x'.repeat({}) }}) |
| 4298 | ctx.commands.register({{ name: 'big-prompt', description: 'd', handler: () => ({{ kind: 'submit', prompt: 'p'.repeat({}) }}) }}) |
| 4299 | ctx.commands.register({{ name: 'blank-prompt', description: 'd', handler: () => ({{ kind: 'submit', prompt: ' \u001b[0m ' }}) }}) |
| 4300 | ctx.commands.register({{ name: 'bad-result', description: 'd', handler: () => ({{ kind: 'approve' }}) }}) |
| 4301 | }} |
| 4302 | "#, |
| 4303 | command::MAX_TEXT_BYTES * 2, |
| 4304 | command::MAX_PROMPT_BYTES + 1 |
| 4305 | ), |
| 4306 | ) |
| 4307 | .unwrap(); |
| 4308 | let _policy = TestPolicyGuard::extension_host(true); |
| 4309 | let fixture = FixturePlugins::new(&[]).await; |
| 4310 | crate::plugins::install::install( |
| 4311 | crate::plugins::install::PluginInstallSource::LocalPath(plugin), |
| 4312 | &fixture.config.user_plugins_dir, |
| 4313 | crate::plugins::install::DEFAULT_MAX_SIZE_BYTES, |
| 4314 | &crate::network_policy::NetworkPolicy::default(), |
| 4315 | false, |
| 4316 | &|_| None, |
| 4317 | ) |
| 4318 | .await |
| 4319 | .unwrap(); |
| 4320 | let mut plugins = discover_with_config(&fixture.config); |
| 4321 | plugins.trust("canned").unwrap(); |
| 4322 | plugins.enable("canned").unwrap(); |
| 4323 | let manager = fixture.manager(node); |
| 4324 | let engine = manager.attach(Arc::new(plugins)); |
| 4325 | engine.sync().await.unwrap(); |
| 4326 | let run = |name: &'static str| { |
| 4327 | let manager = Arc::clone(&manager); |
| 4328 | let plugins = engine.plugin_view(); |
| 4329 | async move { |
| 4330 | let entry = manager |
| 4331 | .commands_for_plugins(plugins.as_ref()) |
| 4332 | .into_iter() |
| 4333 | .find(|entry| entry.registration.name == name) |
| 4334 | .unwrap_or_else(|| panic!("{name} is not live")); |
| 4335 | command::run(&manager.shared, &entry.reference(), "", None).await |
| 4336 | } |
| 4337 | }; |
| 4338 | match run("big-text").await.unwrap() { |
| 4339 | CommandOutcome::Show { text } => { |
| 4340 | assert!( |
| 4341 | text.ends_with("(output truncated)"), |
| 4342 | "{}", |
| 4343 | &text[text.len() - 40..] |
| 4344 | ); |
| 4345 | assert!(text.len() < command::MAX_TEXT_BYTES + 64); |
| 4346 | } |
| 4347 | other => panic!("{other:?}"), |
| 4348 | } |
| 4349 | assert!( |
| 4350 | run("big-prompt") |
| 4351 | .await |
| 4352 | .unwrap_err() |
| 4353 | .contains("not submitted") |
| 4354 | ); |
| 4355 | assert!( |
| 4356 | run("blank-prompt") |
| 4357 | .await |
| 4358 | .unwrap_err() |
| 4359 | .contains("empty prompt") |
| 4360 | ); |
| 4361 | assert!( |
| 4362 | run("bad-result") |
| 4363 | .await |
| 4364 | .unwrap_err() |
| 4365 | .contains("unknown result kind") |
| 4366 | ); |
| 4367 | manager.shutdown().await; |
| 4368 | } |
| 4369 | |
| 4370 | /// With no built-in command catalog installed, a command cannot be checked |
| 4371 | /// against the built-in names, so its registration is refused (never accepted |
| 4372 | /// unchecked); tools are not affected. With one, the catalog decides. |
| 4373 | #[test] |
| 4374 | fn a_command_registration_is_refused_when_no_built_in_catalog_is_installed() { |
| 4375 | let owner = |registry: &mut OwnerRegistry| { |
| 4376 | registry |
| 4377 | .begin_owner( |
| 4378 | HostTier::Plugin, |
| 4379 | "a", |
| 4380 | "a", |
| 4381 | Some(fake_authority("a")), |
| 4382 | "hash-a", |
| 4383 | ) |
| 4384 | .unwrap() |
| 4385 | }; |
| 4386 | let _absent = BuiltinCommandsGuard::absent(); |
| 4387 | let mut registry = OwnerRegistry::new(); |
| 4388 | let a = owner(&mut registry); |
| 4389 | let refused = register_command(&mut registry, &a, "fine-name", None).unwrap_err(); |
| 4390 | assert!( |
| 4391 | refused.contains("no built-in command catalog is installed"), |
| 4392 | "{refused}" |
| 4393 | ); |
| 4394 | assert!(registry.live_commands().is_empty()); |
| 4395 | // A name the grammar refuses is refused for its own reason first. |
| 4396 | assert!( |
| 4397 | register_command(&mut registry, &a, "Not Valid", None) |
| 4398 | .unwrap_err() |
| 4399 | .contains("invalid") |
| 4400 | ); |
| 4401 | // Tools do not consult the command table. |
| 4402 | register(&mut registry, &a, "a_tool").unwrap(); |
| 4403 | |
| 4404 | let _stub = stub_builtin_commands(); |
| 4405 | register_command(&mut registry, &a, "fine-name", None).unwrap(); |
| 4406 | let clash = register_command(&mut registry, &a, "help", None).unwrap_err(); |
| 4407 | assert!( |
| 4408 | clash.contains("collides with a built-in command"), |
| 4409 | "{clash}" |
| 4410 | ); |
| 4411 | } |
| 4412 | |
| 4413 | /// A command from a host that is down says so immediately instead of |
| 4414 | /// hanging, and a revoked registration cannot be run. |
| 4415 | #[tokio::test] |
| 4416 | async fn a_command_from_a_dead_host_reports_host_down() { |
| 4417 | let _catalog = stub_builtin_commands(); |
| 4418 | let policy = TestPolicyGuard::extension_host(true); |
| 4419 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions::default())); |
| 4420 | let reference = { |
| 4421 | let mut registry = manager.shared.registry.lock().unwrap(); |
| 4422 | let owner = registry |
| 4423 | .begin_owner( |
| 4424 | HostTier::Plugin, |
| 4425 | "probe", |
| 4426 | "probe", |
| 4427 | Some(fake_authority("probe")), |
| 4428 | "hash", |
| 4429 | ) |
| 4430 | .unwrap(); |
| 4431 | registry.mark_active(&owner); |
| 4432 | let handle = register_command(&mut registry, &owner, "probe-cmd", None).unwrap(); |
| 4433 | super::command::ExtensionCommandRef { |
| 4434 | selection: None, |
| 4435 | scope: None, |
| 4436 | content_hash: String::new(), |
| 4437 | handle, |
| 4438 | plugin_id: "probe".into(), |
| 4439 | generation: owner.generation, |
| 4440 | origin: "extension:probe".into(), |
| 4441 | workspace: PathBuf::from("/w"), |
| 4442 | } |
| 4443 | }; |
| 4444 | let refused = "start failed: no runtime"; |
| 4445 | for (slot, why) in [ |
| 4446 | ( |
| 4447 | super::HostSlot::Restarting { |
| 4448 | reason: "exited with signal: 9 (SIGKILL)".into(), |
| 4449 | }, |
| 4450 | "restarting after: exited with signal: 9 (SIGKILL)".to_string(), |
| 4451 | ), |
| 4452 | ( |
| 4453 | super::HostSlot::Failed { |
| 4454 | reason: refused.into(), |
| 4455 | stderr_tail: String::new(), |
| 4456 | }, |
| 4457 | format!("{refused} (change or reload a plugin to retry)"), |
| 4458 | ), |
| 4459 | (super::HostSlot::Idle, "not started".to_string()), |
| 4460 | ] { |
| 4461 | *manager.shared.plugin.host.lock().unwrap() = slot; |
| 4462 | let started = Instant::now(); |
| 4463 | let error = super::command::run(&manager.shared, &reference, "", None) |
| 4464 | .await |
| 4465 | .unwrap_err(); |
| 4466 | assert!( |
| 4467 | error.starts_with(&format!("extension host is down: {why}")), |
| 4468 | "{error}" |
| 4469 | ); |
| 4470 | assert!(started.elapsed() < Duration::from_secs(1)); |
| 4471 | } |
| 4472 | drop(policy); |
| 4473 | let _off = TestPolicyGuard::extension_host(false); |
| 4474 | let error = super::command::run(&manager.shared, &reference, "", None) |
| 4475 | .await |
| 4476 | .unwrap_err(); |
| 4477 | assert_eq!( |
| 4478 | error, |
| 4479 | "extension host is down: disabled by config ([features] extension_host is off)" |
| 4480 | ); |
| 4481 | assert!(super::live_commands_for(Path::new("/w")).is_empty()); |
| 4482 | } |
| 4483 | |
| 4484 | /// The call is bounded by `command_run_deadline` and cancelled in the host, |
| 4485 | /// which stays usable; an in-flight command fails with "host down" when the |
| 4486 | /// host is killed, and after the restart only a fresh reference works. |
| 4487 | #[tokio::test] |
| 4488 | async fn a_slow_command_is_cancelled_and_a_killed_host_fails_it_as_down() { |
| 4489 | let Some(node) = node_for_tests("slow_command") else { |
| 4490 | return; |
| 4491 | }; |
| 4492 | let _policy = TestPolicyGuard::extension_host(true); |
| 4493 | let fixture = FixturePlugins::new(&["ext-commands"]).await; |
| 4494 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 4495 | runtime: NODE, |
| 4496 | node_override: Some(node), |
| 4497 | bun_override: None, |
| 4498 | root: Some(fixture.root.clone()), |
| 4499 | supervision: super::SupervisionOptions { |
| 4500 | command_run_deadline: Duration::from_millis(300), |
| 4501 | ..fast_supervision() |
| 4502 | }, |
| 4503 | })); |
| 4504 | let engine = manager.attach(fixture.registry()); |
| 4505 | engine.sync().await.unwrap(); |
| 4506 | let reference = |name: &str| { |
| 4507 | manager |
| 4508 | .commands_for_plugins(engine.plugin_view().as_ref()) |
| 4509 | .into_iter() |
| 4510 | .find(|entry| entry.registration.name == name) |
| 4511 | .unwrap_or_else(|| panic!("{name} is not live")) |
| 4512 | .reference() |
| 4513 | }; |
| 4514 | let pid = manager.host_pid(); |
| 4515 | let slow = reference("ext-slow"); |
| 4516 | let started = Instant::now(); |
| 4517 | let error = tokio::time::timeout( |
| 4518 | Duration::from_secs(5), |
| 4519 | super::command::run(&manager.shared, &slow, "30000", None), |
| 4520 | ) |
| 4521 | .await |
| 4522 | .expect("the deadline bounds the command") |
| 4523 | .unwrap_err(); |
| 4524 | assert!(error.contains("timed out"), "{error}"); |
| 4525 | assert!(started.elapsed() < Duration::from_secs(2)); |
| 4526 | // The same host takes the next command. |
| 4527 | let echo = reference("ext-echo"); |
| 4528 | assert_eq!( |
| 4529 | super::command::run(&manager.shared, &echo, "again", None).await, |
| 4530 | Ok(super::command::CommandOutcome::Show { |
| 4531 | text: "echo: again".to_string() |
| 4532 | }) |
| 4533 | ); |
| 4534 | assert_eq!(manager.host_pid(), pid); |
| 4535 | |
| 4536 | // Kill the host under a running command with a longer deadline. |
| 4537 | let manager2 = Arc::clone(&manager); |
| 4538 | // The kill lands well inside the 300 ms deadline. |
| 4539 | let running = { |
| 4540 | let slow = slow.clone(); |
| 4541 | tokio::spawn( |
| 4542 | async move { super::command::run(&manager2.shared, &slow, "30000", None).await }, |
| 4543 | ) |
| 4544 | }; |
| 4545 | tokio::time::sleep(Duration::from_millis(100)).await; |
| 4546 | #[cfg(unix)] |
| 4547 | let status = std::process::Command::new("kill") |
| 4548 | .args(["-9", &pid.unwrap().to_string()]) |
| 4549 | .status() |
| 4550 | .unwrap(); |
| 4551 | #[cfg(windows)] |
| 4552 | let status = std::process::Command::new("taskkill") |
| 4553 | .args(["/F", "/PID", &pid.unwrap().to_string()]) |
| 4554 | .status() |
| 4555 | .unwrap(); |
| 4556 | assert!(status.success()); |
| 4557 | let error = tokio::time::timeout(Duration::from_secs(2), running) |
| 4558 | .await |
| 4559 | .expect("the call resolves") |
| 4560 | .unwrap() |
| 4561 | .unwrap_err(); |
| 4562 | assert!( |
| 4563 | error.starts_with("extension host is down: exited:"), |
| 4564 | "{error}" |
| 4565 | ); |
| 4566 | // The supervisor restarts the host and replays the plugin under a new |
| 4567 | // generation: the old reference is stale, a fresh one works. |
| 4568 | wait_host(&manager, || { |
| 4569 | manager.spawn_attempts() == 2 && manager.live_command_names().contains(&"ext-echo".into()) |
| 4570 | }) |
| 4571 | .await; |
| 4572 | let error = super::command::run(&manager.shared, &echo, "old", None) |
| 4573 | .await |
| 4574 | .unwrap_err(); |
| 4575 | assert!(error.contains("no longer registered"), "{error}"); |
| 4576 | let fresh = reference("ext-echo"); |
| 4577 | assert_ne!(fresh.generation, echo.generation); |
| 4578 | assert_eq!( |
| 4579 | super::command::run(&manager.shared, &fresh, "new", None).await, |
| 4580 | Ok(super::command::CommandOutcome::Show { |
| 4581 | text: "echo: new".to_string() |
| 4582 | }) |
| 4583 | ); |
| 4584 | manager.shutdown().await; |
| 4585 | } |
| 4586 | |
| 4587 | // --------------------------------------------------------------------------- |
| 4588 | // Trust tiers: a plugin host and a built-in (tier 0) host, never one process |
| 4589 | // --------------------------------------------------------------------------- |
| 4590 | |
| 4591 | use super::tier::{BuiltinModule, Tier0Tool}; |
| 4592 | use sha2::{Digest, Sha256}; |
| 4593 | |
| 4594 | fn tier_register_params(owner: &OwnerRef, name: &str) -> RegisterParams { |
| 4595 | RegisterParams { |
| 4596 | scope: None, |
| 4597 | owner: owner.clone(), |
| 4598 | kind: RegisterKind::Tool, |
| 4599 | spec: RegisterSpecWire { |
| 4600 | name: name.to_string(), |
| 4601 | description: "d".to_string(), |
| 4602 | input_schema: json!({"type": "object", "properties": {}}) |
| 4603 | .as_object() |
| 4604 | .cloned(), |
| 4605 | argument_hint: None, |
| 4606 | }, |
| 4607 | } |
| 4608 | } |
| 4609 | |
| 4610 | #[test] |
| 4611 | fn owners_are_bound_to_their_tier() { |
| 4612 | let mut registry = OwnerRegistry::new(); |
| 4613 | let plugin_id = "user/0123456789ab/demo"; |
| 4614 | |
| 4615 | // A `host:` id on the plugin tier, and any other id on the builtin tier. |
| 4616 | let refused = registry |
| 4617 | .begin_owner( |
| 4618 | HostTier::Plugin, |
| 4619 | "host:mcp", |
| 4620 | "mcp", |
| 4621 | Some(fake_authority("host:mcp")), |
| 4622 | "h", |
| 4623 | ) |
| 4624 | .unwrap_err(); |
| 4625 | assert!( |
| 4626 | refused.contains("a plugin id can never use it"), |
| 4627 | "{refused}" |
| 4628 | ); |
| 4629 | let refused = registry |
| 4630 | .begin_owner(HostTier::Builtin, plugin_id, "demo", None, "h") |
| 4631 | .unwrap_err(); |
| 4632 | assert!( |
| 4633 | refused.contains("not a built-in host module id"), |
| 4634 | "{refused}" |
| 4635 | ); |
| 4636 | let refused = registry |
| 4637 | .begin_owner(HostTier::Builtin, "demo", "demo", None, "h") |
| 4638 | .unwrap_err(); |
| 4639 | assert!( |
| 4640 | refused.contains("not a built-in host module id"), |
| 4641 | "{refused}" |
| 4642 | ); |
| 4643 | // The authority must fit the tier too. |
| 4644 | assert!( |
| 4645 | registry |
| 4646 | .begin_owner(HostTier::Plugin, plugin_id, "demo", None, "h") |
| 4647 | .unwrap_err() |
| 4648 | .contains("needs its reviewed plugin authority") |
| 4649 | ); |
| 4650 | assert!( |
| 4651 | registry |
| 4652 | .begin_owner( |
| 4653 | HostTier::Builtin, |
| 4654 | "host:mcp", |
| 4655 | "mcp", |
| 4656 | Some(fake_authority("host:mcp")), |
| 4657 | "h" |
| 4658 | ) |
| 4659 | .unwrap_err() |
| 4660 | .contains("has no plugin authority") |
| 4661 | ); |
| 4662 | assert!( |
| 4663 | registry.owners().next().is_none(), |
| 4664 | "a refused owner leaves nothing behind" |
| 4665 | ); |
| 4666 | |
| 4667 | let plugin = registry |
| 4668 | .begin_owner( |
| 4669 | HostTier::Plugin, |
| 4670 | plugin_id, |
| 4671 | "demo", |
| 4672 | Some(fake_authority(plugin_id)), |
| 4673 | "hash-demo", |
| 4674 | ) |
| 4675 | .unwrap(); |
| 4676 | let other_plugin = registry |
| 4677 | .begin_owner( |
| 4678 | HostTier::Plugin, |
| 4679 | "user/0123456789ab/other", |
| 4680 | "other", |
| 4681 | Some(fake_authority("user/0123456789ab/other")), |
| 4682 | "hash-other", |
| 4683 | ) |
| 4684 | .unwrap(); |
| 4685 | let builtin = registry |
| 4686 | .begin_owner(HostTier::Builtin, "host:mcp", "mcp", None, "digest-mcp") |
| 4687 | .unwrap(); |
| 4688 | for owner in [&plugin, &other_plugin, &builtin] { |
| 4689 | registry.mark_active(owner); |
| 4690 | } |
| 4691 | assert_eq!(registry.owner("host:mcp").unwrap().tier, HostTier::Builtin); |
| 4692 | assert_eq!(registry.owner(plugin_id).unwrap().tier, HostTier::Plugin); |
| 4693 | assert_eq!(registry.tier_of(&builtin), Some(HostTier::Builtin)); |
| 4694 | assert!(registry.authority_for(&builtin).is_none()); |
| 4695 | assert!(registry.authority_for(&plugin).is_some()); |
| 4696 | // Plugins share a process with each other; a built-in module does not |
| 4697 | // share one with any plugin. |
| 4698 | assert_eq!(registry.other_active_owners(plugin_id), 1); |
| 4699 | assert_eq!(registry.other_active_owners("host:mcp"), 0); |
| 4700 | |
| 4701 | let plugin_tool = register(&mut registry, &plugin, "zz_plugin_probe").unwrap(); |
| 4702 | let builtin_tool = register(&mut registry, &builtin, "zz_builtin_probe").unwrap(); |
| 4703 | let tiers: Vec<_> = registry |
| 4704 | .live_tools() |
| 4705 | .into_iter() |
| 4706 | .map(|tool| (tool.name, tool.tier)) |
| 4707 | .collect(); |
| 4708 | assert_eq!( |
| 4709 | tiers, |
| 4710 | [ |
| 4711 | ("zz_plugin_probe".to_string(), HostTier::Plugin), |
| 4712 | ("zz_builtin_probe".to_string(), HostTier::Builtin) |
| 4713 | ] |
| 4714 | ); |
| 4715 | assert!(registry.is_live(plugin_tool, &plugin)); |
| 4716 | assert!(registry.is_live(builtin_tool, &builtin)); |
| 4717 | |
| 4718 | // Each tier's host is its own process: a crash of one takes only its own |
| 4719 | // owners and registrations with it. |
| 4720 | registry.host_exited(HostTier::Plugin, "plugin host crashed"); |
| 4721 | assert!(registry.owner(plugin_id).is_none()); |
| 4722 | assert_eq!( |
| 4723 | registry.owner("host:mcp").unwrap().state, |
| 4724 | OwnerState::Active |
| 4725 | ); |
| 4726 | assert!(registry.is_live(builtin_tool, &builtin)); |
| 4727 | assert!(!registry.is_live(plugin_tool, &plugin)); |
| 4728 | // ... and a shutdown of the builtin tier leaves the plugin tier's alone. |
| 4729 | let replay = registry |
| 4730 | .begin_owner( |
| 4731 | HostTier::Plugin, |
| 4732 | plugin_id, |
| 4733 | "demo", |
| 4734 | Some(fake_authority(plugin_id)), |
| 4735 | "hash-demo", |
| 4736 | ) |
| 4737 | .unwrap(); |
| 4738 | registry.mark_active(&replay); |
| 4739 | let again = register(&mut registry, &replay, "zz_plugin_probe").unwrap(); |
| 4740 | registry.revoke_all(HostTier::Builtin, "builtin host shut down"); |
| 4741 | assert!(registry.is_live(again, &replay)); |
| 4742 | assert!(matches!( |
| 4743 | registry.owner("host:mcp").unwrap().state, |
| 4744 | OwnerState::Failed(_) |
| 4745 | )); |
| 4746 | assert_eq!(registry.live_tools().len(), 1); |
| 4747 | } |
| 4748 | |
| 4749 | /// A host answers for its own tier's owners only, even if it somehow named |
| 4750 | /// the other tier's: the registration, and the log line, are dropped. |
| 4751 | #[test] |
| 4752 | fn a_host_registers_and_logs_only_for_owners_of_its_own_tier() { |
| 4753 | use super::supervisor::HostEvents; |
| 4754 | let manager = ExtensionHostManager::new(ExtensionHostOptions::default()); |
| 4755 | let (plugin, builtin) = { |
| 4756 | let mut registry = manager.shared.registry.lock().unwrap(); |
| 4757 | let plugin = registry |
| 4758 | .begin_owner( |
| 4759 | HostTier::Plugin, |
| 4760 | "user/0123456789ab/demo", |
| 4761 | "demo", |
| 4762 | Some(fake_authority("user/0123456789ab/demo")), |
| 4763 | "hash", |
| 4764 | ) |
| 4765 | .unwrap(); |
| 4766 | let builtin = registry |
| 4767 | .begin_owner(HostTier::Builtin, "host:mcp", "mcp", None, "digest") |
| 4768 | .unwrap(); |
| 4769 | registry.mark_active(&plugin); |
| 4770 | registry.mark_active(&builtin); |
| 4771 | (plugin, builtin) |
| 4772 | }; |
| 4773 | let events = |tier| super::Events { |
| 4774 | shared: Arc::downgrade(&manager.shared), |
| 4775 | tier, |
| 4776 | generation: 0, |
| 4777 | }; |
| 4778 | let refused = events(HostTier::Plugin).register(&tier_register_params(&builtin, "zz_probe")); |
| 4779 | assert!( |
| 4780 | matches!(&refused, protocol::RegisterResult::Refused { refused } |
| 4781 | if refused.contains("belongs to the builtin tier")), |
| 4782 | "{refused:?}" |
| 4783 | ); |
| 4784 | let refused = events(HostTier::Builtin).register(&tier_register_params(&plugin, "zz_probe")); |
| 4785 | assert!( |
| 4786 | matches!(&refused, protocol::RegisterResult::Refused { refused } |
| 4787 | if refused.contains("belongs to the plugin tier")), |
| 4788 | "{refused:?}" |
| 4789 | ); |
| 4790 | assert!(manager.live_tool_names().is_empty()); |
| 4791 | assert!(matches!( |
| 4792 | events(HostTier::Builtin).register(&tier_register_params(&builtin, "zz_probe")), |
| 4793 | protocol::RegisterResult::Admitted { .. } |
| 4794 | )); |
| 4795 | assert!(matches!( |
| 4796 | events(HostTier::Plugin).register(&tier_register_params(&plugin, "zz_other")), |
| 4797 | protocol::RegisterResult::Admitted { .. } |
| 4798 | )); |
| 4799 | |
| 4800 | let warn = |plugin_id: &str| protocol::LogParams { |
| 4801 | level: "warn".into(), |
| 4802 | msg: "m".into(), |
| 4803 | plugin_id: Some(plugin_id.into()), |
| 4804 | }; |
| 4805 | events(HostTier::Plugin).log(&warn("host:mcp")); |
| 4806 | events(HostTier::Builtin).log(&warn("user/0123456789ab/demo")); |
| 4807 | assert!( |
| 4808 | manager |
| 4809 | .owner_report("host:mcp") |
| 4810 | .unwrap() |
| 4811 | .diagnostics |
| 4812 | .is_empty() |
| 4813 | ); |
| 4814 | assert!( |
| 4815 | manager |
| 4816 | .owner_report("user/0123456789ab/demo") |
| 4817 | .unwrap() |
| 4818 | .diagnostics |
| 4819 | .is_empty() |
| 4820 | ); |
| 4821 | events(HostTier::Builtin).log(&warn("host:mcp")); |
| 4822 | assert_eq!( |
| 4823 | manager.owner_report("host:mcp").unwrap().diagnostics, |
| 4824 | ["warn: m"] |
| 4825 | ); |
| 4826 | } |
| 4827 | |
| 4828 | /// A manifest name cannot carry a `host:` prefix (names are lower-case ASCII |
| 4829 | /// letters, digits and internal `-`/`.`), so discovery never builds a `host:` |
| 4830 | /// id, and the plugin that tries fails validation instead of reaching the |
| 4831 | /// host. |
| 4832 | #[test] |
| 4833 | fn a_plugin_named_like_a_tier_zero_owner_fails_validation_and_never_reaches_the_host() { |
| 4834 | let _policy = TestPolicyGuard::extension_host(true); |
| 4835 | let temp = tempfile::tempdir().unwrap(); |
| 4836 | let user = temp.path().join("user"); |
| 4837 | native_bundle(&user, "evil", "index.mjs", &["index.mjs"]); |
| 4838 | native_bundle(&user, "fine", "index.mjs", &["index.mjs"]); |
| 4839 | // The directory is harmless; the manifest's name is the claim. |
| 4840 | let manifest = user.join("evil/plugin.json"); |
| 4841 | let rewritten = std::fs::read_to_string(&manifest) |
| 4842 | .unwrap() |
| 4843 | .replace("\"name\": \"evil\"", "\"name\": \"host:evil\""); |
| 4844 | assert!(rewritten.contains("host:evil")); |
| 4845 | std::fs::write(&manifest, rewritten).unwrap(); |
| 4846 | let config = DiscoveryConfig { |
| 4847 | workspace: temp.path().join("project"), |
| 4848 | user_plugins_dir: user, |
| 4849 | workspace_plugins_dir: temp.path().join("project/.codewhale/plugins"), |
| 4850 | builtin_plugin_dirs: Vec::new(), |
| 4851 | state_path: temp.path().join("state/plugin-state.json"), |
| 4852 | }; |
| 4853 | let registry = discover_with_config(&config); |
| 4854 | assert!( |
| 4855 | registry.get("host:evil").is_none(), |
| 4856 | "a plugin named `host:evil` must not be discovered as a plugin" |
| 4857 | ); |
| 4858 | assert!( |
| 4859 | registry |
| 4860 | .diagnostics() |
| 4861 | .iter() |
| 4862 | .any(|diagnostic| diagnostic.message.contains("host:evil") |
| 4863 | && diagnostic.message.contains("name")), |
| 4864 | "{:?}", |
| 4865 | registry.diagnostics() |
| 4866 | ); |
| 4867 | assert!(registry.get("fine").is_some()); |
| 4868 | // Whatever discovery yields, no desired owner is a tier-0 id. |
| 4869 | let (desired, _) = super::desired_owners(®istry); |
| 4870 | assert!( |
| 4871 | desired |
| 4872 | .keys() |
| 4873 | .all(|id| HostTier::Plugin.check_owner_id(id).is_ok()) |
| 4874 | ); |
| 4875 | } |
| 4876 | |
| 4877 | #[test] |
| 4878 | fn launch_plans_carry_their_tier_and_use_a_data_directory_each() { |
| 4879 | #[cfg(not(windows))] |
| 4880 | use crate::dependencies::{HostRuntime, HostRuntimeKind}; |
| 4881 | let home = tempfile::tempdir().unwrap(); |
| 4882 | #[cfg(windows)] |
| 4883 | let bundle = super::materialize_bundle(home.path()).unwrap(); |
| 4884 | #[cfg(not(windows))] |
| 4885 | let bundle = home.path().join("codewhale-extension-host.mjs"); |
| 4886 | // Windows admission verifies the real copied runtime and exact bundle. |
| 4887 | #[cfg(windows)] |
| 4888 | let runtime = crate::dependencies::resolve_extension_host_runtime(NODE, None, None) |
| 4889 | .selected |
| 4890 | .expect("tier planning requires the test Node runtime on Windows"); |
| 4891 | #[cfg(not(windows))] |
| 4892 | let runtime = HostRuntime { |
| 4893 | compiled: false, |
| 4894 | kind: HostRuntimeKind::Node, |
| 4895 | path: PathBuf::from("/opt/runtime/bin/node"), |
| 4896 | version: (22, 19, 0), |
| 4897 | native_code_flags: Vec::new(), |
| 4898 | }; |
| 4899 | let mut dirs = Vec::new(); |
| 4900 | for tier in HostTier::ALL { |
| 4901 | let launch = |
| 4902 | match super::supervisor::plan_launch(tier, &runtime, &bundle, home.path(), 1 << 30) { |
| 4903 | Ok(launch) => launch, |
| 4904 | Err(error) if tier == HostTier::Plugin => { |
| 4905 | assert!( |
| 4906 | error.contains("Native extensions require a verified OS sandbox"), |
| 4907 | "{error}" |
| 4908 | ); |
| 4909 | assert!(matches!( |
| 4910 | super::supervisor::planned_sandbox(tier, &runtime, home.path()), |
| 4911 | Ok(super::supervisor::HostSandbox::Unsandboxed(_)) |
| 4912 | )); |
| 4913 | let data = super::supervisor::tier_data_dir(home.path(), tier); |
| 4914 | assert!(data.is_dir()); |
| 4915 | assert!( |
| 4916 | super::supervisor::tier_data_dir(home.path(), HostTier::Builtin).is_dir() |
| 4917 | ); |
| 4918 | dirs.push(data); |
| 4919 | continue; |
| 4920 | } |
| 4921 | Err(error) => panic!("pinned Builtin plan failed: {error}"), |
| 4922 | }; |
| 4923 | assert_eq!(launch.tier, tier); |
| 4924 | // The runtime's own argv ends `<bundle> --tier=<tier>`, wrapped by the |
| 4925 | // OS sandbox or not. |
| 4926 | let at = launch |
| 4927 | .args |
| 4928 | .iter() |
| 4929 | .position(|arg| Path::new(arg) == bundle) |
| 4930 | .expect("bundle in argv"); |
| 4931 | assert_eq!(launch.args[at + 1], format!("--tier={}", tier.name())); |
| 4932 | assert_eq!(launch.args.len(), at + 2, "the tier is the last argument"); |
| 4933 | assert_eq!( |
| 4934 | launch.cwd, |
| 4935 | super::supervisor::tier_data_dir(home.path(), tier) |
| 4936 | ); |
| 4937 | assert!(launch.cwd.is_dir(), "{}", launch.cwd.display()); |
| 4938 | dirs.push(launch.cwd); |
| 4939 | } |
| 4940 | assert_ne!(dirs[0], dirs[1], "each tier has its own data directory"); |
| 4941 | assert_eq!(HostTier::Plugin.argv_flag(), "--tier=plugin"); |
| 4942 | assert_eq!(HostTier::Builtin.argv_flag(), "--tier=builtin"); |
| 4943 | } |
| 4944 | |
| 4945 | /// A built-in module table that pins `source`, for a test. Production has |
| 4946 | /// none: `BUILTIN_MODULES` is empty. |
| 4947 | fn test_builtin_table(source: &Path) -> &'static [BuiltinModule] { |
| 4948 | let digest = super::hex(Sha256::digest(std::fs::read(source).unwrap())); |
| 4949 | let tools: &'static [Tier0Tool] = Box::leak(Box::new([Tier0Tool { |
| 4950 | name: "zz_tier0_listed", |
| 4951 | approval: ApprovalRequirement::Auto, |
| 4952 | }])); |
| 4953 | Box::leak(Box::new([BuiltinModule { |
| 4954 | id: "tier0-module", |
| 4955 | source_sha256: Box::leak(digest.into_boxed_str()), |
| 4956 | tools, |
| 4957 | }])) |
| 4958 | } |
| 4959 | |
| 4960 | /// Put `bytes` where the core looks for built-in module `id`'s source under |
| 4961 | /// the home `root`. |
| 4962 | fn place_builtin_source(root: &Path, id: &str, bytes: &[u8]) { |
| 4963 | let path = super::supervisor::bundle_dir(root, super::bundle_sha256()) |
| 4964 | .join("builtin") |
| 4965 | .join(format!("{id}.mjs")); |
| 4966 | std::fs::create_dir_all(path.parent().unwrap()).unwrap(); |
| 4967 | std::fs::write(path, bytes).unwrap(); |
| 4968 | } |
| 4969 | |
| 4970 | /// The production table pins MCP, but ordinary plugin attachment never |
| 4971 | /// starts that independent builtin backend until the Engine selects it. |
| 4972 | #[tokio::test] |
| 4973 | async fn plugin_attachment_never_spawns_the_independent_builtin_mcp_backend() { |
| 4974 | let manager = ExtensionHostManager::new(ExtensionHostOptions::default()); |
| 4975 | assert_eq!( |
| 4976 | manager.shared.builtin_modules.len(), |
| 4977 | super::tier::BUILTIN_MODULES.len() |
| 4978 | ); |
| 4979 | assert!( |
| 4980 | manager |
| 4981 | .shared |
| 4982 | .builtin_modules |
| 4983 | .iter() |
| 4984 | .any(|module| module.id == "mcp") |
| 4985 | ); |
| 4986 | let Some(node) = |
| 4987 | node_for_tests("plugin_attachment_never_spawns_the_independent_builtin_mcp_backend") |
| 4988 | else { |
| 4989 | return; |
| 4990 | }; |
| 4991 | let _policy = TestPolicyGuard::extension_host(true); |
| 4992 | let fixture = FixturePlugins::new(&["dsh-workspace-deps"]).await; |
| 4993 | let manager = fixture.manager(node); |
| 4994 | let engine = manager.attach(fixture.registry()); |
| 4995 | engine.sync().await.unwrap(); |
| 4996 | assert!(matches!(manager.status(), HostStatus::Ready { .. })); |
| 4997 | assert_eq!(manager.tier_spawn_attempts(HostTier::Plugin), 1); |
| 4998 | assert_eq!(manager.tier_spawn_attempts(HostTier::Builtin), 0); |
| 4999 | assert_eq!(manager.tier_status(HostTier::Builtin), HostStatus::Idle); |
| 5000 | // Planning creates an empty sibling so the Native sandbox can mask it, |
| 5001 | // including on Linux where bubblewrap requires the denied root to exist. |
| 5002 | // That directory is not evidence of a Builtin process or backend. |
| 5003 | let builtin_data = super::supervisor::tier_data_dir(&fixture.root, HostTier::Builtin); |
| 5004 | assert!(builtin_data.is_dir()); |
| 5005 | assert_eq!(std::fs::read_dir(builtin_data).unwrap().count(), 0); |
| 5006 | let report = super::render_status(&manager); |
| 5007 | assert!(!report.contains("built-in host"), "{report}"); |
| 5008 | manager.shutdown().await; |
| 5009 | } |
| 5010 | |
| 5011 | /// A module whose source is not the digest the table pins, or is missing, is |
| 5012 | /// refused with the reason, and no host is started for it. |
| 5013 | #[tokio::test] |
| 5014 | async fn a_built_in_module_that_is_not_the_pinned_source_is_refused_and_starts_nothing() { |
| 5015 | let _policy = TestPolicyGuard::extension_host(true); |
| 5016 | let source = fixtures_dir().join("tier0-module/module.mjs"); |
| 5017 | let modules = test_builtin_table(&source); |
| 5018 | for (case, bytes, expected) in [ |
| 5019 | ( |
| 5020 | "tampered", |
| 5021 | Some(&b"export const name = 'tier0-module'\nexport function apply() {}\n"[..]), |
| 5022 | "the core pins", |
| 5023 | ), |
| 5024 | ("missing", None, "cannot read the source"), |
| 5025 | ] { |
| 5026 | let temp = tempfile::tempdir().unwrap(); |
| 5027 | let root = temp.path().join("home"); |
| 5028 | if let Some(bytes) = bytes { |
| 5029 | place_builtin_source(&root, "tier0-module", bytes); |
| 5030 | } |
| 5031 | let manager = Arc::new(ExtensionHostManager::with_builtin_modules( |
| 5032 | ExtensionHostOptions { |
| 5033 | root: Some(root), |
| 5034 | ..Default::default() |
| 5035 | }, |
| 5036 | modules, |
| 5037 | )); |
| 5038 | let engine = manager.attach(Arc::new(PluginRegistry::empty(temp.path()))); |
| 5039 | engine.sync().await.unwrap(); |
| 5040 | assert_eq!(manager.spawn_attempts(), 0, "{case}"); |
| 5041 | assert_eq!(manager.tier_status(HostTier::Builtin), HostStatus::Idle); |
| 5042 | let Some(OwnerState::Failed(reason)) = manager.owner_state("host:tier0-module") else { |
| 5043 | panic!("{case}: {:?}", manager.owner_state("host:tier0-module")); |
| 5044 | }; |
| 5045 | assert!(reason.contains(expected), "{case}: {reason}"); |
| 5046 | assert!(manager.live_tool_names().is_empty()); |
| 5047 | // Not retried every turn. |
| 5048 | engine.sync().await.unwrap(); |
| 5049 | assert_eq!( |
| 5050 | manager.diagnostics().len(), |
| 5051 | 1, |
| 5052 | "{:?}", |
| 5053 | manager.diagnostics() |
| 5054 | ); |
| 5055 | } |
| 5056 | } |
| 5057 | |
| 5058 | /// The two tiers are two processes under one manager: a built-in module |
| 5059 | /// activates under `host:<module>` in its own host, its tool's approval is |
| 5060 | /// whatever the table says (and `Required` where it says nothing), plugin |
| 5061 | /// tools stay `Required`, and crashing the plugin host disturbs nothing of |
| 5062 | /// the builtin one. |
| 5063 | #[tokio::test] |
| 5064 | async fn a_tier_zero_host_runs_apart_from_the_plugin_host_and_its_tool_approval_follows_the_table() |
| 5065 | { |
| 5066 | let Some(node) = node_for_tests("a_tier_zero_host_runs_apart") else { |
| 5067 | return; |
| 5068 | }; |
| 5069 | let _policy = TestPolicyGuard::extension_host(true); |
| 5070 | let fixture = FixturePlugins::new(&["dsh-workspace-deps"]).await; |
| 5071 | let source = fixtures_dir().join("tier0-module/module.mjs"); |
| 5072 | place_builtin_source( |
| 5073 | &fixture.root, |
| 5074 | "tier0-module", |
| 5075 | &std::fs::read(&source).unwrap(), |
| 5076 | ); |
| 5077 | let manager = Arc::new(ExtensionHostManager::with_builtin_modules( |
| 5078 | ExtensionHostOptions { |
| 5079 | runtime: NODE, |
| 5080 | node_override: Some(node), |
| 5081 | root: Some(fixture.root.clone()), |
| 5082 | ..Default::default() |
| 5083 | }, |
| 5084 | test_builtin_table(&source), |
| 5085 | )); |
| 5086 | assert_eq!(manager.tier_status(HostTier::Builtin), HostStatus::Idle); |
| 5087 | let engine = manager.attach(fixture.registry()); |
| 5088 | engine.sync().await.unwrap(); |
| 5089 | |
| 5090 | // Two hosts, two processes, each launched once. |
| 5091 | let plugin_pid = manager.host_pid().expect("plugin host running"); |
| 5092 | let builtin_pid = manager |
| 5093 | .shared |
| 5094 | .ready_host(HostTier::Builtin) |
| 5095 | .expect("builtin host running") |
| 5096 | .pid |
| 5097 | .expect("builtin pid"); |
| 5098 | assert_ne!(plugin_pid, builtin_pid); |
| 5099 | assert_eq!(manager.tier_spawn_attempts(HostTier::Plugin), 1); |
| 5100 | assert_eq!(manager.tier_spawn_attempts(HostTier::Builtin), 1); |
| 5101 | let plugin_host_id = plugin_id(&fixture, "dsh-workspace-deps"); |
| 5102 | assert_eq!( |
| 5103 | manager.owner_state("host:tier0-module"), |
| 5104 | Some(OwnerState::Active) |
| 5105 | ); |
| 5106 | assert_eq!( |
| 5107 | manager.owner_state(&plugin_host_id), |
| 5108 | Some(OwnerState::Active) |
| 5109 | ); |
| 5110 | { |
| 5111 | let registry = manager.shared.registry.lock().unwrap(); |
| 5112 | assert_eq!( |
| 5113 | registry.owner("host:tier0-module").unwrap().tier, |
| 5114 | HostTier::Builtin |
| 5115 | ); |
| 5116 | assert_eq!( |
| 5117 | registry.owner(&plugin_host_id).unwrap().tier, |
| 5118 | HostTier::Plugin |
| 5119 | ); |
| 5120 | assert_eq!( |
| 5121 | registry.owner("host:tier0-module").unwrap().content_hash, |
| 5122 | test_builtin_table(&source)[0].source_sha256 |
| 5123 | ); |
| 5124 | } |
| 5125 | // The module's data directory is under the builtin tier's, apart from |
| 5126 | // every plugin's. |
| 5127 | assert!( |
| 5128 | super::supervisor::owner_data_dir( |
| 5129 | &fixture.root, |
| 5130 | HostTier::Builtin, |
| 5131 | "host:tier0-module", |
| 5132 | "tier0-module" |
| 5133 | ) |
| 5134 | .is_dir() |
| 5135 | ); |
| 5136 | |
| 5137 | // An engine installs plugin tools only; tier-0 tools are not offered to |
| 5138 | // the model by a plugin snapshot. |
| 5139 | let installed = installed(&engine, fixture.workspace()); |
| 5140 | assert_eq!(installed, ["load_workspace_dependencies"]); |
| 5141 | |
| 5142 | // Approval follows the table for tier 0, and stays Required for plugins. |
| 5143 | let registrations = manager.shared.registry.lock().unwrap().live_tools(); |
| 5144 | let spec_for = |name: &str| { |
| 5145 | let registration = registrations |
| 5146 | .iter() |
| 5147 | .find(|tool| tool.name == name) |
| 5148 | .unwrap_or_else(|| panic!("{name} not registered")) |
| 5149 | .clone(); |
| 5150 | super::tool::HostToolSpec::new(registration, Arc::clone(&manager.shared)) |
| 5151 | }; |
| 5152 | let listed = spec_for("zz_tier0_listed"); |
| 5153 | let unlisted = spec_for("zz_tier0_unlisted"); |
| 5154 | let plugin_tool = spec_for("load_workspace_dependencies"); |
| 5155 | assert_eq!(listed.registration_origin(), "host:tier0-module"); |
| 5156 | assert_eq!(listed.approval_requirement(), ApprovalRequirement::Auto); |
| 5157 | assert_eq!( |
| 5158 | listed.approval_requirement_for(&json!({})), |
| 5159 | ApprovalRequirement::Auto |
| 5160 | ); |
| 5161 | assert_eq!( |
| 5162 | unlisted.approval_requirement(), |
| 5163 | ApprovalRequirement::Required |
| 5164 | ); |
| 5165 | assert_eq!( |
| 5166 | plugin_tool.approval_requirement(), |
| 5167 | ApprovalRequirement::Required |
| 5168 | ); |
| 5169 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(engine.plugin_view()); |
| 5170 | assert_eq!( |
| 5171 | listed.prepare(json!({}), &context).unwrap().approval, |
| 5172 | ApprovalRequirement::Auto |
| 5173 | ); |
| 5174 | assert_eq!( |
| 5175 | unlisted.prepare(json!({}), &context).unwrap().approval, |
| 5176 | ApprovalRequirement::Required |
| 5177 | ); |
| 5178 | // Even an Auto tool is host code: never read-only, never plan-mode safe. |
| 5179 | assert!(!listed.is_read_only_for(&json!({}))); |
| 5180 | |
| 5181 | // The call goes to the builtin host and comes back. |
| 5182 | let result = listed.execute(json!({}), &context).await.unwrap(); |
| 5183 | assert!(result.success); |
| 5184 | assert_eq!( |
| 5185 | serde_json::from_str::<Value>(&result.content).unwrap(), |
| 5186 | json!({"tier": "zero", "listed": true}) |
| 5187 | ); |
| 5188 | |
| 5189 | // Kill the plugin host: its owners go and come back, the builtin host and |
| 5190 | // its module are untouched, and its crash budget is not spent. |
| 5191 | #[cfg(unix)] |
| 5192 | let status = std::process::Command::new("kill") |
| 5193 | .args(["-9", &plugin_pid.to_string()]) |
| 5194 | .status() |
| 5195 | .unwrap(); |
| 5196 | #[cfg(windows)] |
| 5197 | let status = std::process::Command::new("taskkill") |
| 5198 | .args(["/F", "/PID", &plugin_pid.to_string()]) |
| 5199 | .status() |
| 5200 | .unwrap(); |
| 5201 | assert!(status.success()); |
| 5202 | wait_host(&manager, || { |
| 5203 | manager.tier_spawn_attempts(HostTier::Plugin) == 2 |
| 5204 | && manager |
| 5205 | .live_tool_names() |
| 5206 | .contains(&"load_workspace_dependencies".to_string()) |
| 5207 | }) |
| 5208 | .await; |
| 5209 | assert_eq!( |
| 5210 | manager.shared.ready_host(HostTier::Builtin).unwrap().pid, |
| 5211 | Some(builtin_pid) |
| 5212 | ); |
| 5213 | assert_eq!(manager.tier_spawn_attempts(HostTier::Builtin), 1); |
| 5214 | assert_eq!( |
| 5215 | manager.owner_state("host:tier0-module"), |
| 5216 | Some(OwnerState::Active) |
| 5217 | ); |
| 5218 | assert!( |
| 5219 | manager |
| 5220 | .shared |
| 5221 | .builtin |
| 5222 | .supervision |
| 5223 | .lock() |
| 5224 | .unwrap() |
| 5225 | .crashes |
| 5226 | .is_empty() |
| 5227 | ); |
| 5228 | assert_eq!( |
| 5229 | manager |
| 5230 | .shared |
| 5231 | .plugin |
| 5232 | .supervision |
| 5233 | .lock() |
| 5234 | .unwrap() |
| 5235 | .crashes |
| 5236 | .len(), |
| 5237 | 1 |
| 5238 | ); |
| 5239 | assert!(listed.execute(json!({}), &context).await.unwrap().success); |
| 5240 | |
| 5241 | // The status page shows the builtin host and keeps its tools off the |
| 5242 | // plugin list. |
| 5243 | let report = super::render_status(&manager); |
| 5244 | assert!( |
| 5245 | report.contains("built-in host (tier 0): running"), |
| 5246 | "{report}" |
| 5247 | ); |
| 5248 | assert!(!report.contains("tool zz_tier0_listed"), "{report}"); |
| 5249 | assert!(!report.contains("command /"), "{report}"); |
| 5250 | manager.shutdown().await; |
| 5251 | } |
| 5252 | |
| 5253 | /// The existing installer, registry projection and final spec invocation use |
| 5254 | /// one attachment receipt. Two entry scopes under one owner never union into |
| 5255 | /// either caller, and withdrawing one does not cancel its sibling. |
| 5256 | #[tokio::test(flavor = "current_thread")] |
| 5257 | async fn native_preset_membership_filters_discovery_and_final_tool_invocation() { |
| 5258 | let Some(node) = node_for_tests("native_preset_membership") else { |
| 5259 | return; |
| 5260 | }; |
| 5261 | let _policy = TestPolicyGuard::extension_host(true); |
| 5262 | let _catalog = stub_builtin_commands(); |
| 5263 | let fixture = FixturePlugins::new(&["two-entries"]).await; |
| 5264 | let manager = fixture.manager(node); |
| 5265 | let _manager = super::TestManagerGuard::install(Arc::clone(&manager)); |
| 5266 | let initial = manager.attach(fixture.registry()); |
| 5267 | initial.sync().await.unwrap(); |
| 5268 | let presets = super::native_presets_for_plugins(initial.plugin_view().as_ref()); |
| 5269 | let first = presets |
| 5270 | .iter() |
| 5271 | .find(|(preset, _)| preset.entry.path.ends_with("tools.mjs")) |
| 5272 | .unwrap() |
| 5273 | .0 |
| 5274 | .clone(); |
| 5275 | let second = presets |
| 5276 | .iter() |
| 5277 | .find(|(preset, _)| preset.entry.path.ends_with("commands.mjs")) |
| 5278 | .unwrap() |
| 5279 | .0 |
| 5280 | .clone(); |
| 5281 | let a = manager.attach(Arc::new( |
| 5282 | fixture.registry().with_native_preset(first).unwrap(), |
| 5283 | )); |
| 5284 | let b = manager.attach(Arc::new( |
| 5285 | fixture.registry().with_native_preset(second).unwrap(), |
| 5286 | )); |
| 5287 | drop(initial); |
| 5288 | a.sync().await.unwrap(); |
| 5289 | assert_eq!(installed(&a, fixture.workspace()), ["two_first"]); |
| 5290 | assert_eq!(installed(&b, fixture.workspace()), ["two_second"]); |
| 5291 | assert!( |
| 5292 | manager |
| 5293 | .commands_for_plugins(a.plugin_view().as_ref()) |
| 5294 | .is_empty() |
| 5295 | ); |
| 5296 | assert_eq!( |
| 5297 | manager.commands_for_plugins(b.plugin_view().as_ref())[0] |
| 5298 | .registration |
| 5299 | .name, |
| 5300 | "two-hello" |
| 5301 | ); |
| 5302 | let old = host_tool(&a, fixture.workspace(), "two_first"); |
| 5303 | let a_context = ToolContext::new(fixture.workspace()).with_plugin_registry(a.plugin_view()); |
| 5304 | let b_context = ToolContext::new(fixture.workspace()).with_plugin_registry(b.plugin_view()); |
| 5305 | assert!( |
| 5306 | old.prepare(json!({}), &b_context).is_err(), |
| 5307 | "a retained spec cannot execute under another caller" |
| 5308 | ); |
| 5309 | assert_eq!( |
| 5310 | old.execute(json!({}), &a_context).await.unwrap().content, |
| 5311 | "first" |
| 5312 | ); |
| 5313 | a.set_plugins(Arc::new(PluginRegistry::empty(fixture.workspace()))); |
| 5314 | assert!( |
| 5315 | old.prepare(json!({}), &a_context).is_err(), |
| 5316 | "withdrawal rejects a retained approval/spec receipt before reconciliation" |
| 5317 | ); |
| 5318 | a.sync().await.unwrap(); |
| 5319 | assert_eq!( |
| 5320 | host_tool(&b, fixture.workspace(), "two_second") |
| 5321 | .execute(json!({}), &b_context) |
| 5322 | .await |
| 5323 | .unwrap() |
| 5324 | .content, |
| 5325 | "second" |
| 5326 | ); |
| 5327 | // Rediscovery must preserve a now-invalid narrowed selector. A disabled |
| 5328 | // build never turns a selected caller into a broad default caller. |
| 5329 | manager.refresh_workspace(&fixture.disable("two-entries")); |
| 5330 | assert!(!b.plugin_view().selected_native_entries().is_empty()); |
| 5331 | b.sync().await.unwrap(); |
| 5332 | assert!(installed(&b, fixture.workspace()).is_empty()); |
| 5333 | manager.shutdown().await; |
| 5334 | } |
| 5335 | |
| 5336 | /// A real installed raw roster: initial default and two child snapshots share |
| 5337 | /// one owner but all five discovery paths and final invocation use one receipt. |
| 5338 | #[tokio::test(flavor = "current_thread")] |
| 5339 | async fn raw_agent_presets_use_one_default_and_all_five_caller_views() { |
| 5340 | let Some(node) = node_for_tests("raw_agent_presets_use_one_default_and_all_five_caller_views") |
| 5341 | else { |
| 5342 | return; |
| 5343 | }; |
| 5344 | let _policy = TestPolicyGuard::extension_host(true); |
| 5345 | let _catalog = stub_builtin_commands(); |
| 5346 | let fixture = FixturePlugins::new(&["raw-agent-presets"]).await; |
| 5347 | let original = fixture.registry(); |
| 5348 | assert_eq!(original.selected_native_entries().len(), 1); |
| 5349 | assert!( |
| 5350 | Path::new(&original.selected_native_entries()[0].entry.path).file_name() |
| 5351 | == Some(std::ffi::OsStr::new("a.mjs")) |
| 5352 | ); |
| 5353 | let manager = fixture.manager(node); |
| 5354 | let _manager = super::TestManagerGuard::install(Arc::clone(&manager)); |
| 5355 | let a = manager.attach(Arc::clone(&original)); |
| 5356 | a.sync().await.unwrap(); |
| 5357 | let roster = super::native_presets_for_plugins(a.plugin_view().as_ref()); |
| 5358 | assert_eq!( |
| 5359 | roster.len(), |
| 5360 | 2, |
| 5361 | "roster discovery offers admitted alternatives without activating a union" |
| 5362 | ); |
| 5363 | assert_eq!(a.prompt_sections().await.unwrap()[0].text, "A:a"); |
| 5364 | let selected_b = roster |
| 5365 | .iter() |
| 5366 | .find(|(preset, _)| { |
| 5367 | Path::new(&preset.entry.path).file_name() == Some(std::ffi::OsStr::new("b.mjs")) |
| 5368 | }) |
| 5369 | .unwrap() |
| 5370 | .0 |
| 5371 | .clone(); |
| 5372 | let b = manager.attach(Arc::new(original.with_native_preset(selected_b).unwrap())); |
| 5373 | b.sync().await.unwrap(); |
| 5374 | let check = |view: &HostAttachment, expected: &str| { |
| 5375 | assert_eq!(installed(view, fixture.workspace()), ["preset_echo"]); |
| 5376 | let commands = manager.commands_for_plugins(view.plugin_view().as_ref()); |
| 5377 | assert_eq!(commands.len(), 1); |
| 5378 | assert_eq!(commands[0].registration.name, "preset-echo"); |
| 5379 | let roots = super::skills::roots_for_plugins(view.plugin_view().as_ref()); |
| 5380 | assert_eq!(roots.len(), 1); |
| 5381 | assert!(roots[0].0.path.ends_with(expected)); |
| 5382 | assert_eq!(roots[0].0.snapshots[0].name, "preset-note"); |
| 5383 | }; |
| 5384 | check(&a, "a"); |
| 5385 | check(&b, "b"); |
| 5386 | assert_eq!(a.prompt_sections().await.unwrap()[0].text, "A:a"); |
| 5387 | assert_eq!(b.prompt_sections().await.unwrap()[0].text, "B:b"); |
| 5388 | let hook_payload = protocol::HookCallPayload { |
| 5389 | name: "read".into(), |
| 5390 | call_id: "raw-preset-hook".into(), |
| 5391 | input: json!({}), |
| 5392 | mode: "Agent".into(), |
| 5393 | workspace: fixture.workspace().to_string_lossy().into_owned(), |
| 5394 | model: "fixture".into(), |
| 5395 | }; |
| 5396 | let hooks_a = a.tool_before_hooks(hook_payload.clone()).await; |
| 5397 | let hooks_b = b.tool_before_hooks(hook_payload).await; |
| 5398 | assert_eq!(hooks_a.len(), 1); |
| 5399 | assert_eq!(hooks_b.len(), 1); |
| 5400 | assert_eq!( |
| 5401 | serde_json::from_str::<Value>(&hooks_a[0].stdout).unwrap()["additionalContext"], |
| 5402 | "A:a" |
| 5403 | ); |
| 5404 | assert_eq!( |
| 5405 | serde_json::from_str::<Value>(&hooks_b[0].stdout).unwrap()["additionalContext"], |
| 5406 | "B:b" |
| 5407 | ); |
| 5408 | let a_context = ToolContext::new(fixture.workspace()).with_plugin_registry(a.plugin_view()); |
| 5409 | let b_context = ToolContext::new(fixture.workspace()).with_plugin_registry(b.plugin_view()); |
| 5410 | let retained = host_tool(&a, fixture.workspace(), "preset_echo"); |
| 5411 | assert!(retained.prepare(json!({}), &b_context).is_err()); |
| 5412 | assert_eq!( |
| 5413 | retained |
| 5414 | .execute(json!({}), &a_context) |
| 5415 | .await |
| 5416 | .unwrap() |
| 5417 | .content, |
| 5418 | "A:a" |
| 5419 | ); |
| 5420 | let command = manager.commands_for_plugins(a.plugin_view().as_ref())[0].reference(); |
| 5421 | assert!( |
| 5422 | super::run_command_for_plugins(&command, "", None, b.plugin_view().as_ref()) |
| 5423 | .await |
| 5424 | .is_err() |
| 5425 | ); |
| 5426 | assert_eq!( |
| 5427 | super::run_command_for_plugins(&command, "", None, a.plugin_view().as_ref()) |
| 5428 | .await |
| 5429 | .unwrap(), |
| 5430 | super::command::CommandOutcome::Show { text: "A:a".into() } |
| 5431 | ); |
| 5432 | a.set_plugins(Arc::new(PluginRegistry::empty(fixture.workspace()))); |
| 5433 | assert!(retained.prepare(json!({}), &a_context).is_err()); |
| 5434 | a.sync().await.unwrap(); |
| 5435 | assert!(a.prompt_sections().await.unwrap().is_empty()); |
| 5436 | assert_eq!( |
| 5437 | host_tool(&b, fixture.workspace(), "preset_echo") |
| 5438 | .execute(json!({}), &b_context) |
| 5439 | .await |
| 5440 | .unwrap() |
| 5441 | .content, |
| 5442 | "B:b" |
| 5443 | ); |
| 5444 | let disabled = fixture.disable("raw-agent-presets"); |
| 5445 | manager.refresh_workspace(&disabled); |
| 5446 | assert!(!b.plugin_view().selected_native_entries().is_empty()); |
| 5447 | b.sync().await.unwrap(); |
| 5448 | assert!(installed(&b, fixture.workspace()).is_empty()); |
| 5449 | assert!(b.prompt_sections().await.unwrap().is_empty()); |
| 5450 | assert!(super::skills::roots_for_plugins(b.plugin_view().as_ref()).is_empty()); |
| 5451 | assert!( |
| 5452 | manager |
| 5453 | .commands_for_plugins(b.plugin_view().as_ref()) |
| 5454 | .is_empty() |
| 5455 | ); |
| 5456 | assert!( |
| 5457 | b.tool_before_hooks(protocol::HookCallPayload { |
| 5458 | name: "read".into(), |
| 5459 | call_id: "withdrawn".into(), |
| 5460 | input: json!({}), |
| 5461 | mode: "Agent".into(), |
| 5462 | workspace: fixture.workspace().to_string_lossy().into_owned(), |
| 5463 | model: "fixture".into() |
| 5464 | }) |
| 5465 | .await |
| 5466 | .is_empty() |
| 5467 | ); |
| 5468 | manager.shutdown().await; |
| 5469 | } |
| 5470 | |
| 5471 | /// No default is an upstream fact. Catalog discovery must neither start the |
| 5472 | /// host nor grant any contribution until the exact child receipt is selected. |
| 5473 | #[tokio::test(flavor = "current_thread")] |
| 5474 | async fn raw_agent_presets_without_default_require_explicit_child_selection() { |
| 5475 | let Some(node) = |
| 5476 | node_for_tests("raw_agent_presets_without_default_require_explicit_child_selection") |
| 5477 | else { |
| 5478 | return; |
| 5479 | }; |
| 5480 | let _policy = TestPolicyGuard::extension_host(true); |
| 5481 | let _catalog = stub_builtin_commands(); |
| 5482 | let fixture = FixturePlugins::new(&["raw-agent-presets-no-default"]).await; |
| 5483 | let original = fixture.registry(); |
| 5484 | assert!(original.selected_native_entries().is_empty()); |
| 5485 | let manager = fixture.manager(node); |
| 5486 | let _manager = super::TestManagerGuard::install(Arc::clone(&manager)); |
| 5487 | let roster = super::native_presets_for_plugins(original.as_ref()); |
| 5488 | assert_eq!( |
| 5489 | roster.len(), |
| 5490 | 2, |
| 5491 | "healthy admitted catalog is available before a host exists" |
| 5492 | ); |
| 5493 | let idle = manager.attach(Arc::clone(&original)); |
| 5494 | idle.sync().await.unwrap(); |
| 5495 | assert!(installed(&idle, fixture.workspace()).is_empty()); |
| 5496 | assert!(idle.prompt_sections().await.unwrap().is_empty()); |
| 5497 | assert!( |
| 5498 | manager |
| 5499 | .commands_for_plugins(idle.plugin_view().as_ref()) |
| 5500 | .is_empty() |
| 5501 | ); |
| 5502 | assert!(super::skills::roots_for_plugins(idle.plugin_view().as_ref()).is_empty()); |
| 5503 | assert!( |
| 5504 | idle.tool_before_hooks(protocol::HookCallPayload { |
| 5505 | name: "read".into(), |
| 5506 | call_id: "unselected".into(), |
| 5507 | input: json!({}), |
| 5508 | mode: "Agent".into(), |
| 5509 | workspace: fixture.workspace().to_string_lossy().into_owned(), |
| 5510 | model: "fixture".into() |
| 5511 | }) |
| 5512 | .await |
| 5513 | .is_empty() |
| 5514 | ); |
| 5515 | let selected = roster |
| 5516 | .iter() |
| 5517 | .find(|(preset, _)| { |
| 5518 | Path::new(&preset.entry.path).file_name() == Some(std::ffi::OsStr::new("b.mjs")) |
| 5519 | }) |
| 5520 | .unwrap() |
| 5521 | .0 |
| 5522 | .clone(); |
| 5523 | let child = manager.attach(Arc::new(original.with_native_preset(selected).unwrap())); |
| 5524 | child.sync().await.unwrap(); |
| 5525 | assert_eq!(installed(&child, fixture.workspace()), ["preset_echo"]); |
| 5526 | assert_eq!(child.prompt_sections().await.unwrap()[0].text, "B:b"); |
| 5527 | let context = ToolContext::new(fixture.workspace()).with_plugin_registry(child.plugin_view()); |
| 5528 | assert_eq!( |
| 5529 | host_tool(&child, fixture.workspace(), "preset_echo") |
| 5530 | .execute(json!({}), &context) |
| 5531 | .await |
| 5532 | .unwrap() |
| 5533 | .content, |
| 5534 | "B:b" |
| 5535 | ); |
| 5536 | manager.refresh_workspace(&original); |
| 5537 | idle.sync().await.unwrap(); |
| 5538 | assert!(idle.plugin_view().selected_native_entries().is_empty()); |
| 5539 | assert!( |
| 5540 | installed(&idle, fixture.workspace()).is_empty(), |
| 5541 | "another caller's selected owner cannot broaden the catalog-only caller" |
| 5542 | ); |
| 5543 | assert!(idle.prompt_sections().await.unwrap().is_empty()); |
| 5544 | assert_eq!( |
| 5545 | super::native_presets_for_plugins(idle.plugin_view().as_ref()).len(), |
| 5546 | 2 |
| 5547 | ); |
| 5548 | manager.shutdown().await; |
| 5549 | } |
| 5550 | |
| 5551 | /// The real configured-hook consumers redeem the pinned Builtin across every |
| 5552 | /// existing firepoint, and a changed project receipt cannot fall back. |
| 5553 | #[cfg(unix)] |
| 5554 | #[tokio::test(flavor = "current_thread")] |
| 5555 | async fn all_fifteen_project_hook_events_use_the_pinned_runner_and_reject_changed_receipts() { |
| 5556 | let _env = crate::test_support::lock_test_env(); |
| 5557 | let _policy = TestPolicyGuard::extension_host(true); |
| 5558 | let Some(node) = node_for_tests("all_fifteen_project_hook_events") else { |
| 5559 | return; |
| 5560 | }; |
| 5561 | let home = tempfile::tempdir().unwrap(); |
| 5562 | let workspace = home.path().join("workspace"); |
| 5563 | std::fs::create_dir_all(workspace.join(".codewhale")).unwrap(); |
| 5564 | let workspace = workspace.canonicalize().unwrap(); |
| 5565 | let _config = crate::test_support::EnvVarGuard::set( |
| 5566 | "CODEWHALE_CONFIG_PATH", |
| 5567 | home.path().join("config.toml"), |
| 5568 | ); |
| 5569 | let mut config = crate::hooks::HooksConfig { |
| 5570 | enabled: true, |
| 5571 | ..crate::hooks::HooksConfig::default() |
| 5572 | }; |
| 5573 | for event in crate::hooks::config::ALL_HOOK_EVENTS { |
| 5574 | config.hooks.push(crate::hooks::Hook::new( |
| 5575 | event, |
| 5576 | &format!( |
| 5577 | "printf '%s|%s|%s' '{}' \"$DEEPSEEK_SESSION_ID\" \"$DEEPSEEK_TOOL_CALL_ID\"", |
| 5578 | event.as_str() |
| 5579 | ), |
| 5580 | )); |
| 5581 | } |
| 5582 | let hook_path = workspace.join(".codewhale/hooks.toml"); |
| 5583 | std::fs::write(&hook_path, toml::to_string(&config).unwrap()).unwrap(); |
| 5584 | crate::config::save_workspace_trust(&workspace).unwrap(); |
| 5585 | let (reviewed, _) = crate::hooks::authority::review_project_hooks(&workspace).unwrap(); |
| 5586 | crate::hooks::authority::approve_project_hooks(&workspace, &reviewed.digest).unwrap(); |
| 5587 | let admitted = crate::hooks::HooksConfig::load_with_project( |
| 5588 | crate::hooks::HooksConfig { |
| 5589 | enabled: true, |
| 5590 | ..crate::hooks::HooksConfig::default() |
| 5591 | }, |
| 5592 | &workspace, |
| 5593 | ); |
| 5594 | assert_eq!(admitted.hooks.len(), 15); |
| 5595 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 5596 | node_override: Some(node), |
| 5597 | root: Some(home.path().join("host")), |
| 5598 | ..ExtensionHostOptions::default() |
| 5599 | })); |
| 5600 | manager.bind_engine_handle(tokio::runtime::Handle::current()); |
| 5601 | let _manager = super::TestManagerGuard::install(Arc::clone(&manager)); |
| 5602 | let executor = crate::hooks::HookExecutor::new(admitted, workspace.clone()); |
| 5603 | let context = crate::hooks::HookContext::new() |
| 5604 | .with_session_id("actual-session") |
| 5605 | .with_tool_call_id("actual-call") |
| 5606 | .with_caller(crate::hooks::HookCaller { |
| 5607 | workspace, |
| 5608 | plugins: None, |
| 5609 | session_id: Some("actual-session".into()), |
| 5610 | agent_id: Some("actual-agent".into()), |
| 5611 | origin_turn_id: Some("actual-turn".into()), |
| 5612 | origin_call_id: Some("actual-call".into()), |
| 5613 | }); |
| 5614 | let env_scope = crate::test_support::env_scope_ticket(); |
| 5615 | tokio::task::spawn_blocking(move || { |
| 5616 | let _env = crate::test_support::join_env_scope(env_scope); |
| 5617 | let _policy = crate::plugins::activation::PolicyScope::propagate(true); |
| 5618 | for event in crate::hooks::config::ALL_HOOK_EVENTS { |
| 5619 | let results = executor.execute(event, &context); |
| 5620 | assert_eq!(results.len(), 1, "{event:?}"); |
| 5621 | assert!(results[0].success, "{event:?}: {:?}", results[0]); |
| 5622 | assert_eq!( |
| 5623 | results[0].stdout, |
| 5624 | format!("{}|actual-session|actual-call", event.as_str()), |
| 5625 | "{event:?}" |
| 5626 | ); |
| 5627 | } |
| 5628 | std::fs::write(hook_path, "# changed after admission\n").unwrap(); |
| 5629 | let rejected = executor.execute(crate::hooks::HookEvent::SessionStart, &context); |
| 5630 | assert_eq!(rejected.len(), 1); |
| 5631 | assert!(!rejected[0].success); |
| 5632 | assert!(rejected[0].stdout.is_empty()); |
| 5633 | assert!(rejected[0].exit_code.is_none()); |
| 5634 | }) |
| 5635 | .await |
| 5636 | .unwrap(); |
| 5637 | assert!(matches!( |
| 5638 | manager.tier_status(HostTier::Builtin), |
| 5639 | HostStatus::Ready { .. } |
| 5640 | )); |
| 5641 | assert!(!matches!(manager.status(), HostStatus::Ready { .. })); |
| 5642 | manager.shutdown().await; |
| 5643 | } |
| 5644 |