From 9009f127a9a3e884366424a692e1129ac2cdd80a Mon Sep 17 00:00:00 2001 From: Riz Ashraf Date: Sat, 19 Sep 2026 07:06:14 +0100 Subject: [PATCH] refactor: Extract shared MCP Stdio NDJSON parsing loop into mcp-stdio library --- Cargo.lock | 10 ++++++++ Cargo.toml | 2 +- mcp-stdio/Cargo.toml | 8 ++++++ mcp-stdio/src/lib.rs | 40 ++++++++++++++++++++++++++++++ nvim-core/Cargo.toml | 1 + nvim-core/src/lib.rs | 59 +++++++++++--------------------------------- stub/Cargo.toml | 1 + stub/src/main.rs | 36 ++------------------------- 8 files changed, 77 insertions(+), 80 deletions(-) create mode 100644 mcp-stdio/Cargo.toml create mode 100644 mcp-stdio/src/lib.rs diff --git a/Cargo.lock b/Cargo.lock index ced576d..aa3f8d4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1375,6 +1375,7 @@ dependencies = [ "clap", "dirs 7.0.0", "futures-util", + "mcp-stdio", "reqwest", "serde_json", "tokio", @@ -1404,6 +1405,14 @@ dependencies = [ "tracing-subscriber", ] +[[package]] +name = "mcp-stdio" +version = "0.1.0" +dependencies = [ + "tokio", + "tracing", +] + [[package]] name = "measure_time" version = "0.9.0" @@ -1527,6 +1536,7 @@ name = "nvim-core" version = "0.1.0" dependencies = [ "dirs 7.0.0", + "mcp-stdio", "rmcp", "rmp-serde", "rmpv", diff --git a/Cargo.toml b/Cargo.toml index a55e3be..0276da0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,5 +4,5 @@ members = [ "stub", "win-nvim", "linux-nvim" -, "nvim-core"] +, "nvim-core", "mcp-stdio"] resolver = "2" diff --git a/mcp-stdio/Cargo.toml b/mcp-stdio/Cargo.toml new file mode 100644 index 0000000..6a87582 --- /dev/null +++ b/mcp-stdio/Cargo.toml @@ -0,0 +1,8 @@ +[package] +name = "mcp-stdio" +version = "0.1.0" +edition = "2024" + +[dependencies] +tokio = { version = "1.53.1", features = ["io-util"] } +tracing = "0.1.44" diff --git a/mcp-stdio/src/lib.rs b/mcp-stdio/src/lib.rs new file mode 100644 index 0000000..ca4d20d --- /dev/null +++ b/mcp-stdio/src/lib.rs @@ -0,0 +1,40 @@ +use tokio::io::{AsyncBufReadExt, AsyncReadExt, BufReader}; + +/// Reads an MCP (NDJSON or LSP Content-Length prefixed) message from a buffered async reader. +/// Returns the raw JSON string payload if successful, or None on EOF or error. +pub async fn read_mcp_message( + stdin: &mut BufReader, +) -> Option { + let mut length = 0; + loop { + let mut line = String::new(); + if stdin.read_line(&mut line).await.unwrap_or(0) == 0 { + return None; + } + + if line.starts_with('{') { + return Some(line.trim_end().to_string()); + } + + let line = line.trim_end(); + if line.is_empty() { + break; + } + + let lower_line = line.to_lowercase(); + if let Some(len_str) = lower_line.strip_prefix("content-length:") { + length = len_str.trim().parse().unwrap_or(0); + } + } + + if length == 0 { + return None; + } + + let mut buffer = vec![0; length]; + if stdin.read_exact(&mut buffer).await.is_err() { + return None; + } + + String::from_utf8(buffer).ok() +} diff --git a/nvim-core/Cargo.toml b/nvim-core/Cargo.toml index 1d53f20..774d24f 100644 --- a/nvim-core/Cargo.toml +++ b/nvim-core/Cargo.toml @@ -14,4 +14,5 @@ tracing-appender = "0.2.5" tracing-subscriber = "0.3.23" dirs = "7.0.0" rmcp = { version = "3.4.0", features = ["server"] } +mcp-stdio = { version = "0.1.0", path = "../mcp-stdio" } diff --git a/nvim-core/src/lib.rs b/nvim-core/src/lib.rs index 0d84ffe..70d5304 100644 --- a/nvim-core/src/lib.rs +++ b/nvim-core/src/lib.rs @@ -20,47 +20,7 @@ pub struct JsonRpcResponse { pub error: Option, } -pub async fn read_message( - stdin: &mut BufReader, -) -> Option { - let mut length = 0; - loop { - let mut line = String::new(); - if stdin.read_line(&mut line).await.unwrap_or(0) == 0 { - return None; - } - if line.starts_with('{') { - return match serde_json::from_str::(line.trim_end()) { - Ok(req) => Some(req), - Err(e) => { - tracing::error!( - "Failed to parse JSON-RPC request from JSONL: {}. Payload: {}", - e, - line - ); - None - } - }; - } - - let line = line.trim_end(); - if line.is_empty() { - break; - } - let lower_line = line.to_lowercase(); - if let Some(len_str) = lower_line.strip_prefix("content-length:") { - length = len_str.trim().parse().unwrap_or(0); - } - } - if length == 0 { - return None; - } - let mut buffer = vec![0; length]; - stdin.read_exact(&mut buffer).await.unwrap_or(0); - - serde_json::from_slice(&buffer).ok() -} pub async fn send_response(response: JsonRpcResponse) { let msg = serde_json::to_string(&response).unwrap(); @@ -508,17 +468,25 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { tracing::info!("{} MCP server started", app_name); let mut stdin = tokio::io::BufReader::new(tokio::io::stdin()); loop { - let msg = match read_message(&mut stdin).await { - Some(m) => { - tracing::info!("Received message method: {}", m.method); - m - } + let raw_msg = match mcp_stdio::read_mcp_message(&mut stdin).await { + Some(m) => m, None => { tracing::info!("Stdin closed, exiting loop"); break; } }; + let msg = match serde_json::from_str::(&raw_msg) { + Ok(m) => { + tracing::info!("Received message method: {}", m.method); + m + }, + Err(e) => { + tracing::error!("Failed to parse JSON-RPC request from JSONL: {}. Payload: {}", e, raw_msg); + continue; + } + }; + let app_name = app_name.to_string(); let app_version = app_version.to_string(); @@ -915,3 +883,4 @@ mod tests { assert!(req.is_none()); } } + diff --git a/stub/Cargo.toml b/stub/Cargo.toml index 55339a0..e439649 100644 --- a/stub/Cargo.toml +++ b/stub/Cargo.toml @@ -16,6 +16,7 @@ tracing = "0.1.44" tracing-subscriber = "0.3.23" dirs = "7.0.0" serde_json = "1.0.151" +mcp-stdio = { version = "0.1.0", path = "../mcp-stdio" } diff --git a/stub/src/main.rs b/stub/src/main.rs index d4f9884..137999d 100644 --- a/stub/src/main.rs +++ b/stub/src/main.rs @@ -12,40 +12,7 @@ struct Cli { target: String, } -async fn read_mcp_message(stdin: &mut tokio::io::BufReader) -> Option { - use tokio::io::AsyncReadExt; - let mut length = 0; - loop { - let mut line = String::new(); - let bytes_read = stdin.read_line(&mut line).await.unwrap_or(0); - if bytes_read == 0 { - tracing::info!("stdin EOF reached"); - return None; - } - tracing::info!("Read {} bytes from stdin: {:?}", bytes_read, line); - if line.starts_with('{') { - return Some(line.trim_end().to_string()); - } - - let line = line.trim_end(); - if line.is_empty() { - break; - } - let lower_line = line.to_lowercase(); - if let Some(len_str) = lower_line.strip_prefix("content-length:") { - length = len_str.trim().parse().unwrap_or(0); - } - } - if length == 0 { - return None; - } - let mut buffer = vec![0; length]; - if stdin.read_exact(&mut buffer).await.is_err() { - return None; - } - String::from_utf8(buffer).ok() -} fn init_logging(app_name: &str) -> Option { let mut base_dir = dirs::home_dir().unwrap_or_else(|| std::path::PathBuf::from(".")); @@ -74,7 +41,7 @@ fn main() -> Result<(), Box> { tokio::spawn(async move { let mut stdin = tokio::io::BufReader::new(tokio::io::stdin()); - while let Some(msg) = read_mcp_message(&mut stdin).await { + while let Some(msg) = mcp_stdio::read_mcp_message(&mut stdin).await { let _ = msg_tx.send(msg).await; } let _ = shutdown_tx.send(()).await; @@ -190,3 +157,4 @@ fn main() -> Result<(), Box> { Ok(()) }) } +