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