返回 CodeWhale
process_broker.rs
根目录 / crates / tui / src / mcp / process_broker.rs
1 //! Reviewed stdio process ownership, extracted from `mcp/stdio.rs`.
2 //!
3 //! Both the Rust MCP transport and the selected SDK broker use this lifetime
4 //! owner. This session owns spawn,
5 //! scrubbed environment, reviewed executable handles, stdin, bounded stderr,
6 //! authority cancellation, process-tree containment and graceful teardown.
7 //! `StdioTransport` owns only stdout framing and its cancellation-safe buffer.
8 //!
9 //! It accepts only the existing Rust MCP configuration. The SDK broker checks
10 //! Rust-issued launch and exact-operation authority before exposing its pipe
11 //! operations to the pinned builtin; plugins cannot use this process surface.
12 //! Protocol negotiation, catalog admission, tool approvals and call replay
13 //! policy remain in their owners.
14
15 use std::collections::VecDeque;
16 use std::sync::Arc;
17 use std::time::Duration;
18
19 use anyhow::{Context, Result};
20 use tokio::io::{AsyncBufReadExt, AsyncWriteExt};
21 use tokio::process::{Child, ChildStdin, ChildStdout};
22 use tokio::sync::Mutex as TokioMutex;
23
24 use super::McpServerConfig;
25 use crate::child_env;
26
27 pub(crate) struct BrokerSession {
28 child: Arc<TokioMutex<Child>>,
29 stdin: ChildStdin,
30 stderr_tail: Arc<StderrTail>,
31 /// Idle reviewed children die when their connection authority is revoked.
32 authority_cancel_watch: Option<tokio::task::JoinHandle<()>>,
33 /// Keep the reviewed executable/script descriptors alive through cleanup.
34 _reviewed_launch: Option<super::ReviewedStdioLaunch>,
35 /// Owned until Drop transfers it to the bounded asynchronous cleanup.
36 process_tree: Option<Arc<crate::process_tree::ProcessTree>>,
37 }
38
39 /// How long `BrokerSession::shutdown` waits for the child to exit on SIGTERM
40 /// before `kill_on_drop` fires SIGKILL. Tuned short so a hung MCP server
41 /// can't stall TUI exit; well-behaved servers almost always exit within
42 /// a few hundred ms.
43 pub(super) const STDIO_SHUTDOWN_GRACE: Duration = Duration::from_millis(2_000);
44
45 /// How many lines of MCP-server stderr to keep around for crash diagnostics.
46 /// Bounded so a chatty server can't grow this without limit; large enough to
47 /// catch typical Node/Python startup or panic output.
48 const STDERR_TAIL_CAPACITY: usize = 64;
49
50 /// Bounded ring buffer for the most recent stderr lines from a spawned MCP
51 /// server. Owned by `BrokerSession` to surface server-side context when the
52 /// transport read side fails (server crashed, exited early, etc).
53 #[derive(Default)]
54 struct StderrTail {
55 lines: TokioMutex<VecDeque<String>>,
56 }
57
58 impl StderrTail {
59 fn new() -> Arc<Self> {
60 Arc::new(Self {
61 lines: TokioMutex::new(VecDeque::with_capacity(STDERR_TAIL_CAPACITY)),
62 })
63 }
64
65 async fn push(&self, line: String) {
66 let mut buf = self.lines.lock().await;
67 if buf.len() >= STDERR_TAIL_CAPACITY {
68 buf.pop_front();
69 }
70 buf.push_back(line);
71 }
72
73 async fn snapshot(&self) -> Vec<String> {
74 self.lines.lock().await.iter().cloned().collect()
75 }
76
77 async fn last_line(&self) -> Option<String> {
78 self.lines
79 .lock()
80 .await
81 .iter()
82 .rev()
83 .map(|line| line.trim())
84 .find(|line| !line.is_empty())
85 .map(str::to_string)
86 }
87 }
88
89 impl BrokerSession {
90 pub(crate) fn spawn(
91 server_name: &str,
92 command: &str,
93 config: &McpServerConfig,
94 cancel_token: tokio_util::sync::CancellationToken,
95 ) -> Result<(Self, ChildStdout)> {
96 let reviewed_launch = if let Some(reviewed_plugin) = config.reviewed_plugin.as_ref() {
97 // This is deliberately the last trust check before constructing
98 // and spawning the lazy stdio child. It re-reads only the
99 // Codewhale-owned plugin bundle, never user MCP/provider config or
100 // credential files, and fails closed on any content/capability
101 // drift after pool construction.
102 Some(reviewed_plugin.prepare_stdio_launch(
103 server_name,
104 command,
105 &config.args,
106 config.cwd.as_deref(),
107 )?)
108 } else {
109 None
110 };
111 // Windows: Path::canonicalize yields a verbatim prefix. Node's ESM loader
112 // (and several other Windows tools) refuse that spelling for ordinary
113 // paths under MAX_PATH, so the reviewed peer exits 1 before handshake.
114 // Strip only at the child boundary; containment kept the prefixed form.
115 let mut cmd = reviewed_launch.as_ref().map_or_else(
116 || {
117 let mut command_process =
118 tokio::process::Command::new(super::strip_windows_verbatim_for_child(command));
119 for arg in &config.args {
120 command_process.arg(super::strip_windows_verbatim_for_child(arg.as_str()));
121 }
122 command_process
123 },
124 |launch| {
125 let mut command_process = tokio::process::Command::new(
126 super::strip_windows_verbatim_for_child(&launch.command),
127 );
128 for arg in &launch.args {
129 command_process.arg(super::strip_windows_verbatim_for_child(arg));
130 }
131 command_process
132 },
133 );
134 crate::utils::suppress_tokio_console_window(&mut cmd);
135 cmd.stdin(std::process::Stdio::piped())
136 .stdout(std::process::Stdio::piped())
137 .stderr(std::process::Stdio::piped())
138 .kill_on_drop(true);
139 let launch_cwd = reviewed_launch
140 .as_ref()
141 .and_then(|launch| launch.cwd.as_ref())
142 .or(config.cwd.as_ref().filter(|_| reviewed_launch.is_none()));
143 if let Some(cwd) = launch_cwd {
144 cmd.current_dir(super::strip_windows_verbatim_for_child(cwd.as_os_str()));
145 }
146 #[cfg(unix)]
147 if let Some(cwd_fd) = reviewed_launch
148 .as_ref()
149 .and_then(|launch| launch.cwd_fd.as_ref())
150 {
151 use std::os::fd::AsRawFd as _;
152 let fd = cwd_fd.as_raw_fd();
153 // SAFETY: the closure calls only async-signal-safe `fchdir` on an
154 // inherited directory descriptor before exec.
155 unsafe {
156 cmd.pre_exec(move || {
157 if libc::fchdir(fd) == 0 {
158 Ok(())
159 } else {
160 Err(std::io::Error::last_os_error())
161 }
162 });
163 }
164 }
165
166 // Expand `${NAME}` placeholders so secret env values can be sourced
167 // from the process environment instead of being stored in cleartext
168 // in the MCP config. The child env is allowlist-sanitized below, so
169 // these vars would not otherwise be inherited by the child.
170 let expanded_env = super::expanded_mcp_stdio_env(config)
171 .with_context(|| format!("MCP server '{server_name}' env expansion failed"))?;
172
173 // User-configured MCP keeps the compatibility-oriented Node/Python
174 // bootstrap allowlist (#1244). Reviewed plugins receive only the base
175 // secret-scrubbed child environment plus their explicitly reviewed
176 // mappings, so namespaces such as NPM_CONFIG_* are never inherited
177 // ambiently across the consent boundary.
178 if let Some(reviewed_plugin) = config.reviewed_plugin.as_ref() {
179 cmd.env_clear();
180 for (key, value) in child_env::sanitized_plugin_mcp_env_from(
181 reviewed_plugin.host_environment.entries().iter().cloned(),
182 child_env::string_map_env(&expanded_env),
183 ) {
184 cmd.env(key, value);
185 }
186 } else {
187 child_env::apply_to_tokio_command_mcp(
188 &mut cmd,
189 child_env::string_map_env(&expanded_env),
190 );
191 }
192 // Lead a process group of its own so teardown can reach descendants.
193 #[cfg(unix)]
194 cmd.process_group(0);
195
196 let mut child = cmd.spawn().map_err(|error| {
197 let message = if error.kind() == std::io::ErrorKind::NotFound
198 && super::is_node_command(command)
199 && launch_cwd.is_none_or(|directory| directory.is_dir())
200 {
201 format!("MCP server {server_name} could not start because Node.js was not found. Install Node.js from https://nodejs.org/ and restart Codewhale with node on PATH. Built-in Computer Use requires Node.js 20 or newer.")
202 } else if config.reviewed_plugin.is_some() {
203 format!(
204 "MCP stdio spawn failed (transport=stdio server={server_name} reviewed-plugin argv_count={} env_count={})",
205 config.args.len(),
206 expanded_env.len(),
207 )
208 } else {
209 let env_keys: Vec<&str> = expanded_env.keys().map(String::as_str).collect();
210 format!(
211 "MCP stdio spawn failed (transport=stdio server={server_name} cmd={command:?} args={:?} env_keys={env_keys:?})",
212 config.args,
213 )
214 };
215 anyhow::Error::new(error).context(message)
216 })?;
217
218 let process_tree = match crate::process_tree::ProcessTree::attach_tokio(&child) {
219 Ok(tree) => Arc::new(tree),
220 Err(error) => {
221 let _ = child.start_kill();
222 return Err(anyhow::Error::new(error).context(format!(
223 "MCP server {server_name} could not be contained with its child processes"
224 )));
225 }
226 };
227
228 let stdin = child.stdin.take().context("Failed to get MCP stdin")?;
229 let stdout = child.stdout.take().context("Failed to get MCP stdout")?;
230 let stderr = child.stderr.take().context("Failed to get MCP stderr")?;
231
232 // Drain stderr into a bounded ring buffer so a crash mid-run leaves
233 // diagnostic breadcrumbs instead of disappearing into `Stdio::null`.
234 // The task exits naturally when the child closes its stderr
235 // (kill_on_drop / exit / explicit shutdown).
236 let stderr_tail = StderrTail::new();
237 {
238 let tail = Arc::clone(&stderr_tail);
239 // A reviewed plugin child receives environment-backed values that
240 // are intentionally absent from its manifest. Still drain its
241 // stderr to avoid blocking, but do not retain or surface arbitrary
242 // child output that could echo those credentials into a chat or
243 // persisted transcript.
244 let capture = config.reviewed_plugin.is_none().then_some(tail);
245 tokio::spawn(drain_stderr(stderr, capture));
246 }
247
248 let child = Arc::new(TokioMutex::new(child));
249 let authority_cancel_watch = config.reviewed_plugin.as_ref().map(|_| {
250 let watched_child = Arc::clone(&child);
251 let watched_tree = Arc::clone(&process_tree);
252 tokio::spawn(async move {
253 cancel_token.cancelled().await;
254 terminate_child_for_authority_change(&watched_child, &watched_tree).await;
255 })
256 });
257
258 Ok((
259 Self {
260 child,
261 stdin,
262 stderr_tail,
263 authority_cancel_watch,
264 _reviewed_launch: reviewed_launch,
265 process_tree: Some(process_tree),
266 },
267 stdout,
268 ))
269 }
270 }
271
272 impl BrokerSession {
273 /// Write already-framed bytes; the transport owns the framing policy.
274 pub(crate) async fn write(&mut self, bytes: &[u8]) -> Result<()> {
275 self.stdin.write_all(bytes).await?;
276 self.stdin.flush().await?;
277 Ok(())
278 }
279
280 pub(crate) async fn last_stderr_line(&self) -> Option<String> {
281 // The child can write its reason immediately before the error reply.
282 tokio::task::yield_now().await;
283 if let Some(line) = self.stderr_tail.last_line().await {
284 return Some(line);
285 }
286 tokio::time::sleep(Duration::from_millis(100)).await;
287 self.stderr_tail.last_line().await
288 }
289
290 pub(crate) async fn stderr_context(&self) -> Option<String> {
291 format_stderr_context(&self.stderr_tail).await
292 }
293
294 pub(crate) async fn exit_status(&self) -> Option<std::process::ExitStatus> {
295 self.child.lock().await.try_wait().ok().flatten()
296 }
297
298 /// Never await or spawn for readiness; a contended child still reads live.
299 pub(crate) fn probe_dead(&self) -> bool {
300 match self.child.try_lock() {
301 Ok(mut child) => matches!(child.try_wait(), Ok(Some(_))),
302 Err(_) => false,
303 }
304 }
305
306 /// Reap the direct child and contain descendants through the same grace.
307 pub(crate) async fn shutdown(&mut self) {
308 let mut child = self.child.lock().await;
309 terminate_child(&mut child, self.process_tree.as_deref()).await;
310 }
311
312 #[cfg(all(test, unix))]
313 pub(super) fn child_for_tests(&self) -> Arc<TokioMutex<Child>> {
314 Arc::clone(&self.child)
315 }
316 }
317
318 /// Longest stderr line retained in the tail, in bytes after decoding. An
319 /// overlong line keeps its end, where the error usually is, and the start is
320 /// drained and discarded so a newline-free progress stream cannot grow memory.
321 const STDERR_LINE_CAP: usize = 4 * 1024;
322
323 /// Drain a child's stderr until EOF, retaining lossily decoded, length-capped
324 /// lines in `tail` when given. Stderr is not a protocol channel: a non-UTF-8
325 /// byte or an endless line must never stop the drain, because dropping the
326 /// pipe makes the child's next stderr write fail (EPIPE/SIGPIPE) and hides
327 /// the context this tail exists to keep.
328 async fn drain_stderr<R>(stderr: R, tail: Option<Arc<StderrTail>>)
329 where
330 R: tokio::io::AsyncRead + Unpin,
331 {
332 let mut reader = tokio::io::BufReader::new(stderr);
333 let mut line = Vec::new();
334 loop {
335 let (consumed, line_ended) = {
336 let available = match reader.fill_buf().await {
337 Ok(available) => available,
338 Err(error) if error.kind() == std::io::ErrorKind::Interrupted => continue,
339 Err(_) => break,
340 };
341 if available.is_empty() {
342 break;
343 }
344 let (consumed, line_ended) = match available.iter().position(|&b| b == b'\n') {
345 Some(pos) => (pos + 1, true),
346 None => (available.len(), false),
347 };
348 if tail.is_some() {
349 let content = &available[..consumed - usize::from(line_ended)];
350 let keep = content.len().min(STDERR_LINE_CAP);
351 let overflow = (line.len() + keep).saturating_sub(STDERR_LINE_CAP);
352 line.drain(..overflow);
353 line.extend_from_slice(&content[content.len() - keep..]);
354 }
355 (consumed, line_ended)
356 };
357 reader.consume(consumed);
358 if line_ended {
359 push_stderr_line(tail.as_deref(), &mut line).await;
360 }
361 }
362 if !line.is_empty() {
363 push_stderr_line(tail.as_deref(), &mut line).await;
364 }
365 }
366
367 async fn push_stderr_line(tail: Option<&StderrTail>, line: &mut Vec<u8>) {
368 if let Some(tail) = tail {
369 let bytes: &[u8] = line;
370 let text = String::from_utf8_lossy(bytes.strip_suffix(b"\r").unwrap_or(bytes));
371 // Lossy decoding can triple invalid bytes; hold the retained size to
372 // the cap too, again keeping the end.
373 let mut start = text.len().saturating_sub(STDERR_LINE_CAP);
374 while !text.is_char_boundary(start) {
375 start += 1;
376 }
377 tail.push(text[start..].to_owned()).await;
378 }
379 line.clear();
380 }
381
382 /// Format the captured stderr tail for inclusion in an error message. Empty
383 /// tails return `None` so the caller can fall back to its original message.
384 async fn format_stderr_context(tail: &StderrTail) -> Option<String> {
385 let lines = tail.snapshot().await;
386 if lines.is_empty() {
387 return None;
388 }
389 Some(format!(
390 "MCP server stderr (last {} line{}):\n{}",
391 lines.len(),
392 if lines.len() == 1 { "" } else { "s" },
393 lines.join("\n"),
394 ))
395 }
396
397 /// Best-effort SIGTERM. On Unix uses `libc::kill`, addressed to the child's
398 /// whole process group when it leads one (`contained`); on Windows there's no
399 /// equivalent so we let `kill_on_drop` (TerminateProcess) and the Job Object
400 /// handle it. Returns whether a signal was actually sent.
401 fn send_sigterm(child: &Child, contained: bool) -> bool {
402 #[cfg(unix)]
403 {
404 if let Some(pid) = child.id() {
405 let pid = pid as i32;
406 let target = if contained { -pid } else { pid };
407 // SAFETY: pid was just obtained from `child.id()` of an unreaped
408 // child, so neither it nor the group it leads can have been
409 // recycled. `libc::kill` with `SIGTERM` is async-signal-safe and
410 // never observes invalid memory. ESRCH is deliberately ignored.
411 unsafe {
412 let _ = libc::kill(target, libc::SIGTERM);
413 }
414 return true;
415 }
416 false
417 }
418 #[cfg(not(unix))]
419 {
420 let _ = (child, contained);
421 false
422 }
423 }
424
425 async fn terminate_child_for_authority_change(
426 child: &Arc<TokioMutex<Child>>,
427 tree: &crate::process_tree::ProcessTree,
428 ) {
429 let mut child = child.lock().await;
430 terminate_child(&mut child, Some(tree)).await;
431 }
432
433 async fn terminate_child(child: &mut Child, tree: Option<&crate::process_tree::ProcessTree>) {
434 // Reap an already-exited child before resolving its PID. Until it is
435 // reaped, the OS cannot recycle that identity; after it is reaped there is
436 // nothing left to signal. This avoids a PID-only watcher ever targeting an
437 // unrelated process after rapid PID reuse.
438 if child.try_wait().is_ok_and(|status| status.is_some()) {
439 // The server is gone, but what it started may not be.
440 if let Some(tree) = tree {
441 let _ = tree.kill();
442 }
443 return;
444 }
445
446 #[cfg(unix)]
447 send_sigterm(child, tree.is_some());
448
449 #[cfg(not(unix))]
450 let _ = child.start_kill();
451
452 match tokio::time::timeout(STDIO_SHUTDOWN_GRACE, child.wait()).await {
453 Ok(Ok(_)) => {}
454 Ok(Err(_)) | Err(_) => {
455 // SIGTERM is advisory. Revocation and explicit shutdown must not
456 // leave the reviewed child alive indefinitely.
457 let _ = child.start_kill();
458 let _ = child.wait().await;
459 }
460 }
461 // Descendants that outlived the grace (or ignored SIGTERM) go now.
462 if let Some(tree) = tree {
463 let _ = tree.kill();
464 }
465 }
466
467 /// Session changes can drop a pool without explicitly awaiting shutdown.
468 /// Keep the owned child alive for the same bounded cleanup as explicit
469 /// shutdown, so servers can release input and recording resources. Runtime
470 /// teardown still drops the cleanup future and invokes `kill_on_drop`.
471 impl Drop for BrokerSession {
472 fn drop(&mut self) {
473 if let Some(watch) = self.authority_cancel_watch.take() {
474 watch.abort();
475 }
476 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
477 let child = Arc::clone(&self.child);
478 let reviewed_launch = self._reviewed_launch.take();
479 // The task owns the tree so it is not killed before the grace.
480 let tree = self.process_tree.take();
481 runtime.spawn(async move {
482 let _reviewed_launch = reviewed_launch;
483 let mut child = child.lock().await;
484 terminate_child(&mut child, tree.as_deref()).await;
485 });
486 return;
487 }
488 if let Ok(mut child) = self.child.try_lock()
489 && !child.try_wait().is_ok_and(|status| status.is_some())
490 {
491 send_sigterm(&child, self.process_tree.is_some());
492 }
493 // No runtime: dropping the tree below SIGKILLs the group / closes the
494 // job, the same backstop `kill_on_drop` gives the direct child.
495 }
496 }
497
498 #[cfg(test)]
499 mod stderr_drain_tests {
500 use super::{STDERR_LINE_CAP, StderrTail, drain_stderr};
501 use tokio::io::AsyncWriteExt;
502
503 #[tokio::test]
504 async fn invalid_utf8_and_long_lines_do_not_stop_the_drain() {
505 let (mut writer, reader) = tokio::io::duplex(64);
506 let tail = StderrTail::new();
507 let drain = tokio::spawn(drain_stderr(reader, Some(tail.clone())));
508 writer.write_all(b"starting\r\n").await.unwrap();
509 writer.write_all(b"caf\xe9 \xff\xfe\n").await.unwrap();
510 writer
511 .write_all(&vec![b'#'; STDERR_LINE_CAP * 3])
512 .await
513 .unwrap();
514 writer.write_all(b"\npanic: boom").await.unwrap();
515 drop(writer);
516 drain.await.unwrap();
517 let lines = tail.snapshot().await;
518 assert_eq!(lines.len(), 4, "{lines:?}");
519 assert_eq!(lines[0], "starting");
520 assert_eq!(lines[1], "caf\u{fffd} \u{fffd}\u{fffd}");
521 assert_eq!(lines[2].len(), STDERR_LINE_CAP);
522 assert_eq!(lines[3], "panic: boom");
523 }
524
525 #[tokio::test]
526 async fn overlong_lines_keep_their_end_within_the_cap() {
527 let (mut writer, reader) = tokio::io::duplex(64);
528 let tail = StderrTail::new();
529 let drain = tokio::spawn(drain_stderr(reader, Some(tail.clone())));
530 // A long prefix, then the actual failure at the end of the line.
531 writer
532 .write_all(&vec![b'.'; STDERR_LINE_CAP * 3])
533 .await
534 .unwrap();
535 writer.write_all(b"error: config missing\n").await.unwrap();
536 // Invalid bytes decode to three-byte U+FFFD each.
537 writer
538 .write_all(&vec![0xff; STDERR_LINE_CAP])
539 .await
540 .unwrap();
541 writer.write_all(b"\n").await.unwrap();
542 drop(writer);
543 drain.await.unwrap();
544 let lines = tail.snapshot().await;
545 assert_eq!(lines.len(), 2, "{lines:?}");
546 assert!(
547 lines[0].ends_with("error: config missing"),
548 "{:?}",
549 &lines[0][lines[0].len() - 40..]
550 );
551 assert_eq!(lines[0].len(), STDERR_LINE_CAP);
552 assert!(lines[1].len() <= STDERR_LINE_CAP, "{}", lines[1].len());
553 assert!(lines[1].chars().all(|c| c == '\u{fffd}'));
554 }
555
556 #[tokio::test]
557 async fn uncaptured_stderr_is_still_drained_to_eof() {
558 let (mut writer, reader) = tokio::io::duplex(16);
559 let drain = tokio::spawn(drain_stderr(reader, None));
560 // A non-UTF-8 line, then far more than the pipe buffer: a drain that
561 // stopped at the bad line would fail this write with a broken pipe.
562 writer.write_all(b"\xff\n").await.unwrap();
563 writer.write_all(&[b'x'; 4096]).await.unwrap();
564 drop(writer);
565 drain.await.unwrap();
566 }
567 }
568
568 lines RUST