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