1use 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
15fn 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 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 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 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 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 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(()) }
356
357 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 log::error!(
457 "SQLite 存储中一行数据反序列化失败(已从查询结果中缺失): {}",
458 e
459 );
460 None
461 }
462 })
463 .collect();
464 Ok(chunks)
465 }
466}