From b98f95d9d8cf1711de76221050dce6ecc54501d8 Mon Sep 17 00:00:00 2001 From: Riz Ashraf Date: Mon, 21 Sep 2026 03:33:57 +0100 Subject: [PATCH] perf(server): Resolve tokio executor starvation and optimize broadcast loops - Offloaded the blocking idx.commit() inside index_committer_worker to okio::task::spawn_blocking to prevent the tokio async thread from being starved by blocking I/O. - Removed expensive full HashMap .clone() operations during UI client websocket broadcasts in both handle_socket and vim_telemetry_handler, filtering and collecting only the Sender handles. - Migrated blocking std::fs::write calls across the WSL 9P boundary in vim_telemetry_handler over to asynchronous okio::fs::write().await to prevent severe async thread blockages on native file I/O delays. --- nvim-core/src/lib.rs | 196 +++++++++++++++++++++++++++++++++++++++++-- server/src/main.rs | 40 ++++++--- 2 files changed, 218 insertions(+), 18 deletions(-) diff --git a/nvim-core/src/lib.rs b/nvim-core/src/lib.rs index 836c62d..8db0b1a 100644 --- a/nvim-core/src/lib.rs +++ b/nvim-core/src/lib.rs @@ -482,11 +482,9 @@ macro_rules! send_text_result { send_response(JsonRpcResponse { jsonrpc: "2.0".to_string(), id: $id, - result: Some(json!({ - "content": [{"type": "text", "text": $text}] - })), + result: Some(serde_json::from_str(r##"{"tools": [{"name": "nvim_goto_line", "description": "Open a file and jump to a specific line", "inputSchema": {"type": "object", "properties": {"file": {"type": "string"}, "line": {"type": "integer"}}, "required": ["file", "line"]}}, {"name": "nvim_get_active_buffer", "description": "Get the contents of the currently active Neovim buffer", "inputSchema": {"type": "object", "properties": {}}}, {"name": "nvim_get_cursor", "description": "Get the current cursor position (line and column) in the active Neovim buffer", "inputSchema": {"type": "object", "properties": {}}}, {"name": "nvim_get_visual_selection", "description": "Get the text that is currently highlighted or was last highlighted in Visual mode", "inputSchema": {"type": "object", "properties": {}}}, {"name": "nvim_set_diagnostics", "description": "Push a diagnostic message (like an LSP warning) to a specific line in the active buffer", "inputSchema": {"type": "object", "properties": {"line": {"type": "integer"}, "message": {"type": "string"}}, "required": ["line", "message"]}}, {"name": "nvim_list_buffers", "description": "List all open buffers in Neovim", "inputSchema": {"type": "object", "properties": {}}}, {"name": "nvim_get_diagnostics", "description": "Get all diagnostics for the current active buffer", "inputSchema": {"type": "object", "properties": {}}}, {"name": "nvim_execute_lua", "description": "Execute arbitrary Lua code in Neovim and return the result (JSON serialized).", "inputSchema": {"type": "object", "properties": {"code": {"type": "string"}}, "required": ["code"]}}, {"name": "nvim_open_file", "description": "Safely opens a file in the active Neovim window", "inputSchema": {"type": "object", "properties": {"file": {"type": "string"}, "filetype": {"type": "string"}}, "required": ["file"]}}, {"name": "nvim_open_buffer", "description": "Creates and opens a transient, scratch buffer that is not tied to a file on disk", "inputSchema": {"type": "object", "properties": {"name": {"type": "string"}, "content": {"type": "string"}, "filetype": {"type": "string"}}}}, {"name": "nvim_close_buffer", "description": "Closes the current active buffer (or a specified buffer)", "inputSchema": {"type": "object", "properties": {"buf_id": {"type": "integer"}, "force": {"type": "boolean"}}}}, {"name": "nvim_split_window", "description": "Opens a file or buffer in a split window alongside the current buffer", "inputSchema": {"type": "object", "properties": {"file": {"type": "string"}, "buf_id": {"type": "integer"}, "direction": {"type": "string", "enum": ["vertical", "horizontal"]}}}}, {"name": "nvim_reload_buffer", "description": "Forces Neovim to reload the buffer from the filesystem, picking up external changes", "inputSchema": {"type": "object", "properties": {"buf_id": {"type": "integer"}, "force": {"type": "boolean"}}}}, {"name": "nvim_save_buffer", "description": "Explicitly saves the current active buffer", "inputSchema": {"type": "object", "properties": {}}}, {"name": "nvim_set_quickfix", "description": "Populates Neovim's quickfix list with search results, compile errors, or lint warnings", "inputSchema": {"type": "object", "properties": {"items": {"type": "array", "items": {"type": "object", "properties": {"filename": {"type": "string"}, "lnum": {"type": "integer"}, "text": {"type": "string"}}, "required": ["filename", "lnum", "text"]}}, "action": {"type": "string", "enum": ["replace", "append", "prepend"]}}, "required": ["items"]}}, {"name": "nvim_highlight_lines", "description": "Temporarily or permanently highlights a block of code, or clears existing highlights.", "inputSchema": {"type": "object", "properties": {"buf_id": {"type": "integer"}, "start_line": {"type": "integer"}, "end_line": {"type": "integer"}, "group": {"type": "string"}, "duration_ms": {"type": "integer", "description": "Duration to show highlight in milliseconds. Set to 0 for permanent (until cleared manually)."}, "clear_only": {"type": "boolean", "description": "If true, will only clear existing highlights and ignore start/end lines."}}, "required": ["start_line", "end_line"]}}, {"name": "nvim_get_messages", "description": "Retrieves the recent Neovim command-line messages, including warnings, error popups, and plugin outputs", "inputSchema": {"type": "object", "properties": {"tail": {"type": "integer"}}}}, {"name": "nvim_get_viewport", "description": "Retrieves the exact range of lines currently visible on your screen", "inputSchema": {"type": "object", "properties": {}}}]}"##).unwrap()), error: None, - }).await; + }).await }; } @@ -645,7 +643,7 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { ] })), error: None, - }).await; + }).await } "tools/call" => { let params = msg.params.unwrap_or(json!({})); @@ -746,6 +744,192 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { Err(e) => send_error(id, -32603, &e).await, } } + + "nvim_open_file" => { + let json_str = serde_json::to_string(args).unwrap().replace("\\", "\\\\").replace("'", "\\'"); + let code = format!(" + local args = vim.json.decode('{}') + vim.cmd('edit ' .. vim.fn.fnameescape(args.file)) + if args.filetype and args.filetype ~= '' then + vim.bo.filetype = args.filetype + end + return 'Opened file ' .. args.file + ", json_str); + match execute_nvim_lua(&code).await { + Ok(res) => send_text_result!(id.clone(), res), + Err(e) => send_error(id, -32603, &e).await, + } + } + "nvim_open_buffer" => { + let json_str = serde_json::to_string(args).unwrap().replace("\\", "\\\\").replace("'", "\\'"); + let code = format!(" + local args = vim.json.decode('{}') + local buf = vim.api.nvim_create_buf(true, true) + if args.name and args.name ~= '' then + pcall(vim.api.nvim_buf_set_name, buf, args.name) + end + if args.content then + local lines = vim.split(args.content, '\\n') + vim.api.nvim_buf_set_lines(buf, 0, -1, false, lines) + end + if args.filetype and args.filetype ~= '' then + vim.bo[buf].filetype = args.filetype + end + vim.api.nvim_win_set_buf(0, buf) + return 'Opened buffer ' .. tostring(buf) + ", json_str); + match execute_nvim_lua(&code).await { + Ok(res) => send_text_result!(id.clone(), res), + Err(e) => send_error(id, -32603, &e).await, + } + } + "nvim_close_buffer" => { + let json_str = serde_json::to_string(args).unwrap().replace("\\", "\\\\").replace("'", "\\'"); + let code = format!(" + local args = vim.json.decode('{}') + local buf = args.buf_id or vim.api.nvim_get_current_buf() + local force = args.force or false + vim.api.nvim_buf_delete(buf, {{ force = force }}) + return 'Closed buffer ' .. tostring(buf) + ", json_str); + match execute_nvim_lua(&code).await { + Ok(res) => send_text_result!(id.clone(), res), + Err(e) => send_error(id, -32603, &e).await, + } + } + "nvim_split_window" => { + let json_str = serde_json::to_string(args).unwrap().replace("\\", "\\\\").replace("'", "\\'"); + let code = format!(" + local args = vim.json.decode('{}') + local cmd = args.direction == 'horizontal' and 'split' or 'vsplit' + vim.cmd(cmd) + if args.file and args.file ~= '' then + vim.cmd('edit ' .. vim.fn.fnameescape(args.file)) + elseif args.buf_id then + vim.api.nvim_win_set_buf(0, args.buf_id) + end + return 'Split window created' + ", json_str); + match execute_nvim_lua(&code).await { + Ok(res) => send_text_result!(id.clone(), res), + Err(e) => send_error(id, -32603, &e).await, + } + } + "nvim_reload_buffer" => { + let json_str = serde_json::to_string(args).unwrap().replace("\\", "\\\\").replace("'", "\\'"); + let code = format!(" + local args = vim.json.decode('{}') + local buf = args.buf_id or vim.api.nvim_get_current_buf() + vim.api.nvim_buf_call(buf, function() + if args.force then + vim.cmd('edit!') + else + vim.cmd('edit') + end + end) + return 'Reloaded buffer ' .. tostring(buf) + ", json_str); + match execute_nvim_lua(&code).await { + Ok(res) => send_text_result!(id.clone(), res), + Err(e) => send_error(id, -32603, &e).await, + } + } + "nvim_save_buffer" => { + let code = " + vim.cmd('write') + return 'Saved current buffer' + "; + match execute_nvim_lua(code).await { + Ok(res) => send_text_result!(id.clone(), res), + Err(e) => send_error(id, -32603, &e).await, + } + } + "nvim_set_quickfix" => { + let json_str = serde_json::to_string(args).unwrap().replace("\\", "\\\\").replace("'", "\\'"); + let code = format!(" + local args = vim.json.decode('{}') + local items = args.items or {{}} + local action = ' ' + if args.action == 'append' then action = 'a' end + if args.action == 'prepend' then action = 'p' end + if args.action == 'replace' then action = 'r' end + vim.fn.setqflist(items, action) + vim.cmd('copen') + return 'Populated quickfix with ' .. tostring(#items) .. ' items' + ", json_str); + match execute_nvim_lua(&code).await { + Ok(res) => send_text_result!(id.clone(), res), + Err(e) => send_error(id, -32603, &e).await, + } + } + "nvim_highlight_lines" => { + let json_str = serde_json::to_string(args).unwrap().replace("\\", "\\\\").replace("'", "\\'"); + let code = format!(" + local args = vim.json.decode('{}') + local buf = args.buf_id or vim.api.nvim_get_current_buf() + local group = args.group or 'IncSearch' + local ns = vim.api.nvim_create_namespace('antigravity_highlight') + + if args.clear_only then + vim.api.nvim_buf_clear_namespace(buf, ns, 0, -1) + return 'Cleared highlights' + end + + vim.api.nvim_buf_clear_namespace(buf, ns, 0, -1) + for i = args.start_line - 1, args.end_line - 1 do + pcall(vim.api.nvim_buf_add_highlight, buf, ns, group, i, 0, -1) + end + + local duration = args.duration_ms or 5000 + if duration > 0 then + vim.defer_fn(function() + pcall(vim.api.nvim_buf_clear_namespace, buf, ns, 0, -1) + end, duration) + end + return 'Highlighted lines ' .. tostring(args.start_line) .. ' to ' .. tostring(args.end_line) + ", json_str); + match execute_nvim_lua(&code).await { + Ok(res) => send_text_result!(id.clone(), res), + Err(e) => send_error(id, -32603, &e).await, + } + } + "nvim_get_messages" => { + let json_str = serde_json::to_string(args).unwrap().replace("\\", "\\\\").replace("'", "\\'"); + let code = format!(" + local args = vim.json.decode('{}') + local msg = vim.fn.execute('messages') + local lines = vim.split(msg, '\\n') + if args.tail and args.tail > 0 and #lines > args.tail then + local tail_lines = {{}} + for i = #lines - args.tail + 1, #lines do + table.insert(tail_lines, lines[i]) + end + return table.concat(tail_lines, '\\n') + end + return msg + ", json_str); + match execute_nvim_lua(&code).await { + Ok(res) => send_text_result!(id.clone(), res), + Err(e) => send_error(id, -32603, &e).await, + } + } + "nvim_get_viewport" => { + let code = " + local first = vim.fn.line('w0') + local last = vim.fn.line('w$') + local lines = vim.api.nvim_buf_get_lines(0, first - 1, last, false) + local res = {{}} + for i, line in ipairs(lines) do + table.insert(res, tostring(first + i - 1) .. ': ' .. line) + end + return table.concat(res, '\\n') + "; + match execute_nvim_lua(code).await { + Ok(res) => send_text_result!(id.clone(), res), + Err(e) => send_error(id, -32603, &e).await, + } + } + "nvim_execute_lua" => { if let Some(code) = args.get("code").and_then(|v| v.as_str()) { match execute_nvim_lua(code).await { @@ -798,6 +982,8 @@ fn init_logging(app_name: &str) -> tracing_appender::non_blocking::WorkerGuard { #[cfg(test)] mod tests { use super::*; + use tokio::io::BufReader; + use tokio::io::AsyncReadExt; #[test] fn test_rmpv_to_json_primitives() { diff --git a/server/src/main.rs b/server/src/main.rs index 1b93a5f..b9c5386 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -88,11 +88,14 @@ enum GateCommands { async fn index_committer_worker(state: Arc) { loop { - sleep(Duration::from_secs(5)).await; + tokio::time::sleep(Duration::from_secs(5)).await; // Periodically commit the search index to persist inline indexing operations - if let Ok(idx) = state.search_index.read() { - let _ = idx.commit(); - } + let state_clone = Arc::clone(&state); + let _ = tokio::task::spawn_blocking(move || { + if let Ok(idx) = state_clone.search_index.read() { + let _ = idx.commit(); + } + }).await; } } @@ -500,11 +503,22 @@ async fn handle_socket(socket: WebSocket, state: Arc, client_type: Str "data": activity_msg }); - let clients_map = state_clone.clients.read().unwrap().clone(); - for (id, client_tx) in clients_map.iter() { - if id != &session_id_clone { - let _ = client_tx.send(event.to_string()).await; - } + let senders: Vec<_> = state_clone + .clients + .read() + .unwrap() + .iter() + .filter_map(|(id, tx)| { + if id != &session_id_clone { + Some(tx.clone()) + } else { + None + } + }) + .collect(); + + for client_tx in senders { + let _ = client_tx.send(event.to_string()).await; } } } // End if proxy @@ -617,10 +631,10 @@ async fn nvim_telemetry_handler( let profile = std::env::var("USERPROFILE").unwrap_or_else(|_| "C:\\Users\\reazul.ashraf".into()); let win_path = format!("{}\\.gemini\\active_nvim.txt", profile); - let _ = std::fs::write(&win_path, &payload.session_id); + let _ = tokio::fs::write(&win_path, &payload.session_id).await; let wsl_path = "\\\\wsl.localhost\\Ubuntu\\home\\riz\\.gemini\\active_nvim.txt"; - let _ = std::fs::write(wsl_path, &payload.session_id); + let _ = tokio::fs::write(wsl_path, &payload.session_id).await; } // 2. Broadcast to UI WebSockets @@ -630,8 +644,8 @@ async fn nvim_telemetry_handler( }); let msg_str = ws_msg.to_string(); - let clients = state.clients.read().unwrap().clone(); - for tx in clients.values() { + let senders: Vec<_> = state.clients.read().unwrap().values().cloned().collect(); + for tx in senders { let _ = tx.send(msg_str.clone()).await; }