From 7443a5406c00f68176559bdb3ba26d2c10a7c492 Mon Sep 17 00:00:00 2001 From: adamv-symbolica Date: Sat, 3 Oct 2026 13:23:43 +0100 Subject: [PATCH] Bound ACT streaming line reuse with fingerprint-to-offset caching --- examples/act_line_cache.rs | 39 +++++++++++++++++ src/arena_compact.rs | 86 ++++++++++++++++++++++++++++++++------ tests/act_stream_cache.rs | 56 +++++++++++++++++++++++++ 3 files changed, 169 insertions(+), 12 deletions(-) create mode 100644 examples/act_line_cache.rs create mode 100644 tests/act_stream_cache.rs diff --git a/examples/act_line_cache.rs b/examples/act_line_cache.rs new file mode 100644 index 00000000..41af71ab --- /dev/null +++ b/examples/act_line_cache.rs @@ -0,0 +1,39 @@ +//! Compare ACT line-cache policies on the same ordered .paths input. +//! cargo run --release --example act_line_cache --features arena_compact,act_counters -- INPUT.paths OUTPUT.act POLICY +//! POLICY is none, bounded (4194304 fingerprint/offset entries), or unbounded. +use pathmap::{arena_compact::ACTOutputStream, paths_serialization::for_each_deserialized_path}; +use std::{ + fs::File, + io::{self, BufReader}, + time::Instant, +}; + +fn main() -> io::Result<()> { + let args: Vec<_> = std::env::args().skip(1).collect(); + if args.len() != 3 { + return Err(io::Error::other( + "expected INPUT.paths OUTPUT.act none|bounded|unbounded", + )); + } + let entries = match args[2].as_str() { + "none" => 0, + "bounded" => 4194304, + "unbounded" => usize::MAX, + _ => return Err(io::Error::other("unknown cache policy")), + }; + let start = Instant::now(); + let mut output = ACTOutputStream::with_cache_limit(&args[1], entries)?; + let stats = for_each_deserialized_path(BufReader::new(File::open(&args[0])?), |_, path| { + output.push(path) + })?; + let tree = output.finish()?; + println!( + "policy={} paths={} elapsed={:.3}s file_bytes={}\n{:?}", + args[2], + stats.path_count, + start.elapsed().as_secs_f64(), + tree.get_data().len(), + tree.counters() + ); + Ok(()) +} diff --git a/src/arena_compact.rs b/src/arena_compact.rs index 70397698..d07b62bd 100644 --- a/src/arena_compact.rs +++ b/src/arena_compact.rs @@ -1491,8 +1491,57 @@ use std::fs::{File, OpenOptions}; pub struct FileDumper { buf_writer: BufWriter, - line_buf: Vec, - line_map: HashMap::, + line_cache_entries: usize, +} + +impl FileDumper { + fn line_matches(&mut self, offset: u64, end: u64, path: &[u8]) -> std::io::Result { + let mut header = [0; MAX_VARINT_SIZE]; + let header_len = push_varint_u64(&mut &mut header[..], path.len() as u64)?; + let total = header_len + path.len(); + if total as u64 > end - offset { return Ok(false); } + let buffered_start = end - self.buf_writer.buffer().len() as u64; + if offset >= buffered_start { + let start = (offset - buffered_start) as usize; + let stored = &self.buf_writer.buffer()[start..start + total]; + return Ok(stored[..header_len] == header[..header_len] && stored[header_len..] == *path); + } + let mut bytes = [0; 8192]; + for start in (0..total).step_by(bytes.len()) { + let len = (total - start).min(bytes.len()); + let position = offset + start as u64; + let on_disk = buffered_start.saturating_sub(position).min(len as u64) as usize; + if on_disk != 0 { + #[cfg(unix)] + { + use std::os::unix::fs::FileExt; + self.buf_writer.get_ref().read_exact_at(&mut bytes[..on_disk], position)?; + } + #[cfg(not(unix))] + { + // Preserve the append position without flushing the buffered suffix. + let file = self.buf_writer.get_mut(); + file.seek(SeekFrom::Start(position))?; + let read = std::io::Read::read_exact(file, &mut bytes[..on_disk]); + let restored = file.seek(SeekFrom::Start(buffered_start)); + read?; + restored?; + } + } + if on_disk < len { + let buffered_offset = (position + on_disk as u64 - buffered_start) as usize; + bytes[on_disk..len].copy_from_slice( + &self.buf_writer.buffer()[buffered_offset..buffered_offset + len - on_disk]); + } + let prefix = header_len.saturating_sub(start).min(len); + if prefix != 0 && bytes[..prefix] != header[start..start + prefix] { return Ok(false); } + if prefix < len { + let path_start = start + prefix - header_len; + if bytes[prefix..len] != path[path_start..path_start + len - prefix] { return Ok(false); } + } + } + Ok(true) + } } impl Write for FileDumper { @@ -1520,8 +1569,7 @@ impl ArenaCompactTree { let buf_writer = BufWriter::with_capacity(DUMPER_BUFFER_SIZE, file); let storage = FileDumper { buf_writer, - line_buf: Default::default(), - line_map: Default::default(), + line_cache_entries: usize::MAX, }; let act = ArenaCompactTree { storage, @@ -1550,16 +1598,17 @@ impl ArenaCompactTree { let mut hasher = self.hasher.clone(); hasher.write(path); let hash = hasher.finish(); - if let Some(&(start, len, prev)) = self.storage.line_map.get(&hash) { - let buf = &self.storage.line_buf[start..start+len]; - if buf == path { + if let Some(&prev) = self.line_map.get(&hash) { + if self.storage.line_matches(prev.0, self.position, path)? { self.counters.add_line_data_reuse(path.len()); return Ok(prev); } } + // Only fingerprints and offsets are cached; stored bytes verify collisions. + if self.line_map.len() >= self.storage.line_cache_entries { + self.line_map.clear(); + } let line_id = LineId(self.position); - let line_start = self.storage.line_buf.len(); - self.storage.line_buf.extend_from_slice(path); let lenlen = push_varint_u64( &mut self.storage, path.len() as u64 )? as u64; @@ -1567,7 +1616,9 @@ impl ArenaCompactTree { self.storage.write_all(path)?; self.position += path.len() as u64; self.counters.add_line_data(lenlen as usize + path.len()); - self.storage.line_map.insert(hash, (line_start, path.len(), line_id)); + if self.storage.line_cache_entries != 0 { + self.line_map.insert(hash, line_id); + } Ok(line_id) } } @@ -1792,10 +1843,21 @@ pub struct ACTOutputStream { } impl ACTOutputStream { - /// Create (or truncate) the file at `path` and start streaming a trie into it + /// Create (or truncate) the file at `path` and start streaming a trie into it. + /// The reuse cache holds at most 262144 fingerprint-to-offset entries. pub fn new(path: impl AsRef) -> Result { + Self::with_cache_limit(path, 262144) + } + + /// Create a stream with at most `entries` cached fingerprints and ACT offsets. + /// Line bytes stay in the output file or its write buffer. Hash-table capacity + /// and metadata are additional. A zero entry limit disables caching. + /// Eviction may increase file size but never changes the tree's contents. + pub fn with_cache_limit(path: impl AsRef, entries: usize) -> Result { + let mut act = ArenaCompactTree::::open(path)?; + act.storage.line_cache_entries = entries; Ok(ACTOutputStream { - act: ArenaCompactTree::::open(path)?, + act, stack: Vec::from([StreamFrame::empty()]), prev_path: Vec::new(), count: 0, diff --git a/tests/act_stream_cache.rs b/tests/act_stream_cache.rs new file mode 100644 index 00000000..49d667b9 --- /dev/null +++ b/tests/act_stream_cache.rs @@ -0,0 +1,56 @@ +#![cfg(feature = "arena_compact")] +use pathmap::arena_compact::{ACTOutputStream, ArenaCompactTree}; + +#[test] +fn bounded_cache_eviction_preserves_disk_queries() -> std::io::Result<()> { + let paths: Vec<_> = (0u32..1000).map(|n| { + let mut path = n.to_be_bytes().to_vec(); + path.extend_from_slice(format!("unique suffix {n:08}").as_bytes()); + path + }).collect(); + // Disabled cache and several bounded entry counts. + for entries in [0, 1, 2, 4096] { + let dir = tempfile::tempdir()?; + let file = dir.path().join("bounded.act"); + let mut out = ACTOutputStream::with_cache_limit(&file, entries)?; + for path in &paths { out.push(path)?; } + drop(out.finish()?); + let tree = ArenaCompactTree::open_mmap(&file)?; + for path in &paths { assert_eq!(tree.get_val_at(path), Some(0)); } + assert_eq!(tree.get_val_at(b"absent"), None); + assert_eq!(tree.iter().map(|(p, _)| p).collect::>(), paths); + } + Ok(()) +} + +#[test] +fn cached_offsets_reuse_buffered_and_flushed_lines() -> std::io::Result<()> { + let paths: Vec<_> = (0u32..24000).map(|n| { + let suffix = if n >= 23900 { (n - 23900) / 2 } else { n / 2 }; + let mut path = n.to_be_bytes().to_vec(); + path.extend_from_slice(&suffix.to_be_bytes()); + path.extend_from_slice(&[42; 508]); + path + }).collect(); + let dir = tempfile::tempdir()?; + let mut sizes = vec![]; + for entries in [0, 2, 30000] { + let file = dir.path().join(format!("cache-{entries}.act")); + let mut out = ACTOutputStream::with_cache_limit(&file, entries)?; + for path in &paths { out.push(path)?; } + let tree = out.finish()?; + sizes.push(std::fs::metadata(&file)?.len()); + // Exceeds the 4 MiB write buffer; the final paths reuse early on-disk lines. + assert!(sizes.last().unwrap() > &(4 * 1024 * 1024)); + assert_eq!(tree.iter().map(|(p, _)| p).collect::>(), paths); + for path in paths.iter().step_by(257) { assert_eq!(tree.get_val_at(path), Some(0)); } + assert_eq!(tree.get_val_at(b"absent"), None); + } + assert!(sizes[1] < sizes[0]); + assert!(sizes[2] < sizes[1]); + // The ordinary zipper dumper uses the same file-backed line cache. + let map = pathmap::PathMap::from_iter(paths.iter().map(|path| (path, ()))); + let tree = ArenaCompactTree::dump_from_zipper(map.read_zipper(), |_| 0, dir.path().join("zipper.act"))?; + assert_eq!(tree.iter().map(|(p, _)| p).collect::>(), paths); + Ok(()) +}