Skip to repository content

tenant.openagents/omega

No repository description is available.

OpenAgents Git authority 2026-07-28T04:07:05.625Z Public web read
NIP-34 coordinate30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omega
MaintainersHidden in public view
References2 branches · 1 tag
Read-only clonegit clone https://openagents.com/git/tenant.openagents/omega.git
Browse files

stdio_transport.rs

149 lines · 4.4 KB · rust
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
Served at tenant.openagents/omega Member data and write actions are omitted.