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 = crate::event_writer::extract_ts(event);
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 = event.seq() as 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 text = match std::fs::read_to_string(&jsonl) {
219                Ok(t) => t,
220                Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
221                Err(e) => return Err(e).context(format!("read {}", jsonl.display())),
222            };
223            for line in text.lines() {
224                let trimmed = line.trim();
225                if trimmed.is_empty() {
226                    continue;
227                }
228                let Ok(value) = serde_json::from_str::<serde_json::Value>(trimmed) else {
229                    stats.skipped += 1;
230                    continue;
231                };
232                let seq = value.get("seq").and_then(|v| v.as_i64()).unwrap_or(0);
233                let ts = value.get("ts").and_then(|v| v.as_str()).unwrap_or("");
234                let kind = value
235                    .get("type")
236                    .and_then(|v| v.as_str())
237                    .unwrap_or("unknown");
238                let turn_id = value.get("turn_id").and_then(|v| v.as_str());
239                let flow_run_id = value
240                    .get("run_id")
241                    .or_else(|| value.get("flow_run_id"))
242                    .and_then(|v| v.as_str());
243                let text_content = value
244                    .get("message")
245                    .and_then(|m| m.get("parts"))
246                    .and_then(|p| p.as_array())
247                    .map(|parts| {
248                        parts
249                            .iter()
250                            .filter_map(|p| p.get("text").and_then(|t| t.as_str()))
251                            .collect::<Vec<_>>()
252                            .join("")
253                    })
254                    .unwrap_or_default();
255                self.insert_project_event_raw(ProjectEventInsert {
256                    session_id: &sid,
257                    seq,
258                    ts,
259                    kind,
260                    turn_id,
261                    flow_run_id,
262                    text_content: &text_content,
263                    payload_json: trimmed,
264                })?;
265                stats.rebuilt += 1;
266            }
267        }
268        Ok(stats)
269    }
270
271    pub fn find_by_anchor(&self, kind: AnchorKind, id: &str) -> Result<Vec<(String, String)>> {
272        let sql =
273            "SELECT subject_kind, subject_id FROM anchors WHERE kind = ? AND ref = ? ORDER BY id";
274        let conn = self.conn();
275        let mut stmt = conn.prepare(sql)?;
276        let rows = stmt.query_map(rusqlite::params![kind.anchor_tag(), id], |row| {
277            Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
278        })?;
279        let mut out = Vec::new();
280        for r in rows {
281            out.push(r?);
282        }
283        Ok(out)
284    }
285}
286
287#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
288pub struct RebuildStats {
289    pub rebuilt: usize,
290    pub skipped: usize,
291}
292
293#[derive(Debug, Clone, Copy, PartialEq, Eq)]
294pub enum AnchorKind {
295    TurnId,
296    FlowRunId,
297}
298
299impl AnchorKind {
300    fn events_column(self) -> &'static str {
301        match self {
302            AnchorKind::TurnId => "turn_id",
303            AnchorKind::FlowRunId => "flow_run_id",
304        }
305    }
306
307    fn anchor_tag(self) -> &'static str {
308        match self {
309            AnchorKind::TurnId => "turn",
310            AnchorKind::FlowRunId => "flow_run",
311        }
312    }
313}
314
315pub struct ProjectEventInsert<'a> {
316    pub session_id: &'a str,
317    pub seq: i64,
318    pub ts: &'a str,
319    pub kind: &'a str,
320    pub turn_id: Option<&'a str>,
321    pub flow_run_id: Option<&'a str>,
322    pub text_content: &'a str,
323    pub payload_json: &'a str,
324}
325
326#[derive(Debug, Clone, PartialEq, Eq)]
327pub struct ProjectEventRow {
328    pub session_id: String,
329    pub seq: u64,
330    pub ts: String,
331    pub kind: String,
332    pub turn_id: Option<String>,
333    pub flow_run_id: Option<String>,
334    pub payload: String,
335}
336
337fn project_event_row_from(row: &rusqlite::Row<'_>) -> rusqlite::Result<ProjectEventRow> {
338    Ok(ProjectEventRow {
339        session_id: row.get(0)?,
340        seq: row.get::<_, i64>(1)? as u64,
341        ts: row.get(2)?,
342        kind: row.get(3)?,
343        turn_id: row.get(4)?,
344        flow_run_id: row.get(5)?,
345        payload: row.get(6)?,
346    })
347}
348
349fn collect<T>(
350    iter: rusqlite::MappedRows<'_, impl FnMut(&rusqlite::Row<'_>) -> rusqlite::Result<T>>,
351) -> Result<Vec<T>> {
352    let mut out = Vec::new();
353    for r in iter {
354        out.push(r.map_err(|e| anyhow::anyhow!(e))?);
355    }
356    Ok(out)
357}
358
359const PROJECT_SCHEMA: &str = r#"
360CREATE TABLE IF NOT EXISTS events (
361    id          INTEGER PRIMARY KEY AUTOINCREMENT,
362    session_id  TEXT    NOT NULL,
363    seq         INTEGER NOT NULL,
364    ts          TEXT    NOT NULL,
365    kind        TEXT    NOT NULL,
366    turn_id     TEXT,
367    flow_run_id TEXT,
368    payload     TEXT    NOT NULL,
369    UNIQUE (session_id, seq)
370);
371CREATE INDEX IF NOT EXISTS events_session ON events(session_id);
372CREATE INDEX IF NOT EXISTS events_kind    ON events(kind);
373CREATE INDEX IF NOT EXISTS events_turn    ON events(turn_id);
374CREATE INDEX IF NOT EXISTS events_flow    ON events(flow_run_id);
375
376CREATE VIRTUAL TABLE IF NOT EXISTS events_fts USING fts5(
377    text_content,
378    tokenize='porter unicode61'
379);
380
381CREATE TABLE IF NOT EXISTS anchors (
382    id            INTEGER PRIMARY KEY AUTOINCREMENT,
383    kind          TEXT NOT NULL,
384    ref           TEXT NOT NULL,
385    subject_kind  TEXT NOT NULL,
386    subject_id    TEXT NOT NULL,
387    session_id    TEXT,
388    created_at    TEXT NOT NULL
389);
390CREATE INDEX IF NOT EXISTS anchors_lookup  ON anchors(kind, ref);
391CREATE INDEX IF NOT EXISTS anchors_subject ON anchors(subject_kind, subject_id);
392
393CREATE TABLE IF NOT EXISTS confessions (
394    id            TEXT PRIMARY KEY,
395    trigger       TEXT NOT NULL,
396    rule_violated TEXT NOT NULL,
397    what_i_did    TEXT NOT NULL,
398    why           TEXT NOT NULL,
399    mitigation    TEXT NOT NULL,
400    body          TEXT NOT NULL,
401    created_at    TEXT NOT NULL
402);
403CREATE VIRTUAL TABLE IF NOT EXISTS confessions_fts USING fts5(
404    trigger, rule_violated, what_i_did, why, mitigation, body,
405    tokenize='porter unicode61'
406);
407
408CREATE TABLE IF NOT EXISTS spec_entries (
409    id      TEXT PRIMARY KEY,
410    feature TEXT NOT NULL,
411    phase   TEXT NOT NULL,
412    content TEXT NOT NULL,
413    ts      TEXT NOT NULL
414);
415CREATE INDEX IF NOT EXISTS spec_entries_feature ON spec_entries(feature);
416CREATE VIRTUAL TABLE IF NOT EXISTS spec_entries_fts USING fts5(
417    content, tokenize='porter unicode61'
418);
419
420CREATE TABLE IF NOT EXISTS spec_deviations (
421    id      TEXT PRIMARY KEY,
422    feature TEXT NOT NULL,
423    section TEXT NOT NULL,
424    delta   TEXT NOT NULL,
425    reason  TEXT NOT NULL,
426    ts      TEXT NOT NULL
427);
428CREATE INDEX IF NOT EXISTS spec_deviations_feature ON spec_deviations(feature);
429CREATE VIRTUAL TABLE IF NOT EXISTS spec_deviations_fts USING fts5(
430    delta, reason, tokenize='porter unicode61'
431);
432"#;
433
434#[cfg(test)]
435mod tests {
436    use super::*;
437    use rusqlite::params;
438
439    fn tables_in(idx: &AnchorIndex) -> Vec<String> {
440        let conn = idx.conn();
441        let mut stmt = conn
442            .prepare("SELECT name FROM sqlite_master WHERE type IN ('table', 'view') ORDER BY name")
443            .unwrap();
444        stmt.query_map(params![], |row| row.get::<_, String>(0))
445            .unwrap()
446            .filter_map(|r| r.ok())
447            .collect()
448    }
449
450    #[test]
451    fn open_project_creates_all_four_tables() {
452        let dir = tempfile::tempdir().unwrap();
453        let idx = AnchorIndex::open_project(dir.path()).unwrap();
454        let tables = tables_in(&idx);
455        for expected in [
456            "anchors",
457            "confessions",
458            "confessions_fts",
459            "events",
460            "events_fts",
461            "spec_entries",
462            "spec_entries_fts",
463            "spec_deviations",
464            "spec_deviations_fts",
465        ] {
466            assert!(
467                tables.iter().any(|t| t == expected),
468                "missing {expected} in {tables:?}"
469            );
470        }
471    }
472
473    fn seed_project_event(
474        idx: &AnchorIndex,
475        sid: &str,
476        seq: i64,
477        kind: &str,
478        turn: Option<&str>,
479        flow: Option<&str>,
480        text: &str,
481    ) {
482        let payload = format!("{{\"seq\":{seq},\"session_id\":\"{sid}\"}}");
483        idx.insert_project_event_raw(ProjectEventInsert {
484            session_id: sid,
485            seq,
486            ts: "2026-07-05T00:00:00Z",
487            kind,
488            turn_id: turn,
489            flow_run_id: flow,
490            text_content: text,
491            payload_json: &payload,
492        })
493        .unwrap();
494    }
495
496    #[test]
497    fn fts_search_project_events_filters_by_session_when_requested() {
498        let dir = tempfile::tempdir().unwrap();
499        let idx = AnchorIndex::open_project(dir.path()).unwrap();
500        seed_project_event(
501            &idx,
502            "sess-a",
503            1,
504            "user_msg",
505            Some("t1"),
506            None,
507            "hello sqlite fts",
508        );
509        seed_project_event(
510            &idx,
511            "sess-b",
512            1,
513            "user_msg",
514            Some("t2"),
515            None,
516            "sqlite from other session",
517        );
518        seed_project_event(
519            &idx,
520            "sess-a",
521            2,
522            "assistant_msg",
523            Some("t1"),
524            None,
525            "no match here",
526        );
527
528        let scoped = idx
529            .fts_search_project_events("sqlite", Some("sess-a"), 10)
530            .unwrap();
531        assert_eq!(scoped.len(), 1);
532        assert_eq!(scoped[0].session_id, "sess-a");
533        assert_eq!(scoped[0].seq, 1);
534
535        let all = idx.fts_search_project_events("sqlite", None, 10).unwrap();
536        assert_eq!(all.len(), 2);
537    }
538
539    #[test]
540    fn insert_project_event_upsert_is_idempotent() {
541        let dir = tempfile::tempdir().unwrap();
542        let idx = AnchorIndex::open_project(dir.path()).unwrap();
543        seed_project_event(&idx, "sess-a", 42, "user_msg", None, None, "uniquetoken");
544        seed_project_event(&idx, "sess-a", 42, "user_msg", None, None, "uniquetoken");
545        let hits = idx
546            .fts_search_project_events("uniquetoken", None, 10)
547            .unwrap();
548        assert_eq!(hits.len(), 1, "same (session_id, seq) must upsert");
549    }
550
551    #[test]
552    fn find_project_events_around_scopes_to_session() {
553        let dir = tempfile::tempdir().unwrap();
554        let idx = AnchorIndex::open_project(dir.path()).unwrap();
555        seed_project_event(&idx, "sess-a", 1, "flow_start", None, Some("r"), "");
556        seed_project_event(&idx, "sess-a", 2, "user_msg", Some("t"), None, "hello");
557        seed_project_event(
558            &idx,
559            "sess-a",
560            3,
561            "assistant_msg",
562            Some("t"),
563            Some("r"),
564            "world",
565        );
566        seed_project_event(&idx, "sess-a", 4, "flow_end", None, Some("r"), "");
567        seed_project_event(&idx, "sess-b", 3, "user_msg", None, None, "unrelated");
568
569        let rows = idx.find_project_events_around("sess-a", 3, 1).unwrap();
570        let seqs: Vec<u64> = rows.iter().map(|r| r.seq).collect();
571        assert_eq!(seqs, vec![2, 3, 4]);
572        assert!(rows.iter().all(|r| r.session_id == "sess-a"));
573    }
574
575    #[test]
576    fn find_project_events_by_anchor_scopes_to_session() {
577        let dir = tempfile::tempdir().unwrap();
578        let idx = AnchorIndex::open_project(dir.path()).unwrap();
579        seed_project_event(&idx, "sess-a", 1, "flow_start", None, Some("run-x"), "");
580        seed_project_event(&idx, "sess-a", 2, "flow_end", None, Some("run-x"), "");
581        seed_project_event(&idx, "sess-b", 1, "flow_start", None, Some("run-x"), "");
582
583        let rows = idx
584            .find_project_events_by_anchor("sess-a", AnchorKind::FlowRunId, "run-x")
585            .unwrap();
586        assert_eq!(rows.len(), 2);
587        assert!(rows.iter().all(|r| r.session_id == "sess-a"));
588    }
589
590    #[test]
591    fn rebuild_events_from_sessions_filters_by_fingerprint() {
592        let sessions = tempfile::tempdir().unwrap();
593        let index_dir = tempfile::tempdir().unwrap();
594
595        let sess_ours = sessions.path().join("sess-a");
596        let sess_theirs = sessions.path().join("sess-b");
597        std::fs::create_dir_all(&sess_ours).unwrap();
598        std::fs::create_dir_all(&sess_theirs).unwrap();
599
600        crate::session_meta::SessionMeta {
601            project_fingerprint: Some("ours".into()),
602            ..Default::default()
603        }
604        .save(&sess_ours)
605        .unwrap();
606        crate::session_meta::SessionMeta {
607            project_fingerprint: Some("theirs".into()),
608            ..Default::default()
609        }
610        .save(&sess_theirs)
611        .unwrap();
612
613        std::fs::write(
614            sess_ours.join("events.jsonl"),
615            concat!(
616                r#"{"type":"user_msg","seq":1,"turn_id":"t1","ts":"2026-07-05T00:00:00Z","message":{"role":"user","parts":[{"type":"text","text":"needle in ours"}]}}"#,
617                "\n",
618                r#"{"type":"flow_end","seq":2,"run_id":"r1","ts":"2026-07-05T00:00:01Z"}"#,
619                "\n",
620            ),
621        )
622        .unwrap();
623        std::fs::write(
624            sess_theirs.join("events.jsonl"),
625            r#"{"type":"user_msg","seq":1,"turn_id":"t2","ts":"2026-07-05T00:00:00Z","message":{"role":"user","parts":[{"type":"text","text":"needle in theirs"}]}}"#,
626        )
627        .unwrap();
628
629        let idx = AnchorIndex::open_project(index_dir.path()).unwrap();
630        let stats = idx
631            .rebuild_events_from_sessions(sessions.path(), "ours")
632            .unwrap();
633        assert_eq!(stats.rebuilt, 2);
634
635        let hits = idx.fts_search_project_events("needle", None, 10).unwrap();
636        assert_eq!(hits.len(), 1);
637        assert_eq!(hits[0].session_id, "sess-a");
638    }
639
640    #[test]
641    fn reopening_the_same_db_is_idempotent() {
642        let dir = tempfile::tempdir().unwrap();
643        let _one = AnchorIndex::open_project(dir.path()).unwrap();
644        let two = AnchorIndex::open_project(dir.path()).unwrap();
645        let tables = tables_in(&two);
646        assert!(tables.iter().any(|t| t == "confessions"));
647    }
648
649    #[test]
650    fn find_by_anchor_returns_subject_kinds_from_project_db() {
651        let dir = tempfile::tempdir().unwrap();
652        let idx = AnchorIndex::open_project(dir.path()).unwrap();
653        {
654            let conn = idx.conn();
655            conn.execute(
656                "INSERT INTO anchors (kind, ref, subject_kind, subject_id, created_at) VALUES (?, ?, ?, ?, ?)",
657                rusqlite::params!["flow_run", "run-xyz", "confession", "cid-1", "2026-07-05T00:00:00Z"],
658            )
659            .unwrap();
660            conn.execute(
661                "INSERT INTO anchors (kind, ref, subject_kind, subject_id, created_at) VALUES (?, ?, ?, ?, ?)",
662                rusqlite::params!["flow_run", "run-xyz", "spec_entry", "sid-1", "2026-07-05T00:00:00Z"],
663            )
664            .unwrap();
665            conn.execute(
666                "INSERT INTO anchors (kind, ref, subject_kind, subject_id, created_at) VALUES (?, ?, ?, ?, ?)",
667                rusqlite::params!["turn_id", "turn-99", "confession", "cid-2", "2026-07-05T00:00:00Z"],
668            )
669            .unwrap();
670        }
671        let hits = idx
672            .find_by_anchor(AnchorKind::FlowRunId, "run-xyz")
673            .unwrap();
674        assert_eq!(
675            hits,
676            vec![
677                ("confession".to_string(), "cid-1".to_string()),
678                ("spec_entry".to_string(), "sid-1".to_string()),
679            ]
680        );
681    }
682
683    #[test]
684    fn wal_journal_mode_is_set() {
685        let dir = tempfile::tempdir().unwrap();
686        let idx = AnchorIndex::open_project(dir.path()).unwrap();
687        let conn = idx.conn();
688        let mode: String = conn
689            .query_row("PRAGMA journal_mode", params![], |row| row.get(0))
690            .unwrap();
691        assert_eq!(mode.to_lowercase(), "wal");
692    }
693}