use origin_domain::{AppError, Result};
use origin_mcp_core::McpServer;
use tokio::io::{AsyncBufReadExt, AsyncWrite, AsyncWriteExt, BufReader};
pub async fn serve(server: &McpServer) -> Result<()> {
let input = BufReader::new(tokio::io::stdin());
let output = tokio::io::stdout();
serve_streams(server, input, output).await
}
pub async fn serve_streams<R, W>(
server: &McpServer,
input: BufReader<R>,
mut output: W,
) -> Result<()>
where
R: tokio::io::AsyncRead + Unpin,
W: AsyncWrite + Unpin,
{
tracing::debug!("mcp stdio transport started");
let mut lines = input.lines();
loop {
let line = match lines.next_line().await {
Ok(Some(line)) => line,
Ok(None) => break,
Err(error) => {
return Err(AppError::internal(format!(
"cannot read from stdin: {error}"
)));
}
};
if line.trim().is_empty() {
continue;
}
let Some(response) = server.handle_line(&line).await else {
continue;
};
let encoded = serde_json::to_string(&response)
.map_err(|error| AppError::internal(format!("cannot encode response: {error}")))?;
output
.write_all(encoded.as_bytes())
.await
.map_err(write_error)?;
output.write_all(b"\n").await.map_err(write_error)?;
output.flush().await.map_err(write_error)?;
}
tracing::debug!("mcp stdio transport stopped");
Ok(())
}
fn write_error(error: std::io::Error) -> AppError {
AppError::internal(format!("cannot write to stdout: {error}"))
}