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}