Skip to main content

p_memory/
storage.rs

1use crate::index::TextIndex;
2use crate::{schema, text, types::*, Error, Result};
3use fs4::fs_std::FileExt;
4use parking_lot::{Mutex, RwLock};
5use rusqlite::{params, params_from_iter, types::Value as SqlValue, Connection, OptionalExtension, Transaction, TransactionBehavior};
6use serde::de::DeserializeOwned;
7use serde_json::Value;
8use std::{collections::{BTreeMap, BTreeSet, HashMap, HashSet}, fs::{File, OpenOptions}, path::{Path, PathBuf}, sync::{atomic::{AtomicU64, Ordering}, Arc}};
9
10/// 读路径自愈索引时最多试几次拿写锁(每次退避 1ms,合计约 1s)。
11/// 只在库内补向量那条线程正在补齐时才试:它提交的正是读者要看的待办,
12/// 而那份提交一落地待办就空了,循环随即结束。别的时候(例如批量导入的写者占着写锁)
13/// 一律不试,读路径绝不为索引排队——那正是当初吞吐塌方的成因。
14
15
16/// 写者:独占的写连接 + 跨进程文件锁。只挡其他写者,不挡读。
17pub(crate) struct Writer { pub conn: Connection, _file_lock: File }
18
19/// 只读连接池。`rusqlite::Connection` 不是 `Sync`,并发读必须各持一条独立连接;
20/// 池锁只在取出与归还时短暂持有,读的整个过程不占任何全局锁。
21pub(crate) struct Readers { pub idle: Vec<Connection> }
22
23/// 向量分区缓存。失效按知识领域分开做:一个领域写了,别的领域已经载入的分区照旧留着。
24/// 条目留在表里就意味着可用——失效是就地删条目,不靠条目自带的版本号比对。
25/// 版本号仍要按领域记:载入期间该领域若被写过,手上这份数据已经过期,不能再写进缓存。
26/// 版本号也必须独立于全局 revision:向量落盘不推进 revision,只盯 revision 会漏失效。
27pub(crate) struct VectorCache {
28    /// 全局代次,整体失效时 +1。所有领域的版本号随之上移,因此它只增不减。
29    generation: AtomicU64,
30    /// 领域自己的递增计数,只在这个领域被写时 +1。
31    bumps: Mutex<HashMap<String, u64>>,
32    entries: Mutex<HashMap<(String, String, String), Option<Arc<crate::embeddings::Partition>>>>,
33}
34
35impl VectorCache {
36    fn new() -> Self {
37        Self { generation: AtomicU64::new(0), bumps: Mutex::new(HashMap::new()), entries: Mutex::new(HashMap::new()) }
38    }
39    /// 某个领域当前的版本号。没写过的领域就是全局代次本身。
40    pub fn epoch_of(&self, namespace: &str) -> u64 {
41        self.generation.load(Ordering::SeqCst) + self.bumps.lock().get(namespace).copied().unwrap_or(0)
42    }
43    /// 整体失效:清空所有条目,并让在途的载入结果不再被采用。
44    /// 清空是为了让内存有界——重新载入本来就是按需的。
45    pub fn invalidate(&self) {
46        self.generation.fetch_add(1, Ordering::SeqCst);
47        self.entries.lock().clear();
48    }
49    /// 只失效这些领域:删掉它们的条目、各自推进版本号,别的领域的缓存原样留着。
50    pub fn invalidate_namespaces(&self, namespaces: &HashSet<String>) {
51        {
52            let mut bumps = self.bumps.lock();
53            for namespace in namespaces { *bumps.entry(namespace.clone()).or_insert(0) += 1; }
54        }
55        self.entries.lock().retain(|(_, namespace, _), _| !namespaces.contains(namespace));
56    }
57}
58
59thread_local! {
60    /// 本次写入事务动过的知识领域。写路径自己登记(`touch_namespace`),
61    /// `mutate` 装上它、提交后取走,用来精准失效对应领域的向量分区缓存。
62    /// 事务之外的写入(迁移、导入进度回写)登记不生效,那些路径本来就走整体失效。
63    static TOUCHED_NAMESPACES: std::cell::RefCell<Option<HashSet<String>>> = const { std::cell::RefCell::new(None) };
64}
65
66/// 登记本次事务动过的领域。只在事务里有效,事务外是空操作。
67/// 名字用归一化后的写法:检索侧的缓存键就是归一化后的领域名。
68pub(crate) fn touch_namespace(namespace: &str) {
69    TOUCHED_NAMESPACES.with(|slot| {
70        if let Some(touched) = slot.borrow_mut().as_mut() { touched.insert(text::normalized_tag(namespace)); }
71    });
72}
73
74/// 登记某条记录所属的领域。记录还在库里时才取得出名字。
75pub(crate) fn touch_record_namespace(conn: &Connection, record_id: i64) -> Result<()> {
76    if let Some(namespace) = namespace_of(conn, record_id)? { touch_namespace(&namespace); }
77    Ok(())
78}
79
80/// 某条记录所属的领域名;记录不存在时是 `None`。
81pub(crate) fn namespace_of(conn: &Connection, record_id: i64) -> Result<Option<String>> {
82    Ok(conn.query_row("SELECT s.text FROM records r JOIN strings s ON s.id=r.namespace_id WHERE r.id=?1",
83        [record_id], |r| r.get(0)).optional()?)
84}
85
86/// 事务期间装上的领域登记表。装与还原都走它,中途报错提前返回也不会把登记表留在外面。
87struct TouchLog(Option<HashSet<String>>);
88
89impl TouchLog {
90    fn install() -> Self {
91        Self(TOUCHED_NAMESPACES.with(|slot| slot.borrow_mut().replace(HashSet::new())))
92    }
93    /// 取走这次事务登记的领域;登记表本身由 `Drop` 还原。
94    fn take(&self) -> HashSet<String> {
95        TOUCHED_NAMESPACES.with(|slot| slot.borrow_mut().take()).unwrap_or_default()
96    }
97}
98
99impl Drop for TouchLog {
100    fn drop(&mut self) {
101        let previous = self.0.take();
102        TOUCHED_NAMESPACES.with(|slot| *slot.borrow_mut() = previous);
103    }
104}
105
106pub(crate) struct Engine {
107    pub writer: Mutex<Option<Writer>>,
108    pub readers: Mutex<Option<Readers>>,
109    /// 向量外挂库 `vectors.sqlite3` 的写连接与读连接池:与主库同目录、独立 WAL、独立锁。
110    /// 向量是派生检索索引,读写不走主库通道,补齐线程的写不会占用业务写的写锁。
111    pub vector_writer: Mutex<Option<Writer>>,
112    pub vector_readers: Mutex<Option<Readers>>,
113    /// `TextIndex` 自带内部写锁、`IndexReader` 可并发检索,用 `Arc` 共享给所有读线程。
114    /// 放进 `Option` 是为了 `close` 时能真正销毁它——Tantivy 的写锁由 `IndexWriter` 持有,
115    /// 不销毁就无法释放,目录也重开不了。
116    pub index: RwLock<Option<Arc<TextIndex>>>,
117    pub vectors: VectorCache,
118    /// 宿主注册的模型能力。它们是运行时状态(闭包 / Python 函数无法序列化),不落盘。
119    pub embedders: crate::embeddings::EmbedderRegistry,
120    pub rerankers: crate::search::RerankerRegistry,
121    /// 宿主注册的事件接收位。不注册就什么都不产出。
122    pub events: crate::events::EventRegistry,
123    /// 库内的向量化线程。它在 `open` 里启动,线程自己持弱引用,因此只能在这里设一次。
124    pub vectorizer: std::sync::OnceLock<Arc<crate::embeddings::Vectorizer>>,
125    /// 最近观察到的降级档位,供健康检查读出「结果为什么变差」。
126    pub degraded: Mutex<Vec<Degrade>>,
127    pub root: PathBuf,
128}
129
130/// 只读连接的新建:`journal_mode` 是库文件上的持久属性,无需在每条连接上重设。
131fn open_reader(root: &Path) -> Result<Connection> {
132    let conn = Connection::open(root.join("store.sqlite3"))?;
133    conn.execute_batch("PRAGMA busy_timeout=5000; PRAGMA synchronous=NORMAL; PRAGMA foreign_keys=ON;")?;
134    // 向量库外挂:读路径也要看得见它(缺口核对、指纹比对都靠 ATTACH 读)。
135    conn.execute("ATTACH DATABASE ?1 AS vectors", [root.join("vectors.sqlite3").to_string_lossy().to_string()])?;
136    Ok(conn)
137}
138
139/// 向量外挂库的连接:独立文件、独立 WAL;`foreign_keys` 关掉——它只存自己的路由快照,
140/// 不跟主库做任何外键级联。
141fn open_vector_writer(root: &Path) -> Result<Connection> {
142    let conn = Connection::open(root.join("vectors.sqlite3"))?;
143    conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA synchronous=NORMAL; PRAGMA busy_timeout=5000;")?;
144    conn.execute_batch(include_str!("vectors_schema.sql"))?;
145    Ok(conn)
146}
147fn open_vector_reader(root: &Path) -> Result<Connection> {
148    let conn = Connection::open(root.join("vectors.sqlite3"))?;
149    conn.execute_batch("PRAGMA busy_timeout=5000; PRAGMA synchronous=NORMAL;")?;
150    Ok(conn)
151}
152
153/// Clone shares one process-local engine. Close invalidates all its handles.
154#[derive(Clone)]
155pub struct KnowledgeBase { pub(crate) engine: Arc<Engine> }
156
157/// 一次只读访问:独占一条连接,析构时归还池中。
158pub(crate) struct ReadGuard<'a> { engine: &'a Engine, conn: Option<Connection> }
159
160impl ReadGuard<'_> {
161    pub fn conn(&self) -> &Connection { self.conn.as_ref().expect("read connection lives until drop") }
162}
163
164impl Drop for ReadGuard<'_> {
165    fn drop(&mut self) {
166        let Some(conn) = self.conn.take() else { return };
167        // 关闭后归还无处安放,直接丢弃即可。
168        if let Some(readers) = self.engine.readers.lock().as_mut() { readers.idle.push(conn); }
169    }
170}
171
172impl KnowledgeBase {
173    pub fn open(directory: impl AsRef<Path>) -> Result<Self> {
174        std::fs::create_dir_all(directory.as_ref())?;
175        let root = std::fs::canonicalize(directory.as_ref())?;
176        let file_lock = OpenOptions::new().create(true).truncate(false).read(true).write(true).open(root.join("writer.lock"))?;
177        if !file_lock.try_lock_exclusive()? { return Err(Error::Locked(root.display().to_string())); }
178        let mut write_conn = Connection::open(root.join("store.sqlite3"))?;
179        schema::initialize(&mut write_conn)?;
180        // 向量外挂库:主库有 embeddings 表说明是旧库,先把它整表搬进 vectors.sqlite3。
181        // 搬迁在主库写连接上做,搬完再 ATTACH,避免挂一个还不存在的文件。
182        let needs_vector_migration: bool = write_conn.query_row(
183            "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='embeddings'", [], |r| r.get::<_, i64>(0)).map(|n| n > 0)?;
184        if needs_vector_migration { crate::embeddings::migrate_vectors(&root, &mut write_conn)?; }
185        // 向量库文件此刻必须已存在(迁移会建,新库也要建):否则 ATTACH 查不到表。
186        let vector_writer = open_vector_writer(&root)?;
187        // 向量库外挂:写路径的删除与指纹核对也要看得见它。
188        write_conn.execute("ATTACH DATABASE ?1 AS vectors", [root.join("vectors.sqlite3").to_string_lossy().to_string()])?;
189        let index = Arc::new(TextIndex::open(&root)?);
190        index.recover(&write_conn)?;
191        // 只读连接在 schema 建好之后再开,保证它看到的是完整结构。
192        let reader = open_reader(&root)?;
193        let vector_reader = open_vector_reader(&root)?;
194        let engine = Arc::new(Engine {
195            writer: Mutex::new(Some(Writer { conn: write_conn, _file_lock: file_lock })),
196            readers: Mutex::new(Some(Readers { idle: vec![reader] })),
197            vector_writer: Mutex::new(Some(Writer { conn: vector_writer, _file_lock: OpenOptions::new().create(true).truncate(false).read(true).write(true).open(root.join("vector_writer.lock"))? })),
198            vector_readers: Mutex::new(Some(Readers { idle: vec![vector_reader] })),
199            index: RwLock::new(Some(index)), vectors: VectorCache::new(),
200            embedders: crate::embeddings::EmbedderRegistry::new(),
201            rerankers: crate::search::RerankerRegistry::new(),
202            events: crate::events::EventRegistry::default(),
203            vectorizer: std::sync::OnceLock::new(),
204            degraded: Mutex::new(Vec::new()), root,
205        });
206        // 向量化线程在开库时启动,由 `close` 停下并等它收尾。
207        engine.vectorizer.set(crate::embeddings::Vectorizer::start(&engine)?).unwrap_or_else(|_| unreachable!("vectorizer starts once"));
208        Ok(Self { engine })
209    }
210
211    pub fn directory(&self) -> &Path { &self.engine.root }
212
213    /// 取文本索引的共享句柄。只在这一瞬间持有索引锁,拿到 `Arc` 后即可并发使用。
214    pub(crate) fn index(&self) -> Result<Arc<TextIndex>> {
215        self.engine.index.read().clone().ok_or(Error::Closed)
216    }
217
218    /// 把写入流程就地交过来的索引文档写进 Tantivy:此刻只是写进 writer,
219    /// 对搜索不可见,commit 仍由 `update_index`(或关闭时的收尾)一次做完。
220    pub(crate) fn index_documents(&self, docs: &[crate::index::IndexDocument]) -> Result<()> {
221        self.index()?.stage(docs)
222    }
223
224    pub fn close(&self) -> Result<()> {
225        // 先叫停向量化线程并等它收尾:它可能正在补一批向量,不能在库拆到一半时还在写。
226        if let Some(vectorizer) = self.engine.vectorizer.get() { vectorizer.stop(); }
227        let mut guard = self.engine.writer.lock();
228        let result = match guard.as_ref() {
229            Some(writer) => self.index()?.sync(&writer.conn),
230            None => Ok(()),
231        };
232        *guard = None;
233        *self.engine.vector_writer.lock() = None;
234        *self.engine.vector_readers.lock() = None;
235        // 必须真正销毁索引:Tantivy 的目录写锁由 IndexWriter 持有,不销毁就释放不掉。
236        *self.engine.index.write() = None;
237        *self.engine.readers.lock() = None;
238        result
239    }
240
241    /// 取一条独占的只读连接。并发读各拿各的,互不等待。
242    pub(crate) fn read(&self) -> Result<ReadGuard<'_>> {
243        let conn = {
244            let mut readers = self.engine.readers.lock();
245            match readers.as_mut() {
246                Some(readers) => match readers.idle.pop() {
247                    Some(conn) => conn,
248                    None => open_reader(&self.engine.root)?,
249                },
250                None => return Err(Error::Closed),
251            }
252        };
253        Ok(ReadGuard { engine: &self.engine, conn: Some(conn) })
254    }
255
256    /// 取一个向量分区:缓存里有就直接复用,否则从向量外挂库按需载入。
257    /// 载入过程不持缓存锁——否则一个慢分区会挡住所有其他分区的查询。
258    /// 向量库是独立文件、独立连接,载入期间不碰主库。
259    pub(crate) fn partition(&self, space: &crate::embeddings::EmbeddingSpace,
260        namespace: &str, scope: &str) -> Result<Option<Arc<crate::embeddings::Partition>>> {
261        let key = (space.id.clone(), namespace.to_string(), scope.to_string());
262        let epoch = self.engine.vectors.epoch_of(namespace);
263        let cached = self.engine.vectors.entries.lock().get(&key).cloned();
264        if let Some(partition) = cached { return Ok(partition); }
265        let loaded = {
266            let conn = {
267                let mut readers = self.engine.vector_readers.lock();
268                match readers.as_mut() {
269                    Some(readers) => readers.idle.pop().unwrap_or_else(|| open_vector_reader(&self.engine.root).unwrap_or_else(|_| unreachable!())),
270                    None => return Err(Error::Closed),
271                }
272            };
273            let result = crate::embeddings::Partition::load(&conn, space, namespace, scope)?;
274            if let Some(readers) = self.engine.vector_readers.lock().as_mut() { readers.idle.push(conn); }
275            result.map(Arc::new)
276        };
277        // 载入期间这个领域可能被写过:版本号变了就说明手上这份已经过期,索性不写缓存。
278        // 复查版本号与写条目要在同一把条目锁里:否则失效正好落在两步之间时,
279        // 一份过期分区会被永久留在表里,直到这个领域下次被写。
280        {
281            let mut entries = self.engine.vectors.entries.lock();
282            if self.engine.vectors.epoch_of(namespace) == epoch { entries.insert(key, loaded.clone()); }
283        }
284        Ok(loaded)
285    }
286
287    /// 只在索引确实落后、且写者空闲时追平。
288    /// 稳态下待办队列是空的(写路径提交后已就地清空),这一次预检就让读完全不碰写锁。
289    /// 待办非空说明写入正在进行或上一次索引提交失败过:此时**只尝试、不等待**。
290    /// 抢不到写锁就说明写者正在提交索引,读路径绝不能为此排队——一次提交是二十毫秒量级,
291    /// 排队会把所有读都堵在门外。等待空闲时再补平,兼顾「上次提交失败后自愈」。
292    pub(crate) fn sync_index_if_behind(&self, conn: &Connection) -> Result<()> {
293        let pending: i64 = conn.query_row("SELECT COUNT(*) FROM index_updates", [], |r| r.get(0))?;
294        if pending == 0 { return Ok(()); }
295        // 抢得到写锁就顺手追平;抢不到就走,不为任何写者排队。
296        // 向量补齐已拆到外挂库,不再占主库写锁,读者没有需要等的对象。
297        if let Some(mut guard) = self.engine.writer.try_lock() {
298            if let Some(writer) = guard.as_mut() { self.index()?.sync(&writer.conn)?; }
299        }
300        Ok(())
301    }
302
303    /// 业务写入:改记录、改向量。提交后按领域精准失效向量分区缓存。
304    pub(crate) fn mutate<T>(&self, f: impl FnOnce(&Transaction<'_>) -> Result<T>) -> Result<WriteReceipt<T>> {
305        let mut guard = self.engine.writer.lock();
306        let writer = guard.as_mut().ok_or(Error::Closed)?;
307        let changed_before = writer.conn.total_changes();
308        let log = TouchLog::install();
309        let tx = writer.conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
310        let value = f(&tx)?;
311        let revision = current_revision(&tx)?;
312        tx.commit()?;
313        let touched = log.take();
314        drop(log);
315        // 索引不在写入路径上追平——写入只把待办记进 index_updates,
316        // 由使用方稍后调用 `update_index` 一趟索引完、只提交一次。
317        // 向量缓存按领域失效:动过的领域就地清掉,没动过的照旧留着。
318        if touched.is_empty() {
319            // 登记为空却确实改过行:有写路径没登记,或者记录已删、领域无从查起。
320            // 宁可整体失效,也不留下一个来源说不清的陈旧分区。
321            if writer.conn.total_changes() > changed_before { self.engine.vectors.invalidate(); }
322        } else {
323            self.engine.vectors.invalidate_namespaces(&touched);
324        }
325        // 记录一变,「应向量化集合」就变了:受影响领域的就绪标记就地作废,
326        // 等补齐核对过再重新标。这样检索侧读到的永远是「核对过的那一份」。
327        self.invalidate_readiness(&writer.conn);
328        // 也给库内线程说一声:切片更新了就有一批该补的,叫它起来干活。
329        if let Some(vectorizer) = self.engine.vectorizer.get() { vectorizer.notify_work(); }
330        Ok(WriteReceipt { value, revision })
331    }
332
333    /// 只写那些不喂向量分区的表的写入:就绪标记、向量空间登记、各档开关、谓词规则。
334    /// 这类写入既不改记录也不改向量,向量分区缓存因此不动。
335    pub(crate) fn mutate_meta<T>(&self, f: impl FnOnce(&Transaction<'_>) -> Result<T>) -> Result<WriteReceipt<T>> {
336        let mut guard = self.engine.writer.lock();
337        let writer = guard.as_mut().ok_or(Error::Closed)?;
338        let tx = writer.conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
339        let value = f(&tx)?;
340        let revision = current_revision(&tx)?;
341        tx.commit()?;
342        Ok(WriteReceipt { value, revision })
343    }
344
345    /// 作废待办记录所属领域的就绪标记。待办队列记着这次动了哪些记录,
346    /// 按记录所属领域取名字;已经删掉的记录查不到领域,不影响(删记录不会造出缺口)。
347    fn invalidate_readiness(&self, conn: &Connection) {
348        let Ok(mut stmt) = conn.prepare("SELECT DISTINCT s.text FROM index_updates u
349            JOIN records r ON r.id=u.record_id JOIN strings s ON s.id=r.namespace_id") else {
350            return;
351        };
352        let Ok(namespaces) = stmt.query_map([], |r| r.get::<_, String>(0)) else {
353            return;
354        };
355        let namespaces: Vec<String> = namespaces.flatten().collect();
356        if namespaces.is_empty() { return; }
357        let mut guard = self.engine.vector_writer.lock();
358        let Some(vector_writer) = guard.as_mut() else { return; };
359        for namespace in namespaces {
360            let _ = crate::embeddings::clear_vector_ready(&vector_writer.conn, &namespace);
361        }
362    }
363
364    pub fn memories(&self) -> crate::memory::MemoryStore { crate::memory::MemoryStore(self.clone()) }
365    pub fn graph(&self) -> crate::graph::GraphStore { crate::graph::GraphStore(self.clone()) }
366    pub fn notes(&self) -> crate::notes::NoteStore { crate::notes::NoteStore(self.clone()) }
367    pub fn embeddings(&self) -> crate::embeddings::EmbeddingStore { crate::embeddings::EmbeddingStore(self.clone()) }
368
369    /// 需要写连接但不走事务的极少数场景(如导入进度回写)。
370    /// 走写锁,并保守失效向量缓存:调用方改了库里的东西,缓存不能不知情。
371    pub(crate) fn write<T>(&self, f: impl FnOnce(&Writer) -> Result<T>) -> Result<T> {
372        let mut guard = self.engine.writer.lock();
373        let value = f(guard.as_mut().ok_or(Error::Closed)?)?;
374        self.engine.vectors.invalidate();
375        Ok(value)
376    }
377
378    /// 记下一个降级档位,去重保留少量,供健康检查读出。
379    pub(crate) fn note_degrade(&self, degrade: Degrade) {
380        let mut observed = self.engine.degraded.lock();
381        if !observed.contains(&degrade) {
382            observed.push(degrade);
383            if observed.len() > 8 { observed.remove(0); }
384        }
385    }
386
387    /// 追平待索引队列:把写入累积的待办一趟索引完、只提交一次。
388    /// 写入不再就地索引,使用方(尤其批量导入)在合适时机调用本方法即可;
389    /// 期间读取走 `sync_index_if_behind` 自愈兜底。
390    pub fn update_index(&self) -> Result<HealthReport> {
391        self.catch_up_index(true)?;
392        if let Some(vectorizer) = self.engine.vectorizer.get() { vectorizer.notify_work(); }
393        self.health()
394    }
395
396    /// 只追平索引、不做健康检查。向量化前必须走一次:
397    /// 切片正文只存在索引里,索引没追平就取不到正文,只能拿到空串。
398    /// 调用方主动 `sync` 时(blocking)追平是它语义的一部分,阻塞拿锁做完;
399    /// 后台线程(非 blocking)只在拿得到时顺手做,抢不到说明前台正在写,待办留给下一轮。
400    pub(crate) fn catch_up_index(&self, blocking: bool) -> Result<()> {
401        if blocking {
402            let mut guard = self.engine.writer.lock();
403            let Some(writer) = guard.as_mut() else { return Ok(()) };
404            return self.index()?.sync(&writer.conn);
405        }
406        let Some(mut guard) = self.engine.writer.try_lock() else { return Ok(()) };
407        let Some(writer) = guard.as_mut() else { return Ok(()) };
408        self.index()?.sync(&writer.conn)
409    }
410
411    /// 注册事件接收回调。库在检索线程里同步调用它,所以回调必须非阻塞——
412    /// 在里面做同步 IO 或网络上报,会把检索拖住,和慢的重排回调一样。
413    /// 回调抛错只丢这一条事件,不影响检索。不注册就完全不产出事件。
414    pub fn register_event_sink<F: Fn(&crate::events::LogEvent) + Send + Sync + 'static>(&self, sink: F) {
415        self.engine.events.set(Arc::new(sink));
416    }
417
418    /// 注销事件接收回调,返回此前是否有注册。
419    pub fn unregister_event_sink(&self) -> bool { self.engine.events.clear() }
420
421    /// 当前是否注册了事件接收回调。
422    pub fn event_sink_registered(&self) -> bool { self.engine.events.is_registered() }
423
424    pub fn rebuild_indexes(&self) -> Result<HealthReport> {
425        let sink = self.engine.events.get();
426        let started = std::time::Instant::now();
427        {
428            let mut guard = self.engine.writer.lock();
429            let writer = guard.as_mut().ok_or(Error::Closed)?;
430            self.index()?.rebuild(&writer.conn)?;
431            self.engine.vectors.invalidate();
432        }
433        let ms = started.elapsed().as_millis() as u64;
434        let report = self.health()?;
435        if let Some(sink) = sink {
436            let mut event = crate::events::LogEvent::new("index_rebuild");
437            event.ms = ms;
438            event.documents = Some(report.index_document_count);
439            event.format = Some(crate::index::FORMAT.to_string());
440            let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| sink(&event)));
441        }
442        Ok(report)
443    }
444
445    /// 读一次全文索引重建的进度快照。纯原子读、不碰写锁,可在重建进行时从另一线程轮询。
446    pub fn rebuild_progress(&self) -> Result<RebuildProgressReport> {
447        Ok(self.index()?.rebuild_progress())
448    }
449
450    pub fn health(&self) -> Result<HealthReport> {
451        let state = self.read()?;
452        let conn = state.conn();
453        let mut counts = BTreeMap::new();
454        let mut stmt = conn.prepare("SELECT kind, COUNT(*) FROM records GROUP BY kind")?;
455        for row in stmt.query_map([], |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?)))? {
456            let (code, count) = row?;
457            let name = RecordKind::from_code(code).map(|k| k.as_str().to_string()).unwrap_or_else(|| code.to_string());
458            counts.insert(name, count as usize);
459        }
460        let record_count = counts.values().sum();
461        let mut foreign = conn.prepare("PRAGMA foreign_key_check")?;
462        let mut foreign_key_errors = 0;
463        let mut rows = foreign.query([])?;
464        while rows.next()?.is_some() { foreign_key_errors += 1; }
465        Ok(HealthReport {
466            schema_version: schema::SCHEMA_VERSION,
467            revision: current_revision(conn)?,
468            indexed_revision: meta(conn, "indexed_revision")?, record_count,
469            index_document_count: self.index()?.document_count(),
470            pending_index_updates: conn.query_row("SELECT COUNT(*) FROM index_updates", [], |r| r.get::<_, i64>(0))? as usize,
471            sqlite_integrity: conn.query_row("PRAGMA quick_check", [], |r| r.get(0))?,
472            foreign_key_errors, counts,
473            embedder_spaces: self.engine.embedders.space_ids(),
474            reranker_registered: self.engine.rerankers.is_registered(),
475            last_degraded: self.engine.degraded.lock().clone(),
476        })
477    }
478
479    /// Consistent SQLite snapshot, including embeddings. Refuses to overwrite a file.
480    pub fn backup(&self, target: impl AsRef<Path>) -> Result<()> {
481        let target = target.as_ref();
482        let state = self.read()?;
483        let reservation = OpenOptions::new().write(true).create_new(true).open(target)?;
484        drop(reservation);
485        if let Err(err) = state.conn().backup(rusqlite::MAIN_DB, target, None) {
486            let _ = std::fs::remove_file(target);
487            return Err(err.into());
488        }
489        // 向量外挂库也备份:它是派生索引,但备份恢复后不该要求重新跑模型。
490        let vectors_target = target.with_file_name(format!("{}.vectors", target.file_name().unwrap().to_string_lossy()));
491        if let Err(err) = state.conn().backup("vectors", &vectors_target, None) {
492            let _ = std::fs::remove_file(&vectors_target);
493            return Err(err.into());
494        }
495        Ok(())
496    }
497
498    /// Restores to a new directory; search indexes are rebuilt from the snapshot.
499    pub fn restore(snapshot: impl AsRef<Path>, directory: impl AsRef<Path>) -> Result<Self> {
500        let snapshot = snapshot.as_ref();
501        let source = Connection::open_with_flags(snapshot, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY)?;
502        let app: i64 = source.pragma_query_value(None, "application_id", |r| r.get(0))?;
503        let version: i64 = source.pragma_query_value(None, "user_version", |r| r.get(0))?;
504        if app != schema::APPLICATION_ID { return Err(Error::Validation("snapshot is not a p-memory database".into())); }
505        if version != schema::SCHEMA_VERSION { return Err(Error::SchemaVersion { found: version, supported: schema::SCHEMA_VERSION }); }
506        std::fs::create_dir(directory.as_ref())?;
507        source.backup(rusqlite::MAIN_DB, directory.as_ref().join("store.sqlite3"), None)?;
508        // 向量库备份在同目录、同名加 .vectors 后缀;有就恢复,没有就由补齐线程重建。
509        let vectors_snapshot = snapshot.with_file_name(format!("{}.vectors", snapshot.file_name().unwrap().to_string_lossy()));
510        if vectors_snapshot.exists() {
511            let vectors_source = Connection::open_with_flags(&vectors_snapshot, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY)?;
512            vectors_source.backup(rusqlite::MAIN_DB, directory.as_ref().join("vectors.sqlite3"), None)?;
513        }
514        Self::open(directory)
515    }
516}
517
518pub(crate) fn now_us() -> i64 { chrono::Utc::now().timestamp_micros() }
519pub(crate) fn meta(conn: &Connection, key: &str) -> Result<i64> {
520    Ok(conn.query_row("SELECT value FROM meta WHERE key=?1", [key], |r| r.get(0))?)
521}
522
523/// 读一个可能不存在的 meta 键;不存在返回 `None`。用于重建游标这类运行期临时状态。
524pub(crate) fn meta_opt(conn: &Connection, key: &str) -> Result<Option<i64>> {
525    Ok(conn.query_row("SELECT value FROM meta WHERE key=?1", [key], |r| r.get(0)).optional()?)
526}
527
528/// 写入(存在则更新)一个 meta 键。键值表结构固定,新增键不需要 schema 迁移。
529pub(crate) fn set_meta(conn: &Connection, key: &str, value: i64) -> Result<()> {
530    conn.execute("INSERT INTO meta(key,value) VALUES (?1,?2) ON CONFLICT(key) DO UPDATE SET value=excluded.value", params![key, value])?;
531    Ok(())
532}
533
534/// 删除一个运行时 meta 键,不存在时无副作用。
535pub(crate) fn clear_meta(conn: &Connection, key: &str) -> Result<()> {
536    conn.execute("DELETE FROM meta WHERE key=?1", [key])?;
537    Ok(())
538}
539pub(crate) fn current_revision(conn: &Connection) -> Result<i64> { meta(conn, "revision") }
540pub(crate) fn next_revision(conn: &Connection, record_id: i64) -> Result<i64> {
541    conn.execute("UPDATE meta SET value=value+1 WHERE key='revision'", [])?;
542    let revision = current_revision(conn)?;
543    conn.execute("INSERT INTO index_updates(revision,record_id) VALUES (?1,?2)", params![revision, record_id])?;
544    Ok(revision)
545}
546
547/// 登记一个开放标记并返回它的整数 id;字符串只在此表出现一次。
548pub(crate) fn term_id(conn: &Connection, text_value: &str) -> Result<i64> {
549    let normalized = text::normalized_tag(text_value);
550    conn.execute("INSERT OR IGNORE INTO strings(text) VALUES (?1)", [&normalized])?;
551    Ok(conn.query_row("SELECT id FROM strings WHERE text=?1", [&normalized], |r| r.get(0))?)
552}
553
554pub(crate) fn term_text(conn: &Connection, id: i64) -> Result<String> {
555    Ok(conn.query_row("SELECT text FROM strings WHERE id=?1", [id], |r| r.get(0))?)
556}
557
558/// 库里实际出现过的知识领域(记录用到的 namespace),按字典序。
559/// 就绪核对按它逐个领域做:没有记录的领域没有缺口可言。
560pub(crate) fn record_namespaces(conn: &Connection) -> Result<Vec<String>> {
561    let mut stmt = conn.prepare("SELECT DISTINCT s.text FROM records r JOIN strings s ON s.id=r.namespace_id ORDER BY s.text")?;
562    let mut namespaces = Vec::new();
563    for row in stmt.query_map([], |r| r.get::<_, String>(0))? { namespaces.push(row?); }
564    Ok(namespaces)
565}
566
567pub(crate) fn validate_identity(label: &str, value: &str) -> Result<()> {
568    if value.trim().is_empty() || value != value.trim() || value.chars().any(char::is_control) {
569        return Err(Error::Validation(format!("{label} must be nonempty, trimmed, and contain no control characters")));
570    }
571    Ok(())
572}
573pub(crate) fn validate_filter(filter: &ReadFilter) -> Result<()> {
574    validate_identity("namespace", &filter.namespace)?;
575    if filter.scopes.is_empty() { return Err(Error::Validation("at least one explicit read scope is required".into())); }
576    for scope in &filter.scopes { validate_identity("scope", scope)?; }
577    Ok(())
578}
579pub(crate) fn validate_limit(limit: usize) -> Result<()> {
580    if !(1..=10_000).contains(&limit) { return Err(Error::Validation("limit must be between 1 and 10000".into())); }
581    Ok(())
582}
583
584/// 归一化、去重后的标签文本(排序)。标签字符串统一落在 strings 表。
585pub(crate) fn normalize_tags(tags: &[String]) -> Vec<String> {
586    tags.iter().map(|label| text::normalized_tag(label)).filter(|tag| !tag.is_empty()).collect::<BTreeSet<_>>().into_iter().collect()
587}
588
589/// 标签集拼进哪一条的可搜正文:笔记这条自己不占索引文档,由它的第一片承载;
590/// 其余记录没有切片,就是它自己。返回空格连接的标签文本,不带标签时是空串。
591/// `exclude` 是已经升格成独立索引列的那些标签(笔记的目录段与文件名)——它们靠专门的列
592/// 参与匹配,不能再留在正文里,否则搜目录名会命中该目录下每一篇。
593pub(crate) fn tags_prefix(kind: RecordKind, tags: &[String], exclude: &[String], payload: &Value) -> String {
594    let carries = match kind {
595        RecordKind::Note => false,
596        RecordKind::Chunk => payload.get("ordinal").and_then(Value::as_u64) == Some(0),
597        _ => true,
598    };
599    if !carries { return String::new(); }
600    tags.iter().filter(|tag| !exclude.contains(tag)).cloned().collect::<Vec<_>>().join(" ")
601}
602
603/// 笔记相对路径拆成「目录段」与「文件名(去扩展名)」。
604/// 目录段原样,只有最后一段去后缀(目录名里的点不是扩展名)。写入与重建共用同一套规则。
605pub(crate) fn split_note_path(relative: &str) -> (Vec<String>, String) {
606    let segments: Vec<&str> = relative.split('/').filter(|segment| !segment.is_empty()).collect();
607    let Some((last, dirs)) = segments.split_last() else { return (Vec::new(), String::new()); };
608    let stem = last.rsplit_once('.').map(|(stem, _)| stem).unwrap_or(last).trim();
609    (dirs.iter().map(|segment| segment.to_string()).collect(), stem.to_string())
610}
611
612/// 一篇笔记的「目录段 / 文件名」。文件名直接读写入时存下的 `notes.name`(不事后拆路径,
613/// 平台的路径分隔符与扩展名边界都由 `Path` 在写入那一刻定好);目录段取相对路径去掉末段。
614/// 登记根目录后库里存的必是相对路径;迁移前留下的绝对路径没有相对目录可言,只给名字。
615pub(crate) fn note_path_parts(conn: &Connection, note_id: i64) -> (Vec<String>, String) {
616    let Ok((path, name)) = conn.query_row("SELECT path,name FROM notes WHERE record_id=?1", [note_id],
617        |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))) else {
618        return (Vec::new(), String::new());
619    };
620    if Path::new(&path).is_absolute() { return (Vec::new(), name); }
621    (split_note_path(&path).0, name)
622}
623
624/// 一条记录在索引里的「名字列 / 目录列 / 从可搜前缀里摘除的标签」。
625/// 切片的名字列放所属笔记的文件名、目录列放它的目录段;这两样都已升格成独立列,
626/// 就从可搜前缀里摘掉——否则搜目录名会命中该目录下每一篇,文件名也会被正文列弱命中重复计一次。
627/// 其余记录只有名字列(实体规范名),目录列为空、无摘除。
628pub(crate) fn index_columns(conn: &Connection, kind: RecordKind, payload: &Value) -> (String, String, Vec<String>) {
629    if kind == RecordKind::Chunk {
630        // 名字列与目录列只挂在这一篇的第一片上:否则一篇的每一片都会命中同一查询、把结果刷屏。
631        if payload.get("ordinal").and_then(Value::as_u64).unwrap_or(0) != 0 { return (String::new(), String::new(), Vec::new()); }
632        let note_id = payload.get("note_id").and_then(Value::as_i64).unwrap_or(0);
633        let (dirs, stem) = note_path_parts(conn, note_id);
634        let mut exclude = dirs.clone();
635        if !stem.is_empty() { exclude.push(stem.clone()); }
636        return (stem, dirs.join(" "), exclude);
637    }
638    (record_name(kind, payload), String::new(), Vec::new())
639}
640
641pub(crate) fn put_record(conn: &Connection, kind: RecordKind, input: &RecordInput,
642    payload: &Value, text: &str) -> Result<(RecordHeader, crate::index::IndexDocument)> {
643    validate_identity("namespace", &input.namespace)?;
644    validate_identity("scope", &input.scope)?;
645    for evidence in &input.evidence {
646        if evidence.source.trim().is_empty() { return Err(Error::Validation("evidence source is required".into())); }
647        match (evidence.offset, evidence.limit) {
648            (None, None) => {},
649            (Some(offset), Some(limit)) if offset >= 1 && limit >= 1 => {},
650            _ => return Err(Error::Validation("evidence offset/limit must be a 1-based start and a positive line count".into())),
651        }
652    }
653    let namespace_id = term_id(conn, &input.namespace)?;
654    // 记录一落库,这个领域的向量分区就变了:登记它,缓存按领域精准失效。
655    touch_namespace(&input.namespace);
656    let scope_id = term_id(conn, &input.scope)?;
657    let existing = match input.id {
658        Some(id) => Some(conn.query_row("SELECT created_at_us,updated_at_us,revision,scope_id FROM records WHERE id=?1", [id],
659            |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?, r.get::<_, i64>(2)?, r.get::<_, i64>(3)?))).optional()?
660            .ok_or_else(|| Error::NotFound(id.to_string()))?),
661        None => None,
662    };
663    if let (Some(id), true) = (input.id, existing.as_ref().is_some_and(|v| v.3 != scope_id)) {
664        // 换作用域会连带影响引用它的关系与事件,它们的端点必须与记录同域。
665        // 无引用的记录(记忆、笔记、孤立实体)直接换;有引用的先清掉引用再换。
666        let blocking: i64 = conn.query_row(
667            "SELECT (SELECT COUNT(*) FROM relations WHERE subject_id=?1 OR object_id=?1) \
668             + (SELECT COUNT(*) FROM event_participants WHERE entity_id=?1)", [id], |r| r.get(0))?;
669        if blocking > 0 {
670            return Err(Error::Conflict(format!("record {id} is referenced by {blocking} relation(s) or event participant row(s); remove those references before changing its scope")));
671        }
672    }
673    if let Some(expected) = input.expected_revision {
674        if existing.as_ref().map(|v| v.2) != Some(expected) { return Err(Error::StaleRevision(input.id.map(|v| v.to_string()).unwrap_or_default())); }
675    }
676    let now = now_us();
677    let created = existing.as_ref().map(|v| v.0).unwrap_or(input.created_at_us.unwrap_or(now));
678    let updated = input.updated_at_us.unwrap_or_else(|| now.max(existing.as_ref().map(|v| v.1).unwrap_or(created)));
679    if updated < created { return Err(Error::Validation("updated_at_us precedes created_at_us".into())); }
680    // 指纹跟着「送进搜索的文本 + 标签」走:正文或标签一变,各空间的旧向量立即失效。
681    let tags = normalize_tags(&input.tags);
682    let fingerprint = record_fingerprint(text, &tags);
683    let metadata_json = serde_json::to_string(&input.metadata)?;
684    let evidence_json = serde_json::to_string(&input.evidence)?;
685    let payload_json = serde_json::to_string(payload)?;
686    let (id, revision) = match input.id {
687        Some(id) => {
688            let revision = next_revision(conn, id)?;
689            conn.execute("UPDATE records SET namespace_id=?2,kind=?3,scope_id=?4,updated_at_us=?5,revision=?6,metadata_json=?7,
690                evidence_json=?8,fingerprint=?9,payload_json=?10 WHERE id=?1",
691                params![id, namespace_id, kind.code(), scope_id, updated, revision, metadata_json, evidence_json,
692                    fingerprint, payload_json])?;
693            // Updating text invalidates every space's embedding in the same transaction.
694            conn.execute("DELETE FROM vectors.embeddings WHERE record_id=?1 AND fingerprint<>?2", params![id, fingerprint])?;
695            (id, revision)
696        }
697        None => {
698            conn.execute("INSERT INTO records(namespace_id,kind,scope_id,created_at_us,updated_at_us,revision,metadata_json,evidence_json,
699                fingerprint,payload_json) VALUES (?1,?2,?3,?4,?5,0,?6,?7,?8,?9)",
700                params![namespace_id, kind.code(), scope_id, created, updated, metadata_json, evidence_json,
701                    fingerprint, payload_json])?;
702            let id = conn.last_insert_rowid();
703            let revision = next_revision(conn, id)?;
704            conn.execute("UPDATE records SET revision=?2 WHERE id=?1", params![id, revision])?;
705            (id, revision)
706        }
707    };
708    let tag_ids = set_record_tags(conn, id, &tags)?;
709    // 索引文档就地拼好交回调用方:正文来自本次写入手上的那一份,索引阶段不再回源。
710    // 标记一律带整数 id 给索引:namespace、scope、kind、tags 都不写第二份文本。
711    let (name, path, exclude) = index_columns(conn, kind, payload);
712    let document = crate::index::IndexDocument { id, namespace_id, scope_id, kind,
713        text: text.to_string(), name, path,
714        note_id: if kind == RecordKind::Chunk { payload.get("note_id").and_then(Value::as_i64).unwrap_or(0) } else { 0 },
715        tags_prefix: tags_prefix(kind, &tags, &exclude, payload), tag_ids };
716    Ok((RecordHeader { id, namespace: input.namespace.clone(), kind, scope: input.scope.clone(),
717        created_at_us: created, updated_at_us: updated, revision, tags,
718        evidence: input.evidence.clone(), metadata: input.metadata.clone() }, document))
719}
720
721/// 记录指纹:正文 + 标签一起算。两者任一变,各空间的旧向量立即失效。
722pub(crate) fn record_fingerprint(text: &str, tags: &[String]) -> String {
723    text::digest(&format!("text-v1\n{text}\n{}", tags.join(" ")))
724}
725
726/// 换掉一条记录的标签,返回这次的 tag id(按标签文本排序,与 `normalize_tags` 同序)。
727pub(crate) fn set_record_tags(conn: &Connection, id: i64, tags: &[String]) -> Result<Vec<i64>> {
728    conn.execute("DELETE FROM record_tags WHERE record_id=?1", [id])?;
729    let mut tag_ids = Vec::with_capacity(tags.len());
730    for tag in tags {
731        let tag_id = term_id(conn, tag)?;
732        conn.execute("INSERT OR IGNORE INTO record_tags(record_id,tag_id) VALUES (?1,?2)", params![id, tag_id])?;
733        tag_ids.push(tag_id);
734    }
735    Ok(tag_ids)
736}
737
738/// 一条记录的 (tag id, 标签文本),按文本排序。索引写 id 列与文本列都从这里取。
739pub(crate) fn record_tag_pairs(conn: &Connection, id: i64) -> Result<Vec<(i64, String)>> {
740    let mut stmt = conn.prepare("SELECT t.id,t.text FROM record_tags rt JOIN strings t ON t.id=rt.tag_id WHERE rt.record_id=?1 ORDER BY t.text")?;
741    let mut pairs = Vec::new();
742    for row in stmt.query_map([id], |r| Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?)))? { pairs.push(row?); }
743    Ok(pairs)
744}
745
746/// 按库里现有状态给一条记录拼索引文档:标记(namespace / scope / tags)一律取 strings 表的 id。
747pub(crate) fn index_document(conn: &Connection, id: i64, kind: RecordKind, text: String) -> Result<crate::index::IndexDocument> {
748    let (namespace_id, scope_id, payload_json): (i64, i64, String) = conn.query_row("SELECT namespace_id,scope_id,payload_json FROM records WHERE id=?1",
749        [id], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?;
750    let payload: Value = serde_json::from_str(&payload_json)?;
751    let pairs = record_tag_pairs(conn, id)?;
752    let tags: Vec<String> = pairs.iter().map(|(_, tag)| tag.clone()).collect();
753    let (name, path, exclude) = index_columns(conn, kind, &payload);
754    Ok(crate::index::IndexDocument { id, namespace_id, scope_id, kind, text,
755        name, path,
756        note_id: if kind == RecordKind::Chunk { payload.get("note_id").and_then(Value::as_i64).unwrap_or(0) } else { 0 },
757        tags_prefix: tags_prefix(kind, &tags, &exclude, &payload),
758        tag_ids: pairs.into_iter().map(|(tag_id, _)| tag_id).collect() })
759}
760
761/// 一批记录 id 里哪些是切片,各自属于哪一篇笔记、在原文里从第几行起。非切片不出现在结果里。
762/// 「一篇笔记有多少片段命中」要按笔记精确统计,折叠后的代表命中要挂上「同一篇的其余片段」,
763/// 两件事靠的都是这层「切片 → 笔记」的归属。
764pub(crate) fn chunk_notes(conn: &Connection, ids: &[i64]) -> Result<BTreeMap<i64, (i64, usize)>> {
765    let mut out = BTreeMap::new();
766    if ids.is_empty() { return Ok(out); }
767    let placeholders = ids.iter().map(|_| "?").collect::<Vec<_>>().join(",");
768    let mut stmt = conn.prepare(&format!("SELECT record_id,note_id,\"offset\" FROM chunks WHERE record_id IN ({placeholders})"))?;
769    for row in stmt.query_map(params_from_iter(ids.iter().copied()), |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?, r.get::<_, i64>(2)?)))? {
770        let (id, note_id, offset) = row?;
771        out.insert(id, (note_id, offset.max(0) as usize));
772    }
773    Ok(out)
774}
775
776/// 该知识领域登记的笔记根目录;没登记就是 `None`。
777pub(crate) fn namespace_root(conn: &Connection, namespace_id: i64) -> Result<Option<String>> {
778    Ok(conn.query_row("SELECT root FROM namespace_roots WHERE namespace_id=?1", [namespace_id], |r| r.get(0)).optional()?)
779}
780
781/// 库里存的相对路径 → 实际文件路径:登记过根目录就拼回去,没登记就是原样那条。
782pub(crate) fn absolute_note_path(conn: &Connection, namespace_id: i64, stored: &str) -> String {
783    match namespace_root(conn, namespace_id) {
784        Ok(Some(root)) => Path::new(&root).join(stored.replace('/', std::path::MAIN_SEPARATOR_STR)).to_string_lossy().into_owned(),
785        _ => stored.to_string(),
786    }
787}
788
789pub(crate) fn record_value(conn: &Connection, key: &RecordKey) -> Result<Option<Value>> {
790    Ok(record_values(conn, &[key.id])?.remove(&key.id))
791}
792
793/// 批量读取记录本体:一次查 records、一次查 tags,再由同一套装配逻辑还原。
794/// 语义等同于对每个 id 调用 `record_value`,只是把逐行往返压成固定两次查询。
795pub(crate) fn record_values(conn: &Connection, ids: &[i64]) -> Result<BTreeMap<i64, Value>> {
796    let mut out = BTreeMap::new();
797    if ids.is_empty() { return Ok(out); }
798    let placeholders = vec!["?"; ids.len()].join(",");
799    let params = ids.iter().map(|id| SqlValue::Integer(*id)).collect::<Vec<_>>();
800    let mut stmt = conn.prepare(&format!("SELECT r.id,r.namespace_id,n.text,r.kind,s.text,r.created_at_us,r.updated_at_us,r.revision,
801        r.metadata_json,r.evidence_json,r.payload_json FROM records r
802        JOIN strings n ON n.id=r.namespace_id JOIN strings s ON s.id=r.scope_id WHERE r.id IN ({placeholders}) ORDER BY r.id"))?;
803    let mut rows: Vec<(i64, i64, String, i64, String, i64, i64, i64, String, String, String)> = Vec::new();
804    for row in stmt.query_map(params_from_iter(params.iter().cloned()), |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?,
805        r.get::<_, String>(2)?, r.get::<_, i64>(3)?, r.get::<_, String>(4)?, r.get::<_, i64>(5)?, r.get::<_, i64>(6)?,
806        r.get::<_, i64>(7)?, r.get::<_, String>(8)?, r.get::<_, String>(9)?, r.get::<_, String>(10)?)))? {
807        rows.push(row?);
808    }
809    let mut tags_stmt = conn.prepare(&format!("SELECT rt.record_id,t.text FROM record_tags rt JOIN strings t ON t.id=rt.tag_id \
810        WHERE rt.record_id IN ({placeholders}) ORDER BY rt.record_id,t.text"))?;
811    let mut tags: BTreeMap<i64, Vec<String>> = BTreeMap::new();
812    for row in tags_stmt.query_map(params_from_iter(params.iter().cloned()), |r| Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?)))? {
813        let (id, tag) = row?;
814        tags.entry(id).or_default().push(tag);
815    }
816    // 笔记的路径与文件名都是它自己的列(原样):读取时按 record_id 补回,
817    // 路径不进标签字典,payload 里也不重复存。
818    let mut note_meta: BTreeMap<i64, (String, String)> = BTreeMap::new();
819    {
820        let mut stmt = conn.prepare(&format!("SELECT n.record_id,n.path,n.name FROM notes n \
821            WHERE n.record_id IN ({placeholders})"))?;
822        for row in stmt.query_map(params_from_iter(params.iter().cloned()), |r| Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?, r.get::<_, String>(2)?)))? {
823            let (id, path, name) = row?;
824            note_meta.insert(id, (path, name));
825        }
826    }
827    for (id, namespace_id, namespace, kind_code, scope, created, updated, revision, metadata, evidence, payload) in rows {
828        let kind = RecordKind::from_code(kind_code).ok_or_else(|| Error::Validation("invalid stored record kind".into()))?;
829        let header = RecordHeader { id, namespace, kind, scope,
830            created_at_us: created, updated_at_us: updated, revision, tags: tags.remove(&id).unwrap_or_default(),
831            metadata: serde_json::from_str(&metadata)?, evidence: serde_json::from_str(&evidence)? };
832        let mut value = serde_json::to_value(header)?;
833        let object = value.as_object_mut().ok_or_else(|| Error::Validation("invalid stored header".into()))?;
834        let mut payload: Metadata = serde_json::from_str(&payload)?;
835        // memory_type 也走 strings 文本表:payload 里只存 id,读取时还原文本。
836        if let Some(type_id) = payload.get("memory_type_id").and_then(Value::as_i64) {
837            payload.insert("memory_type".into(), Value::String(term_text(conn, type_id)?));
838            payload.remove("memory_type_id");
839        }
840        if kind == RecordKind::Note {
841            // 库里存的是相对路径(登记过领域根目录时),对外取回时拼回绝对路径。
842            let (stored, name) = note_meta.remove(&id).unwrap_or_default();
843            let source = absolute_note_path(conn, namespace_id, &stored);
844            payload.insert("source".into(), Value::String(source));
845            payload.insert("title".into(), Value::String(name));
846        }
847        object.extend(payload);
848        out.insert(id, value);
849    }
850    Ok(out)
851}
852
853pub(crate) fn matches_filter(conn: &Connection, key: &RecordKey, filter: &ReadFilter) -> Result<bool> {
854    validate_filter(filter)?;
855    let row: Option<(i64, i64)> = conn.query_row("SELECT namespace_id,scope_id FROM records WHERE id=?1", [key.id],
856        |r| Ok((r.get(0)?, r.get(1)?))).optional()?;
857    let Some((namespace_id, scope_id)) = row else { return Ok(false) };
858    if term_text(conn, namespace_id)? != text::normalized_tag(&filter.namespace) { return Ok(false); }
859    let scope = term_text(conn, scope_id)?;
860    if !filter.scopes.iter().any(|s| text::normalized_tag(s) == scope) { return Ok(false); }
861    for tag in &filter.tags {
862        let exists: bool = conn.query_row("SELECT EXISTS(SELECT 1 FROM record_tags rt JOIN strings t ON t.id=rt.tag_id WHERE rt.record_id=?1 AND t.text=?2)",
863            params![key.id, text::normalized_tag(tag)], |r| r.get(0))?;
864        if !exists { return Ok(false); }
865    }
866    Ok(true)
867}
868
869pub(crate) fn get<T: DeserializeOwned>(conn: &Connection, key: &RecordKey, filter: &ReadFilter) -> Result<T> {
870    if !matches_filter(conn, key, filter)? { return Err(Error::NotFound(key.id.to_string())); }
871    serde_json::from_value(record_value(conn, key)?.ok_or_else(|| Error::NotFound(key.id.to_string()))?).map_err(Error::from)
872}
873
874/// 批量读取一批记录,只返回满足 `filter` 的那些。
875/// 语义等同于对每个 id 依次调用 `get`,但把过滤压成一条 SQL,避免逐条回库。
876pub(crate) fn load_many<T: DeserializeOwned>(conn: &Connection, ids: &[i64], filter: &ReadFilter) -> Result<BTreeMap<i64, T>> {
877    let mut out = BTreeMap::new();
878    if ids.is_empty() { return Ok(out); }
879    validate_filter(filter)?;
880    let (condition, values) = filter_sql(filter, &[], true)?;
881    let placeholders = vec!["?"; ids.len()].join(",");
882    let mut stmt = conn.prepare(&format!("SELECT r.id FROM records r WHERE r.id IN ({placeholders}) AND {condition} ORDER BY r.id"))?;
883    let params = ids.iter().map(|id| SqlValue::Integer(*id)).chain(values).collect::<Vec<_>>();
884    let allowed = stmt.query_map(params_from_iter(params), |r| r.get::<_, i64>(0))?.collect::<std::result::Result<Vec<_>, _>>()?;
885    for (id, value) in record_values(conn, &allowed)? {
886        out.insert(id, serde_json::from_value(value)?);
887    }
888    Ok(out)
889}
890
891/// 生成 `records` 表上的过滤条件。
892///
893/// `by_ids` 用于「手里已经有一批 record_id、只在这些 id 内做属性过滤」的查询:给
894/// `namespace_id` 加一元 `+`(SQLite 的 no-op,唯一作用是禁止该表达式用索引),
895/// 逼优化器以 IN 列表逐行点查主键。不加的话它会拿 `records_scope` 去扫该 namespace 下的
896/// 全部行——23998 行的库上取 10 条要 1167 µs,加 `+` 只要 12 µs,而且成本随库规模线性增长。
897/// 反之,`select_keys` 那类「按条件翻页、手里没有 id 列表」的查询必须靠索引来缩小范围,不能禁。
898pub(crate) fn filter_sql(filter: &ReadFilter, kinds: &[RecordKind], by_ids: bool) -> Result<(String, Vec<SqlValue>)> {
899    validate_filter(filter)?;
900    let mut query = if by_ids {
901        "+r.namespace_id=(SELECT id FROM strings WHERE text=?)".to_string()
902    } else {
903        "r.namespace_id=(SELECT id FROM strings WHERE text=?)".to_string()
904    };
905    let mut values = vec![SqlValue::Text(text::normalized_tag(&filter.namespace))];
906    query.push_str(" AND r.scope_id IN (SELECT id FROM strings WHERE text IN (");
907    query.push_str(&vec!["?"; filter.scopes.len()].join(",")); query.push_str("))");
908    values.extend(filter.scopes.iter().map(|s| SqlValue::Text(text::normalized_tag(s))));
909    if !kinds.is_empty() {
910        query.push_str(" AND r.kind IN ("); query.push_str(&vec!["?"; kinds.len()].join(",")); query.push(')');
911        values.extend(kinds.iter().map(|k| SqlValue::Integer(k.code())));
912    }
913    for tag in &filter.tags {
914        query.push_str(" AND EXISTS(SELECT 1 FROM record_tags rt JOIN strings t ON t.id=rt.tag_id WHERE rt.record_id=r.id AND t.text=?)");
915        values.push(SqlValue::Text(text::normalized_tag(tag)));
916    }
917    Ok((query, values))
918}
919
920/// 按过滤条件选出记录 id(升序),供批量写入路径使用。
921/// 与翻页查询不同,这里要的是全量命中,所以不带 limit。
922pub(crate) fn select_ids(conn: &Connection, filter: &ReadFilter, kinds: &[RecordKind]) -> Result<Vec<i64>> {
923    let (condition, values) = filter_sql(filter, kinds, false)?;
924    let mut stmt = conn.prepare(&format!("SELECT r.id FROM records r WHERE {condition} ORDER BY r.id"))?;
925    let ids = stmt.query_map(params_from_iter(values), |r| r.get::<_, i64>(0))?
926        .collect::<std::result::Result<Vec<_>, _>>()?;
927    Ok(ids)
928}
929
930/// 过滤条件下的匹配总数(截断之前)。`with_total` 打开时才调用。
931pub(crate) fn count_matches(conn: &Connection, filter: &ReadFilter, kinds: &[RecordKind]) -> Result<usize> {
932    let (condition, values) = filter_sql(filter, kinds, false)?;
933    let count: i64 = conn.query_row(&format!("SELECT COUNT(*) FROM records r WHERE {condition}"), params_from_iter(values), |r| r.get(0))?;
934    Ok(count as usize)
935}
936
937/// 记录的正文列内容:这条记录自己的文本。
938/// 切片正文来自写入时切好的那一段,不在这里算,所以这条纯函数只覆盖其余四种记录。
939/// 实体的正文不含规范名——规范名单独走 `record_name` 的 name 列,不在正文里占位。
940pub(crate) fn record_text(kind: RecordKind, payload: &Value) -> String {
941    let field = |key: &str| payload.get(key).and_then(Value::as_str).unwrap_or("").to_string();
942    match kind {
943        RecordKind::Memory => field("judgment"),
944        RecordKind::Entity => entity_body(payload),
945        RecordKind::Relation => format!("{} {} {} {}", field("subject_name"), field("predicate"), field("object_name"), field("reason")),
946        RecordKind::Event => format!("{} {} {} {}", field("name"), field("summary"), name_list(payload), field("reason")),
947        // 笔记与切片的正文另有来源:笔记的检索面交给切片,切片正文由写入流程就地提供。
948        RecordKind::Note | RecordKind::Chunk => String::new(),
949    }
950}
951
952/// 记录的名字列内容:只有实体有规范名,其它记录为空串(检索时该列不参与)。
953pub(crate) fn record_name(kind: RecordKind, payload: &Value) -> String {
954    match kind {
955        RecordKind::Entity => payload.get("name").and_then(Value::as_str).unwrap_or("").to_string(),
956        _ => String::new(),
957    }
958}
959
960/// 一批事件记录的正文长度(字符数),只含库里存在的 id。
961///
962/// 与 `record_text(RecordKind::Event, ..)` 是同一套拼法,只是走 SQL 算:事件那一路常常
963/// 一次取出比字符预算大一到两个数量级的一批 id,先知道长度就能只装配真正要留下的几条。
964/// 非字符串字段与缺失字段都按空串计,与 Rust 侧 `Value::as_str` 的口径一致;
965/// `participant_names` 不是数组时同样按空串计。
966pub(crate) fn event_text_lengths(conn: &Connection, ids: &[i64]) -> Result<BTreeMap<i64, usize>> {
967    let mut out = BTreeMap::new();
968    if ids.is_empty() { return Ok(out); }
969    let placeholders = vec!["?"; ids.len()].join(",");
970    let params = ids.iter().map(|id| SqlValue::Integer(*id)).collect::<Vec<_>>();
971    let field = |name: &str| format!(
972        "CASE WHEN json_type(r.payload_json,'$.{name}')='text' THEN json_extract(r.payload_json,'$.{name}') ELSE '' END");
973    // 参与者名:不是数组时按空串计,是数组时只收文本元素(`j.type` 是元素的 JSON 类型名,
974    // 用 `json_type(j.value)` 会把已解包的文本再当 JSON 解析一次,直接报 malformed JSON)。
975    let names = "COALESCE(CASE WHEN json_type(r.payload_json,'$.participant_names')='array' \
976        THEN (SELECT group_concat(j.value,' ') FROM json_each(r.payload_json,'$.participant_names') j \
977            WHERE j.type='text') ELSE '' END,'')";
978    let mut stmt = conn.prepare(&format!(
979        "SELECT r.id, LENGTH({} || ' ' || {} || ' ' || {names} || ' ' || {}) \
980         FROM records r WHERE r.id IN ({placeholders}) ORDER BY r.id",
981        field("name"), field("summary"), field("reason")))?;
982    for row in stmt.query_map(params_from_iter(params.iter().cloned()), |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?)))? {
983        let (id, length) = row?;
984        out.insert(id, length.max(0) as usize);
985    }
986    Ok(out)
987}
988
989/// 一批记录里属于实体的那些的规范名(按 record_id)。重排取文档时用:正文列已不含规范名,
990/// 纯名实体(别名、摘要、属性全空)的正文是空串,得把规范名拼回去才能让重排看到名字。
991pub(crate) fn entity_names(conn: &Connection, ids: &[i64]) -> Result<BTreeMap<i64, String>> {
992    let mut out = BTreeMap::new();
993    if ids.is_empty() { return Ok(out); }
994    let placeholders = vec!["?"; ids.len()].join(",");
995    let mut stmt = conn.prepare(&format!("SELECT record_id,name FROM entities WHERE record_id IN ({placeholders})"))?;
996    let params = ids.iter().map(|id| SqlValue::Integer(*id)).collect::<Vec<_>>();
997    for row in stmt.query_map(params_from_iter(params), |r| Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?)))? {
998        let (id, name) = row?;
999        out.insert(id, name);
1000    }
1001    Ok(out)
1002}
1003
1004fn name_list(payload: &Value) -> String {
1005    payload.get("participant_names").and_then(Value::as_array)
1006        .map(|names| names.iter().filter_map(Value::as_str).collect::<Vec<_>>().join(" "))
1007        .unwrap_or_default()
1008}
1009
1010fn entity_body(payload: &Value) -> String {
1011    let aliases = payload.get("aliases").and_then(Value::as_array)
1012        .map(|a| a.iter().filter_map(Value::as_str).collect::<Vec<_>>().join(" ")).unwrap_or_default();
1013    let summary = payload.get("summary").and_then(Value::as_str).unwrap_or("");
1014    let attr_text = payload.get("attributes").and_then(Value::as_object).map(|attrs| {
1015        attrs.iter().map(|(key, values)| {
1016            let joined = values.as_array().map(|v| v.iter().filter_map(Value::as_str).collect::<Vec<_>>().join(" ")).unwrap_or_default();
1017            format!("{key} {joined}")
1018        }).collect::<Vec<_>>().join(" ")
1019    }).unwrap_or_default();
1020    format!("{aliases} {summary} {attr_text}")
1021}
1022
1023pub(crate) fn select_keys(conn: &Connection, filter: &ReadFilter, kinds: &[RecordKind], limit: usize, after: Option<&str>) -> Result<Vec<RecordKey>> {    let (mut condition, mut values) = filter_sql(filter, kinds, false)?;
1024    if let Some(cursor) = after {
1025        let id: i64 = cursor.parse().map_err(|_| Error::Validation("invalid page cursor".into()))?;
1026        condition.push_str(" AND r.id>?");
1027        values.push(SqlValue::Integer(id));
1028    }
1029    values.push(SqlValue::Integer(limit.min(i64::MAX as usize) as i64));
1030    let mut stmt = conn.prepare(&format!("SELECT r.id FROM records r WHERE {condition} ORDER BY r.id LIMIT ?"))?;
1031    let rows = stmt.query_map(params_from_iter(values), |r| r.get::<_, i64>(0))?;
1032    let mut keys = Vec::new();
1033    for row in rows { keys.push(RecordKey { id: row? }); }
1034    Ok(keys)
1035}
1036
1037pub(crate) fn list<T: DeserializeOwned>(conn: &Connection, kind: RecordKind, request: &PageRequest) -> Result<Page<T>> {
1038    validate_limit(request.limit)?;
1039    let mut keys = select_keys(conn, &request.filter, &[kind], request.limit + 1, request.after.as_deref())?;
1040    let has_more = keys.len() > request.limit;
1041    keys.truncate(request.limit);
1042    let next_cursor = if has_more { keys.last().map(RecordKey::index_key) } else { None };
1043    // 一次批量取回:逐条 `get` 会为每条记录各跑一遍过滤与装配,limit=100 就是几百次回库。
1044    let ids: Vec<i64> = keys.iter().map(|key| key.id).collect();
1045    let mut loaded: BTreeMap<i64, T> = load_many(conn, &ids, &request.filter)?;
1046    let items = keys.into_iter().filter_map(|key| loaded.remove(&key.id)).collect::<Vec<_>>();
1047    Ok(Page { items, next_cursor })
1048}
1049
1050pub(crate) fn delete_record(conn: &Connection, key: &RecordKey) -> Result<bool> {
1051    // Graph references are RESTRICT, so callers explicitly remove edges first.
1052    // 领域名要在删之前取:删完这条记录就查不到它属于哪个领域了。
1053    let namespace = namespace_of(conn, key.id)?;
1054    let changed = conn.execute("DELETE FROM records WHERE id=?1", [key.id]);
1055    let changed = match changed {
1056        Err(rusqlite::Error::SqliteFailure(err, _)) if err.code == rusqlite::ErrorCode::ConstraintViolation =>
1057            return Err(Error::Conflict(format!("record {} is still referenced", key.id))),
1058        other => other?,
1059    };
1060    if changed > 0 {
1061        // 向量外挂库没有外键级联,删记录时显式清掉它的向量。
1062        conn.execute("DELETE FROM vectors.embeddings WHERE record_id=?1", [key.id])?;
1063        next_revision(conn, key.id)?;
1064        if let Some(namespace) = namespace { touch_namespace(&namespace); }
1065    }
1066    Ok(changed > 0)
1067}
1068
1069#[cfg(test)]
1070mod tests {
1071    use super::*;
1072
1073    /// 按 id 批量取回时优化器选中的计划。
1074    fn batched_load_plan(conn: &Connection, ids: &[i64], by_ids: bool) -> String {
1075        let (condition, values) = filter_sql(&ReadFilter::default(), &[], by_ids).unwrap();
1076        let placeholders = vec!["?"; ids.len()].join(",");
1077        let sql = format!("EXPLAIN QUERY PLAN SELECT r.id FROM records r WHERE r.id IN ({placeholders}) AND {condition} ORDER BY r.id");
1078        let params: Vec<SqlValue> = ids.iter().map(|id| SqlValue::Integer(*id)).chain(values).collect();
1079        let mut stmt = conn.prepare(&sql).unwrap();
1080        let plans: Vec<String> = stmt.query_map(params_from_iter(params), |row| row.get::<_, String>(3))
1081            .unwrap().map(|row| row.unwrap()).collect();
1082        plans.join(" | ")
1083    }
1084
1085    fn seed(kb: &KnowledgeBase, rows: i64) {
1086        let inputs: Vec<crate::MemoryInput> = (1..=rows).map(|i| crate::MemoryInput::new(format!("记录 {i}"))).collect();
1087        kb.memories().upsert_many(&inputs).unwrap();
1088    }
1089
1090    fn sample_ids() -> Vec<i64> { (1..=10).collect() }
1091
1092    /// 手里已有 id 列表时,取回必须按主键点查。少了 `+`,优化器会去扫 records_scope 索引,
1093    /// 成本随库规模线性增长,而每一行都是额外的读放大。
1094    #[test]
1095    fn batched_load_stays_on_the_primary_key() {
1096        let dir = tempfile::tempdir().unwrap();
1097        let kb = KnowledgeBase::open(dir.path()).unwrap();
1098        seed(&kb, 100);
1099        let guard = kb.read().unwrap();
1100        let plan = batched_load_plan(guard.conn(), &sample_ids(), true);
1101        assert!(plan.contains("INTEGER PRIMARY KEY"), "批量取回退化为扫索引:{plan}");
1102    }
1103
1104    /// 事件正文长度走 SQL 算,必须与 Rust 侧 `record_text` 逐条一致:
1105    /// 两者一旦漂移,事件那一路的字符预算就截在别的地方,返回的事件跟着变。
1106    #[test]
1107    fn event_text_lengths_match_record_text() {
1108        let dir = tempfile::tempdir().unwrap();
1109        let kb = KnowledgeBase::open(dir.path()).unwrap();
1110        let entity = |name: &str| crate::EntityInput { record: Default::default(), name: name.into(),
1111            entity_type: "person".into(), aliases: vec![], attributes: BTreeMap::new(), summary: String::new() };
1112        let created = kb.graph().apply_batch(&crate::GraphBatch {
1113            entities: vec![entity("甲"), entity("乙")], ..Default::default()
1114        }).unwrap().value;
1115        let (first, second) = (created.entities[0].header.id, created.entities[1].header.id);
1116        let created = kb.graph().apply_batch(&crate::GraphBatch {
1117            events: vec![
1118                crate::EventInput { record: Default::default(), name: "别鹤典仪".into(), summary: "两人同去".into(),
1119                    participants: vec![first, second], confidence: 1.0, reason: "有人证".into() },
1120                crate::EventInput { record: Default::default(), name: "堂中自语".into(), summary: String::new(),
1121                    participants: vec![first], confidence: 1.0, reason: String::new() },
1122            ], ..Default::default()
1123        }).unwrap().value;
1124        let ids: Vec<i64> = created.events.iter().map(|event| event.header.id).collect();
1125
1126        // 再把两条改成边角形态:非字符串字段、缺字段、participant_names 不是数组。
1127        {
1128            let raw = Connection::open(dir.path().join("store.sqlite3")).unwrap();
1129            let payloads = [
1130                r#"{"name":7,"summary":"只剩数字名","participant_names":"甲 乙","reason":null}"#,
1131                r#"{"name":"正常","summary":null,"participant_names":["甲",7,"乙"],"reason":"理由"}"#,
1132            ];
1133            for (id, payload) in ids.iter().zip(payloads) {
1134                raw.execute("UPDATE records SET payload_json=?1 WHERE id=?2", params![payload, id]).unwrap();
1135            }
1136        }
1137
1138        let guard = kb.read().unwrap();
1139        let conn = guard.conn();
1140        let lengths = event_text_lengths(conn, &ids).unwrap();
1141        assert_eq!(lengths.len(), ids.len(), "每条事件都该有长度");
1142        for id in ids {
1143            let payload = record_values(conn, &[id]).unwrap().remove(&id).unwrap();
1144            assert_eq!(lengths[&id], record_text(RecordKind::Event, &payload).chars().count(),
1145                "事件 {id} 的 SQL 长度与 record_text 不一致");
1146        }
1147    }
1148
1149    // ── 向量分区缓存的按领域失效 ──────────────────────────────────────
1150    fn fixture_space() -> crate::embeddings::EmbeddingSpace {
1151        crate::embeddings::EmbeddingSpace { id: "v".into(), model: "fixture/v1".into(),
1152            dimension: 2, text_version: 1, encoding: "f32".into() }
1153    }
1154
1155    /// 让某个领域的向量分区进缓存。空领域也会被缓存成空分区,所以不必真造向量。
1156    fn cache_partition(kb: &KnowledgeBase, space: &crate::embeddings::EmbeddingSpace, namespace: &str) {
1157        kb.partition(space, namespace, "public").unwrap();
1158    }
1159
1160    /// 当前缓存着哪些领域的向量分区。
1161    fn cached_namespaces(kb: &KnowledgeBase) -> BTreeSet<String> {
1162        kb.engine.vectors.entries.lock().keys().map(|(_, namespace, _)| namespace.clone()).collect()
1163    }
1164
1165    fn namespace_filter(namespace: &str) -> ReadFilter {
1166        ReadFilter { namespace: namespace.into(), scopes: vec!["public".into()], tags: vec![], note_ids: vec![] }
1167    }
1168
1169    /// 缓存失效按领域分开做:动过的领域清条目并推进版本号,没动过的条目与版本号都不动。
1170    #[test]
1171    fn invalidating_one_namespace_leaves_the_others_alone() {
1172        let cache = VectorCache::new();
1173        let key = |namespace: &str| ("v".to_string(), namespace.to_string(), "public".to_string());
1174        cache.entries.lock().insert(key("a"), None);
1175        cache.entries.lock().insert(key("b"), None);
1176        let epoch_b = cache.epoch_of("b");
1177
1178        cache.invalidate_namespaces(&HashSet::from(["a".to_string()]));
1179
1180        assert!(cache.entries.lock().get(&key("a")).is_none(), "写过的领域要清掉条目");
1181        assert!(cache.entries.lock().get(&key("b")).is_some(), "没写过的领域不该被牵连");
1182        assert_eq!(cache.epoch_of("b"), epoch_b, "没写过的领域版本号不动");
1183        assert_ne!(cache.epoch_of("a"), epoch_b, "写过的领域版本号要前进,在途载入才会作废");
1184
1185        // 整体失效:所有领域一起作废,版本号只增不减。
1186        let epoch_a = cache.epoch_of("a");
1187        cache.invalidate();
1188        assert!(cache.entries.lock().is_empty());
1189        assert!(cache.epoch_of("a") > epoch_a && cache.epoch_of("b") > epoch_b);
1190    }
1191
1192    /// 写一个领域,只该清掉那个领域已经载入的分区。
1193    #[test]
1194    fn writing_one_namespace_keeps_other_vector_partitions_cached() {
1195        let dir = tempfile::tempdir().unwrap();
1196        let kb = KnowledgeBase::open(dir.path()).unwrap();
1197        let space = fixture_space();
1198        for namespace in ["a", "b"] { cache_partition(&kb, &space, namespace); }
1199        assert_eq!(cached_namespaces(&kb), BTreeSet::from(["a".to_string(), "b".to_string()]));
1200
1201        let mut input = crate::MemoryInput::new("写在 a 领域的一条");
1202        input.record.namespace = "a".into();
1203        kb.memories().upsert(input).unwrap();
1204
1205        assert_eq!(cached_namespaces(&kb), BTreeSet::from(["b".to_string()]), "只该清掉被写的那个领域");
1206    }
1207
1208    /// 删记录同样精准失效。领域名必须在删之前取出来——删完这条记录就查不到它属于哪个领域了。
1209    #[test]
1210    fn deleting_a_record_evicts_only_its_own_namespace() {
1211        let dir = tempfile::tempdir().unwrap();
1212        let kb = KnowledgeBase::open(dir.path()).unwrap();
1213        let mut input = crate::MemoryInput::new("要被删掉的一条");
1214        input.record.namespace = "a".into();
1215        let id = kb.memories().upsert(input).unwrap().value.header.id;
1216
1217        let space = fixture_space();
1218        for namespace in ["a", "b"] { cache_partition(&kb, &space, namespace); }
1219        kb.memories().delete(id, &namespace_filter("a")).unwrap();
1220
1221        assert_eq!(cached_namespaces(&kb), BTreeSet::from(["b".to_string()]), "删掉的领域要清,别的领域留着");
1222    }
1223
1224    /// 补齐向量跑一轮,只该作废它写了向量的领域。
1225    /// 补齐收尾会落就绪标记,那是只写不喂向量分区的写入,不能顺手把整个缓存清掉。
1226    #[test]
1227    fn filling_vectors_only_evicts_the_namespaces_it_wrote() {
1228        let dir = tempfile::tempdir().unwrap();
1229        let kb = KnowledgeBase::open(dir.path()).unwrap();
1230        let space = fixture_space();
1231        kb.embeddings().register_space(space.clone()).unwrap();
1232        kb.embeddings().register_embedder("v", |texts: &[String]| -> std::result::Result<Vec<Vec<f32>>, crate::EmbedCallbackError> {
1233            Ok(texts.iter().map(|_| vec![1.0f32, 0.0]).collect())
1234        }).unwrap();
1235        for namespace in ["a", "b"] { cache_partition(&kb, &space, namespace); }
1236        kb.memories().upsert(crate::MemoryInput::new("补齐用的一条")).unwrap();
1237        kb.embeddings().sync("v", 32).unwrap();
1238
1239        let cached = cached_namespaces(&kb);
1240        assert!(cached.contains("a") && cached.contains("b"),
1241            "补齐只写了 default 领域,a 与 b 的分区缓存不该被牵连:{cached:?}");
1242    }
1243
1244    /// 有写路径没登记领域时(按 id 删掉记录、或将来新增的写路径),宁可整体失效,
1245    /// 也不留下一个来源说不清的陈旧分区。
1246    #[test]
1247    fn an_unregistered_write_falls_back_to_invalidating_everything() {
1248        let dir = tempfile::tempdir().unwrap();
1249        let kb = KnowledgeBase::open(dir.path()).unwrap();
1250        let space = fixture_space();
1251        for namespace in ["a", "b"] { cache_partition(&kb, &space, namespace); }
1252
1253        // 直接改一行、不登记领域:模拟一条没接上登记的写路径。
1254        kb.mutate(|tx| Ok(tx.execute("INSERT INTO meta(key,value) VALUES ('cache_probe',1)
1255            ON CONFLICT(key) DO UPDATE SET value=excluded.value", [])?)).unwrap();
1256
1257        assert!(cached_namespaces(&kb).is_empty(), "登记为空却改过行时必须整体失效");
1258    }
1259}