1use std::collections::HashMap;
26use std::fs::File;
27use std::io::Write;
28use std::mem::size_of;
29use std::path::{Path, PathBuf};
30
31use anyhow::{Context, Result, ensure};
32use bytemuck::{Pod, Zeroable};
33use memmap2::Mmap;
34
35const MAGIC: [u8; 8] = *b"PONDRMM5";
36const BLOCK_ROWS: usize = 256;
37const ZSTD_LEVEL: i32 = 3;
38
39const ROW_HEADER_LEN: usize = 8 + 7 * 4;
44
45type BlockCache = Option<(usize, Vec<u8>)>;
49
50#[repr(C)]
51#[derive(Clone, Copy, Pod, Zeroable)]
52struct Header {
53 magic: [u8; 8],
54 version: u64,
55 count: u64,
56 session_count: u64,
57 project_count: u64,
58 agent_count: u64,
59 role_count: u64,
60 block_count: u64,
61 blob_offset: u64,
62}
63
64#[repr(C)]
65#[derive(Clone, Copy, Pod, Zeroable)]
66struct Record {
67 row_id: u64,
68 blob_off: u64,
69}
70
71#[repr(C)]
72#[derive(Clone, Copy, Pod, Zeroable)]
73struct SessionEntry {
74 sid_off: u64,
75 max_ts_micros: i64,
76 sid_len: u32,
77 count: u32,
78}
79
80#[repr(C)]
81#[derive(Clone, Copy, Pod, Zeroable)]
82struct DictEntry {
83 off: u64,
84 len: u32,
85 _pad: u32,
86}
87
88#[repr(C)]
89#[derive(Clone, Copy, Pod, Zeroable)]
90struct BlockEntry {
91 comp_off: u64,
92 comp_len: u32,
93 decomp_len: u32,
94}
95
96#[derive(Clone)]
98pub struct RowMetaEntry {
99 pub row_id: u64,
100 pub session_id: String,
101 pub message_id: String,
102 pub role: String,
103 pub project: String,
104 pub source_agent: String,
105 pub timestamp_micros: i64,
106 pub search_text: String,
107}
108
109pub struct RowMeta<'a> {
112 pub session_id: &'a str,
113 pub message_id: &'a str,
114 pub role: &'a str,
115 pub project: &'a str,
116 pub source_agent: &'a str,
117 pub timestamp_micros: i64,
118 pub search_text: String,
119}
120
121pub struct RowMetaMap {
124 mmap: Mmap,
125 version: u64,
126 count: usize,
127 session_count: usize,
128 project_count: usize,
129 agent_count: usize,
130 role_count: usize,
131 block_count: usize,
132 sessions_off: usize,
133 projects_off: usize,
134 agents_off: usize,
135 roles_off: usize,
136 blocks_off: usize,
137 blob_offset: usize,
138}
139
140impl std::fmt::Debug for RowMetaMap {
141 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
142 formatter
143 .debug_struct("RowMetaMap")
144 .field("version", &self.version)
145 .field("count", &self.count)
146 .field("session_count", &self.session_count)
147 .finish_non_exhaustive()
148 }
149}
150
151impl RowMetaMap {
152 pub fn path_for(cache_dir: &Path, store_key: &str, version: u64) -> PathBuf {
154 cache_dir.join(format!("rowmetamap-{store_key}-v{version}.rmm"))
155 }
156
157 pub fn delta_path(cache_dir: &Path, store_key: &str, version: u64) -> PathBuf {
160 cache_dir.join(format!("rowmetamap-{store_key}-d{version}.rmm"))
161 }
162
163 pub fn build(path: &Path, version: u64, mut entries: Vec<RowMetaEntry>) -> Result<()> {
164 entries.sort_unstable_by_key(|entry| entry.row_id);
165
166 let mut session_agg: HashMap<&str, (u32, i64)> = HashMap::new();
171 for entry in &entries {
172 let agg = session_agg
173 .entry(entry.session_id.as_str())
174 .or_insert((0, i64::MIN));
175 agg.0 += 1;
176 agg.1 = agg.1.max(entry.timestamp_micros);
177 }
178 let mut sessions: Vec<(&str, u32, i64)> = session_agg
179 .into_iter()
180 .map(|(sid, (count, max_ts))| (sid, count, max_ts))
181 .collect();
182 sessions.sort_unstable_by(|left, right| left.0.cmp(right.0));
183 let session_index = index_of(sessions.iter().map(|(value, _, _)| *value));
184
185 let projects = distinct_sorted(entries.iter().map(|entry| entry.project.as_str()));
186 let project_index = index_of(projects.iter().copied());
187 let agents = distinct_sorted(entries.iter().map(|entry| entry.source_agent.as_str()));
188 let agent_index = index_of(agents.iter().copied());
189 let roles = distinct_sorted(entries.iter().map(|entry| entry.role.as_str()));
190 let role_index = index_of(roles.iter().copied());
191
192 let mut blob: Vec<u8> = Vec::new();
193
194 let mut block_entries = Vec::with_capacity(entries.len().div_ceil(BLOCK_ROWS));
197 let mut spans: Vec<(u32, u32)> = Vec::with_capacity(entries.len());
198 for chunk in entries.chunks(BLOCK_ROWS) {
199 let mut plain = Vec::new();
200 for entry in chunk {
201 let off = u32::try_from(plain.len()).context("block too large")?;
202 let len = u32::try_from(entry.search_text.len()).context("search_text too long")?;
203 plain.extend_from_slice(entry.search_text.as_bytes());
204 spans.push((off, len));
205 }
206 let compressed = zstd::bulk::compress(&plain, ZSTD_LEVEL).context("zstd compress")?;
207 block_entries.push(BlockEntry {
208 comp_off: blob.len() as u64,
209 comp_len: u32::try_from(compressed.len()).context("compressed block too large")?,
210 decomp_len: u32::try_from(plain.len()).context("block too large")?,
211 });
212 blob.extend_from_slice(&compressed);
213 }
214
215 let mut records = Vec::with_capacity(entries.len());
216 for (entry, (text_off, text_len)) in entries.iter().zip(&spans) {
217 let blob_off = blob.len() as u64;
218 blob.extend_from_slice(&entry.timestamp_micros.to_le_bytes());
219 blob.extend_from_slice(&session_index[entry.session_id.as_str()].to_le_bytes());
220 blob.extend_from_slice(&project_index[entry.project.as_str()].to_le_bytes());
221 blob.extend_from_slice(&agent_index[entry.source_agent.as_str()].to_le_bytes());
222 blob.extend_from_slice(&role_index[entry.role.as_str()].to_le_bytes());
223 let mid_len = u32::try_from(entry.message_id.len()).context("message_id too long")?;
224 blob.extend_from_slice(&mid_len.to_le_bytes());
225 blob.extend_from_slice(&text_off.to_le_bytes());
226 blob.extend_from_slice(&text_len.to_le_bytes());
227 blob.extend_from_slice(entry.message_id.as_bytes());
228 records.push(Record {
229 row_id: entry.row_id,
230 blob_off,
231 });
232 }
233
234 let session_entries = sessions
235 .iter()
236 .map(|(sid, count, max_ts_micros)| {
237 let off = blob.len() as u64;
238 blob.extend_from_slice(sid.as_bytes());
239 Ok(SessionEntry {
240 sid_off: off,
241 max_ts_micros: *max_ts_micros,
242 sid_len: u32::try_from(sid.len()).context("session_id too long")?,
243 count: *count,
244 })
245 })
246 .collect::<Result<Vec<_>>>()?;
247 let project_entries = dict_entries(&mut blob, &projects)?;
248 let agent_entries = dict_entries(&mut blob, &agents)?;
249 let role_entries = dict_entries(&mut blob, &roles)?;
250
251 let blob_offset = (size_of::<Header>()
252 + records.len() * size_of::<Record>()
253 + session_entries.len() * size_of::<SessionEntry>()
254 + (project_entries.len() + agent_entries.len() + role_entries.len())
255 * size_of::<DictEntry>()
256 + block_entries.len() * size_of::<BlockEntry>()) as u64;
257 let header = Header {
258 magic: MAGIC,
259 version,
260 count: records.len() as u64,
261 session_count: session_entries.len() as u64,
262 project_count: project_entries.len() as u64,
263 agent_count: agent_entries.len() as u64,
264 role_count: role_entries.len() as u64,
265 block_count: block_entries.len() as u64,
266 blob_offset,
267 };
268
269 let tmp = path.with_extension(format!(
273 "tmp-{}-{:016x}",
274 std::process::id(),
275 fastrand::u64(..)
276 ));
277 {
278 let mut file = File::create(&tmp)
279 .with_context(|| format!("create row meta map temp {}", tmp.display()))?;
280 file.write_all(bytemuck::bytes_of(&header))?;
281 file.write_all(bytemuck::cast_slice(&records))?;
282 file.write_all(bytemuck::cast_slice(&session_entries))?;
283 file.write_all(bytemuck::cast_slice(&project_entries))?;
284 file.write_all(bytemuck::cast_slice(&agent_entries))?;
285 file.write_all(bytemuck::cast_slice(&role_entries))?;
286 file.write_all(bytemuck::cast_slice(&block_entries))?;
287 file.write_all(&blob)?;
288 file.sync_all()?;
289 }
290 std::fs::rename(&tmp, path)
291 .with_context(|| format!("rename row meta map into place {}", path.display()))?;
292 Ok(())
293 }
294
295 pub fn open(path: &Path) -> Result<Self> {
296 let file =
297 File::open(path).with_context(|| format!("open row meta map {}", path.display()))?;
298 #[allow(unsafe_code)]
301 let mmap = unsafe { Mmap::map(&file)? };
302 ensure!(
303 mmap.len() >= size_of::<Header>(),
304 "row meta map {} too small for header",
305 path.display()
306 );
307 let header: Header = *bytemuck::from_bytes(&mmap[..size_of::<Header>()]);
308 ensure!(
309 header.magic == MAGIC,
310 "row meta map {} bad magic",
311 path.display()
312 );
313 let count = usize::try_from(header.count).context("count overflow")?;
314 let session_count = usize::try_from(header.session_count).context("session_count")?;
315 let project_count = usize::try_from(header.project_count).context("project_count")?;
316 let agent_count = usize::try_from(header.agent_count).context("agent_count")?;
317 let role_count = usize::try_from(header.role_count).context("role_count")?;
318 let block_count = usize::try_from(header.block_count).context("block_count")?;
319 let blob_offset = usize::try_from(header.blob_offset).context("blob_offset overflow")?;
320
321 let sessions_off = size_of::<Header>() + count * size_of::<Record>();
322 let projects_off = sessions_off + session_count * size_of::<SessionEntry>();
323 let agents_off = projects_off + project_count * size_of::<DictEntry>();
324 let roles_off = agents_off + agent_count * size_of::<DictEntry>();
325 let blocks_off = roles_off + role_count * size_of::<DictEntry>();
326 let blob_offset_expected = blocks_off + block_count * size_of::<BlockEntry>();
327 ensure!(
328 blob_offset == blob_offset_expected && mmap.len() >= blob_offset,
329 "row meta map {} layout mismatch",
330 path.display()
331 );
332 Ok(Self {
333 mmap,
334 version: header.version,
335 count,
336 session_count,
337 project_count,
338 agent_count,
339 role_count,
340 block_count,
341 sessions_off,
342 projects_off,
343 agents_off,
344 roles_off,
345 blocks_off,
346 blob_offset,
347 })
348 }
349
350 pub fn version(&self) -> u64 {
351 self.version
352 }
353
354 pub fn len(&self) -> usize {
355 self.count
356 }
357
358 pub fn is_empty(&self) -> bool {
359 self.count == 0
360 }
361
362 pub fn max_row_id(&self) -> Option<u64> {
365 self.records().last().map(|record| record.row_id)
366 }
367
368 fn records(&self) -> &[Record] {
369 let start = size_of::<Header>();
370 let end = start + self.count * size_of::<Record>();
371 bytemuck::cast_slice(&self.mmap[start..end])
372 }
373
374 fn session_entries(&self) -> &[SessionEntry] {
375 let end = self.sessions_off + self.session_count * size_of::<SessionEntry>();
376 bytemuck::cast_slice(&self.mmap[self.sessions_off..end])
377 }
378
379 fn dict_entries(&self, start: usize, count: usize) -> &[DictEntry] {
380 let end = start + count * size_of::<DictEntry>();
381 bytemuck::cast_slice(&self.mmap[start..end])
382 }
383
384 fn block_entries(&self) -> &[BlockEntry] {
385 let end = self.blocks_off + self.block_count * size_of::<BlockEntry>();
386 bytemuck::cast_slice(&self.mmap[self.blocks_off..end])
387 }
388
389 fn blob_str(&self, off: u64, len: u32) -> &str {
391 let base = self.blob_offset.saturating_add(off as usize);
392 let end = base.saturating_add(len as usize);
393 self.mmap
394 .get(base..end)
395 .and_then(|bytes| std::str::from_utf8(bytes).ok())
396 .unwrap_or_default()
397 }
398
399 fn session_str(&self, index: usize) -> &str {
400 match self.session_entries().get(index) {
401 Some(entry) => self.blob_str(entry.sid_off, entry.sid_len),
402 None => "",
403 }
404 }
405
406 fn dict_str(&self, start: usize, count: usize, index: usize) -> &str {
407 match self.dict_entries(start, count).get(index) {
408 Some(entry) => self.blob_str(entry.off, entry.len),
409 None => "",
410 }
411 }
412
413 fn locate(&self, row_id: u64) -> Option<(usize, usize)> {
415 let records = self.records();
416 let idx = records
417 .binary_search_by_key(&row_id, |record| record.row_id)
418 .ok()?;
419 let base = self
420 .blob_offset
421 .checked_add(usize::try_from(records[idx].blob_off).ok()?)?;
422 Some((idx, base))
423 }
424
425 pub fn lookup(&self, row_id: u64) -> Option<(&str, &str)> {
429 let (_, base) = self.locate(row_id)?;
430 let header = self.mmap.get(base..base.checked_add(ROW_HEADER_LEN)?)?;
431 let session_idx = read_u32(header, 8)?;
432 let mid_len = read_u32(header, 24)?;
433 let mut at = base + ROW_HEADER_LEN;
434 let mid = self.slice_str(&mut at, mid_len)?;
435 Some((self.session_str(session_idx), mid))
436 }
437
438 pub fn lookup_meta(&self, row_id: u64, cache: &mut BlockCache) -> Option<RowMeta<'_>> {
442 let (idx, base) = self.locate(row_id)?;
443 let header = self.mmap.get(base..base.checked_add(ROW_HEADER_LEN)?)?;
444 let timestamp_micros = i64::from_le_bytes(header.get(0..8)?.try_into().ok()?);
445 let session_idx = read_u32(header, 8)?;
446 let project_idx = read_u32(header, 12)?;
447 let agent_idx = read_u32(header, 16)?;
448 let role_idx = read_u32(header, 20)?;
449 let mid_len = read_u32(header, 24)?;
450 let text_off = read_u32(header, 28)?;
451 let text_len = read_u32(header, 32)?;
452 let mut at = base + ROW_HEADER_LEN;
453 let message_id = self.slice_str(&mut at, mid_len)?;
454 let search_text = self.decompress_text(idx, text_off, text_len, cache)?;
455 Some(RowMeta {
456 session_id: self.session_str(session_idx),
457 message_id,
458 role: self.dict_str(self.roles_off, self.role_count, role_idx),
459 project: self.dict_str(self.projects_off, self.project_count, project_idx),
460 source_agent: self.dict_str(self.agents_off, self.agent_count, agent_idx),
461 timestamp_micros,
462 search_text,
463 })
464 }
465
466 fn decompress_block(&self, block_idx: usize) -> Option<Vec<u8>> {
467 let block = self.block_entries().get(block_idx)?;
468 if block.decomp_len == 0 {
469 return Some(Vec::new());
470 }
471 let comp_base = self.blob_offset.checked_add(block.comp_off as usize)?;
472 let comp = self
473 .mmap
474 .get(comp_base..comp_base.checked_add(block.comp_len as usize)?)?;
475 zstd::bulk::decompress(comp, block.decomp_len as usize).ok()
476 }
477
478 fn decompress_text(
479 &self,
480 idx: usize,
481 text_off: usize,
482 text_len: usize,
483 cache: &mut BlockCache,
484 ) -> Option<String> {
485 if text_len == 0 {
486 return Some(String::new());
487 }
488 let block_idx = idx / BLOCK_ROWS;
489 if cache.as_ref().map(|(block, _)| *block) != Some(block_idx) {
490 *cache = Some((block_idx, self.decompress_block(block_idx)?));
491 }
492 let plain = &cache.as_ref()?.1;
493 let value = plain.get(text_off..text_off.checked_add(text_len)?)?;
494 String::from_utf8(value.to_vec()).ok()
495 }
496
497 pub fn entries(&self) -> Vec<RowMetaEntry> {
501 let records = self.records();
502 let mut out = Vec::with_capacity(records.len());
503 let mut current_block = usize::MAX;
504 let mut plain: Vec<u8> = Vec::new();
505 for (idx, record) in records.iter().enumerate() {
506 let block_idx = idx / BLOCK_ROWS;
507 if block_idx != current_block {
508 plain = self.decompress_block(block_idx).unwrap_or_default();
509 current_block = block_idx;
510 }
511 let Some(base) = self.blob_offset.checked_add(record.blob_off as usize) else {
512 continue;
513 };
514 if let Some(entry) = self.entry_at(record.row_id, base, &plain) {
515 out.push(entry);
516 }
517 }
518 out
519 }
520
521 fn entry_at(&self, row_id: u64, base: usize, block_plain: &[u8]) -> Option<RowMetaEntry> {
522 let header = self.mmap.get(base..base.checked_add(ROW_HEADER_LEN)?)?;
523 let timestamp_micros = i64::from_le_bytes(header.get(0..8)?.try_into().ok()?);
524 let session_idx = read_u32(header, 8)?;
525 let project_idx = read_u32(header, 12)?;
526 let agent_idx = read_u32(header, 16)?;
527 let role_idx = read_u32(header, 20)?;
528 let mid_len = read_u32(header, 24)?;
529 let text_off = read_u32(header, 28)?;
530 let text_len = read_u32(header, 32)?;
531 let mut at = base + ROW_HEADER_LEN;
532 let message_id = self.slice_str(&mut at, mid_len)?.to_owned();
533 let search_text = if text_len == 0 {
534 String::new()
535 } else {
536 let bytes = block_plain.get(text_off..text_off.checked_add(text_len)?)?;
537 String::from_utf8(bytes.to_vec()).ok()?
538 };
539 Some(RowMetaEntry {
540 row_id,
541 session_id: self.session_str(session_idx).to_owned(),
542 message_id,
543 role: self
544 .dict_str(self.roles_off, self.role_count, role_idx)
545 .to_owned(),
546 project: self
547 .dict_str(self.projects_off, self.project_count, project_idx)
548 .to_owned(),
549 source_agent: self
550 .dict_str(self.agents_off, self.agent_count, agent_idx)
551 .to_owned(),
552 timestamp_micros,
553 search_text,
554 })
555 }
556
557 pub fn lookup_count(&self, session_id: &str) -> Option<usize> {
560 let idx = self.session_index(session_id)?;
561 Some(self.session_entries()[idx].count as usize)
562 }
563
564 pub fn lookup_max_ts(&self, session_id: &str) -> Option<i64> {
568 let idx = self.session_index(session_id)?;
569 Some(self.session_entries()[idx].max_ts_micros)
570 }
571
572 fn session_index(&self, session_id: &str) -> Option<usize> {
575 self.session_entries()
576 .binary_search_by(|entry| self.blob_str(entry.sid_off, entry.sid_len).cmp(session_id))
577 .ok()
578 }
579
580 fn header_at(&self, record: &Record) -> Option<&[u8]> {
583 let base = self.blob_offset.checked_add(record.blob_off as usize)?;
584 self.mmap.get(base..base.checked_add(ROW_HEADER_LEN)?)
585 }
586
587 pub fn lookup_session_for_message(&self, message_id: &str) -> Option<&str> {
594 let needle = message_id.as_bytes();
595 for record in self.records().iter().rev() {
596 let header = self.header_at(record)?;
597 if read_u32(header, 24)? != needle.len() {
598 continue;
599 }
600 let base = self.blob_offset.checked_add(record.blob_off as usize)?;
601 let start = base + ROW_HEADER_LEN;
602 if self.mmap.get(start..start.checked_add(needle.len())?)? == needle {
603 let session_id = self.session_str(read_u32(header, 8)?);
604 return (!session_id.is_empty()).then_some(session_id);
605 }
606 }
607 None
608 }
609
610 pub fn session_row_ids(&self, session_id: &str) -> Option<Vec<u64>> {
616 let Some(session_idx) = self.session_index(session_id) else {
617 return Some(Vec::new());
618 };
619 let count = self.session_entries()[session_idx].count as usize;
620 let mut out = Vec::with_capacity(count);
621 for record in self.records() {
622 let header = self.header_at(record)?;
623 if read_u32(header, 8)? == session_idx {
624 out.push(record.row_id);
625 if out.len() == count {
626 break;
627 }
628 }
629 }
630 Some(out)
631 }
632
633 fn slice_str(&self, at: &mut usize, len: usize) -> Option<&str> {
636 let end = at.checked_add(len)?;
637 let bytes = self.mmap.get(*at..end)?;
638 *at = end;
639 std::str::from_utf8(bytes).ok()
640 }
641}
642
643pub struct ChainPaths {
646 pub base: PathBuf,
647 pub base_version: u64,
648 pub deltas: Vec<(u64, PathBuf)>,
649}
650
651impl ChainPaths {
652 pub fn version(&self) -> u64 {
654 self.deltas
655 .last()
656 .map(|(version, _)| *version)
657 .unwrap_or(self.base_version)
658 }
659}
660
661pub fn discover_chain(cache_dir: &Path, store_key: &str) -> Option<ChainPaths> {
665 let prefix = format!("rowmetamap-{store_key}-");
666 let mut bases: Vec<(u64, PathBuf)> = Vec::new();
667 let mut deltas: Vec<(u64, PathBuf)> = Vec::new();
668 for entry in std::fs::read_dir(cache_dir).ok()?.flatten() {
669 let name = entry.file_name();
670 let Some(rest) = name
671 .to_str()
672 .and_then(|name| name.strip_prefix(&prefix))
673 .and_then(|rest| rest.strip_suffix(".rmm"))
674 else {
675 continue;
676 };
677 if let Some(version) = rest.strip_prefix('v').and_then(|d| d.parse::<u64>().ok()) {
678 bases.push((version, entry.path()));
679 } else if let Some(version) = rest.strip_prefix('d').and_then(|d| d.parse::<u64>().ok()) {
680 deltas.push((version, entry.path()));
681 }
682 }
683 let (base_version, base) = bases.into_iter().max_by_key(|(version, _)| *version)?;
684 let mut deltas: Vec<(u64, PathBuf)> = deltas
685 .into_iter()
686 .filter(|(version, _)| *version > base_version)
687 .collect();
688 deltas.sort_by_key(|(version, _)| *version);
689 Some(ChainPaths {
690 base,
691 base_version,
692 deltas,
693 })
694}
695
696pub struct RowMetaSet {
701 segments: Vec<RowMetaMap>,
702}
703
704impl std::fmt::Debug for RowMetaSet {
705 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
706 formatter
707 .debug_struct("RowMetaSet")
708 .field("segments", &self.segments.len())
709 .field("version", &self.version())
710 .finish()
711 }
712}
713
714impl RowMetaSet {
715 pub fn open(paths: &ChainPaths) -> Result<Self> {
717 let mut segments = Vec::with_capacity(1 + paths.deltas.len());
718 segments.push(RowMetaMap::open(&paths.base)?);
719 for (_, delta) in &paths.deltas {
720 segments.push(RowMetaMap::open(delta)?);
721 }
722 Ok(Self { segments })
723 }
724
725 pub fn version(&self) -> u64 {
726 self.segments
727 .iter()
728 .map(RowMetaMap::version)
729 .max()
730 .unwrap_or(0)
731 }
732
733 pub fn delta_count(&self) -> usize {
735 self.segments.len().saturating_sub(1)
736 }
737
738 pub fn is_empty(&self) -> bool {
741 self.segments.iter().all(RowMetaMap::is_empty)
742 }
743
744 pub fn len(&self) -> usize {
747 self.segments.iter().map(RowMetaMap::len).sum()
748 }
749
750 pub fn max_row_id(&self) -> Option<u64> {
753 self.segments
754 .iter()
755 .filter_map(RowMetaMap::max_row_id)
756 .max()
757 }
758
759 pub fn lookup(&self, row_id: u64) -> Option<(&str, &str)> {
761 self.segments
762 .iter()
763 .rev()
764 .find_map(|seg| seg.lookup(row_id))
765 }
766
767 pub fn hydrate(&self, rowids: &[u64]) -> (Vec<RowMetaEntry>, Vec<u64>) {
773 let mut sorted = rowids.to_vec();
774 sorted.sort_unstable();
775 let mut caches: Vec<BlockCache> = vec![None; self.segments.len()];
776 let mut hits = Vec::with_capacity(sorted.len());
777 let mut misses = Vec::new();
778 for row_id in sorted {
779 let hit = self
780 .segments
781 .iter()
782 .enumerate()
783 .rev()
784 .find_map(|(segment, map)| {
785 let meta = map.lookup_meta(row_id, &mut caches[segment])?;
786 Some(RowMetaEntry {
787 row_id,
788 session_id: meta.session_id.to_owned(),
789 message_id: meta.message_id.to_owned(),
790 role: meta.role.to_owned(),
791 project: meta.project.to_owned(),
792 source_agent: meta.source_agent.to_owned(),
793 timestamp_micros: meta.timestamp_micros,
794 search_text: meta.search_text,
795 })
796 });
797 match hit {
798 Some(entry) => hits.push(entry),
799 None => misses.push(row_id),
800 }
801 }
802 (hits, misses)
803 }
804
805 pub fn lookup_count(&self, session_id: &str) -> Option<usize> {
807 let mut total = 0;
808 let mut found = false;
809 for seg in &self.segments {
810 if let Some(count) = seg.lookup_count(session_id) {
811 total += count;
812 found = true;
813 }
814 }
815 found.then_some(total)
816 }
817
818 pub fn lookup_max_ts(&self, session_id: &str) -> Option<i64> {
822 self.segments
823 .iter()
824 .filter_map(|seg| seg.lookup_max_ts(session_id))
825 .max()
826 }
827
828 pub fn lookup_session_for_message(&self, message_id: &str) -> Option<&str> {
832 self.segments
833 .iter()
834 .rev()
835 .find_map(|seg| seg.lookup_session_for_message(message_id))
836 }
837
838 pub fn session_row_ids(&self, session_id: &str) -> Option<Vec<u64>> {
842 let mut out = Vec::new();
843 for seg in &self.segments {
844 out.extend(seg.session_row_ids(session_id)?);
845 }
846 Some(out)
847 }
848
849 pub fn merged_entries(&self) -> Vec<RowMetaEntry> {
852 let mut by_row: HashMap<u64, RowMetaEntry> = HashMap::new();
853 for seg in &self.segments {
854 for entry in seg.entries() {
855 by_row.insert(entry.row_id, entry);
856 }
857 }
858 by_row.into_values().collect()
859 }
860}
861
862fn distinct_sorted<'a>(values: impl Iterator<Item = &'a str>) -> Vec<&'a str> {
863 let mut distinct: Vec<&str> = values.collect();
864 distinct.sort_unstable();
865 distinct.dedup();
866 distinct
867}
868
869fn index_of<'a>(values: impl Iterator<Item = &'a str>) -> HashMap<&'a str, u32> {
870 values
871 .enumerate()
872 .map(|(index, value)| (value, index as u32))
873 .collect()
874}
875
876fn dict_entries(blob: &mut Vec<u8>, values: &[&str]) -> Result<Vec<DictEntry>> {
877 values
878 .iter()
879 .map(|value| {
880 let off = blob.len() as u64;
881 blob.extend_from_slice(value.as_bytes());
882 Ok(DictEntry {
883 off,
884 len: u32::try_from(value.len()).context("dictionary value too long")?,
885 _pad: 0,
886 })
887 })
888 .collect()
889}
890
891fn read_u32(bytes: &[u8], at: usize) -> Option<usize> {
892 let slice = bytes.get(at..at.checked_add(4)?)?;
893 Some(u32::from_le_bytes(slice.try_into().ok()?) as usize)
894}
895
896#[cfg(test)]
897mod tests {
898 #![allow(clippy::expect_used, clippy::unwrap_used)]
899 use super::*;
900
901 fn entry(
902 row_id: u64,
903 session_id: &str,
904 message_id: &str,
905 timestamp_micros: i64,
906 search_text: &str,
907 ) -> RowMetaEntry {
908 RowMetaEntry {
909 row_id,
910 session_id: session_id.to_owned(),
911 message_id: message_id.to_owned(),
912 role: "user".to_owned(),
913 project: "/proj".to_owned(),
914 source_agent: "claude-code".to_owned(),
915 timestamp_micros,
916 search_text: search_text.to_owned(),
917 }
918 }
919
920 #[test]
921 fn message_and_conversational_lookups_cover_the_chain() {
922 let dir = tempfile::tempdir().unwrap();
923 let base_path = RowMetaMap::path_for(dir.path(), "s", 1);
924 RowMetaMap::build(
925 &base_path,
926 1,
927 vec![
928 entry(1, "sess-a", "msg-1", 1_000, "hello"),
929 entry(2, "sess-a", "msg-2", 2_000, ""), entry(3, "sess-b", "msg-3", 3_000, "there"),
931 ],
932 )
933 .unwrap();
934 let delta_path = RowMetaMap::delta_path(dir.path(), "s", 2);
935 RowMetaMap::build(
936 &delta_path,
937 2,
938 vec![entry(9, "sess-a", "msg-9", 9_000, "newest")],
939 )
940 .unwrap();
941 let set = RowMetaSet::open(&ChainPaths {
942 base: base_path,
943 base_version: 1,
944 deltas: vec![(2, delta_path)],
945 })
946 .unwrap();
947
948 assert_eq!(set.lookup_session_for_message("msg-1"), Some("sess-a"));
949 assert_eq!(
950 set.lookup_session_for_message("msg-9"),
951 Some("sess-a"),
952 "delta hit"
953 );
954 assert_eq!(set.lookup_session_for_message("msg-3"), Some("sess-b"));
955 assert_eq!(set.lookup_session_for_message("absent"), None);
956
957 let mut ids = set.session_row_ids("sess-a").expect("intact map");
958 ids.sort_unstable();
959 assert_eq!(ids, vec![1, 2, 9], "all roles, base and delta");
960 assert_eq!(
961 set.session_row_ids("missing").expect("intact map"),
962 Vec::<u64>::new(),
963 "absent session is empty, not a corruption signal"
964 );
965 }
966
967 #[test]
968 fn build_open_lookup_roundtrip() {
969 let dir = tempfile::tempdir().unwrap();
970 let path = RowMetaMap::path_for(dir.path(), "teststore", 7);
971 let mut three = entry(99, "sess-a", "msg-3", 3_000, "third");
972 three.role = "assistant".to_owned();
973 three.project = "/other".to_owned();
974 let entries = vec![
975 entry(10, "sess-a", "msg-1", 1_000, "first message text"),
976 entry(3, "sess-b/agent-x", "msg-2", 2_000, ""),
977 three,
978 ];
979 RowMetaMap::build(&path, 7, entries).unwrap();
980
981 let map = RowMetaMap::open(&path).unwrap();
982 assert_eq!(map.version(), 7);
983 assert_eq!(map.len(), 3);
984 assert_eq!(map.lookup(10), Some(("sess-a", "msg-1")));
985 assert_eq!(map.lookup(3), Some(("sess-b/agent-x", "msg-2")));
986 assert_eq!(map.lookup(99), Some(("sess-a", "msg-3")));
987 assert_eq!(map.lookup(42), None);
988
989 let meta = map.lookup_meta(10, &mut None).expect("row 10 present");
990 assert_eq!(meta.session_id, "sess-a");
991 assert_eq!(meta.message_id, "msg-1");
992 assert_eq!(meta.role, "user");
993 assert_eq!(meta.project, "/proj");
994 assert_eq!(meta.source_agent, "claude-code");
995 assert_eq!(meta.timestamp_micros, 1_000);
996 assert_eq!(meta.search_text, "first message text");
997
998 let assistant = map.lookup_meta(99, &mut None).expect("row 99 present");
999 assert_eq!(assistant.role, "assistant");
1000 assert_eq!(assistant.project, "/other");
1001 assert_eq!(assistant.search_text, "third");
1002
1003 let empty_text = map.lookup_meta(3, &mut None).expect("row 3 present");
1004 assert_eq!(empty_text.search_text, "");
1005 assert!(map.lookup_meta(42, &mut None).is_none());
1006
1007 assert_eq!(map.lookup_count("sess-a"), Some(2));
1008 assert_eq!(map.lookup_count("sess-b/agent-x"), Some(1));
1009 assert_eq!(map.lookup_count("missing"), None);
1010
1011 assert_eq!(map.lookup_max_ts("sess-a"), Some(3_000));
1013 assert_eq!(map.lookup_max_ts("sess-b/agent-x"), Some(2_000));
1014 assert_eq!(map.lookup_max_ts("missing"), None);
1015 }
1016
1017 #[test]
1018 fn max_ts_is_the_session_high_water_mark() {
1019 let dir = tempfile::tempdir().unwrap();
1020 let path = RowMetaMap::path_for(dir.path(), "ts", 1);
1021 let entries = vec![
1023 entry(1, "s", "msg-a", 5_000, "a"),
1024 entry(2, "s", "msg-b", 9_000, "b"),
1025 entry(3, "s", "msg-c", 7_000, "c"),
1026 ];
1027 RowMetaMap::build(&path, 1, entries).unwrap();
1028 let map = RowMetaMap::open(&path).unwrap();
1029 assert_eq!(map.lookup_max_ts("s"), Some(9_000));
1030 }
1031
1032 #[test]
1033 fn many_blocks_roundtrip() {
1034 let dir = tempfile::tempdir().unwrap();
1035 let path = RowMetaMap::path_for(dir.path(), "blocks", 1);
1036 let entries: Vec<RowMetaEntry> = (0..(BLOCK_ROWS as u64 * 2 + 5))
1037 .map(|i| {
1038 entry(
1039 i,
1040 "sess",
1041 &format!("msg-{i}"),
1042 i as i64,
1043 &format!("text body {i}"),
1044 )
1045 })
1046 .collect();
1047 RowMetaMap::build(&path, 1, entries).unwrap();
1048
1049 let map = RowMetaMap::open(&path).unwrap();
1050 let mut cache = None;
1054 for i in [0u64, 1, 255, 256, 257, 511, 512, 516] {
1055 let meta = map.lookup_meta(i, &mut cache).expect("row present");
1056 assert_eq!(meta.message_id, format!("msg-{i}"));
1057 assert_eq!(meta.search_text, format!("text body {i}"));
1058 }
1059 }
1060
1061 #[test]
1062 fn lsm_set_layers_delta_over_base() {
1063 let dir = tempfile::tempdir().unwrap();
1064 let base = vec![
1065 entry(10, "sess-a", "m10", 1, "base ten"),
1066 entry(11, "sess-a", "m11", 2, "base eleven"),
1067 entry(12, "sess-b", "m12", 3, "base twelve"),
1068 ];
1069 RowMetaMap::build(&RowMetaMap::path_for(dir.path(), "k", 1), 1, base).unwrap();
1070 let delta = vec![
1071 entry(20, "sess-a", "m20", 4, "delta twenty"),
1072 entry(21, "sess-c", "m21", 5, "delta twentyone"),
1073 ];
1074 RowMetaMap::build(&RowMetaMap::delta_path(dir.path(), "k", 2), 2, delta).unwrap();
1075
1076 let chain = discover_chain(dir.path(), "k").expect("chain present");
1077 assert_eq!(chain.base_version, 1);
1078 assert_eq!(chain.deltas.len(), 1);
1079 assert_eq!(chain.version(), 2);
1080
1081 let set = RowMetaSet::open(&chain).unwrap();
1082 assert_eq!(set.version(), 2);
1083 assert_eq!(set.delta_count(), 1);
1084
1085 assert_eq!(set.lookup(10), Some(("sess-a", "m10")));
1086 assert_eq!(set.lookup(20), Some(("sess-a", "m20")));
1087 assert_eq!(set.lookup(99), None);
1088
1089 let (mut hits, misses) = set.hydrate(&[21, 10, 99]);
1091 assert_eq!(misses, vec![99]);
1092 hits.sort_by_key(|entry| entry.row_id);
1093 assert_eq!(hits.len(), 2);
1094 assert_eq!(hits[0].search_text, "base ten");
1095 assert_eq!(hits[1].search_text, "delta twentyone");
1096
1097 assert_eq!(set.lookup_count("sess-a"), Some(3));
1099 assert_eq!(set.lookup_count("sess-b"), Some(1));
1100 assert_eq!(set.lookup_count("sess-c"), Some(1));
1101 assert_eq!(set.lookup_count("missing"), None);
1102
1103 assert_eq!(set.lookup_max_ts("sess-a"), Some(4));
1105 assert_eq!(set.lookup_max_ts("sess-b"), Some(3));
1106 assert_eq!(set.lookup_max_ts("sess-c"), Some(5));
1107 assert_eq!(set.lookup_max_ts("missing"), None);
1108
1109 let mut merged = set.merged_entries();
1111 merged.sort_by_key(|entry| entry.row_id);
1112 assert_eq!(merged.len(), 5);
1113 assert_eq!(merged[0].row_id, 10);
1114 assert_eq!(merged[4].row_id, 21);
1115 assert_eq!(merged[4].search_text, "delta twentyone");
1116 }
1117
1118 #[test]
1119 fn empty_map_roundtrips() {
1120 let dir = tempfile::tempdir().unwrap();
1121 let path = RowMetaMap::path_for(dir.path(), "empty", 1);
1122 RowMetaMap::build(&path, 1, Vec::new()).unwrap();
1123 let map = RowMetaMap::open(&path).unwrap();
1124 assert!(map.is_empty());
1125 assert_eq!(map.lookup(0), None);
1126 assert!(map.lookup_meta(0, &mut None).is_none());
1127 assert_eq!(map.lookup_count("anything"), None);
1128 assert_eq!(map.lookup_max_ts("anything"), None);
1129 }
1130}