use std::collections::{HashMap, HashSet};
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use std::sync::Mutex;
use anyhow::{Context, Result, bail};
use rusqlite::{Connection, OptionalExtension, params};
use uuid::Uuid;
pub const CACHE_SCHEMA: u32 = 2;
#[derive(Debug, Clone, PartialEq)]
pub struct IndexedMessage {
pub id: Uuid,
pub sub_id: Option<Uuid>,
pub role: String,
pub ts: String,
pub text: String,
}
#[derive(Debug, Clone, PartialEq)]
pub struct MessageHit {
pub chat_id: Uuid,
pub sub_id: Option<Uuid>,
pub message_id: Uuid,
pub role: String,
pub ts: String,
pub text: String,
}
impl MessageHit {
pub fn scope_id(&self) -> Uuid {
self.sub_id.unwrap_or(self.chat_id)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IndexScope {
Chat(Uuid),
Transcript(Uuid),
}
pub struct CacheDb {
conn: Mutex<Connection>,
}
impl CacheDb {
pub fn open(path: &Path) -> Result<Self> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).ok();
}
match Self::try_open(path) {
Ok(db) => Ok(db),
Err(err) => {
tracing::info!(
path = %path.display(),
error = %format!("{err:#}"),
"cache.db is unusable — deleting it and starting empty (derived data, rebuilt in the background)"
);
remove_db_files(path);
Self::try_open(path).with_context(|| format!("recreating {}", path.display()))
}
}
}
#[cfg(test)]
pub fn open_in_memory() -> Result<Self> {
Self::from_conn(Connection::open_in_memory()?)
}
fn try_open(path: &Path) -> Result<Self> {
let conn = Connection::open(path).with_context(|| format!("opening {}", path.display()))?;
Self::from_conn(conn)
}
fn from_conn(conn: Connection) -> Result<Self> {
migrate(&conn)?;
Ok(Self {
conn: Mutex::new(conn),
})
}
pub fn indexed_state(&self) -> Result<HashMap<Uuid, (i64, u64)>> {
let conn = self.conn.lock().unwrap();
let mut stmt = conn.prepare("SELECT chat_id, mtime_ms, size FROM indexed_chats")?;
let rows = stmt
.query_map([], |r| {
Ok((
parse_uuid(r.get::<_, String>(0)?),
(r.get::<_, i64>(1)?, r.get::<_, i64>(2)? as u64),
))
})?
.collect::<rusqlite::Result<HashMap<_, _>>>()?;
Ok(rows)
}
pub fn index_chat(
&self,
chat_id: Uuid,
mtime_ms: i64,
size: u64,
messages: &[IndexedMessage],
) -> Result<()> {
self.index_chat_guarded(chat_id, None, mtime_ms, size, messages)
.map(|_| ())
}
pub fn index_chat_if_unchanged(
&self,
chat_id: Uuid,
expected: Option<(i64, u64)>,
mtime_ms: i64,
size: u64,
messages: &[IndexedMessage],
) -> Result<bool> {
self.index_chat_guarded(chat_id, Some(expected), mtime_ms, size, messages)
}
fn index_chat_guarded(
&self,
chat_id: Uuid,
guard: Option<Option<(i64, u64)>>,
mtime_ms: i64,
size: u64,
messages: &[IndexedMessage],
) -> Result<bool> {
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction()?;
let chat = chat_id.to_string();
if let Some(expected) = guard {
let current: Option<(i64, u64)> = tx
.query_row(
"SELECT mtime_ms, size FROM indexed_chats WHERE chat_id = ?1",
params![chat],
|r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)? as u64)),
)
.optional()?;
if current != expected {
return Ok(false);
}
}
let existing: HashMap<String, (i64, i64)> = {
let mut stmt =
tx.prepare("SELECT message_id, id, text_hash FROM messages WHERE chat_id = ?1")?;
stmt.query_map(params![chat], |r| {
Ok((r.get::<_, String>(0)?, (r.get::<_, i64>(1)?, r.get(2)?)))
})?
.collect::<rusqlite::Result<HashMap<_, _>>>()?
};
let mut seen: HashSet<String> = HashSet::with_capacity(messages.len());
for msg in messages {
let message_id = msg.id.to_string();
if !seen.insert(message_id.clone()) {
continue;
}
let hash = text_hash(&msg.text);
match existing.get(&message_id) {
Some((_, old)) if *old == hash => continue, Some((rowid, _)) => {
tx.execute("DELETE FROM messages WHERE id = ?1", params![rowid])?;
}
None => {}
}
tx.execute(
"INSERT INTO messages(chat_id, sub_id, message_id, text_hash, role, ts, text)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
params![
chat,
msg.sub_id.map(|id| id.to_string()),
message_id,
hash,
msg.role,
msg.ts,
msg.text
],
)?;
}
for (message_id, (rowid, _)) in &existing {
if !seen.contains(message_id) {
tx.execute("DELETE FROM messages WHERE id = ?1", params![rowid])?;
}
}
tx.execute(
"INSERT INTO indexed_chats(chat_id, mtime_ms, size) VALUES (?1, ?2, ?3)
ON CONFLICT(chat_id) DO UPDATE SET mtime_ms = excluded.mtime_ms, size = excluded.size",
params![chat, mtime_ms, size as i64],
)?;
tx.commit()?;
Ok(true)
}
pub fn forget_chat(&self, chat_id: Uuid) -> Result<()> {
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction()?;
let chat = chat_id.to_string();
tx.execute("DELETE FROM messages WHERE chat_id = ?1", params![chat])?;
tx.execute(
"DELETE FROM indexed_chats WHERE chat_id = ?1",
params![chat],
)?;
tx.commit()?;
Ok(())
}
pub fn search_chats(&self, fts_query: &str) -> Result<Vec<Uuid>> {
let conn = self.conn.lock().unwrap();
let mut stmt = conn.prepare(
"SELECT DISTINCT COALESCE(m.sub_id, m.chat_id)
FROM messages_fts f
JOIN messages m ON m.id = f.rowid
WHERE messages_fts MATCH ?1",
)?;
let ids = (|| -> rusqlite::Result<Vec<Uuid>> {
stmt.query_map(params![fts_query], |r| {
Ok(parse_uuid(r.get::<_, String>(0)?))
})?
.collect()
})()
.with_context(|| format!("full-text query {fts_query:?}"))?;
Ok(ids)
}
pub fn search_messages(&self, fts_query: &str, limit: usize) -> Result<Vec<MessageHit>> {
self.search_messages_where(fts_query, None, limit)
}
fn search_messages_where(
&self,
fts_query: &str,
scope: Option<&[Uuid]>,
limit: usize,
) -> Result<Vec<MessageHit>> {
let scope_clause = match scope {
Some(ids) => format!(" AND {SCOPE_ID} IN ({})", vec!["?"; ids.len()].join(", ")),
None => String::new(),
};
let conn = self.conn.lock().unwrap();
let mut stmt = conn.prepare(&format!(
"SELECT m.chat_id, m.sub_id, m.message_id, m.role, m.ts, m.text
FROM messages_fts f
JOIN messages m ON m.id = f.rowid
WHERE messages_fts MATCH ?{scope_clause}
ORDER BY m.chat_id, m.sub_id, m.id
LIMIT ?"
))?;
let params = scoped_params(fts_query, scope.unwrap_or_default(), Some(limit));
let hits = (|| -> rusqlite::Result<Vec<MessageHit>> {
stmt.query_map(rusqlite::params_from_iter(params.iter()), |r| {
Ok(MessageHit {
chat_id: parse_uuid(r.get::<_, String>(0)?),
sub_id: r.get::<_, Option<String>>(1)?.map(parse_uuid),
message_id: parse_uuid(r.get::<_, String>(2)?),
role: r.get(3)?,
ts: r.get(4)?,
text: r.get(5)?,
})
})?
.collect()
})()
.with_context(|| format!("full-text message query {fts_query:?}"))?;
Ok(hits)
}
pub fn count_matching_messages(&self, fts_query: &str) -> Result<usize> {
let conn = self.conn.lock().unwrap();
let n: i64 = conn
.query_row(
"SELECT count(*) FROM messages_fts WHERE messages_fts MATCH ?1",
params![fts_query],
|r| r.get(0),
)
.with_context(|| format!("full-text message count {fts_query:?}"))?;
Ok(n as usize)
}
pub fn matching_messages_in_chat(
&self,
fts_query: &str,
scope: IndexScope,
) -> Result<Vec<Uuid>> {
let (clause, id) = match scope {
IndexScope::Chat(id) => ("m.chat_id = ?2 AND m.sub_id IS NULL", id),
IndexScope::Transcript(id) => ("m.sub_id = ?2", id),
};
let conn = self.conn.lock().unwrap();
let mut stmt = conn.prepare(&format!(
"SELECT m.message_id
FROM messages_fts f
JOIN messages m ON m.id = f.rowid
WHERE messages_fts MATCH ?1 AND {clause}"
))?;
let ids = (|| -> rusqlite::Result<Vec<Uuid>> {
stmt.query_map(params![fts_query, id.to_string()], |r| {
Ok(parse_uuid(r.get::<_, String>(0)?))
})?
.collect()
})()
.with_context(|| format!("full-text query {fts_query:?} in one chat"))?;
Ok(ids)
}
pub fn search_messages_in(
&self,
fts_query: &str,
chat_ids: &[Uuid],
limit: usize,
) -> Result<Vec<MessageHit>> {
if chat_ids.is_empty() {
return Ok(Vec::new());
}
self.search_messages_where(fts_query, Some(chat_ids), limit)
}
pub fn count_matching_messages_in(&self, fts_query: &str, chat_ids: &[Uuid]) -> Result<usize> {
if chat_ids.is_empty() {
return Ok(0);
}
let conn = self.conn.lock().unwrap();
let placeholders = vec!["?"; chat_ids.len()].join(", ");
let params = scoped_params(fts_query, chat_ids, None);
let n: i64 = conn
.query_row(
&format!(
"SELECT count(*)
FROM messages_fts f
JOIN messages m ON m.id = f.rowid
WHERE messages_fts MATCH ? AND {SCOPE_ID} IN ({placeholders})"
),
rusqlite::params_from_iter(params.iter()),
|r| r.get(0),
)
.with_context(|| format!("scoped full-text message count {fts_query:?}"))?;
Ok(n as usize)
}
#[cfg(test)]
pub fn message_count(&self) -> Result<usize> {
let conn = self.conn.lock().unwrap();
let n: i64 = conn.query_row("SELECT count(*) FROM messages", [], |r| r.get(0))?;
Ok(n as usize)
}
}
const SCOPE_ID: &str = "COALESCE(m.sub_id, m.chat_id)";
fn scoped_params(
fts_query: &str,
chat_ids: &[Uuid],
limit: Option<usize>,
) -> Vec<rusqlite::types::Value> {
let mut params: Vec<rusqlite::types::Value> = Vec::with_capacity(chat_ids.len() + 2);
params.push(fts_query.to_string().into());
params.extend(chat_ids.iter().map(|id| id.to_string().into()));
if let Some(limit) = limit {
params.push((limit as i64).into());
}
params
}
fn migrate(conn: &Connection) -> Result<()> {
let version = read_user_version(conn)?;
if version != 0 && version != CACHE_SCHEMA {
bail!("cache.db schema is {version}, this build indexes {CACHE_SCHEMA}");
}
baseline_ddl(conn)?;
if version == 0 {
set_user_version(conn, CACHE_SCHEMA)?;
}
Ok(())
}
fn baseline_ddl(conn: &Connection) -> Result<()> {
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS indexed_chats (
chat_id TEXT PRIMARY KEY,
mtime_ms INTEGER NOT NULL,
size INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS messages (
id INTEGER PRIMARY KEY,
chat_id TEXT NOT NULL,
sub_id TEXT,
message_id TEXT NOT NULL,
text_hash INTEGER NOT NULL,
role TEXT NOT NULL,
ts TEXT NOT NULL,
text TEXT NOT NULL,
UNIQUE(chat_id, message_id)
);
CREATE INDEX IF NOT EXISTS messages_sub_id ON messages(sub_id);
CREATE VIRTUAL TABLE IF NOT EXISTS messages_fts USING fts5(
text,
content='messages',
content_rowid='id',
tokenize='trigram'
);
CREATE TRIGGER IF NOT EXISTS messages_ai AFTER INSERT ON messages BEGIN
INSERT INTO messages_fts(rowid, text) VALUES (new.id, new.text);
END;
CREATE TRIGGER IF NOT EXISTS messages_ad AFTER DELETE ON messages BEGIN
INSERT INTO messages_fts(messages_fts, rowid, text)
VALUES ('delete', old.id, old.text);
END;",
)?;
Ok(())
}
fn read_user_version(conn: &Connection) -> Result<u32> {
Ok(conn.pragma_query_value(None, "user_version", |r| r.get::<_, i64>(0))? as u32)
}
fn set_user_version(conn: &Connection, v: u32) -> Result<()> {
conn.execute_batch(&format!("PRAGMA user_version = {v};"))?;
Ok(())
}
fn remove_db_files(path: &Path) {
let _ = std::fs::remove_file(path);
for suffix in ["-wal", "-shm"] {
let mut name = OsString::from(path.as_os_str());
name.push(suffix);
let _ = std::fs::remove_file(PathBuf::from(name));
}
}
fn text_hash(text: &str) -> i64 {
const OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
const PRIME: u64 = 0x0000_0100_0000_01b3;
let mut hash = OFFSET_BASIS;
for byte in text.as_bytes() {
hash ^= *byte as u64;
hash = hash.wrapping_mul(PRIME);
}
hash as i64
}
fn parse_uuid(s: String) -> Uuid {
Uuid::parse_str(&s).unwrap_or(Uuid::nil())
}
#[cfg(test)]
mod tests {
use super::*;
fn cache() -> CacheDb {
CacheDb::open_in_memory().unwrap()
}
fn msg(text: &str) -> IndexedMessage {
IndexedMessage {
id: Uuid::new_v4(),
sub_id: None,
role: "user".into(),
ts: "2026-07-29T10:00:00Z".into(),
text: text.into(),
}
}
fn rowids(db: &CacheDb, chat: Uuid) -> HashMap<Uuid, i64> {
let conn = db.conn.lock().unwrap();
let mut stmt = conn
.prepare("SELECT message_id, id FROM messages WHERE chat_id = ?1")
.unwrap();
stmt.query_map(params![chat.to_string()], |r| {
Ok((parse_uuid(r.get::<_, String>(0)?), r.get::<_, i64>(1)?))
})
.unwrap()
.collect::<rusqlite::Result<HashMap<_, _>>>()
.unwrap()
}
fn orphan_fts_rows(db: &CacheDb) -> i64 {
let conn = db.conn.lock().unwrap();
let orphans: i64 = conn
.query_row(
"SELECT count(*) FROM messages_fts_docsize d
LEFT JOIN messages m ON m.id = d.id
WHERE m.id IS NULL",
[],
|r| r.get(0),
)
.unwrap();
conn.execute_batch(
"INSERT INTO messages_fts(messages_fts, rank) VALUES('integrity-check', 1);",
)
.expect("the FTS index must match the content table");
orphans
}
#[test]
fn round_trip_index_and_search() {
let db = cache();
let chat = Uuid::new_v4();
db.index_chat(chat, 42, 7, &[msg("hello world"), msg("second message")])
.unwrap();
assert_eq!(db.message_count().unwrap(), 2);
assert_eq!(db.search_chats("\"hello\"").unwrap(), vec![chat]);
assert!(db.search_chats("\"nothing here\"").unwrap().is_empty());
}
#[test]
fn scoped_search_sees_only_the_given_chats() {
let db = cache();
let (mine, foreign) = (Uuid::new_v4(), Uuid::new_v4());
db.index_chat(mine, 1, 1, &[msg("shared password phrase")])
.unwrap();
db.index_chat(foreign, 1, 1, &[msg("shared password phrase")])
.unwrap();
let hits = db.search_messages_in("\"password\"", &[mine], 10).unwrap();
assert_eq!(hits.len(), 1);
assert_eq!(hits[0].chat_id, mine, "the foreign chat must not surface");
assert_eq!(
db.count_matching_messages_in("\"password\"", &[mine])
.unwrap(),
1
);
assert!(
db.search_messages_in("\"password\"", &[], 10)
.unwrap()
.is_empty()
);
assert_eq!(
db.count_matching_messages_in("\"password\"", &[]).unwrap(),
0
);
}
#[test]
fn scoped_search_cap_cuts_within_the_scope_not_before_it() {
let db = cache();
let (mine, foreign) = (Uuid::new_v4(), Uuid::new_v4());
db.index_chat(
foreign,
1,
1,
&(0..20)
.map(|i| msg(&format!("needle {i}")))
.collect::<Vec<_>>(),
)
.unwrap();
db.index_chat(
mine,
1,
1,
&[msg("needle a"), msg("needle b"), msg("needle c")],
)
.unwrap();
let hits = db.search_messages_in("\"needle\"", &[mine], 2).unwrap();
assert_eq!(hits.len(), 2);
assert!(hits.iter().all(|h| h.chat_id == mine));
assert_eq!(
db.count_matching_messages_in("\"needle\"", &[mine])
.unwrap(),
3
);
}
#[test]
fn trigram_matches_cyrillic_infix_and_folds_case() {
let db = cache();
let chat = Uuid::new_v4();
db.index_chat(chat, 1, 1, &[msg("это тестовое сообщение")])
.unwrap();
assert_eq!(
db.search_chats("\"естов\"").unwrap(),
vec![chat],
"infix match, which a word-based tokenizer cannot do"
);
assert_eq!(
db.search_chats("\"ЕСТОВ\"").unwrap(),
vec![chat],
"trigram folds case for Cyrillic on this SQLite build"
);
}
#[test]
fn unchanged_messages_keep_their_rowids_on_reindex() {
let db = cache();
let chat = Uuid::new_v4();
let (a, b, c) = (msg("first"), msg("second"), msg("streaming original"));
db.index_chat(chat, 1, 1, &[a.clone(), b.clone(), c.clone()])
.unwrap();
let before = rowids(&db, chat);
let grown = IndexedMessage {
text: "streaming replaced".into(),
..c.clone()
};
db.index_chat(chat, 2, 2, &[a.clone(), b.clone(), grown])
.unwrap();
let after = rowids(&db, chat);
assert_eq!(after[&a.id], before[&a.id], "untouched message rewritten");
assert_eq!(after[&b.id], before[&b.id], "untouched message rewritten");
assert_eq!(db.message_count().unwrap(), 3);
assert_eq!(db.search_chats("\"replaced\"").unwrap(), vec![chat]);
assert!(
db.search_chats("\"original\"").unwrap().is_empty(),
"the superseded text is still searchable"
);
assert_eq!(orphan_fts_rows(&db), 0);
}
#[test]
fn reindex_adds_and_removes_messages() {
let db = cache();
let chat = Uuid::new_v4();
let (a, b) = (msg("alpha content"), msg("beta content"));
db.index_chat(chat, 1, 1, &[a.clone(), b.clone()]).unwrap();
assert_eq!(db.search_chats("\"beta\"").unwrap(), vec![chat]);
db.index_chat(chat, 2, 2, std::slice::from_ref(&a)).unwrap();
assert_eq!(db.message_count().unwrap(), 1);
assert!(db.search_chats("\"beta\"").unwrap().is_empty());
assert_eq!(orphan_fts_rows(&db), 0);
let c = msg("gamma content");
db.index_chat(chat, 3, 3, &[a, c]).unwrap();
assert_eq!(db.search_chats("\"gamma\"").unwrap(), vec![chat]);
assert_eq!(db.message_count().unwrap(), 2);
assert_eq!(orphan_fts_rows(&db), 0);
}
#[test]
fn a_stale_writer_cannot_clobber_a_fresher_index() {
let db = cache();
let chat = Uuid::new_v4();
let seen_by_reconcile = db.indexed_state().unwrap().get(&chat).copied();
assert_eq!(seen_by_reconcile, None);
db.index_chat(chat, 200, 2000, &[msg("настоящая переписка")])
.unwrap();
let wrote = db
.index_chat_if_unchanged(chat, seen_by_reconcile, 100, 500, &[])
.unwrap();
assert!(!wrote, "the stale write must be refused");
assert_eq!(db.message_count().unwrap(), 1, "the fresh index was wiped");
assert_eq!(db.search_chats("\"настоящая\"").unwrap(), vec![chat]);
assert_eq!(
db.indexed_state().unwrap().get(&chat),
Some(&(200, 2000)),
"and the bookkeeping still describes the fresh write"
);
}
#[test]
fn a_guarded_write_goes_through_when_nothing_moved() {
let db = cache();
let chat = Uuid::new_v4();
db.index_chat(chat, 1, 10, &[msg("старое содержимое")])
.unwrap();
let seen = db.indexed_state().unwrap().get(&chat).copied();
let wrote = db
.index_chat_if_unchanged(chat, seen, 2, 20, &[msg("новое содержимое")])
.unwrap();
assert!(wrote);
assert_eq!(db.search_chats("\"новое\"").unwrap(), vec![chat]);
assert!(db.search_chats("\"старое\"").unwrap().is_empty());
assert_eq!(db.indexed_state().unwrap().get(&chat), Some(&(2, 20)));
}
#[test]
fn forget_chat_drops_rows_and_bookkeeping_of_that_chat_only() {
let db = cache();
let (one, two) = (Uuid::new_v4(), Uuid::new_v4());
db.index_chat(one, 1, 10, &[msg("shared word here")])
.unwrap();
db.index_chat(two, 2, 20, &[msg("shared word too")])
.unwrap();
db.forget_chat(one).unwrap();
assert_eq!(db.search_chats("\"shared\"").unwrap(), vec![two]);
assert_eq!(db.message_count().unwrap(), 1);
let state = db.indexed_state().unwrap();
assert!(!state.contains_key(&one), "bookkeeping left behind");
assert_eq!(state.get(&two), Some(&(2, 20)));
assert_eq!(orphan_fts_rows(&db), 0);
}
#[test]
fn search_isolates_chats() {
let db = cache();
let (one, two) = (Uuid::new_v4(), Uuid::new_v4());
db.index_chat(one, 1, 1, &[msg("apples and pears")])
.unwrap();
db.index_chat(two, 1, 1, &[msg("oranges and lemons")])
.unwrap();
assert_eq!(db.search_chats("\"apples\"").unwrap(), vec![one]);
assert_eq!(db.search_chats("\"oranges\"").unwrap(), vec![two]);
let both = db.search_chats("\"and\"").unwrap();
assert_eq!(both.len(), 2);
}
#[test]
fn search_messages_returns_rows_per_message_and_isolates_chats() {
let db = cache();
let (one, two) = (Uuid::new_v4(), Uuid::new_v4());
let a = IndexedMessage {
id: Uuid::new_v4(),
sub_id: None,
role: "assistant".into(),
ts: "2026-07-29T10:00:00+00:00".into(),
text: "первое упоминание маркера".into(),
};
let b = msg("второе упоминание маркера");
db.index_chat(one, 1, 1, &[a.clone(), b.clone(), msg("ничего")])
.unwrap();
db.index_chat(two, 1, 1, &[msg("маркера тут тоже")])
.unwrap();
let hits = db.search_messages("\"маркера\"", 100).unwrap();
assert_eq!(hits.len(), 3);
let chats: Vec<Uuid> = hits.iter().map(|h| h.chat_id).collect();
let mut deduped = chats.clone();
deduped.dedup();
assert_eq!(
deduped.len(),
2,
"hits of one chat must be adjacent: {chats:?}"
);
let mine: Vec<&MessageHit> = hits.iter().filter(|h| h.chat_id == one).collect();
assert_eq!(mine.len(), 2);
let first = mine.iter().find(|h| h.message_id == a.id).unwrap();
assert_eq!(first.role, "assistant");
assert_eq!(first.ts, a.ts);
assert_eq!(
first.text, a.text,
"the whole text — the snippet is built from it"
);
assert!(mine.iter().any(|h| h.message_id == b.id));
assert!(
db.search_messages("\"отсутствует\"", 100)
.unwrap()
.is_empty()
);
}
#[test]
fn search_messages_respects_the_limit_and_the_count_stays_honest() {
let db = cache();
let chat = Uuid::new_v4();
let msgs: Vec<IndexedMessage> = (0..10).map(|i| msg(&format!("совпадение {i}"))).collect();
db.index_chat(chat, 1, 1, &msgs).unwrap();
assert_eq!(db.search_messages("\"совпадение\"", 3).unwrap().len(), 3);
assert_eq!(db.search_messages("\"совпадение\"", 100).unwrap().len(), 10);
assert_eq!(db.count_matching_messages("\"совпадение\"").unwrap(), 10);
assert_eq!(db.count_matching_messages("\"нет\"").unwrap(), 0);
}
#[test]
fn matching_messages_in_chat_is_scoped_to_that_chat() {
let db = cache();
let (one, two) = (Uuid::new_v4(), Uuid::new_v4());
let (a, b) = (msg("общее слово раз"), msg("общее слово два"));
db.index_chat(one, 1, 1, &[a.clone(), msg("прочее"), b.clone()])
.unwrap();
db.index_chat(two, 1, 1, &[msg("общее слово чужое")])
.unwrap();
let mut ids = db
.matching_messages_in_chat("\"общее\"", IndexScope::Chat(one))
.unwrap();
ids.sort();
let mut want = vec![a.id, b.id];
want.sort();
assert_eq!(ids, want);
assert_eq!(
db.matching_messages_in_chat("\"общее\"", IndexScope::Chat(two))
.unwrap()
.len(),
1
);
assert!(
db.matching_messages_in_chat("\"прочее\"", IndexScope::Chat(two))
.unwrap()
.is_empty()
);
}
#[test]
fn malformed_message_queries_are_errors_not_panics() {
let db = cache();
assert!(db.search_messages("\"unterminated", 10).is_err());
assert!(db.count_matching_messages("\"unterminated").is_err());
assert!(
db.matching_messages_in_chat("\"unterminated", IndexScope::Chat(Uuid::new_v4()))
.is_err()
);
}
fn sub_msg(sub: Uuid, text: &str) -> IndexedMessage {
IndexedMessage {
sub_id: Some(sub),
..msg(text)
}
}
#[test]
fn a_transcript_is_its_own_conversation_in_every_query() {
let db = cache();
let (parent, run) = (Uuid::new_v4(), Uuid::new_v4());
let own = msg("родительское слово");
let inner = sub_msg(run, "дочернее слово");
db.index_chat(parent, 1, 1, &[own.clone(), inner.clone()])
.unwrap();
assert_eq!(db.search_chats("\"дочернее\"").unwrap(), vec![run]);
assert_eq!(db.search_chats("\"родительское\"").unwrap(), vec![parent]);
let mut both = db.search_chats("\"слово\"").unwrap();
both.sort();
let mut want = vec![parent, run];
want.sort();
assert_eq!(both, want);
let hits = db.search_messages("\"дочернее\"", 10).unwrap();
assert_eq!(hits.len(), 1);
assert_eq!((hits[0].chat_id, hits[0].sub_id), (parent, Some(run)));
assert_eq!(hits[0].scope_id(), run);
let scoped = db.search_messages_in("\"слово\"", &[run], 10).unwrap();
assert_eq!(scoped.len(), 1);
assert_eq!(scoped[0].message_id, inner.id);
assert_eq!(
db.count_matching_messages_in("\"слово\"", &[run]).unwrap(),
1
);
let scoped = db.search_messages_in("\"слово\"", &[parent], 10).unwrap();
assert_eq!(scoped.len(), 1);
assert_eq!(scoped[0].message_id, own.id);
assert_eq!(
db.matching_messages_in_chat("\"слово\"", IndexScope::Chat(parent))
.unwrap(),
vec![own.id]
);
assert_eq!(
db.matching_messages_in_chat("\"слово\"", IndexScope::Transcript(run))
.unwrap(),
vec![inner.id]
);
}
#[test]
fn transcript_rows_live_and_die_with_the_parent_file() {
let db = cache();
let (parent, run) = (Uuid::new_v4(), Uuid::new_v4());
let inner = sub_msg(run, "дочернее слово");
db.index_chat(parent, 1, 1, &[msg("своё"), inner.clone()])
.unwrap();
let before = rowids(&db, parent);
db.index_chat(parent, 2, 2, &[msg("своё"), inner.clone()])
.unwrap();
assert_eq!(rowids(&db, parent)[&inner.id], before[&inner.id]);
db.index_chat(parent, 3, 3, &[msg("своё")]).unwrap();
assert!(db.search_chats("\"дочернее\"").unwrap().is_empty());
assert_eq!(orphan_fts_rows(&db), 0);
db.index_chat(parent, 4, 4, &[inner]).unwrap();
db.forget_chat(parent).unwrap();
assert!(db.search_chats("\"дочернее\"").unwrap().is_empty());
assert_eq!(db.message_count().unwrap(), 0);
}
#[test]
fn indexed_state_round_trip() {
let db = cache();
assert!(db.indexed_state().unwrap().is_empty());
let chat = Uuid::new_v4();
db.index_chat(chat, 1_700_000_000_123, 4096, &[msg("x y z")])
.unwrap();
assert_eq!(
db.indexed_state().unwrap().get(&chat),
Some(&(1_700_000_000_123, 4096))
);
db.index_chat(chat, 1_700_000_999_999, 8192, &[msg("x y z")])
.unwrap();
let state = db.indexed_state().unwrap();
assert_eq!(state.len(), 1);
assert_eq!(state.get(&chat), Some(&(1_700_000_999_999, 8192)));
}
#[test]
fn a_malformed_query_is_an_error_not_a_panic() {
let db = cache();
assert!(db.search_chats("\"unterminated").is_err());
}
#[test]
fn open_wipes_a_database_from_another_schema() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("cache.db");
let chat = Uuid::new_v4();
{
let db = CacheDb::open(&path).unwrap();
db.index_chat(chat, 1, 1, &[msg("indexed under the old schema")])
.unwrap();
assert_eq!(db.message_count().unwrap(), 1);
}
{
let conn = Connection::open(&path).unwrap();
conn.execute_batch("PRAGMA user_version = 99;").unwrap();
}
let db = CacheDb::open(&path).unwrap();
assert_eq!(db.message_count().unwrap(), 0, "stale index kept");
assert!(db.indexed_state().unwrap().is_empty());
db.index_chat(chat, 2, 2, &[msg("indexed again")]).unwrap();
assert_eq!(db.search_chats("\"again\"").unwrap(), vec![chat]);
assert_eq!(
read_user_version(&db.conn.lock().unwrap()).unwrap(),
CACHE_SCHEMA
);
}
#[test]
fn open_wipes_a_file_that_is_not_a_database() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("cache.db");
std::fs::write(&path, b"this is not a database, it is a note to self").unwrap();
let db = CacheDb::open(&path).unwrap();
let chat = Uuid::new_v4();
db.index_chat(chat, 1, 1, &[msg("recovered")]).unwrap();
assert_eq!(db.search_chats("\"recovered\"").unwrap(), vec![chat]);
}
#[test]
fn duplicate_message_ids_do_not_fail_the_whole_chat() {
let db = cache();
let chat = Uuid::new_v4();
let one = msg("only once please");
db.index_chat(chat, 1, 1, &[one.clone(), one.clone()])
.unwrap();
assert_eq!(db.message_count().unwrap(), 1);
assert_eq!(db.search_chats("\"once please\"").unwrap(), vec![chat]);
}
#[test]
fn text_hash_is_stable_and_distinguishes() {
assert_eq!(text_hash(""), 0xcbf2_9ce4_8422_2325_u64 as i64);
assert_eq!(text_hash("a"), 0xaf63_dc4c_8601_ec8c_u64 as i64);
assert_ne!(text_hash("hello"), text_hash("hellp"));
}
#[test]
fn orphan_check_catches_a_missing_delete_trigger() {
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch(
"CREATE TABLE messages (
id INTEGER PRIMARY KEY, chat_id TEXT, message_id TEXT,
text_hash INTEGER, role TEXT, ts TEXT, text TEXT);
CREATE VIRTUAL TABLE messages_fts USING fts5(
text, content='messages', content_rowid='id', tokenize='trigram');
CREATE TRIGGER messages_ai AFTER INSERT ON messages BEGIN
INSERT INTO messages_fts(rowid, text) VALUES (new.id, new.text);
END;
-- deliberately no AFTER DELETE trigger
INSERT INTO messages(chat_id, message_id, text_hash, role, ts, text)
VALUES ('c', 'm', 0, 'user', 't', 'orphan me');
DELETE FROM messages;",
)
.unwrap();
let count = |sql: &str| conn.query_row(sql, [], |r| r.get::<_, i64>(0)).unwrap();
assert_eq!(
count("SELECT count(*) FROM messages_fts WHERE messages_fts MATCH '\"orphan\"'"),
1
);
assert_eq!(
count(
"SELECT count(*) FROM messages_fts_docsize d
LEFT JOIN messages m ON m.id = d.id WHERE m.id IS NULL"
),
1
);
assert!(
conn.execute_batch(
"INSERT INTO messages_fts(messages_fts, rank) VALUES('integrity-check', 1);"
)
.is_err(),
"integrity-check(1) must notice the index outliving its content"
);
assert_eq!(
count(
"SELECT count(*) FROM messages_fts f
LEFT JOIN messages m ON m.id = f.rowid WHERE m.id IS NULL"
),
0,
"a plain scan of an external-content FTS table reads through to the content table"
);
assert!(
conn.execute_batch("INSERT INTO messages_fts(messages_fts) VALUES('integrity-check');")
.is_ok(),
"the bare integrity-check only verifies the index's internal consistency"
);
}
}