use std::io::{self, BufRead, Write};
use tracing::{debug, error, info};
use super::{Connection, McpServer};
#[derive(Debug, Clone)]
pub enum Transport {
Stdio,
Tcp(std::net::SocketAddr),
#[cfg(unix)]
Unix(std::path::PathBuf),
}
impl McpServer {
pub fn run_stdio(&self) -> io::Result<()> {
self.run(io::stdin().lock(), io::stdout().lock())
}
pub fn run_tcp(&self, addr: impl tokio::net::ToSocketAddrs) -> io::Result<()> {
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
rt.block_on(self.listen_tcp(addr))
}
#[cfg(unix)]
pub fn run_unix(&self, path: impl AsRef<std::path::Path>) -> io::Result<()> {
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
rt.block_on(self.listen_unix(path))
}
pub fn serve(&self, transport: Transport) -> io::Result<()> {
match transport {
Transport::Stdio => self.run_stdio(),
Transport::Tcp(addr) => self.run_tcp(addr),
#[cfg(unix)]
Transport::Unix(path) => self.run_unix(path),
}
}
pub fn run(&self, reader: impl BufRead, mut writer: impl Write) -> io::Result<()> {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
self.run_with_runtime(&rt, reader, &mut writer)
}
pub fn run_with_runtime(
&self,
rt: &tokio::runtime::Runtime,
reader: impl BufRead,
writer: &mut impl Write,
) -> io::Result<()> {
let mut out_buf: Vec<u8> = Vec::new();
let mut conn = Connection::default();
for line_result in reader.lines() {
let line = line_result?;
if line.trim().is_empty() {
continue;
}
debug!(request = %line, "mcp request");
let Some(outcome) = rt.block_on(self.handle_message_conn(&line, &mut conn)) else {
debug!("dropping notification response");
continue;
};
out_buf.clear();
outcome.write_json(&mut out_buf);
debug!(response = %String::from_utf8_lossy(&out_buf), "mcp response");
out_buf.push(b'\n');
writer.write_all(&out_buf)?;
writer.flush()?;
}
info!("input stream closed — shutting down");
Ok(())
}
pub async fn run_async(
&self,
reader: impl tokio::io::AsyncBufRead + Unpin,
mut writer: impl tokio::io::AsyncWrite + Unpin,
) -> io::Result<()> {
use tokio::io::{AsyncBufReadExt, AsyncWriteExt};
let mut lines = reader.lines();
let mut out_buf: Vec<u8> = Vec::new();
let mut conn = Connection::default();
while let Some(line) = lines.next_line().await? {
if line.trim().is_empty() {
continue;
}
debug!(request = %line, "mcp request");
let Some(outcome) = self.handle_message_conn(&line, &mut conn).await else {
debug!("dropping notification response");
continue;
};
out_buf.clear();
outcome.write_json(&mut out_buf);
debug!(response = %String::from_utf8_lossy(&out_buf), "mcp response");
out_buf.push(b'\n');
writer.write_all(&out_buf).await?;
writer.flush().await?;
}
info!("input stream closed — shutting down");
Ok(())
}
pub async fn listen_tcp(&self, addr: impl tokio::net::ToSocketAddrs) -> io::Result<()> {
let listener = tokio::net::TcpListener::bind(addr).await?;
info!(addr = ?listener.local_addr()?, "listening on TCP for MCP connections");
self.run_tcp_listener(listener).await
}
pub async fn run_tcp_listener(&self, listener: tokio::net::TcpListener) -> io::Result<()> {
loop {
let (mut socket, peer_addr) = listener.accept().await?;
info!(peer = %peer_addr, "accepted MCP TCP connection");
let server = self.clone();
tokio::spawn(async move {
let (reader, writer) = socket.split();
let reader = tokio::io::BufReader::new(reader);
if let Err(e) = server.run_async(reader, writer).await {
error!(peer = %peer_addr, error = %e, "MCP TCP connection error");
}
info!(peer = %peer_addr, "MCP TCP connection closed");
});
}
}
#[cfg(unix)]
pub async fn listen_unix(&self, path: impl AsRef<std::path::Path>) -> io::Result<()> {
let path = path.as_ref();
if path.exists() {
if tokio::net::UnixStream::connect(path).await.is_ok() {
return Err(io::Error::new(
io::ErrorKind::AddrInUse,
format!(
"another process is already listening on Unix socket {}",
path.display()
),
));
}
if let Err(e) = std::fs::remove_file(path) {
tracing::warn!(path = ?path, error = %e, "failed to remove stale socket file");
}
}
let listener = tokio::net::UnixListener::bind(path)?;
info!(path = ?path, "listening on Unix domain socket for MCP connections");
self.run_unix_listener(listener).await
}
#[cfg(unix)]
pub async fn run_unix_listener(&self, listener: tokio::net::UnixListener) -> io::Result<()> {
loop {
let (mut socket, _) = listener.accept().await?;
info!("accepted MCP Unix domain socket connection");
let server = self.clone();
tokio::spawn(async move {
let (reader, writer) = socket.split();
let reader = tokio::io::BufReader::new(reader);
if let Err(e) = server.run_async(reader, writer).await {
error!(error = %e, "MCP Unix connection error");
}
info!("MCP Unix connection closed");
});
}
}
}