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
4 changes: 4 additions & 0 deletions src/paths_serialization.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,10 @@ mod paths_serialization_nightly;
#[cfg(feature="nightly")]
pub use paths_serialization_nightly::*;

#[path="paths_serialization_sort.rs"]
mod paths_serialization_sort;
pub use paths_serialization_sort::sort_paths;

/// Statistics from a `serialize` operation
#[derive(Debug, Clone, Copy)]
pub struct SerializationStats {
Expand Down
235 changes: 235 additions & 0 deletions src/paths_serialization_sort.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,235 @@
//! External bytewise sorting for unordered `.paths` streams.
use super::{SerializationStats, for_each_deserialized_path, serialize_paths_from_funcs};
use std::{
fs::{self, File},
io::{self, BufReader, BufWriter, Read, Write},
path::{Path, PathBuf},
};

const BUFFER: usize = 16 * 1024;

fn new_run(temp_dir: &Path, next_run: &mut u64) -> io::Result<(PathBuf, File)> {
let path = temp_dir.join(format!("pathmap-sort-{next_run}"));
let file = File::create_new(&path)?;
*next_run += 1;
Ok((path, file))
}

// Runs are uncompressed length-prefixed paths, avoiding repeated compression
// during merges. These functions share no codec state or in-memory path index.
fn read_path(source: &mut impl Read, path: &mut Vec<u8>, limit: usize) -> io::Result<bool> {
let mut length = [0; 4];
loop {
match source.read(&mut length[..1]) {
Ok(0) => return Ok(false),
Ok(_) => break,
Err(e) if e.kind() == io::ErrorKind::Interrupted => continue,
Err(e) => return Err(e),
}
}
source.read_exact(&mut length[1..])?;
let length = u32::from_le_bytes(length) as usize;
if length > limit {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"path exceeds sort memory budget",
));
}
path.resize(length, 0);
source.read_exact(path)?;
Ok(true)
}
fn write_path(target: &mut impl Write, path: &[u8]) -> io::Result<()> {
let length = u32::try_from(path.len())
.map_err(|_| io::Error::other("path exceeds .paths length limit"))?;
target.write_all(&length.to_le_bytes())?;
target.write_all(path)
}

struct Chunk {
bytes: Vec<u8>,
entries: Vec<(usize, usize)>,
byte_limit: usize,
entry_limit: usize,
}
impl Chunk {
fn new(memory: usize) -> Self {
let byte_limit = memory / 3;
let entry_limit = memory / 6 / std::mem::size_of::<(usize, usize)>();
Self {
bytes: Vec::with_capacity(byte_limit),
entries: Vec::with_capacity(entry_limit),
byte_limit,
entry_limit,
}
}
fn fits(&self, path: &[u8]) -> bool {
self.bytes.len() + path.len() <= self.byte_limit && self.entries.len() < self.entry_limit
}
fn push(&mut self, path: &[u8]) {
self.entries.push((self.bytes.len(), path.len()));
self.bytes.extend_from_slice(path);
}
fn spill(&mut self, temp_dir: &Path, next_run: &mut u64) -> io::Result<PathBuf> {
let bytes = &self.bytes;
self.entries
.sort_unstable_by(|&(a, al), &(b, bl)| bytes[a..a + al].cmp(&bytes[b..b + bl]));
let (path, file) = new_run(temp_dir, next_run)?;
let mut writer = BufWriter::with_capacity(BUFFER, file);
let mut previous = None;
for &(start, len) in &self.entries {
let path = &bytes[start..start + len];
if previous != Some(path) {
write_path(&mut writer, path)?;
}
previous = Some(path);
}
writer.flush()?;
drop(writer);
self.entries.clear();
self.bytes.clear();
Ok(path)
}
}

fn merge(
left: PathBuf,
right: PathBuf,
temp_dir: &Path,
max_path: usize,
next_run: &mut u64,
) -> io::Result<PathBuf> {
let mut a = BufReader::with_capacity(BUFFER, File::open(&left)?);
let mut b = BufReader::with_capacity(BUFFER, File::open(&right)?);
let mut ap = Vec::new();
let mut bp = Vec::new();
let (path, output) = new_run(temp_dir, next_run)?;
let mut writer = BufWriter::with_capacity(BUFFER, output);
let mut has_a = read_path(&mut a, &mut ap, max_path)?;
let mut has_b = read_path(&mut b, &mut bp, max_path)?;
while has_a || has_b {
let order = match (has_a, has_b) {
(true, true) => ap.cmp(&bp),
(true, false) => std::cmp::Ordering::Less,
_ => std::cmp::Ordering::Greater,
};
write_path(&mut writer, if order.is_le() { &ap } else { &bp })?;
if order.is_le() {
has_a = read_path(&mut a, &mut ap, max_path)?;
}
if order.is_ge() {
has_b = read_path(&mut b, &mut bp, max_path)?;
}
}
writer.flush()?;
drop(writer);
drop((a, b));
fs::remove_file(left)?;
fs::remove_file(right)?;
Ok(path)
}

// Binary carry merging: at most one run per level and no unbounded run list.
fn add_run(
mut run: PathBuf,
levels: &mut [Option<PathBuf>; 64],
temp_dir: &Path,
max_path: usize,
next_run: &mut u64,
) -> io::Result<()> {
for slot in levels {
match slot.take() {
None => {
*slot = Some(run);
return Ok(());
}
Some(other) => run = merge(other, run, temp_dir, max_path, next_run)?,
}
}
Err(io::Error::other("too many sort runs"))
}

/// Sort and deduplicate an unordered `.paths` stream without building a PathMap.
///
/// `memory_bytes` (at least 1 MiB) budgets the sorting arena, indices and merge
/// buffers; codec/runtime overhead and the decoder's current path are additional.
/// Individual paths must fit in 1/16 of the budget (checked after decoding).
/// The caller supplies an existing scratch directory in `temp_dir`, exclusive to
/// this sort. Run files are created with `create_new` and removed on success or
/// error; the caller owns creation and deletion of the directory. Allow up to
/// twice the uncompressed input size for runs.
/// The output uses the usual `.paths` format, in strictly increasing byte order.
pub fn sort_paths<R: Read, W: Write>(
source: R,
target: &mut W,
memory_bytes: usize,
temp_dir: &Path,
) -> io::Result<SerializationStats> {
if memory_bytes < 1024 * 1024 {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"sort memory budget must be at least 1 MiB",
));
}
let mut next_run = 0;
let result = (|| {
let max_path = memory_bytes / 16;
let mut chunk = Chunk::new(memory_bytes);
let mut levels = std::array::from_fn(|_| None);
for_each_deserialized_path(source, |_, path| {
if path.len() > max_path {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"path exceeds sort memory budget",
));
}
if !chunk.fits(path) {
add_run(
chunk.spill(temp_dir, &mut next_run)?,
&mut levels,
temp_dir,
max_path,
&mut next_run,
)?;
}
chunk.push(path);
Ok(())
})?;
if !chunk.entries.is_empty() {
add_run(
chunk.spill(temp_dir, &mut next_run)?,
&mut levels,
temp_dir,
max_path,
&mut next_run,
)?;
}
drop(chunk);
let mut final_run = None;
for run in levels.into_iter().flatten() {
final_run = Some(match final_run {
None => run,
Some(other) => merge(other, run, temp_dir, max_path, &mut next_run)?,
});
}
let mut source: Box<dyn Read> = match &final_run {
Some(path) => Box::new(BufReader::with_capacity(BUFFER, File::open(path)?)),
None => Box::new(io::empty()),
};
serialize_paths_from_funcs(
target,
&mut Vec::new(),
|path| read_path(&mut source, path, max_path),
|path| Some(path),
)
})();
let mut cleanup = Ok(());
for index in 0..next_run {
if let Err(error) = fs::remove_file(temp_dir.join(format!("pathmap-sort-{index}"))) {
if error.kind() != io::ErrorKind::NotFound {
cleanup = Err(error);
}
}
}
result.and_then(|stats| cleanup.map(|()| stats))
}
102 changes: 102 additions & 0 deletions tests/paths_sort.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
#![cfg(feature = "serialization")]
use pathmap::paths_serialization::{
for_each_deserialized_path, serialize_paths_from_funcs, sort_paths,
};
use std::{
fs::{self, File},
io,
path::Path,
};
fn write_paths(path: &Path, paths: &[Vec<u8>]) {
let mut source = (0usize, paths);
serialize_paths_from_funcs(
&mut File::create(path).unwrap(),
&mut source,
|s| {
s.0 += 1;
Ok(s.0 <= s.1.len())
},
|s| Some(s.1[s.0 - 1].as_slice()),
)
.unwrap();
}
fn read_paths(path: &Path) -> Vec<Vec<u8>> {
let mut paths = vec![];
for_each_deserialized_path(File::open(path).unwrap(), |_, p| {
paths.push(p.to_vec());
Ok(())
})
.unwrap();
paths
}
fn sort_file(input: &Path, output: &Path, memory: usize, temp: &Path) -> io::Result<usize> {
let mut target = tempfile::NamedTempFile::new_in(temp)?;
let count = sort_paths(File::open(input)?, &mut target, memory, temp)?.path_count;
target.persist(output).map_err(|e| e.error)?;
Ok(count)
}
#[test]
fn external_sort_many_runs_matches_bytewise_set() {
let scratch = tempfile::tempdir().unwrap();
let input = scratch.path().join("input.upaths");
let output = scratch.path().join("output.paths");
let mut paths = vec![vec![], vec![0], vec![0, 0], vec![255], vec![]];
for n in (0u32..35000).rev() {
let mut path = n.to_be_bytes().to_vec();
path.extend_from_slice(&[0, 255, 0]);
path.extend(std::iter::repeat_n((n % 251) as u8, (n % 123) as usize));
paths.push(path.clone());
if n % 3 == 0 {
paths.push(path);
}
}
write_paths(&input, &paths);
paths.sort();
paths.dedup();
assert_eq!(
sort_file(&input, &output, 1 << 20, scratch.path()).unwrap(),
paths.len()
);
assert_eq!(read_paths(&output), paths);
// A sort-budget error after multiple spills must clean every temporary run.
let before = fs::read(&output).unwrap();
paths.push(vec![0; (1 << 16) + 1]);
write_paths(&input, &paths);
assert!(sort_file(&input, &output, 1 << 20, scratch.path()).is_err());
assert_eq!(fs::read(&output).unwrap(), before);
assert_eq!(fs::read_dir(scratch.path()).unwrap().count(), 2);
}

#[test]
fn sort_empty_duplicates_and_memory_limit() {
let scratch = tempfile::tempdir().unwrap();
let input = scratch.path().join("input.upaths");
let output = scratch.path().join("output.paths");
for paths in [vec![], vec![vec![]; 3], vec![b"same".to_vec(); 25000]] {
write_paths(&input, &paths);
let mut expected = paths;
expected.sort();
expected.dedup();
sort_file(&input, &output, 1 << 20, scratch.path()).unwrap();
assert_eq!(read_paths(&output), expected);
}
write_paths(&input, &[vec![0; (1 << 16) + 1]]);
assert!(sort_file(&input, &output, 1 << 20, scratch.path()).is_err());
assert_eq!(fs::read_dir(scratch.path()).unwrap().count(), 2);
}

#[test]
fn caller_scratch_files_survive_run_name_collisions() {
let scratch = tempfile::tempdir().unwrap();
let input = scratch.path().join("input.upaths");
let output = scratch.path().join("output.paths");
let existing = scratch.path().join("pathmap-sort-1");
fs::write(&existing, b"caller-owned file").unwrap();
fs::write(&output, b"existing destination").unwrap();
write_paths(&input, &vec![vec![42; 64]; 20000]);
let error = sort_file(&input, &output, 1 << 20, scratch.path()).unwrap_err();
assert_eq!(error.kind(), io::ErrorKind::AlreadyExists);
assert_eq!(fs::read(&existing).unwrap(), b"caller-owned file");
assert_eq!(fs::read(&output).unwrap(), b"existing destination");
assert_eq!(fs::read_dir(scratch.path()).unwrap().count(), 3);
}