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        let (sql, params): (String, Vec<Box<dyn rusqlite::ToSql>>) = match session_filter {
86            Some(sid) => (
87                "SELECT e.session_id, e.seq, e.ts, e.kind, e.turn_id, e.flow_run_id, e.payload \
88                 FROM events e JOIN events_fts f ON f.rowid = e.id \
89                 WHERE f.events_fts MATCH ?1 AND e.session_id = ?2 \
90                 ORDER BY e.id DESC LIMIT ?3"
91                    .into(),
92                vec![
93                    Box::new(query.to_string()),
94                    Box::new(sid.to_string()),
95                    Box::new(limit as i64),
96                ],
97            ),
98            None => (
99                "SELECT e.session_id, e.seq, e.ts, e.kind, e.turn_id, e.flow_run_id, e.payload \
100                 FROM events e JOIN events_fts f ON f.rowid = e.id \
101                 WHERE f.events_fts MATCH ?1 \
102                 ORDER BY e.id DESC LIMIT ?2"
103                    .into(),
104                vec![Box::new(query.to_string()), Box::new(limit as i64)],
105            ),
106        };
107        let mut stmt = conn.prepare(&sql)?;
108        let param_refs: Vec<&dyn rusqlite::ToSql> = params.iter().map(|b| b.as_ref()).collect();
109        let rows = stmt.query_map(param_refs.as_slice(), project_event_row_from)?;
110        collect(rows)
111    }
112
113    pub fn delete_events_for_session(&self, session_id: &str) -> rusqlite::Result<usize> {
114        let mut conn = self.conn.lock().unwrap();
115        let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
116        tx.execute(
117            "DELETE FROM events_fts WHERE rowid IN (SELECT id FROM events WHERE session_id = ?)",
118            rusqlite::params![session_id],
119        )?;
120        let n = tx.execute(
121            "DELETE FROM events WHERE session_id = ?",
122            rusqlite::params![session_id],
123        )?;
124        tx.commit()?;
125        Ok(n)
126    }
127
128    pub fn insert_project_event_raw(&self, row: ProjectEventInsert<'_>) -> rusqlite::Result<i64> {
129        let mut conn = self.conn.lock().unwrap();
130        let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
131        tx.execute(
132            "INSERT OR REPLACE INTO events \
133             (session_id, seq, ts, kind, turn_id, flow_run_id, payload) \
134             VALUES (?, ?, ?, ?, ?, ?, ?)",
135            rusqlite::params![
136                row.session_id,
137                row.seq,
138                row.ts,
139                row.kind,
140                row.turn_id,
141                row.flow_run_id,
142                row.payload_json,
143            ],
144        )?;
145        let id = tx.last_insert_rowid();
146        tx.execute(
147            "INSERT OR REPLACE INTO events_fts (rowid, text_content) VALUES (?, ?)",
148            rusqlite::params![id, row.text_content],
149        )?;
150        tx.commit()?;
151        Ok(id)
152    }
153
154    pub fn find_project_events_around(
155        &self,
156        session_id: &str,
157        seq: u64,
158        window: usize,
159    ) -> Result<Vec<ProjectEventRow>> {
160        let low = seq.saturating_sub(window as u64) as i64;
161        let high = seq.saturating_add(window as u64) as i64;
162        let conn = self.conn();
163        let mut stmt = conn.prepare(
164            "SELECT session_id, seq, ts, kind, turn_id, flow_run_id, payload FROM events \
165             WHERE session_id = ? AND seq BETWEEN ? AND ? ORDER BY seq",
166        )?;
167        let rows = stmt.query_map(
168            rusqlite::params![session_id, low, high],
169            project_event_row_from,
170        )?;
171        collect(rows)
172    }
173
174    pub fn find_project_events_by_anchor(
175        &self,
176        session_id: &str,
177        kind: AnchorKind,
178        id: &str,
179    ) -> Result<Vec<ProjectEventRow>> {
180        let sql = format!(
181            "SELECT session_id, seq, ts, kind, turn_id, flow_run_id, payload FROM events \
182             WHERE session_id = ? AND {} = ? ORDER BY seq",
183            kind.events_column()
184        );
185        let conn = self.conn();
186        let mut stmt = conn.prepare(&sql)?;
187        let rows = stmt.query_map(rusqlite::params![session_id, id], project_event_row_from)?;
188        collect(rows)
189    }
190
191    pub fn rebuild_events_from_sessions(
192        &self,
193        sessions_root: &Path,
194        fingerprint: &str,
195    ) -> Result<RebuildStats> {
196        let mut stats = RebuildStats::default();
197        let read = match std::fs::read_dir(sessions_root) {
198            Ok(r) => r,
199            Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(stats),
200            Err(e) => return Err(e).context(format!("read_dir {}", sessions_root.display())),
201        };
202        for entry in read.flatten() {
203            let dir = entry.path();
204            if !dir.is_dir() {
205                continue;
206            }
207            let Some(meta) = crate::session_meta::SessionMeta::load(&dir) else {
208                continue;
209            };
210            if meta.project_fingerprint.as_deref() != Some(fingerprint) {
211                continue;
212            }
213            let sid = dir
214                .file_name()
215                .map(|s| s.to_string_lossy().to_string())
216                .unwrap_or_default();
217            let jsonl = dir.join("events.jsonl");
218            let envelopes = match crate::event_log::reader::read_event_envelopes(&jsonl) {
219                Ok(envelopes) => envelopes,
220                Err(_) => continue,
221            };
222            for envelope in &envelopes {
223                let value = serde_json::to_value(envelope)?;
224                let seq = envelope.seq as i64;
225                let ts = envelope.ts.to_rfc3339();
226                let kind = value.get("type").and_then(|v| v.as_str()).unwrap_or("");
227                let turn_id = value.get("turn_id").and_then(|v| v.as_str());
228                let flow_run_id = value
229                    .get("run_id")
230                    .or_else(|| value.get("flow_run_id"))
231                    .and_then(|v| v.as_str());
232                let text_content = value
233                    .get("message")
234                    .and_then(|m| m.get("parts"))
235                    .and_then(|p| p.as_array())
236                    .map(|parts| {
237                        parts
238                            .iter()
239                            .filter_map(|p| p.get("text").and_then(|t| t.as_str()))
240                            .collect::<Vec<_>>()
241                            .join("")
242                    })
243                    .unwrap_or_default();
244                let payload_json = serde_json::to_string(&value).unwrap_or_default();
245                self.insert_project_event_raw(ProjectEventInsert {
246                    session_id: &sid,
247                    seq,
248                    ts: &ts,
249                    kind,
250                    turn_id,
251                    flow_run_id,
252                    text_content: &text_content,
253                    payload_json: &payload_json,
254                })?;
255                stats.rebuilt += 1;
256            }
257        }
258        Ok(stats)
259    }
260
261    pub fn find_by_anchor(&self, kind: AnchorKind, id: &str) -> Result<Vec<(String, String)>> {
262        let sql =
263            "SELECT subject_kind, subject_id FROM anchors WHERE kind = ? AND ref = ? ORDER BY id";
264        let conn = self.conn();
265        let mut stmt = conn.prepare(sql)?;
266        let rows = stmt.query_map(rusqlite::params![kind.anchor_tag(), id], |row| {
267            Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
268        })?;
269        let mut out = Vec::new();
270        for r in rows {
271            out.push(r?);
272        }
273        Ok(out)
274    }
275}
276
277#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
278pub struct RebuildStats {
279    pub rebuilt: usize,
280    pub skipped: usize,
281}
282
283#[derive(Debug, Clone, Copy, PartialEq, Eq)]
284pub enum AnchorKind {
285    TurnId,
286    FlowRunId,
287}
288
289impl AnchorKind {
290    fn events_column(self) -> &'static str {
291        match self {
292            AnchorKind::TurnId => "turn_id",
293            AnchorKind::FlowRunId => "flow_run_id",
294        }
295    }
296
297    fn anchor_tag(self) -> &'static str {
298        match self {
299            AnchorKind::TurnId => "turn",
300            AnchorKind::FlowRunId => "flow_run",
301        }
302    }
303}
304
305pub struct ProjectEventInsert<'a> {
306    pub session_id: &'a str,
307    pub seq: i64,
308    pub ts: &'a str,
309    pub kind: &'a str,
310    pub turn_id: Option<&'a str>,
311    pub flow_run_id: Option<&'a str>,
312    pub text_content: &'a str,
313    pub payload_json: &'a str,
314}
315
316#[derive(Debug, Clone, PartialEq, Eq)]
317pub struct ProjectEventRow {
318    pub session_id: String,
319    pub seq: u64,
320    pub ts: String,
321    pub kind: String,
322    pub turn_id: Option<String>,
323    pub flow_run_id: Option<String>,
324    pub payload: String,
325}
326
327fn project_event_row_from(row: &rusqlite::Row<'_>) -> rusqlite::Result<ProjectEventRow> {
328    Ok(ProjectEventRow {
329        session_id: row.get(0)?,
330        seq: row.get::<_, i64>(1)? as u64,
331        ts: row.get(2)?,
332        kind: row.get(3)?,
333        turn_id: row.get(4)?,
334        flow_run_id: row.get(5)?,
335        payload: row.get(6)?,
336    })
337}
338
339fn collect<T>(
340    iter: rusqlite::MappedRows<'_, impl FnMut(&rusqlite::Row<'_>) -> rusqlite::Result<T>>,
341) -> Result<Vec<T>> {
342    let mut out = Vec::new();
343    for r in iter {
344        out.push(r.map_err(|e| anyhow::anyhow!(e))?);
345    }
346    Ok(out)
347}
348
349const PROJECT_SCHEMA: &str = r#"
350CREATE TABLE IF NOT EXISTS events (
351    id          INTEGER PRIMARY KEY AUTOINCREMENT,
352    session_id  TEXT    NOT NULL,
353    seq         INTEGER NOT NULL,
354    ts          TEXT    NOT NULL,
355    kind        TEXT    NOT NULL,
356    turn_id     TEXT,
357    flow_run_id TEXT,
358    payload     TEXT    NOT NULL,
359    UNIQUE (session_id, seq)
360);
361CREATE INDEX IF NOT EXISTS events_session ON events(session_id);
362CREATE INDEX IF NOT EXISTS events_kind    ON events(kind);
363CREATE INDEX IF NOT EXISTS events_turn    ON events(turn_id);
364CREATE INDEX IF NOT EXISTS events_flow    ON events(flow_run_id);
365
366CREATE VIRTUAL TABLE IF NOT EXISTS events_fts USING fts5(
367    text_content,
368    tokenize='porter unicode61'
369);
370
371CREATE TABLE IF NOT EXISTS anchors (
372    id            INTEGER PRIMARY KEY AUTOINCREMENT,
373    kind          TEXT NOT NULL,
374    ref           TEXT NOT NULL,
375    subject_kind  TEXT NOT NULL,
376    subject_id    TEXT NOT NULL,
377    session_id    TEXT,
378    created_at    TEXT NOT NULL
379);
380CREATE INDEX IF NOT EXISTS anchors_lookup  ON anchors(kind, ref);
381CREATE INDEX IF NOT EXISTS anchors_subject ON anchors(subject_kind, subject_id);
382
383CREATE TABLE IF NOT EXISTS confessions (
384    id            TEXT PRIMARY KEY,
385    trigger       TEXT NOT NULL,
386    rule_violated TEXT NOT NULL,
387    what_i_did    TEXT NOT NULL,
388    why           TEXT NOT NULL,
389    mitigation    TEXT NOT NULL,
390    body          TEXT NOT NULL,
391    created_at    TEXT NOT NULL
392);
393CREATE VIRTUAL TABLE IF NOT EXISTS confessions_fts USING fts5(
394    trigger, rule_violated, what_i_did, why, mitigation, body,
395    tokenize='porter unicode61'
396);
397
398CREATE TABLE IF NOT EXISTS spec_entries (
399    id      TEXT PRIMARY KEY,
400    feature TEXT NOT NULL,
401    phase   TEXT NOT NULL,
402    content TEXT NOT NULL,
403    ts      TEXT NOT NULL
404);
405CREATE INDEX IF NOT EXISTS spec_entries_feature ON spec_entries(feature);
406CREATE VIRTUAL TABLE IF NOT EXISTS spec_entries_fts USING fts5(
407    content, tokenize='porter unicode61'
408);
409
410CREATE TABLE IF NOT EXISTS spec_deviations (
411    id      TEXT PRIMARY KEY,
412    feature TEXT NOT NULL,
413    section TEXT NOT NULL,
414    delta   TEXT NOT NULL,
415    reason  TEXT NOT NULL,
416    ts      TEXT NOT NULL
417);
418CREATE INDEX IF NOT EXISTS spec_deviations_feature ON spec_deviations(feature);
419CREATE VIRTUAL TABLE IF NOT EXISTS spec_deviations_fts USING fts5(
420    delta, reason, tokenize='porter unicode61'
421);
422"#;
423
424#[cfg(test)]
425mod tests {
426    use super::*;
427    use rusqlite::params;
428
429    fn tables_in(idx: &AnchorIndex) -> Vec<String> {
430        let conn = idx.conn();
431        let mut stmt = conn
432            .prepare("SELECT name FROM sqlite_master WHERE type IN ('table', 'view') ORDER BY name")
433            .unwrap();
434        stmt.query_map(params![], |row| row.get::<_, String>(0))
435            .unwrap()
436            .filter_map(|r| r.ok())
437            .collect()
438    }
439
440    #[test]
441    fn open_project_creates_all_four_tables() {
442        let dir = tempfile::tempdir().unwrap();
443        let idx = AnchorIndex::open_project(dir.path()).unwrap();
444        let tables = tables_in(&idx);
445        for expected in [
446            "anchors",
447            "confessions",
448            "confessions_fts",
449            "events",
450            "events_fts",
451            "spec_entries",
452            "spec_entries_fts",
453            "spec_deviations",
454            "spec_deviations_fts",
455        ] {
456            assert!(
457                tables.iter().any(|t| t == expected),
458                "missing {expected} in {tables:?}"
459            );
460        }
461    }
462
463    fn seed_project_event(
464        idx: &AnchorIndex,
465        sid: &str,
466        seq: i64,
467        kind: &str,
468        turn: Option<&str>,
469        flow: Option<&str>,
470        text: &str,
471    ) {
472        let payload = format!("{{\"seq\":{seq},\"session_id\":\"{sid}\"}}");
473        idx.insert_project_event_raw(ProjectEventInsert {
474            session_id: sid,
475            seq,
476            ts: "2026-07-05T00:00:00Z",
477            kind,
478            turn_id: turn,
479            flow_run_id: flow,
480            text_content: text,
481            payload_json: &payload,
482        })
483        .unwrap();
484    }
485
486    #[test]
487    fn fts_search_project_events_filters_by_session_when_requested() {
488        let dir = tempfile::tempdir().unwrap();
489        let idx = AnchorIndex::open_project(dir.path()).unwrap();
490        seed_project_event(
491            &idx,
492            "sess-a",
493            1,
494            "user_msg",
495            Some("t1"),
496            None,
497            "hello sqlite fts",
498        );
499        seed_project_event(
500            &idx,
501            "sess-b",
502            1,
503            "user_msg",
504            Some("t2"),
505            None,
506            "sqlite from other session",
507        );
508        seed_project_event(
509            &idx,
510            "sess-a",
511            2,
512            "assistant_msg",
513            Some("t1"),
514            None,
515            "no match here",
516        );
517
518        let scoped = idx
519            .fts_search_project_events("sqlite", Some("sess-a"), 10)
520            .unwrap();
521        assert_eq!(scoped.len(), 1);
522        assert_eq!(scoped[0].session_id, "sess-a");
523        assert_eq!(scoped[0].seq, 1);
524
525        let all = idx.fts_search_project_events("sqlite", None, 10).unwrap();
526        assert_eq!(all.len(), 2);
527    }
528
529    #[test]
530    fn insert_project_event_upsert_is_idempotent() {
531        let dir = tempfile::tempdir().unwrap();
532        let idx = AnchorIndex::open_project(dir.path()).unwrap();
533        seed_project_event(&idx, "sess-a", 42, "user_msg", None, None, "uniquetoken");
534        seed_project_event(&idx, "sess-a", 42, "user_msg", None, None, "uniquetoken");
535        let hits = idx
536            .fts_search_project_events("uniquetoken", None, 10)
537            .unwrap();
538        assert_eq!(hits.len(), 1, "same (session_id, seq) must upsert");
539    }
540
541    #[test]
542    fn find_project_events_around_scopes_to_session() {
543        let dir = tempfile::tempdir().unwrap();
544        let idx = AnchorIndex::open_project(dir.path()).unwrap();
545        seed_project_event(&idx, "sess-a", 1, "flow_start", None, Some("r"), "");
546        seed_project_event(&idx, "sess-a", 2, "user_msg", Some("t"), None, "hello");
547        seed_project_event(
548            &idx,
549            "sess-a",
550            3,
551            "assistant_msg",
552            Some("t"),
553            Some("r"),
554            "world",
555        );
556        seed_project_event(&idx, "sess-a", 4, "flow_end", None, Some("r"), "");
557        seed_project_event(&idx, "sess-b", 3, "user_msg", None, None, "unrelated");
558
559        let rows = idx.find_project_events_around("sess-a", 3, 1).unwrap();
560        let seqs: Vec<u64> = rows.iter().map(|r| r.seq).collect();
561        assert_eq!(seqs, vec![2, 3, 4]);
562        assert!(rows.iter().all(|r| r.session_id == "sess-a"));
563    }
564
565    #[test]
566    fn find_project_events_by_anchor_scopes_to_session() {
567        let dir = tempfile::tempdir().unwrap();
568        let idx = AnchorIndex::open_project(dir.path()).unwrap();
569        seed_project_event(&idx, "sess-a", 1, "flow_start", None, Some("run-x"), "");
570        seed_project_event(&idx, "sess-a", 2, "flow_end", None, Some("run-x"), "");
571        seed_project_event(&idx, "sess-b", 1, "flow_start", None, Some("run-x"), "");
572
573        let rows = idx
574            .find_project_events_by_anchor("sess-a", AnchorKind::FlowRunId, "run-x")
575            .unwrap();
576        assert_eq!(rows.len(), 2);
577        assert!(rows.iter().all(|r| r.session_id == "sess-a"));
578    }
579
580    #[test]
581    fn rebuild_events_from_sessions_filters_by_fingerprint() {
582        let sessions = tempfile::tempdir().unwrap();
583        let index_dir = tempfile::tempdir().unwrap();
584
585        let sess_ours = sessions.path().join("sess-a");
586        let sess_theirs = sessions.path().join("sess-b");
587        std::fs::create_dir_all(&sess_ours).unwrap();
588        std::fs::create_dir_all(&sess_theirs).unwrap();
589
590        crate::session_meta::SessionMeta {
591            project_fingerprint: Some("ours".into()),
592            ..Default::default()
593        }
594        .save(&sess_ours)
595        .unwrap();
596        crate::session_meta::SessionMeta {
597            project_fingerprint: Some("theirs".into()),
598            ..Default::default()
599        }
600        .save(&sess_theirs)
601        .unwrap();
602
603        let tid1 = uuid::Uuid::now_v7();
604        let rid1 = uuid::Uuid::now_v7();
605        let events_ours = format!(
606            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}"}}}}
607{{"type":"flow_end","seq":2,"run_id":"{rid1}","flow_name":"demo","status":{{"kind":"ok"}},"ts":"2026-07-05T00:00:01Z"}}"#
608        );
609        std::fs::write(sess_ours.join("events.jsonl"), events_ours).unwrap();
610
611        let tid2 = uuid::Uuid::now_v7();
612        let events_theirs = format!(
613            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}"}}}}"#
614        );
615        std::fs::write(sess_theirs.join("events.jsonl"), events_theirs).unwrap();
616
617        let idx = AnchorIndex::open_project(index_dir.path()).unwrap();
618        let stats = idx
619            .rebuild_events_from_sessions(sessions.path(), "ours")
620            .unwrap();
621        assert_eq!(stats.rebuilt, 2);
622
623        let hits = idx.fts_search_project_events("needle", None, 10).unwrap();
624        assert_eq!(hits.len(), 1);
625        assert_eq!(hits[0].session_id, "sess-a");
626    }
627
628    #[test]
629    fn reopening_the_same_db_is_idempotent() {
630        let dir = tempfile::tempdir().unwrap();
631        let _one = AnchorIndex::open_project(dir.path()).unwrap();
632        let two = AnchorIndex::open_project(dir.path()).unwrap();
633        let tables = tables_in(&two);
634        assert!(tables.iter().any(|t| t == "confessions"));
635    }
636
637    #[test]
638    fn find_by_anchor_returns_subject_kinds_from_project_db() {
639        let dir = tempfile::tempdir().unwrap();
640        let idx = AnchorIndex::open_project(dir.path()).unwrap();
641        {
642            let conn = idx.conn();
643            conn.execute(
644                "INSERT INTO anchors (kind, ref, subject_kind, subject_id, created_at) VALUES (?, ?, ?, ?, ?)",
645                rusqlite::params!["flow_run", "run-xyz", "confession", "cid-1", "2026-07-05T00:00:00Z"],
646            )
647            .unwrap();
648            conn.execute(
649                "INSERT INTO anchors (kind, ref, subject_kind, subject_id, created_at) VALUES (?, ?, ?, ?, ?)",
650                rusqlite::params!["flow_run", "run-xyz", "spec_entry", "sid-1", "2026-07-05T00:00:00Z"],
651            )
652            .unwrap();
653            conn.execute(
654                "INSERT INTO anchors (kind, ref, subject_kind, subject_id, created_at) VALUES (?, ?, ?, ?, ?)",
655                rusqlite::params!["turn_id", "turn-99", "confession", "cid-2", "2026-07-05T00:00:00Z"],
656            )
657            .unwrap();
658        }
659        let hits = idx
660            .find_by_anchor(AnchorKind::FlowRunId, "run-xyz")
661            .unwrap();
662        assert_eq!(
663            hits,
664            vec![
665                ("confession".to_string(), "cid-1".to_string()),
666                ("spec_entry".to_string(), "sid-1".to_string()),
667            ]
668        );
669    }
670
671    #[test]
672    fn wal_journal_mode_is_set() {
673        let dir = tempfile::tempdir().unwrap();
674        let idx = AnchorIndex::open_project(dir.path()).unwrap();
675        let conn = idx.conn();
676        let mode: String = conn
677            .query_row("PRAGMA journal_mode", params![], |row| row.get(0))
678            .unwrap();
679        assert_eq!(mode.to_lowercase(), "wal");
680    }
681}