| 1 | //! Experimental TypeScript extension host — phase 1 (`[features] extension_host`). |
| 2 | //! |
| 3 | //! Codewhale's Rust core stays closed and authoritative: one turn loop, one |
| 4 | //! event authority, one store, one prompt authority, one approval gate. The |
| 5 | //! extension host (`crates/tui/extension-host`, a Bun or Node process running |
| 6 | //! the embedded bundle) is the only extensible surface, and in phase 1 it can do |
| 7 | //! exactly one thing: contribute tools, which become ordinary registry |
| 8 | //! `ToolSpec`s ([`tool::HostToolSpec`]) behind the existing gate. |
| 9 | //! |
| 10 | //! Lifecycle: a reviewed, enabled plugin with a `native` entry (activation |
| 11 | //! policy v4, selected by the flag) makes [`ExtensionHostManager::reconcile_in_background`] |
| 12 | //! spawn the host in the background — never on the first-prompt path — and |
| 13 | //! activate one owner per plugin. Tools join the per-turn registry at the next |
| 14 | //! rebuild, deferred. Disabling, revoking or updating the plugin revokes its |
| 15 | //! registrations synchronously before the host is asked to tear down. |
| 16 | //! |
| 17 | //! `core/call` ([`core_call`], tickets in [`ticket`]): an extension tool that |
| 18 | //! the model called directly may ask the core to run a core tool for it, |
| 19 | //! through the same gate as the model's own calls. The host decides nothing: |
| 20 | //! planning, approval and the card's text are Rust's, MCP and anything that |
| 21 | //! changes the session's authority are refused, and shell and network calls |
| 22 | //! force a prompt. |
| 23 | //! |
| 24 | //! Engines: the manager and its one host are process-wide, but each engine |
| 25 | //! holds its own [`HostAttachment`] carrying the plugin snapshot of its |
| 26 | //! workspace. Reconcile activates the union of what every attached snapshot |
| 27 | //! desires and revokes only owners no attachment desires. Each snapshot is |
| 28 | //! re-verified against persisted plugin state on every reconcile, so a |
| 29 | //! disable or revoke made through any registry revokes the plugin for every |
| 30 | //! engine. An engine installs only the tools of owners its own snapshot |
| 31 | //! desires, so a workspace never sees another workspace's project plugins. |
| 32 | //! Dropping an attachment detaches without revoking; the next reconcile |
| 33 | //! revokes whatever no remaining attachment desires. Engines without a |
| 34 | //! plugin snapshot of their own (isolated chats, the empty fallback) never |
| 35 | //! attach. |
| 36 | //! |
| 37 | //! Commands: a plugin may also contribute slash commands ([`command`]). They |
| 38 | //! are owned registrations like tools, loaded into the user command registry |
| 39 | //! at the lowest precedence, and run by `command/run` only when the user |
| 40 | //! invokes them. |
| 41 | //! |
| 42 | //! Plugin context: a plugin's `apply(ctx, config)` receives the settings from |
| 43 | //! `[plugins."<name>".config]` ([`plugin_config`]), and its tools and commands |
| 44 | //! are told, per call, the workspace the call comes from and the plugin's own |
| 45 | //! data directory. A plugin with several `native` entries activates them in |
| 46 | //! order under one owner. |
| 47 | //! |
| 48 | //! Known limitations (by design — see the design doc §8 and its "As built" |
| 49 | //! sections): |
| 50 | //! * Tools, slash commands, programmable pre-execute admission hooks, additive |
| 51 | //! prompt sections and owner-local storage. Skills and MCP are not host services yet. |
| 52 | //! The one thing the host may ask the core to do is a `core/call` from a tool |
| 53 | //! under the turn's gate; commands, timers and activation code ask for |
| 54 | //! nothing. |
| 55 | //! * Heartbeat and bounded automatic restart preserve the shared crash budget |
| 56 | //! across engine creation and replay. Dead-host calls fail with a typed |
| 57 | //! error and are never replayed. Three crashes in five minutes require an |
| 58 | //! explicit plugin change/reload to retry. Two dirty teardowns within ten |
| 59 | //! minutes retire the process once non-heartbeat calls are idle, without |
| 60 | //! resetting or consuming the unexpected-crash budget. |
| 61 | //! * Every awaited core→host request is bounded by its method's deadline |
| 62 | //! (`protocol::CoreRequest::deadline`) and then cancelled with `$/cancel` |
| 63 | //! (`supervisor::HostProcess::call`). Cancellation is a request: a plugin |
| 64 | //! that ignores its abort signal keeps running in the host until the host |
| 65 | //! is torn down, though the core has already failed the call. |
| 66 | //! * One host per engine process and trust tier ([`tier`]): the plugin tier |
| 67 | //! hosts reviewed third-party plugins, the builtin tier (tier 0) is for |
| 68 | //! Codewhale's own host code and starts only for a row of |
| 69 | //! [`tier::BUILTIN_MODULES`]. The MCP SDK backend starts its pinned module |
| 70 | //! only when explicitly selected; plugin attachment does not start it. |
| 71 | //! Under its OS sandbox |
| 72 | //! (Seatbelt on macOS; bubblewrap on Linux when a launch-time probe shows it |
| 73 | //! works) the host has no direct network, and cannot read the Codewhale |
| 74 | //! home (except the bundle, its data dir and plugin code), the Codex and |
| 75 | //! DSH credential homes, or the default credential stores |
| 76 | //! (`supervisor::plan_launch`, whose module docs list what bubblewrap does |
| 77 | //! not cover). Other files the user can read — including project `.env` |
| 78 | //! files — stay readable, and Mach services are not restricted. On Windows, |
| 79 | //! and on Linux where bwrap is missing or cannot start, Native launch is |
| 80 | //! refused. The pinned Builtin exception is explicitly diagnosed and still |
| 81 | //! ticket-bound. The flag remains Experimental. |
| 82 | //! * The runtime (`[extension_host] runtime`) defaults to Node. Bun is an |
| 83 | //! opt-in (`bun`, or `auto`, which prefers a Bun >= 1.4.0 and uses Node |
| 84 | //! when none is found) and is qualified on macOS only. The runtime is |
| 85 | //! pinned once a host on it completes the handshake; restarts reuse it |
| 86 | //! without probing it again, and the handshake refuses a host reporting a |
| 87 | //! different runtime or version. Under `auto`, a Bun that fails to launch or |
| 88 | //! handshake before anything is pinned is reported once and Node is used |
| 89 | //! for the rest of the session; an explicit `bun`/`node` never falls back. |
| 90 | //! The 1 GiB memory cap is applied through `RLIMIT_DATA` on Linux, the Job |
| 91 | //! Object on Windows and, for a Bun host on macOS, a jetsam limit the host |
| 92 | //! applies to itself; a macOS Node host is checked at each heartbeat |
| 93 | //! instead. The Rust host tests have not run a Bun host on Linux or |
| 94 | //! Windows (`supervisor` module docs say what was measured where). |
| 95 | //! * In-process native code is taken away from plugins by the host |
| 96 | //! (`extension-host/src/runtime.ts`), for the entry points found so far; a |
| 97 | //! native-code entry point a newer runtime adds is not covered until it is |
| 98 | //! added there. A process a plugin starts is outside that policy and runs |
| 99 | //! under the same verified OS sandbox required for Native admission. |
| 100 | //! * The owner token is a bug/staleness guard, not a boundary between |
| 101 | //! plugins that share the process: one plugin can alter another's |
| 102 | //! behaviour, which the approval card discloses. |
| 103 | //! * Extension tool names that any name-keyed approval table special-cases |
| 104 | //! are refused (`registry::core_special_case`), so an extension tool never |
| 105 | //! shares an approval key, summary or category with a built-in. |
| 106 | //! * The host's process tree (Unix process group / Windows Job Object, |
| 107 | //! shared with hooks via `crate::process_tree`) is killed as a whole. On |
| 108 | //! Windows the host is assigned to its job just after spawn, not created |
| 109 | //! suspended as hooks are. On Unix a plugin child that calls `setsid` leaves |
| 110 | //! the group and is not killed with it. When the core goes away, the host |
| 111 | //! kills its own group at stdin EOF, and a watchdog thread does the same |
| 112 | //! when its parent process changes, even if a plugin blocks the event loop. |
| 113 | //! * A running engine's snapshot is replaced only by its own workspace |
| 114 | //! switch or by [`plugins_changed`] for the same workspace. A plugin newly |
| 115 | //! enabled through another workspace's registry reaches an engine at its |
| 116 | //! next snapshot, not at once; disables and revokes always reach it at the |
| 117 | //! next reconcile, through the persisted-state check. |
| 118 | //! * The host re-hashes each `native` entry file before importing it; other |
| 119 | //! files in the staged snapshot are covered by Rust's per-call receipt |
| 120 | //! check, not re-hashed by the host. |
| 121 | |
| 122 | pub(crate) mod command; |
| 123 | pub(crate) mod composition_review; |
| 124 | pub mod composition_scope; |
| 125 | pub(crate) mod core_call; |
| 126 | pub(crate) mod execution; |
| 127 | pub(crate) use execution::StockOperation; |
| 128 | mod hooks; |
| 129 | pub(crate) mod mcp; |
| 130 | pub(crate) mod native_mcp; |
| 131 | pub(crate) mod plugin_config; |
| 132 | pub(crate) mod prompt; |
| 133 | pub(crate) mod protocol; |
| 134 | pub(crate) mod registry; |
| 135 | pub(crate) mod skills; |
| 136 | pub(crate) mod supervisor; |
| 137 | pub(crate) mod ticket; |
| 138 | pub(crate) mod tier; |
| 139 | pub(crate) mod tool; |
| 140 | #[cfg(windows)] |
| 141 | mod windows; |
| 142 | |
| 143 | #[cfg(test)] |
| 144 | mod core_call_tests; |
| 145 | #[cfg(test)] |
| 146 | pub(crate) mod tests; |
| 147 | |
| 148 | use std::collections::{BTreeMap, BTreeSet, HashSet, VecDeque}; |
| 149 | use std::fmt; |
| 150 | use std::path::{Path, PathBuf}; |
| 151 | use std::sync::atomic::{AtomicU64, Ordering}; |
| 152 | use std::sync::{Arc, Mutex, OnceLock, Weak}; |
| 153 | use std::time::{Duration, Instant}; |
| 154 | |
| 155 | use async_trait::async_trait; |
| 156 | use sha2::{Digest, Sha256}; |
| 157 | use tokio_util::sync::CancellationToken; |
| 158 | |
| 159 | use self::protocol::{ |
| 160 | ActivateParams, ActivateResult, CoreRequest, DeactivateParams, DeactivateResult, EntryRef, |
| 161 | OwnerRef, RegisterKind, RegisterResult, |
| 162 | }; |
| 163 | use self::registry::{CommandRegistration, OwnerRegistry, OwnerState, ToolRegistration}; |
| 164 | use self::supervisor::{HostEvents, HostProcess, HostRequestContext}; |
| 165 | use self::tier::{BuiltinModule, HostTier}; |
| 166 | use crate::plugins::PluginRegistry; |
| 167 | use crate::plugins::activation::{self, PluginActivationCapability}; |
| 168 | use crate::plugins::types::PluginAuthority; |
| 169 | |
| 170 | /// The host bundle, embedded so every distribution channel carries it. |
| 171 | const BUNDLE: &[u8] = include_bytes!("../../extension-host/dist/codewhale-extension-host.mjs"); |
| 172 | const BUNDLE_FILE_NAME: &str = "codewhale-extension-host.mjs"; |
| 173 | /// The licence notices of the third-party code inside [`BUNDLE`], generated by |
| 174 | /// the host build from the bundler's metafile (`extension-host/build.mjs`). The |
| 175 | /// bundle's banner promises "LICENSES.txt next to this file", so every |
| 176 | /// materialised copy of the bundle carries one. |
| 177 | const NOTICES: &[u8] = include_bytes!("../../extension-host/dist/LICENSES.txt"); |
| 178 | const NOTICES_FILE_NAME: &str = "LICENSES.txt"; |
| 179 | const MAX_DIAGNOSTICS: usize = 64; |
| 180 | const MAX_DIAGNOSTIC_BYTES: usize = 2048; |
| 181 | const DIRTY_RESTART_REASON: &str = "planned restart after repeated dirty teardowns"; |
| 182 | |
| 183 | struct Diagnostic { |
| 184 | plugin_id: Option<String>, |
| 185 | message: String, |
| 186 | } |
| 187 | |
| 188 | fn bounded_diagnostic(mut message: String) -> String { |
| 189 | if message.len() > MAX_DIAGNOSTIC_BYTES { |
| 190 | let mut end = MAX_DIAGNOSTIC_BYTES; |
| 191 | while !message.is_char_boundary(end) { |
| 192 | end -= 1; |
| 193 | } |
| 194 | message.truncate(end); |
| 195 | message.push('…'); |
| 196 | } |
| 197 | message |
| 198 | } |
| 199 | |
| 200 | fn hex(bytes: impl AsRef<[u8]>) -> String { |
| 201 | bytes |
| 202 | .as_ref() |
| 203 | .iter() |
| 204 | .map(|byte| format!("{byte:02x}")) |
| 205 | .collect() |
| 206 | } |
| 207 | |
| 208 | /// SHA-256 of the embedded bundle. |
| 209 | #[must_use] |
| 210 | pub fn bundle_sha256() -> &'static str { |
| 211 | static DIGEST: OnceLock<String> = OnceLock::new(); |
| 212 | DIGEST.get_or_init(|| hex(Sha256::digest(BUNDLE))) |
| 213 | } |
| 214 | |
| 215 | /// Write the embedded bundle, and the licence notices that belong with it, to |
| 216 | /// `<root>/extension-host/<sha256>/` unless identical copies are already there, |
| 217 | /// then re-verify the bytes on disk. A file is never overwritten in place: a |
| 218 | /// different build writes a different directory. Blocking. |
| 219 | fn materialize_bundle(root: &Path) -> Result<PathBuf, String> { |
| 220 | let dir = supervisor::bundle_dir(root, bundle_sha256()); |
| 221 | let bundle = materialize_file(&dir, BUNDLE_FILE_NAME, BUNDLE)?; |
| 222 | materialize_file(&dir, NOTICES_FILE_NAME, NOTICES)?; |
| 223 | materialize_file( |
| 224 | &dir.join("builtin"), |
| 225 | "mcp.mjs", |
| 226 | include_bytes!("../../extension-host/dist/builtin/mcp.mjs"), |
| 227 | )?; |
| 228 | materialize_file( |
| 229 | &dir.join("builtin"), |
| 230 | "harness.mjs", |
| 231 | include_bytes!("../../extension-host/dist/builtin/harness.mjs"), |
| 232 | )?; |
| 233 | Ok(bundle) |
| 234 | } |
| 235 | |
| 236 | /// One embedded file in `dir`: written to a private staging name, made |
| 237 | /// read-only, renamed into place (so a reader never sees a partial file and a |
| 238 | /// symlink at the final name is replaced, not followed), and read back. |
| 239 | fn materialize_file(dir: &Path, file_name: &str, bytes: &[u8]) -> Result<PathBuf, String> { |
| 240 | let path = dir.join(file_name); |
| 241 | let matches = |path: &Path| std::fs::read(path).is_ok_and(|on_disk| on_disk == bytes); |
| 242 | if !matches(&path) { |
| 243 | std::fs::create_dir_all(dir) |
| 244 | .map_err(|error| format!("cannot create {}: {error}", dir.display()))?; |
| 245 | let staging = dir.join(format!(".{file_name}.{}", uuid::Uuid::new_v4().simple())); |
| 246 | std::fs::write(&staging, bytes) |
| 247 | .map_err(|error| format!("cannot write {}: {error}", staging.display()))?; |
| 248 | #[cfg(unix)] |
| 249 | { |
| 250 | use std::os::unix::fs::PermissionsExt; |
| 251 | let _ = std::fs::set_permissions(&staging, std::fs::Permissions::from_mode(0o400)); |
| 252 | } |
| 253 | if let Err(error) = std::fs::rename(&staging, &path) { |
| 254 | let _ = std::fs::remove_file(&staging); |
| 255 | if !matches(&path) { |
| 256 | return Err(format!("cannot publish {}: {error}", path.display())); |
| 257 | } |
| 258 | } |
| 259 | } |
| 260 | // Re-check what will actually be used. |
| 261 | if !matches(&path) { |
| 262 | return Err(format!( |
| 263 | "{} does not match the embedded copy (sha256 {})", |
| 264 | path.display(), |
| 265 | hex(Sha256::digest(bytes)) |
| 266 | )); |
| 267 | } |
| 268 | Ok(path) |
| 269 | } |
| 270 | |
| 271 | /// Options fixed for the life of one manager. |
| 272 | #[derive(Debug, Clone, Default)] |
| 273 | pub struct ExtensionHostOptions { |
| 274 | /// `[extension_host] runtime`: `node` (the default), `bun`, or `auto` |
| 275 | /// (Bun first). |
| 276 | pub runtime: crate::config::ExtensionHostRuntime, |
| 277 | /// `[extension_host] node`: tried before every `node` on `PATH`. |
| 278 | pub node_override: Option<PathBuf>, |
| 279 | /// `[extension_host] bun`: tried before every `bun` on `PATH`. |
| 280 | pub bun_override: Option<PathBuf>, |
| 281 | /// Where the bundle is materialized; defaults to the Codewhale home. |
| 282 | pub root: Option<PathBuf>, |
| 283 | /// Per-manager timings; tests can shorten them without global state. |
| 284 | pub supervision: SupervisionOptions, |
| 285 | } |
| 286 | |
| 287 | #[derive(Debug, Clone)] |
| 288 | pub struct SupervisionOptions { |
| 289 | pub heartbeat_interval: Duration, |
| 290 | pub ping_timeout: Duration, |
| 291 | pub hang_timeout: Duration, |
| 292 | pub restart_backoff: Duration, |
| 293 | pub crash_window: Duration, |
| 294 | pub crash_limit: usize, |
| 295 | pub start_retry_cooldown: Duration, |
| 296 | pub dirty_window: Duration, |
| 297 | pub dirty_limit: usize, |
| 298 | /// Host memory cap in bytes (`supervisor::HOST_MEMORY_CAP`). |
| 299 | pub memory_cap: u64, |
| 300 | /// How long one extension tool call may take before the core cancels it |
| 301 | /// (`tool::TOOL_CALL_DEADLINE`); sent to the host as `deadline_ms`. The |
| 302 | /// other methods' deadlines are fixed (`CoreRequest::deadline`). |
| 303 | pub tool_call_deadline: Duration, |
| 304 | /// How long the user waits for one extension command before the core |
| 305 | /// cancels it (`command::COMMAND_RUN_DEADLINE`); sent to the host as |
| 306 | /// `deadline_ms`. |
| 307 | pub command_run_deadline: Duration, |
| 308 | } |
| 309 | |
| 310 | impl ExtensionHostOptions { |
| 311 | /// Options for the `[extension_host]` table (paths `~`-expanded). |
| 312 | #[must_use] |
| 313 | pub fn from_config(table: Option<&crate::config::ExtensionHostConfig>) -> Self { |
| 314 | let expand = |path: Option<&String>| { |
| 315 | path.map(|path| PathBuf::from(shellexpand::tilde(path).as_ref())) |
| 316 | }; |
| 317 | Self { |
| 318 | runtime: table.map_or_else(Default::default, |table| table.effective_runtime()), |
| 319 | node_override: expand(table.and_then(|table| table.node.as_ref())), |
| 320 | bun_override: expand(table.and_then(|table| table.bun.as_ref())), |
| 321 | ..Default::default() |
| 322 | } |
| 323 | } |
| 324 | } |
| 325 | |
| 326 | impl Default for SupervisionOptions { |
| 327 | fn default() -> Self { |
| 328 | Self { |
| 329 | heartbeat_interval: Duration::from_secs(3), |
| 330 | ping_timeout: Duration::from_secs(3), |
| 331 | hang_timeout: supervisor::PING_DEADLINE, |
| 332 | restart_backoff: Duration::from_millis(250), |
| 333 | crash_window: Duration::from_secs(5 * 60), |
| 334 | crash_limit: 3, |
| 335 | start_retry_cooldown: Duration::from_secs(60), |
| 336 | dirty_window: Duration::from_secs(10 * 60), |
| 337 | dirty_limit: 2, |
| 338 | memory_cap: supervisor::HOST_MEMORY_CAP, |
| 339 | tool_call_deadline: tool::TOOL_CALL_DEADLINE, |
| 340 | command_run_deadline: command::COMMAND_RUN_DEADLINE, |
| 341 | } |
| 342 | } |
| 343 | } |
| 344 | |
| 345 | #[derive(Default)] |
| 346 | struct SupervisionState { |
| 347 | crashes: VecDeque<Instant>, |
| 348 | last_start: Option<Instant>, |
| 349 | launch_failed: bool, |
| 350 | retry_ticket: u64, |
| 351 | policy: bool, |
| 352 | dirty_teardowns: VecDeque<Instant>, |
| 353 | dirty_restart_pending: bool, |
| 354 | planned_restart: Option<u64>, |
| 355 | } |
| 356 | |
| 357 | impl SupervisionState { |
| 358 | fn record_dirty_teardown(&mut self, now: Instant, options: &SupervisionOptions) { |
| 359 | while self |
| 360 | .dirty_teardowns |
| 361 | .front() |
| 362 | .is_some_and(|at| now.saturating_duration_since(*at) >= options.dirty_window) |
| 363 | { |
| 364 | self.dirty_teardowns.pop_front(); |
| 365 | } |
| 366 | self.dirty_teardowns.push_back(now); |
| 367 | let limit = options.dirty_limit.max(1); |
| 368 | self.dirty_restart_pending |= self.dirty_teardowns.len() >= limit; |
| 369 | while self.dirty_teardowns.len() > limit { |
| 370 | self.dirty_teardowns.pop_front(); |
| 371 | } |
| 372 | } |
| 373 | |
| 374 | fn record_crash(&mut self, now: Instant, options: &SupervisionOptions) -> bool { |
| 375 | while self |
| 376 | .crashes |
| 377 | .front() |
| 378 | .is_some_and(|at| now.saturating_duration_since(*at) >= options.crash_window) |
| 379 | { |
| 380 | self.crashes.pop_front(); |
| 381 | } |
| 382 | self.crashes.push_back(now); |
| 383 | self.launch_failed = false; |
| 384 | self.retry_ticket += 1; |
| 385 | self.crashes.len() < options.crash_limit |
| 386 | } |
| 387 | } |
| 388 | |
| 389 | /// Observable host state: the one source for `/plugin` and for the error a |
| 390 | /// tool call routed to the host gets while it is down ([`fmt::Display`]). |
| 391 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 392 | pub enum HostStatus { |
| 393 | /// `[features] extension_host` is off for this process. |
| 394 | Disabled, |
| 395 | /// Not started: nothing has needed it yet. |
| 396 | Idle, |
| 397 | Starting, |
| 398 | Ready { |
| 399 | pid: Option<u32>, |
| 400 | /// `bun` or `node`. |
| 401 | runtime: &'static str, |
| 402 | /// As reported by the host (`host/hello`). |
| 403 | runtime_version: String, |
| 404 | sandbox: supervisor::HostSandbox, |
| 405 | /// How the memory cap is enforced for this process. |
| 406 | memory: supervisor::MemoryEnforcement, |
| 407 | }, |
| 408 | Unresponsive { |
| 409 | pid: Option<u32>, |
| 410 | }, |
| 411 | /// Crashed (or retired for maintenance); the supervisor restarts it. |
| 412 | Restarting { |
| 413 | reason: String, |
| 414 | }, |
| 415 | /// Start refused (`start failed: …`) or crash budget exhausted; only an |
| 416 | /// explicit plugin change or reload retries. |
| 417 | Failed { |
| 418 | reason: String, |
| 419 | stderr_tail: String, |
| 420 | }, |
| 421 | } |
| 422 | |
| 423 | /// Why the host can or cannot take a call, in one phrase. |
| 424 | impl fmt::Display for HostStatus { |
| 425 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 426 | match self { |
| 427 | Self::Disabled => f.write_str("disabled by config ([features] extension_host is off)"), |
| 428 | Self::Idle => f.write_str( |
| 429 | "not started (it starts in the background when a reviewed plugin with host code is enabled)", |
| 430 | ), |
| 431 | Self::Starting => f.write_str("starting"), |
| 432 | Self::Ready { .. } => f.write_str("running"), |
| 433 | Self::Unresponsive { pid } => write!( |
| 434 | f, |
| 435 | "unresponsive (pid {} missed its heartbeat; the supervisor restarts it if it stays silent)", |
| 436 | pid.map_or_else(|| "?".to_string(), |pid| pid.to_string()) |
| 437 | ), |
| 438 | Self::Restarting { reason } => write!(f, "restarting after: {reason}"), |
| 439 | Self::Failed { reason, .. } => { |
| 440 | write!(f, "{reason} (change or reload a plugin to retry)") |
| 441 | } |
| 442 | } |
| 443 | } |
| 444 | } |
| 445 | |
| 446 | enum HostSlot { |
| 447 | Idle, |
| 448 | Starting, |
| 449 | Ready(Arc<HostProcess>), |
| 450 | Unresponsive(Arc<HostProcess>), |
| 451 | Restarting { reason: String }, |
| 452 | Failed { reason: String, stderr_tail: String }, |
| 453 | } |
| 454 | |
| 455 | impl HostSlot { |
| 456 | fn status(&self) -> HostStatus { |
| 457 | match self { |
| 458 | Self::Idle => HostStatus::Idle, |
| 459 | Self::Starting => HostStatus::Starting, |
| 460 | Self::Ready(host) if host.is_retiring() => HostStatus::Restarting { |
| 461 | reason: DIRTY_RESTART_REASON.to_string(), |
| 462 | }, |
| 463 | // The exit watcher reports the exit a moment later. |
| 464 | Self::Ready(host) | Self::Unresponsive(host) if host.has_exited() => { |
| 465 | HostStatus::Restarting { |
| 466 | reason: "the host process exited".to_string(), |
| 467 | } |
| 468 | } |
| 469 | Self::Ready(host) => HostStatus::Ready { |
| 470 | pid: host.pid, |
| 471 | runtime: host.runtime.kind.name(), |
| 472 | runtime_version: host.runtime_version.get().cloned().unwrap_or_default(), |
| 473 | sandbox: host.sandbox.clone(), |
| 474 | memory: host.memory(), |
| 475 | }, |
| 476 | Self::Unresponsive(host) => HostStatus::Unresponsive { pid: host.pid }, |
| 477 | Self::Restarting { reason } => HostStatus::Restarting { |
| 478 | reason: reason.clone(), |
| 479 | }, |
| 480 | Self::Failed { |
| 481 | reason, |
| 482 | stderr_tail, |
| 483 | } => HostStatus::Failed { |
| 484 | reason: reason.clone(), |
| 485 | stderr_tail: stderr_tail.clone(), |
| 486 | }, |
| 487 | } |
| 488 | } |
| 489 | } |
| 490 | |
| 491 | /// One tier's host process, and everything that goes with supervising it: the |
| 492 | /// slot, the generation that makes old callbacks harmless, the spawn count and |
| 493 | /// the crash/restart bookkeeping. The two tiers each own one, so a crash, |
| 494 | /// restart or failed launch of one never touches the other. The owner |
| 495 | /// registry, the runtime pin and the attachments stay shared. |
| 496 | struct TierRuntime { |
| 497 | host: Mutex<HostSlot>, |
| 498 | host_generation: AtomicU64, |
| 499 | spawn_attempts: AtomicU64, |
| 500 | /// Lock order: see [`ManagerShared`]. |
| 501 | supervision: Mutex<SupervisionState>, |
| 502 | } |
| 503 | |
| 504 | impl TierRuntime { |
| 505 | fn new() -> Self { |
| 506 | Self { |
| 507 | host: Mutex::new(HostSlot::Idle), |
| 508 | host_generation: AtomicU64::new(0), |
| 509 | spawn_attempts: AtomicU64::new(0), |
| 510 | supervision: Mutex::new(SupervisionState::default()), |
| 511 | } |
| 512 | } |
| 513 | } |
| 514 | |
| 515 | struct DesiredOwner { |
| 516 | plugin_name: String, |
| 517 | authority: PluginAuthority, |
| 518 | entries: Vec<(PathBuf, String)>, |
| 519 | } |
| 520 | |
| 521 | /// What [`ExtensionHostManager::activate_owner`] activates, on either tier: |
| 522 | /// a plugin's reviewed bytes, or a built-in module's pinned source. |
| 523 | struct Activation { |
| 524 | tier: HostTier, |
| 525 | owner_id: String, |
| 526 | /// The plugin's manifest name, or the module's name. |
| 527 | name: String, |
| 528 | /// The reviewed plugin authority; `None` for a built-in module. |
| 529 | authority: Option<PluginAuthority>, |
| 530 | /// What the owner is bound to: the plugin's content hash, or the module's |
| 531 | /// pinned source digest. |
| 532 | content_hash: String, |
| 533 | entries: Vec<(PathBuf, String)>, |
| 534 | } |
| 535 | |
| 536 | /// One engine's view: its plugin snapshot and the owners (plugin id → |
| 537 | /// reviewed content hash) that snapshot desired at the last reconcile. |
| 538 | struct AttachmentState { |
| 539 | plugins: Arc<PluginRegistry>, |
| 540 | desired: BTreeMap<String, String>, |
| 541 | selection: composition_scope::CompositionSelection, |
| 542 | revision: u64, |
| 543 | selection_cancel: CancellationToken, |
| 544 | session_id: Option<String>, |
| 545 | agent_id: Option<String>, |
| 546 | } |
| 547 | |
| 548 | /// Lock order, outermost first, everywhere in this module: |
| 549 | /// |
| 550 | /// 1. `sync_lock` (async; the only lock held across an await, and it only |
| 551 | /// serializes reconciliation); |
| 552 | /// 2. one tier's `host` slot, then that same tier's `supervision`; |
| 553 | /// 3. `registry`. |
| 554 | /// |
| 555 | /// No code holds both tiers' slots or both tiers' supervision state at once: |
| 556 | /// a path that must visit both tiers visits one, releases it, then the other. |
| 557 | /// `attachments`, `runtime`, `plugin_configs` and `diagnostics` are leaf |
| 558 | /// locks, never held together with another lock (`diagnostics` is only ever |
| 559 | /// taken last, from the `diagnostic` helpers). No lock but `sync_lock` is held |
| 560 | /// across an await. |
| 561 | pub(crate) struct ManagerShared { |
| 562 | options: ExtensionHostOptions, |
| 563 | /// The built-in host modules this manager may start the builtin tier for: |
| 564 | /// [`tier::BUILTIN_MODULES`] in production; a test substitutes its |
| 565 | /// own through [`ExtensionHostManager::with_builtin_modules`]. |
| 566 | builtin_modules: &'static [BuiltinModule], |
| 567 | attachments: Mutex<BTreeMap<u64, AttachmentState>>, |
| 568 | next_attachment: AtomicU64, |
| 569 | /// One registry for both tiers: an owner's entry records which. |
| 570 | registry: Mutex<OwnerRegistry>, |
| 571 | /// Capability tickets and the invocations they belong to (`core/call`). |
| 572 | /// A leaf lock domain: its mutexes are never held with another lock. |
| 573 | core_calls: Arc<core_call::CoreCalls>, |
| 574 | engine_handle: OnceLock<tokio::runtime::Handle>, |
| 575 | execution_broker: execution::Broker, |
| 576 | harness_users: AtomicU64, |
| 577 | stock_users: AtomicU64, |
| 578 | mcp_broker: mcp::Broker, |
| 579 | mcp_users: AtomicU64, |
| 580 | /// Tier 1: reviewed third-party plugins. |
| 581 | plugin: TierRuntime, |
| 582 | /// Tier 0: Codewhale's own host code. Starts only for a built-in module. |
| 583 | builtin: TierRuntime, |
| 584 | sync_lock: tokio::sync::Mutex<()>, |
| 585 | /// Bound native root parsing independently of reconciliation (activation |
| 586 | /// itself calls admission while reconciliation owns sync_lock). |
| 587 | skill_admission: Arc<tokio::sync::Semaphore>, |
| 588 | diagnostics: Mutex<VecDeque<Diagnostic>>, |
| 589 | /// Which runtime every host (of either tier) runs on. |
| 590 | runtime: Mutex<RuntimePin>, |
| 591 | /// User settings per plugin name. |
| 592 | plugin_configs: Mutex<plugin_config::PluginConfigs>, |
| 593 | } |
| 594 | |
| 595 | /// Which runtime this manager's hosts run on. |
| 596 | #[derive(Default)] |
| 597 | struct RuntimePin { |
| 598 | /// Set by the first host that completes the handshake and reused by |
| 599 | /// every restart, so a session never switches runtime (a `bun` install |
| 600 | /// or removal mid-session changes nothing until the next process; a |
| 601 | /// binary replaced at the pinned path is refused at the handshake). With |
| 602 | /// the one-line summary of why it was chosen, for `/plugin`. |
| 603 | pinned: Option<(crate::dependencies::HostRuntime, String)>, |
| 604 | /// `runtime = "auto"` only: a Bun host failed to launch or handshake |
| 605 | /// before anything was pinned, so every later launch this session |
| 606 | /// resolves Node. |
| 607 | bun_failed: bool, |
| 608 | } |
| 609 | |
| 610 | /// Resolve the runtime unless one is pinned, materialize the bundle and plan |
| 611 | /// the launch. Returns the summary to pin once a host on this runtime |
| 612 | /// completes the handshake, or `None` when the runtime was already pinned. |
| 613 | /// Blocking. |
| 614 | fn prepare_launch( |
| 615 | options: &ExtensionHostOptions, |
| 616 | tier: HostTier, |
| 617 | pinned: Option<(crate::dependencies::HostRuntime, String)>, |
| 618 | bun_failed: bool, |
| 619 | ) -> Result<(supervisor::HostLaunch, Option<String>), String> { |
| 620 | let (runtime, summary) = match pinned { |
| 621 | Some((runtime, _)) => (runtime, None), |
| 622 | None => { |
| 623 | let choice = if bun_failed { |
| 624 | crate::config::ExtensionHostRuntime::Node |
| 625 | } else { |
| 626 | options.runtime |
| 627 | }; |
| 628 | let mut resolution = crate::dependencies::resolve_extension_host_runtime( |
| 629 | choice, |
| 630 | options.node_override.as_deref(), |
| 631 | options.bun_override.as_deref(), |
| 632 | ); |
| 633 | // Report the configured choice, not the Node-only fallback. |
| 634 | resolution.choice = options.runtime; |
| 635 | let note = if bun_failed { |
| 636 | "; Bun failed to start this session (see diagnostics)" |
| 637 | } else { |
| 638 | "" |
| 639 | }; |
| 640 | let Some(runtime) = resolution.selected.clone() else { |
| 641 | return Err(format!("{}{note}", resolution.failure())); |
| 642 | }; |
| 643 | (runtime, Some(format!("{}{note}", resolution.summary()))) |
| 644 | } |
| 645 | }; |
| 646 | let root = host_root(options)?; |
| 647 | let bundle = materialize_bundle(&root)?; |
| 648 | let launch = supervisor::plan_launch( |
| 649 | tier, |
| 650 | &runtime, |
| 651 | &bundle, |
| 652 | &root, |
| 653 | options.supervision.memory_cap, |
| 654 | )?; |
| 655 | Ok((launch, summary)) |
| 656 | } |
| 657 | |
| 658 | /// The Codewhale home the host lives under. |
| 659 | fn host_root(options: &ExtensionHostOptions) -> Result<PathBuf, String> { |
| 660 | match options.root.clone() { |
| 661 | Some(root) => Ok(root), |
| 662 | None => codewhale_config::codewhale_home() |
| 663 | .map_err(|error| format!("Codewhale home unavailable: {error}")), |
| 664 | } |
| 665 | } |
| 666 | |
| 667 | /// The selected tier's sandbox, planned exactly as a launch does (for doctor). |
| 668 | /// Blocking; creates the host's data dir and probes bubblewrap on Linux. |
| 669 | pub(crate) fn planned_sandbox( |
| 670 | options: &ExtensionHostOptions, |
| 671 | runtime: &crate::dependencies::HostRuntime, |
| 672 | tier: HostTier, |
| 673 | ) -> Result<supervisor::HostSandbox, String> { |
| 674 | supervisor::planned_sandbox(tier, runtime, &host_root(options)?) |
| 675 | } |
| 676 | |
| 677 | impl ManagerShared { |
| 678 | pub(crate) fn engine_handle(&self) -> Result<tokio::runtime::Handle, String> { |
| 679 | self.engine_handle.get().cloned().ok_or_else(|| "enabled execution requires the Engine scheduler; no new runtime or legacy fallback is allowed".to_string()) |
| 680 | } |
| 681 | |
| 682 | /// The supervision state of `tier`'s host. |
| 683 | fn tier_runtime(&self, tier: HostTier) -> &TierRuntime { |
| 684 | match tier { |
| 685 | HostTier::Plugin => &self.plugin, |
| 686 | HostTier::Builtin => &self.builtin, |
| 687 | } |
| 688 | } |
| 689 | |
| 690 | async fn deactivate_owner(&self, host: &Arc<HostProcess>, owner: &OwnerRef) { |
| 691 | self.deactivate_entry(host, owner, None).await; |
| 692 | } |
| 693 | /// A failed or retired Native entry must produce the same bounded teardown |
| 694 | /// receipt as a whole owner, while its other selected entries stay live. |
| 695 | async fn deactivate_entry( |
| 696 | &self, |
| 697 | host: &Arc<HostProcess>, |
| 698 | owner: &OwnerRef, |
| 699 | entry: Option<EntryRef>, |
| 700 | ) { |
| 701 | let diagnostic = match host |
| 702 | .call( |
| 703 | CoreRequest::Deactivate(DeactivateParams { |
| 704 | entry, |
| 705 | owner: owner.clone(), |
| 706 | }), |
| 707 | None, |
| 708 | ) |
| 709 | .await |
| 710 | .map(serde_json::from_value::<DeactivateResult>) |
| 711 | { |
| 712 | Ok(Ok(ack)) if ack.disposed && ack.leaked.is_empty() => return, |
| 713 | Ok(Ok(ack)) => format!( |
| 714 | "extension `{}` teardown incomplete (disposed: {}, leaked: {:?})", |
| 715 | owner.plugin_id, ack.disposed, ack.leaked |
| 716 | ), |
| 717 | Ok(Err(error)) => format!( |
| 718 | "extension `{}` teardown answer malformed: {error}", |
| 719 | owner.plugin_id |
| 720 | ), |
| 721 | Err(error) => format!("extension `{}` teardown failed: {error}", owner.plugin_id), |
| 722 | }; |
| 723 | self.plugin_diagnostic(&owner.plugin_id, diagnostic); |
| 724 | self.record_dirty_teardown(host); |
| 725 | } |
| 726 | |
| 727 | fn record_dirty_teardown(&self, host: &Arc<HostProcess>) { |
| 728 | let runtime = self.tier_runtime(host.tier); |
| 729 | let slot = runtime.host.lock().expect("host lock"); |
| 730 | if !matches!(&*slot, HostSlot::Ready(current) | HostSlot::Unresponsive(current) |
| 731 | if Arc::ptr_eq(current, host)) |
| 732 | { |
| 733 | return; |
| 734 | } |
| 735 | runtime |
| 736 | .supervision |
| 737 | .lock() |
| 738 | .expect("supervision lock") |
| 739 | .record_dirty_teardown(Instant::now(), &self.options.supervision); |
| 740 | } |
| 741 | |
| 742 | fn diagnostic(&self, message: String) { |
| 743 | self.record_diagnostic(None, message); |
| 744 | } |
| 745 | |
| 746 | fn plugin_diagnostic(&self, plugin_id: &str, message: String) { |
| 747 | // Refused registrations can contain an arbitrary host-supplied id. |
| 748 | // Keep an oversized id global instead of retaining unbounded metadata |
| 749 | // or truncating it into another plugin's identity. |
| 750 | let plugin_id = (plugin_id.len() <= MAX_DIAGNOSTIC_BYTES).then(|| plugin_id.to_string()); |
| 751 | self.record_diagnostic(plugin_id, message); |
| 752 | } |
| 753 | |
| 754 | fn record_diagnostic(&self, plugin_id: Option<String>, message: String) { |
| 755 | let message = bounded_diagnostic(message); |
| 756 | tracing::info!(target: "extension_host", "{message}"); |
| 757 | let mut diagnostics = self.diagnostics.lock().expect("diagnostics lock"); |
| 758 | diagnostics.push_back(Diagnostic { plugin_id, message }); |
| 759 | while diagnostics.len() > MAX_DIAGNOSTICS { |
| 760 | diagnostics.pop_front(); |
| 761 | } |
| 762 | } |
| 763 | |
| 764 | /// `tier`'s host, when it can take a call; otherwise why not. |
| 765 | fn ready_host(&self, tier: HostTier) -> Result<Arc<HostProcess>, HostStatus> { |
| 766 | match &*self.tier_runtime(tier).host.lock().expect("host lock") { |
| 767 | HostSlot::Ready(host) if !host.has_exited() && !host.is_retiring() => { |
| 768 | Ok(Arc::clone(host)) |
| 769 | } |
| 770 | slot => Err(slot.status()), |
| 771 | } |
| 772 | } |
| 773 | |
| 774 | /// The shared fast liveness guard: current policy, running host, exact |
| 775 | /// owner generation/tier and its reviewed authority. `live_host` follows |
| 776 | /// with Native receipt validation and another host check; root admission |
| 777 | /// uses the same authority inside its one bounded blocking disk job. |
| 778 | /// |
| 779 | /// `owner_if_live` runs under the registry lock and names the owner |
| 780 | /// generation the call belongs to, or why it is no longer registered. |
| 781 | fn live_owner_authority( |
| 782 | &self, |
| 783 | tier: HostTier, |
| 784 | owner_if_live: impl FnOnce(&OwnerRegistry) -> Result<OwnerRef, String>, |
| 785 | ) -> Result<Option<PluginAuthority>, String> { |
| 786 | let native_policy = activation::extension_host_policy_enabled(); |
| 787 | if tier == HostTier::Plugin && !native_policy { |
| 788 | return Err(host_down(tier, &HostStatus::Disabled)); |
| 789 | } |
| 790 | self.ready_host(tier) |
| 791 | .map_err(|status| host_down(tier, &status))?; |
| 792 | { |
| 793 | let registry = self.registry.lock().expect("registry lock"); |
| 794 | let owner = owner_if_live(®istry)?; |
| 795 | // Independent Core demand admits only the exact pinned module. |
| 796 | // Script/hook jobs still require Native policy in Caller::check. |
| 797 | let core_stock = |
| 798 | owner.plugin_id == "host:harness" && self.stock_users.load(Ordering::SeqCst) > 0; |
| 799 | if !native_policy && owner.plugin_id != "host:mcp" && !core_stock { |
| 800 | return Err(host_down(tier, &HostStatus::Disabled)); |
| 801 | } |
| 802 | let owner_tier = registry |
| 803 | .tier_of(&owner) |
| 804 | .ok_or_else(|| "extension owner has no authority".to_string())?; |
| 805 | if owner_tier != tier { |
| 806 | return Err(format!( |
| 807 | "extension owner `{}` belongs to the {} tier, not the {} tier", |
| 808 | owner.plugin_id, |
| 809 | owner_tier.name(), |
| 810 | tier.name() |
| 811 | )); |
| 812 | } |
| 813 | match (tier, registry.authority_for(&owner)) { |
| 814 | (HostTier::Plugin, Some(authority)) => Ok(Some(authority)), |
| 815 | (HostTier::Builtin, None) => Ok(None), |
| 816 | _ => Err("extension owner has no authority".to_string()), |
| 817 | } |
| 818 | } |
| 819 | } |
| 820 | |
| 821 | async fn live_host( |
| 822 | &self, |
| 823 | tier: HostTier, |
| 824 | owner_if_live: impl FnOnce(&OwnerRegistry) -> Result<OwnerRef, String>, |
| 825 | ) -> Result<Arc<HostProcess>, String> { |
| 826 | let policy = activation::extension_host_policy_enabled(); |
| 827 | let authority = self.live_owner_authority(tier, owner_if_live)?; |
| 828 | if let Some(authority) = authority { |
| 829 | #[cfg(test)] |
| 830 | let env_scope = crate::test_support::env_scope_ticket(); |
| 831 | tokio::task::spawn_blocking(move || { |
| 832 | #[cfg(test)] |
| 833 | let _env_scope = crate::test_support::join_env_scope(env_scope); |
| 834 | let _scope = activation::PolicyScope::propagate(policy); |
| 835 | crate::plugins::registry::verify_plugin_component_authority( |
| 836 | &authority, |
| 837 | PluginActivationCapability::Native, |
| 838 | ) |
| 839 | }) |
| 840 | .await |
| 841 | .map_err(|error| format!("authority check failed: {error}"))??; |
| 842 | } |
| 843 | self.ready_host(tier) |
| 844 | .map_err(|status| host_down(tier, &status)) |
| 845 | } |
| 846 | |
| 847 | pub(crate) async fn live_host_for( |
| 848 | &self, |
| 849 | registration: &ToolRegistration, |
| 850 | ) -> Result<Arc<HostProcess>, String> { |
| 851 | self.live_host(registration.tier, |registry| { |
| 852 | if registry.is_live(registration.handle, ®istration.owner) { |
| 853 | Ok(registration.owner.clone()) |
| 854 | } else { |
| 855 | Err(format!( |
| 856 | "extension tool `{}` from `{}` is no longer registered", |
| 857 | registration.name, registration.plugin_name |
| 858 | )) |
| 859 | } |
| 860 | }) |
| 861 | .await |
| 862 | } |
| 863 | |
| 864 | /// [`Self::live_host`] for a command the user just invoked: the exact |
| 865 | /// registration (handle and owner generation) must still be admitted. |
| 866 | pub(crate) async fn live_host_for_command( |
| 867 | &self, |
| 868 | command: &command::ExtensionCommandRef, |
| 869 | ) -> Result<(Arc<HostProcess>, CommandRegistration), String> { |
| 870 | let mut found = None; |
| 871 | let host = self |
| 872 | .live_host(HostTier::Plugin, |registry| { |
| 873 | let registration = registry |
| 874 | .live_command(command.handle, &command.plugin_id, command.generation) |
| 875 | .ok_or_else(|| { |
| 876 | format!( |
| 877 | "extension command from `{}` is no longer registered (the plugin was reloaded, disabled or lost trust)", |
| 878 | command.plugin_id |
| 879 | ) |
| 880 | })?; |
| 881 | let owner = registration.owner.clone(); |
| 882 | found = Some(registration); |
| 883 | Ok(owner) |
| 884 | }) |
| 885 | .await?; |
| 886 | Ok((host, found.expect("set when the owner was found"))) |
| 887 | } |
| 888 | } |
| 889 | |
| 890 | /// The error a call gets while the host of `tier` cannot take it. |
| 891 | fn host_down(tier: HostTier, status: &HostStatus) -> String { |
| 892 | format!("{} is down: {status}", tier.host_label()) |
| 893 | } |
| 894 | |
| 895 | /// Channel callbacks. Holds a `Weak` so the host process (which owns the |
| 896 | /// callbacks) never keeps the manager alive. Each belongs to one tier's host |
| 897 | /// process and generation. |
| 898 | struct Events { |
| 899 | shared: Weak<ManagerShared>, |
| 900 | tier: HostTier, |
| 901 | generation: u64, |
| 902 | } |
| 903 | |
| 904 | #[async_trait] |
| 905 | impl HostEvents for Events { |
| 906 | fn responded(&self) { |
| 907 | if let Some(shared) = self.shared.upgrade() { |
| 908 | // Completing a live request is evidence of recovery before its |
| 909 | // caller's post-await liveness checks. The pending ping retains |
| 910 | // its original hang deadline; this creates no new health clock. |
| 911 | set_host_health(&shared, self.tier, self.generation, false); |
| 912 | } |
| 913 | } |
| 914 | |
| 915 | async fn host_request( |
| 916 | &self, |
| 917 | request: protocol::HostRequest, |
| 918 | cx: HostRequestContext, |
| 919 | ) -> Result<serde_json::Value, protocol::RpcErrorWire> { |
| 920 | let refuse = |message: &str| protocol::RpcErrorWire { |
| 921 | code: protocol::error_code::REFUSED, |
| 922 | message: message.to_string(), |
| 923 | data: None, |
| 924 | }; |
| 925 | match request { |
| 926 | protocol::HostRequest::ExecutionRedeem(params) => { |
| 927 | let shared = self |
| 928 | .shared |
| 929 | .upgrade() |
| 930 | .ok_or_else(|| refuse("extension host manager is gone"))?; |
| 931 | if self.tier != HostTier::Builtin { |
| 932 | return Err(refuse("execution broker is builtin-only")); |
| 933 | } |
| 934 | shared |
| 935 | .execution_broker |
| 936 | .serve(&shared, self.generation, params, cx) |
| 937 | .await |
| 938 | } |
| 939 | protocol::HostRequest::CoreCall(params) => { |
| 940 | let Some(shared) = self.shared.upgrade() else { |
| 941 | return Err(refuse("extension host manager is gone")); |
| 942 | }; |
| 943 | if shared |
| 944 | .tier_runtime(self.tier) |
| 945 | .host_generation |
| 946 | .load(Ordering::SeqCst) |
| 947 | != self.generation |
| 948 | { |
| 949 | return Err(refuse("stale host generation")); |
| 950 | } |
| 951 | shared |
| 952 | .core_calls |
| 953 | .serve(&shared, self.tier, self.generation, params, cx) |
| 954 | .await |
| 955 | } |
| 956 | protocol::HostRequest::Register(params) if params.kind == RegisterKind::McpServer => { |
| 957 | let shared = self |
| 958 | .shared |
| 959 | .upgrade() |
| 960 | .ok_or_else(|| refuse("extension host manager is gone"))?; |
| 961 | let result = |
| 962 | native_mcp::admit(&shared, self.tier, self.generation, params, &cx).await; |
| 963 | Ok(serde_json::to_value(result).expect("registration result is JSON")) |
| 964 | } |
| 965 | protocol::HostRequest::Register(params) if params.kind == RegisterKind::SkillRoot => { |
| 966 | let Some(shared) = self.shared.upgrade() else { |
| 967 | return Err(refuse("extension host manager is gone")); |
| 968 | }; |
| 969 | let result = |
| 970 | skills::admit_root(&shared, self.tier, self.generation, params, &cx).await; |
| 971 | Ok(serde_json::to_value(result).expect("registration result is JSON")) |
| 972 | } |
| 973 | request @ (protocol::HostRequest::ProcLaunch(_) |
| 974 | | protocol::HostRequest::ProcRead(_) |
| 975 | | protocol::HostRequest::ProcWrite(_) |
| 976 | | protocol::HostRequest::ProcClose(_) |
| 977 | | protocol::HostRequest::NetStart(_) |
| 978 | | protocol::HostRequest::NetFetch(_) |
| 979 | | protocol::HostRequest::NetRead(_) |
| 980 | | protocol::HostRequest::NetRelease(_) |
| 981 | | protocol::HostRequest::NetClose(_)) => { |
| 982 | let shared = self |
| 983 | .shared |
| 984 | .upgrade() |
| 985 | .ok_or_else(|| refuse("extension host manager is gone"))?; |
| 986 | if self.tier != HostTier::Builtin { |
| 987 | return Err(refuse("MCP broker is builtin-only")); |
| 988 | } |
| 989 | shared |
| 990 | .mcp_broker |
| 991 | .serve(&shared, self.generation, request, cx) |
| 992 | .await |
| 993 | } |
| 994 | other => supervisor::registry_host_request(self, other, &cx), |
| 995 | } |
| 996 | } |
| 997 | |
| 998 | fn register(&self, params: &protocol::RegisterParams) -> RegisterResult { |
| 999 | let Some(shared) = self.shared.upgrade() else { |
| 1000 | return RegisterResult::Refused { |
| 1001 | refused: "extension host manager is gone".to_string(), |
| 1002 | }; |
| 1003 | }; |
| 1004 | let runtime = shared.tier_runtime(self.tier); |
| 1005 | let slot = runtime.host.lock().expect("host lock"); |
| 1006 | if runtime.host_generation.load(Ordering::SeqCst) != self.generation { |
| 1007 | return RegisterResult::Refused { |
| 1008 | refused: "stale host generation".to_string(), |
| 1009 | }; |
| 1010 | } |
| 1011 | let mut foreign_owner = false; |
| 1012 | let result = { |
| 1013 | let mut registry = shared.registry.lock().expect("registry lock"); |
| 1014 | // A host registers only for owners it was asked to activate: one |
| 1015 | // of its own tier. The owner token already makes another tier's |
| 1016 | // owner unreachable; this says so out loud. |
| 1017 | match registry.owner(¶ms.owner.plugin_id) { |
| 1018 | Some(entry) if entry.tier != self.tier => { |
| 1019 | foreign_owner = true; |
| 1020 | Err(format!( |
| 1021 | "owner `{}` belongs to the {} tier, not this host's {} tier", |
| 1022 | params.owner.plugin_id, |
| 1023 | entry.tier.name(), |
| 1024 | self.tier.name() |
| 1025 | )) |
| 1026 | } |
| 1027 | _ => registry.register(params), |
| 1028 | } |
| 1029 | }; |
| 1030 | drop(slot); |
| 1031 | match result { |
| 1032 | Ok(handle) => RegisterResult::Admitted { handle }, |
| 1033 | Err(reason) => { |
| 1034 | let kind = match params.kind { |
| 1035 | RegisterKind::Tool => "tool", |
| 1036 | RegisterKind::Command => "command", |
| 1037 | RegisterKind::Hook => "hook", |
| 1038 | RegisterKind::PromptSection => "prompt section", |
| 1039 | RegisterKind::PromptTemplate => "prompt template", |
| 1040 | RegisterKind::SkillRoot => "skill root", |
| 1041 | RegisterKind::ShellHook => "shell hook", |
| 1042 | RegisterKind::McpServer => "MCP server", |
| 1043 | }; |
| 1044 | let message = format!( |
| 1045 | "extension `{}` {kind} `{}` refused: {reason}", |
| 1046 | params.owner.plugin_id, params.spec.name |
| 1047 | ); |
| 1048 | // A host naming the other tier's owner does not get to write |
| 1049 | // under that owner's id: the refusal is a global diagnostic. |
| 1050 | if foreign_owner { |
| 1051 | shared.diagnostic(message); |
| 1052 | } else { |
| 1053 | shared.plugin_diagnostic(¶ms.owner.plugin_id, message); |
| 1054 | } |
| 1055 | RegisterResult::Refused { refused: reason } |
| 1056 | } |
| 1057 | } |
| 1058 | } |
| 1059 | |
| 1060 | fn unregister(&self, params: &protocol::UnregisterParams) { |
| 1061 | if let Some(shared) = self.shared.upgrade() { |
| 1062 | shared |
| 1063 | .registry |
| 1064 | .lock() |
| 1065 | .expect("registry lock") |
| 1066 | .unregister(¶ms.owner, params.handle); |
| 1067 | } |
| 1068 | } |
| 1069 | |
| 1070 | fn faulted(&self, params: &protocol::FaultedParams) { |
| 1071 | let Some(shared) = self.shared.upgrade() else { |
| 1072 | return; |
| 1073 | }; |
| 1074 | let runtime = shared.tier_runtime(self.tier); |
| 1075 | let slot = runtime.host.lock().expect("host lock"); |
| 1076 | if runtime.host_generation.load(Ordering::SeqCst) != self.generation { |
| 1077 | return; |
| 1078 | } |
| 1079 | if !shared |
| 1080 | .registry |
| 1081 | .lock() |
| 1082 | .expect("registry lock") |
| 1083 | .mark_failed(¶ms.owner, OwnerState::Faulted(params.error.clone())) |
| 1084 | { |
| 1085 | return; |
| 1086 | } |
| 1087 | if let HostSlot::Ready(host) | HostSlot::Unresponsive(host) = &*slot { |
| 1088 | host.revoke_calls_of(¶ms.owner.plugin_id); |
| 1089 | } |
| 1090 | drop(slot); |
| 1091 | shared.core_calls.revoke_owner(¶ms.owner.plugin_id); |
| 1092 | shared |
| 1093 | .execution_broker |
| 1094 | .revoke_owner(¶ms.owner.plugin_id); |
| 1095 | shared.mcp_users.fetch_sub( |
| 1096 | shared.mcp_broker.revoke_owner(¶ms.owner.plugin_id), |
| 1097 | Ordering::SeqCst, |
| 1098 | ); |
| 1099 | shared.plugin_diagnostic( |
| 1100 | ¶ms.owner.plugin_id, |
| 1101 | format!( |
| 1102 | "extension `{}` faulted and was disposed: {}", |
| 1103 | params.owner.plugin_id, params.error |
| 1104 | ), |
| 1105 | ); |
| 1106 | } |
| 1107 | |
| 1108 | fn exited(&self, host_generation: u64, reason: String, stderr_tail: String) { |
| 1109 | let Some(shared) = self.shared.upgrade() else { |
| 1110 | return; |
| 1111 | }; |
| 1112 | // Whatever this process was asked for, nothing it held is good again. |
| 1113 | shared.core_calls.revoke_host(self.tier, host_generation); |
| 1114 | shared |
| 1115 | .execution_broker |
| 1116 | .revoke_host(self.tier, host_generation); |
| 1117 | shared.mcp_users.fetch_sub( |
| 1118 | shared.mcp_broker.revoke_host(self.tier, host_generation), |
| 1119 | Ordering::SeqCst, |
| 1120 | ); |
| 1121 | let runtime = shared.tier_runtime(self.tier); |
| 1122 | let retry = { |
| 1123 | let mut slot = runtime.host.lock().expect("host lock"); |
| 1124 | if runtime.host_generation.load(Ordering::SeqCst) != host_generation |
| 1125 | || !matches!(&*slot, HostSlot::Ready(_) | HostSlot::Unresponsive(_)) |
| 1126 | { |
| 1127 | return; |
| 1128 | } |
| 1129 | let mut supervision = runtime.supervision.lock().expect("supervision lock"); |
| 1130 | let planned = supervision.planned_restart.take() == Some(host_generation) |
| 1131 | && reason == DIRTY_RESTART_REASON; |
| 1132 | let restart = if planned { |
| 1133 | supervision.retry_ticket += 1; |
| 1134 | true |
| 1135 | } else { |
| 1136 | supervision.record_crash(Instant::now(), &shared.options.supervision) |
| 1137 | }; |
| 1138 | shared |
| 1139 | .registry |
| 1140 | .lock() |
| 1141 | .expect("registry lock") |
| 1142 | .host_exited(self.tier, &reason); |
| 1143 | *slot = if restart { |
| 1144 | HostSlot::Restarting { |
| 1145 | reason: reason.clone(), |
| 1146 | } |
| 1147 | } else { |
| 1148 | HostSlot::Failed { |
| 1149 | reason: format!("crash budget exhausted: {reason}"), |
| 1150 | stderr_tail, |
| 1151 | } |
| 1152 | }; |
| 1153 | restart.then_some((supervision.retry_ticket, supervision.policy)) |
| 1154 | }; |
| 1155 | shared.diagnostic(format!("{} {reason}", self.tier.host_label())); |
| 1156 | if let Some((ticket, policy)) = retry { |
| 1157 | schedule_restart(&shared, self.tier, host_generation, ticket, policy); |
| 1158 | } |
| 1159 | } |
| 1160 | |
| 1161 | fn log(&self, params: &protocol::LogParams) { |
| 1162 | if !matches!(params.level.as_str(), "warn" | "error") { |
| 1163 | return; |
| 1164 | } |
| 1165 | let Some(plugin_id) = params.plugin_id.as_deref() else { |
| 1166 | return; |
| 1167 | }; |
| 1168 | let Some(shared) = self.shared.upgrade() else { |
| 1169 | return; |
| 1170 | }; |
| 1171 | let runtime = shared.tier_runtime(self.tier); |
| 1172 | let _slot = runtime.host.lock().expect("host lock"); |
| 1173 | // Only for an owner of this host's own tier: a host cannot put |
| 1174 | // diagnostics under the other tier's owner ids. |
| 1175 | if runtime.host_generation.load(Ordering::SeqCst) != self.generation |
| 1176 | || !shared |
| 1177 | .registry |
| 1178 | .lock() |
| 1179 | .expect("registry lock") |
| 1180 | .owner(plugin_id) |
| 1181 | .is_some_and(|entry| entry.tier == self.tier) |
| 1182 | { |
| 1183 | return; |
| 1184 | } |
| 1185 | shared.plugin_diagnostic(plugin_id, format!("{}: {}", params.level, params.msg)); |
| 1186 | } |
| 1187 | } |
| 1188 | |
| 1189 | /// One scheduled retry owns a ticket, so explicit retry/shutdown and a newer |
| 1190 | /// host generation invalidate it. The existing reconcile lock owns replay. |
| 1191 | fn schedule_restart( |
| 1192 | shared: &Arc<ManagerShared>, |
| 1193 | tier: HostTier, |
| 1194 | generation: u64, |
| 1195 | ticket: u64, |
| 1196 | policy: bool, |
| 1197 | ) { |
| 1198 | let weak = Arc::downgrade(shared); |
| 1199 | let backoff = shared.options.supervision.restart_backoff; |
| 1200 | tokio::spawn(async move { |
| 1201 | tokio::time::sleep(backoff).await; |
| 1202 | let Some(shared) = weak.upgrade() else { |
| 1203 | return; |
| 1204 | }; |
| 1205 | { |
| 1206 | let _serial = shared.sync_lock.lock().await; |
| 1207 | let runtime = shared.tier_runtime(tier); |
| 1208 | let mut slot = runtime.host.lock().expect("host lock"); |
| 1209 | if runtime.host_generation.load(Ordering::SeqCst) != generation |
| 1210 | || runtime |
| 1211 | .supervision |
| 1212 | .lock() |
| 1213 | .expect("supervision lock") |
| 1214 | .retry_ticket |
| 1215 | != ticket |
| 1216 | || !matches!(&*slot, HostSlot::Restarting { .. }) |
| 1217 | { |
| 1218 | return; |
| 1219 | } |
| 1220 | *slot = HostSlot::Idle; |
| 1221 | } |
| 1222 | let manager = ExtensionHostManager { shared }; |
| 1223 | if let Err(error) = manager.reconcile_with_policy(policy).await { |
| 1224 | manager.shared.diagnostic(error); |
| 1225 | } |
| 1226 | }); |
| 1227 | } |
| 1228 | |
| 1229 | fn set_host_health( |
| 1230 | shared: &ManagerShared, |
| 1231 | tier: HostTier, |
| 1232 | generation: u64, |
| 1233 | unresponsive: bool, |
| 1234 | ) -> bool { |
| 1235 | let runtime = shared.tier_runtime(tier); |
| 1236 | let mut slot = runtime.host.lock().expect("host lock"); |
| 1237 | if runtime.host_generation.load(Ordering::SeqCst) != generation { |
| 1238 | return false; |
| 1239 | } |
| 1240 | let host = match &*slot { |
| 1241 | HostSlot::Ready(host) | HostSlot::Unresponsive(host) => Arc::clone(host), |
| 1242 | _ => return false, |
| 1243 | }; |
| 1244 | *slot = if unresponsive { |
| 1245 | HostSlot::Unresponsive(host) |
| 1246 | } else { |
| 1247 | HostSlot::Ready(host) |
| 1248 | }; |
| 1249 | true |
| 1250 | } |
| 1251 | |
| 1252 | fn restart_dirty_host_when_idle( |
| 1253 | shared: &ManagerShared, |
| 1254 | host: &Arc<HostProcess>, |
| 1255 | generation: u64, |
| 1256 | ) -> bool { |
| 1257 | // Reconciliation owns activation/deactivation between wire requests too. |
| 1258 | let Ok(_serial) = shared.sync_lock.try_lock() else { |
| 1259 | return false; |
| 1260 | }; |
| 1261 | let runtime = shared.tier_runtime(host.tier); |
| 1262 | let slot = runtime.host.lock().expect("host lock"); |
| 1263 | if runtime.host_generation.load(Ordering::SeqCst) != generation |
| 1264 | || !matches!(&*slot, HostSlot::Ready(current) if Arc::ptr_eq(current, host)) |
| 1265 | { |
| 1266 | return false; |
| 1267 | } |
| 1268 | let mut supervision = runtime.supervision.lock().expect("supervision lock"); |
| 1269 | if !supervision.dirty_restart_pending || !host.terminate_if_idle(DIRTY_RESTART_REASON) { |
| 1270 | return false; |
| 1271 | } |
| 1272 | supervision.planned_restart = Some(generation); |
| 1273 | true |
| 1274 | } |
| 1275 | |
| 1276 | /// A monitor never owns the manager. Dropping the manager or changing host |
| 1277 | /// generation stops its monitor; pending calls are never retried here. |
| 1278 | fn monitor_host(shared: &Arc<ManagerShared>, host: &Arc<HostProcess>, generation: u64) { |
| 1279 | let tier = host.tier; |
| 1280 | let weak = Arc::downgrade(shared); |
| 1281 | let host = Arc::clone(host); |
| 1282 | let options = shared.options.supervision.clone(); |
| 1283 | tokio::spawn(async move { |
| 1284 | let mut blocked_since = None; |
| 1285 | loop { |
| 1286 | tokio::time::sleep(options.heartbeat_interval).await; |
| 1287 | let Some(shared) = weak.upgrade() else { |
| 1288 | return; |
| 1289 | }; |
| 1290 | if shared |
| 1291 | .tier_runtime(tier) |
| 1292 | .host_generation |
| 1293 | .load(Ordering::SeqCst) |
| 1294 | != generation |
| 1295 | || host.has_exited() |
| 1296 | { |
| 1297 | return; |
| 1298 | } |
| 1299 | if restart_dirty_host_when_idle(&shared, &host, generation) { |
| 1300 | return; |
| 1301 | } |
| 1302 | // Only where no kernel limit is planned (a macOS Node host). |
| 1303 | if host.memory() == supervisor::MemoryEnforcement::Heartbeat |
| 1304 | && let Some(resident) = host.pid.and_then(supervisor::resident_bytes) |
| 1305 | && resident > host.memory_cap |
| 1306 | { |
| 1307 | shared.diagnostic(format!( |
| 1308 | "{} exceeded its memory cap ({} MiB resident, cap {} MiB); killed", |
| 1309 | tier.host_label(), |
| 1310 | resident / (1024 * 1024), |
| 1311 | host.memory_cap / (1024 * 1024) |
| 1312 | )); |
| 1313 | drop(shared); |
| 1314 | host.terminate("exceeded the memory cap".into()); |
| 1315 | return; |
| 1316 | } |
| 1317 | drop(shared); |
| 1318 | let (id, mut answer) = match host.start_request(CoreRequest::Ping, None) { |
| 1319 | Ok(request) => { |
| 1320 | blocked_since = None; |
| 1321 | request |
| 1322 | } |
| 1323 | Err(supervisor::HostCallError::Busy) => { |
| 1324 | // A full outbound queue can itself be caused by a hung |
| 1325 | // host. Bound that wait too instead of skipping forever. |
| 1326 | let elapsed = blocked_since.get_or_insert_with(Instant::now).elapsed(); |
| 1327 | if elapsed >= options.hang_timeout { |
| 1328 | host.terminate("heartbeat queue remained blocked".into()); |
| 1329 | return; |
| 1330 | } |
| 1331 | if elapsed >= options.ping_timeout { |
| 1332 | let Some(shared) = weak.upgrade() else { |
| 1333 | return; |
| 1334 | }; |
| 1335 | if !set_host_health(&shared, tier, generation, true) { |
| 1336 | return; |
| 1337 | } |
| 1338 | } |
| 1339 | continue; |
| 1340 | } |
| 1341 | Err(_) => return, |
| 1342 | }; |
| 1343 | let result = match tokio::time::timeout(options.ping_timeout, &mut answer).await { |
| 1344 | Ok(result) => result, |
| 1345 | Err(_) => { |
| 1346 | let Some(shared) = weak.upgrade() else { |
| 1347 | host.forget(id); |
| 1348 | return; |
| 1349 | }; |
| 1350 | if !set_host_health(&shared, tier, generation, true) { |
| 1351 | host.forget(id); |
| 1352 | return; |
| 1353 | } |
| 1354 | drop(shared); |
| 1355 | let remaining = options.hang_timeout.saturating_sub(options.ping_timeout); |
| 1356 | match tokio::time::timeout(remaining, answer).await { |
| 1357 | Ok(result) => result, |
| 1358 | Err(_) => { |
| 1359 | host.forget(id); |
| 1360 | host.terminate("heartbeat timed out".into()); |
| 1361 | return; |
| 1362 | } |
| 1363 | } |
| 1364 | } |
| 1365 | }; |
| 1366 | host.forget(id); |
| 1367 | if !matches!(result, Ok(Ok(serde_json::Value::Object(ref object))) if object.is_empty()) |
| 1368 | { |
| 1369 | if !host.has_exited() { |
| 1370 | host.terminate("invalid heartbeat response".into()); |
| 1371 | } |
| 1372 | return; |
| 1373 | } |
| 1374 | let Some(shared) = weak.upgrade() else { |
| 1375 | return; |
| 1376 | }; |
| 1377 | if !set_host_health(&shared, tier, generation, false) { |
| 1378 | return; |
| 1379 | } |
| 1380 | } |
| 1381 | }); |
| 1382 | } |
| 1383 | |
| 1384 | /// Supervises at most one extension host per trust tier for this engine |
| 1385 | /// process: the plugin tier's, and the builtin tier's once a built-in module |
| 1386 | /// asks for it (none does yet). |
| 1387 | pub struct ExtensionHostManager { |
| 1388 | pub(crate) shared: Arc<ManagerShared>, |
| 1389 | } |
| 1390 | |
| 1391 | impl ExtensionHostManager { |
| 1392 | #[must_use] |
| 1393 | pub fn new(options: ExtensionHostOptions) -> Self { |
| 1394 | Self::with_builtin_modules(options, tier::BUILTIN_MODULES) |
| 1395 | } |
| 1396 | |
| 1397 | /// A manager with its own built-in module table. Production passes the |
| 1398 | /// const table through [`Self::new`]; a test passes one of its own to |
| 1399 | /// exercise tier 0 without a production row. |
| 1400 | fn with_builtin_modules( |
| 1401 | options: ExtensionHostOptions, |
| 1402 | builtin_modules: &'static [BuiltinModule], |
| 1403 | ) -> Self { |
| 1404 | Self { |
| 1405 | shared: Arc::new(ManagerShared { |
| 1406 | options, |
| 1407 | builtin_modules, |
| 1408 | attachments: Mutex::new(BTreeMap::new()), |
| 1409 | next_attachment: AtomicU64::new(0), |
| 1410 | registry: Mutex::new(OwnerRegistry::new()), |
| 1411 | core_calls: Arc::default(), |
| 1412 | engine_handle: OnceLock::new(), |
| 1413 | execution_broker: execution::Broker::default(), |
| 1414 | harness_users: AtomicU64::new(0), |
| 1415 | stock_users: AtomicU64::new(0), |
| 1416 | mcp_broker: mcp::Broker::default(), |
| 1417 | mcp_users: AtomicU64::new(0), |
| 1418 | plugin: TierRuntime::new(), |
| 1419 | builtin: TierRuntime::new(), |
| 1420 | sync_lock: tokio::sync::Mutex::new(()), |
| 1421 | skill_admission: Arc::new(tokio::sync::Semaphore::new(1)), |
| 1422 | diagnostics: Mutex::new(VecDeque::new()), |
| 1423 | runtime: Mutex::new(RuntimePin::default()), |
| 1424 | plugin_configs: Mutex::new(plugin_config::PluginConfigs::default()), |
| 1425 | }), |
| 1426 | } |
| 1427 | } |
| 1428 | |
| 1429 | /// Reuse the Engine's existing runtime; structural commands may leave it unbound. |
| 1430 | pub(crate) fn bind_engine_handle(&self, handle: tokio::runtime::Handle) { |
| 1431 | let _ = self.shared.engine_handle.set(handle); |
| 1432 | } |
| 1433 | |
| 1434 | /// The plugin host's state: what `/plugin` and the error a plugin's tool |
| 1435 | /// call gets while it is down report. |
| 1436 | #[must_use] |
| 1437 | pub fn status(&self) -> HostStatus { |
| 1438 | self.tier_status(HostTier::Plugin) |
| 1439 | } |
| 1440 | |
| 1441 | /// The state of `tier`'s host. |
| 1442 | #[must_use] |
| 1443 | fn tier_status(&self, tier: HostTier) -> HostStatus { |
| 1444 | self.shared |
| 1445 | .tier_runtime(tier) |
| 1446 | .host |
| 1447 | .lock() |
| 1448 | .expect("host lock") |
| 1449 | .status() |
| 1450 | } |
| 1451 | |
| 1452 | /// The pinned runtime's one-line summary, once a host on it has completed |
| 1453 | /// the handshake. |
| 1454 | #[must_use] |
| 1455 | pub fn runtime_summary(&self) -> Option<String> { |
| 1456 | self.shared |
| 1457 | .runtime |
| 1458 | .lock() |
| 1459 | .expect("runtime lock") |
| 1460 | .pinned |
| 1461 | .as_ref() |
| 1462 | .map(|(_, summary)| summary.clone()) |
| 1463 | } |
| 1464 | |
| 1465 | /// How many times this manager has tried to start a host process, both |
| 1466 | /// tiers together. |
| 1467 | #[must_use] |
| 1468 | pub fn spawn_attempts(&self) -> u64 { |
| 1469 | HostTier::ALL |
| 1470 | .into_iter() |
| 1471 | .map(|tier| self.tier_spawn_attempts(tier)) |
| 1472 | .sum() |
| 1473 | } |
| 1474 | |
| 1475 | /// How many times this manager has tried to start `tier`'s host. |
| 1476 | #[must_use] |
| 1477 | fn tier_spawn_attempts(&self, tier: HostTier) -> u64 { |
| 1478 | self.shared |
| 1479 | .tier_runtime(tier) |
| 1480 | .spawn_attempts |
| 1481 | .load(Ordering::SeqCst) |
| 1482 | } |
| 1483 | |
| 1484 | #[must_use] |
| 1485 | pub fn diagnostics(&self) -> Vec<String> { |
| 1486 | self.shared |
| 1487 | .diagnostics |
| 1488 | .lock() |
| 1489 | .expect("diagnostics lock") |
| 1490 | .iter() |
| 1491 | .map(|entry| entry.message.clone()) |
| 1492 | .collect() |
| 1493 | } |
| 1494 | |
| 1495 | /// Recent retained diagnostics for exactly this plugin, not a persistent |
| 1496 | /// log. The command renderer escapes each field before displaying it. |
| 1497 | #[must_use] |
| 1498 | pub fn owner_report(&self, plugin_id: &str) -> Option<OwnerReport> { |
| 1499 | let registry = self.shared.registry.lock().expect("registry lock"); |
| 1500 | let state = registry.owner(plugin_id).map(|entry| match &entry.state { |
| 1501 | OwnerState::Failed(reason) => OwnerState::Failed(bounded_diagnostic(reason.clone())), |
| 1502 | OwnerState::Faulted(reason) => OwnerState::Faulted(bounded_diagnostic(reason.clone())), |
| 1503 | state => state.clone(), |
| 1504 | }); |
| 1505 | let tools = registry |
| 1506 | .live_tools() |
| 1507 | .into_iter() |
| 1508 | .filter(|tool| tool.owner.plugin_id == plugin_id) |
| 1509 | .map(|tool| tool.name) |
| 1510 | .collect(); |
| 1511 | drop(registry); |
| 1512 | let mut diagnostics: Vec<_> = self |
| 1513 | .shared |
| 1514 | .diagnostics |
| 1515 | .lock() |
| 1516 | .expect("diagnostics lock") |
| 1517 | .iter() |
| 1518 | .rev() |
| 1519 | .filter(|entry| entry.plugin_id.as_deref() == Some(plugin_id)) |
| 1520 | .take(20) |
| 1521 | .map(|entry| entry.message.clone()) |
| 1522 | .collect(); |
| 1523 | diagnostics.reverse(); |
| 1524 | if state.is_none() && diagnostics.is_empty() { |
| 1525 | return None; |
| 1526 | } |
| 1527 | Some(OwnerReport { |
| 1528 | state, |
| 1529 | tools, |
| 1530 | diagnostics, |
| 1531 | }) |
| 1532 | } |
| 1533 | |
| 1534 | #[cfg(test)] |
| 1535 | #[must_use] |
| 1536 | pub fn live_tool_names(&self) -> Vec<String> { |
| 1537 | self.shared |
| 1538 | .registry |
| 1539 | .lock() |
| 1540 | .expect("registry lock") |
| 1541 | .live_tools() |
| 1542 | .into_iter() |
| 1543 | .map(|tool| tool.name) |
| 1544 | .collect() |
| 1545 | } |
| 1546 | |
| 1547 | #[cfg(test)] |
| 1548 | #[must_use] |
| 1549 | pub fn live_command_names(&self) -> Vec<String> { |
| 1550 | self.shared |
| 1551 | .registry |
| 1552 | .lock() |
| 1553 | .expect("registry lock") |
| 1554 | .live_commands() |
| 1555 | .into_iter() |
| 1556 | .map(|command| command.name) |
| 1557 | .collect() |
| 1558 | } |
| 1559 | |
| 1560 | #[cfg(test)] |
| 1561 | #[must_use] |
| 1562 | pub fn owner_state(&self, plugin_id: &str) -> Option<OwnerState> { |
| 1563 | self.shared |
| 1564 | .registry |
| 1565 | .lock() |
| 1566 | .expect("registry lock") |
| 1567 | .owner(plugin_id) |
| 1568 | .map(|entry| entry.state.clone()) |
| 1569 | } |
| 1570 | |
| 1571 | /// Install the user's `[plugins]` settings (boot), and the config file a |
| 1572 | /// later `/plugin reload` re-reads them from. Takes effect at the next |
| 1573 | /// reconcile: a plugin whose settings changed is re-activated. |
| 1574 | pub fn set_plugin_settings( |
| 1575 | &self, |
| 1576 | settings: &BTreeMap<String, crate::config::PluginSettings>, |
| 1577 | source: Option<PathBuf>, |
| 1578 | ) { |
| 1579 | self.shared |
| 1580 | .plugin_configs |
| 1581 | .lock() |
| 1582 | .expect("plugin configs lock") |
| 1583 | .replace(settings, source); |
| 1584 | } |
| 1585 | |
| 1586 | /// Read `[plugins]` from the config file again, if boot named one. A file |
| 1587 | /// that cannot be read or parsed keeps the previous settings. |
| 1588 | async fn reload_plugin_settings(&self) { |
| 1589 | let Some(path) = self |
| 1590 | .shared |
| 1591 | .plugin_configs |
| 1592 | .lock() |
| 1593 | .expect("plugin configs lock") |
| 1594 | .source() |
| 1595 | else { |
| 1596 | return; |
| 1597 | }; |
| 1598 | let read_path = path.clone(); |
| 1599 | match tokio::task::spawn_blocking(move || crate::config::read_plugin_settings(&read_path)) |
| 1600 | .await |
| 1601 | { |
| 1602 | Ok(Ok(settings)) => self |
| 1603 | .shared |
| 1604 | .plugin_configs |
| 1605 | .lock() |
| 1606 | .expect("plugin configs lock") |
| 1607 | .replace_reloaded(&settings), |
| 1608 | Ok(Err(reason)) => self.shared.diagnostic(format!( |
| 1609 | "plugin settings were not reloaded ({reason}); keeping the previous ones" |
| 1610 | )), |
| 1611 | Err(error) => self.shared.diagnostic(format!( |
| 1612 | "plugin settings were not reloaded from {}: {error}", |
| 1613 | path.display() |
| 1614 | )), |
| 1615 | } |
| 1616 | } |
| 1617 | |
| 1618 | /// What the user configured for `plugin_name`: the keys (never the values) |
| 1619 | /// or why the config is refused. `None` when nothing is configured. |
| 1620 | #[must_use] |
| 1621 | pub fn plugin_config_summary(&self, plugin_name: &str) -> Option<Result<Vec<String>, String>> { |
| 1622 | self.shared |
| 1623 | .plugin_configs |
| 1624 | .lock() |
| 1625 | .expect("plugin configs lock") |
| 1626 | .summary(plugin_name) |
| 1627 | } |
| 1628 | |
| 1629 | /// An explicit plugin mutation retries failed receipts and clears the |
| 1630 | /// shared crash budget. Merely opening another engine never does this. |
| 1631 | pub fn retry(&self) { |
| 1632 | // One tier at a time, so no two tiers' locks are ever held together. |
| 1633 | for tier in HostTier::ALL { |
| 1634 | let runtime = self.shared.tier_runtime(tier); |
| 1635 | let mut slot = runtime.host.lock().expect("host lock"); |
| 1636 | let mut supervision = runtime.supervision.lock().expect("supervision lock"); |
| 1637 | supervision.crashes.clear(); |
| 1638 | supervision.launch_failed = false; |
| 1639 | supervision.retry_ticket += 1; |
| 1640 | if matches!( |
| 1641 | &*slot, |
| 1642 | HostSlot::Failed { .. } | HostSlot::Restarting { .. } |
| 1643 | ) { |
| 1644 | *slot = HostSlot::Idle; |
| 1645 | } |
| 1646 | } |
| 1647 | self.shared |
| 1648 | .registry |
| 1649 | .lock() |
| 1650 | .expect("registry lock") |
| 1651 | .forget_inactive(); |
| 1652 | } |
| 1653 | |
| 1654 | fn retry_launch_on_attach(&self) { |
| 1655 | for tier in HostTier::ALL { |
| 1656 | let runtime = self.shared.tier_runtime(tier); |
| 1657 | let mut slot = runtime.host.lock().expect("host lock"); |
| 1658 | let mut supervision = runtime.supervision.lock().expect("supervision lock"); |
| 1659 | if matches!(&*slot, HostSlot::Failed { .. }) |
| 1660 | && supervision.launch_failed |
| 1661 | && supervision.last_start.is_some_and(|at| { |
| 1662 | at.elapsed() >= self.shared.options.supervision.start_retry_cooldown |
| 1663 | }) |
| 1664 | { |
| 1665 | *slot = HostSlot::Idle; |
| 1666 | supervision.retry_ticket += 1; |
| 1667 | } |
| 1668 | } |
| 1669 | } |
| 1670 | |
| 1671 | /// Record native tool names from one engine's turn build (the registry |
| 1672 | /// before scripts, plugins and extensions are added) so |
| 1673 | /// `registry/register` refuses collisions. Additive: engines in one |
| 1674 | /// process report different native surfaces and none may shrink the set. |
| 1675 | pub fn note_native_names<'a>(&self, names: impl IntoIterator<Item = &'a str>) { |
| 1676 | self.shared |
| 1677 | .registry |
| 1678 | .lock() |
| 1679 | .expect("registry lock") |
| 1680 | .add_native_names(names); |
| 1681 | } |
| 1682 | |
| 1683 | /// Attach an engine whose workspace plugin snapshot is `plugins`. Nothing |
| 1684 | /// is reconciled until [`HostAttachment::sync`] or a background sync. |
| 1685 | #[must_use] |
| 1686 | pub fn attach(self: &Arc<Self>, plugins: Arc<PluginRegistry>) -> HostAttachment { |
| 1687 | self.retry_launch_on_attach(); |
| 1688 | let id = self.shared.next_attachment.fetch_add(1, Ordering::SeqCst) + 1; |
| 1689 | self.shared |
| 1690 | .attachments |
| 1691 | .lock() |
| 1692 | .expect("attachments lock") |
| 1693 | .insert( |
| 1694 | id, |
| 1695 | AttachmentState { |
| 1696 | plugins: Arc::new(plugins.bind_caller(composition_scope::SelectionRevision { |
| 1697 | attachment_id: id, |
| 1698 | revision: 1, |
| 1699 | })), |
| 1700 | desired: BTreeMap::new(), |
| 1701 | selection: composition_scope::CompositionSelection::default(), |
| 1702 | revision: 1, |
| 1703 | selection_cancel: CancellationToken::new(), |
| 1704 | session_id: None, |
| 1705 | agent_id: None, |
| 1706 | }, |
| 1707 | ); |
| 1708 | command::bump_epoch(); |
| 1709 | HostAttachment { |
| 1710 | id, |
| 1711 | manager: Arc::clone(self), |
| 1712 | } |
| 1713 | } |
| 1714 | |
| 1715 | /// How many engines are attached. |
| 1716 | #[must_use] |
| 1717 | pub fn attached_engines(&self) -> usize { |
| 1718 | self.shared |
| 1719 | .attachments |
| 1720 | .lock() |
| 1721 | .expect("attachments lock") |
| 1722 | .len() |
| 1723 | } |
| 1724 | |
| 1725 | /// Replace the snapshot of every engine attached to `plugins`'s |
| 1726 | /// workspace: a plugin was enabled, disabled, trusted or revoked there. |
| 1727 | fn refresh_workspace(&self, plugins: &Arc<PluginRegistry>) { |
| 1728 | let ids: Vec<_> = self |
| 1729 | .shared |
| 1730 | .attachments |
| 1731 | .lock() |
| 1732 | .expect("attachments lock") |
| 1733 | .iter() |
| 1734 | .filter(|(_, state)| state.plugins.workspace() == plugins.workspace()) |
| 1735 | .map(|(id, _)| *id) |
| 1736 | .collect(); |
| 1737 | for id in ids { |
| 1738 | self.shared.core_calls.revoke_attachment(id); |
| 1739 | } |
| 1740 | for state in self |
| 1741 | .shared |
| 1742 | .attachments |
| 1743 | .lock() |
| 1744 | .expect("attachments lock") |
| 1745 | .values_mut() |
| 1746 | { |
| 1747 | if state.plugins.workspace() == plugins.workspace() { |
| 1748 | state.selection_cancel.cancel(); |
| 1749 | state.selection_cancel = CancellationToken::new(); |
| 1750 | state.revision += 1; |
| 1751 | let view = plugins.retain_native_selection_from(state.plugins.as_ref()); |
| 1752 | let id = state |
| 1753 | .plugins |
| 1754 | .caller_selection() |
| 1755 | .expect("attached caller") |
| 1756 | .attachment_id; |
| 1757 | state.plugins = Arc::new(view.bind_caller(composition_scope::SelectionRevision { |
| 1758 | attachment_id: id, |
| 1759 | revision: state.revision, |
| 1760 | })); |
| 1761 | state.desired.clear(); |
| 1762 | state.selection = Default::default(); |
| 1763 | } |
| 1764 | } |
| 1765 | command::bump_epoch(); |
| 1766 | } |
| 1767 | |
| 1768 | /// Add the live tools of the owners attachment `id` desires to |
| 1769 | /// `tool_registry`, *after* natives and `~/.codewhale/tools` scripts. A |
| 1770 | /// name already present is skipped with a diagnostic — |
| 1771 | /// `ToolRegistry::register` would silently overwrite it. Returns the |
| 1772 | /// names added. |
| 1773 | fn install_tools_for( |
| 1774 | &self, |
| 1775 | id: u64, |
| 1776 | tool_registry: &mut crate::tools::ToolRegistry, |
| 1777 | ) -> Vec<String> { |
| 1778 | let desired = self |
| 1779 | .shared |
| 1780 | .attachments |
| 1781 | .lock() |
| 1782 | .expect("attachments lock") |
| 1783 | .get(&id) |
| 1784 | .map(|state| state.selection.clone()) |
| 1785 | .unwrap_or_default(); |
| 1786 | if desired.desired.is_empty() { |
| 1787 | return Vec::new(); |
| 1788 | } |
| 1789 | let tools: Vec<ToolRegistration> = self |
| 1790 | .shared |
| 1791 | .registry |
| 1792 | .lock() |
| 1793 | .expect("registry lock") |
| 1794 | .live_tools() |
| 1795 | .into_iter() |
| 1796 | .filter(|tool| { |
| 1797 | desired.includes( |
| 1798 | &tool.owner.plugin_id, |
| 1799 | &tool.content_hash, |
| 1800 | tool.scope.as_ref(), |
| 1801 | ) |
| 1802 | }) |
| 1803 | .collect(); |
| 1804 | if tools.is_empty() { |
| 1805 | return Vec::new(); |
| 1806 | } |
| 1807 | let taken: HashSet<String> = tool_registry |
| 1808 | .names() |
| 1809 | .into_iter() |
| 1810 | .map(str::to_ascii_lowercase) |
| 1811 | .collect(); |
| 1812 | let mut installed = Vec::new(); |
| 1813 | let mut counts = BTreeMap::new(); |
| 1814 | for tool in &tools { |
| 1815 | *counts |
| 1816 | .entry(tool.name.to_ascii_lowercase()) |
| 1817 | .or_insert(0usize) += 1; |
| 1818 | } |
| 1819 | for registration in tools { |
| 1820 | if counts[®istration.name.to_ascii_lowercase()] != 1 { |
| 1821 | self.shared.plugin_diagnostic(®istration.owner.plugin_id,format!("Native tool `{}` is ambiguous in this caller; rename the selected definitions",registration.name)); |
| 1822 | continue; |
| 1823 | } |
| 1824 | if taken.contains(®istration.name.to_ascii_lowercase()) { |
| 1825 | let origin = tool_registry |
| 1826 | .get(®istration.name) |
| 1827 | .map(|existing| existing.registration_origin().into_owned()) |
| 1828 | .unwrap_or_else(|| "another tool".to_string()); |
| 1829 | self.shared.plugin_diagnostic(®istration.owner.plugin_id, format!( |
| 1830 | "extension tool `{}` from `{}` skipped: the name is already registered by {origin}", |
| 1831 | registration.name, registration.plugin_name |
| 1832 | )); |
| 1833 | continue; |
| 1834 | } |
| 1835 | installed.push(registration.name.clone()); |
| 1836 | tool_registry.register(Arc::new(tool::HostToolSpec::for_selection( |
| 1837 | registration, |
| 1838 | Arc::clone(&self.shared), |
| 1839 | desired.revision, |
| 1840 | ))); |
| 1841 | } |
| 1842 | installed |
| 1843 | } |
| 1844 | |
| 1845 | /// The live commands of the owners that engines attached for `workspace` |
| 1846 | /// desire, each bound to the reviewed bytes that engine desires: what the |
| 1847 | /// user registry loads for that workspace. A workspace never sees another |
| 1848 | /// workspace's project plugins' commands. |
| 1849 | pub(crate) fn commands_for_workspace( |
| 1850 | &self, |
| 1851 | workspace: &Path, |
| 1852 | ) -> Vec<command::ExtensionCommandEntry> { |
| 1853 | let mut desired: BTreeSet<(String, String)> = BTreeSet::new(); |
| 1854 | for state in self |
| 1855 | .shared |
| 1856 | .attachments |
| 1857 | .lock() |
| 1858 | .expect("attachments lock") |
| 1859 | .values() |
| 1860 | .filter(|state| state.plugins.workspace() == workspace) |
| 1861 | { |
| 1862 | desired.extend( |
| 1863 | state |
| 1864 | .desired |
| 1865 | .iter() |
| 1866 | .map(|(plugin_id, hash)| (plugin_id.clone(), hash.clone())), |
| 1867 | ); |
| 1868 | } |
| 1869 | if desired.is_empty() { |
| 1870 | return Vec::new(); |
| 1871 | } |
| 1872 | let registry = self.shared.registry.lock().expect("registry lock"); |
| 1873 | registry |
| 1874 | .live_commands() |
| 1875 | .into_iter() |
| 1876 | .filter(|command| { |
| 1877 | command.scope.is_none() |
| 1878 | && desired.contains(&( |
| 1879 | command.owner.plugin_id.clone(), |
| 1880 | command.content_hash.clone(), |
| 1881 | )) |
| 1882 | }) |
| 1883 | .filter_map(|registration| { |
| 1884 | let authority = registry.authority_for(®istration.owner)?; |
| 1885 | Some(command::ExtensionCommandEntry { |
| 1886 | registration, |
| 1887 | authority, |
| 1888 | workspace: workspace.to_path_buf(), |
| 1889 | selection: None, |
| 1890 | }) |
| 1891 | }) |
| 1892 | .collect() |
| 1893 | } |
| 1894 | |
| 1895 | pub(crate) fn commands_for_plugins( |
| 1896 | &self, |
| 1897 | plugins: &PluginRegistry, |
| 1898 | ) -> Vec<command::ExtensionCommandEntry> { |
| 1899 | let Some(revision) = plugins.caller_selection() else { |
| 1900 | return self.commands_for_workspace(plugins.workspace()); |
| 1901 | }; |
| 1902 | let selected = self.shared.selection(revision.attachment_id); |
| 1903 | if selected.revision != Some(revision) { |
| 1904 | return Vec::new(); |
| 1905 | } |
| 1906 | let registry = self.shared.registry.lock().expect("registry lock"); |
| 1907 | let commands: Vec<_> = registry |
| 1908 | .live_commands() |
| 1909 | .into_iter() |
| 1910 | .filter(|command| { |
| 1911 | selected.includes( |
| 1912 | &command.owner.plugin_id, |
| 1913 | &command.content_hash, |
| 1914 | command.scope.as_ref(), |
| 1915 | ) |
| 1916 | }) |
| 1917 | .collect(); |
| 1918 | let mut counts = BTreeMap::new(); |
| 1919 | for command in &commands { |
| 1920 | *counts.entry(command.name.clone()).or_insert(0usize) += 1; |
| 1921 | } |
| 1922 | commands |
| 1923 | .into_iter() |
| 1924 | .filter(|command| counts[&command.name] == 1) |
| 1925 | .filter_map(|registration| { |
| 1926 | Some(command::ExtensionCommandEntry { |
| 1927 | authority: registry.authority_for(®istration.owner)?, |
| 1928 | registration, |
| 1929 | workspace: plugins.workspace().to_path_buf(), |
| 1930 | selection: Some(revision), |
| 1931 | }) |
| 1932 | }) |
| 1933 | .collect() |
| 1934 | } |
| 1935 | |
| 1936 | /// Reconcile without waiting (turn builds, session start, |
| 1937 | /// plugin changes). |
| 1938 | pub fn reconcile_in_background(self: &Arc<Self>) { |
| 1939 | if tokio::runtime::Handle::try_current().is_err() { |
| 1940 | return; |
| 1941 | } |
| 1942 | let manager = Arc::clone(self); |
| 1943 | let policy = activation::extension_host_policy_enabled(); |
| 1944 | tokio::spawn(async move { |
| 1945 | if let Err(error) = manager.reconcile_with_policy(policy).await { |
| 1946 | manager.shared.diagnostic(error); |
| 1947 | } |
| 1948 | }); |
| 1949 | } |
| 1950 | |
| 1951 | /// Like [`Self::reconcile_in_background`], after re-reading the plugin |
| 1952 | /// settings from the config file: for an explicit plugin reload. |
| 1953 | fn reload_and_reconcile_in_background(self: &Arc<Self>) { |
| 1954 | if tokio::runtime::Handle::try_current().is_err() { |
| 1955 | return; |
| 1956 | } |
| 1957 | let manager = Arc::clone(self); |
| 1958 | let policy = activation::extension_host_policy_enabled(); |
| 1959 | tokio::spawn(async move { |
| 1960 | manager.reload_plugin_settings().await; |
| 1961 | if let Err(error) = manager.reconcile_with_policy(policy).await { |
| 1962 | manager.shared.diagnostic(error); |
| 1963 | } |
| 1964 | }); |
| 1965 | } |
| 1966 | |
| 1967 | /// Reconcile host owners with the reviewed, enabled plugins that declare |
| 1968 | /// `native` entries in any attached engine's snapshot: revoke what no |
| 1969 | /// attachment desires any more or what changed (synchronously, then ask |
| 1970 | /// the host to tear down), spawn the host if needed, activate the rest. |
| 1971 | pub async fn reconcile(&self) -> Result<(), String> { |
| 1972 | self.reconcile_with_policy(activation::extension_host_policy_enabled()) |
| 1973 | .await |
| 1974 | } |
| 1975 | |
| 1976 | async fn reconcile_with_policy(&self, policy: bool) -> Result<(), String> { |
| 1977 | let shared = &self.shared; |
| 1978 | let _serial = shared.sync_lock.lock().await; |
| 1979 | let (desired, errors) = loop { |
| 1980 | let snapshots: Vec<(u64, Arc<PluginRegistry>)> = shared |
| 1981 | .attachments |
| 1982 | .lock() |
| 1983 | .expect("attachments lock") |
| 1984 | .iter() |
| 1985 | .map(|(id, state)| (*id, Arc::clone(&state.plugins))) |
| 1986 | .collect(); |
| 1987 | #[cfg(test)] |
| 1988 | let env_scope = crate::test_support::env_scope_ticket(); |
| 1989 | let scan = tokio::task::spawn_blocking(move || { |
| 1990 | #[cfg(test)] |
| 1991 | let _env_scope = crate::test_support::join_env_scope(env_scope); |
| 1992 | let _scope = activation::PolicyScope::propagate(policy); |
| 1993 | union_of_desired_owners(snapshots) |
| 1994 | }) |
| 1995 | .await |
| 1996 | .map_err(|error| format!("plugin scan failed: {error}"))?; |
| 1997 | let published = scan.publish(&mut shared.attachments.lock().expect("attachments lock")); |
| 1998 | if let Some(published) = published { |
| 1999 | break published; |
| 2000 | } |
| 2001 | // Attach, detach, or workspace refresh raced the blocking scan. |
| 2002 | // Rescan before changing either engine views or global owners. |
| 2003 | }; |
| 2004 | for error in errors { |
| 2005 | shared.diagnostic(error); |
| 2006 | } |
| 2007 | |
| 2008 | // The settings each desired plugin would be activated with now. Read |
| 2009 | // before the registry lock: the settings lock is never held with it. |
| 2010 | let wanted_config: BTreeMap<String, String> = { |
| 2011 | let configs = shared.plugin_configs.lock().expect("plugin configs lock"); |
| 2012 | desired |
| 2013 | .iter() |
| 2014 | .map(|(plugin_id, want)| { |
| 2015 | ( |
| 2016 | plugin_id.clone(), |
| 2017 | plugin_config::activation_hash(&configs.select(&want.plugin_name)), |
| 2018 | ) |
| 2019 | }) |
| 2020 | .collect() |
| 2021 | }; |
| 2022 | |
| 2023 | // 1. Revoke first — never waits for the host. Plugin-tier owners only: |
| 2024 | // the builtin tier's modules are not plugins, are never desired by an |
| 2025 | // attachment, and the table that wants them is fixed for the |
| 2026 | // manager's life, so no plugin change revokes one. |
| 2027 | let mut revoked: Vec<OwnerRef> = Vec::new(); |
| 2028 | let mut to_activate: Vec<(String, DesiredOwner)> = Vec::new(); |
| 2029 | { |
| 2030 | let mut registry = shared.registry.lock().expect("registry lock"); |
| 2031 | let existing: Vec<(String, PluginAuthority, String)> = registry |
| 2032 | .owners() |
| 2033 | .filter(|entry| entry.tier == HostTier::Plugin) |
| 2034 | .filter_map(|entry| { |
| 2035 | Some(( |
| 2036 | entry.owner.plugin_id.clone(), |
| 2037 | entry.authority.clone()?, |
| 2038 | entry.config_hash.clone(), |
| 2039 | )) |
| 2040 | }) |
| 2041 | .collect(); |
| 2042 | for (plugin_id, authority, config_hash) in &existing { |
| 2043 | // An explicit trust/enable transition revokes the persisted |
| 2044 | // authority even when the bytes stay identical. Refresh that |
| 2045 | // owner instead of retaining a tool that can only fail closed. |
| 2046 | // A changed config is a changed activation: a new generation. |
| 2047 | let keep = desired.get(plugin_id).is_some_and(|want| { |
| 2048 | want.authority.content_hash == authority.content_hash |
| 2049 | && want.authority.capability_hash == authority.capability_hash |
| 2050 | && want.authority.state_generation == authority.state_generation |
| 2051 | && want.authority.state_path == authority.state_path |
| 2052 | && wanted_config.get(plugin_id) == Some(config_hash) |
| 2053 | }); |
| 2054 | if !keep { |
| 2055 | if let Some(owner) = registry.revoke_owner(plugin_id) { |
| 2056 | revoked.push(owner); |
| 2057 | } |
| 2058 | registry.forget_owner(plugin_id); |
| 2059 | } |
| 2060 | } |
| 2061 | for (plugin_id, want) in desired { |
| 2062 | // A failed or faulted activation of these exact bytes is not |
| 2063 | // retried every turn; changed authority or an explicit |
| 2064 | // plugin mutation/reload permits another attempt. |
| 2065 | if registry |
| 2066 | .owner(&plugin_id) |
| 2067 | .is_none_or(|owner| owner.state == OwnerState::Active) |
| 2068 | { |
| 2069 | to_activate.push((plugin_id, want)); |
| 2070 | } |
| 2071 | } |
| 2072 | } |
| 2073 | let host = shared.ready_host(HostTier::Plugin).ok(); |
| 2074 | for owner in revoked { |
| 2075 | shared.core_calls.revoke_owner(&owner.plugin_id); |
| 2076 | shared.execution_broker.revoke_owner(&owner.plugin_id); |
| 2077 | shared.mcp_users.fetch_sub( |
| 2078 | shared.mcp_broker.revoke_owner(&owner.plugin_id), |
| 2079 | Ordering::SeqCst, |
| 2080 | ); |
| 2081 | shared.plugin_diagnostic( |
| 2082 | &owner.plugin_id, |
| 2083 | format!("extension `{}` revoked", owner.plugin_id), |
| 2084 | ); |
| 2085 | if let Some(host) = &host { |
| 2086 | host.revoke_calls_of(&owner.plugin_id); |
| 2087 | shared.deactivate_owner(host, &owner).await; |
| 2088 | } |
| 2089 | } |
| 2090 | |
| 2091 | // The plugin tier first, then the builtin tier; each only when it has |
| 2092 | // something to activate, so with no plugin and no module no host |
| 2093 | // starts. One tier's failure never stops the other's. |
| 2094 | let plugins = async { |
| 2095 | if to_activate.is_empty() { |
| 2096 | return Ok::<(), String>(()); |
| 2097 | } |
| 2098 | let host = self.ensure_host(HostTier::Plugin, policy).await?; |
| 2099 | for (plugin_id, want) in to_activate { |
| 2100 | if host.has_exited() { |
| 2101 | break; |
| 2102 | } |
| 2103 | let content_hash = want.authority.content_hash.clone(); |
| 2104 | self.activate_owner( |
| 2105 | &host, |
| 2106 | Activation { |
| 2107 | tier: HostTier::Plugin, |
| 2108 | owner_id: plugin_id, |
| 2109 | name: want.plugin_name, |
| 2110 | authority: Some(want.authority), |
| 2111 | content_hash, |
| 2112 | entries: want.entries, |
| 2113 | }, |
| 2114 | ) |
| 2115 | .await; |
| 2116 | } |
| 2117 | Ok(()) |
| 2118 | } |
| 2119 | .await; |
| 2120 | // MCP selection admits tier-0 demand independently of Native plugins. |
| 2121 | // Harness demand still requires the captured Native policy below. |
| 2122 | let builtin = if policy |
| 2123 | || shared.mcp_users.load(Ordering::SeqCst) > 0 |
| 2124 | || shared.stock_users.load(Ordering::SeqCst) > 0 |
| 2125 | { |
| 2126 | self.reconcile_builtin(policy).await |
| 2127 | } else { |
| 2128 | Ok(()) |
| 2129 | }; |
| 2130 | match (plugins, builtin) { |
| 2131 | (Err(error), Err(other)) => { |
| 2132 | shared.diagnostic(other); |
| 2133 | Err(error) |
| 2134 | } |
| 2135 | (Err(error), Ok(())) | (Ok(()), Err(error)) => Err(error), |
| 2136 | (Ok(()), Ok(())) => Ok(()), |
| 2137 | } |
| 2138 | } |
| 2139 | |
| 2140 | /// One real host consumer asks for the pinned MCP module. Rust mode never |
| 2141 | /// starts tier 0. This uses the same manager, supervision and owner registry. |
| 2142 | async fn ensure_mcp_builtin( |
| 2143 | &self, |
| 2144 | ) -> Result<(Arc<HostProcess>, protocol::OwnerRef, u64), String> { |
| 2145 | self.ensure_demand_builtin("host:mcp", activation::extension_host_policy_enabled()) |
| 2146 | .await |
| 2147 | } |
| 2148 | /// Pure Core adapters demand the same pinned harness without enabling Native plugins. |
| 2149 | async fn ensure_stock_builtin( |
| 2150 | &self, |
| 2151 | ) -> Result<(Arc<HostProcess>, protocol::OwnerRef, u64), String> { |
| 2152 | self.ensure_demand_builtin("host:harness", activation::extension_host_policy_enabled()) |
| 2153 | .await |
| 2154 | } |
| 2155 | async fn ensure_harness_builtin( |
| 2156 | &self, |
| 2157 | ) -> Result<(Arc<HostProcess>, protocol::OwnerRef, u64), String> { |
| 2158 | let native_policy = activation::extension_host_policy_enabled(); |
| 2159 | if !native_policy { |
| 2160 | return Err( |
| 2161 | "Builtin backend requires features.extension_host; no Rust fallback is allowed" |
| 2162 | .to_string(), |
| 2163 | ); |
| 2164 | } |
| 2165 | self.ensure_demand_builtin("host:harness", native_policy) |
| 2166 | .await |
| 2167 | } |
| 2168 | async fn ensure_demand_builtin( |
| 2169 | &self, |
| 2170 | owner_id: &str, |
| 2171 | native_policy: bool, |
| 2172 | ) -> Result<(Arc<HostProcess>, protocol::OwnerRef, u64), String> { |
| 2173 | let _serial = self.shared.sync_lock.lock().await; |
| 2174 | let root = host_root(&self.shared.options)?; |
| 2175 | materialize_bundle(&root)?; |
| 2176 | // A crashed tier-0 process cannot carry its old owner into a new |
| 2177 | // process. Retire only this builtin before minting a fresh generation. |
| 2178 | let stale = self |
| 2179 | .shared |
| 2180 | .registry |
| 2181 | .lock() |
| 2182 | .expect("registry lock") |
| 2183 | .owner(owner_id) |
| 2184 | .is_some_and(|entry| entry.state != OwnerState::Active); |
| 2185 | if stale { |
| 2186 | self.shared.core_calls.revoke_owner(owner_id); |
| 2187 | self.shared.execution_broker.revoke_owner(owner_id); |
| 2188 | self.shared.mcp_users.fetch_sub( |
| 2189 | self.shared.mcp_broker.revoke_owner(owner_id), |
| 2190 | Ordering::SeqCst, |
| 2191 | ); |
| 2192 | self.shared |
| 2193 | .registry |
| 2194 | .lock() |
| 2195 | .expect("registry lock") |
| 2196 | .forget_owner(owner_id); |
| 2197 | } |
| 2198 | // Replay the actual Native policy after a Builtin crash. MCP demand |
| 2199 | // must never switch optional third-party activation on. |
| 2200 | self.reconcile_builtin(native_policy).await?; |
| 2201 | let host = self |
| 2202 | .shared |
| 2203 | .ready_host(HostTier::Builtin) |
| 2204 | .map_err(|status| host_down(HostTier::Builtin, &status))?; |
| 2205 | let owner = self |
| 2206 | .shared |
| 2207 | .registry |
| 2208 | .lock() |
| 2209 | .expect("registry lock") |
| 2210 | .owner(owner_id) |
| 2211 | .filter(|entry| entry.state == OwnerState::Active) |
| 2212 | .map(|entry| entry.owner.clone()) |
| 2213 | .ok_or_else(|| "demanded builtin failed to activate".to_string())?; |
| 2214 | let generation = self.shared.builtin.host_generation.load(Ordering::SeqCst); |
| 2215 | Ok((host, owner, generation)) |
| 2216 | } |
| 2217 | |
| 2218 | async fn reconcile_builtin(&self, policy: bool) -> Result<(), String> { |
| 2219 | let shared = &self.shared; |
| 2220 | let wanted: Vec<&'static BuiltinModule> = { |
| 2221 | let registry = shared.registry.lock().expect("registry lock"); |
| 2222 | shared |
| 2223 | .builtin_modules |
| 2224 | .iter() |
| 2225 | .filter(|module| { |
| 2226 | (module.id != "mcp" || shared.mcp_users.load(Ordering::SeqCst) > 0) |
| 2227 | && (module.id != "harness" |
| 2228 | || (policy && shared.harness_users.load(Ordering::SeqCst) > 0) |
| 2229 | || shared.stock_users.load(Ordering::SeqCst) > 0) |
| 2230 | && registry.owner(&module.owner_id()).is_none() |
| 2231 | }) |
| 2232 | .collect() |
| 2233 | }; |
| 2234 | if wanted.is_empty() { |
| 2235 | return Ok(()); |
| 2236 | } |
| 2237 | let root = host_root(&shared.options)?; |
| 2238 | let mut activations = Vec::new(); |
| 2239 | for module in wanted { |
| 2240 | // The module's source, accepted only if its SHA-256 is the one the |
| 2241 | // table pins: the path and digest to activate it with. |
| 2242 | let source = { |
| 2243 | let root = root.clone(); |
| 2244 | tokio::task::spawn_blocking(move || { |
| 2245 | let path = builtin_source_path(&root, module); |
| 2246 | let bytes = std::fs::read(&path).map_err(|error| { |
| 2247 | format!( |
| 2248 | "cannot read the source of built-in module `{}` at {}: {error}", |
| 2249 | module.id, |
| 2250 | path.display() |
| 2251 | ) |
| 2252 | })?; |
| 2253 | let digest = hex(Sha256::digest(&bytes)); |
| 2254 | if digest != module.source_sha256 { |
| 2255 | return Err(format!( |
| 2256 | "the source of built-in module `{}` at {} has sha256 {digest}, not the {} the core pins", |
| 2257 | module.id, |
| 2258 | path.display(), |
| 2259 | module.source_sha256 |
| 2260 | )); |
| 2261 | } |
| 2262 | Ok((path, digest)) |
| 2263 | }) |
| 2264 | .await |
| 2265 | .map_err(|error| format!("built-in module scan failed: {error}"))? |
| 2266 | }; |
| 2267 | let owner_id = module.owner_id(); |
| 2268 | match source { |
| 2269 | Ok(entry) => activations.push(Activation { |
| 2270 | tier: HostTier::Builtin, |
| 2271 | owner_id, |
| 2272 | name: module.id.to_string(), |
| 2273 | authority: None, |
| 2274 | content_hash: module.source_sha256.to_string(), |
| 2275 | entries: vec![entry], |
| 2276 | }), |
| 2277 | Err(reason) => { |
| 2278 | let mut registry = shared.registry.lock().expect("registry lock"); |
| 2279 | if let Ok(owner) = registry.begin_owner( |
| 2280 | HostTier::Builtin, |
| 2281 | &owner_id, |
| 2282 | module.id, |
| 2283 | None, |
| 2284 | module.source_sha256, |
| 2285 | ) { |
| 2286 | registry.mark_failed(&owner, OwnerState::Failed(reason.clone())); |
| 2287 | } |
| 2288 | drop(registry); |
| 2289 | shared.plugin_diagnostic( |
| 2290 | &owner_id, |
| 2291 | format!( |
| 2292 | "built-in module `{}` failed to activate: {reason}", |
| 2293 | module.id |
| 2294 | ), |
| 2295 | ); |
| 2296 | } |
| 2297 | } |
| 2298 | } |
| 2299 | if activations.is_empty() { |
| 2300 | return Ok(()); |
| 2301 | } |
| 2302 | let host = self.ensure_host(HostTier::Builtin, policy).await?; |
| 2303 | for activation in activations { |
| 2304 | if host.has_exited() { |
| 2305 | break; |
| 2306 | } |
| 2307 | self.activate_owner(&host, activation).await; |
| 2308 | } |
| 2309 | Ok(()) |
| 2310 | } |
| 2311 | |
| 2312 | async fn activate_owner(&self, host: &Arc<HostProcess>, want: Activation) { |
| 2313 | let shared = &self.shared; |
| 2314 | let plugin_id = want.owner_id.as_str(); |
| 2315 | let begun = { |
| 2316 | let mut registry = shared.registry.lock().expect("registry lock"); |
| 2317 | if let Some(owner) = registry |
| 2318 | .owner(plugin_id) |
| 2319 | .filter(|entry| entry.state == OwnerState::Active) |
| 2320 | { |
| 2321 | Ok(owner.owner.clone()) |
| 2322 | } else { |
| 2323 | registry.begin_owner( |
| 2324 | want.tier, |
| 2325 | plugin_id, |
| 2326 | &want.name, |
| 2327 | want.authority.clone(), |
| 2328 | &want.content_hash, |
| 2329 | ) |
| 2330 | } |
| 2331 | }; |
| 2332 | let owner = match begun { |
| 2333 | Ok(owner) => owner, |
| 2334 | Err(reason) => { |
| 2335 | shared.plugin_diagnostic( |
| 2336 | plugin_id, |
| 2337 | format!("extension `{}` was not activated: {reason}", want.name), |
| 2338 | ); |
| 2339 | return; |
| 2340 | } |
| 2341 | }; |
| 2342 | // What the owner is given: its settings and its own directory. A |
| 2343 | // config that is refused, or a directory that cannot be made, fails the |
| 2344 | // activation with the reason; the plugin is never started with |
| 2345 | // different settings than the user wrote. A built-in module takes no |
| 2346 | // user settings: `[plugins]` is keyed by plugin name and is not read |
| 2347 | // for tier 0. |
| 2348 | let selection = match want.tier { |
| 2349 | HostTier::Plugin => shared |
| 2350 | .plugin_configs |
| 2351 | .lock() |
| 2352 | .expect("plugin configs lock") |
| 2353 | .select(&want.name), |
| 2354 | HostTier::Builtin => Ok(plugin_config::PluginConfig::empty()), |
| 2355 | }; |
| 2356 | shared |
| 2357 | .registry |
| 2358 | .lock() |
| 2359 | .expect("registry lock") |
| 2360 | .set_config_hash(&owner, &plugin_config::activation_hash(&selection)); |
| 2361 | let context = match selection { |
| 2362 | Ok(config) => self |
| 2363 | .owner_data_dir(want.tier, plugin_id, &want.name) |
| 2364 | .await |
| 2365 | .map(|data_dir| (config, data_dir)), |
| 2366 | Err(reason) => Err(reason), |
| 2367 | }; |
| 2368 | let (config, data_dir) = match context { |
| 2369 | Ok(context) => context, |
| 2370 | Err(reason) => { |
| 2371 | shared |
| 2372 | .registry |
| 2373 | .lock() |
| 2374 | .expect("registry lock") |
| 2375 | .mark_failed(&owner, OwnerState::Failed(reason.clone())); |
| 2376 | shared.plugin_diagnostic( |
| 2377 | plugin_id, |
| 2378 | format!("extension `{}` failed to activate: {reason}", want.name), |
| 2379 | ); |
| 2380 | return; |
| 2381 | } |
| 2382 | }; |
| 2383 | // Retire no-longer-wanted entry handles before awaiting host disposal. |
| 2384 | if want.tier == HostTier::Plugin { |
| 2385 | let scopes: Vec<_> = { |
| 2386 | let registry = shared.registry.lock().expect("registry lock"); |
| 2387 | registry |
| 2388 | .owner(plugin_id) |
| 2389 | .map(|entry| { |
| 2390 | entry |
| 2391 | .scopes |
| 2392 | .keys() |
| 2393 | .filter(|scope| { |
| 2394 | !want.entries.iter().any(|(path, hash)| { |
| 2395 | scope.path == path.to_string_lossy() && scope.sha256 == *hash |
| 2396 | }) |
| 2397 | }) |
| 2398 | .cloned() |
| 2399 | .collect() |
| 2400 | }) |
| 2401 | .unwrap_or_default() |
| 2402 | }; |
| 2403 | for scope in scopes { |
| 2404 | let handles = shared |
| 2405 | .registry |
| 2406 | .lock() |
| 2407 | .expect("registry lock") |
| 2408 | .revoke_scope(&owner, &scope); |
| 2409 | shared.core_calls.revoke_scope(plugin_id, &scope); |
| 2410 | shared.execution_broker.revoke_scope(plugin_id, &scope); |
| 2411 | host.revoke_calls_for_handles(plugin_id, &handles); |
| 2412 | shared.deactivate_entry(host, &owner, Some(scope)).await; |
| 2413 | } |
| 2414 | } |
| 2415 | #[cfg(windows)] |
| 2416 | if want.tier == HostTier::Plugin { |
| 2417 | let admission = match want.authority.clone() { |
| 2418 | Some(authority) => host.admit_windows_root(authority).await, |
| 2419 | None => Err("Native owner has no reviewed bundle authority".into()), |
| 2420 | }; |
| 2421 | if let Err(reason) = admission { |
| 2422 | shared |
| 2423 | .registry |
| 2424 | .lock() |
| 2425 | .expect("registry lock") |
| 2426 | .mark_failed(&owner, OwnerState::Failed(reason.clone())); |
| 2427 | shared.plugin_diagnostic( |
| 2428 | plugin_id, |
| 2429 | format!( |
| 2430 | "extension `{}` failed Windows admission: {reason}", |
| 2431 | want.name |
| 2432 | ), |
| 2433 | ); |
| 2434 | shared.deactivate_owner(host, &owner).await; |
| 2435 | return; |
| 2436 | } |
| 2437 | } |
| 2438 | let mut failure = None; |
| 2439 | let mut native_failure = None; |
| 2440 | // A plugin may declare several `native` entries. Each is activated in |
| 2441 | // declaration order under one owner. Each entry has an independent |
| 2442 | // scope: a failing entry withdraws its own handles and cleanup, while |
| 2443 | // already-active sibling entries remain usable. |
| 2444 | let mut tools = BTreeSet::new(); |
| 2445 | let mut commands = BTreeSet::new(); |
| 2446 | for (path, sha256) in &want.entries { |
| 2447 | let scope = (want.tier == HostTier::Plugin).then(|| EntryRef { |
| 2448 | path: path.to_string_lossy().into_owned(), |
| 2449 | sha256: sha256.clone(), |
| 2450 | }); |
| 2451 | if let Some(scope) = &scope { |
| 2452 | let admitted = shared |
| 2453 | .registry |
| 2454 | .lock() |
| 2455 | .expect("registry lock") |
| 2456 | .begin_scope(&owner, scope.clone()); |
| 2457 | if admitted.is_err() { |
| 2458 | continue; |
| 2459 | } |
| 2460 | } |
| 2461 | let request = CoreRequest::Activate(ActivateParams { |
| 2462 | scope: scope.clone(), |
| 2463 | owner: owner.clone(), |
| 2464 | plugin_name: want.name.clone(), |
| 2465 | entry: EntryRef { |
| 2466 | path: path.to_string_lossy().into_owned(), |
| 2467 | sha256: sha256.clone(), |
| 2468 | }, |
| 2469 | config: config.value.clone(), |
| 2470 | data_dir: Some(data_dir.clone()), |
| 2471 | }); |
| 2472 | let outcome = host |
| 2473 | .call(request, Some(plugin_id.to_string())) |
| 2474 | .await |
| 2475 | .map_err(|error| error.to_string()) |
| 2476 | .and_then(|value| { |
| 2477 | serde_json::from_value::<ActivateResult>(value) |
| 2478 | .map_err(|error| format!("malformed activation answer: {error}")) |
| 2479 | }) |
| 2480 | .and_then(|result| match result { |
| 2481 | ActivateResult::Ok { tools, commands } => Ok((tools, commands)), |
| 2482 | ActivateResult::Failed { diagnostic } => Err(diagnostic), |
| 2483 | }); |
| 2484 | match outcome { |
| 2485 | Ok((names, slash)) => { |
| 2486 | tools.extend(names); |
| 2487 | commands.extend(slash); |
| 2488 | if let Some(scope) = &scope { |
| 2489 | shared |
| 2490 | .registry |
| 2491 | .lock() |
| 2492 | .expect("registry lock") |
| 2493 | .mark_scope_active(&owner, scope); |
| 2494 | } |
| 2495 | } |
| 2496 | Err(reason) => { |
| 2497 | if let Some(scope) = scope { |
| 2498 | shared |
| 2499 | .registry |
| 2500 | .lock() |
| 2501 | .expect("registry lock") |
| 2502 | .fail_scope(&owner, &scope); |
| 2503 | shared.core_calls.revoke_scope(plugin_id, &scope); |
| 2504 | shared.execution_broker.revoke_scope(plugin_id, &scope); |
| 2505 | native_failure.get_or_insert_with(|| reason.clone()); |
| 2506 | shared.plugin_diagnostic(plugin_id, format!("Native entry failed to activate: {reason}; other selected entries remain live")); |
| 2507 | shared.deactivate_entry(host, &owner, Some(scope)).await; |
| 2508 | } else { |
| 2509 | failure = Some(reason); |
| 2510 | break; |
| 2511 | } |
| 2512 | } |
| 2513 | } |
| 2514 | } |
| 2515 | if want.tier == HostTier::Plugin |
| 2516 | && shared |
| 2517 | .registry |
| 2518 | .lock() |
| 2519 | .expect("registry lock") |
| 2520 | .owner(plugin_id) |
| 2521 | .is_none_or(|entry| { |
| 2522 | !entry |
| 2523 | .scopes |
| 2524 | .values() |
| 2525 | .any(|state| *state == OwnerState::Active) |
| 2526 | }) |
| 2527 | { |
| 2528 | let reason = |
| 2529 | native_failure.unwrap_or_else(|| "No selected Native entry became active".into()); |
| 2530 | shared |
| 2531 | .registry |
| 2532 | .lock() |
| 2533 | .expect("registry lock") |
| 2534 | .mark_failed(&owner, OwnerState::Failed(reason.clone())); |
| 2535 | shared.plugin_diagnostic( |
| 2536 | plugin_id, |
| 2537 | format!("extension `{}` failed to activate: {reason}", want.name), |
| 2538 | ); |
| 2539 | // Per-entry acknowledgements above account for incomplete cleanup. |
| 2540 | // Also retire the now-empty parent so its owner record cannot linger. |
| 2541 | shared.deactivate_owner(host, &owner).await; |
| 2542 | return; |
| 2543 | } |
| 2544 | match failure { |
| 2545 | None => { |
| 2546 | let active = shared |
| 2547 | .registry |
| 2548 | .lock() |
| 2549 | .expect("registry lock") |
| 2550 | .mark_active(&owner); |
| 2551 | if active { |
| 2552 | shared.plugin_diagnostic( |
| 2553 | plugin_id, |
| 2554 | format!( |
| 2555 | "extension `{}` active (tools: {}; commands: {})", |
| 2556 | want.name, |
| 2557 | tools.iter().cloned().collect::<Vec<_>>().join(", "), |
| 2558 | commands |
| 2559 | .iter() |
| 2560 | .map(|name| format!("/{name}")) |
| 2561 | .collect::<Vec<_>>() |
| 2562 | .join(", ") |
| 2563 | ), |
| 2564 | ); |
| 2565 | } |
| 2566 | } |
| 2567 | Some(reason) => { |
| 2568 | // All-or-nothing: drop anything half-registered on this side, |
| 2569 | // and dispose any entry the host did activate. |
| 2570 | shared |
| 2571 | .registry |
| 2572 | .lock() |
| 2573 | .expect("registry lock") |
| 2574 | .mark_failed(&owner, OwnerState::Failed(reason.clone())); |
| 2575 | shared.plugin_diagnostic( |
| 2576 | plugin_id, |
| 2577 | format!("extension `{}` failed to activate: {reason}", want.name), |
| 2578 | ); |
| 2579 | shared.deactivate_owner(host, &owner).await; |
| 2580 | } |
| 2581 | } |
| 2582 | } |
| 2583 | |
| 2584 | /// Create (if needed) and name an owner's own directory inside its tier's |
| 2585 | /// data dir, the host sandbox's writable root: a plugin's at its |
| 2586 | /// long-standing path, a built-in module's under the builtin tier's. |
| 2587 | async fn owner_data_dir( |
| 2588 | &self, |
| 2589 | tier: HostTier, |
| 2590 | plugin_id: &str, |
| 2591 | plugin_name: &str, |
| 2592 | ) -> Result<String, String> { |
| 2593 | let path = supervisor::owner_data_dir( |
| 2594 | &host_root(&self.shared.options)?, |
| 2595 | tier, |
| 2596 | plugin_id, |
| 2597 | plugin_name, |
| 2598 | ); |
| 2599 | let mut builder = tokio::fs::DirBuilder::new(); |
| 2600 | builder.recursive(true); |
| 2601 | #[cfg(unix)] |
| 2602 | builder.mode(0o700); |
| 2603 | builder.create(&path).await.map_err(|error| { |
| 2604 | format!( |
| 2605 | "cannot create the plugin data directory {}: {error}", |
| 2606 | path.display() |
| 2607 | ) |
| 2608 | })?; |
| 2609 | path.to_str().map(str::to_owned).ok_or_else(|| { |
| 2610 | format!( |
| 2611 | "the plugin data directory {} is not valid UTF-8", |
| 2612 | path.display() |
| 2613 | ) |
| 2614 | }) |
| 2615 | } |
| 2616 | |
| 2617 | /// `tier`'s host, started (lazily, once per generation) if it is idle. |
| 2618 | async fn ensure_host(&self, tier: HostTier, policy: bool) -> Result<Arc<HostProcess>, String> { |
| 2619 | let shared = &self.shared; |
| 2620 | let runtime = shared.tier_runtime(tier); |
| 2621 | let label = tier.host_label(); |
| 2622 | let mut generation = { |
| 2623 | let mut slot = runtime.host.lock().expect("host lock"); |
| 2624 | match &*slot { |
| 2625 | HostSlot::Ready(host) if !host.has_exited() && !host.is_retiring() => { |
| 2626 | return Ok(Arc::clone(host)); |
| 2627 | } |
| 2628 | HostSlot::Idle => {} |
| 2629 | other => return Err(host_down(tier, &other.status())), |
| 2630 | } |
| 2631 | *slot = HostSlot::Starting; |
| 2632 | let mut supervision = runtime.supervision.lock().expect("supervision lock"); |
| 2633 | supervision.last_start = Some(Instant::now()); |
| 2634 | supervision.policy = policy; |
| 2635 | supervision.launch_failed = false; |
| 2636 | supervision.dirty_teardowns.clear(); |
| 2637 | supervision.dirty_restart_pending = false; |
| 2638 | supervision.planned_restart = None; |
| 2639 | runtime.host_generation.fetch_add(1, Ordering::SeqCst) + 1 |
| 2640 | }; |
| 2641 | // At most two launches: under `auto`, before anything is pinned, a |
| 2642 | // Bun that cannot start is reported once and Node is resolved for the |
| 2643 | // rest of the session. An explicit `bun` or `node` never falls back. |
| 2644 | let spawned = loop { |
| 2645 | runtime.spawn_attempts.fetch_add(1, Ordering::SeqCst); |
| 2646 | let options = shared.options.clone(); |
| 2647 | let (pinned, bun_failed) = { |
| 2648 | let runtime = shared.runtime.lock().expect("runtime lock"); |
| 2649 | (runtime.pinned.clone(), runtime.bun_failed) |
| 2650 | }; |
| 2651 | #[cfg(test)] |
| 2652 | let env_scope = crate::test_support::env_scope_ticket(); |
| 2653 | let prepared = tokio::task::spawn_blocking(move || { |
| 2654 | #[cfg(test)] |
| 2655 | let _env_scope = crate::test_support::join_env_scope(env_scope); |
| 2656 | prepare_launch(&options, tier, pinned, bun_failed) |
| 2657 | }) |
| 2658 | .await |
| 2659 | .map_err(|error| format!("extension host preparation failed: {error}")) |
| 2660 | .and_then(|result| result); |
| 2661 | let (launch, summary) = match prepared { |
| 2662 | Ok(prepared) => prepared, |
| 2663 | Err(error) => break Err(error), |
| 2664 | }; |
| 2665 | let events: Arc<dyn HostEvents> = Arc::new(Events { |
| 2666 | shared: Arc::downgrade(shared), |
| 2667 | tier, |
| 2668 | generation, |
| 2669 | }); |
| 2670 | let reason = |
| 2671 | match HostProcess::spawn(generation, &launch, bundle_sha256(), events).await { |
| 2672 | Ok(host) => break Ok((host, summary)), |
| 2673 | Err(reason) => reason, |
| 2674 | }; |
| 2675 | if shared.options.runtime != crate::config::ExtensionHostRuntime::Auto |
| 2676 | || launch.runtime.kind != crate::dependencies::HostRuntimeKind::Bun |
| 2677 | || summary.is_none() |
| 2678 | { |
| 2679 | break Err(reason); |
| 2680 | } |
| 2681 | shared.runtime.lock().expect("runtime lock").bun_failed = true; |
| 2682 | shared.diagnostic(format!( |
| 2683 | "{label}: Bun {} at {} failed to start ({reason}); runtime = \"auto\" uses Node for the rest of this session", |
| 2684 | launch.runtime.version_string(), |
| 2685 | launch.runtime.path.display() |
| 2686 | )); |
| 2687 | // A fresh generation, so the failed Bun host's exit report can |
| 2688 | // never be taken for the Node host's. |
| 2689 | generation = { |
| 2690 | let slot = runtime.host.lock().expect("host lock"); |
| 2691 | if runtime.host_generation.load(Ordering::SeqCst) != generation |
| 2692 | || !matches!(&*slot, HostSlot::Starting) |
| 2693 | { |
| 2694 | break Err(format!("{label} startup was superseded")); |
| 2695 | } |
| 2696 | runtime.host_generation.fetch_add(1, Ordering::SeqCst) + 1 |
| 2697 | }; |
| 2698 | }; |
| 2699 | let mut slot = runtime.host.lock().expect("host lock"); |
| 2700 | if runtime.host_generation.load(Ordering::SeqCst) != generation |
| 2701 | || !matches!(&*slot, HostSlot::Starting) |
| 2702 | { |
| 2703 | drop(slot); |
| 2704 | if let Ok((host, _)) = spawned { |
| 2705 | host.terminate("host startup superseded".into()); |
| 2706 | } |
| 2707 | return Err(format!("{label} startup was superseded")); |
| 2708 | } |
| 2709 | match spawned { |
| 2710 | Ok((host, summary)) => { |
| 2711 | *slot = HostSlot::Ready(Arc::clone(&host)); |
| 2712 | drop(slot); |
| 2713 | // Pinned only once a host on this runtime has completed the |
| 2714 | // handshake; every restart then reuses it. |
| 2715 | if let Some(summary) = summary { |
| 2716 | let mut runtime = shared.runtime.lock().expect("runtime lock"); |
| 2717 | if runtime.pinned.is_none() { |
| 2718 | runtime.pinned = Some((host.runtime.clone(), summary.clone())); |
| 2719 | drop(runtime); |
| 2720 | shared.diagnostic(format!("extension host runtime: {summary}")); |
| 2721 | } |
| 2722 | } |
| 2723 | if host.has_exited() { |
| 2724 | Events { |
| 2725 | shared: Arc::downgrade(shared), |
| 2726 | tier, |
| 2727 | generation, |
| 2728 | } |
| 2729 | .exited( |
| 2730 | generation, |
| 2731 | "exited immediately after handshake".into(), |
| 2732 | host.stderr_tail(), |
| 2733 | ); |
| 2734 | return Err(format!("{label} exited immediately after handshake")); |
| 2735 | } |
| 2736 | monitor_host(shared, &host, generation); |
| 2737 | shared.diagnostic(format!( |
| 2738 | "{label} started (pid {}, {} {}, sandbox {})", |
| 2739 | host.pid |
| 2740 | .map_or_else(|| "?".to_string(), |pid| pid.to_string()), |
| 2741 | host.runtime.kind.name(), |
| 2742 | host.runtime_version.get().map_or("?", String::as_str), |
| 2743 | host.sandbox.label() |
| 2744 | )); |
| 2745 | Ok(host) |
| 2746 | } |
| 2747 | Err(reason) => { |
| 2748 | runtime |
| 2749 | .supervision |
| 2750 | .lock() |
| 2751 | .expect("supervision lock") |
| 2752 | .launch_failed = true; |
| 2753 | *slot = HostSlot::Failed { |
| 2754 | reason: format!("start failed: {reason}"), |
| 2755 | stderr_tail: String::new(), |
| 2756 | }; |
| 2757 | drop(slot); |
| 2758 | shared.diagnostic(format!("{label} failed to start: {reason}")); |
| 2759 | Err(reason) |
| 2760 | } |
| 2761 | } |
| 2762 | } |
| 2763 | |
| 2764 | /// Bounded shutdown of each tier's host process, if one is running. |
| 2765 | /// Production has no such call: the hosts are shared by every engine in |
| 2766 | /// the process, so no single engine's shutdown may stop them. When this |
| 2767 | /// process ends a host sees stdin EOF and kills its own process tree; if a |
| 2768 | /// plugin blocks its event loop, its watchdog thread does so when the |
| 2769 | /// parent changes. |
| 2770 | #[cfg(test)] |
| 2771 | pub async fn shutdown(&self) { |
| 2772 | for tier in HostTier::ALL { |
| 2773 | let runtime = self.shared.tier_runtime(tier); |
| 2774 | let host = { |
| 2775 | let mut slot = runtime.host.lock().expect("host lock"); |
| 2776 | runtime.host_generation.fetch_add(1, Ordering::SeqCst); |
| 2777 | runtime |
| 2778 | .supervision |
| 2779 | .lock() |
| 2780 | .expect("supervision lock") |
| 2781 | .retry_ticket += 1; |
| 2782 | match std::mem::replace(&mut *slot, HostSlot::Idle) { |
| 2783 | HostSlot::Ready(host) | HostSlot::Unresponsive(host) => Some(host), |
| 2784 | _ => None, |
| 2785 | } |
| 2786 | }; |
| 2787 | if let Some(host) = host { |
| 2788 | self.shared |
| 2789 | .registry |
| 2790 | .lock() |
| 2791 | .expect("registry lock") |
| 2792 | .revoke_all(tier, "extension host shut down"); |
| 2793 | host.shutdown().await; |
| 2794 | } |
| 2795 | } |
| 2796 | } |
| 2797 | |
| 2798 | #[cfg(test)] |
| 2799 | pub(crate) fn host_requests_started(&self) -> Option<u64> { |
| 2800 | self.shared |
| 2801 | .ready_host(HostTier::Plugin) |
| 2802 | .ok() |
| 2803 | .map(|host| host.requests_started()) |
| 2804 | } |
| 2805 | |
| 2806 | #[cfg(test)] |
| 2807 | pub(crate) fn host_pid(&self) -> Option<u32> { |
| 2808 | self.shared |
| 2809 | .ready_host(HostTier::Plugin) |
| 2810 | .ok() |
| 2811 | .and_then(|host| host.pid) |
| 2812 | } |
| 2813 | } |
| 2814 | |
| 2815 | /// Human-readable host section for `/plugin`. |
| 2816 | pub(crate) fn render_status(manager: &ExtensionHostManager) -> String { |
| 2817 | use std::fmt::Write as _; |
| 2818 | let mut out = String::from("Extension host (experimental): "); |
| 2819 | let status = manager.status(); |
| 2820 | match status.clone() { |
| 2821 | HostStatus::Ready { |
| 2822 | pid, |
| 2823 | runtime, |
| 2824 | runtime_version, |
| 2825 | sandbox, |
| 2826 | memory: _, |
| 2827 | } => { |
| 2828 | let _ = write!( |
| 2829 | out, |
| 2830 | "running · pid {} · {runtime} {runtime_version} · {sandbox}", |
| 2831 | pid.map_or_else(|| "?".to_string(), |pid| pid.to_string()), |
| 2832 | ); |
| 2833 | } |
| 2834 | HostStatus::Failed { stderr_tail, .. } => { |
| 2835 | let _ = write!(out, "{status}"); |
| 2836 | let tail = stderr_tail.trim(); |
| 2837 | if !tail.is_empty() { |
| 2838 | let start = tail |
| 2839 | .char_indices() |
| 2840 | .rev() |
| 2841 | .nth(599) |
| 2842 | .map_or(0, |(index, _)| index); |
| 2843 | let _ = write!(out, "\n stderr: {}", &tail[start..]); |
| 2844 | } |
| 2845 | } |
| 2846 | down => { |
| 2847 | let _ = write!(out, "{down}"); |
| 2848 | } |
| 2849 | } |
| 2850 | if let Some(summary) = manager.runtime_summary() { |
| 2851 | let _ = write!(out, "\n runtime: {summary}"); |
| 2852 | // How the running host's cap is enforced, as its handshake settled |
| 2853 | // it; with no host running there is nothing enforced to describe. |
| 2854 | if let HostStatus::Ready { memory, .. } = &status { |
| 2855 | let cap = manager.shared.options.supervision.memory_cap; |
| 2856 | let _ = write!(out, " · {}", memory.describe(cap)); |
| 2857 | } |
| 2858 | } |
| 2859 | let _ = write!( |
| 2860 | out, |
| 2861 | "\n spawn attempts: {} · engines attached: {}", |
| 2862 | manager.spawn_attempts(), |
| 2863 | manager.attached_engines() |
| 2864 | ); |
| 2865 | // The built-in host is shown only once something has started it (nothing |
| 2866 | // does yet); the lists and the shared-process count below are the plugin |
| 2867 | // host's. |
| 2868 | match manager.tier_status(HostTier::Builtin) { |
| 2869 | HostStatus::Idle => {} |
| 2870 | HostStatus::Ready { |
| 2871 | pid, |
| 2872 | runtime, |
| 2873 | runtime_version, |
| 2874 | sandbox, |
| 2875 | memory: _, |
| 2876 | } => { |
| 2877 | let _ = write!( |
| 2878 | out, |
| 2879 | "\n built-in host (tier 0): running · pid {} · {runtime} {runtime_version} · {sandbox}", |
| 2880 | pid.map_or_else(|| "?".to_string(), |pid| pid.to_string()), |
| 2881 | ); |
| 2882 | } |
| 2883 | other => { |
| 2884 | let _ = write!(out, "\n built-in host (tier 0): {other}"); |
| 2885 | } |
| 2886 | } |
| 2887 | let (tools, commands, owners) = { |
| 2888 | let registry = manager.shared.registry.lock().expect("registry lock"); |
| 2889 | let owners = registry |
| 2890 | .owners() |
| 2891 | .filter(|entry| entry.tier == HostTier::Plugin && entry.state == OwnerState::Active) |
| 2892 | .count(); |
| 2893 | let tools: Vec<_> = registry |
| 2894 | .live_tools() |
| 2895 | .into_iter() |
| 2896 | .filter(|tool| tool.tier == HostTier::Plugin) |
| 2897 | .collect(); |
| 2898 | let commands: Vec<_> = registry |
| 2899 | .live_commands() |
| 2900 | .into_iter() |
| 2901 | .filter(|command| command.tier == HostTier::Plugin) |
| 2902 | .collect(); |
| 2903 | (tools, commands, owners) |
| 2904 | }; |
| 2905 | if owners > 1 { |
| 2906 | let _ = write!( |
| 2907 | out, |
| 2908 | "\n {owners} plugins share this one host process and can alter each other's behaviour" |
| 2909 | ); |
| 2910 | } |
| 2911 | for tool in tools { |
| 2912 | let _ = write!( |
| 2913 | out, |
| 2914 | "\n tool {} (extension:{}; needs approval, which your approval mode or a session grant for this exact call of this plugin build may give)", |
| 2915 | tool.name, tool.plugin_name |
| 2916 | ); |
| 2917 | } |
| 2918 | for command in commands { |
| 2919 | let _ = write!( |
| 2920 | out, |
| 2921 | "\n command /{} (extension:{}; runs only when you invoke it)", |
| 2922 | command.name, command.plugin_name |
| 2923 | ); |
| 2924 | } |
| 2925 | let diagnostics = manager.diagnostics(); |
| 2926 | for diagnostic in diagnostics.iter().rev().take(5).rev() { |
| 2927 | let _ = write!(out, "\n · {diagnostic}"); |
| 2928 | } |
| 2929 | out |
| 2930 | } |
| 2931 | |
| 2932 | /// A plugin was enabled, disabled, trusted, revoked or removed: reconcile the |
| 2933 | /// host now, so a disabled plugin's calls are cancelled and its code torn |
| 2934 | /// down without waiting for the next turn. No-op with the flag off. |
| 2935 | pub fn plugins_changed(plugins: Arc<PluginRegistry>) { |
| 2936 | if activation::extension_host_policy_enabled() { |
| 2937 | let manager = manager(); |
| 2938 | manager.refresh_workspace(&plugins); |
| 2939 | manager.retry(); |
| 2940 | // An explicit plugin action, so also the moment to pick up an edit to |
| 2941 | // `[plugins."<name>".config]`. |
| 2942 | manager.reload_and_reconcile_in_background(); |
| 2943 | } |
| 2944 | } |
| 2945 | |
| 2946 | /// Install the user's per-plugin settings at boot (the host's other boot |
| 2947 | /// options are `configure`). `source` is the config file `/plugin reload` |
| 2948 | /// re-reads them from. No-op with the flag off. |
| 2949 | pub fn install_plugin_settings( |
| 2950 | settings: Option<&BTreeMap<String, crate::config::PluginSettings>>, |
| 2951 | source: Option<PathBuf>, |
| 2952 | ) { |
| 2953 | if activation::extension_host_policy_enabled() { |
| 2954 | manager().set_plugin_settings(settings.unwrap_or(&BTreeMap::new()), source); |
| 2955 | } |
| 2956 | } |
| 2957 | |
| 2958 | /// What the user configured for `plugin_name` (keys only, or the reason it is |
| 2959 | /// refused), for `/plugin show`. `None` with the flag off or nothing configured. |
| 2960 | #[must_use] |
| 2961 | pub fn plugin_config_summary(plugin_name: &str) -> Option<Result<Vec<String>, String>> { |
| 2962 | if !activation::extension_host_policy_enabled() { |
| 2963 | return None; |
| 2964 | } |
| 2965 | manager().plugin_config_summary(plugin_name) |
| 2966 | } |
| 2967 | |
| 2968 | /// The live extension commands of `workspace`'s reviewed plugins, for the user |
| 2969 | /// command registry. Empty with the flag off. Never starts the host. |
| 2970 | #[must_use] |
| 2971 | pub(crate) fn live_commands_for(workspace: &Path) -> Vec<command::ExtensionCommandEntry> { |
| 2972 | if !activation::extension_host_policy_enabled() { |
| 2973 | return Vec::new(); |
| 2974 | } |
| 2975 | manager().commands_for_workspace(workspace) |
| 2976 | } |
| 2977 | |
| 2978 | pub(crate) fn live_commands_for_plugins( |
| 2979 | plugins: &PluginRegistry, |
| 2980 | ) -> Vec<command::ExtensionCommandEntry> { |
| 2981 | if !activation::extension_host_policy_enabled() { |
| 2982 | return Vec::new(); |
| 2983 | } |
| 2984 | manager().commands_for_plugins(plugins) |
| 2985 | } |
| 2986 | |
| 2987 | /// Run an extension command the user invoked; see [`command::run`]. Errors |
| 2988 | /// are the text to show. |
| 2989 | #[cfg(test)] |
| 2990 | pub async fn run_command( |
| 2991 | command: &command::ExtensionCommandRef, |
| 2992 | raw_input: &str, |
| 2993 | session_id: Option<&str>, |
| 2994 | ) -> Result<command::CommandOutcome, String> { |
| 2995 | command::run(&manager().shared, command, raw_input, session_id).await |
| 2996 | } |
| 2997 | |
| 2998 | /// The command the current UI caller selected. A stale focus cannot consume a |
| 2999 | /// command reference captured from another agent's palette. |
| 3000 | pub async fn run_command_for_plugins( |
| 3001 | command: &command::ExtensionCommandRef, |
| 3002 | raw_input: &str, |
| 3003 | session_id: Option<&str>, |
| 3004 | plugins: &PluginRegistry, |
| 3005 | ) -> Result<command::CommandOutcome, String> { |
| 3006 | if command.scope.is_some() && command.selection != plugins.caller_selection() { |
| 3007 | return Err("Native command belongs to a different caller; select it again".into()); |
| 3008 | } |
| 3009 | command::run(&manager().shared, command, raw_input, session_id).await |
| 3010 | } |
| 3011 | |
| 3012 | /// The `/plugin` section. With the experimental host off it is one line |
| 3013 | /// saying so. |
| 3014 | #[must_use] |
| 3015 | pub fn status_report() -> String { |
| 3016 | if !activation::extension_host_policy_enabled() { |
| 3017 | return format!("Extension host (experimental): {}", HostStatus::Disabled); |
| 3018 | } |
| 3019 | render_status(&manager()) |
| 3020 | } |
| 3021 | |
| 3022 | pub struct OwnerReport { |
| 3023 | pub state: Option<OwnerState>, |
| 3024 | pub tools: Vec<String>, |
| 3025 | pub diagnostics: Vec<String>, |
| 3026 | } |
| 3027 | |
| 3028 | /// The resolved plugin id is supplied by the existing `/plugin show` facet. |
| 3029 | pub fn owner_report(plugin_id: &str) -> Option<OwnerReport> { |
| 3030 | if !activation::extension_host_policy_enabled() { |
| 3031 | return None; |
| 3032 | } |
| 3033 | manager().owner_report(plugin_id) |
| 3034 | } |
| 3035 | |
| 3036 | /// Where a built-in module's source is expected: |
| 3037 | /// `<bundle dir>/builtin/<module>.mjs`, beside the bundle that was |
| 3038 | /// materialized for this build. |
| 3039 | fn builtin_source_path(root: &Path, module: &BuiltinModule) -> PathBuf { |
| 3040 | supervisor::bundle_dir(root, bundle_sha256()) |
| 3041 | .join("builtin") |
| 3042 | .join(format!("{}.mjs", module.id)) |
| 3043 | } |
| 3044 | |
| 3045 | /// Reviewed, enabled plugins with `native` entries, keyed by plugin id, read |
| 3046 | /// from Codewhale's immutable staged snapshot. Blocking. |
| 3047 | fn desired_owners(plugins: &PluginRegistry) -> (BTreeMap<String, DesiredOwner>, Vec<String>) { |
| 3048 | let (sources, mut errors) = crate::plugins::runtime::active_component_sources( |
| 3049 | plugins, |
| 3050 | PluginActivationCapability::Native, |
| 3051 | ); |
| 3052 | let mut desired: BTreeMap<String, DesiredOwner> = BTreeMap::new(); |
| 3053 | let mut broken: BTreeSet<String> = BTreeSet::new(); |
| 3054 | let mut needs_selection: BTreeSet<String> = BTreeSet::new(); |
| 3055 | for source in sources { |
| 3056 | let plugin_id = source.authority.plugin_id.as_str().to_string(); |
| 3057 | if plugins.native_catalog_requires_selection(&plugin_id) { |
| 3058 | if needs_selection.insert(plugin_id.clone()) { |
| 3059 | errors.push(format!("Plugin `{}` raw agent-presets has no default; select a reviewed roster preset before activation", source.plugin_name)); |
| 3060 | } |
| 3061 | continue; |
| 3062 | } |
| 3063 | // A plugin id can never be a tier-0 owner id. Discovery builds ids as |
| 3064 | // `<scope>/<hex>/<name>` and a manifest name cannot hold `:`, so this |
| 3065 | // cannot fire; it is the last check before an id reaches the host. |
| 3066 | if let Err(reason) = HostTier::Plugin.check_owner_id(&plugin_id) { |
| 3067 | errors.push(format!( |
| 3068 | "Plugin `{}` native entry {} was denied: {reason}", |
| 3069 | source.plugin_name, |
| 3070 | source.path.display() |
| 3071 | )); |
| 3072 | broken.insert(plugin_id); |
| 3073 | continue; |
| 3074 | } |
| 3075 | // The rule discovery reports, re-checked on the staged copy: the |
| 3076 | // name here, and file-ness by the read itself. |
| 3077 | let bytes = match crate::plugins::runtime::native_entry_problem(&source.path, true) { |
| 3078 | None => std::fs::read(&source.path).map_err(|error| error.to_string()), |
| 3079 | Some(problem) => Err(problem.to_string()), |
| 3080 | }; |
| 3081 | match bytes { |
| 3082 | Ok(bytes) => { |
| 3083 | let digest = hex(Sha256::digest(&bytes)); |
| 3084 | if !plugins.native_entry_selected( |
| 3085 | &plugin_id, |
| 3086 | &source.path.to_string_lossy(), |
| 3087 | &digest, |
| 3088 | ) { |
| 3089 | continue; |
| 3090 | } |
| 3091 | let owner = desired.entry(plugin_id).or_insert_with(|| DesiredOwner { |
| 3092 | plugin_name: source.plugin_name.clone(), |
| 3093 | authority: source.authority.clone(), |
| 3094 | entries: Vec::new(), |
| 3095 | }); |
| 3096 | // The manifest may name one file twice; it is one entry. |
| 3097 | if !owner.entries.iter().any(|(path, _)| *path == source.path) { |
| 3098 | owner |
| 3099 | .entries |
| 3100 | .push((source.path.clone(), hex(Sha256::digest(&bytes)))); |
| 3101 | } |
| 3102 | } |
| 3103 | Err(reason) => { |
| 3104 | errors.push(format!( |
| 3105 | "Plugin `{}` native entry {} was denied: {reason}", |
| 3106 | source.plugin_name, |
| 3107 | source.path.display() |
| 3108 | )); |
| 3109 | broken.insert(plugin_id); |
| 3110 | } |
| 3111 | } |
| 3112 | } |
| 3113 | // All-or-nothing per plugin: one unusable entry keeps the whole plugin out. |
| 3114 | for plugin_id in broken { |
| 3115 | desired.remove(&plugin_id); |
| 3116 | } |
| 3117 | (desired, errors) |
| 3118 | } |
| 3119 | |
| 3120 | /// One scan of the complete attachment set. Its per-engine views and global |
| 3121 | /// owner union must be published together, against those same snapshots. |
| 3122 | struct DesiredScan { |
| 3123 | attachments: Vec<(u64, Arc<PluginRegistry>, BTreeMap<String, String>)>, |
| 3124 | owners: BTreeMap<String, DesiredOwner>, |
| 3125 | errors: Vec<String>, |
| 3126 | } |
| 3127 | |
| 3128 | impl DesiredScan { |
| 3129 | fn publish( |
| 3130 | self, |
| 3131 | current: &mut BTreeMap<u64, AttachmentState>, |
| 3132 | ) -> Option<(BTreeMap<String, DesiredOwner>, Vec<String>)> { |
| 3133 | if current.len() != self.attachments.len() |
| 3134 | || self.attachments.iter().any(|(id, scanned, _)| { |
| 3135 | !current |
| 3136 | .get(id) |
| 3137 | .is_some_and(|state| Arc::ptr_eq(&state.plugins, scanned)) |
| 3138 | }) |
| 3139 | { |
| 3140 | return None; |
| 3141 | } |
| 3142 | for (id, _, desired) in self.attachments { |
| 3143 | let state = current.get_mut(&id).expect("validated attachment"); |
| 3144 | let entries = self |
| 3145 | .owners |
| 3146 | .iter() |
| 3147 | .flat_map(|(plugin_id, want)| { |
| 3148 | want.entries.iter().filter_map(|(path, sha256)| { |
| 3149 | let entry = EntryRef { |
| 3150 | path: path.to_string_lossy().into_owned(), |
| 3151 | sha256: sha256.clone(), |
| 3152 | }; |
| 3153 | (desired.get(plugin_id) == Some(&want.authority.content_hash) |
| 3154 | && state.plugins.native_entry_selected( |
| 3155 | plugin_id, |
| 3156 | &entry.path, |
| 3157 | &entry.sha256, |
| 3158 | )) |
| 3159 | .then(|| composition_scope::NativePresetRef { |
| 3160 | plugin_id: plugin_id.clone(), |
| 3161 | content_hash: want.authority.content_hash.clone(), |
| 3162 | entry, |
| 3163 | }) |
| 3164 | }) |
| 3165 | }) |
| 3166 | .collect(); |
| 3167 | state.selection = composition_scope::CompositionSelection { |
| 3168 | revision: state.plugins.caller_selection(), |
| 3169 | desired: desired.clone(), |
| 3170 | entries, |
| 3171 | }; |
| 3172 | state.desired = desired; |
| 3173 | } |
| 3174 | command::bump_epoch(); |
| 3175 | Some((self.owners, self.errors)) |
| 3176 | } |
| 3177 | } |
| 3178 | |
| 3179 | /// Scan every attached snapshot (engines sharing one snapshot scan it once) |
| 3180 | /// and merge what they desire. Blocking. |
| 3181 | /// |
| 3182 | /// Two snapshots can disagree about one plugin id only while one of them is |
| 3183 | /// stale; the stale one then fails its persisted-state check and desires |
| 3184 | /// nothing, so the first valid scan wins and the per-attachment hashes keep |
| 3185 | /// each engine's tools to the bytes it desires. |
| 3186 | fn union_of_desired_owners(snapshots: Vec<(u64, Arc<PluginRegistry>)>) -> DesiredScan { |
| 3187 | let mut union: BTreeMap<String, DesiredOwner> = BTreeMap::new(); |
| 3188 | let mut per_attachment = Vec::with_capacity(snapshots.len()); |
| 3189 | let mut scanned: Vec<(Arc<PluginRegistry>, BTreeMap<String, String>)> = Vec::new(); |
| 3190 | let mut errors: Vec<String> = Vec::new(); |
| 3191 | for (id, plugins) in snapshots { |
| 3192 | if let Some((_, hashes)) = scanned.iter().find(|(seen, _)| Arc::ptr_eq(seen, &plugins)) { |
| 3193 | per_attachment.push((id, Arc::clone(&plugins), hashes.clone())); |
| 3194 | continue; |
| 3195 | } |
| 3196 | let (desired, scan_errors) = desired_owners(&plugins); |
| 3197 | for error in scan_errors { |
| 3198 | if !errors.contains(&error) { |
| 3199 | errors.push(error); |
| 3200 | } |
| 3201 | } |
| 3202 | let hashes: BTreeMap<String, String> = desired |
| 3203 | .iter() |
| 3204 | .map(|(plugin_id, want)| (plugin_id.clone(), want.authority.content_hash.clone())) |
| 3205 | .collect(); |
| 3206 | for (plugin_id, want) in desired { |
| 3207 | match union.entry(plugin_id) { |
| 3208 | std::collections::btree_map::Entry::Vacant(entry) => { |
| 3209 | entry.insert(want); |
| 3210 | } |
| 3211 | std::collections::btree_map::Entry::Occupied(mut entry) => { |
| 3212 | for native in want.entries { |
| 3213 | if !entry.get().entries.contains(&native) { |
| 3214 | entry.get_mut().entries.push(native); |
| 3215 | } |
| 3216 | } |
| 3217 | } |
| 3218 | } |
| 3219 | } |
| 3220 | per_attachment.push((id, Arc::clone(&plugins), hashes.clone())); |
| 3221 | scanned.push((plugins, hashes)); |
| 3222 | } |
| 3223 | DesiredScan { |
| 3224 | attachments: per_attachment, |
| 3225 | owners: union, |
| 3226 | errors, |
| 3227 | } |
| 3228 | } |
| 3229 | |
| 3230 | /// One engine's hold on the process-wide extension host. |
| 3231 | /// |
| 3232 | /// The engine publishes its workspace plugin snapshot here and installs only |
| 3233 | /// the tools of owners that snapshot desires. Dropping it detaches without |
| 3234 | /// revoking anything: the next reconcile revokes owners no remaining |
| 3235 | /// attachment desires, so an engine being replaced never tears down plugins |
| 3236 | /// its successor is about to use. |
| 3237 | pub struct HostAttachment { |
| 3238 | id: u64, |
| 3239 | manager: Arc<ExtensionHostManager>, |
| 3240 | } |
| 3241 | |
| 3242 | impl HostAttachment { |
| 3243 | #[must_use] |
| 3244 | pub fn manager(&self) -> &Arc<ExtensionHostManager> { |
| 3245 | &self.manager |
| 3246 | } |
| 3247 | |
| 3248 | /// The engine switched workspace: publish the new snapshot. |
| 3249 | pub fn set_plugins(&self, plugins: Arc<PluginRegistry>) { |
| 3250 | self.manager.shared.core_calls.revoke_attachment(self.id); |
| 3251 | self.manager |
| 3252 | .shared |
| 3253 | .execution_broker |
| 3254 | .revoke_attachment(self.id); |
| 3255 | if let Some(state) = self |
| 3256 | .manager |
| 3257 | .shared |
| 3258 | .attachments |
| 3259 | .lock() |
| 3260 | .expect("attachments lock") |
| 3261 | .get_mut(&self.id) |
| 3262 | { |
| 3263 | state.selection_cancel.cancel(); |
| 3264 | state.selection_cancel = CancellationToken::new(); |
| 3265 | state.revision += 1; |
| 3266 | state.plugins = Arc::new(plugins.bind_caller(composition_scope::SelectionRevision { |
| 3267 | attachment_id: self.id, |
| 3268 | revision: state.revision, |
| 3269 | })); |
| 3270 | state.desired.clear(); |
| 3271 | state.selection = Default::default(); |
| 3272 | } |
| 3273 | command::bump_epoch(); |
| 3274 | } |
| 3275 | |
| 3276 | pub fn set_identity(&self, session_id: Option<String>, agent_id: Option<String>) { |
| 3277 | if let Some(state) = self |
| 3278 | .manager |
| 3279 | .shared |
| 3280 | .attachments |
| 3281 | .lock() |
| 3282 | .expect("attachments lock") |
| 3283 | .get_mut(&self.id) |
| 3284 | { |
| 3285 | state.session_id = session_id; |
| 3286 | state.agent_id = agent_id; |
| 3287 | } |
| 3288 | } |
| 3289 | pub fn plugin_view(&self) -> Arc<PluginRegistry> { |
| 3290 | self.manager |
| 3291 | .shared |
| 3292 | .attachments |
| 3293 | .lock() |
| 3294 | .expect("attachments lock") |
| 3295 | .get(&self.id) |
| 3296 | .map(|state| Arc::clone(&state.plugins)) |
| 3297 | .unwrap_or_else(|| Arc::new(PluginRegistry::new())) |
| 3298 | } |
| 3299 | pub async fn reconcile(&self) -> Result<(), String> { |
| 3300 | self.manager.reconcile().await |
| 3301 | } |
| 3302 | |
| 3303 | /// Reconcile the host against every attachment, waiting for it. |
| 3304 | #[cfg(test)] |
| 3305 | pub async fn sync(&self) -> Result<(), String> { |
| 3306 | self.manager.reconcile().await |
| 3307 | } |
| 3308 | |
| 3309 | /// Reconcile the host against every attachment, without waiting. |
| 3310 | pub fn sync_in_background(&self) { |
| 3311 | self.manager.reconcile_in_background(); |
| 3312 | } |
| 3313 | |
| 3314 | /// Add this engine's live extension tools to `tool_registry`. |
| 3315 | pub fn install_tools(&self, tool_registry: &mut crate::tools::ToolRegistry) -> Vec<String> { |
| 3316 | self.manager.install_tools_for(self.id, tool_registry) |
| 3317 | } |
| 3318 | } |
| 3319 | |
| 3320 | impl ManagerShared { |
| 3321 | pub(crate) fn plugins_for_selection( |
| 3322 | &self, |
| 3323 | revision: composition_scope::SelectionRevision, |
| 3324 | session_id: Option<&str>, |
| 3325 | ) -> Option<(Arc<PluginRegistry>, Option<String>)> { |
| 3326 | self.attachments |
| 3327 | .lock() |
| 3328 | .expect("attachments lock") |
| 3329 | .get(&revision.attachment_id) |
| 3330 | .filter(|state| { |
| 3331 | state.revision == revision.revision && state.session_id.as_deref() == session_id |
| 3332 | }) |
| 3333 | .map(|state| (Arc::clone(&state.plugins), state.agent_id.clone())) |
| 3334 | } |
| 3335 | pub(super) fn selection_cancellation( |
| 3336 | &self, |
| 3337 | revision: composition_scope::SelectionRevision, |
| 3338 | ) -> Option<CancellationToken> { |
| 3339 | self.attachments |
| 3340 | .lock() |
| 3341 | .expect("attachments lock") |
| 3342 | .get(&revision.attachment_id) |
| 3343 | .filter(|state| state.revision == revision.revision) |
| 3344 | .map(|state| state.selection_cancel.clone()) |
| 3345 | } |
| 3346 | pub(crate) fn selection(&self, id: u64) -> composition_scope::CompositionSelection { |
| 3347 | self.attachments |
| 3348 | .lock() |
| 3349 | .expect("attachments lock") |
| 3350 | .get(&id) |
| 3351 | .map(|state| state.selection.clone()) |
| 3352 | .unwrap_or_default() |
| 3353 | } |
| 3354 | pub(crate) fn selection_current( |
| 3355 | &self, |
| 3356 | selection: composition_scope::SelectionRevision, |
| 3357 | plugin_id: &str, |
| 3358 | hash: &str, |
| 3359 | scope: Option<&EntryRef>, |
| 3360 | ) -> bool { |
| 3361 | self.attachments |
| 3362 | .lock() |
| 3363 | .expect("attachments lock") |
| 3364 | .get(&selection.attachment_id) |
| 3365 | .is_some_and(|state| { |
| 3366 | state.revision == selection.revision |
| 3367 | && state.selection.includes(plugin_id, hash, scope) |
| 3368 | }) |
| 3369 | } |
| 3370 | pub(crate) fn check_selection( |
| 3371 | &self, |
| 3372 | selection: Option<composition_scope::SelectionRevision>, |
| 3373 | plugins: Option<&PluginRegistry>, |
| 3374 | plugin_id: &str, |
| 3375 | hash: &str, |
| 3376 | scope: Option<&EntryRef>, |
| 3377 | ) -> Result<(), String> { |
| 3378 | if scope.is_none() { |
| 3379 | return Ok(()); |
| 3380 | } |
| 3381 | let selection = selection.ok_or("Native contribution has no caller selection")?; |
| 3382 | if plugins.and_then(PluginRegistry::caller_selection) != Some(selection) |
| 3383 | || !self.selection_current(selection, plugin_id, hash, scope) |
| 3384 | { |
| 3385 | return Err("Native contribution is no longer selected for this caller".into()); |
| 3386 | } |
| 3387 | Ok(()) |
| 3388 | } |
| 3389 | } |
| 3390 | /// Last-consumer guard for any tool executed under a caller-bound composition. |
| 3391 | pub(crate) fn validate_caller_plugins(plugins: Option<&PluginRegistry>) -> Result<(), String> { |
| 3392 | let Some(revision) = plugins.and_then(PluginRegistry::caller_selection) else { |
| 3393 | return Ok(()); |
| 3394 | }; |
| 3395 | let manager = manager(); |
| 3396 | let attachments = manager.shared.attachments.lock().expect("attachments lock"); |
| 3397 | if attachments |
| 3398 | .get(&revision.attachment_id) |
| 3399 | .is_some_and(|state| state.revision == revision.revision) |
| 3400 | { |
| 3401 | Ok(()) |
| 3402 | } else { |
| 3403 | Err("Native caller composition changed or was detached; prepare the call again".into()) |
| 3404 | } |
| 3405 | } |
| 3406 | /// Raw catalog choices are revalidated against the reviewed staged bundle; |
| 3407 | /// discovery never activates them. Generic Native entries retain active-scope |
| 3408 | /// discovery. Child preparation revalidates before attaching the chosen entry. |
| 3409 | pub(crate) fn native_presets_for_plugins( |
| 3410 | plugins: &PluginRegistry, |
| 3411 | ) -> Vec<(composition_scope::NativePresetRef, PluginAuthority)> { |
| 3412 | if !activation::extension_host_policy_enabled() { |
| 3413 | return Vec::new(); |
| 3414 | } |
| 3415 | let manager = manager(); |
| 3416 | let mut presets: Vec<_> = crate::plugins::native_presets::admitted_entries(plugins) |
| 3417 | .into_iter() |
| 3418 | .map(|(preset, _, authority)| (preset, authority)) |
| 3419 | .collect(); |
| 3420 | if manager.shared.ready_host(HostTier::Plugin).is_err() { |
| 3421 | return presets; |
| 3422 | } |
| 3423 | let registry = manager.shared.registry.lock().expect("registry lock"); |
| 3424 | for owner in registry.owners() { |
| 3425 | let Some(authority) = owner.authority.as_ref() else { |
| 3426 | continue; |
| 3427 | }; |
| 3428 | if owner.tier != HostTier::Plugin |
| 3429 | || owner.state != OwnerState::Active |
| 3430 | || authority.workspace != plugins.workspace() |
| 3431 | || !plugins.get(&owner.owner.plugin_id).is_some_and(|plugin| { |
| 3432 | plugin.component_active(PluginActivationCapability::Native) |
| 3433 | && plugin.content_hash == owner.content_hash |
| 3434 | && plugin.state_generation == authority.state_generation |
| 3435 | }) |
| 3436 | { |
| 3437 | continue; |
| 3438 | } |
| 3439 | for (entry, state) in &owner.scopes { |
| 3440 | if *state == OwnerState::Active |
| 3441 | && !presets.iter().any(|(preset, _)| { |
| 3442 | preset.plugin_id == owner.owner.plugin_id && preset.entry == *entry |
| 3443 | }) |
| 3444 | { |
| 3445 | presets.push(( |
| 3446 | composition_scope::NativePresetRef { |
| 3447 | plugin_id: owner.owner.plugin_id.clone(), |
| 3448 | content_hash: owner.content_hash.clone(), |
| 3449 | entry: entry.clone(), |
| 3450 | }, |
| 3451 | authority.clone(), |
| 3452 | )); |
| 3453 | } |
| 3454 | } |
| 3455 | } |
| 3456 | presets.sort_by(|(left, _), (right, _)| { |
| 3457 | left.plugin_id |
| 3458 | .cmp(&right.plugin_id) |
| 3459 | .then_with(|| left.entry.path.cmp(&right.entry.path)) |
| 3460 | }); |
| 3461 | presets |
| 3462 | } |
| 3463 | |
| 3464 | pub(crate) fn caller_view( |
| 3465 | workspace: &Path, |
| 3466 | session_id: Option<&str>, |
| 3467 | agent_id: Option<&str>, |
| 3468 | ) -> Option<Arc<PluginRegistry>> { |
| 3469 | let manager = manager(); |
| 3470 | let attachments = manager.shared.attachments.lock().expect("attachments lock"); |
| 3471 | let mut found = attachments.values().filter(|state| { |
| 3472 | state.plugins.workspace() == workspace |
| 3473 | && state.session_id.as_deref() == session_id |
| 3474 | && state.agent_id.as_deref() == agent_id |
| 3475 | }); |
| 3476 | let view = found.next().map(|state| Arc::clone(&state.plugins)); |
| 3477 | if found.next().is_some() { None } else { view } |
| 3478 | } |
| 3479 | |
| 3480 | impl Drop for HostAttachment { |
| 3481 | fn drop(&mut self) { |
| 3482 | self.manager.shared.core_calls.revoke_attachment(self.id); |
| 3483 | if let Ok(mut attachments) = self.manager.shared.attachments.lock() |
| 3484 | && let Some(state) = attachments.remove(&self.id) |
| 3485 | { |
| 3486 | state.selection_cancel.cancel(); |
| 3487 | } |
| 3488 | command::bump_epoch(); |
| 3489 | } |
| 3490 | } |
| 3491 | |
| 3492 | impl std::fmt::Debug for HostAttachment { |
| 3493 | fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { |
| 3494 | f.debug_struct("HostAttachment") |
| 3495 | .field("id", &self.id) |
| 3496 | .finish() |
| 3497 | } |
| 3498 | } |
| 3499 | |
| 3500 | static GLOBAL: OnceLock<Arc<ExtensionHostManager>> = OnceLock::new(); |
| 3501 | |
| 3502 | #[cfg(test)] |
| 3503 | thread_local! { |
| 3504 | static TEST_MANAGER: std::cell::RefCell<Option<Arc<ExtensionHostManager>>> = |
| 3505 | const { std::cell::RefCell::new(None) }; |
| 3506 | } |
| 3507 | |
| 3508 | /// Configure the process-wide manager once, at boot, from user config. |
| 3509 | pub(crate) fn configure_with_handle( |
| 3510 | options: ExtensionHostOptions, |
| 3511 | handle: Option<tokio::runtime::Handle>, |
| 3512 | ) { |
| 3513 | let manager = Arc::new(ExtensionHostManager::new(options)); |
| 3514 | if let Some(handle) = handle { |
| 3515 | manager.bind_engine_handle(handle); |
| 3516 | } |
| 3517 | let _ = GLOBAL.set(manager); |
| 3518 | } |
| 3519 | |
| 3520 | /// The manager for this engine process (one host per process and tier). |
| 3521 | #[must_use] |
| 3522 | pub fn manager() -> Arc<ExtensionHostManager> { |
| 3523 | #[cfg(test)] |
| 3524 | if let Some(manager) = TEST_MANAGER.with(|cell| cell.borrow().clone()) { |
| 3525 | return manager; |
| 3526 | } |
| 3527 | Arc::clone( |
| 3528 | GLOBAL.get_or_init(|| Arc::new(ExtensionHostManager::new(ExtensionHostOptions::default()))), |
| 3529 | ) |
| 3530 | } |
| 3531 | |
| 3532 | /// Test-only: route [`manager`] on this thread to `manager`. |
| 3533 | #[cfg(test)] |
| 3534 | pub(crate) struct TestManagerGuard(Option<Arc<ExtensionHostManager>>); |
| 3535 | |
| 3536 | #[cfg(test)] |
| 3537 | impl TestManagerGuard { |
| 3538 | pub(crate) fn install(manager: Arc<ExtensionHostManager>) -> Self { |
| 3539 | Self(TEST_MANAGER.with(|cell| cell.replace(Some(manager)))) |
| 3540 | } |
| 3541 | } |
| 3542 | |
| 3543 | #[cfg(test)] |
| 3544 | impl Drop for TestManagerGuard { |
| 3545 | fn drop(&mut self) { |
| 3546 | TEST_MANAGER.with(|cell| *cell.borrow_mut() = self.0.take()); |
| 3547 | } |
| 3548 | } |
| 3549 | |
| 3550 | impl ExtensionHostManager { |
| 3551 | pub(crate) fn has_shell_hooks(&self, event: crate::hooks::HookEvent) -> bool { |
| 3552 | activation::extension_host_policy_enabled() |
| 3553 | && self |
| 3554 | .shared |
| 3555 | .registry |
| 3556 | .lock() |
| 3557 | .expect("registry lock") |
| 3558 | .live_shell_hooks() |
| 3559 | .iter() |
| 3560 | .any(|r| r.hook.event == event) |
| 3561 | } |
| 3562 | pub(crate) fn shell_hooks( |
| 3563 | &self, |
| 3564 | caller: &crate::hooks::HookCaller, |
| 3565 | event: crate::hooks::HookEvent, |
| 3566 | ) -> Vec<crate::hooks::Hook> { |
| 3567 | if !activation::extension_host_policy_enabled() { |
| 3568 | return Vec::new(); |
| 3569 | } |
| 3570 | let Some(plugins) = caller.plugins.as_ref() else { |
| 3571 | return Vec::new(); |
| 3572 | }; |
| 3573 | let Some(revision) = plugins.caller_selection() else { |
| 3574 | return Vec::new(); |
| 3575 | }; |
| 3576 | let Some((current, agent)) = self |
| 3577 | .shared |
| 3578 | .plugins_for_selection(revision, caller.session_id.as_deref()) |
| 3579 | else { |
| 3580 | return Vec::new(); |
| 3581 | }; |
| 3582 | if agent != caller.agent_id |
| 3583 | || !Arc::ptr_eq(¤t, plugins) |
| 3584 | || caller.workspace != plugins.workspace() |
| 3585 | { |
| 3586 | return Vec::new(); |
| 3587 | } |
| 3588 | self.shared |
| 3589 | .registry |
| 3590 | .lock() |
| 3591 | .expect("registry lock") |
| 3592 | .live_shell_hooks() |
| 3593 | .into_iter() |
| 3594 | .filter(|r| { |
| 3595 | r.hook.event == event |
| 3596 | && plugins.selected_native_entries().iter().any(|entry| { |
| 3597 | entry.plugin_id == r.owner.plugin_id |
| 3598 | && entry.content_hash == r.content_hash |
| 3599 | && r.scope.as_ref().is_none_or(|scope| entry.entry == *scope) |
| 3600 | }) |
| 3601 | }) |
| 3602 | .map(|r| r.hook) |
| 3603 | .collect() |
| 3604 | } |
| 3605 | } |
| 3606 |