Skip to main content

pond/
rowmap.rs

1//! Memory-mapped, process-shareable per-message meta map keyed by stable
2//! `row_id`. Resolves FTS/vector `_rowid`s to `(session_id, message_id)` and
3//! hydrates hit meta (`role`, `project`, `source_agent`, `timestamp`,
4//! `search_text`) in memory, with a `take_rows` miss-fallback for rows appended
5//! since the build. Also carries a `session_id -> message count` aggregate.
6//!
7//! `session_id`/`project`/`source_agent`/`role` are dictionary-encoded (each
8//! distinct value stored once, referenced by `u32` index); `search_text` is
9//! block-compressed ([`BLOCK_ROWS`] rows per zstd block). On the real 2M-message
10//! corpus that takes the map from ~655 MB (flat) to ~270 MB.
11//!
12//! Built via temp + atomic rename and `mmap`'d read-only, so N pond processes on
13//! the box share one physical copy in the OS page cache and a restart re-`open`s
14//! instantly. Stable row ids (`enable_stable_row_ids`) keep a built map valid
15//! across compaction; it only rebuilds when the dataset version advances.
16//!
17//! Layout: `Header | [Record; count] | [SessionEntry] | [DictEntry; project] |
18//! [DictEntry; agent] | [DictEntry; role] | [BlockEntry] | blob`. Records are
19//! sorted by `row_id` (binary search); a row's block is `record_index /
20//! BLOCK_ROWS`. The blob holds the compressed `search_text` blocks, then per-row
21//! `{32-byte header + message_id}`, then per-session `session_id` bytes, then the
22//! dict value bytes. Each session entry also carries its max message timestamp,
23//! the watermark the `pond sync` skip oracle compares against the source.
24
25use 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
39/// Per-row blob header: `timestamp_micros` (i64 LE) then seven `u32` LE fields -
40/// the four dictionary indices (`session`, `project`, `source_agent`, `role`),
41/// the `message_id` length, and the `search_text` offset+length within its
42/// decompressed block.
43const ROW_HEADER_LEN: usize = 8 + 7 * 4;
44
45/// One-slot decompressed-block cache `(block_index, plaintext)`, threaded
46/// through a batch of `lookup_meta` calls so rows sharing a block (a session's
47/// hits are row-id-adjacent) decompress it once, not once per row.
48type 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/// Owned input row for [`RowMetaMap::build`].
97#[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
109/// Borrowed view of one row's meta. The dictionary-encoded fields borrow the
110/// mmap; `search_text` is owned (decompressed from its block).
111pub 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
121/// An open, memory-mapped row meta map. `lookup`, `lookup_meta`, and
122/// `lookup_count` are lock-free and reentrant.
123pub 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    /// Base segment path (`-v{version}`): the foot of the LSM chain.
153    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    /// Delta segment path (`-d{version}`): rows appended since the previous
158    /// segment, layered over the base.
159    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        // Per session: message count plus the max message timestamp - the
167        // watermark the sync skip oracle compares against the source's latest
168        // message timestamp (spec.md#adapters; deterministic, rebuilt from the
169        // store, no local cursor).
170        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        // Compressed search_text blocks first; spans record each row's
195        // offset+length within its decompressed block.
196        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        // Unique temp name per builder (pid + nonce): two processes prewarming
270        // the same store+version must not share one temp inode, or the second's
271        // create would mutate the file the first is mapping.
272        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        // SAFETY: the file is immutable once renamed into place, so the mapping
299        // never sees concurrent truncation/mutation.
300        #[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    /// Highest `row_id` in this segment (records are row_id-sorted), or `None`
363    /// when empty. With stable row ids this is the append high-water mark.
364    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    /// `""` on a corrupt/truncated extent - treated as a miss, never a panic.
390    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    /// `(record index, blob header offset)` for `row_id`, `None` if unmapped.
414    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    /// Resolve a `row_id` to `(session_id, message_id)` - the cheap arm-
426    /// resolution path. Never decompresses a block. `None` for rows appended
427    /// after the build (caller falls back to a data take).
428    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    /// Resolve a `row_id` to its full hydration meta, decompressing the row's
439    /// `search_text` block (reusing `cache` if it already holds that block).
440    /// `None` if unmapped (caller falls back to take_rows).
441    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    /// Reconstruct every row as an owned [`RowMetaEntry`], decompressing each
498    /// block once. Used to merge segments into a fresh base at compaction
499    /// without re-reading the store.
500    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    /// Whole-session message count for `session_id`. `None` if the session is
558    /// not in this map (caller falls back to the `session_id IN (...)` scan).
559    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    /// Max message timestamp (micros) stored for `session_id` - the watermark the
565    /// sync skip oracle compares against the source's latest message timestamp.
566    /// `None` if the session is not in this map.
567    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    /// Index into `session_entries` for `session_id` - the shared binary
573    /// search behind every per-session accessor.
574    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    /// A record's blob header slice. Checked so a corrupt map yields `None`
581    /// (-> caller falls back to the store), not a panic.
582    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    /// Resolve a message id to its session id by scanning the record headers -
588    /// resident memory only, never a store read. Newest rows first: recent
589    /// messages are the likely targets, and records are laid out in ascending
590    /// `row_id` order. Length-check first, so most rows cost one integer
591    /// compare. `None` is a miss, including on a corrupt map - an empty or
592    /// unresolvable session string must never suppress the store scan.
593    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    /// Row ids of every `session_id` row in this segment, `Some(empty)` when
611    /// the session is not here. `None` on any malformed record: a silently
612    /// dropped row would serve an incomplete page with no fallback, so
613    /// corruption aborts the map path instead. The session entry's row count
614    /// bounds the walk - it stops at the session's last row.
615    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    /// Slice `len` UTF-8 bytes at `*at`, advancing `*at`. Checked so a corrupt
634    /// map yields `None` (-> take fallback), not a panic.
635    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
643/// The on-disk LSM chain for a store: the highest-version base plus every
644/// delta layered above it, ascending.
645pub struct ChainPaths {
646    pub base: PathBuf,
647    pub base_version: u64,
648    pub deltas: Vec<(u64, PathBuf)>,
649}
650
651impl ChainPaths {
652    /// Version the chain covers - the newest segment's version.
653    pub fn version(&self) -> u64 {
654        self.deltas
655            .last()
656            .map(|(version, _)| *version)
657            .unwrap_or(self.base_version)
658    }
659}
660
661/// Discover the chain under `cache_dir` for `store_key`: the highest-version
662/// base (`-v{V}`) plus every delta (`-d{V}`) above it, ascending. `None` if no
663/// base exists yet.
664pub 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
696/// An LSM chain of immutable segment maps (base + ascending deltas) viewed as
697/// one logical map. Rows are partitioned across segments by `row_id` (append is
698/// disjoint; compaction rebuilds the base), so key/meta lookups take the newest
699/// hit and counts sum across segments.
700pub 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    /// Open every segment in `paths` (base first, then deltas ascending).
716    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    /// Number of delta segments layered on the base.
734    pub fn delta_count(&self) -> usize {
735        self.segments.len().saturating_sub(1)
736    }
737
738    /// No rows in any segment - the first-ingest hint that lets the sync oracle
739    /// short-circuit the per-session source last-id read.
740    pub fn is_empty(&self) -> bool {
741        self.segments.iter().all(RowMetaMap::is_empty)
742    }
743
744    /// Total row entries across every segment. Rows are disjoint across segments
745    /// (append-only deltas), so this sums - the live row count the base covers.
746    pub fn len(&self) -> usize {
747        self.segments.iter().map(RowMetaMap::len).sum()
748    }
749
750    /// Highest `row_id` across all segments - the append high-water mark a delta
751    /// extends past. Stable row ids keep it monotonic under fragment churn.
752    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    /// Newest segment wins (a row lives in exactly one segment).
760    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    /// Hydrate `rowids` to owned metas, splitting out the ones no segment holds
768    /// (appended since the build) for the caller's take_rows fallback. Output
769    /// order is unspecified (caller indexes by key). Rowids are visited in
770    /// sorted order with a per-segment block cache, so the common case - many
771    /// hits from a few row-id-adjacent sessions - decompresses each block once.
772    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    /// A session's rows are split across segments, so its count is the sum.
806    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    /// Max message timestamp (micros) for `session_id` across the chain - the max
819    /// over every segment that holds it (a session's rows can be split across
820    /// base and deltas). `None` if no segment has it.
821    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    /// Resolve a message id to its session id, newest segment first (recent
829    /// messages - the likely targets - live in the small deltas). `None` is a
830    /// miss the caller resolves against the store.
831    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    /// A session's row ids across the chain (rows are disjoint across
839    /// segments, so this concatenates). `None` if any segment reports a
840    /// malformed record - the caller falls back to the store.
841    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    /// Every row across all segments, newest-segment-wins on `row_id`
850    /// collision - the input to a base rebuild at compaction.
851    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, ""), // bare tool call: not conversational
930                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        // Watermark = max timestamp: sess-a's msg-3 (ts 3000) over msg-1.
1012        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        // Out-of-row-order timestamps: the max wins regardless of row order.
1022        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        // One cache reused across rows spanning several blocks - same-block hits
1051        // must reuse it, block crossings must refill it, both yielding the right
1052        // text.
1053        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        // hydrate spans both segments and splits out the absent row.
1090        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        // Counts sum across base + delta: sess-a = 2 + 1.
1098        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        // Max timestamp across segments: sess-a spans base (ts<=2) + delta (ts 4).
1104        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        // Compaction input: all 5 distinct rows reconstructed.
1110        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}