Skip to main content

lc_vector_stores/
sqlite_store.rs

1// lc-vector-stores/src/sqlite_store.rs
2//! SQLite 文档存储实现
3
4use async_trait::async_trait;
5use lc_shared::splitter::{RecursiveCharacterSplitter, TextSplitter};
6use rusqlite::Connection;
7use std::path::Path;
8use std::sync::Arc;
9use tokio::sync::Mutex;
10use uuid::Uuid;
11
12use crate::document_store::{ChunkDocument, ChunkedDocumentStoreTrait, DocumentStore};
13use crate::{Document, VectorStoreError};
14
15/// 解析 metadata 列;损坏时记日志并回退空映射,不再静默吞错。
16fn parse_metadata_or_default(raw: &str) -> std::collections::HashMap<String, String> {
17    serde_json::from_str(raw).unwrap_or_else(|e| {
18        log::warn!("SQLite 存储中 metadata 损坏,已回退为空映射: {}", e);
19        std::collections::HashMap::new()
20    })
21}
22
23#[derive(Debug, Clone)]
24pub struct SQLiteStoreConfig {
25    pub db_path: String,
26}
27
28impl Default for SQLiteStoreConfig {
29    fn default() -> Self {
30        Self {
31            db_path: "langchainrust.db".to_string(),
32        }
33    }
34}
35
36impl SQLiteStoreConfig {
37    pub fn new(path: impl Into<String>) -> Self {
38        Self {
39            db_path: path.into(),
40        }
41    }
42}
43
44pub struct SQLiteDocumentStore {
45    conn: Arc<Mutex<Connection>>,
46}
47
48impl SQLiteDocumentStore {
49    pub fn new(config: SQLiteStoreConfig) -> Result<Self, VectorStoreError> {
50        let conn = Connection::open(&config.db_path)
51            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
52        conn.execute_batch(
53            "CREATE TABLE IF NOT EXISTS documents (
54                id TEXT PRIMARY KEY, content TEXT NOT NULL,
55                metadata TEXT NOT NULL DEFAULT '{}'
56            );
57            CREATE TABLE IF NOT EXISTS chunks (
58                chunk_id TEXT PRIMARY KEY, parent_id TEXT NOT NULL,
59                content TEXT NOT NULL, segment INTEGER NOT NULL,
60                metadata TEXT NOT NULL DEFAULT '{}'
61            );
62            CREATE INDEX IF NOT EXISTS idx_chunks_parent ON chunks(parent_id);
63            CREATE INDEX IF NOT EXISTS idx_chunks_segment ON chunks(parent_id, segment);",
64        )
65        .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
66        Ok(Self {
67            conn: Arc::new(Mutex::new(conn)),
68        })
69    }
70}
71
72#[async_trait]
73impl DocumentStore for SQLiteDocumentStore {
74    async fn add_document(&self, document: Document) -> Result<String, VectorStoreError> {
75        let id = document
76            .id
77            .clone()
78            .unwrap_or_else(|| Uuid::new_v4().to_string());
79        let meta = serde_json::to_string(&document.metadata).unwrap_or_else(|_| "{}".to_string());
80        let conn = self.conn.lock().await;
81        conn.execute(
82            "INSERT OR REPLACE INTO documents (id, content, metadata) VALUES (?1, ?2, ?3)",
83            rusqlite::params![id, document.content, meta],
84        )
85        .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
86        Ok(id)
87    }
88
89    async fn add_documents(
90        &self,
91        documents: Vec<Document>,
92    ) -> Result<Vec<String>, VectorStoreError> {
93        let conn = self.conn.lock().await;
94        let mut ids = Vec::new();
95        for doc in documents {
96            let id = doc.id.clone().unwrap_or_else(|| Uuid::new_v4().to_string());
97            let meta = serde_json::to_string(&doc.metadata).unwrap_or_else(|_| "{}".to_string());
98            conn.execute(
99                "INSERT OR REPLACE INTO documents (id, content, metadata) VALUES (?1, ?2, ?3)",
100                rusqlite::params![id, doc.content, meta],
101            )
102            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
103            ids.push(id);
104        }
105        Ok(ids)
106    }
107
108    async fn get_document(&self, id: &str) -> Result<Option<Document>, VectorStoreError> {
109        let conn = self.conn.lock().await;
110        let mut stmt = conn
111            .prepare("SELECT id, content, metadata FROM documents WHERE id = ?1")
112            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
113        let result = stmt.query_row(rusqlite::params![id], |row| {
114            let id: String = row.get(0)?;
115            let content: String = row.get(1)?;
116            let meta_str: String = row.get(2)?;
117            Ok(Document {
118                id: Some(id),
119                content,
120                metadata: parse_metadata_or_default(&meta_str),
121            })
122        });
123        match result {
124            Ok(doc) => Ok(Some(doc)),
125            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
126            Err(e) => Err(VectorStoreError::StorageError(e.to_string())),
127        }
128    }
129
130    async fn delete_document(&self, id: &str) -> Result<(), VectorStoreError> {
131        let conn = self.conn.lock().await;
132        conn.execute(
133            "DELETE FROM chunks WHERE parent_id = ?1",
134            rusqlite::params![id],
135        )
136        .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
137        conn.execute("DELETE FROM documents WHERE id = ?1", rusqlite::params![id])
138            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
139        Ok(())
140    }
141
142    async fn count(&self) -> usize {
143        let conn = self.conn.lock().await;
144        match conn.query_row("SELECT COUNT(*) FROM documents", [], |r| r.get(0)) {
145            Ok(count) => count,
146            // M3: trait 返回 usize 无法传播错误——不再静默吞错返回 0,
147            // 记 error 暴露存储故障,与文件内"不再静默丢行"的既有模式一致。
148            Err(e) => {
149                log::error!("SQLite count(documents) 查询失败,返回 0: {}", e);
150                0
151            }
152        }
153    }
154
155    async fn clear(&self) -> Result<(), VectorStoreError> {
156        let conn = self.conn.lock().await;
157        conn.execute_batch("DELETE FROM chunks; DELETE FROM documents;")
158            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
159        Ok(())
160    }
161}
162
163#[async_trait]
164impl ChunkedDocumentStoreTrait for SQLiteDocumentStore {
165    async fn add_parent_document(
166        &self,
167        document: Document,
168        chunk_size: usize,
169    ) -> Result<(String, Vec<String>), VectorStoreError> {
170        let splitter = RecursiveCharacterSplitter::new(chunk_size, chunk_size / 10);
171        let chunks_text = splitter.split_text(&document.content);
172        let parent_id = document
173            .id
174            .clone()
175            .unwrap_or_else(|| Uuid::new_v4().to_string());
176        let meta = serde_json::to_string(&document.metadata).unwrap_or_else(|_| "{}".to_string());
177
178        let conn = self.conn.lock().await;
179        conn.execute(
180            "INSERT OR REPLACE INTO documents (id, content, metadata) VALUES (?1, ?2, ?3)",
181            rusqlite::params![parent_id, document.content, meta],
182        )
183        .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
184
185        let mut chunk_ids = Vec::new();
186        for (i, text) in chunks_text.iter().enumerate() {
187            let cid = format!("{}:chunk:{}", parent_id, i);
188            conn.execute("INSERT OR REPLACE INTO chunks (chunk_id, parent_id, content, segment, metadata) VALUES (?1, ?2, ?3, ?4, ?5)",
189                rusqlite::params![cid, parent_id, text, i, "{}"])
190                .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
191            chunk_ids.push(cid);
192        }
193        Ok((parent_id, chunk_ids))
194    }
195
196    async fn add_parent_documents(
197        &self,
198        documents: Vec<Document>,
199        chunk_size: usize,
200    ) -> Result<Vec<(String, Vec<String>)>, VectorStoreError> {
201        let mut results = Vec::new();
202        for doc in documents {
203            results.push(self.add_parent_document(doc, chunk_size).await?);
204        }
205        Ok(results)
206    }
207
208    async fn get_parent_document(
209        &self,
210        parent_id: &str,
211    ) -> Result<Option<Document>, VectorStoreError> {
212        self.get_document(parent_id).await
213    }
214
215    async fn get_chunk(&self, chunk_id: &str) -> Result<Option<ChunkDocument>, VectorStoreError> {
216        let conn = self.conn.lock().await;
217        let mut stmt = conn.prepare("SELECT chunk_id, parent_id, content, segment, metadata FROM chunks WHERE chunk_id = ?1")
218            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
219        let result = stmt.query_row(rusqlite::params![chunk_id], |row| {
220            Ok(ChunkDocument {
221                chunk_id: row.get(0)?,
222                parent_id: row.get(1)?,
223                content: row.get(2)?,
224                segment: row.get(3)?,
225                metadata: parse_metadata_or_default(&row.get::<_, String>(4)?),
226            })
227        });
228        match result {
229            Ok(chunk) => Ok(Some(chunk)),
230            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
231            Err(e) => Err(VectorStoreError::StorageError(e.to_string())),
232        }
233    }
234
235    async fn get_chunk_document(
236        &self,
237        chunk_id: &str,
238    ) -> Result<Option<Document>, VectorStoreError> {
239        Ok(self.get_chunk(chunk_id).await?.map(|c| c.to_document()))
240    }
241
242    async fn get_chunks_for_parent(
243        &self,
244        parent_id: &str,
245    ) -> Result<Vec<ChunkDocument>, VectorStoreError> {
246        let conn = self.conn.lock().await;
247        let mut stmt = conn.prepare("SELECT chunk_id, parent_id, content, segment, metadata FROM chunks WHERE parent_id = ?1 ORDER BY segment")
248            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
249        let chunks = stmt
250            .query_map(rusqlite::params![parent_id], |row| {
251                Ok(ChunkDocument {
252                    chunk_id: row.get(0)?,
253                    parent_id: row.get(1)?,
254                    content: row.get(2)?,
255                    segment: row.get(3)?,
256                    metadata: parse_metadata_or_default(&row.get::<_, String>(4)?),
257                })
258            })
259            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?
260            .filter_map(|r| match r {
261                Ok(chunk) => Some(chunk),
262                Err(e) => {
263                    // 不再静默丢行:记录到日志,暴露存储降级
264                    log::error!(
265                        "SQLite 存储中一行数据反序列化失败(已从查询结果中缺失): {}",
266                        e
267                    );
268                    None
269                }
270            })
271            .collect();
272        Ok(chunks)
273    }
274
275    async fn get_chunk_documents_for_parent(
276        &self,
277        parent_id: &str,
278    ) -> Result<Vec<Document>, VectorStoreError> {
279        Ok(self
280            .get_chunks_for_parent(parent_id)
281            .await?
282            .into_iter()
283            .map(|c| c.to_document())
284            .collect())
285    }
286
287    async fn delete_parent_document(&self, parent_id: &str) -> Result<(), VectorStoreError> {
288        self.delete_document(parent_id).await
289    }
290
291    async fn parent_count(&self) -> usize {
292        let conn = self.conn.lock().await;
293        match conn.query_row("SELECT COUNT(*) FROM documents", [], |r| r.get(0)) {
294            Ok(count) => count,
295            // M3: 不静默吞错返回 0,记 error 暴露存储故障。
296            Err(e) => {
297                log::error!("SQLite parent_count 查询失败,返回 0: {}", e);
298                0
299            }
300        }
301    }
302
303    async fn chunk_count(&self) -> usize {
304        let conn = self.conn.lock().await;
305        match conn.query_row("SELECT COUNT(*) FROM chunks", [], |r| r.get(0)) {
306            Ok(count) => count,
307            // M3: 不静默吞错返回 0,记 error 暴露存储故障。
308            Err(e) => {
309                log::error!("SQLite chunk_count 查询失败,返回 0: {}", e);
310                0
311            }
312        }
313    }
314
315    async fn get_all_chunks(&self) -> Result<Vec<ChunkDocument>, VectorStoreError> {
316        let conn = self.conn.lock().await;
317        let mut stmt = conn
318            .prepare("SELECT chunk_id, parent_id, content, segment, metadata FROM chunks")
319            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
320        let chunks = stmt
321            .query_map([], |row| {
322                Ok(ChunkDocument {
323                    chunk_id: row.get(0)?,
324                    parent_id: row.get(1)?,
325                    content: row.get(2)?,
326                    segment: row.get(3)?,
327                    metadata: parse_metadata_or_default(&row.get::<_, String>(4)?),
328                })
329            })
330            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?
331            .filter_map(|r| match r {
332                Ok(chunk) => Some(chunk),
333                Err(e) => {
334                    // 不再静默丢行:记录到日志,暴露存储降级
335                    log::error!(
336                        "SQLite 存储中一行数据反序列化失败(已从查询结果中缺失): {}",
337                        e
338                    );
339                    None
340                }
341            })
342            .collect();
343        Ok(chunks)
344    }
345
346    async fn clear(&self) -> Result<(), VectorStoreError> {
347        let conn = self.conn.lock().await;
348        conn.execute_batch("DELETE FROM chunks; DELETE FROM documents;")
349            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
350        Ok(())
351    }
352
353    async fn save(&self, _path: impl AsRef<Path> + Send) -> Result<(), VectorStoreError> {
354        Ok(()) // SQLite auto-saves
355    }
356
357    // Blocking methods
358    fn add_parent_document_blocking(
359        &self,
360        document: Document,
361        chunk_size: usize,
362    ) -> Result<(String, Vec<String>), VectorStoreError> {
363        let conn = self.conn.blocking_lock();
364        let splitter = RecursiveCharacterSplitter::new(chunk_size, chunk_size / 10);
365        let chunks_text = splitter.split_text(&document.content);
366        let parent_id = document
367            .id
368            .clone()
369            .unwrap_or_else(|| Uuid::new_v4().to_string());
370        let meta = serde_json::to_string(&document.metadata).unwrap_or_else(|_| "{}".to_string());
371
372        conn.execute(
373            "INSERT OR REPLACE INTO documents (id, content, metadata) VALUES (?1, ?2, ?3)",
374            rusqlite::params![parent_id, document.content, meta],
375        )
376        .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
377
378        let mut chunk_ids = Vec::new();
379        for (i, text) in chunks_text.iter().enumerate() {
380            let cid = format!("{}:chunk:{}", parent_id, i);
381            conn.execute("INSERT OR REPLACE INTO chunks (chunk_id, parent_id, content, segment, metadata) VALUES (?1, ?2, ?3, ?4, ?5)",
382                rusqlite::params![cid, parent_id, text, i, "{}"])
383                .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
384            chunk_ids.push(cid);
385        }
386        Ok((parent_id, chunk_ids))
387    }
388
389    fn get_parent_document_blocking(
390        &self,
391        parent_id: &str,
392    ) -> Result<Option<Document>, VectorStoreError> {
393        let conn = self.conn.blocking_lock();
394        let mut stmt = conn
395            .prepare("SELECT id, content, metadata FROM documents WHERE id = ?1")
396            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
397        let result = stmt.query_row(rusqlite::params![parent_id], |row| {
398            Ok(Document {
399                id: Some(row.get(0)?),
400                content: row.get(1)?,
401                metadata: parse_metadata_or_default(&row.get::<_, String>(2)?),
402            })
403        });
404        match result {
405            Ok(doc) => Ok(Some(doc)),
406            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
407            Err(e) => Err(VectorStoreError::StorageError(e.to_string())),
408        }
409    }
410
411    fn get_chunk_blocking(
412        &self,
413        chunk_id: &str,
414    ) -> Result<Option<ChunkDocument>, VectorStoreError> {
415        let conn = self.conn.blocking_lock();
416        let mut stmt = conn.prepare("SELECT chunk_id, parent_id, content, segment, metadata FROM chunks WHERE chunk_id = ?1")
417            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
418        let result = stmt.query_row(rusqlite::params![chunk_id], |row| {
419            Ok(ChunkDocument {
420                chunk_id: row.get(0)?,
421                parent_id: row.get(1)?,
422                content: row.get(2)?,
423                segment: row.get(3)?,
424                metadata: parse_metadata_or_default(&row.get::<_, String>(4)?),
425            })
426        });
427        match result {
428            Ok(chunk) => Ok(Some(chunk)),
429            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
430            Err(e) => Err(VectorStoreError::StorageError(e.to_string())),
431        }
432    }
433
434    fn blocking_get_chunks_for_parent(
435        &self,
436        parent_id: &str,
437    ) -> Result<Vec<ChunkDocument>, VectorStoreError> {
438        let conn = self.conn.blocking_lock();
439        let mut stmt = conn.prepare("SELECT chunk_id, parent_id, content, segment, metadata FROM chunks WHERE parent_id = ?1 ORDER BY segment")
440            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?;
441        let chunks = stmt
442            .query_map(rusqlite::params![parent_id], |row| {
443                Ok(ChunkDocument {
444                    chunk_id: row.get(0)?,
445                    parent_id: row.get(1)?,
446                    content: row.get(2)?,
447                    segment: row.get(3)?,
448                    metadata: parse_metadata_or_default(&row.get::<_, String>(4)?),
449                })
450            })
451            .map_err(|e| VectorStoreError::StorageError(e.to_string()))?
452            .filter_map(|r| match r {
453                Ok(chunk) => Some(chunk),
454                Err(e) => {
455                    // 不再静默丢行:记录到日志,暴露存储降级
456                    log::error!(
457                        "SQLite 存储中一行数据反序列化失败(已从查询结果中缺失): {}",
458                        e
459                    );
460                    None
461                }
462            })
463            .collect();
464        Ok(chunks)
465    }
466}