1use std::path::{Path, PathBuf};
2use std::sync::Mutex;
3
4use anyhow::{Context, Result};
5use rusqlite::Connection;
6
7pub struct AnchorIndex {
8 path: PathBuf,
9 conn: Mutex<Connection>,
10}
11
12impl AnchorIndex {
13 pub fn open_project(project_dir: &Path) -> Result<Self> {
14 Self::open_with_schema(&project_dir.join("index.db"), PROJECT_SCHEMA)
15 }
16
17 fn open_with_schema(path: &Path, schema: &str) -> Result<Self> {
18 if let Some(parent) = path.parent() {
19 std::fs::create_dir_all(parent)
20 .with_context(|| format!("mkdir {}", parent.display()))?;
21 }
22 let conn = Connection::open(path).with_context(|| format!("open {}", path.display()))?;
23 conn.pragma_update(None, "journal_mode", "WAL")?;
24 conn.pragma_update(None, "busy_timeout", 5000)?;
25 conn.pragma_update(None, "synchronous", "NORMAL")?;
26 conn.execute_batch(schema)
27 .with_context(|| format!("apply schema on {}", path.display()))?;
28 Ok(Self {
29 path: path.to_path_buf(),
30 conn: Mutex::new(conn),
31 })
32 }
33
34 pub fn path(&self) -> &Path {
35 &self.path
36 }
37
38 pub fn conn(&self) -> std::sync::MutexGuard<'_, Connection> {
39 self.conn.lock().unwrap()
40 }
41
42 pub fn insert_project_event(
43 &self,
44 session_id: &str,
45 event: &crate::event::Event,
46 payload_json: &str,
47 ) -> rusqlite::Result<i64> {
48 let ts = chrono::Utc::now().to_rfc3339();
49 let kind = crate::event_writer::event_kind(event);
50 let (turn_id, flow_run_id) = crate::event_writer::extract_anchors(event);
51 let text_content = crate::event_writer::extract_text_content(event).unwrap_or_default();
52 let seq = 0_i64;
53 let mut conn = self.conn.lock().unwrap();
54 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
55 tx.execute(
56 "INSERT OR REPLACE INTO events \
57 (session_id, seq, ts, kind, turn_id, flow_run_id, payload) \
58 VALUES (?, ?, ?, ?, ?, ?, ?)",
59 rusqlite::params![
60 session_id,
61 seq,
62 ts,
63 kind,
64 turn_id,
65 flow_run_id,
66 payload_json
67 ],
68 )?;
69 let id = tx.last_insert_rowid();
70 tx.execute(
71 "INSERT OR REPLACE INTO events_fts (rowid, text_content) VALUES (?, ?)",
72 rusqlite::params![id, text_content],
73 )?;
74 tx.commit()?;
75 Ok(id)
76 }
77
78 pub fn fts_search_project_events(
79 &self,
80 query: &str,
81 session_filter: Option<&str>,
82 limit: usize,
83 ) -> Result<Vec<ProjectEventRow>> {
84 let conn = self.conn();
85 if let Some(pattern) = parse_regex_query(query) {
86 return self.search_events_regex(&pattern, session_filter, limit, &conn);
87 }
88 self.search_events_like(query, session_filter, limit, &conn)
89 }
90
91 fn search_events_like(
92 &self,
93 query: &str,
94 session_filter: Option<&str>,
95 limit: usize,
96 conn: &std::sync::MutexGuard<'_, rusqlite::Connection>,
97 ) -> Result<Vec<ProjectEventRow>> {
98 let pattern = format!("%{query}%");
99 let (sql, params): (&str, Vec<Box<dyn rusqlite::ToSql>>) = match session_filter {
100 Some(sid) => (
101 "SELECT e.session_id, e.seq, e.ts, e.kind, e.turn_id, e.flow_run_id, e.payload \
102 FROM events e JOIN events_fts f ON f.rowid = e.id \
103 WHERE f.text_content LIKE ?1 AND e.session_id = ?2 \
104 ORDER BY e.id DESC LIMIT ?3",
105 vec![
106 Box::new(pattern),
107 Box::new(sid.to_string()),
108 Box::new(limit as i64),
109 ],
110 ),
111 None => (
112 "SELECT e.session_id, e.seq, e.ts, e.kind, e.turn_id, e.flow_run_id, e.payload \
113 FROM events e JOIN events_fts f ON f.rowid = e.id \
114 WHERE f.text_content LIKE ?1 \
115 ORDER BY e.id DESC LIMIT ?2",
116 vec![Box::new(pattern), Box::new(limit as i64)],
117 ),
118 };
119 let mut stmt = conn.prepare(sql)?;
120 let param_refs: Vec<&dyn rusqlite::ToSql> = params.iter().map(|b| b.as_ref()).collect();
121 let rows = stmt.query_map(param_refs.as_slice(), project_event_row_from)?;
122 collect(rows)
123 }
124
125 pub fn count_events(&self, session_id: &str, kinds: Option<&[&str]>) -> Result<u64> {
126 let conn = self.conn();
127 let (sql, params): (String, Vec<Box<dyn rusqlite::ToSql>>) = match kinds {
128 Some(ks) if !ks.is_empty() => {
129 let placeholders = ks.iter().map(|_| "?").collect::<Vec<_>>().join(",");
130 let sql = format!(
131 "SELECT COUNT(*) FROM events WHERE session_id = ? AND kind IN ({placeholders})"
132 );
133 let mut params: Vec<Box<dyn rusqlite::ToSql>> =
134 vec![Box::new(session_id.to_string())];
135 for k in ks {
136 params.push(Box::new(k.to_string()));
137 }
138 (sql, params)
139 }
140 _ => (
141 "SELECT COUNT(*) FROM events WHERE session_id = ?".into(),
142 vec![Box::new(session_id.to_string())],
143 ),
144 };
145 let mut stmt = conn.prepare(&sql)?;
146 let param_refs: Vec<&dyn rusqlite::ToSql> = params.iter().map(|b| b.as_ref()).collect();
147 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
148 Ok(count as u64)
149 }
150
151 fn search_events_regex(
152 &self,
153 pattern: &str,
154 session_filter: Option<&str>,
155 limit: usize,
156 conn: &std::sync::MutexGuard<'_, rusqlite::Connection>,
157 ) -> Result<Vec<ProjectEventRow>> {
158 let re = regex::RegexBuilder::new(pattern)
159 .case_insensitive(true)
160 .build()
161 .map_err(|e| anyhow::anyhow!("invalid regex: {e}"))?;
162 let (sql, params): (&str, Vec<Box<dyn rusqlite::ToSql>>) = match session_filter {
163 Some(sid) => (
164 "SELECT e.session_id, e.seq, e.ts, e.kind, e.turn_id, e.flow_run_id, e.payload, \
165 f.text_content \
166 FROM events e JOIN events_fts f ON f.rowid = e.id \
167 WHERE e.session_id = ?1 ORDER BY e.id DESC LIMIT 2000",
168 vec![Box::new(sid.to_string())],
169 ),
170 None => (
171 "SELECT e.session_id, e.seq, e.ts, e.kind, e.turn_id, e.flow_run_id, e.payload, \
172 f.text_content \
173 FROM events e JOIN events_fts f ON f.rowid = e.id \
174 ORDER BY e.id DESC LIMIT 2000",
175 vec![],
176 ),
177 };
178 let mut stmt = conn.prepare(sql)?;
179 let param_refs: Vec<&dyn rusqlite::ToSql> = params.iter().map(|b| b.as_ref()).collect();
180 let rows = stmt.query_map(param_refs.as_slice(), |row| {
181 Ok((
182 ProjectEventRow {
183 session_id: row.get(0)?,
184 seq: row.get::<_, i64>(1)? as u64,
185 ts: row.get(2)?,
186 kind: row.get(3)?,
187 turn_id: row.get(4)?,
188 flow_run_id: row.get(5)?,
189 payload: row.get(6)?,
190 },
191 row.get::<_, String>(7)?,
192 ))
193 })?;
194 let mut hits = Vec::new();
195 for row in rows {
196 let (event, text) = row?;
197 if re.is_match(&text) {
198 hits.push(event);
199 if hits.len() >= limit {
200 break;
201 }
202 }
203 }
204 Ok(hits)
205 }
206
207 pub fn read_events_paginated(
208 &self,
209 session_id: &str,
210 offset: usize,
211 limit: usize,
212 kinds: Option<&[&str]>,
213 ) -> Result<Vec<ProjectEventRow>> {
214 let conn = self.conn();
215 let (sql, params): (String, Vec<Box<dyn rusqlite::ToSql>>) = match kinds {
216 Some(ks) if !ks.is_empty() => {
217 let placeholders = ks.iter().map(|_| "?").collect::<Vec<_>>().join(",");
218 let sql = format!(
219 "SELECT session_id, seq, ts, kind, turn_id, flow_run_id, payload \
220 FROM events WHERE session_id = ? AND kind IN ({placeholders}) \
221 ORDER BY seq LIMIT ? OFFSET ?"
222 );
223 let mut params: Vec<Box<dyn rusqlite::ToSql>> =
224 vec![Box::new(session_id.to_string())];
225 for k in ks {
226 params.push(Box::new(k.to_string()));
227 }
228 params.push(Box::new(limit as i64));
229 params.push(Box::new(offset as i64));
230 (sql, params)
231 }
232 _ => (
233 "SELECT session_id, seq, ts, kind, turn_id, flow_run_id, payload \
234 FROM events WHERE session_id = ? ORDER BY seq LIMIT ? OFFSET ?"
235 .into(),
236 vec![
237 Box::new(session_id.to_string()),
238 Box::new(limit as i64),
239 Box::new(offset as i64),
240 ],
241 ),
242 };
243 let mut stmt = conn.prepare(&sql)?;
244 let param_refs: Vec<&dyn rusqlite::ToSql> = params.iter().map(|b| b.as_ref()).collect();
245 let rows = stmt.query_map(param_refs.as_slice(), project_event_row_from)?;
246 collect(rows)
247 }
248
249 pub fn count_search_hits(&self, query: &str, session_filter: Option<&str>) -> Result<u64> {
250 let conn = self.conn();
251 if let Some(pattern) = parse_regex_query(query) {
252 return Ok(self
253 .search_events_regex(&pattern, session_filter, 10000, &conn)?
254 .len() as u64);
255 }
256 if query.chars().any(is_cjk_char) {
257 return self.count_search_hits_like(query, session_filter, &conn);
258 }
259 let (sql, params): (String, Vec<Box<dyn rusqlite::ToSql>>) = match session_filter {
260 Some(sid) => (
261 "SELECT COUNT(*) FROM events e JOIN events_fts f ON f.rowid = e.id \
262 WHERE f.events_fts MATCH ?1 AND e.session_id = ?2"
263 .into(),
264 vec![Box::new(query.to_string()), Box::new(sid.to_string())],
265 ),
266 None => (
267 "SELECT COUNT(*) FROM events e JOIN events_fts f ON f.rowid = e.id \
268 WHERE f.events_fts MATCH ?1"
269 .into(),
270 vec![Box::new(query.to_string())],
271 ),
272 };
273 let mut stmt = conn.prepare(&sql)?;
274 let param_refs: Vec<&dyn rusqlite::ToSql> = params.iter().map(|b| b.as_ref()).collect();
275 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
276 Ok(count as u64)
277 }
278
279 fn count_search_hits_like(
280 &self,
281 query: &str,
282 session_filter: Option<&str>,
283 conn: &std::sync::MutexGuard<'_, rusqlite::Connection>,
284 ) -> Result<u64> {
285 let pattern = format!("%{query}%");
286 let (sql, params): (&str, Vec<Box<dyn rusqlite::ToSql>>) = match session_filter {
287 Some(sid) => (
288 "SELECT COUNT(*) FROM events e JOIN events_fts f ON f.rowid = e.id \
289 WHERE f.text_content LIKE ?1 AND e.session_id = ?2",
290 vec![Box::new(pattern), Box::new(sid.to_string())],
291 ),
292 None => (
293 "SELECT COUNT(*) FROM events e JOIN events_fts f ON f.rowid = e.id \
294 WHERE f.text_content LIKE ?1",
295 vec![Box::new(pattern)],
296 ),
297 };
298 let mut stmt = conn.prepare(sql)?;
299 let param_refs: Vec<&dyn rusqlite::ToSql> = params.iter().map(|b| b.as_ref()).collect();
300 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
301 Ok(count as u64)
302 }
303
304 pub fn delete_events_for_session(&self, session_id: &str) -> rusqlite::Result<usize> {
305 let mut conn = self.conn.lock().unwrap();
306 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
307 tx.execute(
308 "DELETE FROM events_fts WHERE rowid IN (SELECT id FROM events WHERE session_id = ?)",
309 rusqlite::params![session_id],
310 )?;
311 let n = tx.execute(
312 "DELETE FROM events WHERE session_id = ?",
313 rusqlite::params![session_id],
314 )?;
315 tx.commit()?;
316 Ok(n)
317 }
318
319 pub fn insert_project_event_raw(&self, row: ProjectEventInsert<'_>) -> rusqlite::Result<i64> {
320 let mut conn = self.conn.lock().unwrap();
321 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
322 tx.execute(
323 "INSERT OR REPLACE INTO events \
324 (session_id, seq, ts, kind, turn_id, flow_run_id, payload) \
325 VALUES (?, ?, ?, ?, ?, ?, ?)",
326 rusqlite::params![
327 row.session_id,
328 row.seq,
329 row.ts,
330 row.kind,
331 row.turn_id,
332 row.flow_run_id,
333 row.payload_json,
334 ],
335 )?;
336 let id = tx.last_insert_rowid();
337 tx.execute(
338 "INSERT OR REPLACE INTO events_fts (rowid, text_content) VALUES (?, ?)",
339 rusqlite::params![id, row.text_content],
340 )?;
341 tx.commit()?;
342 Ok(id)
343 }
344
345 pub fn find_project_events_around(
346 &self,
347 session_id: &str,
348 seq: u64,
349 window: usize,
350 ) -> Result<Vec<ProjectEventRow>> {
351 let low = seq.saturating_sub(window as u64) as i64;
352 let high = seq.saturating_add(window as u64) as i64;
353 let conn = self.conn();
354 let mut stmt = conn.prepare(
355 "SELECT session_id, seq, ts, kind, turn_id, flow_run_id, payload FROM events \
356 WHERE session_id = ? AND seq BETWEEN ? AND ? ORDER BY seq",
357 )?;
358 let rows = stmt.query_map(
359 rusqlite::params![session_id, low, high],
360 project_event_row_from,
361 )?;
362 collect(rows)
363 }
364
365 pub fn find_project_events_by_anchor(
366 &self,
367 session_id: &str,
368 kind: AnchorKind,
369 id: &str,
370 ) -> Result<Vec<ProjectEventRow>> {
371 let sql = format!(
372 "SELECT session_id, seq, ts, kind, turn_id, flow_run_id, payload FROM events \
373 WHERE session_id = ? AND {} = ? ORDER BY seq",
374 kind.events_column()
375 );
376 let conn = self.conn();
377 let mut stmt = conn.prepare(&sql)?;
378 let rows = stmt.query_map(rusqlite::params![session_id, id], project_event_row_from)?;
379 collect(rows)
380 }
381
382 pub fn rebuild_events_from_sessions(
383 &self,
384 sessions_root: &Path,
385 fingerprint: &str,
386 ) -> Result<RebuildStats> {
387 let mut stats = RebuildStats::default();
388 let read = match std::fs::read_dir(sessions_root) {
389 Ok(r) => r,
390 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(stats),
391 Err(e) => return Err(e).context(format!("read_dir {}", sessions_root.display())),
392 };
393 for entry in read.flatten() {
394 let dir = entry.path();
395 if !dir.is_dir() {
396 continue;
397 }
398 let Some(meta) = crate::session_meta::SessionMeta::load(&dir) else {
399 continue;
400 };
401 if meta.project_fingerprint.as_deref() != Some(fingerprint) {
402 continue;
403 }
404 let sid = dir
405 .file_name()
406 .map(|s| s.to_string_lossy().to_string())
407 .unwrap_or_default();
408 let jsonl = dir.join("events.jsonl");
409 let envelopes = match crate::event_log::reader::read_event_envelopes(&jsonl) {
410 Ok(envelopes) => envelopes,
411 Err(_) => continue,
412 };
413 for envelope in &envelopes {
414 let value = serde_json::to_value(envelope)?;
415 let seq = envelope.seq as i64;
416 let ts = envelope.ts.to_rfc3339();
417 let kind = value.get("type").and_then(|v| v.as_str()).unwrap_or("");
418 let turn_id = value.get("turn_id").and_then(|v| v.as_str());
419 let flow_run_id = value
420 .get("run_id")
421 .or_else(|| value.get("flow_run_id"))
422 .and_then(|v| v.as_str());
423 let text_content = value
424 .get("message")
425 .and_then(|m| m.get("parts"))
426 .and_then(|p| p.as_array())
427 .map(|parts| {
428 parts
429 .iter()
430 .filter_map(|p| p.get("text").and_then(|t| t.as_str()))
431 .collect::<Vec<_>>()
432 .join("")
433 })
434 .unwrap_or_default();
435 let payload_json = serde_json::to_string(&value).unwrap_or_default();
436 self.insert_project_event_raw(ProjectEventInsert {
437 session_id: &sid,
438 seq,
439 ts: &ts,
440 kind,
441 turn_id,
442 flow_run_id,
443 text_content: &text_content,
444 payload_json: &payload_json,
445 })?;
446 stats.rebuilt += 1;
447 }
448 }
449 Ok(stats)
450 }
451
452 pub fn find_by_anchor(&self, kind: AnchorKind, id: &str) -> Result<Vec<(String, String)>> {
453 let sql =
454 "SELECT subject_kind, subject_id FROM anchors WHERE kind = ? AND ref = ? ORDER BY id";
455 let conn = self.conn();
456 let mut stmt = conn.prepare(sql)?;
457 let rows = stmt.query_map(rusqlite::params![kind.anchor_tag(), id], |row| {
458 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
459 })?;
460 let mut out = Vec::new();
461 for r in rows {
462 out.push(r?);
463 }
464 Ok(out)
465 }
466}
467
468#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
469pub struct RebuildStats {
470 pub rebuilt: usize,
471 pub skipped: usize,
472}
473
474#[derive(Debug, Clone, Copy, PartialEq, Eq)]
475pub enum AnchorKind {
476 TurnId,
477 FlowRunId,
478}
479
480impl AnchorKind {
481 fn events_column(self) -> &'static str {
482 match self {
483 AnchorKind::TurnId => "turn_id",
484 AnchorKind::FlowRunId => "flow_run_id",
485 }
486 }
487
488 fn anchor_tag(self) -> &'static str {
489 match self {
490 AnchorKind::TurnId => "turn",
491 AnchorKind::FlowRunId => "flow_run",
492 }
493 }
494}
495
496pub struct ProjectEventInsert<'a> {
497 pub session_id: &'a str,
498 pub seq: i64,
499 pub ts: &'a str,
500 pub kind: &'a str,
501 pub turn_id: Option<&'a str>,
502 pub flow_run_id: Option<&'a str>,
503 pub text_content: &'a str,
504 pub payload_json: &'a str,
505}
506
507#[derive(Debug, Clone, PartialEq, Eq)]
508pub struct ProjectEventRow {
509 pub session_id: String,
510 pub seq: u64,
511 pub ts: String,
512 pub kind: String,
513 pub turn_id: Option<String>,
514 pub flow_run_id: Option<String>,
515 pub payload: String,
516}
517
518fn project_event_row_from(row: &rusqlite::Row<'_>) -> rusqlite::Result<ProjectEventRow> {
519 Ok(ProjectEventRow {
520 session_id: row.get(0)?,
521 seq: row.get::<_, i64>(1)? as u64,
522 ts: row.get(2)?,
523 kind: row.get(3)?,
524 turn_id: row.get(4)?,
525 flow_run_id: row.get(5)?,
526 payload: row.get(6)?,
527 })
528}
529
530fn collect<T>(
531 iter: rusqlite::MappedRows<'_, impl FnMut(&rusqlite::Row<'_>) -> rusqlite::Result<T>>,
532) -> Result<Vec<T>> {
533 let mut out = Vec::new();
534 for r in iter {
535 out.push(r.map_err(|e| anyhow::anyhow!(e))?);
536 }
537 Ok(out)
538}
539
540fn parse_regex_query(query: &str) -> Option<String> {
541 let q = query.trim();
542 if q.len() >= 2 && q.starts_with('/') && q.ends_with('/') {
543 Some(q[1..q.len() - 1].to_string())
544 } else {
545 None
546 }
547}
548
549fn is_cjk_char(ch: char) -> bool {
550 matches!(ch as u32,
551 0x4E00..=0x9FFF | 0x3400..=0x4DBF | 0x20000..=0x2A6DF | 0x3040..=0x309F | 0x30A0..=0x30FF | 0xAC00..=0xD7AF )
558}
559
560const PROJECT_SCHEMA: &str = r#"
561CREATE TABLE IF NOT EXISTS events (
562 id INTEGER PRIMARY KEY AUTOINCREMENT,
563 session_id TEXT NOT NULL,
564 seq INTEGER NOT NULL,
565 ts TEXT NOT NULL,
566 kind TEXT NOT NULL,
567 turn_id TEXT,
568 flow_run_id TEXT,
569 payload TEXT NOT NULL,
570 UNIQUE (session_id, seq)
571);
572CREATE INDEX IF NOT EXISTS events_session ON events(session_id);
573CREATE INDEX IF NOT EXISTS events_kind ON events(kind);
574CREATE INDEX IF NOT EXISTS events_turn ON events(turn_id);
575CREATE INDEX IF NOT EXISTS events_flow ON events(flow_run_id);
576
577CREATE VIRTUAL TABLE IF NOT EXISTS events_fts USING fts5(
578 text_content,
579 tokenize='porter unicode61'
580);
581
582CREATE TABLE IF NOT EXISTS anchors (
583 id INTEGER PRIMARY KEY AUTOINCREMENT,
584 kind TEXT NOT NULL,
585 ref TEXT NOT NULL,
586 subject_kind TEXT NOT NULL,
587 subject_id TEXT NOT NULL,
588 session_id TEXT,
589 created_at TEXT NOT NULL
590);
591CREATE INDEX IF NOT EXISTS anchors_lookup ON anchors(kind, ref);
592CREATE INDEX IF NOT EXISTS anchors_subject ON anchors(subject_kind, subject_id);
593
594CREATE TABLE IF NOT EXISTS confessions (
595 id TEXT PRIMARY KEY,
596 trigger TEXT NOT NULL,
597 rule_violated TEXT NOT NULL,
598 what_i_did TEXT NOT NULL,
599 why TEXT NOT NULL,
600 mitigation TEXT NOT NULL,
601 body TEXT NOT NULL,
602 created_at TEXT NOT NULL
603);
604CREATE VIRTUAL TABLE IF NOT EXISTS confessions_fts USING fts5(
605 trigger, rule_violated, what_i_did, why, mitigation, body,
606 tokenize='porter unicode61'
607);
608
609CREATE TABLE IF NOT EXISTS spec_entries (
610 id TEXT PRIMARY KEY,
611 feature TEXT NOT NULL,
612 phase TEXT NOT NULL,
613 content TEXT NOT NULL,
614 ts TEXT NOT NULL
615);
616CREATE INDEX IF NOT EXISTS spec_entries_feature ON spec_entries(feature);
617CREATE VIRTUAL TABLE IF NOT EXISTS spec_entries_fts USING fts5(
618 content, tokenize='porter unicode61'
619);
620
621CREATE TABLE IF NOT EXISTS spec_deviations (
622 id TEXT PRIMARY KEY,
623 feature TEXT NOT NULL,
624 section TEXT NOT NULL,
625 delta TEXT NOT NULL,
626 reason TEXT NOT NULL,
627 ts TEXT NOT NULL
628);
629CREATE INDEX IF NOT EXISTS spec_deviations_feature ON spec_deviations(feature);
630CREATE VIRTUAL TABLE IF NOT EXISTS spec_deviations_fts USING fts5(
631 delta, reason, tokenize='porter unicode61'
632);
633"#;
634
635#[cfg(test)]
636mod tests {
637 use super::*;
638 use rusqlite::params;
639
640 fn tables_in(idx: &AnchorIndex) -> Vec<String> {
641 let conn = idx.conn();
642 let mut stmt = conn
643 .prepare("SELECT name FROM sqlite_master WHERE type IN ('table', 'view') ORDER BY name")
644 .unwrap();
645 stmt.query_map(params![], |row| row.get::<_, String>(0))
646 .unwrap()
647 .filter_map(|r| r.ok())
648 .collect()
649 }
650
651 #[test]
652 fn fts_search_finds_cjk_substring() {
653 let dir = tempfile::tempdir().unwrap();
654 let idx = AnchorIndex::open_project(dir.path()).unwrap();
655 seed_project_event(
656 &idx,
657 "sess-a",
658 1,
659 "user_msg",
660 None,
661 None,
662 "这是一个浮动面板的设计方案",
663 );
664 let hits = idx
666 .fts_search_project_events("浮动", Some("sess-a"), 10)
667 .unwrap();
668 assert_eq!(hits.len(), 1, "2-char CJK substring should match");
669 let hits = idx
671 .fts_search_project_events("浮动面板", Some("sess-a"), 10)
672 .unwrap();
673 assert_eq!(hits.len(), 1, "4-char CJK phrase should match");
674 }
675
676 #[test]
677 fn open_project_creates_all_four_tables() {
678 let dir = tempfile::tempdir().unwrap();
679 let idx = AnchorIndex::open_project(dir.path()).unwrap();
680 let tables = tables_in(&idx);
681 for expected in [
682 "anchors",
683 "confessions",
684 "confessions_fts",
685 "events",
686 "events_fts",
687 "spec_entries",
688 "spec_entries_fts",
689 "spec_deviations",
690 "spec_deviations_fts",
691 ] {
692 assert!(
693 tables.iter().any(|t| t == expected),
694 "missing {expected} in {tables:?}"
695 );
696 }
697 }
698
699 fn seed_project_event(
700 idx: &AnchorIndex,
701 sid: &str,
702 seq: i64,
703 kind: &str,
704 turn: Option<&str>,
705 flow: Option<&str>,
706 text: &str,
707 ) {
708 let payload = format!("{{\"seq\":{seq},\"session_id\":\"{sid}\"}}");
709 idx.insert_project_event_raw(ProjectEventInsert {
710 session_id: sid,
711 seq,
712 ts: "2026-07-05T00:00:00Z",
713 kind,
714 turn_id: turn,
715 flow_run_id: flow,
716 text_content: text,
717 payload_json: &payload,
718 })
719 .unwrap();
720 }
721
722 #[test]
723 fn fts_search_project_events_filters_by_session_when_requested() {
724 let dir = tempfile::tempdir().unwrap();
725 let idx = AnchorIndex::open_project(dir.path()).unwrap();
726 seed_project_event(
727 &idx,
728 "sess-a",
729 1,
730 "user_msg",
731 Some("t1"),
732 None,
733 "hello sqlite fts",
734 );
735 seed_project_event(
736 &idx,
737 "sess-b",
738 1,
739 "user_msg",
740 Some("t2"),
741 None,
742 "sqlite from other session",
743 );
744 seed_project_event(
745 &idx,
746 "sess-a",
747 2,
748 "assistant_msg",
749 Some("t1"),
750 None,
751 "no match here",
752 );
753
754 let scoped = idx
755 .fts_search_project_events("sqlite", Some("sess-a"), 10)
756 .unwrap();
757 assert_eq!(scoped.len(), 1);
758 assert_eq!(scoped[0].session_id, "sess-a");
759 assert_eq!(scoped[0].seq, 1);
760
761 let all = idx.fts_search_project_events("sqlite", None, 10).unwrap();
762 assert_eq!(all.len(), 2);
763 }
764
765 #[test]
766 fn insert_project_event_upsert_is_idempotent() {
767 let dir = tempfile::tempdir().unwrap();
768 let idx = AnchorIndex::open_project(dir.path()).unwrap();
769 seed_project_event(&idx, "sess-a", 42, "user_msg", None, None, "uniquetoken");
770 seed_project_event(&idx, "sess-a", 42, "user_msg", None, None, "uniquetoken");
771 let hits = idx
772 .fts_search_project_events("uniquetoken", None, 10)
773 .unwrap();
774 assert_eq!(hits.len(), 1, "same (session_id, seq) must upsert");
775 }
776
777 #[test]
778 fn find_project_events_around_scopes_to_session() {
779 let dir = tempfile::tempdir().unwrap();
780 let idx = AnchorIndex::open_project(dir.path()).unwrap();
781 seed_project_event(&idx, "sess-a", 1, "flow_start", None, Some("r"), "");
782 seed_project_event(&idx, "sess-a", 2, "user_msg", Some("t"), None, "hello");
783 seed_project_event(
784 &idx,
785 "sess-a",
786 3,
787 "assistant_msg",
788 Some("t"),
789 Some("r"),
790 "world",
791 );
792 seed_project_event(&idx, "sess-a", 4, "flow_end", None, Some("r"), "");
793 seed_project_event(&idx, "sess-b", 3, "user_msg", None, None, "unrelated");
794
795 let rows = idx.find_project_events_around("sess-a", 3, 1).unwrap();
796 let seqs: Vec<u64> = rows.iter().map(|r| r.seq).collect();
797 assert_eq!(seqs, vec![2, 3, 4]);
798 assert!(rows.iter().all(|r| r.session_id == "sess-a"));
799 }
800
801 #[test]
802 fn find_project_events_by_anchor_scopes_to_session() {
803 let dir = tempfile::tempdir().unwrap();
804 let idx = AnchorIndex::open_project(dir.path()).unwrap();
805 seed_project_event(&idx, "sess-a", 1, "flow_start", None, Some("run-x"), "");
806 seed_project_event(&idx, "sess-a", 2, "flow_end", None, Some("run-x"), "");
807 seed_project_event(&idx, "sess-b", 1, "flow_start", None, Some("run-x"), "");
808
809 let rows = idx
810 .find_project_events_by_anchor("sess-a", AnchorKind::FlowRunId, "run-x")
811 .unwrap();
812 assert_eq!(rows.len(), 2);
813 assert!(rows.iter().all(|r| r.session_id == "sess-a"));
814 }
815
816 #[test]
817 fn rebuild_events_from_sessions_filters_by_fingerprint() {
818 let sessions = tempfile::tempdir().unwrap();
819 let index_dir = tempfile::tempdir().unwrap();
820
821 let sess_ours = sessions.path().join("sess-a");
822 let sess_theirs = sessions.path().join("sess-b");
823 std::fs::create_dir_all(&sess_ours).unwrap();
824 std::fs::create_dir_all(&sess_theirs).unwrap();
825
826 crate::session_meta::SessionMeta {
827 project_fingerprint: Some("ours".into()),
828 ..Default::default()
829 }
830 .save(&sess_ours)
831 .unwrap();
832 crate::session_meta::SessionMeta {
833 project_fingerprint: Some("theirs".into()),
834 ..Default::default()
835 }
836 .save(&sess_theirs)
837 .unwrap();
838
839 let tid1 = uuid::Uuid::now_v7();
840 let rid1 = uuid::Uuid::now_v7();
841 let events_ours = format!(
842 r#"{{"type":"user_msg","seq":1,"turn_id":"{tid1}","ts":"2026-07-05T00:00:00Z","message":{{"role":"user","parts":[{{"type":"text","text":"needle in ours"}}],"turn_id":"{tid1}"}}}}
843{{"type":"flow_end","seq":2,"run_id":"{rid1}","flow_name":"demo","status":{{"kind":"ok"}},"ts":"2026-07-05T00:00:01Z"}}"#
844 );
845 std::fs::write(sess_ours.join("events.jsonl"), events_ours).unwrap();
846
847 let tid2 = uuid::Uuid::now_v7();
848 let events_theirs = format!(
849 r#"{{"type":"user_msg","seq":1,"turn_id":"{tid2}","ts":"2026-07-05T00:00:00Z","message":{{"role":"user","parts":[{{"type":"text","text":"needle in theirs"}}],"turn_id":"{tid2}"}}}}"#
850 );
851 std::fs::write(sess_theirs.join("events.jsonl"), events_theirs).unwrap();
852
853 let idx = AnchorIndex::open_project(index_dir.path()).unwrap();
854 let stats = idx
855 .rebuild_events_from_sessions(sessions.path(), "ours")
856 .unwrap();
857 assert_eq!(stats.rebuilt, 2);
858
859 let hits = idx.fts_search_project_events("needle", None, 10).unwrap();
860 assert_eq!(hits.len(), 1);
861 assert_eq!(hits[0].session_id, "sess-a");
862 }
863
864 #[test]
865 fn reopening_the_same_db_is_idempotent() {
866 let dir = tempfile::tempdir().unwrap();
867 let _one = AnchorIndex::open_project(dir.path()).unwrap();
868 let two = AnchorIndex::open_project(dir.path()).unwrap();
869 let tables = tables_in(&two);
870 assert!(tables.iter().any(|t| t == "confessions"));
871 }
872
873 #[test]
874 fn find_by_anchor_returns_subject_kinds_from_project_db() {
875 let dir = tempfile::tempdir().unwrap();
876 let idx = AnchorIndex::open_project(dir.path()).unwrap();
877 {
878 let conn = idx.conn();
879 conn.execute(
880 "INSERT INTO anchors (kind, ref, subject_kind, subject_id, created_at) VALUES (?, ?, ?, ?, ?)",
881 rusqlite::params!["flow_run", "run-xyz", "confession", "cid-1", "2026-07-05T00:00:00Z"],
882 )
883 .unwrap();
884 conn.execute(
885 "INSERT INTO anchors (kind, ref, subject_kind, subject_id, created_at) VALUES (?, ?, ?, ?, ?)",
886 rusqlite::params!["flow_run", "run-xyz", "spec_entry", "sid-1", "2026-07-05T00:00:00Z"],
887 )
888 .unwrap();
889 conn.execute(
890 "INSERT INTO anchors (kind, ref, subject_kind, subject_id, created_at) VALUES (?, ?, ?, ?, ?)",
891 rusqlite::params!["turn_id", "turn-99", "confession", "cid-2", "2026-07-05T00:00:00Z"],
892 )
893 .unwrap();
894 }
895 let hits = idx
896 .find_by_anchor(AnchorKind::FlowRunId, "run-xyz")
897 .unwrap();
898 assert_eq!(
899 hits,
900 vec![
901 ("confession".to_string(), "cid-1".to_string()),
902 ("spec_entry".to_string(), "sid-1".to_string()),
903 ]
904 );
905 }
906
907 #[test]
908 fn wal_journal_mode_is_set() {
909 let dir = tempfile::tempdir().unwrap();
910 let idx = AnchorIndex::open_project(dir.path()).unwrap();
911 let conn = idx.conn();
912 let mode: String = conn
913 .query_row("PRAGMA journal_mode", params![], |row| row.get(0))
914 .unwrap();
915 assert_eq!(mode.to_lowercase(), "wal");
916 }
917
918 #[test]
919 fn count_events_returns_zero_for_empty_session() {
920 let dir = tempfile::tempdir().unwrap();
921 let idx = AnchorIndex::open_project(dir.path()).unwrap();
922 assert_eq!(idx.count_events("no-such-session", None).unwrap(), 0);
923 }
924
925 #[test]
926 fn count_events_counts_all_kinds() {
927 let dir = tempfile::tempdir().unwrap();
928 let idx = AnchorIndex::open_project(dir.path()).unwrap();
929 seed_project_event(&idx, "s1", 1, "user_msg", None, None, "hello");
930 seed_project_event(&idx, "s1", 2, "assistant_msg", None, None, "hi");
931 seed_project_event(&idx, "s1", 3, "tool_result_msg", None, None, "ok");
932 seed_project_event(&idx, "s2", 1, "user_msg", None, None, "other");
933 assert_eq!(idx.count_events("s1", None).unwrap(), 3);
934 assert_eq!(idx.count_events("s2", None).unwrap(), 1);
935 }
936
937 #[test]
938 fn count_events_filters_by_kind() {
939 let dir = tempfile::tempdir().unwrap();
940 let idx = AnchorIndex::open_project(dir.path()).unwrap();
941 seed_project_event(&idx, "s1", 1, "user_msg", None, None, "a");
942 seed_project_event(&idx, "s1", 2, "assistant_msg", None, None, "b");
943 seed_project_event(&idx, "s1", 3, "tool_result_msg", None, None, "c");
944 seed_project_event(&idx, "s1", 4, "user_msg", None, None, "d");
945 assert_eq!(idx.count_events("s1", Some(&["user_msg"])).unwrap(), 2);
946 assert_eq!(
947 idx.count_events("s1", Some(&["user_msg", "assistant_msg"]))
948 .unwrap(),
949 3
950 );
951 assert_eq!(idx.count_events("s1", Some(&["system_msg"])).unwrap(), 0);
952 }
953
954 #[test]
955 fn read_events_paginated_returns_ordered_by_seq() {
956 let dir = tempfile::tempdir().unwrap();
957 let idx = AnchorIndex::open_project(dir.path()).unwrap();
958 seed_project_event(&idx, "s1", 3, "user_msg", None, None, "third");
959 seed_project_event(&idx, "s1", 1, "user_msg", None, None, "first");
960 seed_project_event(&idx, "s1", 2, "user_msg", None, None, "second");
961 let rows = idx.read_events_paginated("s1", 0, 10, None).unwrap();
962 assert_eq!(rows.len(), 3);
963 assert_eq!(rows[0].seq, 1);
964 assert_eq!(rows[1].seq, 2);
965 assert_eq!(rows[2].seq, 3);
966 }
967
968 #[test]
969 fn read_events_paginated_respects_offset_and_limit() {
970 let dir = tempfile::tempdir().unwrap();
971 let idx = AnchorIndex::open_project(dir.path()).unwrap();
972 for i in 1..=5 {
973 seed_project_event(&idx, "s1", i, "user_msg", None, None, &format!("msg {i}"));
974 }
975 let rows = idx.read_events_paginated("s1", 1, 2, None).unwrap();
976 assert_eq!(rows.len(), 2);
977 assert_eq!(rows[0].seq, 2);
978 assert_eq!(rows[1].seq, 3);
979 }
980
981 #[test]
982 fn read_events_paginated_filters_by_kind() {
983 let dir = tempfile::tempdir().unwrap();
984 let idx = AnchorIndex::open_project(dir.path()).unwrap();
985 seed_project_event(&idx, "s1", 1, "user_msg", None, None, "u1");
986 seed_project_event(&idx, "s1", 2, "assistant_msg", None, None, "a1");
987 seed_project_event(&idx, "s1", 3, "user_msg", None, None, "u2");
988 let rows = idx
989 .read_events_paginated("s1", 0, 10, Some(&["user_msg"]))
990 .unwrap();
991 assert_eq!(rows.len(), 2);
992 assert_eq!(rows[0].seq, 1);
993 assert_eq!(rows[1].seq, 3);
994 }
995
996 #[test]
997 fn read_events_paginated_empty_session() {
998 let dir = tempfile::tempdir().unwrap();
999 let idx = AnchorIndex::open_project(dir.path()).unwrap();
1000 let rows = idx.read_events_paginated("no-such", 0, 10, None).unwrap();
1001 assert!(rows.is_empty());
1002 }
1003
1004 #[test]
1005 fn count_search_hits_fts() {
1006 let dir = tempfile::tempdir().unwrap();
1007 let idx = AnchorIndex::open_project(dir.path()).unwrap();
1008 seed_project_event(&idx, "s1", 1, "user_msg", None, None, "hello world");
1009 seed_project_event(&idx, "s1", 2, "user_msg", None, None, "hello again");
1010 seed_project_event(&idx, "s1", 3, "user_msg", None, None, "goodbye");
1011 seed_project_event(&idx, "s2", 1, "user_msg", None, None, "hello other");
1012 assert_eq!(idx.count_search_hits("hello", None).unwrap(), 3);
1013 assert_eq!(idx.count_search_hits("hello", Some("s1")).unwrap(), 2);
1014 assert_eq!(idx.count_search_hits("goodbye", None).unwrap(), 1);
1015 assert_eq!(idx.count_search_hits("nomatch", None).unwrap(), 0);
1016 }
1017
1018 #[test]
1019 fn count_search_hits_cjk() {
1020 let dir = tempfile::tempdir().unwrap();
1021 let idx = AnchorIndex::open_project(dir.path()).unwrap();
1022 seed_project_event(&idx, "s1", 1, "user_msg", None, None, "浮动面板设计");
1023 seed_project_event(&idx, "s1", 2, "user_msg", None, None, "另一个消息");
1024 assert_eq!(idx.count_search_hits("浮动", None).unwrap(), 1);
1025 assert_eq!(idx.count_search_hits("浮动", Some("s1")).unwrap(), 1);
1026 assert_eq!(idx.count_search_hits("不存在", None).unwrap(), 0);
1027 }
1028}