From c5b1cdc2f9ce017991dc89c5757ca03528a9e967 Mon Sep 17 00:00:00 2001 From: "ilblackdragon@gmail.com" Date: Sun, 8 Mar 2026 17:34:11 -0700 Subject: [PATCH] fix: address PR review comments (round 2) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Human reviewer (zmanian): - H1: Accept configurable embedding dimension in LanceDbVectorStore::new() instead of hardcoding 1536. Dimension is sourced from EmbeddingProvider::dimension() at init time. - H2: Skip double-write of embeddings to DB when external vector store is active (pass None to insert_chunk for embedding column). - H3: Update PR title from "refactor" to "feat" (net-new feature). - H4: Document non-atomic update_embedding in struct doc comment. Bot reviewer (Copilot): - Cache LanceDB table handle via tokio::sync::OnceCell (avoid open_table per operation). - Cache Arc in struct (avoid rebuilding per insert). - Fix error variants: ChunkingFailed → EmbeddingFailed for LanceDB store/delete operations. - Propagate store_embedding errors in reindex_document instead of warn-only (prevents silent data loss). - Prefetch document metadata map in backfill_embeddings to avoid N+1 queries. - Add lancedb feature + protoc to CI test matrix so LanceDB tests actually run on Linux. Co-Authored-By: Claude Opus 4.6 --- .github/workflows/test.yml | 7 ++- src/app.rs | 3 +- src/workspace/lancedb_store.rs | 110 +++++++++++++++++---------------- src/workspace/mod.rs | 74 ++++++++++++++-------- tests/lancedb_integration.rs | 2 +- 5 files changed, 114 insertions(+), 82 deletions(-) diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index c0b0d49a..e73ca77f 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -14,7 +14,7 @@ jobs: matrix: include: - name: all-features - flags: "--features postgres,libsql,html-to-markdown" + flags: "--features postgres,libsql,lancedb,html-to-markdown" - name: default flags: "" - name: libsql-only @@ -26,6 +26,11 @@ jobs: uses: dtolnay/rust-toolchain@stable with: targets: wasm32-wasip2 + - name: Install protoc (for lancedb) + if: matrix.name == 'all-features' + uses: arduino/setup-protoc@v3 + with: + repo-token: ${{ secrets.GITHUB_TOKEN }} - uses: Swatinem/rust-cache@v2 with: key: ${{ matrix.name }} diff --git a/src/app.rs b/src/app.rs index 43630c32..cf8e251d 100644 --- a/src/app.rs +++ b/src/app.rs @@ -397,7 +397,8 @@ impl AppBuilder { .lancedb_path .clone() .unwrap_or_else(crate::config::default_lancedb_path); - match crate::workspace::LanceDbVectorStore::new(path).await { + let dim = embeddings.as_ref().map(|p| p.dimension()); + match crate::workspace::LanceDbVectorStore::new(path, dim).await { Ok(store) => { tracing::info!("LanceDB vector store connected for workspace search"); Some(Arc::new(store) as Arc) diff --git a/src/workspace/lancedb_store.rs b/src/workspace/lancedb_store.rs index 9b53243d..98f7bf6c 100644 --- a/src/workspace/lancedb_store.rs +++ b/src/workspace/lancedb_store.rs @@ -8,8 +8,8 @@ //! LANCEDB_PATH=~/.ironclaw/lancedb # Default //! VECTOR_BACKEND=lancedb # Use LanceDB for vector search -/// Embedding dimension (text-embedding-3-small default). -/// Must match the embedding model used. +/// Default embedding dimension (text-embedding-3-small). +/// Override by passing the actual provider dimension to `LanceDbVectorStore::new()`. pub const DEFAULT_EMBEDDING_DIM: i32 = 1536; #[cfg(feature = "lancedb")] @@ -38,15 +38,28 @@ mod impl_lancedb { } /// LanceDB-backed vector store. + /// + /// The `update_embedding` method uses delete-then-insert (not atomic). + /// LanceDB does not support transactions, so a crash between the two + /// operations can lose the embedding for that chunk. This is acceptable + /// for personal workspace sizes where data can be reindexed. pub struct LanceDbVectorStore { db: Arc, table_name: String, embedding_dim: i32, + schema: Arc, + table: tokio::sync::OnceCell, } impl LanceDbVectorStore { /// Create a new LanceDB store at the given path. - pub async fn new(path: impl AsRef) -> Result { + /// + /// `embedding_dim` should match `EmbeddingProvider::dimension()`. + /// Pass `None` to use the default (1536, text-embedding-3-small). + pub async fn new( + path: impl AsRef, + embedding_dim: Option, + ) -> Result { let path_str = path .as_ref() .to_str() @@ -60,10 +73,15 @@ mod impl_lancedb { } })?; + let dim = embedding_dim.unwrap_or(DEFAULT_EMBEDDING_DIM as usize) as i32; + let schema = Arc::new(Self::build_schema(dim)); + let store = Self { db: Arc::new(db), table_name: TABLE_NAME.to_string(), - embedding_dim: DEFAULT_EMBEDDING_DIM, + embedding_dim: dim, + schema, + table: tokio::sync::OnceCell::new(), }; store.ensure_table().await?; @@ -81,9 +99,8 @@ mod impl_lancedb { return Ok(()); } - let schema = Arc::new(self.schema()); self.db - .create_empty_table(&self.table_name, schema.clone()) + .create_empty_table(&self.table_name, self.schema.clone()) .execute() .await .map_err(|e| WorkspaceError::SearchFailed { @@ -96,7 +113,22 @@ mod impl_lancedb { Ok(()) } - fn schema(&self) -> Schema { + /// Get or open the cached table handle. + async fn table(&self) -> Result<&lancedb::Table, WorkspaceError> { + self.table + .get_or_try_init(|| async { + self.db + .open_table(&self.table_name) + .execute() + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to open table: {}", e), + }) + }) + .await + } + + fn build_schema(embedding_dim: i32) -> Schema { Schema::new(vec![ Field::new("chunk_id", DataType::Utf8, false), Field::new("document_id", DataType::Utf8, false), @@ -108,7 +140,7 @@ mod impl_lancedb { "vector", DataType::FixedSizeList( Arc::new(Field::new("item", DataType::Float32, true)), - self.embedding_dim, + embedding_dim, ), false, ), @@ -138,14 +170,7 @@ mod impl_lancedb { }); } - let table = self - .db - .open_table(&self.table_name) - .execute() - .await - .map_err(|e| WorkspaceError::SearchFailed { - reason: format!("Failed to open table: {}", e), - })?; + let table = self.table().await?; let chunk_ids = StringArray::from(vec![chunk_id.to_string()]); let document_ids = StringArray::from(vec![document_id.to_string()]); @@ -160,7 +185,7 @@ mod impl_lancedb { ); let batch = RecordBatch::try_new( - Arc::new(self.schema()), + self.schema.clone(), vec![ Arc::new(chunk_ids), Arc::new(document_ids), @@ -171,19 +196,19 @@ mod impl_lancedb { Arc::new(vectors), ], ) - .map_err(|e| WorkspaceError::ChunkingFailed { + .map_err(|e| WorkspaceError::EmbeddingFailed { reason: format!("Failed to create record batch: {}", e), })?; let batches = - RecordBatchIterator::new(vec![Ok(batch)].into_iter(), Arc::new(self.schema())); + RecordBatchIterator::new(vec![Ok(batch)].into_iter(), self.schema.clone()); table .add(Box::new(batches) as Box) .execute() .await - .map_err(|e| WorkspaceError::ChunkingFailed { - reason: format!("Failed to insert chunk: {}", e), + .map_err(|e| WorkspaceError::EmbeddingFailed { + reason: format!("Failed to store embedding: {}", e), })?; Ok(()) @@ -199,14 +224,7 @@ mod impl_lancedb { content: &str, embedding: &[f32], ) -> Result<(), WorkspaceError> { - let table = self - .db - .open_table(&self.table_name) - .execute() - .await - .map_err(|e| WorkspaceError::SearchFailed { - reason: format!("Failed to open table: {}", e), - })?; + let table = self.table().await?; table .delete(&format!( @@ -231,14 +249,7 @@ mod impl_lancedb { } async fn delete_embeddings(&self, document_id: Uuid) -> Result<(), WorkspaceError> { - let table = self - .db - .open_table(&self.table_name) - .execute() - .await - .map_err(|e| WorkspaceError::SearchFailed { - reason: format!("Failed to open table: {}", e), - })?; + let table = self.table().await?; table .delete(&format!( @@ -246,8 +257,8 @@ mod impl_lancedb { escape_predicate_value(&document_id.to_string()) )) .await - .map_err(|e| WorkspaceError::ChunkingFailed { - reason: format!("Failed to delete chunks: {}", e), + .map_err(|e| WorkspaceError::EmbeddingFailed { + reason: format!("Failed to delete embeddings: {}", e), })?; Ok(()) @@ -260,14 +271,7 @@ mod impl_lancedb { embedding: &[f32], limit: usize, ) -> Result, WorkspaceError> { - let table = self - .db - .open_table(&self.table_name) - .execute() - .await - .map_err(|e| WorkspaceError::SearchFailed { - reason: format!("Failed to open table: {}", e), - })?; + let table = self.table().await?; let filter = if let Some(aid) = agent_id { format!( @@ -407,7 +411,7 @@ mod tests { #[tokio::test] async fn test_insert_and_vector_search() { let dir = TempDir::new().unwrap(); - let store = LanceDbVectorStore::new(dir.path()).await.unwrap(); + let store = LanceDbVectorStore::new(dir.path(), None).await.unwrap(); let chunk_id = Uuid::new_v4(); let document_id = Uuid::new_v4(); @@ -443,7 +447,7 @@ mod tests { #[tokio::test] async fn test_insert_multiple_and_search_returns_ordered() { let dir = TempDir::new().unwrap(); - let store = LanceDbVectorStore::new(dir.path()).await.unwrap(); + let store = LanceDbVectorStore::new(dir.path(), None).await.unwrap(); let doc_id = Uuid::new_v4(); let user_id = "user1"; @@ -479,7 +483,7 @@ mod tests { #[tokio::test] async fn test_delete_chunks() { let dir = TempDir::new().unwrap(); - let store = LanceDbVectorStore::new(dir.path()).await.unwrap(); + let store = LanceDbVectorStore::new(dir.path(), None).await.unwrap(); let doc_id = Uuid::new_v4(); let user_id = "user1"; @@ -515,7 +519,7 @@ mod tests { #[tokio::test] async fn test_update_chunk_embedding() { let dir = TempDir::new().unwrap(); - let store = LanceDbVectorStore::new(dir.path()).await.unwrap(); + let store = LanceDbVectorStore::new(dir.path(), None).await.unwrap(); let chunk_id = Uuid::new_v4(); let doc_id = Uuid::new_v4(); @@ -560,7 +564,7 @@ mod tests { #[tokio::test] async fn test_vector_search_filters_by_user_and_agent() { let dir = TempDir::new().unwrap(); - let store = LanceDbVectorStore::new(dir.path()).await.unwrap(); + let store = LanceDbVectorStore::new(dir.path(), None).await.unwrap(); let doc_id = Uuid::new_v4(); let embedding = make_embedding(1.0); @@ -615,7 +619,7 @@ mod tests { #[tokio::test] async fn test_insert_rejects_wrong_embedding_dim() { let dir = TempDir::new().unwrap(); - let store = LanceDbVectorStore::new(dir.path()).await.unwrap(); + let store = LanceDbVectorStore::new(dir.path(), None).await.unwrap(); let wrong_dim: Vec = vec![1.0; 100]; diff --git a/src/workspace/mod.rs b/src/workspace/mod.rs index 3680116e..5e294ee4 100644 --- a/src/workspace/mod.rs +++ b/src/workspace/mod.rs @@ -868,26 +868,32 @@ impl Workspace { None }; + // When an external vector store is active, skip writing embeddings + // to the DB (they'd never be queried from there). + let db_embedding = if self.vector_store.is_some() { + None + } else { + embedding.as_deref() + }; + let chunk_id = self .storage - .insert_chunk(document_id, index as i32, &content, embedding.as_deref()) + .insert_chunk(document_id, index as i32, &content, db_embedding) .await?; - // Sync embedding to external vector store - if let (Some(vs), Some(emb)) = (&self.vector_store, &embedding) - && let Err(e) = vs - .store_embedding( - chunk_id, - document_id, - &doc.path, - &doc.user_id, - doc.agent_id, - &content, - emb, - ) - .await - { - tracing::warn!("Failed to store embedding in vector store: {}", e); + // Sync embedding to external vector store (propagate errors to + // avoid leaving a document with deleted-then-missing embeddings). + if let (Some(vs), Some(emb)) = (&self.vector_store, &embedding) { + vs.store_embedding( + chunk_id, + document_id, + &doc.path, + &doc.user_id, + doc.agent_id, + &content, + emb, + ) + .await?; } } @@ -1133,6 +1139,23 @@ impl Workspace { .get_chunks_without_embeddings(&self.user_id, self.agent_id, 100) .await?; + // Prefetch document metadata to avoid N+1 queries when syncing to vector store + let doc_map: std::collections::HashMap = + if self.vector_store.is_some() { + let mut map = std::collections::HashMap::new(); + for chunk in &chunks { + if !map.contains_key(&chunk.document_id) + && let Ok(doc) = + self.storage.get_document_by_id(chunk.document_id).await + { + map.insert(doc.id, doc); + } + } + map + } else { + std::collections::HashMap::new() + }; + let mut count = 0; for chunk in chunks { match provider.embed(&chunk.content).await { @@ -1142,9 +1165,9 @@ impl Workspace { .await?; // Sync to external vector store - if let Some(ref vs) = self.vector_store { - let doc = self.storage.get_document_by_id(chunk.document_id).await?; - if let Err(e) = vs + if let Some(ref vs) = self.vector_store + && let Some(doc) = doc_map.get(&chunk.document_id) + && let Err(e) = vs .update_embedding( chunk.id, chunk.document_id, @@ -1155,13 +1178,12 @@ impl Workspace { &embedding, ) .await - { - tracing::warn!( - "Failed to sync embedding to vector store for chunk {}: {}", - chunk.id, - e - ); - } + { + tracing::warn!( + "Failed to sync embedding to vector store for chunk {}: {}", + chunk.id, + e + ); } count += 1; diff --git a/tests/lancedb_integration.rs b/tests/lancedb_integration.rs index c96c0d20..185bce20 100644 --- a/tests/lancedb_integration.rs +++ b/tests/lancedb_integration.rs @@ -57,7 +57,7 @@ async fn setup_workspace() -> (Workspace, TempDir, TempDir) { libsql.run_migrations().await.unwrap(); let lancedb_dir = TempDir::new().unwrap(); - let store = LanceDbVectorStore::new(lancedb_dir.path()).await.unwrap(); + let store = LanceDbVectorStore::new(lancedb_dir.path(), None).await.unwrap(); let embedding = make_embedding(1.0); let ws = Workspace::new_with_db("test_user", Arc::new(libsql) as Arc)