From 61b03bc6e39f9936c4f9c80ae63d1b2d10879638 Mon Sep 17 00:00:00 2001 From: Riz Ashraf Date: Tue, 22 Sep 2026 22:26:00 +0100 Subject: [PATCH] perf(nvim-core): implement zero-copy BytesMut stream buffer and DashMap for lock contention --- nvim-core/Cargo.toml | 2 ++ nvim-core/src/lib.rs | 61 +++++++++++++++----------------------------- 2 files changed, 23 insertions(+), 40 deletions(-) diff --git a/nvim-core/Cargo.toml b/nvim-core/Cargo.toml index f2362d1..cc65fd5 100644 --- a/nvim-core/Cargo.toml +++ b/nvim-core/Cargo.toml @@ -14,4 +14,6 @@ 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" } +bytes = "1.12.1" +dashmap = "6.2.1" diff --git a/nvim-core/src/lib.rs b/nvim-core/src/lib.rs index d27d936..12f3693 100644 --- a/nvim-core/src/lib.rs +++ b/nvim-core/src/lib.rs @@ -114,7 +114,6 @@ async fn get_socket_path() -> Result { } Err("Could not find Neovim socket".to_string()) } -use std::collections::HashMap; use std::sync::Arc; use std::sync::LazyLock; use tokio::sync::{mpsc, oneshot}; @@ -162,8 +161,8 @@ async fn get_nvim_connection() -> Result, String> { let (mut read_half, mut write_half) = tokio::io::split(stream); let (tx, mut rx) = mpsc::channel::(32); type PendingRequestsMap = - Arc>>>>; - let pending_requests: PendingRequestsMap = Arc::new(std::sync::Mutex::new(HashMap::new())); + Arc>>>; + let pending_requests: PendingRequestsMap = Arc::new(dashmap::DashMap::new()); // Write task let pending_clone = Arc::clone(&pending_requests); @@ -175,10 +174,7 @@ async fn get_nvim_connection() -> Result, String> { continue; } - pending_clone - .lock() - .unwrap_or_else(|e| e.into_inner()) - .insert(req.msgid, req.reply); + pending_clone.insert(req.msgid, req.reply); if write_half.write_all(&buf).await.is_err() { tracing::error!("Failed to write to Neovim socket"); @@ -190,15 +186,15 @@ async fn get_nvim_connection() -> Result, String> { // Read task let pending_clone2 = Arc::clone(&pending_requests); tokio::spawn(async move { - let mut resp_buf = Vec::new(); - let mut chunk = vec![0u8; 65536]; - let mut offset = 0; + use bytes::{Buf, BytesMut}; + let mut resp_buf = BytesMut::with_capacity(65536); loop { - let mut cursor = std::io::Cursor::new(&resp_buf[offset..]); + let mut cursor = std::io::Cursor::new(&resp_buf[..]); match rmpv::decode::read_value(&mut cursor) { Ok(val) => { - offset += cursor.position() as usize; + let parsed_len = cursor.position() as usize; + resp_buf.advance(parsed_len); if let rmpv::Value::Array(ref arr) = val && arr.len() >= 4 @@ -209,21 +205,10 @@ async fn get_nvim_connection() -> Result, String> { _ => 0, }; - if let Some(reply_sender) = pending_clone2 - .lock() - .unwrap_or_else(|e| e.into_inner()) - .remove(&msgid) - { + if let Some((_, reply_sender)) = pending_clone2.remove(&msgid) { let _ = reply_sender.send(Ok(val)); } } - if offset == resp_buf.len() { - resp_buf.clear(); - offset = 0; - } else if offset > 1024 * 1024 { - resp_buf.drain(..offset); - offset = 0; - } continue; } Err(e) @@ -237,16 +222,13 @@ async fn get_nvim_connection() -> Result, String> { _ => false, } => { - resp_buf.drain(..offset); - offset = 0; - - let read_future = read_half.read(&mut chunk); - match tokio::time::timeout(tokio::time::Duration::from_secs(60), read_future) - .await + match tokio::time::timeout( + tokio::time::Duration::from_secs(60), + read_half.read_buf(&mut resp_buf), + ) + .await { - Ok(Ok(n)) if n > 0 => { - resp_buf.extend_from_slice(&chunk[..n]); - } + Ok(Ok(n)) if n > 0 => {} _ => { tracing::error!("Neovim socket read loop closed or timeout"); break; @@ -261,9 +243,11 @@ async fn get_nvim_connection() -> Result, String> { } // Cleanup pending requests on disconnect - let mut pending = pending_clone2.lock().unwrap_or_else(|e| e.into_inner()); - for (_, sender) in pending.drain() { - let _ = sender.send(Err("Connection closed".to_string())); + let keys: Vec<_> = pending_clone2.iter().map(|kv| *kv.key()).collect(); + for k in keys { + if let Some((_, sender)) = pending_clone2.remove(&k) { + let _ = sender.send(Err("Connection closed".to_string())); + } } }); @@ -276,10 +260,7 @@ async fn get_nvim_connection() -> Result, String> { if Arc::strong_count(&pending_clone3) <= 1 { break; // Socket closed and other tasks finished, no need to keep cleaning up } - pending_clone3 - .lock() - .unwrap_or_else(|e| e.into_inner()) - .retain(|_, sender| !sender.is_closed()); + pending_clone3.retain(|_, sender| !sender.is_closed()); } });