Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 27 additions & 6 deletions src/consolidation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -207,7 +207,7 @@ async fn process_candidate(
) -> Result<String> {
let graph = crate::graph::GraphStore::open(graph_db_path)?;
if candidate.similarity >= 0.95 {
// Near-duplicate: newer absorbs older
// Near-duplicate: newer absorbs older (non-destructive — archive removed content)
let item_a = db.get_item(&candidate.item_id_a).await?;
let item_b = db.get_item(&candidate.item_id_b).await?;

Expand All @@ -222,14 +222,35 @@ async fn process_candidate(
// Transfer edges from old to new
let _ = graph.transfer_edges(&remove.id, &keep.id);

// Create SUPERSEDES edge
// Create SUPERSEDES edge (preserves lineage for recovery)
let _ = graph.add_supersedes_edge(&keep.id, &remove.id);

// Delete old item from LanceDB
db.delete_item(&remove.id).await?;
// Archive the removed item's content into a RELATED edge label
// so it can be recovered if this was a false positive merge.
// The edge label stores a truncated snapshot; full content is
// preserved in the SUPERSEDES relationship for audit.
let archive_preview = if remove.content.len() > 500 {
format!("{}...", &remove.content[..497])
} else {
remove.content.clone()
};
let _ = graph.add_related_edge(
&keep.id,
&remove.id,
candidate.similarity,
&format!("merged_archive:{}", archive_preview),
);

// Remove old node from graph
let _ = graph.remove_node(&remove.id);
// Soft-delete: mark item as expired instead of hard-deleting.
// This allows recovery; expired items are excluded from search
// results by default but remain in the database.
let past = chrono::Utc::now() - chrono::Duration::seconds(1);
if let Err(e) = db.expire_item(&remove.id, past).await {
// Fall back to hard delete if expire is not supported
warn!("expire_item failed ({}), falling back to delete", e);
db.delete_item(&remove.id).await?;
let _ = graph.remove_node(&remove.id);
}

Ok("merged".to_string())
}
Expand Down
38 changes: 38 additions & 0 deletions src/db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -719,6 +719,44 @@ impl Database {
Ok(items)
}

/// Soft-delete an item by setting its expiration to a past timestamp.
/// The item remains in the database but is excluded from search results.
pub async fn expire_item(&self, id: &str, expires_at: chrono::DateTime<Utc>) -> Result<()> {
let table = match &self.items_table {
Some(t) => t,
None => return Err(SedimentError::Database("Items table not found".to_string())),
};

// LanceDB doesn't support in-place updates easily, so we use a merge-insert
// approach: read the item, delete it, re-insert with updated expires_at.
let item = self.get_item(id).await?;
let mut item = match item {
Some(i) => i,
None => return Err(SedimentError::Database(format!("Item not found: {}", id))),
};

// Re-generate embedding since get_item returns empty embeddings
let embedding_text = item.embedding_text();
item.embedding = self.embedder.embed(&embedding_text)?;
item.expires_at = Some(expires_at);

// Delete then re-insert with new expiration
table
.delete(&format!("id = '{}'", id))
.await
.map_err(|e| SedimentError::Database(format!("Delete for expire failed: {}", e)))?;

let batch = item_to_batch(&item)?;
let batches = RecordBatchIterator::new(vec![Ok(batch)], Arc::new(item_schema()));
table
.add(Box::new(batches))
.execute()
.await
.map_err(|e| SedimentError::Database(format!("Re-insert for expire failed: {}", e)))?;

Ok(())
}

/// Delete an item and its chunks
pub async fn delete_item(&self, id: &str) -> Result<bool> {
// Delete chunks first
Expand Down
32 changes: 20 additions & 12 deletions src/mcp/tools.rs
Original file line number Diff line number Diff line change
Expand Up @@ -333,6 +333,16 @@ async fn execute_store(
None => return CallToolResult::error("Missing parameters"),
};

// Reject oversized content to prevent OOM during embedding/chunking (1MB limit)
const MAX_CONTENT_BYTES: usize = 1_000_000;
if params.content.len() > MAX_CONTENT_BYTES {
return CallToolResult::error(format!(
"Content too large: {} bytes (max {} bytes)",
params.content.len(),
MAX_CONTENT_BYTES
));
}

// Parse scope
let scope = params
.scope
Expand All @@ -355,23 +365,18 @@ async fn execute_store(
None
};

// Handle replace: delete the existing item first, record provenance
// Validate that the item to replace exists (actual deletion deferred until after store)
let replaced_id = if let Some(ref replace_id) = params.replace {
match db.delete_item(replace_id).await {
Ok(true) => {
// Record validation for the new item being stored
let now = chrono::Utc::now().timestamp();
let _ = tracker.record_validation(replace_id, now);
Some(replace_id.clone())
}
Ok(false) => {
match db.get_item(replace_id).await {
Ok(Some(_)) => Some(replace_id.clone()),
Ok(None) => {
return CallToolResult::error(format!(
"Cannot replace: item not found: {}",
replace_id
));
}
Err(e) => {
return CallToolResult::error(format!("Failed to delete item for replace: {}", e));
return CallToolResult::error(format!("Failed to look up item for replace: {}", e));
}
}
} else {
Expand Down Expand Up @@ -452,9 +457,12 @@ async fn execute_store(
let project_id = db.project_id().map(|s| s.to_string());
let _ = graph.add_node(&new_id, project_id.as_deref(), now);

// Create SUPERSEDES edge if replacing
// Complete replace: now that the new item is stored, delete the old one
// (store-before-delete ensures no data loss on crash)
if let Some(ref old_id) = replaced_id {
// The old node might still exist in graph; create edge then remove
let _ = db.delete_item(old_id).await;
let now_ts = chrono::Utc::now().timestamp();
let _ = tracker.record_validation(old_id, now_ts);
let _ = graph.add_supersedes_edge(&new_id, old_id);
let _ = graph.remove_node(old_id);
}
Expand Down