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        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   |  // CJK Unified Ideographs
552        0x3400..=0x4DBF   |  // CJK Extension A
553        0x20000..=0x2A6DF |  // CJK Extension B
554        0x3040..=0x309F   |  // Hiragana
555        0x30A0..=0x30FF   |  // Katakana
556        0xAC00..=0xD7AF      // Hangul Syllables
557    )
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        // 2-char CJK substring — would fail with FTS5 unicode61
665        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        // 4-char CJK phrase
670        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}