返回 CodeWhale
mod.rs
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(&registry)?;
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, &registration.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(&params.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(&params.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(&params.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(&params.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(&params.owner.plugin_id);
1089 }
1090 drop(slot);
1091 shared.core_calls.revoke_owner(&params.owner.plugin_id);
1092 shared
1093 .execution_broker
1094 .revoke_owner(&params.owner.plugin_id);
1095 shared.mcp_users.fetch_sub(
1096 shared.mcp_broker.revoke_owner(&params.owner.plugin_id),
1097 Ordering::SeqCst,
1098 );
1099 shared.plugin_diagnostic(
1100 &params.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[&registration.name.to_ascii_lowercase()] != 1 {
1821 self.shared.plugin_diagnostic(&registration.owner.plugin_id,format!("Native tool `{}` is ambiguous in this caller; rename the selected definitions",registration.name));
1822 continue;
1823 }
1824 if taken.contains(&registration.name.to_ascii_lowercase()) {
1825 let origin = tool_registry
1826 .get(&registration.name)
1827 .map(|existing| existing.registration_origin().into_owned())
1828 .unwrap_or_else(|| "another tool".to_string());
1829 self.shared.plugin_diagnostic(&registration.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(&registration.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(&registration.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(&current, 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
3606 lines RUST