diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index c0909db..77ee661 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -58,6 +58,27 @@ jobs: ${{ matrix.name }}.tar.gz ${{ matrix.name }}.tar.gz.sha256 + checksums: + needs: build + runs-on: ubuntu-latest + steps: + - name: Download SHA256 files + env: + GH_TOKEN: ${{ secrets.GITHUB_TOKEN }} + run: | + gh release download ${{ github.ref_name }} \ + --repo ${{ github.repository }} \ + --pattern "*.sha256" + + - name: Generate checksums.txt + run: | + cat *.sha256 > checksums.txt + + - name: Upload checksums.txt + uses: softprops/action-gh-release@v1 + with: + files: checksums.txt + publish-crate: needs: build runs-on: ubuntu-latest diff --git a/CLAUDE.md b/CLAUDE.md index 8f620a0..1593dcf 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -113,13 +113,15 @@ CREATE TABLE graph_nodes ( -- Graph edges (RELATED, SUPERSEDES, CO_ACCESSED, CLUSTER_SIBLING) CREATE TABLE graph_edges ( - source_id TEXT, target_id TEXT, - rel_type TEXT, -- 'related', 'supersedes', 'co_accessed', 'cluster_sibling' - strength REAL, - count INTEGER, - label TEXT, - created_at INTEGER, - UNIQUE(source_id, target_id, rel_type) + from_id TEXT NOT NULL, + to_id TEXT NOT NULL, + edge_type TEXT NOT NULL, -- 'related', 'supersedes', 'co_accessed', 'cluster_sibling' + strength REAL NOT NULL DEFAULT 0.0, + rel_type TEXT NOT NULL DEFAULT '', + count INTEGER NOT NULL DEFAULT 0, + last_at INTEGER NOT NULL DEFAULT 0, + created_at INTEGER NOT NULL, + UNIQUE(from_id, to_id, edge_type) ); -- Access tracking and decay scoring diff --git a/Cargo.toml b/Cargo.toml index 85054a3..e4165fa 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -38,10 +38,7 @@ serde = { version = "1.0", features = ["derive"] } serde_json = "1.0" # Async -tokio = { version = "1", features = ["full"] } - -# MCP -jsonrpc-core = "18" +tokio = { version = "1", features = ["rt", "rt-multi-thread", "sync", "time", "macros"] } # CLI clap = { version = "4", features = ["derive", "env"] } @@ -53,6 +50,7 @@ rusqlite = { version = "0.31", features = ["bundled"] } # Utils uuid = { version = "1", features = ["v4", "serde"] } chrono = { version = "0.4", features = ["serde"] } +sha2 = "0.10" thiserror = "2" anyhow = "1" tracing = "0.1" diff --git a/install.sh b/install.sh index 75b9671..cab5909 100755 --- a/install.sh +++ b/install.sh @@ -46,6 +46,33 @@ TMPDIR="$(mktemp -d)" trap 'rm -rf "$TMPDIR"' EXIT curl -fsSL "$URL" -o "${TMPDIR}/${TARBALL}" + +# Verify SHA256 checksum if checksums.txt is available +CHECKSUMS_URL="https://github.com/${REPO}/releases/download/v${VERSION}/checksums.txt" +if curl -fsSL "$CHECKSUMS_URL" -o "${TMPDIR}/checksums.txt" 2>/dev/null; then + EXPECTED="$(grep "${TARBALL}" "${TMPDIR}/checksums.txt" | awk '{print $1}')" + if [ -n "$EXPECTED" ]; then + if command -v sha256sum &>/dev/null; then + ACTUAL="$(sha256sum "${TMPDIR}/${TARBALL}" | awk '{print $1}')" + elif command -v shasum &>/dev/null; then + ACTUAL="$(shasum -a 256 "${TMPDIR}/${TARBALL}" | awk '{print $1}')" + else + echo "Warning: no sha256sum or shasum found, skipping checksum verification" >&2 + ACTUAL="" + fi + if [ -n "$ACTUAL" ] && [ "$ACTUAL" != "$EXPECTED" ]; then + echo "Checksum verification failed!" >&2 + echo " Expected: $EXPECTED" >&2 + echo " Actual: $ACTUAL" >&2 + exit 1 + elif [ -n "$ACTUAL" ]; then + echo "Checksum verified." + fi + fi +else + echo "No checksums.txt available, skipping verification." +fi + tar -xzf "${TMPDIR}/${TARBALL}" -C "$TMPDIR" # Install diff --git a/src/access.rs b/src/access.rs index 80bc3d0..e941c0a 100644 --- a/src/access.rs +++ b/src/access.rs @@ -30,6 +30,8 @@ impl AccessTracker { SedimentError::Database(format!("Failed to open access database: {}", e)) })?; + conn.execute_batch("PRAGMA journal_mode=WAL;").ok(); + conn.execute_batch( "CREATE TABLE IF NOT EXISTS access_log ( item_id TEXT PRIMARY KEY, diff --git a/src/chunker.rs b/src/chunker.rs index e87ff92..0329cfa 100644 --- a/src/chunker.rs +++ b/src/chunker.rs @@ -111,40 +111,40 @@ pub fn chunk_content( /// Split text at sentence boundaries fn split_at_sentences(text: &str) -> Vec<&str> { let mut sentences = Vec::new(); - let bytes = text.as_bytes(); let mut start = 0; - let mut i = 0; - - while i < bytes.len() { - // Look for sentence-ending punctuation followed by space or end - if matches!(bytes[i], b'.' | b'?' | b'!') { - let next_idx = i + 1; - if next_idx >= bytes.len() - || bytes[next_idx] == b' ' - || bytes[next_idx] == b'\n' - || bytes[next_idx] == b'\t' - { - // Include the punctuation - let end = i + 1; - if start < end && end <= bytes.len() { + let mut char_indices = text.char_indices().peekable(); + + while let Some((i, ch)) = char_indices.next() { + // Look for sentence-ending punctuation (ASCII and Unicode) + if matches!(ch, '.' | '?' | '!' | '。' | '?' | '!') { + let end = i + ch.len_utf8(); + // Check if followed by whitespace or end of text + let at_end_or_ws = match char_indices.peek() { + None => true, + Some(&(_, next_ch)) => next_ch == ' ' || next_ch == '\n' || next_ch == '\t', + }; + if at_end_or_ws { + if start < end { sentences.push(&text[start..end]); } // Skip whitespace after punctuation - i += 1; - while i < bytes.len() - && (bytes[i] == b' ' || bytes[i] == b'\n' || bytes[i] == b'\t') - { - i += 1; + while let Some(&(_, next_ch)) = char_indices.peek() { + if next_ch == ' ' || next_ch == '\n' || next_ch == '\t' { + char_indices.next(); + } else { + break; + } } - start = i; - continue; + start = match char_indices.peek() { + Some(&(idx, _)) => idx, + None => text.len(), + }; } } - i += 1; } // Add remaining text - if start < bytes.len() { + if start < text.len() { sentences.push(&text[start..]); } diff --git a/src/consolidation.rs b/src/consolidation.rs index 1175fee..c2e9abd 100644 --- a/src/consolidation.rs +++ b/src/consolidation.rs @@ -32,6 +32,8 @@ impl ConsolidationQueue { SedimentError::Database(format!("Failed to open consolidation database: {}", e)) })?; + conn.execute_batch("PRAGMA journal_mode=WAL;").ok(); + conn.execute_batch( "CREATE TABLE IF NOT EXISTS consolidation_queue ( item_id_a TEXT NOT NULL, @@ -180,7 +182,11 @@ async fn run_consolidation_batch( let result = process_candidate(&mut db, access_db_path, candidate).await; match result { Ok(status) => { - let _ = queue.mark_processed(&candidate.item_id_a, &candidate.item_id_b, &status); + if let Err(e) = + queue.mark_processed(&candidate.item_id_a, &candidate.item_id_b, &status) + { + tracing::warn!("mark_processed failed: {}", e); + } info!( "Consolidated {} <-> {}: {} (similarity: {:.2})", candidate.item_id_a, candidate.item_id_b, status, candidate.similarity @@ -220,26 +226,38 @@ async fn process_candidate( }; // Transfer edges from old to new - let _ = graph.transfer_edges(&remove.id, &keep.id); + if let Err(e) = graph.transfer_edges(&remove.id, &keep.id) { + tracing::warn!("transfer_edges failed: {}", e); + } // Create SUPERSEDES edge (preserves lineage for recovery) - let _ = graph.add_supersedes_edge(&keep.id, &remove.id); + if let Err(e) = graph.add_supersedes_edge(&keep.id, &remove.id) { + tracing::warn!("add_supersedes_edge failed: {}", e); + } // 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]) + let archive_preview = if remove.content.chars().count() > 500 { + let cut = remove + .content + .char_indices() + .nth(497) + .map(|(i, _)| i) + .unwrap_or(remove.content.len()); + format!("{}...", &remove.content[..cut]) } else { remove.content.clone() }; - let _ = graph.add_related_edge( + if let Err(e) = graph.add_related_edge( &keep.id, &remove.id, candidate.similarity, &format!("merged_archive:{}", archive_preview), - ); + ) { + tracing::warn!("add_related_edge failed: {}", e); + } // Soft-delete: mark item as expired instead of hard-deleting. // This allows recovery; expired items are excluded from search @@ -249,7 +267,9 @@ async fn process_candidate( // 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); + if let Err(e) = graph.remove_node(&remove.id) { + tracing::warn!("remove_node failed: {}", e); + } } Ok("merged".to_string()) @@ -261,12 +281,14 @@ async fn process_candidate( } } else { // Similar but distinct: create RELATED edge - let _ = graph.add_related_edge( + if let Err(e) = graph.add_related_edge( &candidate.item_id_a, &candidate.item_id_b, candidate.similarity, "similar", - ); + ) { + tracing::warn!("add_related_edge failed: {}", e); + } Ok("linked".to_string()) } } diff --git a/src/db.rs b/src/db.rs index d09cc44..75cbc4f 100644 --- a/src/db.rs +++ b/src/db.rs @@ -24,6 +24,7 @@ use lancedb::connect; use lancedb::query::{ExecutableQuery, QueryBase}; use tracing::{debug, info}; +use crate::boost_similarity; use crate::chunker::{ChunkingConfig, chunk_content}; use crate::document::ContentType; use crate::embedder::{EMBEDDING_DIM, Embedder}; @@ -288,16 +289,6 @@ impl Database { item.project_id = self.project_id.clone(); } - // Check for potential conflicts before storing - let potential_conflicts = self - .find_similar_items( - &item.content, - CONFLICT_SIMILARITY_THRESHOLD, - CONFLICT_SEARCH_LIMIT, - ) - .await - .unwrap_or_default(); - // Determine if we need to chunk let should_chunk = item.content.len() > CHUNK_THRESHOLD; item.is_chunked = should_chunk; @@ -360,6 +351,19 @@ impl Database { debug!("Stored item: {} (no chunking)", item.id); } + // Detect conflicts after storing (informational only, avoids TOCTOU race) + let potential_conflicts = self + .find_similar_items( + &item.content, + CONFLICT_SIMILARITY_THRESHOLD, + CONFLICT_SEARCH_LIMIT, + ) + .await + .unwrap_or_default() + .into_iter() + .filter(|c| c.id != item.id) + .collect(); + Ok(StoreResult { id: item.id, potential_conflicts, @@ -751,7 +755,7 @@ impl Database { // Delete then re-insert with new expiration table - .delete(&format!("id = '{}'", id)) + .delete(&format!("id = '{}'", sanitize_sql_string(id))) .await .map_err(|e| SedimentError::Database(format!("Delete for expire failed: {}", e)))?; @@ -810,6 +814,31 @@ impl Database { Ok(stats) } + + /// Delete items whose expires_at timestamp is in the past. + pub async fn cleanup_expired(&self) -> Result { + let table = match &self.items_table { + Some(t) => t, + None => return Ok(0), + }; + + let now = Utc::now().timestamp(); + let filter = format!("expires_at IS NOT NULL AND expires_at < {}", now); + + // Count how many will be deleted + let count = table.count_rows(Some(filter.clone())).await.unwrap_or(0); + + if count > 0 { + table + .delete(&filter) + .await + .map_err(|e| SedimentError::Database(format!("Expired cleanup failed: {}", e)))?; + + info!("Cleaned up {} expired items", count); + } + + Ok(count) + } } // ==================== Decay Scoring ==================== @@ -841,15 +870,6 @@ pub fn score_with_decay( // ==================== Helper Functions ==================== -/// Apply similarity boosting based on project context. -fn boost_similarity(base: f32, item_project: Option<&str>, current_project: Option<&str>) -> f32 { - match (item_project, current_project) { - (Some(m), Some(c)) if m == c => (base * 1.15).min(1.0), // Same project: boost - (Some(_), Some(_)) => base * 0.95, // Different project: slight penalty - _ => base, // Global or no context - } -} - /// Detect content type for smart chunking fn detect_content_type(content: &str) -> ContentType { let trimmed = content.trim(); @@ -1037,12 +1057,20 @@ fn batch_to_items(batch: &RecordBatch) -> Result> { if c.is_null(i) { None } else { - Some(Utc.timestamp_opt(c.value(i), 0).unwrap()) + Some( + Utc.timestamp_opt(c.value(i), 0) + .single() + .unwrap_or_else(Utc::now), + ) } }); let created_at = created_at_col - .map(|c| Utc.timestamp_opt(c.value(i), 0).unwrap()) + .map(|c| { + Utc.timestamp_opt(c.value(i), 0) + .single() + .unwrap_or_else(Utc::now) + }) .unwrap_or_else(Utc::now); let item = Item { diff --git a/src/embedder.rs b/src/embedder.rs index dfb9e4f..1981cf6 100644 --- a/src/embedder.rs +++ b/src/embedder.rs @@ -1,9 +1,10 @@ -use std::path::PathBuf; +use std::path::{Path, PathBuf}; use candle_core::{DType, Device, Tensor}; use candle_nn::VarBuilder; use candle_transformers::models::bert::{BertModel, Config, DTYPE}; use hf_hub::{Repo, RepoType, api::sync::ApiBuilder}; +use sha2::{Digest, Sha256}; use tokenizers::{PaddingParams, Tokenizer, TruncationParams}; use tracing::info; @@ -194,7 +195,11 @@ fn download_model(model_id: &str) -> Result<(PathBuf, PathBuf, PathBuf)> { .build() .map_err(|e| SedimentError::ModelLoading(format!("Failed to create HF API: {}", e)))?; - let repo = api.repo(Repo::new(model_id.to_string(), RepoType::Model)); + let repo = api.repo(Repo::with_revision( + model_id.to_string(), + RepoType::Model, + "e4ce9877abf3edfe10b0d82785e83bdcb973e22e".to_string(), + )); let model_path = repo .get("model.safetensors") @@ -208,9 +213,49 @@ fn download_model(model_id: &str) -> Result<(PathBuf, PathBuf, PathBuf)> { .get("config.json") .map_err(|e| SedimentError::ModelLoading(format!("Failed to download config: {}", e)))?; + // TOFU: verify model integrity + verify_tofu_hash(&model_path, "model.safetensors")?; + Ok((model_path, tokenizer_path, config_path)) } +/// Compute SHA256 hash of a file. +fn sha256_file(path: &Path) -> Result { + let data = std::fs::read(path) + .map_err(|e| SedimentError::ModelLoading(format!("Failed to read file for hash: {}", e)))?; + let hash = Sha256::digest(&data); + Ok(format!("{:x}", hash)) +} + +/// Trust-on-first-use hash verification for model files. +/// +/// On first download, stores the SHA256 hash in a `.sha256` sidecar file next to the model. +/// On subsequent loads, verifies the file still matches the stored hash. +fn verify_tofu_hash(file_path: &Path, label: &str) -> Result<()> { + let hash_path = file_path.with_extension("sha256"); + let current_hash = sha256_file(file_path)?; + + if hash_path.exists() { + let stored_hash = std::fs::read_to_string(&hash_path) + .map_err(|e| SedimentError::ModelLoading(format!("Failed to read hash file: {}", e)))?; + let stored_hash = stored_hash.trim(); + if stored_hash != current_hash { + return Err(SedimentError::ModelLoading(format!( + "TOFU integrity check failed for {}: expected {}, got {}", + label, stored_hash, current_hash + ))); + } + info!("TOFU hash verified for {}", label); + } else { + std::fs::write(&hash_path, ¤t_hash).map_err(|e| { + SedimentError::ModelLoading(format!("Failed to write hash file: {}", e)) + })?; + info!("TOFU hash recorded for {}: {}", label, current_hash); + } + + Ok(()) +} + /// L2 normalize a tensor fn normalize_l2(tensor: &Tensor) -> Result { let norm = tensor diff --git a/src/lib.rs b/src/lib.rs index 2531b59..a0c0213 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,7 +6,7 @@ //! //! - **Embedded storage** - LanceDB-powered, directory-based, no server required //! - **Local embeddings** - Uses `all-MiniLM-L6-v2` locally, no API keys needed -//! - **MCP-native** - 4 tools for seamless LLM integration +//! - **MCP-native** - 5 tools for seamless LLM integration //! - **Project-aware** - Scoped memories with automatic project detection //! - **Auto-chunking** - Long content is automatically chunked for better search @@ -168,7 +168,7 @@ pub fn get_or_create_project_id(project_root: &Path) -> std::io::Result // Save config let content = serde_json::to_string_pretty(&config).map_err(|e| std::io::Error::other(e.to_string()))?; - std::fs::write(&config_path, content)?; + std::fs::write(&config_path, &content)?; Ok(config.project_id) } @@ -202,7 +202,13 @@ pub fn find_project_root(start: &Path) -> Option { current = current.parent()?.to_path_buf(); } + let mut depth = 0; loop { + if depth >= 100 { + return None; + } + depth += 1; + // Check for .sediment directory first (explicit project marker) if current.join(".sediment").is_dir() { return Some(current); @@ -213,8 +219,9 @@ pub fn find_project_root(start: &Path) -> Option { return Some(current); } - // Move to parent directory + // Move to parent directory; stop at filesystem root match current.parent() { + Some(parent) if parent == current => return None, Some(parent) => current = parent.to_path_buf(), None => return None, } diff --git a/src/main.rs b/src/main.rs index dcf468b..bc0b397 100644 --- a/src/main.rs +++ b/src/main.rs @@ -336,7 +336,11 @@ fn run_list(db_override: Option, limit: usize) -> Result<()> { .take(80) .collect::() .replace('\n', " "); - let ellipsis = if item.content.len() > 80 { "..." } else { "" }; + let ellipsis = if item.content.chars().count() > 80 { + "..." + } else { + "" + }; println!(" Content: {}{}", content_preview, ellipsis); println!(); } @@ -351,12 +355,13 @@ fn generate_claude_md_instructions() -> String { Use the Sediment MCP tools for persistent memory storage. -## Tools (4 total) +## Tools (5 total) - `mcp__sediment__store` - Store content for later retrieval - `mcp__sediment__recall` - Search by semantic similarity - `mcp__sediment__list` - List stored items - `mcp__sediment__forget` - Delete an item by ID +- `mcp__sediment__connections` - Show relationship graph for an item ## When to Store diff --git a/src/mcp/server.rs b/src/mcp/server.rs index 4b0c8a0..5dadbc4 100644 --- a/src/mcp/server.rs +++ b/src/mcp/server.rs @@ -5,6 +5,7 @@ use serde_json::{Value, json}; use std::io::{self, BufRead, Write}; use std::path::{Path, PathBuf}; use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; use tokio::runtime::Runtime; use tokio::sync::Semaphore; @@ -33,6 +34,10 @@ pub struct ServerContext { pub consolidation_semaphore: Arc, /// Counter for consolidation runs (for periodic clustering) pub consolidation_run_count: std::sync::atomic::AtomicU64, + /// Rate limiter: start of the current 60-second window (unix timestamp in millis) + pub rate_limit_window: AtomicU64, + /// Rate limiter: number of tool calls in the current window + pub rate_limit_count: AtomicU64, } /// Run the MCP server @@ -61,6 +66,8 @@ pub fn run(db_path: &Path, project_id: Option) -> Result<()> { cwd, consolidation_semaphore: Arc::new(Semaphore::new(1)), consolidation_run_count: std::sync::atomic::AtomicU64::new(0), + rate_limit_window: AtomicU64::new(0), + rate_limit_count: AtomicU64::new(0), }; let stdin = io::stdin(); @@ -191,6 +198,28 @@ fn handle_call_tool( tracing::info!("Calling tool: {}", params.name); + // Rate limiting: 60 calls per 60-second window + { + const MAX_CALLS_PER_MINUTE: u64 = 60; + let now_ms = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as u64; + let window = ctx.rate_limit_window.load(Ordering::Relaxed); + if now_ms - window > 60_000 { + // Reset window + ctx.rate_limit_window.store(now_ms, Ordering::Relaxed); + ctx.rate_limit_count.store(1, Ordering::Relaxed); + } else { + let count = ctx.rate_limit_count.fetch_add(1, Ordering::Relaxed) + 1; + if count > MAX_CALLS_PER_MINUTE { + let result = + super::protocol::CallToolResult::error("Rate limit exceeded, try again later"); + return Response::success(id, serde_json::to_value(result).unwrap()); + } + } + } + // Execute tool async with fresh DB connection and retry logic let result = rt.block_on(execute_tool(ctx, ¶ms.name, params.arguments)); diff --git a/src/mcp/tools.rs b/src/mcp/tools.rs index 3bf8b18..0002e7d 100644 --- a/src/mcp/tools.rs +++ b/src/mcp/tools.rs @@ -7,6 +7,7 @@ use std::sync::Arc; use chrono::DateTime; use serde::Deserialize; use serde_json::{Value, json}; +use tracing::Instrument; use crate::access::AccessTracker; use crate::consolidation::{ConsolidationQueue, spawn_consolidation}; @@ -455,22 +456,34 @@ async fn execute_store( // Create graph node let now = chrono::Utc::now().timestamp(); let project_id = db.project_id().map(|s| s.to_string()); - let _ = graph.add_node(&new_id, project_id.as_deref(), now); + if let Err(e) = graph.add_node(&new_id, project_id.as_deref(), now) { + tracing::warn!("graph add_node failed: {}", e); + } // 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 { - let _ = db.delete_item(old_id).await; + if let Err(e) = db.delete_item(old_id).await { + tracing::warn!("delete_item failed: {}", e); + } 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); + if let Err(e) = tracker.record_validation(old_id, now_ts) { + tracing::warn!("record_validation failed: {}", e); + } + if let Err(e) = graph.add_supersedes_edge(&new_id, old_id) { + tracing::warn!("add_supersedes_edge failed: {}", e); + } + if let Err(e) = graph.remove_node(old_id) { + tracing::warn!("remove_node failed: {}", e); + } } // Create RELATED edges if specified if let Some(ref related_ids) = params.related { for rid in related_ids { - let _ = graph.add_related_edge(&new_id, rid, 1.0, "user_linked"); + if let Err(e) = graph.add_related_edge(&new_id, rid, 1.0, "user_linked") { + tracing::warn!("add_related_edge failed: {}", e); + } } } @@ -479,7 +492,10 @@ async fn execute_store( && let Ok(queue) = ConsolidationQueue::open(&ctx.access_db_path) { for conflict in &store_result.potential_conflicts { - let _ = queue.enqueue(&new_id, &conflict.id, conflict.similarity as f64); + if let Err(e) = queue.enqueue(&new_id, &conflict.id, conflict.similarity as f64) + { + tracing::warn!("enqueue consolidation failed: {}", e); + } } } @@ -539,11 +555,13 @@ pub async fn recall_pipeline( // Lazy graph backfill (uses project_id from SearchResult, no extra queries) if config.enable_graph_backfill { for result in &results { - let _ = graph.ensure_node_exists( + if let Err(e) = graph.ensure_node_exists( &result.id, result.project_id.as_deref(), result.created_at.timestamp(), - ); + ) { + tracing::warn!("ensure_node_exists failed: {}", e); + } } } @@ -573,7 +591,7 @@ pub async fn recall_pipeline( let trust_bonus = 1.0 + 0.05 * (1.0 + validation_count as f64).ln() as f32 + 0.02 * edge_count as f32; - result.similarity = base_score * trust_bonus; + result.similarity = (base_score * trust_bonus).min(1.0); } results.sort_by(|a, b| b.similarity.partial_cmp(&a.similarity).unwrap()); @@ -582,7 +600,9 @@ pub async fn recall_pipeline( // Record access for result in &results { let created_at = result.created_at.timestamp(); - let _ = tracker.record_access(&result.id, created_at); + if let Err(e) = tracker.record_access(&result.id, created_at) { + tracing::warn!("record_access failed: {}", e); + } } // Graph expansion @@ -766,33 +786,68 @@ async fn execute_recall( // Fire-and-forget: co-access recording (Phase 3a) let result_ids: Vec = results.iter().map(|r| r.id.clone()).collect(); let access_db_path = ctx.access_db_path.clone(); - tokio::spawn(async move { - if let Ok(g) = GraphStore::open(&access_db_path) { - let _ = g.record_co_access(&result_ids); + tokio::spawn( + async move { + if let Ok(g) = GraphStore::open(&access_db_path) { + if let Err(e) = g.record_co_access(&result_ids) { + tracing::warn!("record_co_access failed: {}", e); + } + } else { + tracing::warn!("co_access: failed to open graph store"); + } } - }); + .instrument(tracing::info_span!("co_access")), + ); - // Periodic clustering (Phase 4b): every 10th consolidation run + // Periodic expired item cleanup: every 10th recall let run_count = ctx .consolidation_run_count .fetch_add(1, std::sync::atomic::Ordering::Relaxed); if run_count % 10 == 9 { + // Clustering let access_db_path = ctx.access_db_path.clone(); - tokio::spawn(async move { - if let Ok(g) = GraphStore::open(&access_db_path) - && let Ok(clusters) = g.detect_clusters() - { - for (a, b, c) in &clusters { - let label = format!("cluster-{}", &a[..8.min(a.len())]); - let _ = g.add_related_edge(a, b, 0.8, &label); - let _ = g.add_related_edge(b, c, 0.8, &label); - let _ = g.add_related_edge(a, c, 0.8, &label); + tokio::spawn( + async move { + if let Ok(g) = GraphStore::open(&access_db_path) + && let Ok(clusters) = g.detect_clusters() + { + for (a, b, c) in &clusters { + let label = format!("cluster-{}", &a[..8.min(a.len())]); + if let Err(e) = g.add_related_edge(a, b, 0.8, &label) { + tracing::warn!("cluster add_related_edge failed: {}", e); + } + if let Err(e) = g.add_related_edge(b, c, 0.8, &label) { + tracing::warn!("cluster add_related_edge failed: {}", e); + } + if let Err(e) = g.add_related_edge(a, c, 0.8, &label) { + tracing::warn!("cluster add_related_edge failed: {}", e); + } + } + if !clusters.is_empty() { + tracing::info!("Detected {} clusters", clusters.len()); + } } - if !clusters.is_empty() { - tracing::info!("Detected {} clusters", clusters.len()); + } + .instrument(tracing::info_span!("clustering")), + ); + + // Expired item cleanup + let db_path = ctx.db_path.clone(); + let project_id = ctx.project_id.clone(); + let embedder = ctx.embedder.clone(); + tokio::spawn( + async move { + match Database::open_with_embedder(&db_path, project_id, embedder).await { + Ok(db) => { + if let Err(e) = db.cleanup_expired().await { + tracing::warn!("cleanup_expired failed: {}", e); + } + } + Err(e) => tracing::warn!("cleanup_expired: failed to open db: {}", e), } } - }); + .instrument(tracing::info_span!("cleanup_expired")), + ); } CallToolResult::success(serde_json::to_string_pretty(&result_json).unwrap()) @@ -883,7 +938,9 @@ async fn execute_forget( match db.delete_item(¶ms.id).await { Ok(true) => { // Remove from graph - let _ = graph.remove_node(¶ms.id); + if let Err(e) = graph.remove_node(¶ms.id) { + tracing::warn!("remove_node failed: {}", e); + } let result = json!({ "success": true, @@ -959,9 +1016,14 @@ async fn execute_connections( // ========== Utilities ========== fn truncate(s: &str, max_len: usize) -> String { - if s.len() <= max_len { + if s.chars().count() <= max_len { s.to_string() } else { - format!("{}...", &s[..max_len - 3]) + let cut = s + .char_indices() + .nth(max_len - 3) + .map(|(i, _)| i) + .unwrap_or(s.len()); + format!("{}...", &s[..cut]) } }