- 1
//! Bulk import and incremental file-level indexing. - 2
- 3
use std::path::Path; - 4
- 5
use walkdir::WalkDir; - 6
- 7
use crate::{ImportStats, RebuildStats, Store, StoreError}; - 8
use vak_session::types::Entry; - 9
- 10
impl Store { - 11
/// Scan every `<home>/sessions/<hash>/*.jsonl` file and index entries. - 12
pub(crate) fn import_all( - 13
&self, - 14
conn: &rusqlite::Connection, - 15
sessions_home: &Path, - 16
) -> Result<RebuildStats, StoreError> { - 17
let mut roots = Vec::new(); - 18
let direct = sessions_home.join("sessions"); - 19
if direct.exists() { - 20
roots.push((sessions_home.to_path_buf(), direct)); - 21
} - 22
if let Ok(agents) = std::fs::read_dir(sessions_home.join("agents")) { - 23
for agent in agents.flatten() { - 24
let s = agent.path().join("sessions"); - 25
if s.exists() { - 26
roots.push((agent.path(), s)); - 27
} - 28
} - 29
} - 30
if let Some(siblings) = sessions_home - 31
.parent() - 32
.filter(|p| p.file_name().and_then(|s| s.to_str()) == Some("agents")) - 33
.and_then(|p| std::fs::read_dir(p).ok()) - 34
{ - 35
for sibling in siblings.flatten() { - 36
let s = sibling.path().join("sessions"); - 37
if s.exists() && !roots.iter().any(|(_, r)| r == &s) { - 38
roots.push((sibling.path(), s)); - 39
} - 40
} - 41
} - 42
if roots.is_empty() { - 43
return Ok(RebuildStats { - 44
files_scanned: 0, - 45
entries_indexed: 0, - 46
fts_rows: 0, - 47
}); - 48
} - 49
let mut files_scanned = 0usize; - 50
let mut entries_indexed = 0usize; - 51
let mut fts_rows = 0usize; - 52
- 53
for (home, root) in roots { - 54
for entry in WalkDir::new(&root) - 55
.min_depth(2) - 56
.max_depth(2) - 57
.follow_links(false) - 58
.into_iter() - 59
.filter_entry(|e| e.file_type().is_file()) - 60
{ - 61
let entry = entry.map_err(|e| std::io::Error::other(e.to_string()))?; - 62
let path = entry.path(); - 63
if path.extension().and_then(|e| e.to_str()) != Some("jsonl") { - 64
continue; - 65
} - 66
let stats = self.import_file(conn, &home, path)?; - 67
entries_indexed += stats.entries_indexed; - 68
fts_rows += stats.fts_rows; - 69
files_scanned += 1; - 70
} - 71
} - 72
- 73
Ok(RebuildStats { - 74
files_scanned, - 75
entries_indexed, - 76
fts_rows, - 77
}) - 78
} - 79
- 80
/// Parse a single JSONL file and insert all entries into the index. - 81
/// Already-indexed entry IDs are skipped. - 82
pub(crate) fn import_file( - 83
&self, - 84
conn: &rusqlite::Connection, - 85
_sessions_home: &Path, - 86
jsonl_path: &Path, - 87
) -> Result<ImportStats, StoreError> { - 88
let session_id = jsonl_path - 89
.file_stem() - 90
.and_then(|s| s.to_str()) - 91
.unwrap_or_default() - 92
.to_string(); - 93
- 94
// Derive project hash from the parent directory name. - 95
// Directory structure: <home>/sessions/<hash>/<session>.jsonl - 96
let project_hash = jsonl_path - 97
.parent() - 98
.and_then(|p| p.file_name()) - 99
.and_then(|n| n.to_str()) - 100
.unwrap_or_default() - 101
.to_string(); - 102
- 103
let content = std::fs::read_to_string(jsonl_path)?; - 104
let mut entries_indexed = 0usize; - 105
let mut fts_rows = 0usize; - 106
let mut skipped = 0usize; - 107
- 108
// Prepare check and insert statements once. - 109
let mut check_stmt = conn.prepare("SELECT 1 FROM entries WHERE entry_id = ?1 LIMIT 1")?; - 110
let mut insert_entry = conn.prepare( - 111
"INSERT OR IGNORE INTO entries - 112
(entry_id, session_id, project_hash, parent_id, ts, kind, role, - 113
provider, model, tool_name, content_text, is_error) - 114
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)", - 115
)?; - 116
let mut insert_fts = conn.prepare( - 117
"INSERT INTO entries_fts(entry_id, session_id, project_hash, ts, kind, role, content) - 118
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)", - 119
)?; - 120
- 121
for (i, line) in content.lines().enumerate() { - 122
let line = line.trim(); - 123
if line.is_empty() { - 124
continue; - 125
} - 126
let entry: Entry = match serde_json::from_str(line) { - 127
Ok(e) => e, - 128
Err(_) => continue, - 129
}; - 130
- 131
// Skip already indexed. - 132
let already: bool = check_stmt - 133
.query_row(rusqlite::params![entry.id], |row| row.get::<_, i32>(0)) - 134
.is_ok(); - 135
if already { - 136
skipped += 1; - 137
continue; - 138
} - 139
- 140
if let Some(meta) = self.extract_meta_with_hash(&session_id, &project_hash, &entry) { - 141
let has_text = !meta.content_text.trim().is_empty(); - 142
- 143
insert_entry.execute(rusqlite::params![ - 144
meta.entry_id, - 145
meta.session_id, - 146
meta.project_hash, - 147
meta.parent_id, - 148
meta.ts, - 149
meta.kind.as_str(), - 150
meta.role, - 151
meta.provider, - 152
meta.model, - 153
meta.tool_name, - 154
meta.content_text, - 155
meta.is_error as i32, - 156
])?; - 157
- 158
if has_text { - 159
insert_fts.execute(rusqlite::params![ - 160
meta.entry_id, - 161
meta.session_id, - 162
meta.project_hash, - 163
meta.ts, - 164
meta.kind.as_str(), - 165
meta.role.as_deref().unwrap_or(""), - 166
meta.content_text, - 167
])?; - 168
fts_rows += 1; - 169
} - 170
- 171
entries_indexed += 1; - 172
} - 173
- 174
if i % 1000 == 999 { - 175
conn.execute("PRAGMA optimize", [])?; - 176
} - 177
} - 178
- 179
Ok(ImportStats { - 180
entries_indexed, - 181
fts_rows, - 182
skipped, - 183
}) - 184
} - 185
- 186
/// Like `extract_meta` but attaches project_hash. - 187
fn extract_meta_with_hash( - 188
&self, - 189
session_id: &str, - 190
project_hash: &str, - 191
entry: &Entry, - 192
) -> Option<crate::IndexedEntry> { - 193
let mut meta = self.extract_meta(session_id, entry)?; - 194
meta.project_hash = project_hash.to_string(); - 195
Some(meta) - 196
} - 197
} - 198
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.