diff --git a/server/src/handlers_v2/graph.rs b/server/src/handlers_v2/graph.rs index 1be73e1..b237e54 100644 --- a/server/src/handlers_v2/graph.rs +++ b/server/src/handlers_v2/graph.rs @@ -126,7 +126,7 @@ impl McpTool for CreateEntitiesHandler { }); let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); for entity in inserted { - let _ = idx.index_entity(&entity).await; + let _ = idx.index_entity(&entity); } Ok("Entities created".to_string()) } @@ -208,7 +208,7 @@ impl McpTool for DeleteEntitiesHandler { let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); for name in to_delete { - let _ = idx.delete_document(&name).await; + let _ = idx.delete_document(&name); } Ok("Entities deleted".to_string()) } diff --git a/server/src/handlers_v2/meta.rs b/server/src/handlers_v2/meta.rs index 2456120..a7d9e8e 100644 --- a/server/src/handlers_v2/meta.rs +++ b/server/src/handlers_v2/meta.rs @@ -21,13 +21,14 @@ 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 mut adr_id = String::new(); - let mut new_adr = None; + + let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); + let mut final_id = String::new(); state.adrs.modify(|adrs| { - adr_id = format!("ADR-{:04}", adrs.len() + 1); + final_id = format!("ADR-{:04}", adrs.len() + 1); let a = Adr { - id: adr_id.clone(), + id: final_id.clone(), title: req.title, context: req.context, decision: req.decision, @@ -37,16 +38,12 @@ impl McpTool for LogDecisionHandler { .unwrap_or_default() .as_secs(), }; - new_adr = Some(a.clone()); + + let _ = idx.index_adr(&a); adrs.push(a); }); - if let Some(adr) = new_adr { - let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); - let _ = idx.index_adr(&adr).await; - } - - Ok(format!("Decision logged as {}", adr_id).to_string()) + Ok(format!("Decision logged as {}", final_id)) } } diff --git a/server/src/handlers_v2/tasks.rs b/server/src/handlers_v2/tasks.rs index 0e237cd..c0cfc7c 100644 --- a/server/src/handlers_v2/tasks.rs +++ b/server/src/handlers_v2/tasks.rs @@ -27,8 +27,7 @@ impl McpTool for AddTaskHandler { .as_secs(); let task_id = uuid::Uuid::new_v4().to_string(); - let parent_id = req.parent_id.clone(); - let deps = req.dependencies.clone().unwrap_or_default(); + let deps = req.dependencies.unwrap_or_default(); let task = Task { id: task_id.clone(), @@ -38,12 +37,12 @@ impl McpTool for AddTaskHandler { created_at: now, updated_at: now, git_branch: req.git_branch, - parent_id, + parent_id: req.parent_id, dependencies: deps, acceptance_criteria: vec![], }; let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); - let _ = idx.index_task(&task).await; + let _ = idx.index_task(&task); state.tasks.modify(|tasks| { tasks.push(task); }); @@ -69,33 +68,40 @@ impl McpTool for DeleteTaskHandler { let mut actually_deleted = Vec::new(); state.tasks.modify(|tasks| { let initial_len = tasks.len(); - // Collect IDs of tasks to delete (this task + all its recursive children) - let mut to_delete = std::collections::HashSet::new(); - to_delete.insert(req.id.clone()); - - let mut children_map: std::collections::HashMap> = - std::collections::HashMap::new(); - for t in tasks.iter() { + + // Build index-based children map + 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 { - children_map - .entry(pid.clone()) - .or_default() - .push(t.id.clone()); + if let Some(&parent_idx) = id_to_index.get(pid.as_str()) { + children_map.entry(parent_idx).or_default().push(idx); + } } } - let mut queue = std::collections::VecDeque::new(); - queue.push_back(req.id.clone()); + let mut to_delete_idx = std::collections::HashSet::new(); + if let Some(&start_idx) = id_to_index.get(req.id.as_str()) { + let mut queue = std::collections::VecDeque::new(); + queue.push_back(start_idx); - while let Some(curr) = queue.pop_front() { - if to_delete.insert(curr.clone()) - && let Some(children) = children_map.get(&curr) - { - queue.extend(children.iter().cloned()); + while let Some(curr) = queue.pop_front() { + if to_delete_idx.insert(curr) { + if let Some(children) = children_map.get(&curr) { + queue.extend(children.iter().copied()); + } + } } } - actually_deleted = to_delete.into_iter().collect(); + 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(); }); @@ -103,7 +109,7 @@ impl McpTool for DeleteTaskHandler { if deleted_count > 0 { let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); for id in actually_deleted { - let _ = idx.delete_document(&id).await; + let _ = idx.delete_document(&id); } Ok(vec![ format!("Deleted task and its children ({} total).", deleted_count).to_string(), @@ -139,21 +145,18 @@ impl McpTool for UpdateTaskStatusHandler { state.tasks.modify(|tasks| { // Find target task - let mut target_id = String::new(); - if let Some(t) = tasks.iter().find(|t| t.id == req.id || t.title == req.id) { - target_id = t.id.clone(); - } - - if target_id.is_empty() { - return; - } + 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 let Some(t) = tasks.iter().find(|t| t.id == target_id) - && t.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(); } @@ -161,53 +164,41 @@ impl McpTool for UpdateTaskStatusHandler { // 2. Check dependencies if !blocked { let mut uncompleted_deps = Vec::new(); - if let Some(t) = tasks.iter().find(|t| t.id == target_id) { - for dep_id in &t.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.clone()); + for dep_id in &tasks[target_idx].dependencies { + if let Some(dep_task) = tasks.iter().find(|dt| dt.id == *dep_id) { + if 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.clone()); + 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(", ")); } } } if !blocked { // Apply update - if let Some(t) = tasks.iter_mut().find(|t| t.id == target_id) { - t.status = target_status.clone(); - t.updated_at = SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(); - } + tasks[target_idx].status = target_status.clone(); + tasks[target_idx].updated_at = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_secs(); // Cascade cancellation to children if target_status == "cancelled" || target_status == "abandoned" { diff --git a/server/src/handlers_v2/workspaces.rs b/server/src/handlers_v2/workspaces.rs index c54053b..df68dfa 100644 --- a/server/src/handlers_v2/workspaces.rs +++ b/server/src/handlers_v2/workspaces.rs @@ -107,8 +107,9 @@ impl McpTool for StoreSnippetHandler { async fn execute(&self, args: Value, state: Arc) -> Result { let req: StoreSnippetTool = serde_json::from_value(args).map_err(|e| e.to_string())?; + let req_name = req.name.clone(); // Keep for the OK message and retain closure let snippet = Snippet { - name: req.name.clone(), + name: req.name, language: req.language, code: req.code, description: req.description, @@ -118,16 +119,15 @@ impl McpTool for StoreSnippetHandler { .as_secs(), }; - let s_clone = snippet.clone(); + let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); + let _ = idx.index_snippet(&snippet); + state.snippets.modify(|snippets| { - snippets.retain(|s| s.name != req.name); - snippets.push(s_clone); + snippets.retain(|s| s.name != req_name); + snippets.push(snippet); }); - let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); - let _ = idx.index_snippet(&snippet).await; - - Ok(format!("Snippet '{}' stored.", req.name).to_string()) + Ok(format!("Snippet '{}' stored.", req_name).to_string()) } } @@ -180,7 +180,7 @@ impl McpTool for DeleteSnippetHandler { }); if deleted { let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone(); - let _ = idx.delete_document(&req.name).await; + let _ = idx.delete_document(&req.name); Ok("Snippet deleted.".to_string()) } else { Ok("Snippet not found.".to_string()) diff --git a/server/src/search.rs b/server/src/search.rs index e01ad98..a45bce7 100644 --- a/server/src/search.rs +++ b/server/src/search.rs @@ -10,6 +10,7 @@ pub struct MemoryIndex { index: Index, reader: IndexReader, writer: Arc>, + needs_commit: Arc, // Schema fields pub id_field: Field, @@ -47,6 +48,7 @@ impl MemoryIndex { index, reader, writer: Arc::new(Mutex::new(writer)), + needs_commit: Arc::new(std::sync::atomic::AtomicBool::new(false)), id_field, title_field, body_field, @@ -59,6 +61,7 @@ impl MemoryIndex { let writer = Arc::clone(&self.writer); 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(), @@ -72,6 +75,7 @@ impl MemoryIndex { let writer = writer.lock().unwrap_or_else(|e| e.into_inner()); writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val)); writer.add_document(doc)?; + needs_commit.store(true, std::sync::atomic::Ordering::SeqCst); Ok(()) }) } @@ -80,6 +84,7 @@ impl MemoryIndex { let writer = Arc::clone(&self.writer); let id_field = self.id_field; let id_val = t.id.clone(); + let needs_commit = Arc::clone(&self.needs_commit); let doc = doc!( self.id_field => t.id.as_str(), @@ -93,6 +98,7 @@ impl MemoryIndex { let writer = writer.lock().unwrap_or_else(|e| e.into_inner()); writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val)); writer.add_document(doc)?; + needs_commit.store(true, std::sync::atomic::Ordering::SeqCst); Ok(()) }) } @@ -101,19 +107,24 @@ impl MemoryIndex { let writer = Arc::clone(&self.writer); 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)); + needs_commit.store(true, std::sync::atomic::Ordering::SeqCst); Ok(()) }) } pub async fn commit(&self) -> tantivy::Result<()> { let writer = Arc::clone(&self.writer); + let needs_commit = Arc::clone(&self.needs_commit); tokio::task::spawn_blocking(move || { - let mut writer = writer.lock().unwrap_or_else(|e| e.into_inner()); - writer.commit()?; + if needs_commit.swap(false, std::sync::atomic::Ordering::SeqCst) { + let mut writer = writer.lock().unwrap_or_else(|e| e.into_inner()); + writer.commit()?; + } Ok(()) }) .await @@ -182,6 +193,7 @@ impl MemoryIndex { let writer = Arc::clone(&self.writer); let id_field = self.id_field; let id_val = s.name.clone(); + let needs_commit = Arc::clone(&self.needs_commit); let doc = doc!( self.id_field => s.name.as_str(), @@ -195,6 +207,7 @@ impl MemoryIndex { let writer = writer.lock().unwrap_or_else(|e| e.into_inner()); writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val)); writer.add_document(doc)?; + needs_commit.store(true, std::sync::atomic::Ordering::SeqCst); Ok(()) }) } @@ -203,6 +216,7 @@ impl MemoryIndex { let writer = Arc::clone(&self.writer); let id_field = self.id_field; let id_val = a.id.clone(); + let needs_commit = Arc::clone(&self.needs_commit); let doc = doc!( self.id_field => a.id.as_str(), @@ -216,6 +230,7 @@ impl MemoryIndex { let writer = writer.lock().unwrap_or_else(|e| e.into_inner()); writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val)); writer.add_document(doc)?; + needs_commit.store(true, std::sync::atomic::Ordering::SeqCst); Ok(()) }) } diff --git a/server/src/state.rs b/server/src/state.rs index 1c73958..c6c6648 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -56,8 +56,9 @@ impl MemoryState { }); let payload = serde_json::json!({ - "type": "activity", - "data": item + "jsonrpc": "2.0", + "method": "notifications/activity", + "params": item }) .to_string(); let _ = self.activity_tx.send(payload);