From 478655698e742fe6c91a60c49b00b8b2125eecab Mon Sep 17 00:00:00 2001 From: Riz Ashraf Date: Tue, 22 Sep 2026 03:19:51 +0100 Subject: [PATCH] fix(server): safely drop detached tokio join handles and refactor timestamp boilerplate --- mcp-stdio/src/lib.rs | 1 - nvim-core/src/lib.rs | 131 ++++++++++++------ server/build.rs | 3 +- server/src/handlers_v2/env.rs | 25 ++-- server/src/handlers_v2/graph.rs | 70 ++++++---- server/src/handlers_v2/meta.rs | 190 +++++++++++++++++---------- server/src/handlers_v2/notes.rs | 47 +++---- server/src/handlers_v2/tasks.rs | 143 +++++++++++--------- server/src/handlers_v2/utils.rs | 6 + server/src/handlers_v2/workspaces.rs | 93 ++++++------- server/src/main.rs | 43 +++--- server/src/mcp.rs | 1 - server/src/search.rs | 4 +- server/src/state.rs | 10 +- stub/src/main.rs | 2 - 15 files changed, 457 insertions(+), 312 deletions(-) diff --git a/mcp-stdio/src/lib.rs b/mcp-stdio/src/lib.rs index 00b59f7..77d968e 100644 --- a/mcp-stdio/src/lib.rs +++ b/mcp-stdio/src/lib.rs @@ -38,4 +38,3 @@ pub async fn read_mcp_message( String::from_utf8(buffer).ok() } - diff --git a/nvim-core/src/lib.rs b/nvim-core/src/lib.rs index 56208c3..f5257f6 100644 --- a/nvim-core/src/lib.rs +++ b/nvim-core/src/lib.rs @@ -199,8 +199,10 @@ async fn get_nvim_connection() -> Result, String> { let msgid = &arr[1]; let msgid_str = format!("{msgid:?}"); - if let Some(reply_sender) = - pending_clone2.lock().unwrap_or_else(|e| e.into_inner()).remove(&msgid_str) + if let Some(reply_sender) = pending_clone2 + .lock() + .unwrap_or_else(|e| e.into_inner()) + .remove(&msgid_str) { let _ = reply_sender.send(Ok(val)); } @@ -555,7 +557,9 @@ 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 raw_msg = if let Some(m) = mcp_stdio::read_mcp_message(&mut stdin).await { m } else { + let raw_msg = if let Some(m) = mcp_stdio::read_mcp_message(&mut stdin).await { + m + } else { tracing::info!("Stdin closed, exiting loop"); break; }; @@ -595,7 +599,6 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { }; match msg.method.as_str() { - "initialize" => { let init = rmcp::model::InitializeResult::new( rmcp::model::ServerCapabilities::builder() @@ -719,7 +722,13 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { let cmd = format!("e {escaped_file} | {line} | normal! zz"); match send_nvim_command(&cmd).await { Ok(()) => { - send_text_result!(id.clone(), format!("Successfully jumped to {} line {}", file, line)); + send_text_result!( + id.clone(), + format!( + "Successfully jumped to {} line {}", + file, line + ) + ); } Err(e) => send_error(id, -32603, &e).await, } @@ -752,7 +761,10 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { ) { match set_nvim_diagnostics(line, message).await { Ok(()) => { - send_text_result!(id.clone(), format!("Successfully set diagnostic on line {}", line)); + send_text_result!( + id.clone(), + format!("Successfully set diagnostic on line {}", line) + ); } Err(e) => send_error(id, -32603, &e).await, } @@ -804,23 +816,32 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { } "nvim_open_file" => { - let json_str = serde_json::to_string(args).unwrap_or_else(|_| "{}".to_string()).replace('\\', "\\\\").replace('\'', "\\'"); - let code = format!(" + let json_str = serde_json::to_string(args) + .unwrap_or_else(|_| "{}".to_string()) + .replace('\\', "\\\\") + .replace('\'', "\\'"); + let code = format!( + " local args = vim.json.decode('{json_str}') 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 - "); + " + ); 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_or_else(|_| "{}".to_string()).replace('\\', "\\\\").replace('\'', "\\'"); - let code = format!(" + let json_str = serde_json::to_string(args) + .unwrap_or_else(|_| "{}".to_string()) + .replace('\\', "\\\\") + .replace('\'', "\\'"); + let code = format!( + " local args = vim.json.decode('{json_str}') local buf = vim.api.nvim_create_buf(true, true) if args.name and args.name ~= '' then @@ -835,29 +856,39 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { end vim.api.nvim_win_set_buf(0, buf) return 'Opened buffer ' .. tostring(buf) - "); + " + ); 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_or_else(|_| "{}".to_string()).replace('\\', "\\\\").replace('\'', "\\'"); - let code = format!(" + let json_str = serde_json::to_string(args) + .unwrap_or_else(|_| "{}".to_string()) + .replace('\\', "\\\\") + .replace('\'', "\\'"); + let code = format!( + " local args = vim.json.decode('{json_str}') 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) - "); + " + ); 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_or_else(|_| "{}".to_string()).replace('\\', "\\\\").replace('\'', "\\'"); - let code = format!(" + let json_str = serde_json::to_string(args) + .unwrap_or_else(|_| "{}".to_string()) + .replace('\\', "\\\\") + .replace('\'', "\\'"); + let code = format!( + " local args = vim.json.decode('{json_str}') local cmd = args.direction == 'horizontal' and 'split' or 'vsplit' vim.cmd(cmd) @@ -867,15 +898,20 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { vim.api.nvim_win_set_buf(0, args.buf_id) end return 'Split window created' - "); + " + ); 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_or_else(|_| "{}".to_string()).replace('\\', "\\\\").replace('\'', "\\'"); - let code = format!(" + let json_str = serde_json::to_string(args) + .unwrap_or_else(|_| "{}".to_string()) + .replace('\\', "\\\\") + .replace('\'', "\\'"); + let code = format!( + " local args = vim.json.decode('{json_str}') local buf = args.buf_id or vim.api.nvim_get_current_buf() vim.api.nvim_buf_call(buf, function() @@ -886,7 +922,8 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { end end) return 'Reloaded buffer ' .. tostring(buf) - "); + " + ); match execute_nvim_lua(&code).await { Ok(res) => send_text_result!(id.clone(), res), Err(e) => send_error(id, -32603, &e).await, @@ -903,8 +940,12 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { } } "nvim_set_quickfix" => { - let json_str = serde_json::to_string(args).unwrap_or_else(|_| "{}".to_string()).replace('\\', "\\\\").replace('\'', "\\'"); - let code = format!(" + let json_str = serde_json::to_string(args) + .unwrap_or_else(|_| "{}".to_string()) + .replace('\\', "\\\\") + .replace('\'', "\\'"); + let code = format!( + " local args = vim.json.decode('{json_str}') local items = args.items or {{}} local action = ' ' @@ -914,14 +955,18 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { vim.fn.setqflist(items, action) vim.cmd('copen') return 'Populated quickfix with ' .. tostring(#items) .. ' items' - "); + " + ); 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_or_else(|_| "{}".to_string()).replace('\\', "\\\\").replace('\'', "\\'"); + let json_str = serde_json::to_string(args) + .unwrap_or_else(|_| "{}".to_string()) + .replace('\\', "\\\\") + .replace('\'', "\\'"); let code = format!(" local args = vim.json.decode('{json_str}') local buf = args.buf_id or vim.api.nvim_get_current_buf() @@ -952,8 +997,12 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { } } "nvim_get_messages" => { - let json_str = serde_json::to_string(args).unwrap_or_else(|_| "{}".to_string()).replace('\\', "\\\\").replace('\'', "\\'"); - let code = format!(" + let json_str = serde_json::to_string(args) + .unwrap_or_else(|_| "{}".to_string()) + .replace('\\', "\\\\") + .replace('\'', "\\'"); + let code = format!( + " local args = vim.json.decode('{json_str}') local msg = vim.fn.execute('messages') local lines = vim.split(msg, '\\n') @@ -965,7 +1014,8 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { return table.concat(tail_lines, '\\n') end return msg - "); + " + ); match execute_nvim_lua(&code).await { Ok(res) => send_text_result!(id.clone(), res), Err(e) => send_error(id, -32603, &e).await, @@ -992,14 +1042,26 @@ pub async fn run_mcp_loop(app_name: &str, app_version: &str) { if let Some(code) = args.get("code").and_then(|v| v.as_str()) { // BAKE IN: Block interactive prompts that cause server deadlocks let lower_code = code.to_lowercase(); - if lower_code.contains("vim.fn.input") || lower_code.contains("vim.ui.select") || lower_code.contains("vim.fn.confirm") || lower_code.contains("vim.ui.input") { + if lower_code.contains("vim.fn.input") + || lower_code.contains("vim.ui.select") + || lower_code.contains("vim.fn.confirm") + || lower_code.contains("vim.ui.input") + { send_error(id, -32600, "CRITICAL ERROR: Interactive Lua functions (input, select, confirm) are strictly forbidden via MCP as they will hang the headless server.").await; return; } // Block the 'c' confirm flag in vim.cmd substitutions - if (code.contains("vim.cmd") || code.contains("vim.api.nvim_command")) && code.contains("%s") && (code.contains("gc'") || code.contains("gc\"") || code.contains("gc\n") || code.contains("c'") || code.contains("c\"")) { - send_error(id, -32600, "CRITICAL ERROR: The 'c' (confirm) flag in Neovim substitutions is strictly forbidden via MCP as it triggers an interactive prompt that hangs the headless server. Use '/g' or '/ge' instead.").await; - return; + if (code.contains("vim.cmd") + || code.contains("vim.api.nvim_command")) + && code.contains("%s") + && (code.contains("gc'") + || code.contains("gc\"") + || code.contains("gc\n") + || code.contains("c'") + || code.contains("c\"")) + { + send_error(id, -32600, "CRITICAL ERROR: The 'c' (confirm) flag in Neovim substitutions is strictly forbidden via MCP as it triggers an interactive prompt that hangs the headless server. Use '/g' or '/ge' instead.").await; + return; } match execute_nvim_lua(code).await { @@ -1116,8 +1178,3 @@ mod tests { assert!(req.is_none()); } } - - - - - diff --git a/server/build.rs b/server/build.rs index 214fcb0..e799f94 100644 --- a/server/build.rs +++ b/server/build.rs @@ -12,7 +12,8 @@ fn main() { let is_dirty = Command::new("git") .args(["status", "--porcelain"]) - .output().is_ok_and(|out| !out.stdout.is_empty()); + .output() + .is_ok_and(|out| !out.stdout.is_empty()); if is_dirty { git_hash.push_str("-dirty"); diff --git a/server/src/handlers_v2/env.rs b/server/src/handlers_v2/env.rs index 2381908..6ef407d 100644 --- a/server/src/handlers_v2/env.rs +++ b/server/src/handlers_v2/env.rs @@ -4,7 +4,6 @@ use crate::tools::*; use async_trait::async_trait; use serde_json::Value; use std::sync::Arc; -use std::time::{SystemTime, UNIX_EPOCH}; pub struct UpdateEnvFingerprintHandler; @@ -32,10 +31,7 @@ impl McpTool for UpdateEnvFingerprintHandler { os: std::env::consts::OS.to_string(), shell: std::env::var("SHELL").unwrap_or_else(|_| "unknown".to_string()), tool_versions: req.tool_versions, - updated_at: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + updated_at: crate::handlers_v2::utils::now_secs(), }, ); }); @@ -61,9 +57,9 @@ impl McpTool for ReadEnvFingerprintHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: ReadEnvFingerprintTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let data = state.env_fingerprints.read_with(|fps| { - fps.get(&req.namespace).cloned() - }); + let data = state + .env_fingerprints + .read_with(|fps| fps.get(&req.namespace).cloned()); if let Some(fp) = data { let data = serde_json::to_string(&fp).unwrap_or_default(); Ok(data.to_string()) @@ -129,10 +125,7 @@ impl McpTool for RegisterEnvironmentHandler { url: req.url, description: req.description, requires_vpn: req.requires_vpn, - updated_at: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + updated_at: crate::handlers_v2::utils::now_secs(), }); }); Ok("Environment registered".to_string()) @@ -158,12 +151,12 @@ impl McpTool for GetEnvironmentDetailsHandler { let req: GetEnvironmentDetailsTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let data = state.environments.read_with(|envs| { - let filtered: Vec<_> = envs.iter().filter(|e| e.namespace == req.namespace).collect(); + let filtered: Vec<_> = envs + .iter() + .filter(|e| e.namespace == req.namespace) + .collect(); serde_json::to_string(&filtered).unwrap_or_default() }); Ok(data.to_string()) } } - - - diff --git a/server/src/handlers_v2/graph.rs b/server/src/handlers_v2/graph.rs index f3395bd..6fdb8f8 100644 --- a/server/src/handlers_v2/graph.rs +++ b/server/src/handlers_v2/graph.rs @@ -6,14 +6,12 @@ use serde_json::Value; use std::collections::HashSet; use std::sync::Arc; -#[derive(serde::Serialize)] -#[derive(Default)] +#[derive(serde::Serialize, Default)] struct BorrowedGraph<'a> { entities: std::collections::HashMap<&'a String, &'a crate::models::Entity>, relations: Vec<&'a crate::models::Relation>, } - pub struct QueryGraphPathHandler; #[async_trait] @@ -31,13 +29,13 @@ impl McpTool for QueryGraphPathHandler { serde_json::from_value(args).map_err(|e| e.to_string())?; state.read_graph(|graph| { let max_depth = req.max_depth.unwrap_or(5); - let mut queue = std::collections::VecDeque::new(); - let mut visited = std::collections::HashSet::new(); - let mut parents: std::collections::HashMap = + let mut queue: std::collections::VecDeque<&str> = std::collections::VecDeque::new(); + let mut visited: std::collections::HashSet<&str> = std::collections::HashSet::new(); + let mut parents: std::collections::HashMap<&str, (&str, std::borrow::Cow<'_, str>)> = std::collections::HashMap::new(); - queue.push_back(req.start_node.clone()); - visited.insert(req.start_node.clone()); + queue.push_back(req.start_node.as_str()); + visited.insert(req.start_node.as_str()); let mut found = false; let mut current_depth = 0; @@ -52,21 +50,30 @@ impl McpTool for QueryGraphPathHandler { nodes_at_current_depth -= 1; if current_depth < max_depth { for rel in &graph.relations { - if rel.from == current && !visited.contains(&rel.to) { - visited.insert(rel.to.clone()); + if rel.from == current && !visited.contains(rel.to.as_str()) { + visited.insert(rel.to.as_str()); parents.insert( - rel.to.clone(), - (current.clone(), rel.relation_type.clone()), + rel.to.as_str(), + ( + current, + std::borrow::Cow::Borrowed(rel.relation_type.as_str()), + ), ); - queue.push_back(rel.to.clone()); + queue.push_back(rel.to.as_str()); nodes_at_next_depth += 1; - } else if rel.to == current && !visited.contains(&rel.from) { - visited.insert(rel.from.clone()); + } else if rel.to == current && !visited.contains(rel.from.as_str()) { + visited.insert(rel.from.as_str()); parents.insert( - rel.from.clone(), - (current.clone(), format!("inverse({})", rel.relation_type)), + rel.from.as_str(), + ( + current, + std::borrow::Cow::Owned(format!( + "inverse({})", + rel.relation_type + )), + ), ); - queue.push_back(rel.from.clone()); + queue.push_back(rel.from.as_str()); nodes_at_next_depth += 1; } } @@ -80,11 +87,11 @@ impl McpTool for QueryGraphPathHandler { if found { let mut path = Vec::new(); - let mut curr = req.end_node.clone(); + let mut curr = req.end_node.as_str(); while curr != req.start_node { if let Some((parent, rel_type)) = parents.get(&curr) { path.push(format!("{} -[{}]-> {}", parent, rel_type, curr)); - curr = parent.clone(); + curr = parent; } else { break; } @@ -124,9 +131,13 @@ impl McpTool for CreateEntitiesHandler { } } }); - let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); + let idx = state + .search_index + .read() + .unwrap_or_else(|e| e.into_inner()) + .clone(); for entity in inserted { - let _ = idx.index_entity(&entity); + drop(idx.index_entity(&entity)); } Ok("Entities created".to_string()) } @@ -205,10 +216,14 @@ impl McpTool for DeleteEntitiesHandler { .relations .retain(|r| !to_delete.contains(&r.from) && !to_delete.contains(&r.to)); }); - - let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); + + let idx = state + .search_index + .read() + .unwrap_or_else(|e| e.into_inner()) + .clone(); for name in to_delete { - let _ = idx.delete_document(&name); + drop(idx.delete_document(&name)); } Ok("Entities deleted".to_string()) } @@ -416,7 +431,10 @@ impl McpTool for VisualizeGraphHandler { { continue; } - if query.is_empty() || included.contains(r.from.as_str()) || included.contains(r.to.as_str()) { + if query.is_empty() + || included.contains(r.from.as_str()) + || included.contains(r.to.as_str()) + { included.insert(r.from.as_str()); included.insert(r.to.as_str()); to_draw.push(r); diff --git a/server/src/handlers_v2/meta.rs b/server/src/handlers_v2/meta.rs index a7d9e8e..85f158b 100644 --- a/server/src/handlers_v2/meta.rs +++ b/server/src/handlers_v2/meta.rs @@ -5,7 +5,6 @@ use crate::tools::*; use async_trait::async_trait; use serde_json::Value; use std::sync::Arc; -use std::time::{SystemTime, UNIX_EPOCH}; pub struct LogDecisionHandler; @@ -21,8 +20,12 @@ impl McpTool for LogDecisionHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: LogDecisionTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - - let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); + + let idx = state + .search_index + .read() + .unwrap_or_else(|e| e.into_inner()) + .clone(); let mut final_id = String::new(); state.adrs.modify(|adrs| { @@ -33,13 +36,10 @@ impl McpTool for LogDecisionHandler { context: req.context, decision: req.decision, consequence: req.consequence, - timestamp: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + timestamp: crate::handlers_v2::utils::now_secs(), }; - - let _ = idx.index_adr(&a); + + drop(idx.index_adr(&a)); adrs.push(a); }); @@ -62,15 +62,18 @@ impl McpTool for QueryDecisionsHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: QueryDecisionsTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let data = state.adrs.read_with(|adrs| { - let filtered: Vec<_> = adrs.iter().filter(|a| { - if let Some(q) = &req.query { - contains_ignore_ascii_case(&a.title, q) - || contains_ignore_ascii_case(&a.context, q) - || contains_ignore_ascii_case(&a.decision, q) - } else { - true - } - }).collect(); + let filtered: Vec<_> = adrs + .iter() + .filter(|a| { + if let Some(q) = &req.query { + contains_ignore_ascii_case(&a.title, q) + || contains_ignore_ascii_case(&a.context, q) + || contains_ignore_ascii_case(&a.decision, q) + } else { + true + } + }) + .collect(); serde_json::to_string(&filtered).unwrap_or_default() }); Ok(data.to_string()) @@ -95,10 +98,7 @@ impl McpTool for LogErrorFixHandler { fixes.push(crate::models::ErrorFix { signature: req.signature, solution: req.solution, - timestamp: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + timestamp: crate::handlers_v2::utils::now_secs(), git_commit: req.git_commit, git_branch: req.git_branch, }) @@ -126,10 +126,13 @@ impl McpTool for SearchErrorFixesHandler { let req: SearchErrorFixesTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let q = req.query; let data = state.error_fixes.read_with(|fixes| { - let filtered: Vec<_> = fixes.iter().filter(|f| { - contains_ignore_ascii_case(&f.signature, &q) - || contains_ignore_ascii_case(&f.solution, &q) - }).collect(); + let filtered: Vec<_> = fixes + .iter() + .filter(|f| { + contains_ignore_ascii_case(&f.signature, &q) + || contains_ignore_ascii_case(&f.solution, &q) + }) + .collect(); serde_json::to_string(&filtered).unwrap_or_default() }); Ok(data.to_string()) @@ -152,10 +155,7 @@ impl McpTool for LogCodeChangeHandler { let req: LogCodeChangeTool = serde_json::from_value(args).map_err(|e| e.to_string())?; state.ledger.modify(|ledger| { ledger.push(CodeChange { - timestamp: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + timestamp: crate::handlers_v2::utils::now_secs(), file_path: req.file_path, description: req.description, git_commit: req.git_commit, @@ -182,7 +182,9 @@ impl McpTool for QueryRecentChangesHandler { } async fn execute(&self, _args: Value, state: Arc) -> Result { - let data = state.ledger.read_with(|l| serde_json::to_string(l).unwrap_or_else(|_| "[]".to_string())); + let data = state + .ledger + .read_with(|l| serde_json::to_string(l).unwrap_or_else(|_| "[]".to_string())); Ok(data.to_string()) } } @@ -207,10 +209,7 @@ impl McpTool for LearnPreferenceHandler { crate::models::Preference { key: req.key.clone(), value: req.value, - updated_at: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + updated_at: crate::handlers_v2::utils::now_secs(), }, ); }); @@ -231,7 +230,9 @@ impl McpTool for ReadPreferencesHandler { } async fn execute(&self, _args: Value, state: Arc) -> Result { - let data = state.prefs.read_with(|prefs| serde_json::to_string(prefs).unwrap_or_default()); + let data = state + .prefs + .read_with(|prefs| serde_json::to_string(prefs).unwrap_or_default()); Ok(data.to_string()) } } @@ -257,10 +258,7 @@ impl McpTool for LogTechDebtHandler { description: req.description, ideal_solution: req.ideal_solution, is_resolved: false, - created_at: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + created_at: crate::handlers_v2::utils::now_secs(), git_commit: req.git_commit, git_branch: req.git_branch, }) @@ -319,9 +317,12 @@ impl McpTool for ListTechDebtHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: ListTechDebtTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let data = state.tech_debts.read_with(|debts| { - let filtered: Vec<_> = debts.iter().filter(|d| { - d.namespace == req.namespace && (req.include_resolved || !d.is_resolved) - }).collect(); + let filtered: Vec<_> = debts + .iter() + .filter(|d| { + d.namespace == req.namespace && (req.include_resolved || !d.is_resolved) + }) + .collect(); serde_json::to_string(&filtered).unwrap_or_default() }); Ok(data.to_string()) @@ -361,49 +362,73 @@ impl McpTool for OmniSearchHandler { }); let tasks_json = state.tasks.read_with(|all_tasks| { - let filtered: Vec<_> = all_tasks.iter().filter(|t| { - matches.iter().any(|(id, typ, _, _, _)| id == &t.id && typ == "task") - }).collect(); + let filtered: Vec<_> = all_tasks + .iter() + .filter(|t| { + matches + .iter() + .any(|(id, typ, _, _, _)| id == &t.id && typ == "task") + }) + .collect(); serde_json::to_value(&filtered).unwrap_or_default() }); let snippets_json = state.snippets.read_with(|all_snippets| { - let filtered: Vec<_> = all_snippets.iter().filter(|s| { - matches.iter().any(|(id, typ, _, _, _)| id == &s.name && typ == "snippet") - }).collect(); + let filtered: Vec<_> = all_snippets + .iter() + .filter(|s| { + matches + .iter() + .any(|(id, typ, _, _, _)| id == &s.name && typ == "snippet") + }) + .collect(); serde_json::to_value(&filtered).unwrap_or_default() }); let adrs_json = state.adrs.read_with(|all_adrs| { - let filtered: Vec<_> = all_adrs.iter().filter(|a| { - matches.iter().any(|(id, typ, _, _, _)| id == &a.id && typ == "adr") - }).collect(); + let filtered: Vec<_> = all_adrs + .iter() + .filter(|a| { + matches + .iter() + .any(|(id, typ, _, _, _)| id == &a.id && typ == "adr") + }) + .collect(); serde_json::to_value(&filtered).unwrap_or_default() }); let q = req.query; let tech_debts_json = state.tech_debts.read_with(|debts| { - let filtered: Vec<_> = debts.iter().filter(|d| { - req.namespace.as_ref().is_none_or(|ns| d.namespace == *ns) - && (contains_ignore_ascii_case(&d.description, &q) - || contains_ignore_ascii_case(&d.ideal_solution, &q)) - }).collect(); + let filtered: Vec<_> = debts + .iter() + .filter(|d| { + req.namespace.as_ref().is_none_or(|ns| d.namespace == *ns) + && (contains_ignore_ascii_case(&d.description, &q) + || contains_ignore_ascii_case(&d.ideal_solution, &q)) + }) + .collect(); serde_json::to_value(&filtered).unwrap_or_default() }); let memos_json = state.handoff_memos.read_with(|memos| { - let filtered: Vec<_> = memos.iter().filter(|m| { - req.namespace.as_ref().is_none_or(|ns| m.namespace == *ns) - && contains_ignore_ascii_case(&m.content, &q) - }).collect(); + let filtered: Vec<_> = memos + .iter() + .filter(|m| { + req.namespace.as_ref().is_none_or(|ns| m.namespace == *ns) + && contains_ignore_ascii_case(&m.content, &q) + }) + .collect(); serde_json::to_value(&filtered).unwrap_or_default() }); let error_fixes_json = state.error_fixes.read_with(|fixes| { - let filtered: Vec<_> = fixes.iter().filter(|f| { - contains_ignore_ascii_case(&f.signature, &q) - || contains_ignore_ascii_case(&f.solution, &q) - }).collect(); + let filtered: Vec<_> = fixes + .iter() + .filter(|f| { + contains_ignore_ascii_case(&f.signature, &q) + || contains_ignore_ascii_case(&f.solution, &q) + }) + .collect(); serde_json::to_value(&filtered).unwrap_or_default() }); @@ -437,11 +462,33 @@ impl McpTool for GetProjectHealthHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: GetProjectHealthTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let active_tasks = state.tasks.read_with(|tasks| tasks.iter().filter(|t| t.status != "done").count()); - let unresolved_debt = state.tech_debts.read_with(|debts| debts.iter().filter(|d| d.namespace == req.namespace && !d.is_resolved).count()); - let unread_memos = state.handoff_memos.read_with(|memos| memos.iter().filter(|m| m.namespace == req.namespace).count()); - let active_milestones = state.milestones.read_with(|milestones| milestones.iter().filter(|m| m.namespace == req.namespace && m.status != "done").count()); - let remaining_checklists = state.pr_checklists.read_with(|checklists| checklists.iter().filter(|c| c.namespace == req.namespace).count()); + let active_tasks = state + .tasks + .read_with(|tasks| tasks.iter().filter(|t| t.status != "done").count()); + let unresolved_debt = state.tech_debts.read_with(|debts| { + debts + .iter() + .filter(|d| d.namespace == req.namespace && !d.is_resolved) + .count() + }); + let unread_memos = state.handoff_memos.read_with(|memos| { + memos + .iter() + .filter(|m| m.namespace == req.namespace) + .count() + }); + let active_milestones = state.milestones.read_with(|milestones| { + milestones + .iter() + .filter(|m| m.namespace == req.namespace && m.status != "done") + .count() + }); + let remaining_checklists = state.pr_checklists.read_with(|checklists| { + checklists + .iter() + .filter(|c| c.namespace == req.namespace) + .count() + }); let report = serde_json::json!({ "active_tasks": active_tasks, @@ -455,6 +502,3 @@ impl McpTool for GetProjectHealthHandler { } use crate::handlers_v2::utils::*; - - - diff --git a/server/src/handlers_v2/notes.rs b/server/src/handlers_v2/notes.rs index e9348bc..bf1cfee 100644 --- a/server/src/handlers_v2/notes.rs +++ b/server/src/handlers_v2/notes.rs @@ -6,7 +6,6 @@ use async_trait::async_trait; use serde_json::Value; use std::collections::HashSet; use std::sync::Arc; -use std::time::{SystemTime, UNIX_EPOCH}; pub struct AddStickyNoteHandler; @@ -24,10 +23,7 @@ impl McpTool for AddStickyNoteHandler { let req: AddStickyNoteTool = serde_json::from_value(args).map_err(|e| e.to_string())?; state.sticky.modify(|notes| { notes.push(StickyNote { - timestamp: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + timestamp: crate::handlers_v2::utils::now_secs(), content: req.content, }); }); @@ -51,7 +47,9 @@ impl McpTool for ReadStickyNotesHandler { } async fn execute(&self, _args: Value, state: Arc) -> Result { - let data = state.sticky.read_with(|s| serde_json::to_string(s).unwrap_or_else(|_| "[]".to_string())); + let data = state + .sticky + .read_with(|s| serde_json::to_string(s).unwrap_or_else(|_| "[]".to_string())); Ok(data.to_string()) } } @@ -134,10 +132,7 @@ impl McpTool for LeaveHandoffMemoHandler { author: "agy".to_string(), content: req.content, namespace: req.namespace, - timestamp: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + timestamp: crate::handlers_v2::utils::now_secs(), }) }); Ok("Handoff memo left".to_string()) @@ -162,13 +157,16 @@ impl McpTool for ReadHandoffMemosHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: ReadHandoffMemosTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let data = state.handoff_memos.read_with(|items| { - let filtered: Vec<_> = items.iter().filter(|i| { - if let Some(ns) = &req.namespace { - &i.namespace == ns - } else { - true - } - }).collect(); + let filtered: Vec<_> = items + .iter() + .filter(|i| { + if let Some(ns) = &req.namespace { + &i.namespace == ns + } else { + true + } + }) + .collect(); serde_json::to_string(&filtered).unwrap_or_default() }); Ok(data.to_string()) @@ -221,10 +219,7 @@ impl McpTool for AddSessionSummaryHandler { summaries.push(crate::models::SessionSummary { summary: req.summary, namespace: req.namespace, - timestamp: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + timestamp: crate::handlers_v2::utils::now_secs(), }) }); Ok("Session summary added".to_string()) @@ -249,12 +244,9 @@ impl McpTool for GenerateStandupReportHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: GenerateStandupReportTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let cutoff = SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs() - .saturating_sub(req.hours_lookback * 3600); - + let cutoff = + crate::handlers_v2::utils::now_secs().saturating_sub(req.hours_lookback * 3600); + let report_str = state.tasks.read_with(|items| { state.ledger.read_with(|changes| { state.session_summaries.read_with(|summaries| { @@ -269,4 +261,3 @@ impl McpTool for GenerateStandupReportHandler { Ok(report_str) } } - diff --git a/server/src/handlers_v2/tasks.rs b/server/src/handlers_v2/tasks.rs index e6f2727..8f125e7 100644 --- a/server/src/handlers_v2/tasks.rs +++ b/server/src/handlers_v2/tasks.rs @@ -5,7 +5,6 @@ use crate::tools::*; use async_trait::async_trait; use serde_json::Value; use std::sync::Arc; -use std::time::{SystemTime, UNIX_EPOCH}; pub struct AddTaskHandler; @@ -21,10 +20,7 @@ impl McpTool for AddTaskHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: AddTaskTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let now = SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(); + let now = crate::handlers_v2::utils::now_secs(); let task_id = uuid::Uuid::new_v4().to_string(); let deps = req.dependencies.unwrap_or_default(); @@ -41,8 +37,12 @@ impl McpTool for AddTaskHandler { dependencies: deps, acceptance_criteria: vec![], }; - let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); - let _ = idx.index_task(&task); + let idx = state + .search_index + .read() + .unwrap_or_else(|e| e.into_inner()) + .clone(); + drop(idx.index_task(&task)); state.tasks.modify(|tasks| { tasks.push(task); }); @@ -68,19 +68,21 @@ impl McpTool for DeleteTaskHandler { let mut actually_deleted = Vec::new(); state.tasks.modify(|tasks| { let initial_len = tasks.len(); - + // Build index-based children map - let mut children_map: std::collections::HashMap> = std::collections::HashMap::new(); + let mut children_map: std::collections::HashMap> = + std::collections::HashMap::new(); let mut id_to_index = std::collections::HashMap::new(); for (idx, t) in tasks.iter().enumerate() { id_to_index.insert(t.id.as_str(), idx); } - + for (idx, t) in tasks.iter().enumerate() { if let Some(pid) = &t.parent_id - && let Some(&parent_idx) = id_to_index.get(pid.as_str()) { - children_map.entry(parent_idx).or_default().push(idx); - } + && let Some(&parent_idx) = id_to_index.get(pid.as_str()) + { + children_map.entry(parent_idx).or_default().push(idx); + } } let mut to_delete_idx = std::collections::HashSet::new(); @@ -90,24 +92,29 @@ impl McpTool for DeleteTaskHandler { while let Some(curr) = queue.pop_front() { if to_delete_idx.insert(curr) - && let Some(children) = children_map.get(&curr) { - queue.extend(children.iter().copied()); - } + && let Some(children) = children_map.get(&curr) + { + queue.extend(children.iter().copied()); + } } } - + for &idx in &to_delete_idx { actually_deleted.push(tasks[idx].id.clone()); } - + tasks.retain(|t| !actually_deleted.contains(&t.id)); deleted_count = initial_len - tasks.len(); }); if deleted_count > 0 { - let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); + let idx = state + .search_index + .read() + .unwrap_or_else(|e| e.into_inner()) + .clone(); for id in actually_deleted { - let _ = idx.delete_document(&id); + drop(idx.delete_document(&id)); } Ok(vec![ format!("Deleted task and its children ({} total).", deleted_count).to_string(), @@ -143,18 +150,24 @@ impl McpTool for UpdateTaskStatusHandler { state.tasks.modify(|tasks| { // Find target task - let target_idx = tasks.iter().position(|t| t.id == req.id || t.title == req.id); + let target_idx = tasks + .iter() + .position(|t| t.id == req.id || t.title == req.id); let target_idx = match target_idx { Some(idx) => idx, None => return, }; - + found = true; let target_id = tasks[target_idx].id.clone(); if target_status == "done" || target_status == "completed" { // 1. Check Acceptance Criteria - if tasks[target_idx].acceptance_criteria.iter().any(|c| !c.is_met) { + if tasks[target_idx] + .acceptance_criteria + .iter() + .any(|c| !c.is_met) + { blocked = true; blocker_details = "Unmet acceptance criteria exist.".to_string(); } @@ -164,27 +177,36 @@ impl McpTool for UpdateTaskStatusHandler { let mut uncompleted_deps = Vec::new(); for dep_id in &tasks[target_idx].dependencies { if let Some(dep_task) = tasks.iter().find(|dt| dt.id == *dep_id) - && dep_task.status != "completed" && dep_task.status != "done" { - uncompleted_deps.push(dep_task.title.as_str()); - } + && dep_task.status != "completed" + && dep_task.status != "done" + { + uncompleted_deps.push(dep_task.title.as_str()); + } } if !uncompleted_deps.is_empty() { blocked = true; - blocker_details = format!("Blocked by dependencies: {}", uncompleted_deps.join(", ")); + blocker_details = + format!("Blocked by dependencies: {}", uncompleted_deps.join(", ")); } } // 3. Check child tasks if !blocked { let mut uncompleted_children = Vec::new(); - for child in tasks.iter().filter(|t| t.parent_id.as_ref() == Some(&target_id)) { + for child in tasks + .iter() + .filter(|t| t.parent_id.as_ref() == Some(&target_id)) + { if child.status != "completed" && child.status != "done" { uncompleted_children.push(child.title.as_str()); } } if !uncompleted_children.is_empty() { blocked = true; - blocker_details = format!("Blocked by child tasks: {}", uncompleted_children.join(", ")); + blocker_details = format!( + "Blocked by child tasks: {}", + uncompleted_children.join(", ") + ); } } } @@ -192,16 +214,13 @@ impl McpTool for UpdateTaskStatusHandler { if !blocked { // Apply update tasks[target_idx].status = target_status.clone(); - tasks[target_idx].updated_at = SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(); + tasks[target_idx].updated_at = crate::handlers_v2::utils::now_secs(); // Cascade cancellation to children if target_status == "cancelled" || target_status == "abandoned" { let mut children_map: std::collections::HashMap> = std::collections::HashMap::new(); - + // First pass: map string ID to index to build the adjacency list by index let mut id_to_idx = std::collections::HashMap::new(); for (idx, t) in tasks.iter().enumerate() { @@ -210,9 +229,10 @@ impl McpTool for UpdateTaskStatusHandler { for (idx, t) in tasks.iter().enumerate() { if let Some(pid) = &t.parent_id - && let Some(&p_idx) = id_to_idx.get(pid.as_str()) { - children_map.entry(p_idx).or_default().push(idx); - } + && let Some(&p_idx) = id_to_idx.get(pid.as_str()) + { + children_map.entry(p_idx).or_default().push(idx); + } } if let Some(&start_idx) = id_to_idx.get(target_id.as_str()) { @@ -267,14 +287,20 @@ impl McpTool for ListActiveTasksHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: ListActiveTasksTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let data = state.tasks.read_with(|tasks| { - let filtered: Vec<_> = tasks.iter().filter(|t| { - let status_match = t.status != "done"; - let branch_match = match &req.git_branch { - Some(branch) => t.git_branch.is_none() || t.git_branch.as_deref() == Some(branch.as_str()), - None => true, - }; - status_match && branch_match - }).collect(); + let filtered: Vec<_> = tasks + .iter() + .filter(|t| { + let status_match = t.status != "done"; + let branch_match = match &req.git_branch { + Some(branch) => { + t.git_branch.is_none() + || t.git_branch.as_deref() == Some(branch.as_str()) + } + None => true, + }; + status_match && branch_match + }) + .collect(); serde_json::to_string(&filtered).unwrap_or_default() }); Ok(data) @@ -311,10 +337,7 @@ impl McpTool for SetAcceptanceCriteriaHandler { is_met: false, }) .collect(); - task.updated_at = SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(); + task.updated_at = crate::handlers_v2::utils::now_secs(); success = true; } }); @@ -358,10 +381,7 @@ impl McpTool for VerifyAcceptanceCriteriaHandler { } else { ac.is_met = true; success = true; - task.updated_at = SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(); + task.updated_at = crate::handlers_v2::utils::now_secs(); } } }); @@ -452,13 +472,16 @@ impl McpTool for ListMilestonesHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: ListMilestonesTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let data = state.milestones.read_with(|items| { - let filtered: Vec<_> = items.iter().filter(|i| { - if let Some(ns) = &req.namespace { - &i.namespace == ns - } else { - true - } - }).collect(); + let filtered: Vec<_> = items + .iter() + .filter(|i| { + if let Some(ns) = &req.namespace { + &i.namespace == ns + } else { + true + } + }) + .collect(); serde_json::to_string(&filtered).unwrap_or_default() }); Ok(data) diff --git a/server/src/handlers_v2/utils.rs b/server/src/handlers_v2/utils.rs index bbfe842..998ddc3 100644 --- a/server/src/handlers_v2/utils.rs +++ b/server/src/handlers_v2/utils.rs @@ -7,3 +7,9 @@ pub fn contains_ignore_ascii_case(haystack: &str, needle: &str) -> bool { .windows(needle.len()) .any(|w| w.eq_ignore_ascii_case(needle.as_bytes())) } +pub fn now_secs() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_secs() +} diff --git a/server/src/handlers_v2/workspaces.rs b/server/src/handlers_v2/workspaces.rs index d739207..35cff51 100644 --- a/server/src/handlers_v2/workspaces.rs +++ b/server/src/handlers_v2/workspaces.rs @@ -5,7 +5,6 @@ use crate::tools::*; use async_trait::async_trait; use serde_json::Value; use std::sync::Arc; -use std::time::{SystemTime, UNIX_EPOCH}; pub struct PinFileHandler; @@ -26,10 +25,7 @@ impl McpTool for PinFileHandler { pinned.push(crate::models::PinnedFile { namespace: req.namespace, file_path: req.file_path, - timestamp: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + timestamp: crate::handlers_v2::utils::now_secs(), git_branch: req.git_branch, }); }); @@ -76,17 +72,23 @@ impl McpTool for ListPinnedFilesHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: ListPinnedFilesTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let data = state.pinned_files.read_with(|pinned| { - let filtered: Vec<_> = pinned.iter().filter(|p| { - let ns_match = match &req.namespace { - Some(ns) => &p.namespace == ns, - std::option::Option::None => true, - }; - let branch_match = match &req.git_branch { - Some(branch) => p.git_branch.is_none() || p.git_branch.as_deref() == Some(branch.as_str()), - std::option::Option::None => true, - }; - ns_match && branch_match - }).collect(); + let filtered: Vec<_> = pinned + .iter() + .filter(|p| { + let ns_match = match &req.namespace { + Some(ns) => &p.namespace == ns, + std::option::Option::None => true, + }; + let branch_match = match &req.git_branch { + Some(branch) => { + p.git_branch.is_none() + || p.git_branch.as_deref() == Some(branch.as_str()) + } + std::option::Option::None => true, + }; + ns_match && branch_match + }) + .collect(); serde_json::to_string(&filtered).unwrap_or_default() }); Ok(data.to_string()) @@ -113,14 +115,15 @@ impl McpTool for StoreSnippetHandler { language: req.language, code: req.code, description: req.description, - updated_at: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + updated_at: crate::handlers_v2::utils::now_secs(), }; - let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); - let _ = idx.index_snippet(&snippet); + let idx = state + .search_index + .read() + .unwrap_or_else(|e| e.into_inner()) + .clone(); + drop(idx.index_snippet(&snippet)); state.snippets.modify(|snippets| { snippets.retain(|s| s.name != req_name); @@ -147,11 +150,14 @@ impl McpTool for SearchSnippetsHandler { let req: SearchSnippetsTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let query = req.query; let data = state.snippets.read_with(|snippets| { - let results: Vec<_> = snippets.iter().filter(|s| { - contains_ignore_ascii_case(&s.name, &query) - || contains_ignore_ascii_case(&s.description, &query) - || contains_ignore_ascii_case(&s.language, &query) - }).collect(); + let results: Vec<_> = snippets + .iter() + .filter(|s| { + contains_ignore_ascii_case(&s.name, &query) + || contains_ignore_ascii_case(&s.description, &query) + || contains_ignore_ascii_case(&s.language, &query) + }) + .collect(); serde_json::to_string(&results).unwrap_or_default() }); Ok(data.to_string()) @@ -179,8 +185,12 @@ impl McpTool for DeleteSnippetHandler { deleted = snippets.len() < orig; }); if deleted { - let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); - let _ = idx.delete_document(&req.name); + let idx = state + .search_index + .read() + .unwrap_or_else(|e| e.into_inner()) + .clone(); + drop(idx.delete_document(&req.name)); Ok("Snippet deleted.".to_string()) } else { Ok("Snippet not found.".to_string()) @@ -213,10 +223,7 @@ impl McpTool for SaveContextWorkspaceHandler { name: req.name, pinned_files: req.pinned_files, active_task_ids: req.active_task_ids, - saved_at: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + saved_at: crate::handlers_v2::utils::now_secs(), }); }); Ok("Context workspace saved".to_string()) @@ -242,9 +249,10 @@ impl McpTool for LoadContextWorkspaceHandler { let req: LoadContextWorkspaceTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let data = state.context_workspaces.read_with(|ws| { - let filtered: Vec<_> = ws.iter().filter(|w| { - w.namespace == req.namespace && w.name == req.name - }).collect(); + let filtered: Vec<_> = ws + .iter() + .filter(|w| w.namespace == req.namespace && w.name == req.name) + .collect(); serde_json::to_string(&filtered.first()).unwrap_or_default() }); Ok(data.to_string()) @@ -270,9 +278,7 @@ impl McpTool for ListContextWorkspacesHandler { let req: ListContextWorkspacesTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let data = state.context_workspaces.read_with(|ws| { - let filtered: Vec<_> = ws.iter().filter(|w| { - w.namespace == req.namespace - }).collect(); + let filtered: Vec<_> = ws.iter().filter(|w| w.namespace == req.namespace).collect(); serde_json::to_string(&filtered).unwrap_or_default() }); Ok(data.to_string()) @@ -323,9 +329,10 @@ impl McpTool for GetPrChecklistHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: GetPrChecklistTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let data = state.pr_checklists.read_with(|items| { - let filtered: Vec<_> = items.iter().filter(|i| { - i.namespace == req.namespace - }).collect(); + let filtered: Vec<_> = items + .iter() + .filter(|i| i.namespace == req.namespace) + .collect(); serde_json::to_string(&filtered).unwrap_or_default() }); Ok(data.to_string()) @@ -357,5 +364,3 @@ impl McpTool for ClearPrChecklistHandler { } use crate::handlers_v2::utils::*; - - diff --git a/server/src/main.rs b/server/src/main.rs index 44cea5b..27c55b0 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -22,7 +22,7 @@ use redb::ReadableTable; use std::fs; use std::path::PathBuf; use std::sync::{Arc, RwLock}; -use std::time::{Duration, SystemTime, UNIX_EPOCH}; +use std::time::Duration; use clap::{Parser, Subcommand}; use std::collections::HashMap; @@ -204,10 +204,7 @@ async fn gate_set_handler( params: body.params.clone(), status, reason: body.reason.clone(), - timestamp: SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(), + timestamp: crate::handlers_v2::utils::now_secs(), }; app_state.handler.state.gates.modify(|gates| { gates.retain(|g| !(g.action == record.action && g.target == record.target)); @@ -430,7 +427,9 @@ async fn run_server(state: Arc) -> Result<(), Box l, @@ -726,14 +725,16 @@ fn main() -> Result<(), Box> { if !cli.daemon { // Just spawn the daemon and exit. We no longer act as a proxy. #[allow(clippy::zombie_processes)] - let _ = std::process::Command::new(std::env::current_exe().expect("Failed to get current executable path")) - .arg("--daemon") - .stdin(std::process::Stdio::null()) - .stdout(std::process::Stdio::null()) - .stderr(std::process::Stdio::null()) - .creation_flags(0x08000000) // CREATE_NO_WINDOW - .spawn() - .expect("Failed to spawn daemon"); + let _ = std::process::Command::new( + std::env::current_exe().expect("Failed to get current executable path"), + ) + .arg("--daemon") + .stdin(std::process::Stdio::null()) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .creation_flags(0x08000000) // CREATE_NO_WINDOW + .spawn() + .expect("Failed to spawn daemon"); return Ok(()); } } @@ -752,7 +753,9 @@ fn main() -> Result<(), Box> { { let write_txn = db.begin_write().expect("Failed to begin write txn on redb"); { - let mut table = write_txn.open_table(crate::store::STORE_TABLE).expect("Failed to open STORE_TABLE"); + let mut table = write_txn + .open_table(crate::store::STORE_TABLE) + .expect("Failed to open STORE_TABLE"); let stores = vec![ ("knowledge_graph_master", "knowledge_graph_master.json"), @@ -777,13 +780,19 @@ fn main() -> Result<(), Box> { ]; for (key, file_name) in stores.iter() { - if table.get(*key).expect("Failed to read from table").is_none() { + if table + .get(*key) + .expect("Failed to read from table") + .is_none() + { let json_path = base.join(file_name); if json_path.exists() && let Ok(data) = fs::read(&json_path) && serde_json::from_slice::(&data).is_ok() { - table.insert(*key, data.as_slice()).expect("Failed to insert migrated data"); + table + .insert(*key, data.as_slice()) + .expect("Failed to insert migrated data"); let _ = fs::rename(&json_path, json_path.with_extension("json.migrated")); } } diff --git a/server/src/mcp.rs b/server/src/mcp.rs index a2850c7..2ea83e7 100644 --- a/server/src/mcp.rs +++ b/server/src/mcp.rs @@ -30,4 +30,3 @@ pub fn tool_def(name: &str, description: &str) -> serde_json::Val "inputSchema": schema_val }) } - diff --git a/server/src/search.rs b/server/src/search.rs index a45bce7..85bb1e9 100644 --- a/server/src/search.rs +++ b/server/src/search.rs @@ -62,7 +62,7 @@ impl MemoryIndex { let id_field = self.id_field; let id_val = e.name.clone(); let needs_commit = Arc::clone(&self.needs_commit); - + let doc = doc!( self.id_field => e.name.as_str(), self.title_field => e.name.as_str(), @@ -108,7 +108,7 @@ impl MemoryIndex { let id_field = self.id_field; let id_val = id.to_string(); let needs_commit = Arc::clone(&self.needs_commit); - + tokio::task::spawn_blocking(move || { let writer = writer.lock().unwrap_or_else(|e| e.into_inner()); writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val)); diff --git a/server/src/state.rs b/server/src/state.rs index 7d6e9e1..7eb3008 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -42,12 +42,12 @@ impl MemoryState { .duration_since(std::time::UNIX_EPOCH) .unwrap_or_default() .as_millis() as u64; - + let item = serde_json::json!({ "time": time, "message": message }); - + self.recent_activities.modify(|activities| { activities.push_back(item.clone()); if activities.len() > 100 { @@ -79,7 +79,7 @@ impl MemoryState { if let Ok(new_idx) = MemoryIndex::new(&self.base_dir) { let state = Arc::clone(self); let idx = new_idx.clone(); - + tokio::task::spawn_blocking(move || { state.graph.read_with(|g| { for e in g.entities.values() { @@ -101,7 +101,9 @@ impl MemoryState { idx.add_adr_sync(a); } }); - }).await.unwrap_or_else(|e| { + }) + .await + .unwrap_or_else(|e| { tracing::error!("Failed to join tantivy index rebuild thread: {}", e); }); diff --git a/stub/src/main.rs b/stub/src/main.rs index 028193b..f615935 100644 --- a/stub/src/main.rs +++ b/stub/src/main.rs @@ -170,5 +170,3 @@ fn main() -> Result<(), Box> { Ok(()) }) } - -