| 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 |