返回 CodeWhale
supervisor.rs
根目录 / crates / tui / src / extension_host / supervisor.rs
1 //! One extension-host process: launch plan, spawn, handshake, channel, exit.
2 //!
3 //! When the process exits, every in-flight call fails with a typed error.
4 //! The existing Rust manager owns heartbeat, crash budget, generation changes,
5 //! and receipt-checked replay; this channel never replays a tool call.
6 //!
7 //! **Host-originated requests** run as tasks, not inline in the reader. Each
8 //! request the host sends is admitted into [`InboundRequests`] (at most
9 //! [`protocol::MAX_INFLIGHT`] at a time, no id twice, both violations end the
10 //! host) and handed to [`HostEvents::host_request`] as its own task. It is
11 //! cancelled by the host's `$/cancel` (whatever the handler later produces is
12 //! dropped), by its owner's revocation (the host is answered `Cancelled`) and
13 //! by the host's exit (nobody is answered); a handler that ignores its token is
14 //! abandoned [`CANCEL_GRACE`] after the cancel. A method reserved for the other
15 //! tier is neither accepted from the host nor sent to it
16 //! (`protocol::MethodSpec::tiers`), and `host/hello` must report the tier the
17 //! core launched and the built-in module digests it pins
18 //! ([`check_hello_identity`]).
19 //!
20 //! **Tiers** ([`HostTier`]). The host is two processes, one per trust tier,
21 //! started from the same bundle with `--tier=plugin|builtin` as the last
22 //! argument. Each has its own data directory ([`tier_data_dir`]) and its own
23 //! sandbox plan; the plugin tier's also denies reads of the builtin tier's
24 //! data directory ([`host_denied_read_paths`]). Only the plugin tier starts
25 //! in production today.
26 //!
27 //! **OS sandbox** ([`HostSandbox`]: the one value `/plugin`, doctor and the
28 //! start diagnostic report). The host runs under Codewhale's command sandbox
29 //! with a workspace-write policy rooted at its tier's data directory (the
30 //! plugin tier's is `$CODEWHALE_HOME/extension-host/data`): no direct network,
31 //! writes only there and in the temp dirs, and **no reads** of the Codewhale
32 //! homes (everything but the bundle, its data dir and plugin code), the Codex
33 //! and DSH credential homes, and the credential-store default deny-list
34 //! (`sandbox::read_guard`). Other user-readable files stay readable —
35 //! including `.env` files, whose filename rule has no Seatbelt subpath or
36 //! bubblewrap mount form — so this is defense-in-depth, not a containment
37 //! boundary.
38 //! * macOS: Seatbelt. Mach services are not restricted.
39 //! * Linux: bubblewrap (`/usr/bin/bwrap`, the shell's builder), used without
40 //! the shell's `prefer_bwrap` opt-in. Every launch first runs the finished
41 //! wrapper around `<runtime> --version` ([`probe_bwrap`]). When bwrap is
42 //! missing or cannot start — e.g. unprivileged user namespaces blocked by
43 //! Ubuntu 24.04's `kernel.apparmor_restrict_unprivileged_userns` — the host
44 //! refuses Native launch and reports the concrete error. The pinned Builtin
45 //! exception is diagnosed and ticket-bound. bwrap can mask only what exists, so each
46 //! Codewhale home is masked whole and its readable entries are bound again
47 //! (`sandbox::bwrap_exception_args`): an entry created after launch is
48 //! denied, as on macOS.
49 //! * Windows: Native uses a freshly created LPAC AppContainer with no network
50 //! capabilities. Rust checks its actual token, attaches/caps the existing Job
51 //! before resuming, and requires a real data-read/write + outside-read/write
52 //! + loopback-network allow/deny probe. Only Core-selected runtime/bundle and
53 //! reviewed staged roots are granted reads. The Job is lifetime/memory only.
54 //! Builtin retains the separately diagnosed Rust-ticket-bound exception.
55 //!
56 //! Known limits under bubblewrap: a default-deny-list credential store
57 //! created after launch stays readable (Seatbelt denies it by name); a
58 //! Codewhale home, or a readable entry such as `plugins/`, that does not exist
59 //! at launch is not masked, or not visible, until the host restarts; bwrap
60 //! anywhere but `/usr/bin/bwrap` is not used; the probe costs one extra
61 //! runtime start per launch. The reported pid is bwrap's. The host runs in
62 //! bwrap's PID namespace, where its parent is bwrap's init and never
63 //! changes, so the host's own parent watchdog (`extension-host/src/main.ts`)
64 //! cannot fire; `--die-with-parent` is what ends it with the core. That is
65 //! a parent-death signal tied to the thread that spawned bwrap, a Tokio
66 //! worker that lives as long as the runtime. Plugin child processes end with
67 //! the namespace when bwrap's init loses the host.
68 //!
69 //! **Runtime.** Node (the default) or Bun (opt-in: `runtime = "bun"`, or
70 //! `"auto"`, which prefers a supported Bun), chosen once per manager
71 //! (`[extension_host] runtime`). Each gets its own flags ([`runtime_args`])
72 //! and environment ([`runtime_env`]); Bun silently ignores Node's heap and
73 //! `__proto__` flags. The host itself takes in-process native code away from
74 //! plugins (`extension-host/src/runtime.ts`: `bun:ffi`, `Bun.FFI`, SQLite
75 //! extension loading, Worker threads, ShadowRealm, `process.dlopen`) and
76 //! refuses to start if a lock does not hold. That lockdown covers the entry
77 //! points found so far (Bun 1.4, Node 22 and 26), not every one a runtime
78 //! may add.
79 //!
80 //! **Memory cap** ([`MemoryEnforcement`]). What was measured where: the
81 //! macOS mechanism on macOS 26.1 arm64 (2026-09-30, Bun 1.4.0 and Node);
82 //! the Linux thresholds in [`HOST_MEMORY_CAP`] in a Linux container
83 //! (2026-09-29). Hosted CI runs the Rust memory-cap test on Linux, macOS and
84 //! Windows with Node only; the Rust host tests have not run a Bun host on
85 //! Linux or Windows (CI's JS host suites run under Bun 1.4.0 on Linux).
86 //! * Linux: `RLIMIT_DATA`, set in the child before exec (clamped to an
87 //! inherited hard limit that is already lower); an allocation past the cap
88 //! fails. Plugin child processes inherit it.
89 //! * Windows: the Job Object's per-process limit; an allocation past the cap
90 //! fails. It applies to each process in the job, plugin children included.
91 //! * macOS: `setrlimit(RLIMIT_AS/RLIMIT_DATA)` below the current mapping size
92 //! fails with `EINVAL`, and `memorystatus_control` needs privilege. A fatal
93 //! jetsam limit set as a `posix_spawn` attribute works unprivileged, but a
94 //! later `exec` clears it, so it cannot be set on `sandbox-exec`. The Bun
95 //! host therefore re-executes itself in place with the limit before any
96 //! plugin loads, and reports it in `host/hello`; past the cap the kernel
97 //! SIGKILLs it. Plugin child processes are not covered. A Bun host that
98 //! cannot apply the requested limit is refused before initialization. Node
99 //! has no FFI to do the same, so a Node host is checked at each heartbeat
100 //! instead: enforced only as often as
101 //! the heartbeat runs. Node also keeps `--max-old-space-size=256`.
102
103 use std::collections::{HashMap, VecDeque};
104 use std::path::{Path, PathBuf};
105 use std::process::Stdio;
106 use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
107 use std::sync::{Arc, Mutex};
108 use std::time::Duration;
109
110 use async_trait::async_trait;
111 use serde_json::{Value, json};
112 use tokio::io::{AsyncReadExt, AsyncWriteExt};
113 use tokio::sync::{mpsc, oneshot};
114 use tokio_util::sync::CancellationToken;
115
116 use crate::dependencies::{HostRuntime, HostRuntimeKind};
117 use crate::tools::codemode::PauseClock;
118
119 use super::protocol::{
120 self, CoreRequest, Direction, HelloParams, HostLimits, HostMessage, HostNotification,
121 HostRequest, InitializeParams, RegisterResult, RpcErrorWire, error_code,
122 };
123 use super::tier::{self, BuiltinModule, HostTier};
124
125 /// Budget for spawn → `host/hello` → `host/initialize` → `host/ready`.
126 ///
127 /// A warm start takes well under 100 ms, but the first start of a freshly
128 /// materialized bundle pays for a cold `node` launch, and on Windows for an
129 /// on-access antivirus scan of both. Windows CI under full test load missed
130 /// the former 5 s budget with the host silent on stderr (four runs on
131 /// 2026-09-29) while the same tests normally finish in about 1 s.
132 /// A miss is sticky: the host is marked failed until the next session, so
133 /// a too-tight budget disables every extension for that session. The
134 /// handshake runs in the background, off the first-prompt path, so a wider
135 /// budget costs nothing when the host is healthy; 30 s matches the MCP stdio
136 /// handshake; it is independent of the active TUI MCP client.
137 pub const HANDSHAKE_DEADLINE: Duration = Duration::from_secs(30);
138 pub const ACTIVATE_DEADLINE: Duration = Duration::from_secs(5);
139 pub const DISPOSE_DEADLINE: Duration = Duration::from_secs(2);
140 /// A `host/ping` answer later than this means a hung host: the heartbeat's
141 /// default `hang_timeout`, and the bound for any other ping.
142 pub const PING_DEADLINE: Duration = Duration::from_secs(10);
143 /// Grace between `$/cancel` and resolving a call as cancelled on this side.
144 pub const CANCEL_GRACE: Duration = Duration::from_millis(500);
145 const STDERR_TAIL_BYTES: usize = 8 * 1024;
146 const OUTBOUND_QUEUE: usize = 128;
147
148 #[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
149 pub enum HostCallError {
150 #[error("{message} (extension error {code})")]
151 Rpc { code: i64, message: String },
152 #[error("extension host exited: {0}")]
153 Exited(String),
154 #[error("cancelled: {0}")]
155 Cancelled(String),
156 #[error("`{method}` timed out after {after:?}; the host was told to cancel it")]
157 Timeout {
158 method: &'static str,
159 after: Duration,
160 },
161 #[error("extension host channel is full")]
162 Busy,
163 }
164
165 /// How the host is started: the argv (wrapped by the OS sandbox when one is
166 /// available), its working directory, and which sandbox applies.
167 #[derive(Debug, Clone)]
168 pub(crate) struct HostLaunch {
169 /// The trust tier this host process serves (`--tier=` in its argv).
170 pub tier: HostTier,
171 pub program: PathBuf,
172 pub args: Vec<String>,
173 pub cwd: PathBuf,
174 pub sandbox: HostSandbox,
175 /// Environment the sandbox wrapper adds (`CODEWHALE_SANDBOX`, …).
176 pub sandbox_env: Vec<(String, String)>,
177 /// The runtime inside the wrapper; `host/hello` must report the same one.
178 pub runtime: HostRuntime,
179 /// Environment the runtime needs ([`runtime_env`]).
180 pub runtime_env: Vec<(String, String)>,
181 /// Bytes; see the module docs for how each platform enforces it.
182 pub memory_cap: u64,
183 /// How the cap is meant to be enforced; a macOS Bun host confirms its
184 /// jetsam limit in `host/hello` or initialization is refused.
185 pub memory: MemoryEnforcement,
186 /// The built-in module digests `host/hello` must report: the table this
187 /// build pins ([`tier::BUILTIN_MODULES`]), which the host bundle embedded
188 /// from the same build must agree with ([`check_hello_identity`]).
189 pub builtin_modules: &'static [BuiltinModule],
190 /// Exact verified Native LPAC plan; Builtin retains its labelled exception.
191 #[cfg(windows)]
192 pub windows: Option<super::windows::NativeSandbox>,
193 }
194
195 /// Whether the host runs under an OS sandbox, and why not when it does not
196 /// (module docs). Settled by [`plan_launch`]; never inferred afterwards.
197 #[derive(Debug, Clone, PartialEq, Eq)]
198 pub enum HostSandbox {
199 /// Under this wrapper, named as `sandbox::SandboxType` names it
200 /// (`macos-seatbelt`, `linux-bwrap`).
201 Wrapped(String),
202 /// With the user's permissions, for this reason.
203 Unsandboxed(String),
204 }
205
206 impl HostSandbox {
207 /// The wrapper's name, or `none: <reason>`, for one-line diagnostics.
208 #[must_use]
209 pub fn label(&self) -> String {
210 match self {
211 Self::Wrapped(name) => name.clone(),
212 Self::Unsandboxed(reason) => format!("none: {reason}"),
213 }
214 }
215 }
216
217 /// What `/plugin` and doctor say about the sandbox.
218 impl std::fmt::Display for HostSandbox {
219 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
220 match self {
221 Self::Wrapped(name) if name == "windows-lpac" => write!(
222 f,
223 "windows-lpac sandbox (no direct network; reads only the selected runtime, canonical bundle and reviewed staged code; writes to Native data and platform-private scratch; Job limits lifetime/memory)"
224 ),
225 Self::Wrapped(name) => write!(
226 f,
227 "{name} sandbox (no direct network; the Codewhale home except plugin code, the Codex and DSH credential homes and the default credential stores are unreadable; other files you can read, such as project .env files, are not protected)"
228 ),
229 Self::Unsandboxed(reason) => write!(
230 f,
231 "UNSANDBOXED ({reason}): host code runs with your user permissions"
232 ),
233 }
234 }
235 }
236
237 /// Environment variable carrying the jetsam limit a macOS Bun host applies to
238 /// itself, in MiB (`extension-host/src/runtime.ts`, `applyMemoryLimit`).
239 pub(crate) const MEMORY_LIMIT_REQUEST_ENV: &str = "CODEWHALE_HOST_MEMORY_LIMIT_MIB";
240
241 /// How the host's memory cap is enforced (module docs).
242 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
243 pub enum MemoryEnforcement {
244 /// Linux: kernel `RLIMIT_DATA`.
245 Rlimit,
246 /// Windows: the Job Object's per-process memory limit.
247 JobObject,
248 /// macOS + Bun: a fatal jetsam limit the host applied to itself.
249 Jetsam,
250 /// macOS otherwise: resident size checked at each heartbeat.
251 Heartbeat,
252 /// No cap on this platform.
253 Unenforced,
254 }
255
256 impl MemoryEnforcement {
257 /// What this platform does for `kind`, before the host has confirmed it.
258 #[must_use]
259 pub fn planned(kind: HostRuntimeKind) -> Self {
260 if cfg!(target_os = "linux") {
261 Self::Rlimit
262 } else if cfg!(windows) {
263 Self::JobObject
264 } else if cfg!(target_os = "macos") {
265 match kind {
266 HostRuntimeKind::Bun => Self::Jetsam,
267 HostRuntimeKind::Node => Self::Heartbeat,
268 }
269 } else {
270 Self::Unenforced
271 }
272 }
273
274 /// One line for `/plugin` and doctor.
275 #[must_use]
276 pub fn describe(self, cap: u64) -> String {
277 let mib = cap / (1024 * 1024);
278 match self {
279 Self::Rlimit => format!(
280 "memory cap {mib} MiB (kernel RLIMIT_DATA; plugin child processes inherit it)"
281 ),
282 Self::JobObject => format!(
283 "memory cap {mib} MiB (Job Object per-process limit; plugin child processes included)"
284 ),
285 Self::Jetsam => format!(
286 "memory cap {mib} MiB (kernel jetsam limit; the host is killed past it; plugin child processes are not covered)"
287 ),
288 Self::Heartbeat => format!(
289 "memory cap {mib} MiB (resident size checked at each heartbeat; on macOS only the Bun host gets a kernel limit)"
290 ),
291 Self::Unenforced => "no memory cap on this platform".to_string(),
292 }
293 }
294 }
295
296 /// The default memory cap. Measured once in a Linux container (2026-09-29),
297 /// not in CI: under `RLIMIT_DATA` Node 24 aborts when it creates the host's
298 /// watchdog Worker at 512 MiB, and Bun 1.4 aborts at startup at 256 MiB.
299 /// Both started and ran at 1 GiB and failed an allocation past it. An idle
300 /// host used 34–67 MB resident. On hosted x64 Linux (Node 24.21) each V8
301 /// isolate's executable code range is charged to `RLIMIT_DATA` in full, and
302 /// a default-sized second isolate aborted the host at this cap; the
303 /// watchdog Worker therefore asks for a 16 MiB code range (`src/main.ts`).
304 /// Plugins cannot start Workers (`denyNativeCode`), so it is the only one.
305 pub const HOST_MEMORY_CAP: u64 = 1 << 30;
306
307 /// Runtime flags, before the bundle path. Keep in sync with
308 /// `extension-host/test/harness.mjs` (`HOST_ARGS`).
309 ///
310 /// - Node: a 256 MB old-space cap, `__proto__` throws, no native addons,
311 /// and the [`crate::dependencies::NODE_NATIVE_CODE_FLAGS`] this Node
312 /// accepts (`node:sqlite`, and `node:ffi` where it exists). On Windows it
313 /// also keeps symlinked paths as given, so module resolution never lstats
314 /// the drive root the LPAC cannot read.
315 /// - Bun ignores all of those except `--no-addons`. Instead it gets
316 /// `--no-install`, because Bun otherwise fetches a missing package from npm
317 /// while plugin code is running. It also gets `--no-env-file` and
318 /// `--config=<null device>`, because Bun otherwise loads `.env` and
319 /// `bunfig.toml` (which can preload code) from the working directory, and
320 /// that directory is the host's writable data dir. `bun:ffi` and the rest
321 /// are locked in the host itself (`src/runtime.ts`).
322 #[must_use]
323 pub(crate) fn runtime_args(runtime: &HostRuntime) -> Vec<String> {
324 if runtime.compiled {
325 return Vec::new();
326 }
327 base_runtime_args(runtime.kind)
328 .iter()
329 .chain(&runtime.native_code_flags)
330 .map(|arg| (*arg).to_string())
331 .collect()
332 }
333
334 /// Environment for the runtime. Keep in sync with `HOST_ENV` in
335 /// `extension-host/test/harness.mjs`.
336 ///
337 /// - Both: `NODE_OPTIONS` is blanked. The child-environment allowlist passes
338 /// the user's value through (shell tools need it), and a `--require` or
339 /// `--import` preload there would run before the host's native-code
340 /// lockdown. Node honours it; Bun 1.4.0 ignored a `--require` in it when
341 /// checked (2026-09-30), and it is blanked for Bun too in case a later
342 /// Bun does not. `BUN_OPTIONS`, which Bun does honour, is not in that
343 /// allowlist.
344 /// - Bun: no ShadowRealm, engine-wide (a realm imports a fresh `bun:ffi`;
345 /// `node:vm` contexts would otherwise hand the constructor out). The host
346 /// refuses to start without it.
347 #[must_use]
348 pub(crate) fn runtime_env(kind: HostRuntimeKind) -> Vec<(String, String)> {
349 let mut env = vec![("NODE_OPTIONS".to_string(), String::new())];
350 if kind == HostRuntimeKind::Bun {
351 env.push(("BUN_JSC_useShadowRealm".to_string(), "0".to_string()));
352 env.push(("BUN_OPTIONS".to_string(), String::new()));
353 // This variable otherwise makes a standalone image act as the Bun CLI.
354 env.push(("BUN_BE_BUN".to_string(), "0".to_string()));
355 }
356 env
357 }
358
359 fn base_runtime_args(kind: HostRuntimeKind) -> &'static [&'static str] {
360 match kind {
361 #[cfg(not(windows))]
362 HostRuntimeKind::Node => &[
363 "--max-old-space-size=256",
364 "--disable-proto=throw",
365 "--no-addons",
366 ],
367 // Node resolves the entry bundle with realpathSync, which lstats every
368 // ancestor from `C:\`. A Windows LPAC cannot read the drive root, so
369 // the host died before its handshake (EPERM, lstat 'C:\'). Keep the
370 // granted path as given; the LPAC still decides every access.
371 #[cfg(windows)]
372 HostRuntimeKind::Node => &[
373 "--max-old-space-size=256",
374 "--disable-proto=throw",
375 "--no-addons",
376 "--preserve-symlinks",
377 "--preserve-symlinks-main",
378 ],
379 #[cfg(windows)]
380 HostRuntimeKind::Bun => &[
381 "--no-install",
382 "--no-env-file",
383 "--config=NUL",
384 "--no-addons",
385 ],
386 #[cfg(not(windows))]
387 HostRuntimeKind::Bun => &[
388 "--no-install",
389 "--no-env-file",
390 "--config=/dev/null",
391 "--no-addons",
392 ],
393 }
394 }
395
396 /// Top-level entries of a Codewhale home the host may read: its own bundle
397 /// and data (`extension-host`), and plugin code (the staged snapshots live
398 /// under `plugins/.runtime`). Everything else in a Codewhale home — secrets,
399 /// tokens, config and its backups, sessions, state, tool outputs, history —
400 /// is denied.
401 const HOST_READABLE_HOME_ENTRIES: &[&str] = &["extension-host", "plugins", "builtin-plugins"];
402
403 /// Codewhale-home entries denied by name even before they exist, so a store
404 /// created after the host started is still covered. Existing entries are
405 /// denied by enumeration (`host_denied_read_paths`).
406 const HOST_DENIED_HOME_ENTRIES: &[&str] = &[
407 "secrets",
408 "credentials",
409 "tokens",
410 "state",
411 "state.db",
412 "sessions",
413 "session-archives",
414 "session_index.jsonl",
415 "tool_outputs",
416 "composer_history.txt",
417 "composer_history.jsonl",
418 "remote-control",
419 "integrations",
420 "audit.log",
421 "logs",
422 "memory",
423 "mcp.json",
424 "mcp.json.bak",
425 "config.toml.bak",
426 "settings.toml",
427 ];
428
429 /// Paths the host process must never read, even though the sandbox otherwise
430 /// grants full-disk read, and the exceptions inside them: the curated
431 /// credential-store defaults; Codewhale's homes (the runtime home, the ambient
432 /// `~/.codewhale`, and the legacy `~/.deepseek`) except
433 /// [`HOST_READABLE_HOME_ENTRIES`]; and the Codex and DSH homes whose
434 /// credential files Codewhale itself reads. Blocking.
435 ///
436 /// With `whole_homes` (bubblewrap, which can mask only what exists) each home
437 /// is denied whole and its readable entries come back as exceptions to bind
438 /// again. Without it (Seatbelt, which matches paths that do not exist yet)
439 /// every other entry is denied by name and there are no exceptions.
440 ///
441 /// The plugin tier also denies the builtin tier's data directory, which lies
442 /// inside the readable `extension-host/` entry: plugin code never reads
443 /// tier-0 state. Planning materializes that sibling before the wrapper masks
444 /// it, because bubblewrap cannot mask a directory that does not exist yet.
445 pub(crate) fn host_denied_read_paths(
446 tier: HostTier,
447 home: &Path,
448 whole_homes: bool,
449 ) -> (Vec<PathBuf>, Vec<PathBuf>) {
450 let mut paths = crate::sandbox::read_guard::ReadDenylist::build(true, &[], &[]).subtree_paths();
451 let mut exceptions: Vec<PathBuf> = Vec::new();
452 let mut push = |path: PathBuf| {
453 if !paths.contains(&path) {
454 paths.push(path);
455 }
456 };
457 let user_home = codewhale_paths::user_home();
458 let mut roots = vec![home.to_path_buf()];
459 roots.extend(codewhale_config::codewhale_home().ok());
460 if let Some(user) = &user_home {
461 roots.push(user.join(codewhale_config::CODEWHALE_APP_DIR));
462 roots.push(user.join(".deepseek"));
463 }
464 for root in roots {
465 // The builtin tier's data directory, inside the readable
466 // `extension-host/` entry of every Codewhale home: not for plugin code.
467 let tier_zero_data =
468 (tier == HostTier::Plugin).then(|| tier_data_dir(&root, HostTier::Builtin));
469 if whole_homes {
470 for entry in HOST_READABLE_HOME_ENTRIES {
471 let readable = root.join(entry);
472 if !exceptions.contains(&readable) {
473 exceptions.push(readable);
474 }
475 }
476 if let Some(dir) = tier_zero_data {
477 push(dir);
478 }
479 push(root);
480 continue;
481 }
482 let mut names: Vec<std::ffi::OsString> = HOST_DENIED_HOME_ENTRIES
483 .iter()
484 .chain(std::iter::once(&codewhale_config::CONFIG_FILE_NAME))
485 .map(std::ffi::OsString::from)
486 .collect();
487 if let Ok(entries) = std::fs::read_dir(&root) {
488 names.extend(
489 entries
490 .filter_map(Result::ok)
491 .map(|entry| entry.file_name()),
492 );
493 }
494 // Seatbelt matches the kernel-resolved path, and a name that does not
495 // exist yet cannot be canonicalized later, so deny it under both the
496 // given and the resolved spelling of its (existing) root.
497 let resolved = std::fs::canonicalize(&root).ok();
498 if let Some(dir) = tier_zero_data {
499 if let Some(resolved) = &resolved {
500 push(tier_data_dir(resolved, HostTier::Builtin));
501 }
502 push(dir);
503 }
504 for name in names {
505 let readable = name
506 .to_str()
507 .is_some_and(|name| HOST_READABLE_HOME_ENTRIES.contains(&name));
508 if !readable {
509 if let Some(resolved) = &resolved {
510 push(resolved.join(&name));
511 }
512 push(root.join(name));
513 }
514 }
515 }
516 // Codex's home (ChatGPT OAuth tokens in `auth.json`) and the DSH home
517 // (`.credentials.yaml`), wherever the environment points them.
518 if let Some(codex_home) = crate::oauth::auth_file_path().parent() {
519 push(codex_home.to_path_buf());
520 }
521 if let Some(user) = &user_home {
522 push(user.join(".codex"));
523 push(user.join(".dsh"));
524 }
525 if let Some(dsh_home) = codewhale_config::default_dsh_credentials_path().parent() {
526 push(dsh_home.to_path_buf());
527 }
528 (paths, exceptions)
529 }
530
531 /// Plan the host launch for `tier`: the runtime's flags, the bundle, then the
532 /// tier (`--tier=plugin|builtin`, which the host reads from its own argv).
533 /// Blocking (creates the tier's data dir, canonicalizes the deny-list, and on
534 /// Linux runs the bwrap probe); call from `spawn_blocking`.
535 pub(crate) fn plan_launch(
536 tier: HostTier,
537 runtime: &HostRuntime,
538 bundle: &Path,
539 home: &Path,
540 memory_cap: u64,
541 ) -> Result<HostLaunch, String> {
542 let data = host_data_dir(home, tier)?;
543 // Bubblewrap cannot mask a root that does not exist yet. Materialize the
544 // sibling before a Native host gets its immutable read-deny projection.
545 if tier == HostTier::Plugin {
546 host_data_dir(home, HostTier::Builtin)?;
547 }
548 let mut args = runtime_args(runtime);
549 if !runtime.compiled {
550 args.push(bundle.to_string_lossy().into_owned());
551 }
552 args.push(tier.argv_flag());
553 #[cfg(windows)]
554 if tier == HostTier::Plugin {
555 let sandbox =
556 super::windows::NativeSandbox::prepare(runtime, bundle, home, &data, memory_cap)?;
557 return Ok(HostLaunch {
558 tier,
559 program: sandbox.program.clone(),
560 args,
561 cwd: data,
562 sandbox: HostSandbox::Wrapped("windows-lpac".into()),
563 sandbox_env: vec![("CODEWHALE_SANDBOX".into(), "windows-lpac".into())],
564 runtime: runtime.clone(),
565 runtime_env: runtime_env(runtime.kind),
566 memory_cap,
567 memory: MemoryEnforcement::JobObject,
568 builtin_modules: tier::BUILTIN_MODULES,
569 windows: Some(sandbox),
570 });
571 }
572 let wrapped = wrap_host(tier, &runtime.path, &args, &data, home);
573 host_launch(tier, runtime, args, data, memory_cap, wrapped)
574 }
575
576 /// The sandbox a `tier` host started now would get, planned (and on Linux
577 /// probed) exactly as [`plan_launch`] does, for doctor. Blocking; creates the
578 /// data dir, as a launch would.
579 pub(crate) fn planned_sandbox(
580 tier: HostTier,
581 runtime: &HostRuntime,
582 home: &Path,
583 ) -> Result<HostSandbox, String> {
584 let data = host_data_dir(home, tier)?;
585 if tier == HostTier::Plugin {
586 host_data_dir(home, HostTier::Builtin)?;
587 }
588 #[cfg(windows)]
589 if tier == HostTier::Plugin {
590 let bundle = super::materialize_bundle(home)?;
591 return Ok(
592 match super::windows::NativeSandbox::prepare(
593 runtime,
594 &bundle,
595 home,
596 &data,
597 HOST_MEMORY_CAP,
598 ) {
599 Ok(_) => HostSandbox::Wrapped("windows-lpac".into()),
600 Err(error) => HostSandbox::Unsandboxed(error),
601 },
602 );
603 }
604 Ok(
605 match wrap_host(tier, &runtime.path, &runtime_args(runtime), &data, home) {
606 Ok(wrapped) => HostSandbox::Wrapped(wrapped.name),
607 Err(reason) => HostSandbox::Unsandboxed(reason),
608 },
609 )
610 }
611
612 /// A tier's data directory: its host's working directory and, under the OS
613 /// sandbox, its only writable root. Pure.
614 ///
615 /// The plugin tier keeps the directory it has always had,
616 /// `extension-host/data`, because the per-plugin directories under it
617 /// ([`plugin_data_dir`]) hold installed plugins' data and moving it would lose
618 /// that. The builtin tier's is a sibling, `extension-host/data-builtin`, not a
619 /// child: a child would lie inside the plugin tier's writable root, and the
620 /// host sandbox has no per-subpath write deny, so plugin code could then write
621 /// tier-0 state.
622 pub(crate) fn tier_data_dir(home: &Path, tier: HostTier) -> PathBuf {
623 let base = home.join("extension-host");
624 match tier {
625 HostTier::Plugin => base.join("data"),
626 HostTier::Builtin => base.join("data-builtin"),
627 }
628 }
629
630 /// The directory an owner's code is given as its own: a plugin's
631 /// ([`plugin_data_dir`], unchanged), or `modules/<module>` under the builtin
632 /// tier's data directory for a built-in module (`plugin_name` is the module
633 /// name, which [`HostTier::check_owner_id`] has restricted to a plain name).
634 /// Pure.
635 pub(crate) fn owner_data_dir(
636 home: &Path,
637 tier: HostTier,
638 plugin_id: &str,
639 plugin_name: &str,
640 ) -> PathBuf {
641 match tier {
642 HostTier::Plugin => plugin_data_dir(home, plugin_id, plugin_name),
643 HostTier::Builtin => tier_data_dir(home, tier).join("modules").join(plugin_name),
644 }
645 }
646
647 /// One plugin's own directory inside the host's data dir, which is the host
648 /// sandbox's writable root: stable for one plugin id (so it survives updates
649 /// and restarts), distinct per plugin, and a single path component under
650 /// `plugins/`. The id contains slashes and the name alone can collide between
651 /// scopes, so the name is joined to a digest of the id. Pure.
652 pub(crate) fn plugin_data_dir(home: &Path, plugin_id: &str, plugin_name: &str) -> PathBuf {
653 use sha2::{Digest, Sha256};
654 let digest = Sha256::digest(plugin_id.as_bytes());
655 let short: String = digest
656 .iter()
657 .take(6)
658 .map(|byte| format!("{byte:02x}"))
659 .collect();
660 let name: String = plugin_name
661 .chars()
662 .map(|c| {
663 if c.is_ascii_alphanumeric() || c == '-' {
664 c
665 } else {
666 '_'
667 }
668 })
669 .take(64)
670 .collect();
671 home.join("extension-host")
672 .join("data")
673 .join("plugins")
674 .join(format!("{name}-{short}"))
675 }
676
677 fn host_data_dir(home: &Path, tier: HostTier) -> Result<PathBuf, String> {
678 let data = tier_data_dir(home, tier);
679 std::fs::create_dir_all(&data)
680 .map_err(|error| format!("cannot create {}: {error}", data.display()))?;
681 Ok(data)
682 }
683
684 /// A host command wrapped in an OS sandbox ([`wrap_host`]).
685 #[derive(Debug)]
686 struct Wrapped {
687 /// `sandbox::SandboxType`'s name for the wrapper.
688 name: String,
689 /// The wrapper's argv, ending with the runtime's.
690 command: Vec<String>,
691 /// Environment the wrapper adds (`CODEWHALE_SANDBOX`, …).
692 env: Vec<(String, String)>,
693 }
694
695 /// Why the host runs on Windows without an OS sandbox.
696 const WINDOWS_UNSANDBOXED: &str = "Windows has no host sandbox yet; its Job Object contains the process tree, which is not isolation";
697
698 /// Wrap `program args` in the host's OS sandbox (module docs), or say why it
699 /// has none here. Blocking.
700 fn wrap_host(
701 tier: HostTier,
702 program: &Path,
703 args: &[String],
704 data: &Path,
705 home: &Path,
706 ) -> Result<Wrapped, String> {
707 use crate::sandbox::{CommandSpec, SandboxManager, SandboxPolicy, SandboxType};
708 if cfg!(windows) {
709 return Err(WINDOWS_UNSANDBOXED.to_string());
710 }
711 let spec = |args: Vec<String>| {
712 CommandSpec::program(
713 &program.to_string_lossy(),
714 args,
715 data.to_path_buf(),
716 Duration::ZERO,
717 )
718 .with_policy(SandboxPolicy::WorkspaceWrite {
719 writable_roots: Vec::new(),
720 network_access: false,
721 exclude_tmpdir: false,
722 exclude_slash_tmp: false,
723 })
724 };
725 let bwrap = cfg!(all(target_os = "linux", not(target_env = "ohos")));
726 let mut manager = SandboxManager::with_bwrap_preference(bwrap);
727 let (denied, exceptions) = host_denied_read_paths(tier, home, bwrap);
728 manager.set_denied_read_subpaths(denied);
729 manager.set_denied_read_exceptions(exceptions);
730 let env = manager.prepare(&spec(args.to_vec()));
731 if matches!(env.sandbox_type, SandboxType::None) {
732 return Err(no_wrapper_reason());
733 }
734 #[cfg(all(target_os = "linux", not(target_env = "ohos")))]
735 if matches!(env.sandbox_type, SandboxType::LinuxBubblewrap) {
736 probe_bwrap(
737 &manager
738 .prepare(&spec(vec!["--version".to_string()]))
739 .command,
740 data,
741 )?;
742 }
743 Ok(Wrapped {
744 name: env.sandbox_type.to_string(),
745 command: env.command,
746 env: env.env.into_iter().collect(),
747 })
748 }
749
750 /// The launch for a [`wrap_host`] outcome: the wrapper's argv, or the
751 /// pinned Builtin runtime itself, carrying the reason `/plugin` shows. Native
752 /// code is refused without the verified wrapper. Pure.
753 fn host_launch(
754 tier: HostTier,
755 runtime: &HostRuntime,
756 args: Vec<String>,
757 data: PathBuf,
758 memory_cap: u64,
759 wrapped: Result<Wrapped, String>,
760 ) -> Result<HostLaunch, String> {
761 let (program, args, sandbox, sandbox_env) = match wrapped {
762 Ok(wrapped) => {
763 let mut command = wrapped.command.into_iter();
764 let program = command
765 .next()
766 .ok_or("sandbox wrapper produced an empty command")?;
767 (
768 PathBuf::from(program),
769 command.collect(),
770 HostSandbox::Wrapped(wrapped.name),
771 wrapped.env,
772 )
773 }
774 Err(reason) if tier == HostTier::Plugin => {
775 return Err(format!(
776 "Native extensions require a verified OS sandbox: {reason}"
777 ));
778 }
779 // Only the pinned Builtin tier may run without filesystem/network
780 // isolation. Every effect is still admitted by Rust operation tickets.
781 Err(reason) => (
782 runtime.path.clone(),
783 args,
784 HostSandbox::Unsandboxed(reason),
785 Vec::new(),
786 ),
787 };
788 Ok(HostLaunch {
789 tier,
790 program,
791 args,
792 cwd: data,
793 sandbox,
794 sandbox_env,
795 runtime: runtime.clone(),
796 runtime_env: runtime_env(runtime.kind),
797 memory_cap,
798 memory: MemoryEnforcement::planned(runtime.kind),
799 builtin_modules: tier::BUILTIN_MODULES,
800 #[cfg(windows)]
801 windows: None,
802 })
803 }
804
805 /// Why [`wrap_host`] found no wrapper on this platform.
806 #[cfg(all(target_os = "linux", not(target_env = "ohos")))]
807 fn no_wrapper_reason() -> String {
808 format!(
809 "bwrap unavailable: {} is not an executable file (install bubblewrap)",
810 crate::sandbox::bwrap::BWRAP_PATH
811 )
812 }
813
814 #[cfg(target_os = "macos")]
815 fn no_wrapper_reason() -> String {
816 "Seatbelt (sandbox-exec) is unavailable".to_string()
817 }
818
819 #[cfg(not(any(
820 target_os = "macos",
821 all(target_os = "linux", not(target_env = "ohos"))
822 )))]
823 fn no_wrapper_reason() -> String {
824 "no OS sandbox for the host on this platform".to_string()
825 }
826
827 /// How long the bwrap probe may take. A working bwrap runs
828 /// `<runtime> --version` in well under a second, even on a loaded CI runner.
829 #[cfg(all(target_os = "linux", not(target_env = "ohos")))]
830 const BWRAP_PROBE_DEADLINE: Duration = Duration::from_secs(10);
831
832 /// Run the finished bwrap wrapper around `<runtime> --version` (`command`),
833 /// with no environment, to learn whether bwrap works on this host: it may be
834 /// installed yet unable to create its namespaces. Blocking.
835 #[cfg(all(target_os = "linux", not(target_env = "ohos")))]
836 fn probe_bwrap(command: &[String], cwd: &Path) -> Result<(), String> {
837 use std::io::Read as _;
838 use wait_timeout::ChildExt as _;
839 let (program, args) = command
840 .split_first()
841 .ok_or("the sandbox wrapper produced an empty command")?;
842 let mut child = std::process::Command::new(program)
843 .args(args)
844 .current_dir(cwd)
845 .env_clear()
846 .stdin(Stdio::null())
847 .stdout(Stdio::null())
848 .stderr(Stdio::piped())
849 .spawn()
850 .map_err(|error| format!("bwrap unavailable: {program} does not start ({error})"))?;
851 let status = match child.wait_timeout(BWRAP_PROBE_DEADLINE) {
852 Ok(Some(status)) => status,
853 outcome => {
854 let _ = child.kill();
855 let _ = child.wait();
856 return Err(match outcome {
857 Err(error) => format!("bwrap unavailable: waiting for the probe failed ({error})"),
858 _ => format!(
859 "bwrap unavailable: the probe did not finish within {BWRAP_PROBE_DEADLINE:?}"
860 ),
861 });
862 }
863 };
864 let mut stderr = String::new();
865 if let Some(mut pipe) = child.stderr.take() {
866 let _ = pipe.read_to_string(&mut stderr);
867 }
868 bwrap_probe_verdict(status.success(), &status.to_string(), &stderr)
869 }
870
871 /// What a wrapped `<runtime> --version` run says about bwrap here: `Ok` when
872 /// it ran, otherwise the reason `/plugin` and doctor show — bwrap's first
873 /// line of stderr (or the exit status), and a hint when that line is about
874 /// the user namespace bwrap could not create. Pure.
875 #[cfg(any(test, all(target_os = "linux", not(target_env = "ohos"))))]
876 fn bwrap_probe_verdict(succeeded: bool, status: &str, stderr: &str) -> Result<(), String> {
877 if succeeded {
878 return Ok(());
879 }
880 let first = stderr.lines().map(str::trim).find(|line| !line.is_empty());
881 let mut detail: String = first.map_or_else(
882 || status.to_string(),
883 |line| line.chars().take(240).collect(),
884 );
885 if first.is_some_and(|line| line.contains("namespace") || line.contains("uid map")) {
886 detail.push_str(
887 "; unprivileged user namespaces look blocked here (on Ubuntu 24.04 and later: the kernel.apparmor_restrict_unprivileged_userns sysctl)",
888 );
889 }
890 Err(format!("bwrap unavailable ({detail})"))
891 }
892
893 /// The kernel-enforced memory cap on Linux: `RLIMIT_DATA`, applied in the
894 /// child between fork and exec, so only the host (and what it starts) is
895 /// limited. Soft and hard limit are both set, so plugin code cannot raise
896 /// it. An unprivileged process cannot raise its hard limit, so when the
897 /// inherited hard limit is already below `cap` the host gets that lower
898 /// limit instead of failing to spawn with `EPERM`.
899 ///
900 /// Known limit: in that case `/plugin` and doctor still name the configured
901 /// cap, not the lower inherited one.
902 #[cfg(target_os = "linux")]
903 pub(super) fn limit_child_memory(command: &mut tokio::process::Command, cap: u64) {
904 // SAFETY: the closure runs in the forked child before exec and calls only
905 // `getrlimit` and `setrlimit`, which are async-signal-safe; it allocates
906 // nothing.
907 unsafe {
908 command.pre_exec(move || {
909 let mut inherited = libc::rlimit {
910 rlim_cur: 0,
911 rlim_max: 0,
912 };
913 if libc::getrlimit(libc::RLIMIT_DATA, &raw mut inherited) != 0 {
914 return Err(std::io::Error::last_os_error());
915 }
916 let cap = (cap as libc::rlim_t).min(inherited.rlim_max);
917 let limit = libc::rlimit {
918 rlim_cur: cap,
919 rlim_max: cap,
920 };
921 if libc::setrlimit(libc::RLIMIT_DATA, &raw const limit) != 0 {
922 return Err(std::io::Error::last_os_error());
923 }
924 Ok(())
925 });
926 }
927 }
928
929 #[cfg(not(target_os = "linux"))]
930 pub(super) fn limit_child_memory(_command: &mut tokio::process::Command, _cap: u64) {}
931
932 /// Resident size of `pid` in bytes, for the macOS memory-cap check.
933 #[cfg(target_os = "macos")]
934 pub(crate) fn resident_bytes(pid: u32) -> Option<u64> {
935 let mut info = std::mem::MaybeUninit::<libc::proc_taskinfo>::zeroed();
936 let size = std::mem::size_of::<libc::proc_taskinfo>() as libc::c_int;
937 // SAFETY: `info` is a correctly sized, writable `proc_taskinfo` buffer.
938 let written = unsafe {
939 libc::proc_pidinfo(
940 pid as libc::c_int,
941 libc::PROC_PIDTASKINFO,
942 0,
943 info.as_mut_ptr().cast(),
944 size,
945 )
946 };
947 // SAFETY: a full-size write initialized the struct.
948 (written == size).then(|| unsafe { info.assume_init() }.pti_resident_size)
949 }
950
951 /// Platforms without a supervisor-side check (Linux uses `RLIMIT_DATA`,
952 /// Windows the Job Object).
953 #[cfg(not(target_os = "macos"))]
954 pub(crate) fn resident_bytes(_pid: u32) -> Option<u64> {
955 None
956 }
957
958 /// Report the observed exit and configured cap. SIGKILL alone cannot identify
959 /// jetsam: an operator or another process can send the same signal.
960 fn exit_reason(
961 status: std::process::ExitStatus,
962 memory: Option<MemoryEnforcement>,
963 cap: u64,
964 ) -> String {
965 let reason = format!("exited with {status}");
966 #[cfg(unix)]
967 {
968 use std::os::unix::process::ExitStatusExt as _;
969 if status.signal() == Some(libc::SIGKILL) && memory == Some(MemoryEnforcement::Jetsam) {
970 return format!(
971 "{reason}; configured kernel memory limit: {} MiB; SIGKILL cause unavailable",
972 cap / (1024 * 1024)
973 );
974 }
975 }
976 #[cfg(not(unix))]
977 let _ = (memory, cap);
978 reason
979 }
980
981 /// What one host-originated request is told about its own life: its id on the
982 /// channel and the token that fires when the request is cancelled (the host's
983 /// `$/cancel`, the owner's revocation, or the host's exit). A handler that
984 /// waits for anything must wait on `cancel` too; one that ignores it is
985 /// abandoned [`CANCEL_GRACE`] after it fires and its answer is dropped.
986 pub(crate) struct HostRequestContext {
987 pub id: u64,
988 pub cancel: CancellationToken,
989 /// The channel's kill switch, for a handler that finds the host in
990 /// violation of the protocol (it ends the host, like a bad frame).
991 kill: mpsc::Sender<String>,
992 }
993
994 #[cfg(test)]
995 impl HostRequestContext {
996 /// A context not attached to any channel, with the receiver its violations
997 /// arrive on and the token that cancels it.
998 pub(crate) fn for_test(id: u64) -> (Self, mpsc::Receiver<String>, CancellationToken) {
999 let (kill, violations) = mpsc::channel(8);
1000 let cancel = CancellationToken::new();
1001 (
1002 Self {
1003 id,
1004 cancel: cancel.clone(),
1005 kill,
1006 },
1007 violations,
1008 cancel,
1009 )
1010 }
1011 }
1012
1013 impl HostRequestContext {
1014 /// Report a protocol violation by the host: the host process is ended.
1015 pub(crate) fn violation(&self, reason: String) {
1016 let _ = self.kill.try_send(format!("protocol violation: {reason}"));
1017 }
1018 }
1019
1020 /// Callbacks from the channel into the manager.
1021 #[async_trait]
1022 pub(crate) trait HostEvents: Send + Sync + 'static {
1023 fn register(&self, params: &protocol::RegisterParams) -> RegisterResult;
1024 fn unregister(&self, params: &protocol::UnregisterParams);
1025 fn faulted(&self, params: &protocol::FaultedParams);
1026 fn log(&self, params: &protocol::LogParams);
1027 fn exited(&self, host_generation: u64, reason: String, stderr_tail: String);
1028
1029 /// A current, non-revoked call received a correlated reply. Late replies
1030 /// and heartbeat answers do not pass here; the monitor validates pongs.
1031 fn responded(&self) {}
1032
1033 /// Answer one host-originated request. Every such request runs as its own
1034 /// task (`start_host_request`), so a handler may take as long as it needs
1035 /// without holding up the reader. The registry requests are quick and
1036 /// synchronous; a request that has to wait (a tool call the host asked the
1037 /// core to make, later) overrides this and observes `cx.cancel`.
1038 async fn host_request(
1039 &self,
1040 request: HostRequest,
1041 cx: HostRequestContext,
1042 ) -> Result<Value, RpcErrorWire> {
1043 registry_host_request(self, request, &cx)
1044 }
1045 }
1046
1047 /// The answer to a registry request (`registry/register`,
1048 /// `registry/unregister`), which is quick and synchronous, and the refusal of a
1049 /// `core/call` an events implementation does not serve. The default
1050 /// [`HostEvents::host_request`], and what an overriding one falls back to.
1051 pub(crate) fn registry_host_request<E: HostEvents + ?Sized>(
1052 events: &E,
1053 request: HostRequest,
1054 cx: &HostRequestContext,
1055 ) -> Result<Value, RpcErrorWire> {
1056 tracing::trace!(target: "extension_host", id = cx.id, "host request");
1057 // A request cancelled before its handler began does nothing.
1058 if cx.cancel.is_cancelled() {
1059 return Err(RpcErrorWire {
1060 code: error_code::CANCELLED,
1061 message: "cancelled".to_string(),
1062 data: None,
1063 });
1064 }
1065 match request {
1066 HostRequest::Register(params) => Ok(serde_json::to_value(events.register(&params))
1067 .unwrap_or_else(|_| json!({"refused": "internal"}))),
1068 HostRequest::Unregister(params) => {
1069 events.unregister(&params);
1070 Ok(json!({}))
1071 }
1072 HostRequest::ExecutionRedeem(_)
1073 | HostRequest::CoreCall(_)
1074 | HostRequest::ProcLaunch(_)
1075 | HostRequest::ProcRead(_)
1076 | HostRequest::ProcWrite(_)
1077 | HostRequest::ProcClose(_)
1078 | HostRequest::NetStart(_)
1079 | HostRequest::NetFetch(_)
1080 | HostRequest::NetRead(_)
1081 | HostRequest::NetRelease(_)
1082 | HostRequest::NetClose(_) => Err(RpcErrorWire {
1083 code: error_code::REFUSED,
1084 message: "core/call is not served here".to_string(),
1085 data: None,
1086 }),
1087 }
1088 }
1089
1090 /// Why an in-flight host request was cancelled; decides what the host is told.
1091 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
1092 enum CancelReason {
1093 /// The host sent `$/cancel`: it has already stopped waiting, so whatever
1094 /// the handler produces afterwards is dropped.
1095 Host,
1096 /// The request's owner was revoked: the host is answered `Cancelled`.
1097 Revoked,
1098 /// The host exited: there is nobody to answer.
1099 Exit,
1100 }
1101
1102 struct InboundRequest {
1103 /// The plugin whose revocation cancels it.
1104 owner: String,
1105 cancel: CancellationToken,
1106 cancelled: Option<CancelReason>,
1107 }
1108
1109 /// The host-originated requests in flight, by the id the host gave them. At
1110 /// most [`protocol::MAX_INFLIGHT`] at a time (the same bound the host holds
1111 /// itself to), no id twice, and nothing admitted once the host has exited.
1112 #[derive(Default)]
1113 pub(crate) struct InboundRequests {
1114 table: Mutex<InboundTable>,
1115 }
1116
1117 #[derive(Default)]
1118 struct InboundTable {
1119 requests: HashMap<u64, InboundRequest>,
1120 closed: bool,
1121 }
1122
1123 impl InboundRequests {
1124 /// Start tracking request `id` of `owner`. `Err` is a protocol violation:
1125 /// the host reused an id still in flight, or has more requests in flight
1126 /// than its own limit allows.
1127 fn admit(&self, id: u64, owner: &str) -> Result<CancellationToken, String> {
1128 let mut table = self.table.lock().expect("inbound lock");
1129 if table.closed {
1130 return Err("host request after the host exited".to_string());
1131 }
1132 if table.requests.contains_key(&id) {
1133 return Err(format!("host request id {id} is already in flight"));
1134 }
1135 if table.requests.len() >= protocol::MAX_INFLIGHT {
1136 return Err(format!(
1137 "more than {} host requests in flight",
1138 protocol::MAX_INFLIGHT
1139 ));
1140 }
1141 let cancel = CancellationToken::new();
1142 table.requests.insert(
1143 id,
1144 InboundRequest {
1145 owner: owner.to_string(),
1146 cancel: cancel.clone(),
1147 cancelled: None,
1148 },
1149 );
1150 Ok(cancel)
1151 }
1152
1153 /// The host's `$/cancel {id}`. An id that is not in flight (already
1154 /// answered, or never sent) is ignored: the cancel raced the answer.
1155 fn cancel_by_host(&self, id: u64) {
1156 if let Some(request) = self
1157 .table
1158 .lock()
1159 .expect("inbound lock")
1160 .requests
1161 .get_mut(&id)
1162 {
1163 request.fire(CancelReason::Host);
1164 }
1165 }
1166
1167 /// Cancel every in-flight request of `plugin_id`: its owner was revoked.
1168 fn cancel_owner(&self, plugin_id: &str) {
1169 for request in self
1170 .table
1171 .lock()
1172 .expect("inbound lock")
1173 .requests
1174 .values_mut()
1175 .filter(|request| request.owner == plugin_id)
1176 {
1177 request.fire(CancelReason::Revoked);
1178 }
1179 }
1180
1181 /// The host exited: cancel everything and admit nothing more.
1182 fn cancel_all(&self) {
1183 let mut table = self.table.lock().expect("inbound lock");
1184 table.closed = true;
1185 for request in table.requests.values_mut() {
1186 request.fire(CancelReason::Exit);
1187 }
1188 }
1189
1190 /// The handler is done (or abandoned): forget the request and say why it
1191 /// was cancelled, if it was.
1192 fn finish(&self, id: u64) -> Option<CancelReason> {
1193 self.table
1194 .lock()
1195 .expect("inbound lock")
1196 .requests
1197 .remove(&id)
1198 .and_then(|request| request.cancelled)
1199 }
1200
1201 #[cfg(test)]
1202 fn in_flight(&self) -> usize {
1203 self.table.lock().expect("inbound lock").requests.len()
1204 }
1205 }
1206
1207 impl InboundRequest {
1208 /// Cancel for `reason`; the first reason stands.
1209 fn fire(&mut self, reason: CancelReason) {
1210 self.cancelled.get_or_insert(reason);
1211 self.cancel.cancel();
1212 }
1213 }
1214
1215 /// Receives one request's outcome.
1216 pub(crate) type CallReceiver = oneshot::Receiver<Result<Value, HostCallError>>;
1217
1218 struct PendingCall {
1219 tx: oneshot::Sender<Result<Value, HostCallError>>,
1220 /// Plugin whose revocation cancels this call.
1221 owner: Option<String>,
1222 /// Exact never-reused registration handle, when this is a contribution call.
1223 handle: Option<u64>,
1224 revoked: bool,
1225 heartbeat: bool,
1226 }
1227
1228 #[derive(Default)]
1229 struct Handshake {
1230 hello: Option<oneshot::Sender<protocol::HelloParams>>,
1231 ready: Option<oneshot::Sender<()>>,
1232 }
1233
1234 pub(crate) struct HostProcess {
1235 /// The tier this process was launched for.
1236 pub tier: HostTier,
1237 /// The generation of its tier's host this process is (what the manager
1238 /// bumps per launch); a capability ticket is bound to it.
1239 pub generation: u64,
1240 pub pid: Option<u32>,
1241 /// The runtime the Rust side launched; `host/hello` must agree.
1242 pub runtime: HostRuntime,
1243 /// Runtime version as reported by the host in `host/hello`.
1244 pub runtime_version: std::sync::OnceLock<String>,
1245 pub memory_cap: u64,
1246 /// How the cap is enforced for this process, settled at the handshake.
1247 memory: Arc<std::sync::OnceLock<MemoryEnforcement>>,
1248 pub sandbox: HostSandbox,
1249 tree: Arc<crate::process_tree::ProcessTree>,
1250 #[cfg(windows)]
1251 windows: Option<super::windows::NativeSandbox>,
1252 outbound: mpsc::Sender<Vec<u8>>,
1253 pending: Arc<Mutex<HashMap<u64, PendingCall>>>,
1254 /// Requests the host sent, running as tasks ([`start_host_request`]).
1255 inbound: Arc<InboundRequests>,
1256 /// Admission checks and sealing hold `pending`; the manager also reads
1257 /// this flag to avoid activation while the exit callback is still pending.
1258 admission_closed: AtomicBool,
1259 next_id: AtomicU64,
1260 /// Heartbeat pings among those ids; tests count only the core's own work.
1261 #[cfg(test)]
1262 heartbeats_sent: AtomicU64,
1263 stderr_tail: Arc<Mutex<VecDeque<u8>>>,
1264 exited: tokio::sync::watch::Receiver<bool>,
1265 kill: mpsc::Sender<String>,
1266 }
1267
1268 fn push_tail(tail: &Mutex<VecDeque<u8>>, bytes: &[u8]) {
1269 let mut tail = tail.lock().expect("stderr tail lock");
1270 tail.extend(bytes);
1271 while tail.len() > STDERR_TAIL_BYTES {
1272 tail.pop_front();
1273 }
1274 }
1275
1276 fn tail_string(tail: &Mutex<VecDeque<u8>>) -> String {
1277 let tail = tail.lock().expect("stderr tail lock");
1278 let bytes: Vec<u8> = tail.iter().copied().collect();
1279 String::from_utf8_lossy(&bytes).into_owned()
1280 }
1281
1282 /// The watcher/protocol remains one implementation for both launch mechanisms.
1283 enum HostChild {
1284 Tokio(tokio::process::Child),
1285 #[cfg(windows)]
1286 Native(super::windows::Child),
1287 }
1288 impl HostChild {
1289 async fn wait(&mut self) -> std::io::Result<std::process::ExitStatus> {
1290 match self {
1291 Self::Tokio(child) => child.wait().await,
1292 #[cfg(windows)]
1293 Self::Native(child) => child.wait().await,
1294 }
1295 }
1296 async fn kill(&mut self) -> std::io::Result<()> {
1297 match self {
1298 Self::Tokio(child) => child.kill().await,
1299 #[cfg(windows)]
1300 Self::Native(child) => child.kill().await,
1301 }
1302 }
1303 }
1304 struct SpawnedHost {
1305 child: HostChild,
1306 pid: Option<u32>,
1307 tree: Arc<crate::process_tree::ProcessTree>,
1308 stdin: tokio::process::ChildStdin,
1309 stdout: tokio::process::ChildStdout,
1310 stderr: tokio::process::ChildStderr,
1311 }
1312 fn spawn_host(
1313 launch: &HostLaunch,
1314 mut command: tokio::process::Command,
1315 _environment: &[(std::ffi::OsString, std::ffi::OsString)],
1316 ) -> Result<SpawnedHost, String> {
1317 #[cfg(windows)]
1318 if let Some(sandbox) = &launch.windows {
1319 let native = sandbox
1320 .spawn(&launch.args, _environment, launch.memory_cap)
1321 .map_err(|error| {
1322 format!("failed to launch the verified Windows Native host: {error}")
1323 })?;
1324 let stdin =
1325 tokio::process::ChildStdin::from_std(std::process::ChildStdin::from(native.stdin))
1326 .map_err(|error| format!("host stdin conversion failed: {error}"))?;
1327 let stdout =
1328 tokio::process::ChildStdout::from_std(std::process::ChildStdout::from(native.stdout))
1329 .map_err(|error| format!("host stdout conversion failed: {error}"))?;
1330 let stderr =
1331 tokio::process::ChildStderr::from_std(std::process::ChildStderr::from(native.stderr))
1332 .map_err(|error| format!("host stderr conversion failed: {error}"))?;
1333 let tree = Arc::clone(&native.child.tree);
1334 let pid = Some(native.child.pid);
1335 return Ok(SpawnedHost {
1336 child: HostChild::Native(native.child),
1337 pid,
1338 tree,
1339 stdin,
1340 stdout,
1341 stderr,
1342 });
1343 }
1344 // On Linux a failure to apply the memory cap in the child surfaces
1345 // here as a spawn error carrying only its errno, indistinguishable
1346 // from a failed exec, so the message names both.
1347 let mut child = command.spawn().map_err(|error| {
1348 if cfg!(target_os = "linux") {
1349 format!(
1350 "failed to start {}, or to apply its {} MiB memory cap (RLIMIT_DATA) before exec: {error}",
1351 launch.program.display(),
1352 launch.memory_cap / (1024 * 1024)
1353 )
1354 } else {
1355 format!("failed to start {}: {error}", launch.program.display())
1356 }
1357 })?;
1358 let pid = child.id();
1359 let tree = match crate::process_tree::ProcessTree::attach_tokio(&child) {
1360 Ok(tree) => Arc::new(tree),
1361 Err(error) => {
1362 let _ = child.start_kill();
1363 return Err(format!("failed to contain the extension host: {error}"));
1364 }
1365 };
1366 #[cfg(windows)]
1367 if let Err(error) = tree.limit_process_memory(launch.memory_cap) {
1368 let _ = tree.kill();
1369 let _ = child.start_kill();
1370 return Err(format!(
1371 "failed to cap the extension host's memory: {error}"
1372 ));
1373 }
1374 let stdin = child.stdin.take().ok_or("host stdin unavailable")?;
1375 let stdout = child.stdout.take().ok_or("host stdout unavailable")?;
1376 let stderr = child.stderr.take().ok_or("host stderr unavailable")?;
1377
1378 Ok(SpawnedHost {
1379 child: HostChild::Tokio(child),
1380 pid,
1381 tree,
1382 stdin,
1383 stdout,
1384 stderr,
1385 })
1386 }
1387
1388 impl HostProcess {
1389 #[cfg(windows)]
1390 pub(crate) async fn admit_windows_root(
1391 &self,
1392 authority: crate::plugins::types::PluginAuthority,
1393 ) -> Result<(), String> {
1394 let sandbox = self
1395 .windows
1396 .clone()
1397 .ok_or("Native host has no verified LPAC plan")?;
1398 let policy = crate::plugins::activation::extension_host_policy_enabled();
1399 if !policy {
1400 return Err("Native extensions are disabled".into());
1401 }
1402 #[cfg(test)]
1403 let env_scope = crate::test_support::env_scope_ticket();
1404 tokio::task::spawn_blocking(move || {
1405 #[cfg(test)]
1406 let _env_scope = crate::test_support::join_env_scope(env_scope);
1407 let _policy = crate::plugins::activation::PolicyScope::propagate(policy);
1408 let capability = crate::plugins::activation::PluginActivationCapability::Native;
1409 crate::plugins::registry::verify_plugin_component_authority(&authority, capability)?;
1410 let root =
1411 crate::plugins::agent_plugin::plugin_root_for_manifest(&authority.staged_manifest)
1412 .ok_or("reviewed runtime manifest has no bundle root")?;
1413 sandbox.admit_root(root)?;
1414 crate::plugins::registry::verify_plugin_component_authority(&authority, capability)
1415 })
1416 .await
1417 .map_err(|error| format!("Windows bundle admission worker failed: {error}"))?
1418 }
1419
1420 /// Spawn the host and complete the handshake. `expected_sha256` is the
1421 /// digest of the bundle this process materialized; the host's
1422 /// self-reported digest must match (a consistency check, not
1423 /// anti-substitution: the control is that Rust chooses what to exec).
1424 pub(crate) async fn spawn(
1425 generation: u64,
1426 launch: &HostLaunch,
1427 expected_sha256: &str,
1428 events: Arc<dyn HostEvents>,
1429 ) -> Result<Arc<Self>, String> {
1430 let mut command = tokio::process::Command::new(&launch.program);
1431 crate::utils::suppress_tokio_console_window(&mut command);
1432 command
1433 .args(&launch.args)
1434 .current_dir(&launch.cwd)
1435 .stdin(Stdio::piped())
1436 .stdout(Stdio::piped())
1437 .stderr(Stdio::piped())
1438 .kill_on_drop(true);
1439 // Scrubbed environment: no credentials, no ambient proxy URLs.
1440 command.env_clear();
1441 let parent_pid = std::process::id().to_string();
1442 // On Unix the host leads its own process group (below), so it may kill
1443 // that group when the core goes away (stdin EOF, or a parent change
1444 // seen by its watchdog thread).
1445 let own_group = if cfg!(unix) { "1" } else { "0" };
1446 let memory_request = (launch.memory == MemoryEnforcement::Jetsam)
1447 .then(|| (launch.memory_cap / (1024 * 1024)).max(1).to_string());
1448 let overrides = launch
1449 .sandbox_env
1450 .iter()
1451 .chain(&launch.runtime_env)
1452 .map(|(key, value)| (key.as_str(), value.as_str()))
1453 .chain(
1454 memory_request
1455 .as_deref()
1456 .map(|mib| (MEMORY_LIMIT_REQUEST_ENV, mib)),
1457 )
1458 .chain([
1459 ("CODEWHALE_HOST_PARENT_PID", parent_pid.as_str()),
1460 ("CODEWHALE_HOST_PROCESS_GROUP", own_group),
1461 ]);
1462 let environment =
1463 crate::child_env::sanitized_plugin_mcp_env_from(std::env::vars_os(), overrides);
1464 command.envs(environment.iter().map(|(key, value)| (key, value)));
1465 #[cfg(unix)]
1466 command.process_group(0);
1467 limit_child_memory(&mut command, launch.memory_cap);
1468
1469 let SpawnedHost {
1470 mut child,
1471 pid,
1472 tree,
1473 stdin,
1474 stdout,
1475 stderr,
1476 } = spawn_host(launch, command, &environment)?;
1477 let memory: Arc<std::sync::OnceLock<MemoryEnforcement>> = Arc::default();
1478
1479 let (outbound, mut outbound_rx) = mpsc::channel::<Vec<u8>>(OUTBOUND_QUEUE);
1480 let pending: Arc<Mutex<HashMap<u64, PendingCall>>> = Arc::default();
1481 let inbound: Arc<InboundRequests> = Arc::default();
1482 let stderr_tail: Arc<Mutex<VecDeque<u8>>> = Arc::default();
1483 let (exited_tx, exited_rx) = tokio::sync::watch::channel(false);
1484 let (hello_tx, hello_rx) = oneshot::channel();
1485 let (ready_tx, ready_rx) = oneshot::channel();
1486 let handshake = Arc::new(Mutex::new(Handshake {
1487 hello: Some(hello_tx),
1488 ready: Some(ready_tx),
1489 }));
1490
1491 // Writer: the only task that touches stdin. Dropping every sender
1492 // closes stdin, which the host treats as "core is gone".
1493 tokio::spawn(async move {
1494 let mut stdin = stdin;
1495 while let Some(frame) = outbound_rx.recv().await {
1496 if stdin.write_all(&frame).await.is_err() || stdin.flush().await.is_err() {
1497 break;
1498 }
1499 }
1500 });
1501
1502 // stderr: a bounded tail for diagnostics, chunks into tracing. Read
1503 // in fixed-size chunks, never by line: a plugin writing endless
1504 // output without a newline must not grow this process's memory.
1505 {
1506 let tail = Arc::clone(&stderr_tail);
1507 tokio::spawn(async move {
1508 let mut stderr = stderr;
1509 let mut chunk = vec![0_u8; 4096];
1510 while let Ok(read) = stderr.read(&mut chunk).await {
1511 if read == 0 {
1512 break;
1513 }
1514 push_tail(&tail, &chunk[..read]);
1515 tracing::debug!(
1516 target: "extension_host",
1517 "host stderr: {}",
1518 String::from_utf8_lossy(&chunk[..read]).trim_end()
1519 );
1520 }
1521 });
1522 }
1523
1524 // Reader: every frame is validated strictly; a framing or protocol
1525 // violation kills the host (a plugin wrote to the channel, or the
1526 // host is not ours).
1527 let (kill_tx, mut kill_rx) = mpsc::channel::<String>(1);
1528 {
1529 let reader = Reader {
1530 kill: kill_tx.clone(),
1531 pending: Arc::clone(&pending),
1532 inbound: Arc::clone(&inbound),
1533 outbound: outbound.clone(),
1534 events: Arc::clone(&events),
1535 handshake: Arc::clone(&handshake),
1536 };
1537 let tier = launch.tier;
1538 let kill_tx = kill_tx.clone();
1539 tokio::spawn(async move {
1540 let mut stdout = stdout;
1541 loop {
1542 let value = match protocol::read_frame(&mut stdout).await {
1543 Ok(Some(value)) => value,
1544 Ok(None) => break,
1545 Err(error) => {
1546 let _ = kill_tx.try_send(format!("channel framing violation: {error}"));
1547 break;
1548 }
1549 };
1550 let message = match protocol::parse_host_message(value, tier) {
1551 Ok(message) => message,
1552 Err(error) => {
1553 let _ = kill_tx.try_send(format!("protocol violation: {error}"));
1554 break;
1555 }
1556 };
1557 if let Err(violation) = reader.handle(message) {
1558 let _ = kill_tx.try_send(format!("protocol violation: {violation}"));
1559 break;
1560 }
1561 }
1562 });
1563 }
1564
1565 // Exit watcher: owns the child. On exit, fail everything and report.
1566 {
1567 let pending = Arc::clone(&pending);
1568 let inbound = Arc::clone(&inbound);
1569 let tail = Arc::clone(&stderr_tail);
1570 let events = Arc::clone(&events);
1571 let tree = Arc::clone(&tree);
1572 let memory = Arc::clone(&memory);
1573 let memory_cap = launch.memory_cap;
1574 tokio::spawn(async move {
1575 let reason = tokio::select! {
1576 biased;
1577 status = child.wait() => match status {
1578 Ok(status) => exit_reason(status, memory.get().copied(), memory_cap),
1579 Err(error) => format!("wait failed: {error}"),
1580 },
1581 Some(reason) = kill_rx.recv() => {
1582 let _ = tree.kill();
1583 let _ = child.kill().await;
1584 reason
1585 }
1586 };
1587 // The leader is gone; take anything it left behind with it.
1588 let _ = tree.kill();
1589 let drained: Vec<PendingCall> = {
1590 let mut pending = pending.lock().expect("pending lock");
1591 // Publish exit under the admission lock before draining:
1592 // waking a failed call must not admit another orphaned call.
1593 let _ = exited_tx.send(true);
1594 pending.drain().map(|(_, call)| call).collect()
1595 };
1596 for call in drained {
1597 let _ = call.tx.send(Err(HostCallError::Exited(reason.clone())));
1598 }
1599 // Requests the host sent have nobody left to answer.
1600 inbound.cancel_all();
1601 // Give the stderr task a moment to capture the last lines.
1602 tokio::time::sleep(Duration::from_millis(50)).await;
1603 events.exited(generation, reason, tail_string(&tail));
1604 });
1605 }
1606
1607 let host = Arc::new(Self {
1608 tier: launch.tier,
1609 generation,
1610 pid,
1611 runtime: launch.runtime.clone(),
1612 runtime_version: std::sync::OnceLock::new(),
1613 memory_cap: launch.memory_cap,
1614 memory,
1615 sandbox: launch.sandbox.clone(),
1616 tree,
1617 #[cfg(windows)]
1618 windows: launch.windows.clone(),
1619 outbound,
1620 pending,
1621 inbound,
1622 admission_closed: AtomicBool::new(false),
1623 next_id: AtomicU64::new(1),
1624 #[cfg(test)]
1625 heartbeats_sent: AtomicU64::new(0),
1626 stderr_tail,
1627 exited: exited_rx,
1628 kill: kill_tx.clone(),
1629 });
1630
1631 let handshake_result = tokio::time::timeout(HANDSHAKE_DEADLINE, async {
1632 let hello = hello_rx
1633 .await
1634 .map_err(|_| "host exited before host/hello".to_string())?;
1635 if hello.protocol.min > protocol::PROTOCOL_VERSION
1636 || hello.protocol.max < protocol::PROTOCOL_VERSION
1637 {
1638 return Err(format!(
1639 "host speaks protocol {}..={}, core speaks {}",
1640 hello.protocol.min,
1641 hello.protocol.max,
1642 protocol::PROTOCOL_VERSION
1643 ));
1644 }
1645 // Never a silent runtime switch: the host must be running on
1646 // the runtime this process chose and launched.
1647 if hello.runtime.name != launch.runtime.kind.name() {
1648 return Err(format!(
1649 "host reports runtime {} but {} was launched",
1650 hello.runtime.name,
1651 launch.runtime.kind.name()
1652 ));
1653 }
1654 check_hello_identity(&hello, launch.tier, launch.builtin_modules)?;
1655 // Restarts reuse the pinned runtime without probing it again, so
1656 // the binary at that path can have been replaced since (an
1657 // upgrade mid-session). Its flags and lockdown were chosen for
1658 // the probed version; refuse rather than run an unprobed one.
1659 if !launch.runtime.reports_version(&hello.runtime.version) {
1660 return Err(format!(
1661 "host reports {} {} but {} {} was probed at {}: the runtime binary changed mid-session; restart Codewhale to use the new version",
1662 hello.runtime.name,
1663 hello.runtime.version,
1664 launch.runtime.kind.name(),
1665 launch.runtime.version_string(),
1666 launch.runtime.path.display()
1667 ));
1668 }
1669 if hello.bundle_sha256 != expected_sha256 {
1670 return Err(format!(
1671 "host bundle digest {} does not match the materialized bundle {}",
1672 hello.bundle_sha256, expected_sha256
1673 ));
1674 }
1675 let requested = memory_request
1676 .as_deref()
1677 .and_then(|mib| mib.parse::<u64>().ok());
1678 let memory = match (hello.memory_limit_mib, requested) {
1679 (Some(applied), Some(requested)) if applied == requested => {
1680 MemoryEnforcement::Jetsam
1681 }
1682 (Some(applied), _) => {
1683 return Err(format!(
1684 "host reports a {applied} MiB memory limit the core did not ask for"
1685 ));
1686 }
1687 // The mandatory kernel cap cannot degrade to a delayed RSS
1688 // observation. Refuse before plugin initialization and keep
1689 // the host's stderr explanation in the existing diagnosis.
1690 (None, Some(requested)) => {
1691 return Err(format!(
1692 "host did not apply the requested {requested} MiB kernel memory limit; initialization refused"
1693 ));
1694 }
1695 (None, None) => launch.memory,
1696 };
1697 let _ = host.memory.set(memory);
1698 let initialize = CoreRequest::Initialize(InitializeParams {
1699 protocol: protocol::PROTOCOL_VERSION,
1700 limits: HostLimits {
1701 max_frame: protocol::MAX_FRAME as u64,
1702 max_inflight: protocol::MAX_INFLIGHT as u64,
1703 dispose_deadline_ms: DISPOSE_DEADLINE.as_millis() as u64,
1704 activate_deadline_ms: ACTIVATE_DEADLINE.as_millis() as u64,
1705 },
1706 });
1707 host.call(initialize, None)
1708 .await
1709 .map_err(|error| format!("host/initialize failed: {error}"))?;
1710 ready_rx
1711 .await
1712 .map_err(|_| "host exited before host/ready".to_string())?;
1713 Ok(hello.runtime.version)
1714 })
1715 .await;
1716 match handshake_result {
1717 Ok(Ok(runtime_version)) => {
1718 let _ = host.runtime_version.set(runtime_version);
1719 Ok(host)
1720 }
1721 Ok(Err(reason)) => {
1722 let _ = kill_tx.try_send(reason.clone());
1723 Err(format!(
1724 "{reason}; stderr: {}",
1725 tail_string(&host.stderr_tail)
1726 ))
1727 }
1728 Err(_) => {
1729 let reason = format!("handshake exceeded {HANDSHAKE_DEADLINE:?}");
1730 let _ = kill_tx.try_send(reason.clone());
1731 Err(format!(
1732 "{reason}; stderr: {}",
1733 tail_string(&host.stderr_tail)
1734 ))
1735 }
1736 }
1737 }
1738
1739 /// How the memory cap is enforced for this process (planned until the
1740 /// handshake settles it).
1741 #[must_use]
1742 pub fn memory(&self) -> MemoryEnforcement {
1743 self.memory
1744 .get()
1745 .copied()
1746 .unwrap_or_else(|| MemoryEnforcement::planned(self.runtime.kind))
1747 }
1748
1749 /// How many requests the core has sent this host (handshake included),
1750 /// excluding the monitor's heartbeat pings, which run on their own timer
1751 /// and would otherwise make a "no request yet" assertion timing-dependent.
1752 #[cfg(test)]
1753 pub(crate) fn requests_started(&self) -> u64 {
1754 self.next_id.load(Ordering::Relaxed) - 1 - self.heartbeats_sent.load(Ordering::Relaxed)
1755 }
1756
1757 #[must_use]
1758 pub fn has_exited(&self) -> bool {
1759 *self.exited.borrow()
1760 }
1761
1762 pub(crate) fn terminate(&self, reason: String) {
1763 let _ = self.kill.try_send(reason);
1764 }
1765
1766 pub(crate) fn is_retiring(&self) -> bool {
1767 self.admission_closed.load(Ordering::Acquire)
1768 }
1769
1770 /// Seal admission and retire this process only if no non-heartbeat call
1771 /// is pending. The same lock guards admission in `start_request`.
1772 pub(crate) fn terminate_if_idle(&self, reason: &str) -> bool {
1773 let pending = self.pending.lock().expect("pending lock");
1774 if self.has_exited()
1775 || self.admission_closed.load(Ordering::Relaxed)
1776 || pending.values().any(|call| !call.heartbeat)
1777 {
1778 return false;
1779 }
1780 if self.kill.try_send(reason.to_string()).is_err() {
1781 return false;
1782 }
1783 self.admission_closed.store(true, Ordering::Release);
1784 true
1785 }
1786
1787 pub(crate) fn stderr_tail(&self) -> String {
1788 tail_string(&self.stderr_tail)
1789 }
1790
1791 fn send_frame(&self, value: &Value) -> Result<(), HostCallError> {
1792 let frame = protocol::encode_frame(value).map_err(|error| HostCallError::Rpc {
1793 code: error_code::INVALID_PARAMS,
1794 message: error.to_string(),
1795 })?;
1796 self.outbound.try_send(frame).map_err(|error| match error {
1797 mpsc::error::TrySendError::Full(_) => HostCallError::Busy,
1798 mpsc::error::TrySendError::Closed(_) => {
1799 HostCallError::Exited("channel closed".to_string())
1800 }
1801 })
1802 }
1803
1804 /// Send a request; the returned id can be cancelled with [`Self::cancel`].
1805 pub(crate) fn start_request(
1806 &self,
1807 request: CoreRequest,
1808 owner: Option<String>,
1809 ) -> Result<(u64, CallReceiver), HostCallError> {
1810 // A method reserved for another tier is never sent to this host.
1811 if !protocol::allowed_on(Direction::CoreToHost, request.method(), self.tier) {
1812 return Err(HostCallError::Rpc {
1813 code: error_code::METHOD_NOT_FOUND,
1814 message: format!(
1815 "`{}` is not allowed on the {} tier",
1816 request.method(),
1817 self.tier.name()
1818 ),
1819 });
1820 }
1821 if self.has_exited() {
1822 return Err(HostCallError::Exited("already exited".to_string()));
1823 }
1824 let id = self.next_id.fetch_add(1, Ordering::Relaxed);
1825 #[cfg(test)]
1826 if matches!(request, CoreRequest::Ping) {
1827 self.heartbeats_sent.fetch_add(1, Ordering::Relaxed);
1828 }
1829 let (tx, rx) = oneshot::channel();
1830 {
1831 let mut pending = self.pending.lock().expect("pending lock");
1832 // Exit may have won after the fast check above. It publishes under
1833 // this same lock, so every admitted call is either live or drained.
1834 if self.has_exited() {
1835 return Err(HostCallError::Exited("already exited".to_string()));
1836 }
1837 if self.admission_closed.load(Ordering::Relaxed) {
1838 return Err(HostCallError::Exited("host is restarting".to_string()));
1839 }
1840 // Reserve one control request for the single heartbeat monitor,
1841 // so saturated tool calls cannot make a healthy host look hung.
1842 if pending.len() >= protocol::MAX_INFLIGHT && !matches!(request, CoreRequest::Ping) {
1843 return Err(HostCallError::Busy);
1844 }
1845 pending.insert(
1846 id,
1847 PendingCall {
1848 tx,
1849 owner,
1850 handle: match &request {
1851 CoreRequest::ToolCall(params) => Some(params.handle),
1852 CoreRequest::CommandRun(params) => Some(params.handle),
1853 CoreRequest::HookEvaluate(params) => Some(params.handle),
1854 _ => None,
1855 },
1856 revoked: false,
1857 heartbeat: matches!(request, CoreRequest::Ping),
1858 },
1859 );
1860 }
1861 if let Err(error) = self.send_frame(&request.to_value(id)) {
1862 self.pending.lock().expect("pending lock").remove(&id);
1863 return Err(error);
1864 }
1865 Ok((id, rx))
1866 }
1867
1868 /// Send `request` and wait at most its method's deadline
1869 /// ([`CoreRequest::deadline`]) for the answer. On expiry the host is sent
1870 /// `$/cancel`, the call is forgotten (a late answer is dropped) and the
1871 /// caller gets [`HostCallError::Timeout`]. Dropping the returned future
1872 /// (a turn interrupt) cancels the same way. Every awaited core→host
1873 /// request goes through here; only the heartbeat drives
1874 /// [`Self::start_request`] itself, under its own timeouts.
1875 pub(crate) async fn call(
1876 &self,
1877 request: CoreRequest,
1878 owner: Option<String>,
1879 ) -> Result<Value, HostCallError> {
1880 self.call_with_clock(request, owner, None).await
1881 }
1882
1883 /// [`Self::call`] with the deadline measured on `clock`, which stops while
1884 /// the caller waits on something that is not the host's time (a
1885 /// `core/call` waiting on an approval card). `None` is a clock nothing
1886 /// pauses: the plain deadline.
1887 pub(crate) async fn call_with_clock(
1888 &self,
1889 request: CoreRequest,
1890 owner: Option<String>,
1891 clock: Option<Arc<Mutex<PauseClock>>>,
1892 ) -> Result<Value, HostCallError> {
1893 let method = request.method();
1894 let deadline = request.deadline();
1895 let clock = clock.unwrap_or_else(|| Arc::new(Mutex::new(PauseClock::new())));
1896 let (id, mut rx) = self.start_request(request, owner)?;
1897 let mut guard = CancelOnDrop {
1898 host: self,
1899 id,
1900 armed: true,
1901 };
1902 let answer = loop {
1903 let remaining = clock
1904 .lock()
1905 .expect("deadline clock lock")
1906 .remaining(deadline);
1907 let Some(remaining) = remaining else {
1908 // `guard` is still armed: dropping it sends `$/cancel`.
1909 return Err(HostCallError::Timeout {
1910 method,
1911 after: deadline,
1912 });
1913 };
1914 tokio::select! {
1915 answer = &mut rx => break answer,
1916 () = tokio::time::sleep(remaining) => {}
1917 }
1918 };
1919 // Answered, drained at exit, or resolved by revocation: nothing left
1920 // to cancel.
1921 guard.armed = false;
1922 answer.unwrap_or_else(|_| Err(HostCallError::Exited("channel closed".to_string())))
1923 }
1924
1925 /// Fire `$/cancel`. Best effort: a full or closed channel is fine, the
1926 /// caller resolves its side on its own schedule.
1927 pub(crate) fn cancel(&self, id: u64) {
1928 let _ = self.send_frame(&protocol::cancel_value(id));
1929 }
1930
1931 /// Forget a request without cancelling it (its answer will be dropped).
1932 pub(crate) fn forget(&self, id: u64) {
1933 self.pending.lock().expect("pending lock").remove(&id);
1934 }
1935
1936 /// Revocation: cancel every in-flight call owned by `plugin_id`; each
1937 /// resolves as cancelled when the host answers or after `CANCEL_GRACE`,
1938 /// whichever is first — revocation never waits on the host.
1939 pub(crate) fn revoke_calls_of(self: &Arc<Self>, plugin_id: &str) {
1940 // Requests the host sent for this owner are cancelled too (and
1941 // answered `Cancelled`), whatever they are waiting for.
1942 self.inbound.cancel_owner(plugin_id);
1943 self.revoke_pending_where(|call| call.owner.as_deref() == Some(plugin_id));
1944 }
1945
1946 /// Retiring one Native entry cancels exactly its admitted contribution
1947 /// calls. Sibling entry requests and owner-wide broker lifecycle survive.
1948 pub(crate) fn revoke_calls_for_handles(self: &Arc<Self>, plugin_id: &str, handles: &[u64]) {
1949 self.revoke_pending_where(|call| {
1950 call.owner.as_deref() == Some(plugin_id)
1951 && call.handle.is_some_and(|handle| handles.contains(&handle))
1952 });
1953 }
1954
1955 fn revoke_pending_where(self: &Arc<Self>, drop_it: impl Fn(&PendingCall) -> bool) {
1956 let ids: Vec<u64> = {
1957 let mut pending = self.pending.lock().expect("pending lock");
1958 pending
1959 .iter_mut()
1960 .filter(|(_, call)| drop_it(call))
1961 .map(|(id, call)| {
1962 call.revoked = true;
1963 *id
1964 })
1965 .collect()
1966 };
1967 for id in ids {
1968 self.cancel(id);
1969 let pending = Arc::clone(&self.pending);
1970 tokio::spawn(async move {
1971 tokio::time::sleep(CANCEL_GRACE).await;
1972 if let Some(call) = pending.lock().expect("pending lock").remove(&id) {
1973 let _ = call.tx.send(Err(HostCallError::Cancelled(
1974 "extension was revoked".to_string(),
1975 )));
1976 }
1977 });
1978 }
1979 }
1980
1981 /// Bounded shutdown: `host/shutdown` (2 s), close stdin, then kill the
1982 /// process group at 3 s total.
1983 #[cfg(test)]
1984 pub(crate) async fn shutdown(&self) {
1985 let _ = self.call(CoreRequest::Shutdown, None).await;
1986 let mut exited = self.exited.clone();
1987 let waited = tokio::time::timeout(Duration::from_secs(1), async {
1988 while !*exited.borrow() {
1989 if exited.changed().await.is_err() {
1990 break;
1991 }
1992 }
1993 })
1994 .await;
1995 if waited.is_err() {
1996 let _ = self.tree.kill();
1997 }
1998 }
1999 }
2000
2001 /// Cancels a [`HostProcess::call`] that stops waiting before its answer:
2002 /// deadline expiry, or the caller's future being dropped.
2003 struct CancelOnDrop<'a> {
2004 host: &'a HostProcess,
2005 id: u64,
2006 armed: bool,
2007 }
2008
2009 impl Drop for CancelOnDrop<'_> {
2010 fn drop(&mut self) {
2011 if self.armed {
2012 self.host.cancel(self.id);
2013 self.host.forget(self.id);
2014 }
2015 }
2016 }
2017
2018 impl Drop for HostProcess {
2019 fn drop(&mut self) {
2020 // The exit watcher also holds the tree; kill explicitly so dropping
2021 // the last handle to a live host never leaves it running.
2022 if !self.has_exited() {
2023 let _ = self.tree.kill();
2024 }
2025 }
2026 }
2027
2028 /// What the channel's reader task holds: everything one decoded host→core
2029 /// message may touch.
2030 struct Reader {
2031 kill: mpsc::Sender<String>,
2032 pending: Arc<Mutex<HashMap<u64, PendingCall>>>,
2033 inbound: Arc<InboundRequests>,
2034 outbound: mpsc::Sender<Vec<u8>>,
2035 events: Arc<dyn HostEvents>,
2036 handshake: Arc<Mutex<Handshake>>,
2037 }
2038
2039 impl Reader {
2040 /// Route one validated host→core message. `Err` is a protocol violation
2041 /// the caller ends the host for.
2042 fn handle(&self, message: HostMessage) -> Result<(), String> {
2043 match message {
2044 HostMessage::Response { id, outcome } => {
2045 let Some(call) = self.pending.lock().expect("pending lock").remove(&id) else {
2046 tracing::debug!(target: "extension_host", id, "dropping late host response");
2047 return Ok(());
2048 };
2049 if !call.revoked && !call.heartbeat {
2050 self.events.responded();
2051 }
2052 let result = if call.revoked {
2053 Err(HostCallError::Cancelled(
2054 "extension was revoked".to_string(),
2055 ))
2056 } else {
2057 match outcome {
2058 Ok(value) => Ok(value),
2059 Err(error) if error.code == error_code::CANCELLED => {
2060 Err(HostCallError::Cancelled(error.message))
2061 }
2062 Err(error) => Err(HostCallError::Rpc {
2063 code: error.code,
2064 message: error.message,
2065 }),
2066 }
2067 };
2068 let _ = call.tx.send(result);
2069 }
2070 HostMessage::Request { id, request } => self.start_host_request(id, request)?,
2071 HostMessage::Notification(notification) => match notification {
2072 HostNotification::Hello(hello) => {
2073 if let Some(tx) = self.handshake.lock().expect("handshake lock").hello.take() {
2074 let _ = tx.send(hello);
2075 }
2076 }
2077 HostNotification::Ready => {
2078 if let Some(tx) = self.handshake.lock().expect("handshake lock").ready.take() {
2079 let _ = tx.send(());
2080 }
2081 }
2082 HostNotification::Faulted(params) => self.events.faulted(&params),
2083 HostNotification::Log(log) => {
2084 self.events.log(&log);
2085 let plugin = log.plugin_id.as_deref().unwrap_or("host");
2086 match log.level.as_str() {
2087 "error" => tracing::warn!(target: "extension_host", plugin, "{}", log.msg),
2088 "warn" => tracing::info!(target: "extension_host", plugin, "{}", log.msg),
2089 _ => tracing::debug!(target: "extension_host", plugin, "{}", log.msg),
2090 }
2091 }
2092 // The host withdrawing a request of its own. One that has
2093 // already been answered is not in flight: the cancel lost the
2094 // race, and nothing is owed.
2095 HostNotification::Cancel(params) => self.inbound.cancel_by_host(params.id),
2096 },
2097 }
2098 Ok(())
2099 }
2100
2101 /// Run one host-originated request as its own task: tracked by id so the
2102 /// host's `$/cancel`, the owner's revocation and the host's exit can
2103 /// cancel it, and answered when the handler finishes unless it was
2104 /// cancelled for a reason that makes the answer moot ([`CancelReason`]).
2105 /// A handler that does not stop within [`CANCEL_GRACE`] of its cancel is
2106 /// abandoned.
2107 fn start_host_request(&self, id: u64, request: HostRequest) -> Result<(), String> {
2108 let cancel = self.inbound.admit(id, request.plugin_id())?;
2109 let events = Arc::clone(&self.events);
2110 let inbound = Arc::clone(&self.inbound);
2111 let outbound = self.outbound.clone();
2112 let kill = self.kill.clone();
2113 tokio::spawn(async move {
2114 let cx = HostRequestContext {
2115 id,
2116 cancel: cancel.clone(),
2117 kill,
2118 };
2119 let mut handler = std::pin::pin!(events.host_request(request, cx));
2120 let outcome = tokio::select! {
2121 outcome = &mut handler => Some(outcome),
2122 () = async {
2123 cancel.cancelled().await;
2124 tokio::time::sleep(CANCEL_GRACE).await;
2125 } => None,
2126 };
2127 let outcome = match (inbound.finish(id), outcome) {
2128 // The host stopped waiting, or is gone: a late answer is dropped.
2129 (Some(CancelReason::Host | CancelReason::Exit), _) => return,
2130 (Some(CancelReason::Revoked), _) | (None, None) => Err(RpcErrorWire {
2131 code: error_code::CANCELLED,
2132 message: "extension was revoked".to_string(),
2133 data: None,
2134 }),
2135 (None, Some(outcome)) => outcome,
2136 };
2137 if let Ok(frame) = protocol::encode_frame(&protocol::response_value(id, &outcome)) {
2138 let _ = outbound.send(frame).await;
2139 }
2140 });
2141 Ok(())
2142 }
2143 }
2144
2145 /// What `host/hello` must say about the host's tier and built-in modules,
2146 /// checked the way its runtime name and version are: against what the core
2147 /// launched. The tier must be the one in the launch plan (`--tier=`), and the
2148 /// module digests the bundle embeds must be exactly the ones `modules` pins, so
2149 /// a host build and a Rust table that disagree about a built-in module are
2150 /// refused before any module can run. Pure.
2151 pub(crate) fn check_hello_identity(
2152 hello: &HelloParams,
2153 tier: HostTier,
2154 modules: &[BuiltinModule],
2155 ) -> Result<(), String> {
2156 if hello.tier != tier {
2157 return Err(format!(
2158 "host reports the {} tier but the {} tier was launched",
2159 hello.tier.name(),
2160 tier.name()
2161 ));
2162 }
2163 let describe = |rows: &[(String, String)]| {
2164 rows.iter()
2165 .map(|(id, digest)| format!("{id}={}", digest.chars().take(12).collect::<String>()))
2166 .collect::<Vec<_>>()
2167 .join(", ")
2168 };
2169 let mut reported: Vec<(String, String)> = hello
2170 .builtin_modules
2171 .iter()
2172 .map(|module| (module.id.clone(), module.sha256.clone()))
2173 .collect();
2174 let mut pinned: Vec<(String, String)> = modules
2175 .iter()
2176 .map(|module| (module.id.to_string(), module.source_sha256.to_string()))
2177 .collect();
2178 reported.sort();
2179 pinned.sort();
2180 if reported != pinned {
2181 return Err(format!(
2182 "host bundle embeds the built-in module digests [{}] but the core pins [{}]",
2183 describe(&reported),
2184 describe(&pinned)
2185 ));
2186 }
2187 Ok(())
2188 }
2189
2190 /// Where the embedded bundle is written: `<root>/extension-host/<sha256>/`.
2191 #[must_use]
2192 pub fn bundle_dir(root: &Path, sha256: &str) -> PathBuf {
2193 root.join("extension-host").join(sha256)
2194 }
2195
2196 #[cfg(test)]
2197 mod tests {
2198 use super::*;
2199
2200 #[test]
2201 fn compiled_host_has_direct_argv_and_cannot_become_a_bun_cli() {
2202 let runtime = HostRuntime {
2203 kind: HostRuntimeKind::Bun,
2204 path: PathBuf::from("/opt/codewhale-extension-host"),
2205 version: (1, 4, 0),
2206 native_code_flags: Vec::new(),
2207 compiled: true,
2208 };
2209 assert!(runtime_args(&runtime).is_empty());
2210 let temp = tempfile::tempdir().unwrap();
2211 let launch = host_launch(
2212 HostTier::Builtin,
2213 &runtime,
2214 vec![HostTier::Builtin.argv_flag()],
2215 temp.path().to_path_buf(),
2216 HOST_MEMORY_CAP,
2217 Err("fixture unsupported platform".to_string()),
2218 )
2219 .unwrap();
2220 assert_eq!(launch.program, runtime.path);
2221 assert_eq!(launch.args, vec!["--tier=builtin".to_string()]);
2222 assert!(
2223 launch
2224 .runtime_env
2225 .contains(&("BUN_OPTIONS".to_string(), String::new()))
2226 );
2227 assert!(
2228 launch
2229 .runtime_env
2230 .contains(&("BUN_BE_BUN".to_string(), "0".to_string()))
2231 );
2232 assert_eq!(
2233 launch.memory,
2234 MemoryEnforcement::planned(HostRuntimeKind::Bun)
2235 );
2236 }
2237
2238 fn node() -> HostRuntime {
2239 HostRuntime {
2240 kind: HostRuntimeKind::Node,
2241 path: PathBuf::from("/opt/node/bin/node"),
2242 version: (22, 19, 0),
2243 native_code_flags: Vec::new(),
2244 compiled: false,
2245 }
2246 }
2247
2248 /// A failed exact sandbox probe refuses Native code on every platform.
2249 /// The pinned Builtin exception keeps the concrete diagnostic; a verified
2250 /// wrapper keeps its command and never silently falls through to raw JS.
2251 #[test]
2252 fn missing_verified_sandbox_refuses_native_but_reports_pinned_builtin_exception() {
2253 let runtime = node();
2254 let data = PathBuf::from("/home/u/.codewhale/extension-host/data");
2255 let mut args = runtime_args(&runtime);
2256 args.push("/home/u/.codewhale/extension-host/abc/host.mjs".to_string());
2257
2258 let refused = bwrap_probe_verdict(
2259 false,
2260 "exit status: 1",
2261 "\nbwrap: setting up uid map: Permission denied\n",
2262 )
2263 .unwrap_err();
2264 assert!(
2265 refused.starts_with("bwrap unavailable (bwrap: setting up uid map: Permission denied;"),
2266 "{refused}"
2267 );
2268 assert!(
2269 refused.contains("kernel.apparmor_restrict_unprivileged_userns"),
2270 "{refused}"
2271 );
2272 let error = host_launch(
2273 HostTier::Plugin,
2274 &runtime,
2275 args.clone(),
2276 data.clone(),
2277 HOST_MEMORY_CAP,
2278 Err(refused.clone()),
2279 )
2280 .unwrap_err();
2281 assert!(error.contains("Native extensions require a verified OS sandbox"));
2282 assert!(error.contains(&refused));
2283 let launch = host_launch(
2284 HostTier::Builtin,
2285 &runtime,
2286 args.clone(),
2287 data.clone(),
2288 HOST_MEMORY_CAP,
2289 Err(refused.clone()),
2290 )
2291 .unwrap();
2292 assert_eq!(launch.program, runtime.path);
2293 assert_eq!(launch.args, args);
2294 assert!(launch.sandbox_env.is_empty());
2295 assert_eq!(launch.sandbox, HostSandbox::Unsandboxed(refused.clone()));
2296 assert_eq!(launch.sandbox.label(), format!("none: {refused}"));
2297 assert!(
2298 launch
2299 .sandbox
2300 .to_string()
2301 .starts_with(&format!("UNSANDBOXED ({refused}): ")),
2302 "{}",
2303 launch.sandbox
2304 );
2305
2306 // No stderr: the exit status is the reason, with no namespace hint.
2307 assert_eq!(
2308 bwrap_probe_verdict(false, "signal: 9 (SIGKILL)", ""),
2309 Err("bwrap unavailable (signal: 9 (SIGKILL))".to_string())
2310 );
2311
2312 // A probe that ran keeps the wrapper.
2313 assert_eq!(bwrap_probe_verdict(true, "exit status: 0", ""), Ok(()));
2314 let command: Vec<String> = ["/usr/bin/bwrap", "--unshare-all", "--die-with-parent", "--"]
2315 .iter()
2316 .map(ToString::to_string)
2317 .chain(std::iter::once(runtime.path.to_string_lossy().into_owned()))
2318 .chain(args.iter().cloned())
2319 .collect();
2320 let launch = host_launch(
2321 HostTier::Plugin,
2322 &runtime,
2323 args.clone(),
2324 data.clone(),
2325 HOST_MEMORY_CAP,
2326 Ok(Wrapped {
2327 name: "linux-bwrap".to_string(),
2328 command: command.clone(),
2329 env: vec![("CODEWHALE_SANDBOX".to_string(), "bwrap".to_string())],
2330 }),
2331 )
2332 .unwrap();
2333 assert_eq!(launch.program, PathBuf::from("/usr/bin/bwrap"));
2334 assert_eq!(launch.args, command[1..].to_vec());
2335 assert_eq!(launch.cwd, data);
2336 assert_eq!(
2337 launch.sandbox,
2338 HostSandbox::Wrapped("linux-bwrap".to_string())
2339 );
2340 assert!(
2341 launch
2342 .sandbox
2343 .to_string()
2344 .starts_with("linux-bwrap sandbox (no direct network;"),
2345 "{}",
2346 launch.sandbox
2347 );
2348 }
2349
2350 /// Each tier has its own data directory, and neither is inside the other:
2351 /// the data directory is the host's only writable root, so a nested
2352 /// builtin directory would be writable by plugin code. The plugin tier
2353 /// keeps the path it has always had, and the plugin tier's sandbox denies
2354 /// reads of the builtin tier's directory.
2355 #[test]
2356 fn each_tier_has_its_own_data_directory_and_the_plugin_tier_cannot_read_the_builtin_one() {
2357 let temp = tempfile::tempdir().unwrap();
2358 let home = temp.path().join("home");
2359 std::fs::create_dir_all(tier_data_dir(&home, HostTier::Builtin)).unwrap();
2360 let plugin = tier_data_dir(&home, HostTier::Plugin);
2361 let builtin = tier_data_dir(&home, HostTier::Builtin);
2362 assert_ne!(plugin, builtin);
2363 assert!(!builtin.starts_with(&plugin) && !plugin.starts_with(&builtin));
2364 assert_eq!(plugin, home.join("extension-host").join("data"));
2365
2366 for whole_homes in [false, true] {
2367 let (denied, _) = host_denied_read_paths(HostTier::Plugin, &home, whole_homes);
2368 assert!(
2369 denied.contains(&builtin),
2370 "plugin tier (whole_homes {whole_homes}) must deny {}",
2371 builtin.display()
2372 );
2373 assert!(!denied.contains(&plugin));
2374 let (denied, _) = host_denied_read_paths(HostTier::Builtin, &home, whole_homes);
2375 assert!(
2376 !denied.contains(&builtin),
2377 "the builtin tier reads its own directory"
2378 );
2379 }
2380 }
2381
2382 /// The per-plugin directory plugin config and context introduced keeps its
2383 /// path, so an installed plugin's data survives the tier split; a built-in
2384 /// module's is under the builtin tier's directory.
2385 #[test]
2386 fn the_per_plugin_data_directory_keeps_its_path() {
2387 let home = PathBuf::from("/home/u/.codewhale");
2388 let id = "user/0123456789ab/demo";
2389 assert_eq!(
2390 plugin_data_dir(&home, id, "demo"),
2391 home.join("extension-host/data/plugins/demo-e9f63c6f1a7f")
2392 );
2393 assert_eq!(
2394 owner_data_dir(&home, HostTier::Plugin, id, "demo"),
2395 plugin_data_dir(&home, id, "demo")
2396 );
2397 assert_eq!(
2398 owner_data_dir(&home, HostTier::Builtin, "host:mcp", "mcp"),
2399 home.join("extension-host/data-builtin/modules/mcp")
2400 );
2401 }
2402
2403 /// bubblewrap can mask only what exists, so under it each Codewhale home
2404 /// is denied whole (an entry created later is then denied too) and its
2405 /// readable entries come back as exceptions; Seatbelt's form is unchanged.
2406 #[test]
2407 fn under_bubblewrap_a_codewhale_home_is_denied_whole_but_its_readable_entries() {
2408 let temp = tempfile::tempdir().unwrap();
2409 let home = temp.path().join("home");
2410 std::fs::create_dir_all(home.join("extension-host")).unwrap();
2411 std::fs::write(home.join("config.toml.bak-1"), "").unwrap();
2412
2413 let (seatbelt, none) = host_denied_read_paths(HostTier::Plugin, &home, false);
2414 assert!(none.is_empty());
2415 assert!(seatbelt.contains(&home.join("config.toml.bak-1")));
2416 assert!(
2417 seatbelt.contains(&home.join("secrets")),
2418 "named before it exists"
2419 );
2420 assert!(!seatbelt.contains(&home));
2421 assert!(!seatbelt.contains(&home.join("extension-host")));
2422
2423 let (bwrap, exceptions) = host_denied_read_paths(HostTier::Plugin, &home, true);
2424 assert!(bwrap.contains(&home));
2425 assert!(!bwrap.contains(&home.join("config.toml.bak-1")));
2426 for entry in HOST_READABLE_HOME_ENTRIES {
2427 assert!(exceptions.contains(&home.join(entry)), "{entry}");
2428 assert!(!bwrap.contains(&home.join(entry)), "{entry}");
2429 }
2430 }
2431
2432 // -----------------------------------------------------------------------
2433 // The host's tier and built-in module digests in `host/hello`
2434 // -----------------------------------------------------------------------
2435
2436 const DEMO_DIGEST: &str = "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef";
2437 const OTHER_DIGEST: &str = "fedcba9876543210fedcba9876543210fedcba9876543210fedcba9876543210";
2438 const PINNED: &[BuiltinModule] = &[
2439 BuiltinModule {
2440 id: "demo",
2441 source_sha256: DEMO_DIGEST,
2442 tools: &[],
2443 },
2444 BuiltinModule {
2445 id: "other",
2446 source_sha256: OTHER_DIGEST,
2447 tools: &[],
2448 },
2449 ];
2450
2451 fn hello(tier: HostTier, modules: &[(&str, &str)]) -> HelloParams {
2452 HelloParams {
2453 protocol: protocol::ProtocolRange { min: 1, max: 1 },
2454 host_version: "0.1.0".to_string(),
2455 bundle_sha256: "0".repeat(64),
2456 runtime: protocol::HelloRuntime {
2457 name: "node".to_string(),
2458 version: "22.20.0".to_string(),
2459 },
2460 tier,
2461 builtin_modules: modules
2462 .iter()
2463 .map(|(id, sha256)| protocol::ModuleDigestWire {
2464 id: (*id).to_string(),
2465 sha256: (*sha256).to_string(),
2466 })
2467 .collect(),
2468 memory_limit_mib: None,
2469 }
2470 }
2471
2472 #[test]
2473 fn hello_must_report_the_launched_tier_and_exactly_the_pinned_module_digests() {
2474 // Production: no module, so none may be reported, on either tier.
2475 for tier in HostTier::ALL {
2476 assert_eq!(check_hello_identity(&hello(tier, &[]), tier, &[]), Ok(()));
2477 }
2478 // Order is not part of the claim.
2479 let reported = [("other", OTHER_DIGEST), ("demo", DEMO_DIGEST)];
2480 assert_eq!(
2481 check_hello_identity(
2482 &hello(HostTier::Builtin, &reported),
2483 HostTier::Builtin,
2484 PINNED
2485 ),
2486 Ok(())
2487 );
2488
2489 let refused = |hello: HelloParams, tier: HostTier| {
2490 check_hello_identity(&hello, tier, PINNED).unwrap_err()
2491 };
2492 assert_eq!(
2493 refused(hello(HostTier::Plugin, &reported), HostTier::Builtin),
2494 "host reports the plugin tier but the builtin tier was launched"
2495 );
2496 assert_eq!(
2497 refused(hello(HostTier::Builtin, &reported), HostTier::Plugin),
2498 "host reports the builtin tier but the plugin tier was launched"
2499 );
2500 let pinned = "[demo=0123456789ab, other=fedcba987654]";
2501 for (what, rows, shown) in [
2502 ("none reported", vec![], "[]"),
2503 (
2504 "one missing",
2505 vec![("demo", DEMO_DIGEST)],
2506 "[demo=0123456789ab]",
2507 ),
2508 (
2509 "an unpinned module",
2510 vec![
2511 ("demo", DEMO_DIGEST),
2512 ("other", OTHER_DIGEST),
2513 ("extra", DEMO_DIGEST),
2514 ],
2515 "[demo=0123456789ab, extra=0123456789ab, other=fedcba987654]",
2516 ),
2517 (
2518 "a changed digest",
2519 vec![("demo", OTHER_DIGEST), ("other", OTHER_DIGEST)],
2520 "[demo=fedcba987654, other=fedcba987654]",
2521 ),
2522 (
2523 "a row twice",
2524 vec![
2525 ("demo", DEMO_DIGEST),
2526 ("demo", DEMO_DIGEST),
2527 ("other", OTHER_DIGEST),
2528 ],
2529 "[demo=0123456789ab, demo=0123456789ab, other=fedcba987654]",
2530 ),
2531 ] {
2532 assert_eq!(
2533 refused(hello(HostTier::Builtin, &rows), HostTier::Builtin),
2534 format!(
2535 "host bundle embeds the built-in module digests {shown} but the core pins {pinned}"
2536 ),
2537 "{what}"
2538 );
2539 }
2540 // A host that reports modules when the core pins none is refused too.
2541 assert!(
2542 check_hello_identity(
2543 &hello(HostTier::Plugin, &[("demo", DEMO_DIGEST)]),
2544 HostTier::Plugin,
2545 &[]
2546 )
2547 .is_err()
2548 );
2549 }
2550
2551 // -----------------------------------------------------------------------
2552 // Host-originated requests run as their own tracked tasks
2553 // -----------------------------------------------------------------------
2554
2555 use std::sync::atomic::AtomicUsize;
2556 use tokio::sync::Notify;
2557
2558 /// How a stub handler waits.
2559 #[derive(Clone, Copy)]
2560 enum Behaviour {
2561 /// Until `release`, or until its request is cancelled (then it stops).
2562 WaitsUnlessCancelled,
2563 /// Until `release`, whatever happens to its request.
2564 IgnoresCancel,
2565 /// Forever.
2566 Hangs,
2567 }
2568
2569 struct Stub {
2570 behaviour: Behaviour,
2571 release: Arc<Notify>,
2572 started: AtomicUsize,
2573 saw_cancel: AtomicBool,
2574 }
2575
2576 impl Stub {
2577 fn new(behaviour: Behaviour) -> Arc<Self> {
2578 Arc::new(Self {
2579 behaviour,
2580 release: Arc::new(Notify::new()),
2581 started: AtomicUsize::new(0),
2582 saw_cancel: AtomicBool::new(false),
2583 })
2584 }
2585 }
2586
2587 #[async_trait]
2588 impl HostEvents for Stub {
2589 fn register(&self, _: &protocol::RegisterParams) -> RegisterResult {
2590 unreachable!("the stub answers every request itself")
2591 }
2592 fn unregister(&self, _: &protocol::UnregisterParams) {}
2593 fn faulted(&self, _: &protocol::FaultedParams) {}
2594 fn log(&self, _: &protocol::LogParams) {}
2595 fn exited(&self, _: u64, _: String, _: String) {}
2596
2597 async fn host_request(
2598 &self,
2599 _request: HostRequest,
2600 cx: HostRequestContext,
2601 ) -> Result<Value, RpcErrorWire> {
2602 self.started.fetch_add(1, Ordering::SeqCst);
2603 match self.behaviour {
2604 Behaviour::WaitsUnlessCancelled => {
2605 tokio::select! {
2606 () = cx.cancel.cancelled() => {
2607 self.saw_cancel.store(true, Ordering::SeqCst);
2608 Err(RpcErrorWire {
2609 code: error_code::CANCELLED,
2610 message: "stub saw the cancel".to_string(),
2611 data: None,
2612 })
2613 }
2614 () = self.release.notified() => Ok(json!({"answered": cx.id})),
2615 }
2616 }
2617 Behaviour::IgnoresCancel => {
2618 self.release.notified().await;
2619 Ok(json!({"answered": cx.id}))
2620 }
2621 Behaviour::Hangs => std::future::pending().await,
2622 }
2623 }
2624 }
2625
2626 /// An events implementation that keeps the trait's own answers for the
2627 /// registry requests.
2628 struct Plain;
2629
2630 #[async_trait]
2631 impl HostEvents for Plain {
2632 fn register(&self, _: &protocol::RegisterParams) -> RegisterResult {
2633 RegisterResult::Admitted { handle: 9 }
2634 }
2635 fn unregister(&self, _: &protocol::UnregisterParams) {}
2636 fn faulted(&self, _: &protocol::FaultedParams) {}
2637 fn log(&self, _: &protocol::LogParams) {}
2638 fn exited(&self, _: u64, _: String, _: String) {}
2639 }
2640
2641 struct Rig {
2642 reader: Reader,
2643 frames: mpsc::Receiver<Vec<u8>>,
2644 }
2645
2646 fn new_rig(events: Arc<dyn HostEvents>) -> Rig {
2647 let (outbound, frames) = mpsc::channel(OUTBOUND_QUEUE);
2648 Rig {
2649 reader: Reader {
2650 kill: mpsc::channel(1).0,
2651 pending: Arc::default(),
2652 inbound: Arc::default(),
2653 outbound,
2654 events,
2655 handshake: Arc::default(),
2656 },
2657 frames,
2658 }
2659 }
2660
2661 fn owner(plugin: &str) -> protocol::OwnerRef {
2662 protocol::OwnerRef {
2663 plugin_id: plugin.to_string(),
2664 generation: 1,
2665 owner_token: "t".repeat(32),
2666 }
2667 }
2668
2669 fn register_request(plugin: &str) -> HostRequest {
2670 HostRequest::Register(protocol::RegisterParams {
2671 scope: None,
2672 owner: owner(plugin),
2673 kind: protocol::RegisterKind::Tool,
2674 spec: protocol::RegisterSpecWire {
2675 name: "t".to_string(),
2676 description: "d".to_string(),
2677 input_schema: Some(serde_json::Map::new()),
2678 argument_hint: None,
2679 },
2680 })
2681 }
2682
2683 fn request(rig: &Rig, id: u64, plugin: &str) -> Result<(), String> {
2684 rig.reader.handle(HostMessage::Request {
2685 id,
2686 request: register_request(plugin),
2687 })
2688 }
2689
2690 fn cancel(rig: &Rig, id: u64) {
2691 rig.reader
2692 .handle(HostMessage::Notification(HostNotification::Cancel(
2693 protocol::CancelParams { id },
2694 )))
2695 .unwrap();
2696 }
2697
2698 /// The next frame the core wrote to the host, decoded.
2699 async fn next_frame(rig: &mut Rig) -> Value {
2700 let frame = tokio::time::timeout(Duration::from_secs(5), rig.frames.recv())
2701 .await
2702 .expect("a frame within 5 s")
2703 .expect("the channel is open");
2704 serde_json::from_slice(&frame[protocol::HEADER_LEN..]).expect("a JSON frame")
2705 }
2706
2707 /// Nothing is written to the host for `ms`.
2708 async fn no_frame(rig: &mut Rig, ms: u64) {
2709 let frame = tokio::time::timeout(Duration::from_millis(ms), rig.frames.recv()).await;
2710 assert!(frame.is_err(), "unexpected frame: {frame:?}");
2711 }
2712
2713 async fn until(what: &str, mut done: impl FnMut() -> bool) {
2714 for _ in 0..200 {
2715 if done() {
2716 return;
2717 }
2718 tokio::time::sleep(Duration::from_millis(10)).await;
2719 }
2720 panic!("timed out waiting for {what}");
2721 }
2722
2723 #[tokio::test]
2724 async fn registry_requests_are_answered_through_the_generic_handler_as_before() {
2725 let mut rig = new_rig(Arc::new(Plain));
2726 request(&rig, 5, "p").unwrap();
2727 assert_eq!(
2728 next_frame(&mut rig).await,
2729 json!({"jsonrpc": "2.0", "id": 5, "result": {"handle": 9}})
2730 );
2731 rig.reader
2732 .handle(HostMessage::Request {
2733 id: 6,
2734 request: HostRequest::Unregister(protocol::UnregisterParams {
2735 owner: owner("p"),
2736 handle: 9,
2737 }),
2738 })
2739 .unwrap();
2740 assert_eq!(
2741 next_frame(&mut rig).await,
2742 json!({"jsonrpc": "2.0", "id": 6, "result": {}})
2743 );
2744 until("the table to empty", || rig.reader.inbound.in_flight() == 0).await;
2745 }
2746
2747 #[tokio::test]
2748 async fn a_host_cancel_cancels_the_task_and_the_answer_is_dropped() {
2749 // A handler that stops when cancelled answers nothing.
2750 let stub = Stub::new(Behaviour::WaitsUnlessCancelled);
2751 let mut rig = new_rig(stub.clone());
2752 request(&rig, 1, "p").unwrap();
2753 until("the handler to start", || {
2754 stub.started.load(Ordering::SeqCst) == 1
2755 })
2756 .await;
2757 assert_eq!(rig.reader.inbound.in_flight(), 1);
2758 cancel(&rig, 1);
2759 until("the handler to see the cancel", || {
2760 stub.saw_cancel.load(Ordering::SeqCst)
2761 })
2762 .await;
2763 no_frame(&mut rig, 150).await;
2764 until("the table to empty", || rig.reader.inbound.in_flight() == 0).await;
2765 // A cancel for an id that is not in flight (answered, or never sent) is ignored.
2766 cancel(&rig, 1);
2767 cancel(&rig, 4242);
2768
2769 // A handler that ignores the cancel and finishes later: its late answer is dropped.
2770 let stub = Stub::new(Behaviour::IgnoresCancel);
2771 let mut rig = new_rig(stub.clone());
2772 request(&rig, 2, "p").unwrap();
2773 until("the handler to start", || {
2774 stub.started.load(Ordering::SeqCst) == 1
2775 })
2776 .await;
2777 cancel(&rig, 2);
2778 stub.release.notify_one();
2779 no_frame(&mut rig, 150).await;
2780 until("the table to empty", || rig.reader.inbound.in_flight() == 0).await;
2781
2782 // One that never finishes is abandoned CANCEL_GRACE after the cancel,
2783 // and frees its slot.
2784 let mut rig = new_rig(Stub::new(Behaviour::Hangs));
2785 request(&rig, 3, "p").unwrap();
2786 cancel(&rig, 3);
2787 assert_eq!(
2788 rig.reader.inbound.in_flight(),
2789 1,
2790 "held until the grace ends"
2791 );
2792 until("the abandoned request to leave the table", || {
2793 rig.reader.inbound.in_flight() == 0
2794 })
2795 .await;
2796 no_frame(&mut rig, 100).await;
2797 }
2798
2799 #[tokio::test]
2800 async fn revoking_an_owner_answers_its_requests_cancelled_and_leaves_others_alone() {
2801 let stub = Stub::new(Behaviour::WaitsUnlessCancelled);
2802 let mut rig = new_rig(stub.clone());
2803 request(&rig, 1, "a").unwrap();
2804 request(&rig, 2, "b").unwrap();
2805 until("both handlers to start", || {
2806 stub.started.load(Ordering::SeqCst) == 2
2807 })
2808 .await;
2809 rig.reader.inbound.cancel_owner("a");
2810 // The revoked owner's host is told (its own `$/cancel` never came).
2811 let frame = next_frame(&mut rig).await;
2812 assert_eq!(frame["id"], 1);
2813 assert_eq!(frame["error"]["code"], error_code::CANCELLED);
2814 assert_eq!(frame["error"]["message"], "extension was revoked");
2815 // The other owner's request is untouched and answers when released.
2816 assert_eq!(rig.reader.inbound.in_flight(), 1);
2817 stub.release.notify_one();
2818 assert_eq!(
2819 next_frame(&mut rig).await,
2820 json!({"jsonrpc": "2.0", "id": 2, "result": {"answered": 2}})
2821 );
2822 }
2823
2824 #[tokio::test]
2825 async fn a_host_exit_cancels_every_request_answers_none_and_admits_no_more() {
2826 let stub = Stub::new(Behaviour::WaitsUnlessCancelled);
2827 let mut rig = new_rig(stub.clone());
2828 request(&rig, 1, "a").unwrap();
2829 request(&rig, 2, "b").unwrap();
2830 until("both handlers to start", || {
2831 stub.started.load(Ordering::SeqCst) == 2
2832 })
2833 .await;
2834 rig.reader.inbound.cancel_all();
2835 no_frame(&mut rig, 150).await;
2836 until("the table to empty", || rig.reader.inbound.in_flight() == 0).await;
2837 assert!(
2838 request(&rig, 3, "a")
2839 .unwrap_err()
2840 .contains("after the host exited")
2841 );
2842 }
2843
2844 #[tokio::test]
2845 async fn host_requests_are_capped_and_an_id_cannot_be_reused_in_flight() {
2846 let rig = new_rig(Stub::new(Behaviour::Hangs));
2847 for id in 1..=protocol::MAX_INFLIGHT as u64 {
2848 request(&rig, id, "p").unwrap_or_else(|reason| panic!("request {id}: {reason}"));
2849 }
2850 assert_eq!(rig.reader.inbound.in_flight(), protocol::MAX_INFLIGHT);
2851 let over = request(&rig, 9999, "p").unwrap_err();
2852 assert!(
2853 over.contains("more than 256 host requests in flight"),
2854 "{over}"
2855 );
2856 let reused = request(&rig, 1, "p").unwrap_err();
2857 assert!(reused.contains("already in flight"), "{reused}");
2858 // Both are protocol violations (the reader ends the host); neither
2859 // disturbed what was admitted.
2860 assert_eq!(rig.reader.inbound.in_flight(), protocol::MAX_INFLIGHT);
2861 rig.reader.inbound.cancel_all();
2862 }
2863 }
2864
2864 lines RUST