feat(nvim-core): implement bidirectional event streaming and real-time shadow buffer

- Created global thread-safe NvimState to shadow Neovim context.
- Injected Lua autocmds (CursorMoved, TextChanged) on connection to stream events back to MCP.
- Upgraded msgpack RPC loop to intercept and parse 'mcp_event' notifications alongside standard responses.
- Optimized nvim_get_cursor to pull instantly from the local shadow buffer, bypassing RPC round-trips.
This commit is contained in:
Riz Ashraf committed 2026-09-22 22:32:37 +01:00
1 parent 61b03bc6e3
commit ea735003e2
2 files changed
+113 -12

No files matched your search

Generated
+22
View File
@@ -497,6 +497,20 @@ dependencies = [
"syn 3.0.6", "syn 3.0.6",
] ]
[[package]]
name = "dashmap"
version = "6.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e6361d5c062261c78a176addb82d4c821ae42bed6089de0e12603cd25de2059c"
dependencies = [
"cfg-if",
"crossbeam-utils",
"hashbrown 0.14.5",
"lock_api",
"once_cell",
"parking_lot_core",
]
[[package]] [[package]]
name = "data-encoding" name = "data-encoding"
version = "2.11.1" version = "2.11.1"
@@ -807,6 +821,12 @@ dependencies = [
"rand_core 0.10.1", "rand_core 0.10.1",
] ]
[[package]]
name = "hashbrown"
version = "0.14.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1"
[[package]] [[package]]
name = "hashbrown" name = "hashbrown"
version = "0.16.1" version = "0.16.1"
@@ -1375,6 +1395,8 @@ dependencies = [
name = "nvim-core" name = "nvim-core"
version = "0.1.0" version = "0.1.0"
dependencies = [ dependencies = [
"bytes",
"dashmap",
"dirs", "dirs",
"mcp-stdio", "mcp-stdio",
"rmcp", "rmcp",
+91 -12
View File
@@ -129,6 +129,36 @@ static NEXT_MSGID: AtomicU64 = AtomicU64::new(1);
static NVIM_CONN: LazyLock<Arc<std::sync::Mutex<Option<mpsc::Sender<NvimRequest>>>>> = static NVIM_CONN: LazyLock<Arc<std::sync::Mutex<Option<mpsc::Sender<NvimRequest>>>>> =
LazyLock::new(|| Arc::new(std::sync::Mutex::new(None))); LazyLock::new(|| Arc::new(std::sync::Mutex::new(None)));
#[derive(Debug, Default, Clone)]
pub struct NvimState {
pub cursor: String,
pub active_buffer_id: u64,
}
static NVIM_STATE: LazyLock<Arc<std::sync::Mutex<NvimState>>> =
LazyLock::new(|| Arc::new(std::sync::Mutex::new(NvimState::default())));
fn handle_nvim_notification(params: &[rmpv::Value]) {
if params.is_empty() { return; }
if let rmpv::Value::String(event) = &params[0] {
match event.as_str().unwrap_or("") {
"CursorMoved" => {
if params.len() > 1
&& let rmpv::Value::Array(pos) = &params[1]
&& pos.len() >= 4
&& let (rmpv::Value::Integer(row), rmpv::Value::Integer(col)) = (&pos[1], &pos[2]) {
let mut state = NVIM_STATE.lock().unwrap_or_else(|e| e.into_inner());
state.cursor = format!("Line: {}, Column: {}", row, col);
}
},
"TextChanged" => {
tracing::debug!("Shadow buffer text changed event received");
}
_ => {}
}
}
}
async fn get_nvim_connection() -> Result<mpsc::Sender<NvimRequest>, String> { async fn get_nvim_connection() -> Result<mpsc::Sender<NvimRequest>, String> {
{ {
let conn_lock = NVIM_CONN.lock().unwrap_or_else(|e| e.into_inner()); let conn_lock = NVIM_CONN.lock().unwrap_or_else(|e| e.into_inner());
@@ -196,18 +226,21 @@ async fn get_nvim_connection() -> Result<mpsc::Sender<NvimRequest>, String> {
let parsed_len = cursor.position() as usize; let parsed_len = cursor.position() as usize;
resp_buf.advance(parsed_len); resp_buf.advance(parsed_len);
if let rmpv::Value::Array(ref arr) = val if let rmpv::Value::Array(ref arr) = val {
&& arr.len() >= 4 if arr.len() >= 4 && arr[0] == rmpv::Value::Integer(1.into()) {
&& arr[0] == rmpv::Value::Integer(1.into()) let msgid = match &arr[1] {
{ rmpv::Value::Integer(i) => i.as_u64().unwrap_or(0),
let msgid = match &arr[1] { _ => 0,
rmpv::Value::Integer(i) => i.as_u64().unwrap_or(0), };
_ => 0, if let Some((_, reply_sender)) = pending_clone2.remove(&msgid) {
}; let _ = reply_sender.send(Ok(val));
}
if let Some((_, reply_sender)) = pending_clone2.remove(&msgid) { } else if arr.len() >= 3 && arr[0] == rmpv::Value::Integer(2.into())
let _ = reply_sender.send(Ok(val)); && let rmpv::Value::String(method) = &arr[1]
} && method.as_str().unwrap_or("") == "mcp_event"
&& let rmpv::Value::Array(params) = &arr[2] {
handle_nvim_notification(params);
}
} }
continue; continue;
} }
@@ -272,6 +305,45 @@ async fn get_nvim_connection() -> Result<mpsc::Sender<NvimRequest>, String> {
return Ok(existing_sender.clone()); return Ok(existing_sender.clone());
} }
*conn_lock = Some(tx.clone()); *conn_lock = Some(tx.clone());
drop(conn_lock);
let tx_clone = tx.clone();
tokio::spawn(async move {
let setup_code = r#"
local channel = vim.api.nvim_get_api_info()[1]
vim.api.nvim_create_augroup("MCP_Tracking", { clear = true })
vim.api.nvim_create_autocmd({"CursorMoved", "CursorMovedI"}, {
group = "MCP_Tracking",
callback = function()
pcall(vim.rpcnotify, channel, "mcp_event", "CursorMoved", vim.fn.getpos('.'))
end
})
vim.api.nvim_create_autocmd({"TextChanged", "TextChangedI", "BufEnter"}, {
group = "MCP_Tracking",
callback = function()
pcall(vim.rpcnotify, channel, "mcp_event", "TextChanged", vim.api.nvim_get_current_buf())
end
})
"#;
let msgid = NEXT_MSGID.fetch_add(1, Ordering::SeqCst);
let req = rmpv::Value::Array(vec![
rmpv::Value::Integer(0.into()),
rmpv::Value::Integer(msgid.into()),
rmpv::Value::String("nvim_exec_lua".into()),
rmpv::Value::Array(vec![
rmpv::Value::String(setup_code.into()),
rmpv::Value::Array(vec![]),
]),
]);
let (reply_tx, _reply_rx) = tokio::sync::oneshot::channel();
let _ = tx_clone.send(NvimRequest {
msgid,
req,
reply: reply_tx,
}).await;
tracing::info!("Injected bidirectional event tracking autocmds into Neovim");
});
Ok(tx) Ok(tx)
} }
@@ -363,6 +435,13 @@ async fn get_nvim_active_buffer() -> Result<String, String> {
} }
async fn get_nvim_cursor() -> Result<String, String> { async fn get_nvim_cursor() -> Result<String, String> {
{
let state = NVIM_STATE.lock().unwrap_or_else(|e| e.into_inner());
if !state.cursor.is_empty() {
return Ok(state.cursor.clone());
}
}
let result = let result =
call_nvim_method("nvim_win_get_cursor", vec![rmpv::Value::Integer(0.into())]).await?; call_nvim_method("nvim_win_get_cursor", vec![rmpv::Value::Integer(0.into())]).await?;