Skip to main content

agent_client_protocol/
stdio.rs

1//! Stdio transport for connecting ACP components via standard input/output.
2
3use crate::{ByteStreams, ConnectTo, LineDirection, Role};
4use std::sync::Arc;
5
6/// A transport that connects to an ACP peer via standard input/output.
7///
8/// This is useful for building agents or proxies that communicate over stdio,
9/// which is the standard transport for MCP and ACP subprocess communication.
10pub struct Stdio {
11    debug_callback: Option<Arc<dyn Fn(&str, LineDirection) + Send + Sync + 'static>>,
12}
13
14impl std::fmt::Debug for Stdio {
15    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
16        f.debug_struct("Stdio").finish_non_exhaustive()
17    }
18}
19
20impl Stdio {
21    /// Create a new `Stdio` transport.
22    #[must_use]
23    pub fn new() -> Self {
24        Self {
25            debug_callback: None,
26        }
27    }
28
29    /// Add a debug callback that will be invoked for each line sent/received.
30    #[must_use]
31    pub fn with_debug<F>(mut self, callback: F) -> Self
32    where
33        F: Fn(&str, LineDirection) + Send + Sync + 'static,
34    {
35        self.debug_callback = Some(Arc::new(callback));
36        self
37    }
38}
39
40impl Default for Stdio {
41    fn default() -> Self {
42        Self::new()
43    }
44}
45
46impl<Counterpart: Role> ConnectTo<Counterpart> for Stdio {
47    async fn connect_to(
48        self,
49        client: impl ConnectTo<Counterpart::Counterpart>,
50    ) -> Result<(), crate::Error> {
51        let stdin = blocking::Unblock::new(std::io::stdin());
52        let stdout = blocking::Unblock::new(std::io::stdout());
53
54        if let Some(callback) = self.debug_callback {
55            use futures::io::BufReader;
56            use futures::{AsyncBufReadExt, StreamExt};
57
58            let incoming_callback = callback.clone();
59            let incoming_lines = Box::pin(BufReader::new(stdin).lines().inspect(move |result| {
60                if let Ok(line) = result {
61                    incoming_callback(line, LineDirection::Stdin);
62                }
63            }))
64                as std::pin::Pin<Box<dyn futures::Stream<Item = std::io::Result<String>> + Send>>;
65
66            let outgoing_sink = Box::pin(futures::sink::unfold(
67                (stdout, callback),
68                async move |(mut writer, callback), line: String| {
69                    callback(&line, LineDirection::Stdout);
70                    crate::jsonrpc::write_line(&mut writer, line).await?;
71                    Ok::<_, std::io::Error>((writer, callback))
72                },
73            ))
74                as std::pin::Pin<Box<dyn futures::Sink<String, Error = std::io::Error> + Send>>;
75
76            ConnectTo::<Counterpart>::connect_to(
77                crate::Lines::new(outgoing_sink, incoming_lines),
78                client,
79            )
80            .await
81        } else {
82            ConnectTo::<Counterpart>::connect_to(ByteStreams::new(stdout, stdin), client).await
83        }
84    }
85}