返回 CodeWhale
stdio.rs
根目录 / crates / tui / src / mcp / stdio.rs
1 //! MCP newline framing over the broker-owned reviewed stdio process.
2 use anyhow::Result;
3 use tokio::process::ChildStdout;
4
5 use super::process_broker::BrokerSession;
6 use super::wire::{MAX_MCP_RESPONSE_BYTES, read_line_capped};
7 use super::{McpServerConfig, McpTransport};
8
9 pub(super) struct StdioTransport {
10 pub(super) session: BrokerSession,
11 reader: tokio::io::BufReader<ChildStdout>,
12 /// Partial frame bytes survive cancellation of the receive future.
13 pub(super) pending_line: Vec<u8>,
14 }
15
16 impl StdioTransport {
17 pub(super) fn spawn(
18 server_name: &str,
19 command: &str,
20 config: &McpServerConfig,
21 cancel_token: tokio_util::sync::CancellationToken,
22 ) -> Result<Self> {
23 let (session, stdout) = BrokerSession::spawn(server_name, command, config, cancel_token)?;
24 Ok(Self {
25 session,
26 reader: tokio::io::BufReader::new(stdout),
27 pending_line: Vec::new(),
28 })
29 }
30 }
31
32 #[async_trait::async_trait]
33 impl McpTransport for StdioTransport {
34 async fn last_stderr_line(&self) -> Option<String> {
35 self.session.last_stderr_line().await
36 }
37
38 async fn send(&mut self, mut msg: Vec<u8>) -> Result<()> {
39 msg.push(b'\n');
40 self.session.write(&msg).await
41 }
42
43 /// Non-blocking liveness probe: a reaped child means the transport is
44 /// dead even though the `Ready` flag is still set (#6187). The sync
45 /// trait contract forbids awaiting the lock, so a contended lock reads
46 /// as alive — the read side observes the death on the next call.
47 fn probe_dead(&self) -> bool {
48 self.session.probe_dead()
49 }
50
51 async fn recv(&mut self) -> Result<Vec<u8>> {
52 loop {
53 // Bounded read: a server emitting a newline-free multi-GB "line"
54 // must not OOM us (read_line is unbounded).
55 let bytes = match read_line_capped(
56 &mut self.reader,
57 &mut self.pending_line,
58 MAX_MCP_RESPONSE_BYTES,
59 )
60 .await
61 {
62 Ok(b) => b,
63 Err(err) => {
64 if let Some(stderr) = self.session.stderr_context().await {
65 anyhow::bail!("Stdio transport read error: {err}\n{stderr}");
66 }
67 return Err(err.into());
68 }
69 };
70 if bytes == 0 {
71 // Let the stderr drain task catch up before snapshotting, and
72 // name the exit status: a reviewed plugin's stderr is never
73 // retained, so the status is the only reason the operator
74 // gets when the child dies before the handshake (#5916).
75 tokio::task::yield_now().await;
76 let exit = self.session.exit_status().await;
77 let exit = exit.map_or_else(String::new, |status| format!(" ({status})"));
78 if let Some(stderr) = self.session.stderr_context().await {
79 anyhow::bail!("Stdio transport closed{exit}\n{stderr}");
80 }
81 anyhow::bail!("Stdio transport closed{exit}");
82 }
83
84 let line_bytes = std::mem::take(&mut self.pending_line);
85 let line = String::from_utf8_lossy(&line_bytes);
86 let trimmed = line.trim();
87 if trimmed.is_empty() {
88 continue;
89 }
90
91 return Ok(trimmed.as_bytes().to_vec());
92 }
93 }
94
95 /// Send SIGTERM and wait up to `STDIO_SHUTDOWN_GRACE` for graceful exit,
96 /// then force termination and reap the child as the backstop.
97 async fn shutdown(&mut self) {
98 self.session.shutdown().await;
99 }
100 }
101
101 lines RUST