| 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(¶ms)) |
| 1067 | .unwrap_or_else(|_| json!({"refused": "internal"}))), |
| 1068 | HostRequest::Unregister(params) => { |
| 1069 | events.unregister(¶ms); |
| 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(¶ms), |
| 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 |