- 1
//! FTS5 full-text search and structured query API. - 2
- 3
use crate::{IndexedEntry, Store, StoreError}; - 4
- 5
#[derive(Debug, Clone, serde::Serialize)] - 6
pub struct SearchResult { - 7
pub entries: Vec<SearchHit>, - 8
pub total: usize, - 9
} - 10
- 11
#[derive(Debug, Clone, serde::Serialize)] - 12
pub struct SearchHit { - 13
pub entry_id: String, - 14
pub session_id: String, - 15
pub project_hash: String, - 16
pub ts: String, - 17
pub kind: String, - 18
pub role: Option<String>, - 19
pub provider: Option<String>, - 20
pub model: Option<String>, - 21
pub tool_name: Option<String>, - 22
pub score: f64, - 23
pub snippet: String, - 24
} - 25
- 26
#[derive(Debug, Clone, Default)] - 27
pub struct SearchFilter { - 28
pub session_id: Option<String>, - 29
pub project_hash: Option<String>, - 30
pub kind: Option<String>, - 31
pub role: Option<String>, - 32
pub since: Option<String>, - 33
pub until: Option<String>, - 34
/// Sessions no result may come from: the one searching, and the trash. - 35
pub excluded_sessions: Vec<String>, - 36
} - 37
- 38
fn exclude_sessions( - 39
column: &str, - 40
excluded: &[String], - 41
conditions: &mut Vec<String>, - 42
params: &mut Vec<Box<dyn rusqlite::types::ToSql>>, - 43
) { - 44
if excluded.is_empty() { - 45
return; - 46
} - 47
let slots = vec!["?"; excluded.len()].join(", "); - 48
conditions.push(format!("{column} NOT IN ({slots})")); - 49
params.extend( - 50
excluded - 51
.iter() - 52
.map(|id| Box::new(id.clone()) as Box<dyn rusqlite::types::ToSql>), - 53
); - 54
} - 55
- 56
impl Store { - 57
/// Full-text search using FTS5 BM25 ranking. The query is passed - 58
/// directly to FTS5 (supports `AND`, `OR`, `"exact phrase"`, `NEAR`, - 59
/// prefix `*`). - 60
pub fn search( - 61
&self, - 62
query: &str, - 63
limit: usize, - 64
filter: &SearchFilter, - 65
) -> Result<SearchResult, StoreError> { - 66
let conn = self.conn(); - 67
let limit = limit.clamp(1, 100); - 68
- 69
// Build the FTS5 query with optional structured filters. - 70
let mut conditions = vec!["entries_fts MATCH ?1".to_string()]; - 71
let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new(); - 72
params.push(Box::new(query.to_string())); - 73
- 74
if let Some(ref sid) = filter.session_id { - 75
conditions.push("entries_fts.session_id = ?".to_string()); - 76
params.push(Box::new(sid.clone())); - 77
} - 78
if let Some(ref hash) = filter.project_hash { - 79
conditions.push("entries_fts.project_hash = ?".to_string()); - 80
params.push(Box::new(hash.clone())); - 81
} - 82
if let Some(ref kind) = filter.kind { - 83
conditions.push("entries_fts.kind = ?".to_string()); - 84
params.push(Box::new(kind.clone())); - 85
} - 86
if let Some(ref role) = filter.role { - 87
conditions.push("entries_fts.role = ?".to_string()); - 88
params.push(Box::new(role.clone())); - 89
} - 90
if let Some(ref since) = filter.since { - 91
conditions.push("entries_fts.ts >= ?".to_string()); - 92
params.push(Box::new(since.clone())); - 93
} - 94
if let Some(ref until) = filter.until { - 95
conditions.push("entries_fts.ts <= ?".to_string()); - 96
params.push(Box::new(until.clone())); - 97
} - 98
exclude_sessions( - 99
"entries_fts.session_id", - 100
&filter.excluded_sessions, - 101
&mut conditions, - 102
&mut params, - 103
); - 104
- 105
let where_clause = conditions.join(" AND "); - 106
let sql = format!( - 107
"SELECT entries_fts.entry_id, entries_fts.session_id, - 108
entries_fts.project_hash, entries_fts.ts, - 109
entries_fts.kind, entries_fts.role, - 110
e.provider, e.model, e.tool_name, - 111
bm25(entries_fts) AS score, - 112
snippet(entries_fts, 6, '<b>', '</b>', '…', 32) AS snippet - 113
FROM entries_fts - 114
LEFT JOIN entries e ON e.entry_id = entries_fts.entry_id - 115
WHERE {where_clause} - 116
ORDER BY bm25(entries_fts) - 117
LIMIT {limit}" - 118
); - 119
- 120
let param_refs: Vec<&dyn rusqlite::types::ToSql> = - 121
params.iter().map(|p| p.as_ref()).collect(); - 122
- 123
let mut stmt = conn.prepare(&sql)?; - 124
let rows = stmt.query_map(param_refs.as_slice(), |row| { - 125
Ok(SearchHit { - 126
entry_id: row.get(0)?, - 127
session_id: row.get(1)?, - 128
project_hash: row.get(2)?, - 129
ts: row.get(3)?, - 130
kind: row.get(4)?, - 131
role: row.get(5)?, - 132
provider: row.get(6)?, - 133
model: row.get(7)?, - 134
tool_name: row.get(8)?, - 135
score: row.get(9)?, - 136
snippet: row.get(10)?, - 137
}) - 138
})?; - 139
- 140
let mut entries = Vec::new(); - 141
for row in rows { - 142
entries.push(row?); - 143
} - 144
let total = entries.len(); - 145
Ok(SearchResult { entries, total }) - 146
} - 147
- 148
/// Structured query without FTS — filter by metadata only. - 149
pub fn query( - 150
&self, - 151
filter: &SearchFilter, - 152
limit: usize, - 153
) -> Result<Vec<IndexedEntry>, StoreError> { - 154
self.query_page(filter, limit, 0, false) - 155
.map(|(entries, _)| entries) - 156
} - 157
- 158
/// Structured query with pagination, ordering, and total matching count. - 159
pub fn query_page( - 160
&self, - 161
filter: &SearchFilter, - 162
limit: usize, - 163
offset: usize, - 164
ascending: bool, - 165
) -> Result<(Vec<IndexedEntry>, usize), StoreError> { - 166
let conn = self.conn(); - 167
let limit = limit.clamp(1, 500); - 168
- 169
let mut conditions = Vec::new(); - 170
let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new(); - 171
- 172
if let Some(ref sid) = filter.session_id { - 173
conditions.push("session_id = ?".to_string()); - 174
params.push(Box::new(sid.clone())); - 175
} - 176
if let Some(ref hash) = filter.project_hash { - 177
conditions.push("project_hash = ?".to_string()); - 178
params.push(Box::new(hash.clone())); - 179
} - 180
if let Some(ref kind) = filter.kind { - 181
conditions.push("kind = ?".to_string()); - 182
params.push(Box::new(kind.clone())); - 183
} - 184
if let Some(ref role) = filter.role { - 185
conditions.push("role = ?".to_string()); - 186
params.push(Box::new(role.clone())); - 187
} - 188
if let Some(ref since) = filter.since { - 189
conditions.push("ts >= ?".to_string()); - 190
params.push(Box::new(since.clone())); - 191
} - 192
if let Some(ref until) = filter.until { - 193
conditions.push("ts <= ?".to_string()); - 194
params.push(Box::new(until.clone())); - 195
} - 196
exclude_sessions( - 197
"session_id", - 198
&filter.excluded_sessions, - 199
&mut conditions, - 200
&mut params, - 201
); - 202
- 203
let where_clause = if conditions.is_empty() { - 204
"1=1".to_string() - 205
} else { - 206
conditions.join(" AND ") - 207
}; - 208
- 209
let count_sql = format!("SELECT COUNT(*) FROM entries WHERE {where_clause}"); - 210
let param_refs: Vec<&dyn rusqlite::types::ToSql> = - 211
params.iter().map(|p| p.as_ref()).collect(); - 212
- 213
let total: usize = conn.query_row(&count_sql, param_refs.as_slice(), |row| row.get(0))?; - 214
- 215
let order_dir = if ascending { "ASC" } else { "DESC" }; - 216
let sql = format!( - 217
"SELECT entry_id, session_id, project_hash, parent_id, ts, - 218
kind, role, provider, model, tool_name, - 219
content_text, is_error - 220
FROM entries - 221
WHERE {where_clause} - 222
ORDER BY ts {order_dir}, entry_id {order_dir} - 223
LIMIT {limit} OFFSET {offset}" - 224
); - 225
- 226
let mut stmt = conn.prepare(&sql)?; - 227
let rows = stmt.query_map(param_refs.as_slice(), |row| { - 228
let kind_str: String = row.get(5)?; - 229
Ok(IndexedEntry { - 230
entry_id: row.get(0)?, - 231
session_id: row.get(1)?, - 232
project_hash: row.get(2)?, - 233
parent_id: row.get(3)?, - 234
ts: row.get(4)?, - 235
kind: crate::EntryKind::parse_str(&kind_str).unwrap_or(crate::EntryKind::Message), - 236
role: row.get(6)?, - 237
provider: row.get(7)?, - 238
model: row.get(8)?, - 239
tool_name: row.get(9)?, - 240
content_text: row.get(10)?, - 241
is_error: row.get::<_, i32>(11)? != 0, - 242
}) - 243
})?; - 244
- 245
let mut entries = Vec::new(); - 246
for row in rows { - 247
entries.push(row?); - 248
} - 249
Ok((entries, total)) - 250
} - 251
- 252
/// List all distinct session IDs with entry counts. - 253
pub fn list_sessions(&self) -> Result<Vec<SessionInfo>, StoreError> { - 254
let conn = self.conn(); - 255
let mut stmt = conn.prepare( - 256
"SELECT session_id, project_hash, - 257
COUNT(*) as count, - 258
MIN(ts) as first_ts, - 259
MAX(ts) as last_ts - 260
FROM entries - 261
GROUP BY session_id - 262
ORDER BY last_ts DESC", - 263
)?; - 264
- 265
let rows = stmt.query_map([], |row| { - 266
Ok(SessionInfo { - 267
session_id: row.get(0)?, - 268
project_hash: row.get(1)?, - 269
entry_count: row.get(2)?, - 270
first_ts: row.get(3)?, - 271
last_ts: row.get(4)?, - 272
}) - 273
})?; - 274
- 275
Ok(rows.collect::<Result<Vec<_>, _>>()?) - 276
} - 277
} - 278
- 279
#[derive(Debug, Clone, serde::Serialize)] - 280
pub struct SessionInfo { - 281
pub session_id: String, - 282
pub project_hash: String, - 283
pub entry_count: usize, - 284
pub first_ts: String, - 285
pub last_ts: String, - 286
} - 287
- 288
#[cfg(test)] - 289
mod tests { - 290
#![allow(clippy::unwrap_used, clippy::expect_used)] - 291
- 292
use super::*; - 293
use crate::Store; - 294
use std::io::Write; - 295
use std::path::Path; - 296
use vak_llm::types::{ContentBlock, Message}; - 297
use vak_session::log::SessionPath; - 298
use vak_session::types::{Entry, EntryPayload, FrozenContract, MessageRecord, SessionHeader}; - 299
- 300
fn test_header(id: &str) -> SessionHeader { - 301
SessionHeader { - 302
agent: None, - 303
session_id: id.to_string(), - 304
created_at: chrono::Utc::now(), - 305
cwd: std::env::current_dir().unwrap(), - 306
parent_session_id: None, - 307
contract_id: None, - 308
work_item_id: None, - 309
conversation: None, - 310
contract: FrozenContract { - 311
app_version: "test".into(), - 312
provider: "openai".into(), - 313
model: "gpt-4o".into(), - 314
route_ladder: vec![], - 315
route_objective: String::new(), - 316
route_annotations: vec![], - 317
system_prompt: String::new(), - 318
permission_mode: "workspace-write".into(), - 319
capabilities: Vec::new(), - 320
prompt_layers: Vec::new(), - 321
}, - 322
} - 323
} - 324
- 325
fn user_msg(text: &str) -> MessageRecord { - 326
MessageRecord { - 327
message: Message { - 328
role: vak_llm::Role::User, - 329
content: vec![ContentBlock::text(text)], - 330
}, - 331
meta: None, - 332
} - 333
} - 334
- 335
fn assistant_msg(text: &str) -> MessageRecord { - 336
MessageRecord { - 337
message: Message { - 338
role: vak_llm::Role::Assistant, - 339
content: vec![ContentBlock::text(text)], - 340
}, - 341
meta: None, - 342
} - 343
} - 344
- 345
fn write_session(home: &Path, cwd: &Path, id: &str, msgs: &[MessageRecord]) { - 346
let path = SessionPath::new_session_file(home, cwd, id); - 347
std::fs::create_dir_all(path.parent().unwrap()).unwrap(); - 348
let file = std::fs::File::create(&path).unwrap(); - 349
let mut w = std::io::BufWriter::new(file); - 350
let header = Entry::new(None, EntryPayload::Header(test_header(id))); - 351
serde_json::to_writer(&mut w, &header).unwrap(); - 352
w.write_all(b"\n").unwrap(); - 353
let mut parent = Some(header.id.clone()); - 354
for m in msgs { - 355
let entry = Entry { - 356
prev_hash: None, - 357
id: uuid::Uuid::now_v7().to_string(), - 358
parent_id: parent.clone(), - 359
ts: chrono::Utc::now(), - 360
payload: EntryPayload::Message(m.clone()), - 361
}; - 362
parent = Some(entry.id.clone()); - 363
serde_json::to_writer(&mut w, &entry).unwrap(); - 364
w.write_all(b"\n").unwrap(); - 365
} - 366
w.flush().unwrap(); - 367
} - 368
- 369
#[test] - 370
fn fts_search_returns_ranked_results() { - 371
let dir = tempfile::tempdir().unwrap(); - 372
let home = dir.path(); - 373
let cwd = dir.path(); - 374
write_session( - 375
home, - 376
cwd, - 377
"s1", - 378
&[ - 379
user_msg("the deploy pipeline handles rollbacks gracefully"), - 380
assistant_msg("confirmed — rollback windows pause before deploy"), - 381
], - 382
); - 383
write_session( - 384
home, - 385
cwd, - 386
"s2", - 387
&[user_msg("pizza toppings are irrelevant to this query")], - 388
); - 389
- 390
let store = Store::open(home).unwrap(); - 391
store.rebuild(home).unwrap(); - 392
- 393
let result = store - 394
.search("deploy rollback", 10, &SearchFilter::default()) - 395
.unwrap(); - 396
assert!(result.total >= 2, "both deploy messages should match"); - 397
// All hits should be from s1. - 398
assert!(result.entries.iter().all(|h| h.session_id == "s1")); - 399
// First hit should have a snippet with highlight markers. - 400
assert!(result.entries[0].snippet.contains("<b>")); - 401
} - 402
- 403
#[test] - 404
fn fts_search_respects_filters() { - 405
let dir = tempfile::tempdir().unwrap(); - 406
let home = dir.path(); - 407
let cwd = dir.path(); - 408
write_session( - 409
home, - 410
cwd, - 411
"s1", - 412
&[user_msg("rust compiler optimization flags")], - 413
); - 414
write_session( - 415
home, - 416
cwd, - 417
"s2", - 418
&[assistant_msg("rust compiler error messages are cryptic")], - 419
); - 420
- 421
let store = Store::open(home).unwrap(); - 422
store.rebuild(home).unwrap(); - 423
- 424
// Filter by role=user. - 425
let result = store - 426
.search( - 427
"rust compiler", - 428
10, - 429
&SearchFilter { - 430
role: Some("user".into()), - 431
..Default::default() - 432
}, - 433
) - 434
.unwrap(); - 435
assert!( - 436
result - 437
.entries - 438
.iter() - 439
.all(|h| h.role.as_deref() == Some("user")) - 440
); - 441
- 442
// Exclude session s1. - 443
let result = store - 444
.search( - 445
"rust compiler", - 446
10, - 447
&SearchFilter { - 448
excluded_sessions: vec!["s1".into()], - 449
..Default::default() - 450
}, - 451
) - 452
.unwrap(); - 453
assert!(result.entries.iter().all(|h| h.session_id == "s2")); - 454
} - 455
- 456
#[test] - 457
fn structured_query_by_kind() { - 458
let dir = tempfile::tempdir().unwrap(); - 459
let home = dir.path(); - 460
let cwd = dir.path(); - 461
write_session(home, cwd, "s1", &[user_msg("hello world")]); - 462
- 463
let store = Store::open(home).unwrap(); - 464
store.rebuild(home).unwrap(); - 465
- 466
let headers = store - 467
.query( - 468
&SearchFilter { - 469
kind: Some("header".into()), - 470
..Default::default() - 471
}, - 472
10, - 473
) - 474
.unwrap(); - 475
assert_eq!(headers.len(), 1); - 476
assert_eq!(headers[0].kind, crate::EntryKind::Header); - 477
- 478
let messages = store - 479
.query( - 480
&SearchFilter { - 481
kind: Some("message".into()), - 482
..Default::default() - 483
}, - 484
10, - 485
) - 486
.unwrap(); - 487
assert_eq!(messages.len(), 1); - 488
} - 489
- 490
#[test] - 491
fn list_sessions_groups_correctly() { - 492
let dir = tempfile::tempdir().unwrap(); - 493
let home = dir.path(); - 494
let cwd = dir.path(); - 495
write_session( - 496
home, - 497
cwd, - 498
"sess-aaa", - 499
&[user_msg("first"), assistant_msg("second")], - 500
); - 501
write_session(home, cwd, "sess-bbb", &[user_msg("third")]); - 502
write_session(home, cwd, "sess-empty", &[]); - 503
- 504
let store = Store::open(home).unwrap(); - 505
store.rebuild(home).unwrap(); - 506
- 507
let sessions = store.list_sessions().unwrap(); - 508
assert_eq!(sessions.len(), 3); - 509
let aaa = sessions - 510
.iter() - 511
.find(|s| s.session_id == "sess-aaa") - 512
.unwrap(); - 513
assert_eq!(aaa.entry_count, 3); // 1 header + 2 messages - 514
let bbb = sessions - 515
.iter() - 516
.find(|s| s.session_id == "sess-bbb") - 517
.unwrap(); - 518
assert_eq!(bbb.entry_count, 2); // 1 header + 1 message - 519
let empty = sessions - 520
.iter() - 521
.find(|s| s.session_id == "sess-empty") - 522
.unwrap(); - 523
assert_eq!(empty.entry_count, 1); // 1 header - 524
} - 525
- 526
#[test] - 527
fn query_page_paginates_and_orders_correctly() { - 528
let dir = tempfile::tempdir().unwrap(); - 529
let home = dir.path(); - 530
let cwd = dir.path(); - 531
write_session( - 532
home, - 533
cwd, - 534
"sess-p", - 535
&[user_msg("msg 1"), assistant_msg("msg 2"), user_msg("msg 3")], - 536
); - 537
- 538
let store = Store::open(home).unwrap(); - 539
store.rebuild(home).unwrap(); - 540
- 541
let filter = SearchFilter { - 542
session_id: Some("sess-p".into()), - 543
..Default::default() - 544
}; - 545
- 546
// Total entries = 1 header + 3 messages = 4 - 547
let (page1, total) = store.query_page(&filter, 2, 0, true).unwrap(); - 548
assert_eq!(total, 4); - 549
assert_eq!(page1.len(), 2); - 550
assert_eq!(page1[0].kind, crate::EntryKind::Header); - 551
- 552
let (page2, total2) = store.query_page(&filter, 2, 2, true).unwrap(); - 553
assert_eq!(total2, 4); - 554
assert_eq!(page2.len(), 2); - 555
- 556
let (page3, total3) = store.query_page(&filter, 2, 4, true).unwrap(); - 557
assert_eq!(total3, 4); - 558
assert_eq!(page3.len(), 0); - 559
} - 560
- 561
#[test] - 562
fn import_session_is_idempotent() { - 563
let dir = tempfile::tempdir().unwrap(); - 564
let home = dir.path(); - 565
let cwd = dir.path(); - 566
write_session( - 567
home, - 568
cwd, - 569
"s-idem", - 570
&[user_msg("unique term idempotency check")], - 571
); - 572
- 573
let store = Store::open(home).unwrap(); - 574
let path = SessionPath::new_session_file(home, cwd, "s-idem"); - 575
let s1 = store.import_session(home, &path).unwrap(); - 576
assert_eq!(s1.entries_indexed, 2); - 577
let s2 = store.import_session(home, &path).unwrap(); - 578
assert_eq!(s2.skipped, 2); - 579
assert_eq!(s2.entries_indexed, 0); - 580
- 581
// Search still works after double-import. - 582
let result = store - 583
.search("idempotency", 5, &SearchFilter::default()) - 584
.unwrap(); - 585
assert_eq!(result.total, 1); - 586
} - 587
} - 588
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.