Skip to repository content149 lines · 4.4 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T05:06:39.088Z Public web read
NIP-34 coordinate
30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omegaMaintainersHidden in public view
References2 branches · 1 tag
Read-only clone
git clone https://openagents.com/git/tenant.openagents/omega.gitBrowse files
stdio_transport.rs
1use std::path::PathBuf;
2use std::pin::Pin;
3
4use anyhow::Result;
5use async_trait::async_trait;
6use futures::io::{BufReader, BufWriter};
7use futures::{
8 AsyncBufReadExt as _, AsyncRead, AsyncWrite, AsyncWriteExt as _, Stream, StreamExt as _,
9};
10use gpui::AsyncApp;
11
12use util::TryFutureExt as _;
13use util::process::Child;
14use util::shell::Shell;
15use util::shell_builder::ShellBuilder;
16
17use crate::client::ModelContextServerBinary;
18use crate::transport::Transport;
19
20pub struct StdioTransport {
21 stdout_sender: async_channel::Sender<String>,
22 stdin_receiver: async_channel::Receiver<String>,
23 stderr_receiver: async_channel::Receiver<String>,
24 server: Child,
25}
26
27impl StdioTransport {
28 pub fn new(
29 binary: ModelContextServerBinary,
30 working_directory: &Option<PathBuf>,
31 cx: &AsyncApp,
32 ) -> Result<Self> {
33 let builder = ShellBuilder::new(&Shell::System, cfg!(windows)).non_interactive();
34 let mut command =
35 builder.build_std_command(Some(binary.executable.display().to_string()), &binary.args);
36
37 command.envs(binary.env.unwrap_or_default());
38
39 if let Some(working_directory) = working_directory {
40 command.current_dir(working_directory);
41 }
42
43 let mut server = Child::spawn(
44 command,
45 std::process::Stdio::piped(),
46 std::process::Stdio::piped(),
47 std::process::Stdio::piped(),
48 )?;
49
50 let stdin = server.stdin.take().unwrap();
51 let stdout = server.stdout.take().unwrap();
52 let stderr = server.stderr.take().unwrap();
53
54 let (stdin_sender, stdin_receiver) = async_channel::unbounded::<String>();
55 let (stdout_sender, stdout_receiver) = async_channel::unbounded::<String>();
56 let (stderr_sender, stderr_receiver) = async_channel::unbounded::<String>();
57
58 cx.spawn(async move |_| Self::handle_output(stdin, stdout_receiver).log_err().await)
59 .detach();
60
61 cx.spawn(async move |_| Self::handle_input(stdout, stdin_sender).await)
62 .detach();
63
64 cx.spawn(async move |_| Self::handle_err(stderr, stderr_sender).await)
65 .detach();
66
67 Ok(Self {
68 stdout_sender,
69 stdin_receiver,
70 stderr_receiver,
71 server,
72 })
73 }
74
75 async fn handle_input<Stdout>(stdin: Stdout, inbound_rx: async_channel::Sender<String>)
76 where
77 Stdout: AsyncRead + Unpin + Send + 'static,
78 {
79 let mut stdin = BufReader::new(stdin);
80 let mut line = String::new();
81 while let Ok(n) = stdin.read_line(&mut line).await {
82 if n == 0 {
83 break;
84 }
85 if inbound_rx.send(line.clone()).await.is_err() {
86 break;
87 }
88 line.clear();
89 }
90 }
91
92 async fn handle_output<Stdin>(
93 stdin: Stdin,
94 outbound_rx: async_channel::Receiver<String>,
95 ) -> Result<()>
96 where
97 Stdin: AsyncWrite + Unpin + Send + 'static,
98 {
99 let mut stdin = BufWriter::new(stdin);
100 let mut pinned_rx = Box::pin(outbound_rx);
101 while let Some(message) = pinned_rx.next().await {
102 log::trace!("outgoing message: {}", message);
103
104 stdin.write_all(message.as_bytes()).await?;
105 stdin.write_all(b"\n").await?;
106 stdin.flush().await?;
107 }
108 Ok(())
109 }
110
111 async fn handle_err<Stderr>(stderr: Stderr, stderr_tx: async_channel::Sender<String>)
112 where
113 Stderr: AsyncRead + Unpin + Send + 'static,
114 {
115 let mut stderr = BufReader::new(stderr);
116 let mut line = String::new();
117 while let Ok(n) = stderr.read_line(&mut line).await {
118 if n == 0 {
119 break;
120 }
121 if stderr_tx.send(line.clone()).await.is_err() {
122 break;
123 }
124 line.clear();
125 }
126 }
127}
128
129#[async_trait]
130impl Transport for StdioTransport {
131 async fn send(&self, message: String) -> Result<()> {
132 Ok(self.stdout_sender.send(message).await?)
133 }
134
135 fn receive(&self) -> Pin<Box<dyn Stream<Item = String> + Send>> {
136 Box::pin(self.stdin_receiver.clone())
137 }
138
139 fn receive_err(&self) -> Pin<Box<dyn Stream<Item = String> + Send>> {
140 Box::pin(self.stderr_receiver.clone())
141 }
142}
143
144impl Drop for StdioTransport {
145 fn drop(&mut self) {
146 let _ = self.server.kill();
147 }
148}
149