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