- 1
//! Per-ledger recall cache (docs/design/23-memory.md M1): keyed by the - 2
//! canonical ledger path and invalidated when mtime or length differs, so - 3
//! warm queries skip the JSONL rescan while appends stay visible on the - 4
//! very next search. Purely an accelerator — cold behavior is identical to - 5
//! a fresh scan. - 6
- 7
use std::collections::{HashMap, VecDeque}; - 8
use std::io::BufRead; - 9
use std::path::{Path, PathBuf}; - 10
use std::sync::{Arc, Mutex, MutexGuard, OnceLock}; - 11
use std::time::{Instant, SystemTime}; - 12
- 13
use chrono::{DateTime, Utc}; - 14
use vak_llm::Role; - 15
- 16
use crate::SearchError; - 17
use crate::search::normalize_impl; - 18
use crate::types::{Entry, EntryPayload}; - 19
- 20
/// Bounded work: at most the trailing N message lines of a ledger are - 21
/// scanned. Long sessions still recall their recent past. - 22
pub(crate) const MAX_SCAN_LINES: usize = 4000; - 23
const MAX_CACHED_LEDGERS: usize = 128; - 24
- 25
pub(crate) struct CachedMessage { - 26
pub entry_id: String, - 27
pub ts: DateTime<Utc>, - 28
pub role: &'static str, - 29
pub text: String, - 30
pub normalized: String, - 31
} - 32
- 33
#[derive(Clone, Copy, PartialEq, Eq)] - 34
struct Fingerprint { - 35
mtime: SystemTime, - 36
len: u64, - 37
} - 38
- 39
struct CacheSlot { - 40
fingerprint: Fingerprint, - 41
last_used: Instant, - 42
ledger: Arc<Vec<CachedMessage>>, - 43
} - 44
- 45
fn cache() -> &'static Mutex<HashMap<PathBuf, CacheSlot>> { - 46
static CACHE: OnceLock<Mutex<HashMap<PathBuf, CacheSlot>>> = OnceLock::new(); - 47
CACHE.get_or_init(|| Mutex::new(HashMap::new())) - 48
} - 49
- 50
fn lock(map: &Mutex<HashMap<PathBuf, CacheSlot>>) -> MutexGuard<'_, HashMap<PathBuf, CacheSlot>> { - 51
map.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) - 52
} - 53
- 54
fn evict_to_capacity(map: &mut HashMap<PathBuf, CacheSlot>) { - 55
while map.len() >= MAX_CACHED_LEDGERS { - 56
let oldest = map - 57
.iter() - 58
.min_by_key(|(_, slot)| slot.last_used) - 59
.map(|(k, _)| k.clone()); - 60
match oldest { - 61
Some(k) => { - 62
map.remove(&k); - 63
} - 64
None => break, - 65
} - 66
} - 67
} - 68
- 69
/// Tokenized representation of `path`'s trailing window, served from cache - 70
/// when the file's mtime and length are unchanged since the last scan. - 71
pub(crate) fn ledger(path: &Path) -> Result<Arc<Vec<CachedMessage>>, SearchError> { - 72
let key = path.canonicalize()?; - 73
let meta = std::fs::metadata(path)?; - 74
let fingerprint = Fingerprint { - 75
mtime: meta.modified()?, - 76
len: meta.len(), - 77
}; - 78
{ - 79
let mut map = lock(cache()); - 80
if let Some(slot) = map.get_mut(&key) - 81
&& slot.fingerprint == fingerprint - 82
{ - 83
slot.last_used = Instant::now(); - 84
return Ok(Arc::clone(&slot.ledger)); - 85
} - 86
} - 87
let scanned = Arc::new(scan(path)?); - 88
let mut map = lock(cache()); - 89
// A concurrent query may have cached a snapshot of this exact state - 90
// while we scanned; prefer it over re-inserting an equal copy. - 91
if let Some(slot) = map.get(&key) - 92
&& slot.fingerprint == fingerprint - 93
{ - 94
return Ok(Arc::clone(&slot.ledger)); - 95
} - 96
evict_to_capacity(&mut map); - 97
map.insert( - 98
key, - 99
CacheSlot { - 100
fingerprint, - 101
last_used: Instant::now(), - 102
ledger: Arc::clone(&scanned), - 103
}, - 104
); - 105
Ok(scanned) - 106
} - 107
- 108
fn scan(path: &Path) -> Result<Vec<CachedMessage>, SearchError> { - 109
let file = std::fs::File::open(path)?; - 110
let reader = std::io::BufReader::new(file); - 111
// Ring buffer keeps memory bounded on huge ledgers while preserving - 112
// "trailing lines" semantics. - 113
let mut ring: VecDeque<String> = VecDeque::with_capacity(MAX_SCAN_LINES); - 114
for line in reader.lines() { - 115
let line = line?; - 116
if line.trim().is_empty() { - 117
continue; - 118
} - 119
if ring.len() == MAX_SCAN_LINES { - 120
ring.pop_front(); - 121
} - 122
ring.push_back(line); - 123
} - 124
let mut messages = Vec::new(); - 125
for line in ring { - 126
let Ok(entry) = serde_json::from_str::<Entry>(&line) else { - 127
continue; - 128
}; - 129
let EntryPayload::Message(record) = entry.payload else { - 130
continue; - 131
}; - 132
// A runtime nudge is not something the user said or the model - 133
// answered; leaving it out keeps it from surfacing in session search. - 134
if record.control_kind().is_some() { - 135
continue; - 136
} - 137
let role = match record.message.role { - 138
Role::User => "user", - 139
Role::Assistant => "assistant", - 140
}; - 141
let text = record.message.text_content(); - 142
if text.trim().is_empty() { - 143
continue; - 144
} - 145
messages.push(CachedMessage { - 146
normalized: normalize_impl(&text), - 147
text, - 148
role, - 149
entry_id: entry.id, - 150
ts: entry.ts, - 151
}); - 152
} - 153
Ok(messages) - 154
} - 155
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.