Skip to main content

atman_runtime/
index.rs

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        // CJK text doesn't tokenize well under unicode61 — use LIKE for CJK queries.
86        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   |  // CJK Unified Ideographs
508        0x3400..=0x4DBF   |  // CJK Extension A
509        0x20000..=0x2A6DF |  // CJK Extension B
510        0x3040..=0x309F   |  // Hiragana
511        0x30A0..=0x30FF   |  // Katakana
512        0xAC00..=0xD7AF      // Hangul Syllables
513    )
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        // 2-char CJK substring — would fail with FTS5 unicode61
621        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        // 4-char CJK phrase
626        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}