Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions examples/act_line_cache.rs
Original file line number Diff line number Diff line change
@@ -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(())
}
86 changes: 74 additions & 12 deletions src/arena_compact.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1491,8 +1491,57 @@ use std::fs::{File, OpenOptions};

pub struct FileDumper {
buf_writer: BufWriter<File>,
line_buf: Vec<u8>,
line_map: HashMap::<u64, (usize, usize, LineId)>,
line_cache_entries: usize,
}

impl FileDumper {
fn line_matches(&mut self, offset: u64, end: u64, path: &[u8]) -> std::io::Result<bool> {
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 {
Expand Down Expand Up @@ -1520,8 +1569,7 @@ impl ArenaCompactTree<FileDumper> {
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,
Expand Down Expand Up @@ -1550,24 +1598,27 @@ impl ArenaCompactTree<FileDumper> {
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;
self.position += lenlen;
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)
}
}
Expand Down Expand Up @@ -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<Path>) -> Result<Self, std::io::Error> {
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<Path>, entries: usize) -> Result<Self, std::io::Error> {
let mut act = ArenaCompactTree::<FileDumper>::open(path)?;
act.storage.line_cache_entries = entries;
Ok(ACTOutputStream {
act: ArenaCompactTree::<FileDumper>::open(path)?,
act,
stack: Vec::from([StreamFrame::empty()]),
prev_path: Vec::new(),
count: 0,
Expand Down
56 changes: 56 additions & 0 deletions tests/act_stream_cache.rs
Original file line number Diff line number Diff line change
@@ -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::<Vec<_>>(), 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::<Vec<_>>(), 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::<Vec<_>>(), paths);
Ok(())
}
Loading