fix(mcp): wrap websocket activity broadcasts in JSON-RPC notifications
Fixes a critical bug where the Antigravity MCP client would crash with 'invalid message version tag' during long-running tool executions. The server was broadcasting raw JSON activity objects ({'type': 'activity', 'data': ...}) over the WebSocket connection without wrapping them in the required JSON-RPC 2.0 Notification envelope, violating the protocol expectation on the proxy stub.
This commit is contained in:
1 parent
208c5d448f
commit
da1b7cdc9d
6 files changed
+81
-77
No files matched your search
@@ -126,7 +126,7 @@ 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 {
|
for entity in inserted {
|
||||||
let _ = idx.index_entity(&entity).await;
|
let _ = idx.index_entity(&entity);
|
||||||
}
|
}
|
||||||
Ok("Entities created".to_string())
|
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();
|
let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone();
|
||||||
for name in to_delete {
|
for name in to_delete {
|
||||||
let _ = idx.delete_document(&name).await;
|
let _ = idx.delete_document(&name);
|
||||||
}
|
}
|
||||||
Ok("Entities deleted".to_string())
|
Ok("Entities deleted".to_string())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -21,13 +21,14 @@ impl McpTool for LogDecisionHandler {
|
|||||||
|
|
||||||
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> Result<String, String> {
|
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> Result<String, String> {
|
||||||
let req: LogDecisionTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
|
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| {
|
state.adrs.modify(|adrs| {
|
||||||
adr_id = format!("ADR-{:04}", adrs.len() + 1);
|
final_id = format!("ADR-{:04}", adrs.len() + 1);
|
||||||
let a = Adr {
|
let a = Adr {
|
||||||
id: adr_id.clone(),
|
id: final_id.clone(),
|
||||||
title: req.title,
|
title: req.title,
|
||||||
context: req.context,
|
context: req.context,
|
||||||
decision: req.decision,
|
decision: req.decision,
|
||||||
@@ -37,16 +38,12 @@ impl McpTool for LogDecisionHandler {
|
|||||||
.unwrap_or_default()
|
.unwrap_or_default()
|
||||||
.as_secs(),
|
.as_secs(),
|
||||||
};
|
};
|
||||||
new_adr = Some(a.clone());
|
|
||||||
|
let _ = idx.index_adr(&a);
|
||||||
adrs.push(a);
|
adrs.push(a);
|
||||||
});
|
});
|
||||||
|
|
||||||
if let Some(adr) = new_adr {
|
Ok(format!("Decision logged as {}", final_id))
|
||||||
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())
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -27,8 +27,7 @@ impl McpTool for AddTaskHandler {
|
|||||||
.as_secs();
|
.as_secs();
|
||||||
let task_id = uuid::Uuid::new_v4().to_string();
|
let task_id = uuid::Uuid::new_v4().to_string();
|
||||||
|
|
||||||
let parent_id = req.parent_id.clone();
|
let deps = req.dependencies.unwrap_or_default();
|
||||||
let deps = req.dependencies.clone().unwrap_or_default();
|
|
||||||
|
|
||||||
let task = Task {
|
let task = Task {
|
||||||
id: task_id.clone(),
|
id: task_id.clone(),
|
||||||
@@ -38,12 +37,12 @@ impl McpTool for AddTaskHandler {
|
|||||||
created_at: now,
|
created_at: now,
|
||||||
updated_at: now,
|
updated_at: now,
|
||||||
git_branch: req.git_branch,
|
git_branch: req.git_branch,
|
||||||
parent_id,
|
parent_id: req.parent_id,
|
||||||
dependencies: deps,
|
dependencies: deps,
|
||||||
acceptance_criteria: vec![],
|
acceptance_criteria: vec![],
|
||||||
};
|
};
|
||||||
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 _ = idx.index_task(&task).await;
|
let _ = idx.index_task(&task);
|
||||||
state.tasks.modify(|tasks| {
|
state.tasks.modify(|tasks| {
|
||||||
tasks.push(task);
|
tasks.push(task);
|
||||||
});
|
});
|
||||||
@@ -69,33 +68,40 @@ impl McpTool for DeleteTaskHandler {
|
|||||||
let mut actually_deleted = Vec::new();
|
let mut actually_deleted = Vec::new();
|
||||||
state.tasks.modify(|tasks| {
|
state.tasks.modify(|tasks| {
|
||||||
let initial_len = tasks.len();
|
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<String, Vec<String>> =
|
// Build index-based children map
|
||||||
std::collections::HashMap::new();
|
let mut children_map: std::collections::HashMap<usize, Vec<usize>> = std::collections::HashMap::new();
|
||||||
for t in tasks.iter() {
|
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 {
|
if let Some(pid) = &t.parent_id {
|
||||||
children_map
|
if let Some(&parent_idx) = id_to_index.get(pid.as_str()) {
|
||||||
.entry(pid.clone())
|
children_map.entry(parent_idx).or_default().push(idx);
|
||||||
.or_default()
|
}
|
||||||
.push(t.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();
|
let mut queue = std::collections::VecDeque::new();
|
||||||
queue.push_back(req.id.clone());
|
queue.push_back(start_idx);
|
||||||
|
|
||||||
while let Some(curr) = queue.pop_front() {
|
while let Some(curr) = queue.pop_front() {
|
||||||
if to_delete.insert(curr.clone())
|
if to_delete_idx.insert(curr) {
|
||||||
&& let Some(children) = children_map.get(&curr)
|
if let Some(children) = children_map.get(&curr) {
|
||||||
{
|
queue.extend(children.iter().copied());
|
||||||
queue.extend(children.iter().cloned());
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
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));
|
tasks.retain(|t| !actually_deleted.contains(&t.id));
|
||||||
deleted_count = initial_len - tasks.len();
|
deleted_count = initial_len - tasks.len();
|
||||||
});
|
});
|
||||||
@@ -103,7 +109,7 @@ impl McpTool for DeleteTaskHandler {
|
|||||||
if deleted_count > 0 {
|
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 {
|
for id in actually_deleted {
|
||||||
let _ = idx.delete_document(&id).await;
|
let _ = idx.delete_document(&id);
|
||||||
}
|
}
|
||||||
Ok(vec![
|
Ok(vec![
|
||||||
format!("Deleted task and its children ({} total).", deleted_count).to_string(),
|
format!("Deleted task and its children ({} total).", deleted_count).to_string(),
|
||||||
@@ -139,21 +145,18 @@ impl McpTool for UpdateTaskStatusHandler {
|
|||||||
|
|
||||||
state.tasks.modify(|tasks| {
|
state.tasks.modify(|tasks| {
|
||||||
// Find target task
|
// Find target task
|
||||||
let mut target_id = String::new();
|
let target_idx = tasks.iter().position(|t| t.id == req.id || t.title == req.id);
|
||||||
if let Some(t) = tasks.iter().find(|t| t.id == req.id || t.title == req.id) {
|
let target_idx = match target_idx {
|
||||||
target_id = t.id.clone();
|
Some(idx) => idx,
|
||||||
}
|
None => return,
|
||||||
|
};
|
||||||
|
|
||||||
if target_id.is_empty() {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
found = true;
|
found = true;
|
||||||
|
let target_id = tasks[target_idx].id.clone();
|
||||||
|
|
||||||
if target_status == "done" || target_status == "completed" {
|
if target_status == "done" || target_status == "completed" {
|
||||||
// 1. Check Acceptance Criteria
|
// 1. Check Acceptance Criteria
|
||||||
if let Some(t) = tasks.iter().find(|t| t.id == target_id)
|
if tasks[target_idx].acceptance_criteria.iter().any(|c| !c.is_met) {
|
||||||
&& t.acceptance_criteria.iter().any(|c| !c.is_met)
|
|
||||||
{
|
|
||||||
blocked = true;
|
blocked = true;
|
||||||
blocker_details = "Unmet acceptance criteria exist.".to_string();
|
blocker_details = "Unmet acceptance criteria exist.".to_string();
|
||||||
}
|
}
|
||||||
@@ -161,53 +164,41 @@ impl McpTool for UpdateTaskStatusHandler {
|
|||||||
// 2. Check dependencies
|
// 2. Check dependencies
|
||||||
if !blocked {
|
if !blocked {
|
||||||
let mut uncompleted_deps = Vec::new();
|
let mut uncompleted_deps = Vec::new();
|
||||||
if let Some(t) = tasks.iter().find(|t| t.id == target_id) {
|
for dep_id in &tasks[target_idx].dependencies {
|
||||||
for dep_id in &t.dependencies {
|
if let Some(dep_task) = tasks.iter().find(|dt| dt.id == *dep_id) {
|
||||||
if let Some(dep_task) = tasks.iter().find(|dt| dt.id == *dep_id)
|
if dep_task.status != "completed" && dep_task.status != "done" {
|
||||||
&& dep_task.status != "completed"
|
uncompleted_deps.push(dep_task.title.as_str());
|
||||||
&& dep_task.status != "done"
|
|
||||||
{
|
|
||||||
uncompleted_deps.push(dep_task.title.clone());
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if !uncompleted_deps.is_empty() {
|
if !uncompleted_deps.is_empty() {
|
||||||
blocked = true;
|
blocked = true;
|
||||||
blocker_details =
|
blocker_details = format!("Blocked by dependencies: {}", uncompleted_deps.join(", "));
|
||||||
format!("Blocked by dependencies: {}", uncompleted_deps.join(", "));
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 3. Check child tasks
|
// 3. Check child tasks
|
||||||
if !blocked {
|
if !blocked {
|
||||||
let mut uncompleted_children = Vec::new();
|
let mut uncompleted_children = Vec::new();
|
||||||
for child in tasks
|
for child in tasks.iter().filter(|t| t.parent_id.as_ref() == Some(&target_id)) {
|
||||||
.iter()
|
|
||||||
.filter(|t| t.parent_id.as_ref() == Some(&target_id))
|
|
||||||
{
|
|
||||||
if child.status != "completed" && child.status != "done" {
|
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() {
|
if !uncompleted_children.is_empty() {
|
||||||
blocked = true;
|
blocked = true;
|
||||||
blocker_details = format!(
|
blocker_details = format!("Blocked by child tasks: {}", uncompleted_children.join(", "));
|
||||||
"Blocked by child tasks: {}",
|
|
||||||
uncompleted_children.join(", ")
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if !blocked {
|
if !blocked {
|
||||||
// Apply update
|
// Apply update
|
||||||
if let Some(t) = tasks.iter_mut().find(|t| t.id == target_id) {
|
tasks[target_idx].status = target_status.clone();
|
||||||
t.status = target_status.clone();
|
tasks[target_idx].updated_at = SystemTime::now()
|
||||||
t.updated_at = SystemTime::now()
|
|
||||||
.duration_since(UNIX_EPOCH)
|
.duration_since(UNIX_EPOCH)
|
||||||
.unwrap_or_default()
|
.unwrap_or_default()
|
||||||
.as_secs();
|
.as_secs();
|
||||||
}
|
|
||||||
|
|
||||||
// Cascade cancellation to children
|
// Cascade cancellation to children
|
||||||
if target_status == "cancelled" || target_status == "abandoned" {
|
if target_status == "cancelled" || target_status == "abandoned" {
|
||||||
|
|||||||
@@ -107,8 +107,9 @@ impl McpTool for StoreSnippetHandler {
|
|||||||
|
|
||||||
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> Result<String, String> {
|
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> Result<String, String> {
|
||||||
let req: StoreSnippetTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
|
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 {
|
let snippet = Snippet {
|
||||||
name: req.name.clone(),
|
name: req.name,
|
||||||
language: req.language,
|
language: req.language,
|
||||||
code: req.code,
|
code: req.code,
|
||||||
description: req.description,
|
description: req.description,
|
||||||
@@ -118,16 +119,15 @@ impl McpTool for StoreSnippetHandler {
|
|||||||
.as_secs(),
|
.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| {
|
state.snippets.modify(|snippets| {
|
||||||
snippets.retain(|s| s.name != req.name);
|
snippets.retain(|s| s.name != req_name);
|
||||||
snippets.push(s_clone);
|
snippets.push(snippet);
|
||||||
});
|
});
|
||||||
|
|
||||||
let idx = state.search_index.read().unwrap_or_else(|e| e.into_inner()).clone();
|
Ok(format!("Snippet '{}' stored.", req_name).to_string())
|
||||||
let _ = idx.index_snippet(&snippet).await;
|
|
||||||
|
|
||||||
Ok(format!("Snippet '{}' stored.", req.name).to_string())
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -180,7 +180,7 @@ impl McpTool for DeleteSnippetHandler {
|
|||||||
});
|
});
|
||||||
if deleted {
|
if deleted {
|
||||||
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 _ = idx.delete_document(&req.name).await;
|
let _ = idx.delete_document(&req.name);
|
||||||
Ok("Snippet deleted.".to_string())
|
Ok("Snippet deleted.".to_string())
|
||||||
} else {
|
} else {
|
||||||
Ok("Snippet not found.".to_string())
|
Ok("Snippet not found.".to_string())
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ pub struct MemoryIndex {
|
|||||||
index: Index,
|
index: Index,
|
||||||
reader: IndexReader,
|
reader: IndexReader,
|
||||||
writer: Arc<Mutex<IndexWriter>>,
|
writer: Arc<Mutex<IndexWriter>>,
|
||||||
|
needs_commit: Arc<std::sync::atomic::AtomicBool>,
|
||||||
|
|
||||||
// Schema fields
|
// Schema fields
|
||||||
pub id_field: Field,
|
pub id_field: Field,
|
||||||
@@ -47,6 +48,7 @@ impl MemoryIndex {
|
|||||||
index,
|
index,
|
||||||
reader,
|
reader,
|
||||||
writer: Arc::new(Mutex::new(writer)),
|
writer: Arc::new(Mutex::new(writer)),
|
||||||
|
needs_commit: Arc::new(std::sync::atomic::AtomicBool::new(false)),
|
||||||
id_field,
|
id_field,
|
||||||
title_field,
|
title_field,
|
||||||
body_field,
|
body_field,
|
||||||
@@ -59,6 +61,7 @@ impl MemoryIndex {
|
|||||||
let writer = Arc::clone(&self.writer);
|
let writer = Arc::clone(&self.writer);
|
||||||
let id_field = self.id_field;
|
let id_field = self.id_field;
|
||||||
let id_val = e.name.clone();
|
let id_val = e.name.clone();
|
||||||
|
let needs_commit = Arc::clone(&self.needs_commit);
|
||||||
|
|
||||||
let doc = doc!(
|
let doc = doc!(
|
||||||
self.id_field => e.name.as_str(),
|
self.id_field => e.name.as_str(),
|
||||||
@@ -72,6 +75,7 @@ impl MemoryIndex {
|
|||||||
let writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
let writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
||||||
writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val));
|
writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val));
|
||||||
writer.add_document(doc)?;
|
writer.add_document(doc)?;
|
||||||
|
needs_commit.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||||
Ok(())
|
Ok(())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -80,6 +84,7 @@ impl MemoryIndex {
|
|||||||
let writer = Arc::clone(&self.writer);
|
let writer = Arc::clone(&self.writer);
|
||||||
let id_field = self.id_field;
|
let id_field = self.id_field;
|
||||||
let id_val = t.id.clone();
|
let id_val = t.id.clone();
|
||||||
|
let needs_commit = Arc::clone(&self.needs_commit);
|
||||||
|
|
||||||
let doc = doc!(
|
let doc = doc!(
|
||||||
self.id_field => t.id.as_str(),
|
self.id_field => t.id.as_str(),
|
||||||
@@ -93,6 +98,7 @@ impl MemoryIndex {
|
|||||||
let writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
let writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
||||||
writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val));
|
writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val));
|
||||||
writer.add_document(doc)?;
|
writer.add_document(doc)?;
|
||||||
|
needs_commit.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||||
Ok(())
|
Ok(())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -101,19 +107,24 @@ impl MemoryIndex {
|
|||||||
let writer = Arc::clone(&self.writer);
|
let writer = Arc::clone(&self.writer);
|
||||||
let id_field = self.id_field;
|
let id_field = self.id_field;
|
||||||
let id_val = id.to_string();
|
let id_val = id.to_string();
|
||||||
|
let needs_commit = Arc::clone(&self.needs_commit);
|
||||||
|
|
||||||
tokio::task::spawn_blocking(move || {
|
tokio::task::spawn_blocking(move || {
|
||||||
let writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
let writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
||||||
writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val));
|
writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val));
|
||||||
|
needs_commit.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||||
Ok(())
|
Ok(())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn commit(&self) -> tantivy::Result<()> {
|
pub async fn commit(&self) -> tantivy::Result<()> {
|
||||||
let writer = Arc::clone(&self.writer);
|
let writer = Arc::clone(&self.writer);
|
||||||
|
let needs_commit = Arc::clone(&self.needs_commit);
|
||||||
tokio::task::spawn_blocking(move || {
|
tokio::task::spawn_blocking(move || {
|
||||||
|
if needs_commit.swap(false, std::sync::atomic::Ordering::SeqCst) {
|
||||||
let mut writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
let mut writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
||||||
writer.commit()?;
|
writer.commit()?;
|
||||||
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
})
|
})
|
||||||
.await
|
.await
|
||||||
@@ -182,6 +193,7 @@ impl MemoryIndex {
|
|||||||
let writer = Arc::clone(&self.writer);
|
let writer = Arc::clone(&self.writer);
|
||||||
let id_field = self.id_field;
|
let id_field = self.id_field;
|
||||||
let id_val = s.name.clone();
|
let id_val = s.name.clone();
|
||||||
|
let needs_commit = Arc::clone(&self.needs_commit);
|
||||||
|
|
||||||
let doc = doc!(
|
let doc = doc!(
|
||||||
self.id_field => s.name.as_str(),
|
self.id_field => s.name.as_str(),
|
||||||
@@ -195,6 +207,7 @@ impl MemoryIndex {
|
|||||||
let writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
let writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
||||||
writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val));
|
writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val));
|
||||||
writer.add_document(doc)?;
|
writer.add_document(doc)?;
|
||||||
|
needs_commit.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||||
Ok(())
|
Ok(())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -203,6 +216,7 @@ impl MemoryIndex {
|
|||||||
let writer = Arc::clone(&self.writer);
|
let writer = Arc::clone(&self.writer);
|
||||||
let id_field = self.id_field;
|
let id_field = self.id_field;
|
||||||
let id_val = a.id.clone();
|
let id_val = a.id.clone();
|
||||||
|
let needs_commit = Arc::clone(&self.needs_commit);
|
||||||
|
|
||||||
let doc = doc!(
|
let doc = doc!(
|
||||||
self.id_field => a.id.as_str(),
|
self.id_field => a.id.as_str(),
|
||||||
@@ -216,6 +230,7 @@ impl MemoryIndex {
|
|||||||
let writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
let writer = writer.lock().unwrap_or_else(|e| e.into_inner());
|
||||||
writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val));
|
writer.delete_term(tantivy::Term::from_field_text(id_field, &id_val));
|
||||||
writer.add_document(doc)?;
|
writer.add_document(doc)?;
|
||||||
|
needs_commit.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||||
Ok(())
|
Ok(())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
+3
-2
@@ -56,8 +56,9 @@ impl MemoryState {
|
|||||||
});
|
});
|
||||||
|
|
||||||
let payload = serde_json::json!({
|
let payload = serde_json::json!({
|
||||||
"type": "activity",
|
"jsonrpc": "2.0",
|
||||||
"data": item
|
"method": "notifications/activity",
|
||||||
|
"params": item
|
||||||
})
|
})
|
||||||
.to_string();
|
.to_string();
|
||||||
let _ = self.activity_tx.send(payload);
|
let _ = self.activity_tx.send(payload);
|
||||||
|
|||||||
Reference in new issue
Block a user