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}