feat(walker): implemented parallel filesystem traversal with streaming grep
- Introduced parallel streaming grep with windowed result processing to enhance search performance and memory management. - Implemented stateful parallel file walking with buffer pooling and directory entry record caching to minimize memory allocations. - Optimized ignore state derivation by using directory entry names instead of stat probes. - Added comprehensive unit and performance tests covering parallel traversal correctness, early termination, and streaming behavior.
This commit is contained in:
+1
-1
@@ -238,7 +238,7 @@ smallvec = { version = "1.15.1", features = [
|
||||
xxhash-rust = { version = "0.8", features = ["xxh64"] }
|
||||
|
||||
# ──────────────────────────────────────────────────────────────────────────────
|
||||
# Memory Management & Allocators
|
||||
# Memory Mapping
|
||||
# ──────────────────────────────────────────────────────────────────────────────
|
||||
memmap2 = "0.9"
|
||||
|
||||
|
||||
+665
-125
@@ -9,10 +9,11 @@
|
||||
|
||||
use std::{
|
||||
borrow::Cow,
|
||||
cell::RefCell,
|
||||
fs::File,
|
||||
io::{self, Read},
|
||||
path::{Path, PathBuf},
|
||||
sync::atomic::{AtomicU64, AtomicUsize, Ordering},
|
||||
sync::atomic::{AtomicU64, Ordering},
|
||||
};
|
||||
|
||||
use grep_matcher::Matcher;
|
||||
@@ -34,23 +35,6 @@ use crate::{glob_util, iofs, task};
|
||||
const MAX_FILE_BYTES: u64 = 4 * 1024 * 1024;
|
||||
const SMALL_FILE_READ_BYTES: u64 = 128 * 1024;
|
||||
|
||||
static ACTIVE_STREAMING_GREPS: AtomicUsize = AtomicUsize::new(0);
|
||||
|
||||
struct ActiveStreamingGrep;
|
||||
|
||||
impl ActiveStreamingGrep {
|
||||
fn enter() -> (Self, usize) {
|
||||
let active = ACTIVE_STREAMING_GREPS.fetch_add(1, Ordering::Relaxed) + 1;
|
||||
(Self, active)
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for ActiveStreamingGrep {
|
||||
fn drop(&mut self) {
|
||||
ACTIVE_STREAMING_GREPS.fetch_sub(1, Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
/// Output mode for [`search`] and [`grep`] (string values match JS callers).
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
#[napi(string_enum)]
|
||||
@@ -539,7 +523,7 @@ fn resolve_context(
|
||||
// Search engine
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
#[derive(Clone, Copy, PartialEq, Eq)]
|
||||
struct SearchParams {
|
||||
context_before: u32,
|
||||
context_after: u32,
|
||||
@@ -597,6 +581,24 @@ fn build_searcher_for_params(params: SearchParams) -> Searcher {
|
||||
collect_content,
|
||||
)
|
||||
}
|
||||
std::thread_local! {
|
||||
static PARALLEL_GREP_SEARCHER: RefCell<Option<(SearchParams, Searcher)>> =
|
||||
const { RefCell::new(None) };
|
||||
}
|
||||
|
||||
fn with_parallel_grep_searcher<T>(
|
||||
params: SearchParams,
|
||||
search: impl FnOnce(&mut Searcher) -> T,
|
||||
) -> T {
|
||||
PARALLEL_GREP_SEARCHER.with(|cell| {
|
||||
let mut cached = cell.borrow_mut();
|
||||
if !matches!(cached.as_ref(), Some((cached_params, _)) if *cached_params == params) {
|
||||
*cached = Some((params, build_searcher_for_params(params)));
|
||||
}
|
||||
let (_, searcher) = cached.as_mut().expect("parallel grep searcher initialized");
|
||||
search(searcher)
|
||||
})
|
||||
}
|
||||
|
||||
fn build_searcher(
|
||||
context_before: u32,
|
||||
@@ -755,23 +757,23 @@ const fn empty_search_result(error: Option<String>) -> SearchResult {
|
||||
}
|
||||
|
||||
/// Internal configuration for grep, extracted from options.
|
||||
struct GrepConfig {
|
||||
pattern: String,
|
||||
path: String,
|
||||
glob: Option<String>,
|
||||
type_filter: Option<String>,
|
||||
ignore_case: Option<bool>,
|
||||
multiline: Option<bool>,
|
||||
hidden: Option<bool>,
|
||||
gitignore: Option<bool>,
|
||||
max_count: Option<u32>,
|
||||
offset: Option<u32>,
|
||||
context_before: Option<u32>,
|
||||
context_after: Option<u32>,
|
||||
context: Option<u32>,
|
||||
max_columns: Option<u32>,
|
||||
mode: Option<GrepOutputMode>,
|
||||
max_count_per_file: Option<u32>,
|
||||
pub(crate) struct GrepConfig {
|
||||
pub(crate) pattern: String,
|
||||
pub(crate) path: String,
|
||||
pub(crate) glob: Option<String>,
|
||||
pub(crate) type_filter: Option<String>,
|
||||
pub(crate) ignore_case: Option<bool>,
|
||||
pub(crate) multiline: Option<bool>,
|
||||
pub(crate) hidden: Option<bool>,
|
||||
pub(crate) gitignore: Option<bool>,
|
||||
pub(crate) max_count: Option<u32>,
|
||||
pub(crate) offset: Option<u32>,
|
||||
pub(crate) context_before: Option<u32>,
|
||||
pub(crate) context_after: Option<u32>,
|
||||
pub(crate) context: Option<u32>,
|
||||
pub(crate) max_columns: Option<u32>,
|
||||
pub(crate) mode: Option<GrepOutputMode>,
|
||||
pub(crate) max_count_per_file: Option<u32>,
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -956,10 +958,19 @@ fn build_regex_matcher(
|
||||
ignore_case: bool,
|
||||
multiline: bool,
|
||||
) -> std::result::Result<grep_regex::RegexMatcher, grep_regex::Error> {
|
||||
RegexMatcherBuilder::new()
|
||||
.case_insensitive(ignore_case)
|
||||
.multi_line(multiline)
|
||||
.build(pattern)
|
||||
let build = |line_terminated| {
|
||||
let mut builder = RegexMatcherBuilder::new();
|
||||
builder.case_insensitive(ignore_case).multi_line(multiline);
|
||||
if line_terminated {
|
||||
builder.line_terminator(Some(b'\n'));
|
||||
}
|
||||
builder.build(pattern)
|
||||
};
|
||||
|
||||
if !multiline && let Ok(matcher) = build(true) {
|
||||
return Ok(matcher);
|
||||
}
|
||||
build(false)
|
||||
}
|
||||
|
||||
fn build_matcher(
|
||||
@@ -997,6 +1008,7 @@ fn build_matcher(
|
||||
// File / directory search orchestration
|
||||
// ---------------------------------------------------------------------------
|
||||
const ORDERED_STREAMING_STOP_MAX_COUNT: u64 = 64;
|
||||
const GREP_STREAM_WINDOW: usize = 512;
|
||||
|
||||
fn per_file_params(params: SearchParams) -> SearchParams {
|
||||
let file_limit = match params.mode {
|
||||
@@ -1140,7 +1152,6 @@ struct PassState {
|
||||
skipped_oversized: AtomicU64,
|
||||
emitted: AtomicU64,
|
||||
}
|
||||
|
||||
/// Memory-map the first [`MAX_FILE_BYTES`] of a file for searching.
|
||||
///
|
||||
/// Used by the deferred oversized pass: files larger than the cap are searched
|
||||
@@ -1271,19 +1282,13 @@ fn run_pass(
|
||||
ct: &task::CancelToken,
|
||||
) -> Result<Vec<FileSearchResult>> {
|
||||
if parallel_allowed && pi_walker::should_parallelize(candidates.len()) {
|
||||
pi_walker::execute_candidates(candidates, |file| {
|
||||
let mut searcher = build_searcher_for_params(file_params);
|
||||
handle_file(
|
||||
file,
|
||||
&mut searcher,
|
||||
matcher,
|
||||
file_params,
|
||||
policy,
|
||||
stop_after_matches,
|
||||
state,
|
||||
ct,
|
||||
)
|
||||
})?;
|
||||
pi_walker::execute_candidates_init(
|
||||
candidates,
|
||||
|| build_searcher_for_params(file_params),
|
||||
|searcher, file| {
|
||||
handle_file(file, searcher, matcher, file_params, policy, stop_after_matches, state, ct)
|
||||
},
|
||||
)?;
|
||||
} else {
|
||||
let mut searcher = build_searcher_for_params(file_params);
|
||||
ct.heartbeat()?;
|
||||
@@ -1409,7 +1414,7 @@ fn run_sequential_grep(
|
||||
clippy::fn_params_excessive_bools,
|
||||
reason = "matches options structure of underlying walk candidates collector"
|
||||
)]
|
||||
fn try_run_native_walk_grep(
|
||||
fn run_parallel_streaming_grep(
|
||||
search_path: &Path,
|
||||
matcher: &grep_regex::RegexMatcher,
|
||||
glob: Option<&str>,
|
||||
@@ -1419,35 +1424,179 @@ fn try_run_native_walk_grep(
|
||||
use_gitignore: bool,
|
||||
skip_node_modules: bool,
|
||||
ct: &task::CancelToken,
|
||||
stop_after_matches: Option<u64>,
|
||||
parallel_search: bool,
|
||||
) -> Result<Option<(Vec<FileSearchResult>, u64, u64)>> {
|
||||
let requires_path_order = stop_after_matches.is_some() || params.offset != 0;
|
||||
let Some(candidates) = collect_grep_candidates(
|
||||
) -> Result<(Vec<FileSearchResult>, u64, u64)> {
|
||||
let request = build_grep_walk_request(
|
||||
search_path,
|
||||
glob,
|
||||
type_filter,
|
||||
include_hidden,
|
||||
use_gitignore,
|
||||
skip_node_modules,
|
||||
if requires_path_order {
|
||||
pi_walker::WalkOrder::Path
|
||||
} else {
|
||||
pi_walker::WalkOrder::Unordered
|
||||
},
|
||||
ct,
|
||||
)?
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
Ok(Some(process_candidates(
|
||||
candidates,
|
||||
matcher,
|
||||
params,
|
||||
parallel_search,
|
||||
stop_after_matches,
|
||||
ct,
|
||||
)?))
|
||||
pi_walker::WalkOrder::Unordered,
|
||||
)?;
|
||||
let file_params = per_file_params(params);
|
||||
let state = PassState::default();
|
||||
|
||||
request
|
||||
.for_each_file_candidate_parallel(
|
||||
|file| {
|
||||
if let Some(filter) = type_filter
|
||||
&& !matches_type_filter_str(&file.relative, filter)
|
||||
{
|
||||
return Ok(pi_walker::ParallelWalkControl::Continue);
|
||||
}
|
||||
with_parallel_grep_searcher(file_params, |searcher| {
|
||||
handle_file(file, searcher, matcher, file_params, ReadPolicy::Full, None, &state, ct)
|
||||
})?;
|
||||
Ok(pi_walker::ParallelWalkControl::Continue)
|
||||
},
|
||||
|| ct.heartbeat(),
|
||||
)
|
||||
.map_err(iofs::map_walker_error)?;
|
||||
|
||||
let mut results = std::mem::take(&mut *state.results.lock());
|
||||
results.sort_unstable_by(|a, b| a.relative_path.cmp(&b.relative_path));
|
||||
|
||||
let deferred = std::mem::take(&mut *state.deferred.lock());
|
||||
if !deferred.is_empty() {
|
||||
let oversized =
|
||||
run_pass(&deferred, matcher, file_params, ReadPolicy::Prefix, true, None, &state, ct)?;
|
||||
results.extend(oversized);
|
||||
}
|
||||
|
||||
Ok((
|
||||
results,
|
||||
state.skipped_oversized.load(Ordering::Relaxed),
|
||||
state.files_searched.load(Ordering::Relaxed),
|
||||
))
|
||||
}
|
||||
|
||||
fn emitted_content_matches(results: &[FileSearchResult]) -> u64 {
|
||||
results.iter().fold(0, |total, result| {
|
||||
total.saturating_add(u64::try_from(result.matches.len()).unwrap_or(u64::MAX))
|
||||
})
|
||||
}
|
||||
|
||||
fn flush_stream_window(
|
||||
window: &mut Vec<pi_walker::FileCandidate>,
|
||||
results: &mut Vec<FileSearchResult>,
|
||||
matcher: &grep_regex::RegexMatcher,
|
||||
file_params: SearchParams,
|
||||
state: &PassState,
|
||||
ct: &task::CancelToken,
|
||||
stop_after_matches: u64,
|
||||
) -> Result<bool> {
|
||||
if window.is_empty() {
|
||||
return Ok(state.emitted.load(Ordering::Relaxed) >= stop_after_matches);
|
||||
}
|
||||
let window_results =
|
||||
run_pass(window.as_slice(), matcher, file_params, ReadPolicy::Full, true, None, state, ct)?;
|
||||
let emitted = emitted_content_matches(&window_results);
|
||||
let total_emitted = state
|
||||
.emitted
|
||||
.fetch_add(emitted, Ordering::Relaxed)
|
||||
.saturating_add(emitted);
|
||||
results.extend(window_results);
|
||||
window.clear();
|
||||
Ok(total_emitted >= stop_after_matches)
|
||||
}
|
||||
|
||||
#[allow(
|
||||
clippy::fn_params_excessive_bools,
|
||||
reason = "matches options structure of underlying walk candidates collector"
|
||||
)]
|
||||
fn run_windowed_streaming_grep(
|
||||
search_path: &Path,
|
||||
matcher: &grep_regex::RegexMatcher,
|
||||
glob: Option<&str>,
|
||||
type_filter: Option<&TypeFilter>,
|
||||
params: SearchParams,
|
||||
include_hidden: bool,
|
||||
use_gitignore: bool,
|
||||
skip_node_modules: bool,
|
||||
ct: &task::CancelToken,
|
||||
stop_after_matches: u64,
|
||||
) -> Result<(Vec<FileSearchResult>, u64, u64)> {
|
||||
let request = build_grep_walk_request(
|
||||
search_path,
|
||||
glob,
|
||||
include_hidden,
|
||||
use_gitignore,
|
||||
skip_node_modules,
|
||||
pi_walker::WalkOrder::Path,
|
||||
)?;
|
||||
let file_params = per_file_params(params);
|
||||
let state = PassState::default();
|
||||
let mut window = Vec::with_capacity(GREP_STREAM_WINDOW);
|
||||
let mut results = Vec::new();
|
||||
|
||||
request
|
||||
.for_each_entry_with_heartbeat(
|
||||
|| ct.heartbeat(),
|
||||
|entry| {
|
||||
if let Some(filter) = type_filter
|
||||
&& !matches_type_filter_str(entry.relative_path, filter)
|
||||
{
|
||||
return Ok(pi_walker::WalkDecision::Include);
|
||||
}
|
||||
let relative = entry.relative_path.to_owned();
|
||||
window.push(pi_walker::FileCandidate {
|
||||
path: entry.absolute_path.into_owned(),
|
||||
relative,
|
||||
mtime: entry.mtime,
|
||||
size: entry.size,
|
||||
});
|
||||
if window.len() == GREP_STREAM_WINDOW
|
||||
&& flush_stream_window(
|
||||
&mut window,
|
||||
&mut results,
|
||||
matcher,
|
||||
file_params,
|
||||
&state,
|
||||
ct,
|
||||
stop_after_matches,
|
||||
)? {
|
||||
return Ok(pi_walker::WalkDecision::Stop);
|
||||
}
|
||||
Ok(pi_walker::WalkDecision::Include)
|
||||
},
|
||||
|_| Ok(pi_walker::WalkDecision::Include),
|
||||
)
|
||||
.map_err(iofs::map_walker_error)?;
|
||||
|
||||
if state.emitted.load(Ordering::Relaxed) < stop_after_matches {
|
||||
flush_stream_window(
|
||||
&mut window,
|
||||
&mut results,
|
||||
matcher,
|
||||
file_params,
|
||||
&state,
|
||||
ct,
|
||||
stop_after_matches,
|
||||
)?;
|
||||
}
|
||||
|
||||
let mut deferred = std::mem::take(&mut *state.deferred.lock());
|
||||
let limit_satisfied = state.emitted.load(Ordering::Relaxed) >= stop_after_matches;
|
||||
if !deferred.is_empty() && !limit_satisfied {
|
||||
deferred.sort_unstable_by(|a, b| a.relative.cmp(&b.relative));
|
||||
let oversized = run_pass(
|
||||
&deferred,
|
||||
matcher,
|
||||
file_params,
|
||||
ReadPolicy::Prefix,
|
||||
false,
|
||||
Some(stop_after_matches),
|
||||
&state,
|
||||
ct,
|
||||
)?;
|
||||
results.extend(oversized);
|
||||
}
|
||||
|
||||
Ok((
|
||||
results,
|
||||
state.skipped_oversized.load(Ordering::Relaxed),
|
||||
state.files_searched.load(Ordering::Relaxed),
|
||||
))
|
||||
}
|
||||
|
||||
fn run_streaming_grep(
|
||||
@@ -1461,42 +1610,46 @@ fn run_streaming_grep(
|
||||
skip_node_modules: bool,
|
||||
ct: &task::CancelToken,
|
||||
) -> Result<(Vec<FileSearchResult>, u64, u64)> {
|
||||
let (_active_guard, active_greps) = ActiveStreamingGrep::enter();
|
||||
let workers = pi_walker::walk_workers();
|
||||
let stop_after_matches = streaming_stop_after(params);
|
||||
let small_budget = stop_after_matches.is_some_and(|max| max <= ORDERED_STREAMING_STOP_MAX_COUNT);
|
||||
// Sequential path: forced workers, contended pool, or a small first-page
|
||||
// budget where strict path-order matters. The parallel path below also
|
||||
// honors `stop_after_matches`, so larger budgets bound work without losing
|
||||
// parallelism.
|
||||
let use_parallel_fast_path = workers > 1 && !small_budget && active_greps == 1;
|
||||
if let Some(result) = try_run_native_walk_grep(
|
||||
search_path,
|
||||
matcher,
|
||||
glob,
|
||||
type_filter,
|
||||
params,
|
||||
include_hidden,
|
||||
use_gitignore,
|
||||
skip_node_modules,
|
||||
ct,
|
||||
stop_after_matches,
|
||||
use_parallel_fast_path,
|
||||
)? {
|
||||
return Ok(result);
|
||||
match stop_after_matches {
|
||||
None => run_parallel_streaming_grep(
|
||||
search_path,
|
||||
matcher,
|
||||
glob,
|
||||
type_filter,
|
||||
params,
|
||||
include_hidden,
|
||||
use_gitignore,
|
||||
skip_node_modules,
|
||||
ct,
|
||||
),
|
||||
Some(stop) if stop <= ORDERED_STREAMING_STOP_MAX_COUNT || pi_walker::walk_workers() <= 1 => {
|
||||
run_sequential_grep(
|
||||
search_path,
|
||||
matcher,
|
||||
glob,
|
||||
type_filter,
|
||||
params,
|
||||
include_hidden,
|
||||
use_gitignore,
|
||||
skip_node_modules,
|
||||
ct,
|
||||
Some(stop),
|
||||
)
|
||||
},
|
||||
Some(stop) => run_windowed_streaming_grep(
|
||||
search_path,
|
||||
matcher,
|
||||
glob,
|
||||
type_filter,
|
||||
params,
|
||||
include_hidden,
|
||||
use_gitignore,
|
||||
skip_node_modules,
|
||||
ct,
|
||||
stop,
|
||||
),
|
||||
}
|
||||
run_sequential_grep(
|
||||
search_path,
|
||||
matcher,
|
||||
glob,
|
||||
type_filter,
|
||||
params,
|
||||
include_hidden,
|
||||
use_gitignore,
|
||||
skip_node_modules,
|
||||
ct,
|
||||
stop_after_matches,
|
||||
)
|
||||
}
|
||||
|
||||
fn push_count_match(matches: &mut Vec<GrepMatch>, path: String, match_count: u64) {
|
||||
@@ -1668,7 +1821,7 @@ fn search_sync(content: &[u8], options: SearchOptions) -> SearchResult {
|
||||
}
|
||||
}
|
||||
|
||||
fn grep_sync(
|
||||
pub(crate) fn grep_sync(
|
||||
options: GrepConfig,
|
||||
on_match: Option<&ThreadsafeFunction<GrepMatch>>,
|
||||
ct: task::CancelToken,
|
||||
@@ -2188,7 +2341,6 @@ mod tests {
|
||||
assert!(matcher.is_match(b"foooo").unwrap());
|
||||
assert!(!matcher.is_match(b"bar").unwrap());
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn grep_directory_skips_fifo_entries() {
|
||||
@@ -2411,6 +2563,318 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn unlimited_params(mode: super::OutputMode, context: u32) -> super::SearchParams {
|
||||
super::SearchParams {
|
||||
context_before: context,
|
||||
context_after: context,
|
||||
max_columns: None,
|
||||
mode,
|
||||
max_count: None,
|
||||
max_count_per_file: None,
|
||||
offset: 0,
|
||||
multiline: false,
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
struct MatchSnapshot {
|
||||
line_number: u64,
|
||||
line: String,
|
||||
context_before: Vec<(u32, String)>,
|
||||
context_after: Vec<(u32, String)>,
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
struct FileSnapshot {
|
||||
relative_path: String,
|
||||
match_count: u64,
|
||||
limit_reached: bool,
|
||||
matches: Vec<MatchSnapshot>,
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn file_snapshots(results: &[super::FileSearchResult]) -> Vec<FileSnapshot> {
|
||||
results
|
||||
.iter()
|
||||
.map(|result| FileSnapshot {
|
||||
relative_path: result.relative_path.clone(),
|
||||
match_count: result.match_count,
|
||||
limit_reached: result.limit_reached,
|
||||
matches: result
|
||||
.matches
|
||||
.iter()
|
||||
.map(|matched| MatchSnapshot {
|
||||
line_number: matched.line_number,
|
||||
line: matched.line.clone(),
|
||||
context_before: matched
|
||||
.context_before
|
||||
.iter()
|
||||
.map(|line| (line.line_number, line.line.clone()))
|
||||
.collect(),
|
||||
context_after: matched
|
||||
.context_after
|
||||
.iter()
|
||||
.map(|line| (line.line_number, line.line.clone()))
|
||||
.collect(),
|
||||
})
|
||||
.collect(),
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
struct GrepMatchSnapshot {
|
||||
path: String,
|
||||
line_number: u32,
|
||||
line: String,
|
||||
match_count: Option<u32>,
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn grep_match_snapshots(matches: &[super::GrepMatch]) -> Vec<GrepMatchSnapshot> {
|
||||
matches
|
||||
.iter()
|
||||
.map(|matched| GrepMatchSnapshot {
|
||||
path: matched.path.clone(),
|
||||
line_number: matched.line_number,
|
||||
line: matched.line.clone(),
|
||||
match_count: matched.match_count,
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn populate_parallel_parity_tree(root: &Path) {
|
||||
fs::create_dir_all(root.join(".git")).expect("create repo marker");
|
||||
write_file(&root.join("dir_00/.gitignore"), "ignored_match.txt\n");
|
||||
write_file(
|
||||
&root.join("dir_00/ignored_match.txt"),
|
||||
"before ignored\nneedle ignored\nafter ignored\n",
|
||||
);
|
||||
|
||||
for index in 0..300 {
|
||||
let path =
|
||||
root
|
||||
.join(format!("dir_{:02}/nested_{:02}/file_{index:03}.txt", index % 12, index % 5,));
|
||||
let content = if index % 3 == 0 {
|
||||
format!("before {index}\nneedle {index}\nafter {index}\n")
|
||||
} else {
|
||||
format!("before {index}\nhaystack {index}\nafter {index}\n")
|
||||
};
|
||||
write_file(&path, &content);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn sequential_reference_result(root: &Path, params: super::SearchParams) -> super::GrepResult {
|
||||
let matcher = super::build_matcher("needle", false, false).expect("build test matcher");
|
||||
let (results, skipped_oversized, files_searched) = super::run_sequential_grep(
|
||||
root,
|
||||
&matcher,
|
||||
None,
|
||||
None,
|
||||
params,
|
||||
true,
|
||||
true,
|
||||
true,
|
||||
&task::CancelToken::default(),
|
||||
super::streaming_stop_after(params),
|
||||
)
|
||||
.expect("sequential grep should succeed");
|
||||
let (matches, total_matches, files_with_matches, files_searched, limit_reached) =
|
||||
super::aggregate_parallel_results(results, params, files_searched);
|
||||
|
||||
super::GrepResult {
|
||||
matches,
|
||||
total_matches: crate::utils::clamp_u32(total_matches),
|
||||
files_with_matches,
|
||||
files_searched,
|
||||
limit_reached: limit_reached.then_some(true),
|
||||
skipped_oversized: (skipped_oversized > 0)
|
||||
.then(|| crate::utils::clamp_u32(skipped_oversized)),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn assert_same_grep_result(actual: &super::GrepResult, expected: &super::GrepResult) {
|
||||
assert_eq!(actual.total_matches, expected.total_matches);
|
||||
assert_eq!(actual.files_with_matches, expected.files_with_matches);
|
||||
assert_eq!(actual.files_searched, expected.files_searched);
|
||||
assert_eq!(actual.limit_reached, expected.limit_reached);
|
||||
assert_eq!(actual.skipped_oversized, expected.skipped_oversized);
|
||||
assert_eq!(grep_match_snapshots(&actual.matches), grep_match_snapshots(&expected.matches));
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn parallel_streaming_content_matches_sequential_with_context_and_gitignore() {
|
||||
let root = TempDirGuard::new();
|
||||
populate_parallel_parity_tree(root.path());
|
||||
let matcher = super::build_matcher("needle", false, false).expect("build test matcher");
|
||||
let params = unlimited_params(super::OutputMode::Content, 1);
|
||||
|
||||
let parallel = super::run_streaming_grep(
|
||||
root.path(),
|
||||
&matcher,
|
||||
None,
|
||||
None,
|
||||
params,
|
||||
true,
|
||||
true,
|
||||
true,
|
||||
&task::CancelToken::default(),
|
||||
)
|
||||
.expect("parallel streaming grep should succeed");
|
||||
let sequential = super::run_sequential_grep(
|
||||
root.path(),
|
||||
&matcher,
|
||||
None,
|
||||
None,
|
||||
params,
|
||||
true,
|
||||
true,
|
||||
true,
|
||||
&task::CancelToken::default(),
|
||||
None,
|
||||
)
|
||||
.expect("sequential grep should succeed");
|
||||
|
||||
assert_eq!(parallel.1, sequential.1, "oversized skip counts must match");
|
||||
assert_eq!(parallel.2, sequential.2, "searched file counts must match");
|
||||
assert_eq!(file_snapshots(¶llel.0), file_snapshots(&sequential.0));
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn parallel_streaming_count_and_files_modes_match_sequential_reference() {
|
||||
let root = TempDirGuard::new();
|
||||
populate_parallel_parity_tree(root.path());
|
||||
|
||||
for (mode, output_mode) in [
|
||||
(super::OutputMode::Count, GrepOutputMode::Count),
|
||||
(super::OutputMode::FilesWithMatches, GrepOutputMode::FilesWithMatches),
|
||||
] {
|
||||
let mut config = base_grep_config(root.path());
|
||||
config.gitignore = Some(true);
|
||||
config.mode = Some(output_mode);
|
||||
|
||||
let actual = grep_sync(config, None, task::CancelToken::default())
|
||||
.expect("parallel grep should succeed");
|
||||
let expected = sequential_reference_result(root.path(), unlimited_params(mode, 0));
|
||||
|
||||
assert_same_grep_result(&actual, &expected);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn parallel_streaming_type_filter_limits_search_to_matching_extensions() {
|
||||
let root = TempDirGuard::new();
|
||||
let mut expected = Vec::new();
|
||||
for index in 0..150 {
|
||||
let source = format!("dir_{:02}/source_{index:03}.rs", index % 6);
|
||||
expected.push(source.clone());
|
||||
write_file(&root.path().join(&source), "needle\n");
|
||||
write_file(
|
||||
&root
|
||||
.path()
|
||||
.join(format!("dir_{:02}/note_{index:03}.txt", index % 6)),
|
||||
"needle\n",
|
||||
);
|
||||
}
|
||||
expected.sort_unstable();
|
||||
|
||||
let mut config = base_grep_config(root.path());
|
||||
config.type_filter = Some("rs".to_string());
|
||||
|
||||
let result = grep_sync(config, None, task::CancelToken::default())
|
||||
.expect("parallel grep should succeed");
|
||||
let paths: Vec<&str> = result
|
||||
.matches
|
||||
.iter()
|
||||
.map(|matched| matched.path.as_str())
|
||||
.collect();
|
||||
|
||||
assert_eq!(paths, expected);
|
||||
assert_eq!(result.files_searched, 150);
|
||||
assert_eq!(result.files_with_matches, 150);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn parallel_streaming_large_budget_stops_walking_before_scanning_tree() {
|
||||
let root = TempDirGuard::new();
|
||||
let file_count = 3_000;
|
||||
for index in 0..file_count {
|
||||
write_file(&root.path().join(format!("{index:04}.txt")), "needle\n");
|
||||
}
|
||||
|
||||
let mut config = base_grep_config(root.path());
|
||||
config.max_count = Some(100);
|
||||
|
||||
let result = grep_sync(config, None, task::CancelToken::default())
|
||||
.expect("parallel grep should succeed");
|
||||
|
||||
assert_eq!(result.limit_reached, Some(true));
|
||||
assert_eq!(result.matches.len(), 100);
|
||||
for (index, matched) in result.matches.iter().enumerate() {
|
||||
assert_eq!(matched.path, format!("{index:04}.txt"));
|
||||
}
|
||||
assert!(
|
||||
result.files_searched < file_count / 2,
|
||||
"expected budget to stop the walk, scanned {} of {file_count} files",
|
||||
result.files_searched,
|
||||
);
|
||||
}
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn parallel_streaming_defers_oversized_results_until_after_normal_results() {
|
||||
let root = TempDirGuard::new();
|
||||
write_oversized_file(&root.path().join("000_big.txt"), "needle big 0\n");
|
||||
write_oversized_file(&root.path().join("001_big.txt"), "needle big 1\n");
|
||||
for index in 0..300 {
|
||||
write_file(&root.path().join(format!("normal/file_{index:03}.txt")), "needle normal\n");
|
||||
}
|
||||
|
||||
let result = grep_sync(base_grep_config(root.path()), None, task::CancelToken::default())
|
||||
.expect("parallel grep should succeed");
|
||||
let paths: Vec<&str> = result
|
||||
.matches
|
||||
.iter()
|
||||
.map(|matched| matched.path.as_str())
|
||||
.collect();
|
||||
|
||||
assert_eq!(paths.len(), 302);
|
||||
assert!(paths[..300].iter().all(|path| path.starts_with("normal/")));
|
||||
assert_eq!(paths[300..], ["000_big.txt", "001_big.txt"]);
|
||||
assert_eq!(result.files_searched, 302);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn parallel_streaming_respects_cancelled_token_mid_walk() {
|
||||
let root = TempDirGuard::new();
|
||||
for index in 0..1_000 {
|
||||
write_file(&root.path().join(format!("{index:04}.txt")), "needle\n");
|
||||
}
|
||||
|
||||
let ct = task::CancelToken::new(Some(1), None);
|
||||
std::thread::sleep(Duration::from_millis(5));
|
||||
let result = grep_sync(base_grep_config(root.path()), None, ct);
|
||||
|
||||
let Err(err) = result else {
|
||||
panic!("cancelled parallel grep should fail before returning matches");
|
||||
};
|
||||
assert!(
|
||||
err.to_string().contains("Timeout"),
|
||||
"expected timeout cancellation error, got: {err}"
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn streaming_grep_stops_after_first_page_content_budget() {
|
||||
@@ -2480,9 +2944,12 @@ mod tests {
|
||||
fn streaming_grep_quits_parallel_after_large_budget() {
|
||||
let root = TempDirGuard::new();
|
||||
let budget = super::ORDERED_STREAMING_STOP_MAX_COUNT + 1;
|
||||
let file_count = (budget * 3) as usize;
|
||||
for index in 0..file_count {
|
||||
write_file(&root.path().join(format!("{index:05}.txt")), "needle\n");
|
||||
let file_count = super::GREP_STREAM_WINDOW * 3;
|
||||
let expected_paths: Vec<String> = (0..file_count)
|
||||
.map(|index| format!("{index:05}.txt"))
|
||||
.collect();
|
||||
for path in &expected_paths {
|
||||
write_file(&root.path().join(path), "needle\n");
|
||||
}
|
||||
let matcher = super::build_matcher("needle", false, false).expect("build test matcher");
|
||||
let params = content_search_params(budget, None);
|
||||
@@ -2500,18 +2967,30 @@ mod tests {
|
||||
)
|
||||
.expect("streaming grep should succeed");
|
||||
|
||||
let result_paths: Vec<&str> = results
|
||||
.iter()
|
||||
.map(|result| result.relative_path.as_str())
|
||||
.collect();
|
||||
let expected_prefix: Vec<&str> = expected_paths
|
||||
.iter()
|
||||
.take(result_paths.len())
|
||||
.map(String::as_str)
|
||||
.collect();
|
||||
let searched_bound = u64::try_from(file_count / 2).unwrap_or(u64::MAX);
|
||||
let file_count = u64::try_from(file_count).unwrap_or(u64::MAX);
|
||||
|
||||
assert_eq!(skipped_oversized, 0);
|
||||
// The parallel walker honors the budget; tail-race may produce up to one
|
||||
// extra file per worker, but it must not scan the whole tree.
|
||||
assert_eq!(result_paths.first().copied(), Some("00000.txt"));
|
||||
assert_eq!(result_paths, expected_prefix, "results must be a path-ordered prefix");
|
||||
assert!(
|
||||
files_searched < file_count as u64,
|
||||
"expected early stop, scanned {files_searched} of {file_count} files",
|
||||
files_searched < searched_bound,
|
||||
"expected budget to bound work, searched {files_searched} of {file_count} files",
|
||||
);
|
||||
assert!(
|
||||
results.len() < file_count,
|
||||
"expected early stop, collected {} of {file_count} files",
|
||||
results.len(),
|
||||
files_searched >= budget,
|
||||
"early stop must search enough files to satisfy the match budget",
|
||||
);
|
||||
assert!(files_searched < file_count, "early stop must avoid scanning the whole tree");
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
@@ -2526,6 +3005,67 @@ mod tests {
|
||||
write_file(path, &content);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn populate_windowed_oversized_tree(root: &Path) -> Vec<String> {
|
||||
write_oversized_file(&root.join("00_big.txt"), "needle\n");
|
||||
let mut normal_paths = Vec::new();
|
||||
for index in 0..70 {
|
||||
let relative = format!("01_normal_{index:02}.txt");
|
||||
write_file(&root.join(&relative), "needle\n");
|
||||
normal_paths.push(relative);
|
||||
}
|
||||
normal_paths
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn windowed_streaming_skips_oversized_when_normal_matches_satisfy_budget() {
|
||||
let root = TempDirGuard::new();
|
||||
let normal_paths = populate_windowed_oversized_tree(root.path());
|
||||
let mut config = base_grep_config(root.path());
|
||||
config.max_count = Some(65);
|
||||
|
||||
let result = grep_sync(config, None, task::CancelToken::default())
|
||||
.expect("windowed grep should succeed");
|
||||
let paths: Vec<&str> = result
|
||||
.matches
|
||||
.iter()
|
||||
.map(|matched| matched.path.as_str())
|
||||
.collect();
|
||||
let expected: Vec<&str> = normal_paths.iter().take(65).map(String::as_str).collect();
|
||||
|
||||
assert_eq!(paths.len(), 65);
|
||||
assert!(!paths.contains(&"00_big.txt"));
|
||||
assert_eq!(paths, expected, "returned matches must be the normal-file path prefix");
|
||||
assert_eq!(
|
||||
result.files_searched, 70,
|
||||
"oversized file must remain deferred and unsearched once normal files satisfy the budget",
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn windowed_streaming_emits_deferred_oversized_after_normals_when_budget_remains() {
|
||||
let root = TempDirGuard::new();
|
||||
let normal_paths = populate_windowed_oversized_tree(root.path());
|
||||
let mut config = base_grep_config(root.path());
|
||||
config.max_count = Some(100);
|
||||
|
||||
let result = grep_sync(config, None, task::CancelToken::default())
|
||||
.expect("windowed grep should succeed");
|
||||
let paths: Vec<&str> = result
|
||||
.matches
|
||||
.iter()
|
||||
.map(|matched| matched.path.as_str())
|
||||
.collect();
|
||||
let expected_normals: Vec<&str> = normal_paths.iter().map(String::as_str).collect();
|
||||
|
||||
assert_eq!(paths.len(), 71);
|
||||
assert_eq!(paths[..70], expected_normals);
|
||||
assert_eq!(paths.last().copied(), Some("00_big.txt"));
|
||||
assert_eq!(result.files_searched, 71);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn oversized_file_is_searched_over_its_prefix_window() {
|
||||
|
||||
@@ -100,7 +100,7 @@ pub fn walk_workers() -> usize {
|
||||
}
|
||||
|
||||
/// Run parallel traversal-adjacent work on the centralized walker pool.
|
||||
fn with_walk_pool<R>(operation: impl FnOnce() -> R + Send) -> R
|
||||
pub fn with_walk_pool<R>(operation: impl FnOnce() -> R + Send) -> R
|
||||
where
|
||||
R: Send,
|
||||
{
|
||||
@@ -133,6 +133,27 @@ where
|
||||
with_walk_pool(|| items.par_iter().try_for_each(operation))
|
||||
}
|
||||
|
||||
/// Run traversal-adjacent work with per-worker state on the centralized walker
|
||||
/// pool.
|
||||
pub fn parallel_for_each_init<T, S, E>(
|
||||
items: &[T],
|
||||
init: impl Fn() -> S + Send + Sync,
|
||||
operation: impl Fn(&mut S, &T) -> std::result::Result<(), E> + Send + Sync,
|
||||
) -> std::result::Result<(), E>
|
||||
where
|
||||
T: Sync,
|
||||
S: Send,
|
||||
E: Send,
|
||||
{
|
||||
if !should_parallelize(items.len()) {
|
||||
let mut state = init();
|
||||
return items
|
||||
.iter()
|
||||
.try_for_each(|item| operation(&mut state, item));
|
||||
}
|
||||
with_walk_pool(|| items.par_iter().try_for_each_init(init, operation))
|
||||
}
|
||||
|
||||
fn evict_oldest() {
|
||||
if SCAN_CACHE.len() > *MAX_CACHE_ENTRIES
|
||||
&& let Some(oldest_key) = SCAN_CACHE
|
||||
|
||||
+1018
-250
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,380 @@
|
||||
use std::{
|
||||
collections::BTreeMap,
|
||||
convert::Infallible,
|
||||
fs,
|
||||
path::{Path, PathBuf},
|
||||
sync::{
|
||||
Arc, Mutex,
|
||||
atomic::{AtomicUsize, Ordering},
|
||||
},
|
||||
time::{SystemTime, UNIX_EPOCH},
|
||||
};
|
||||
#[cfg(unix)]
|
||||
use std::{
|
||||
ffi::OsString,
|
||||
os::unix::ffi::{OsStrExt, OsStringExt},
|
||||
};
|
||||
|
||||
use pi_walker::{
|
||||
CompiledWalkGlob, Entry, EntryVisitor, FollowLinks, ParallelWalkControl, WalkControl, WalkError,
|
||||
WalkFilter, WalkOptions, WalkOrder, WalkRequest, WalkStatus, walk_entries,
|
||||
};
|
||||
|
||||
struct TempTree {
|
||||
root: PathBuf,
|
||||
}
|
||||
|
||||
impl TempTree {
|
||||
fn new(name: &str) -> Self {
|
||||
let unique = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.expect("system time should be after UNIX epoch")
|
||||
.as_nanos();
|
||||
let root = std::env::temp_dir().join(format!("pi-walker-parallel-{name}-{unique}"));
|
||||
fs::create_dir(&root).expect("temporary root should be created");
|
||||
Self { root }
|
||||
}
|
||||
|
||||
fn path(&self) -> &Path {
|
||||
&self.root
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for TempTree {
|
||||
fn drop(&mut self) {
|
||||
let _ = fs::remove_dir_all(&self.root);
|
||||
}
|
||||
}
|
||||
|
||||
fn write_file(path: impl AsRef<Path>) {
|
||||
let path = path.as_ref();
|
||||
if let Some(parent) = path.parent() {
|
||||
fs::create_dir_all(parent).expect("parent directory should be created");
|
||||
}
|
||||
fs::write(path, b"x").expect("file should be written");
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn write_file_if_supported(path: impl AsRef<Path>) -> std::io::Result<()> {
|
||||
let path = path.as_ref();
|
||||
if let Some(parent) = path.parent() {
|
||||
fs::create_dir_all(parent)?;
|
||||
}
|
||||
fs::write(path, b"x")
|
||||
}
|
||||
|
||||
fn sorted_serial_candidates(request: &WalkRequest) -> Vec<String> {
|
||||
let mut paths = request
|
||||
.collect_file_candidates()
|
||||
.expect("serial candidate collection should succeed")
|
||||
.into_iter()
|
||||
.map(|candidate| candidate.relative)
|
||||
.collect::<Vec<_>>();
|
||||
paths.sort_unstable();
|
||||
paths
|
||||
}
|
||||
|
||||
fn sorted_parallel_candidates(request: &WalkRequest) -> Vec<String> {
|
||||
let paths = Arc::new(Mutex::new(Vec::new()));
|
||||
request
|
||||
.for_each_file_candidate_parallel(
|
||||
{
|
||||
let paths = Arc::clone(&paths);
|
||||
move |candidate| {
|
||||
paths
|
||||
.lock()
|
||||
.expect("candidate list mutex should not be poisoned")
|
||||
.push(candidate.relative.clone());
|
||||
Ok::<_, Infallible>(ParallelWalkControl::Continue)
|
||||
}
|
||||
},
|
||||
|| Ok::<(), Infallible>(()),
|
||||
)
|
||||
.expect("parallel candidate walk should succeed");
|
||||
let mut paths = Arc::into_inner(paths)
|
||||
.expect("candidate list should have no remaining owners")
|
||||
.into_inner()
|
||||
.expect("candidate list mutex should not be poisoned");
|
||||
paths.sort_unstable();
|
||||
paths
|
||||
}
|
||||
|
||||
fn rs_request(root: &Path) -> WalkRequest {
|
||||
WalkRequest::new(root)
|
||||
.hidden(false)
|
||||
.gitignore(true)
|
||||
.skip_node_modules(true)
|
||||
.filter(
|
||||
WalkFilter::files_only()
|
||||
.glob(CompiledWalkGlob::new(["*.rs", "**/*.rs"]).expect("test glob should compile")),
|
||||
)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_candidates_match_serial_with_gitignore_hidden_node_modules_and_glob() {
|
||||
let tree = TempTree::new("candidate-equivalence");
|
||||
write_file(tree.path().join("root.rs"));
|
||||
write_file(tree.path().join("root.txt"));
|
||||
write_file(tree.path().join(".hidden.rs"));
|
||||
write_file(tree.path().join("node_modules/pkg/index.rs"));
|
||||
write_file(tree.path().join("src/visible.rs"));
|
||||
write_file(tree.path().join("src/note.txt"));
|
||||
write_file(tree.path().join("src/nested/keep.rs"));
|
||||
write_file(tree.path().join("src/nested/drop.rs"));
|
||||
write_file(tree.path().join("src/nested/deeper/other.rs"));
|
||||
fs::write(tree.path().join("src/nested/.gitignore"), "*.rs\n!keep.rs\n")
|
||||
.expect("nested .gitignore should be written");
|
||||
|
||||
let request = rs_request(tree.path());
|
||||
|
||||
assert_eq!(
|
||||
sorted_parallel_candidates(&request),
|
||||
sorted_serial_candidates(&request),
|
||||
"parallel traversal should return the same accepted candidate set as serial collection"
|
||||
);
|
||||
assert_eq!(
|
||||
sorted_serial_candidates(&request),
|
||||
vec!["root.rs", "src/nested/keep.rs", "src/visible.rs"],
|
||||
"fixture should exercise the ignore whitelist, hidden-file pruning, node_modules pruning, \
|
||||
and glob filter"
|
||||
);
|
||||
}
|
||||
|
||||
fn create_wide_tree(root: &Path, dirs: usize, files_per_dir: usize) {
|
||||
for dir_index in 0..dirs {
|
||||
let dir = root.join(format!("dir-{dir_index:03}"));
|
||||
fs::create_dir_all(&dir).expect("wide-tree directory should be created");
|
||||
for file_index in 0..files_per_dir {
|
||||
write_file(dir.join(format!("file-{file_index:03}.txt")));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_walk_stops_promptly_when_sink_requests_stop() {
|
||||
let tree = TempTree::new("early-stop");
|
||||
let full_file_count = 2_000;
|
||||
create_wide_tree(tree.path(), 100, full_file_count / 100);
|
||||
let request = WalkRequest::new(tree.path()).filter(WalkFilter::files_only());
|
||||
let invocations = AtomicUsize::new(0);
|
||||
|
||||
let status = request
|
||||
.for_each_file_candidate_parallel(
|
||||
|_| {
|
||||
let seen = invocations.fetch_add(1, Ordering::SeqCst) + 1;
|
||||
if seen >= 5 {
|
||||
Ok::<_, Infallible>(ParallelWalkControl::Stop)
|
||||
} else {
|
||||
Ok(ParallelWalkControl::Continue)
|
||||
}
|
||||
},
|
||||
|| Ok::<(), Infallible>(()),
|
||||
)
|
||||
.expect("parallel walk should stop without an error");
|
||||
|
||||
assert_eq!(status, WalkStatus::Stopped);
|
||||
assert!(
|
||||
invocations.load(Ordering::SeqCst) < full_file_count / 2,
|
||||
"stop should prevent most candidate callbacks after five files, saw {} of {full_file_count}",
|
||||
invocations.load(Ordering::SeqCst)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_walk_returns_sink_error_and_terminates() {
|
||||
let tree = TempTree::new("sink-error");
|
||||
let full_file_count = 800;
|
||||
create_wide_tree(tree.path(), 80, full_file_count / 80);
|
||||
let request = WalkRequest::new(tree.path()).filter(WalkFilter::files_only());
|
||||
let invocations = AtomicUsize::new(0);
|
||||
|
||||
let result = request.for_each_file_candidate_parallel(
|
||||
|_| {
|
||||
let seen = invocations.fetch_add(1, Ordering::SeqCst) + 1;
|
||||
if seen == 3 {
|
||||
Err("sink failed")
|
||||
} else {
|
||||
Ok(ParallelWalkControl::Continue)
|
||||
}
|
||||
},
|
||||
|| Ok(()),
|
||||
);
|
||||
|
||||
match result {
|
||||
Err(WalkError::Interrupted("sink failed")) => {},
|
||||
other => panic!("sink error should be returned as WalkError::Interrupted, got {other:?}"),
|
||||
}
|
||||
assert!(
|
||||
invocations.load(Ordering::SeqCst) < full_file_count / 2,
|
||||
"sink error should terminate traversal instead of visiting most files, saw {} of \
|
||||
{full_file_count}",
|
||||
invocations.load(Ordering::SeqCst)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_walk_returns_heartbeat_error_before_visiting_candidates() {
|
||||
let tree = TempTree::new("heartbeat-error");
|
||||
create_wide_tree(tree.path(), 20, 10);
|
||||
let request = WalkRequest::new(tree.path()).filter(WalkFilter::files_only());
|
||||
let invocations = AtomicUsize::new(0);
|
||||
|
||||
let result = request.for_each_file_candidate_parallel(
|
||||
|_| {
|
||||
invocations.fetch_add(1, Ordering::SeqCst);
|
||||
Ok(ParallelWalkControl::Continue)
|
||||
},
|
||||
|| Err("heartbeat failed"),
|
||||
);
|
||||
|
||||
match result {
|
||||
Err(WalkError::Interrupted("heartbeat failed")) => {},
|
||||
other => {
|
||||
panic!("heartbeat error should be returned as WalkError::Interrupted, got {other:?}")
|
||||
},
|
||||
}
|
||||
assert!(
|
||||
invocations.load(Ordering::SeqCst) < 10,
|
||||
"pre-cancelled heartbeat should allow zero or only a few sink calls, saw {}",
|
||||
invocations.load(Ordering::SeqCst)
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn parallel_follow_links_always_returns_same_candidates_as_serial_collection() {
|
||||
let tree = TempTree::new("follow-links");
|
||||
write_file(tree.path().join("target/child.txt"));
|
||||
write_file(tree.path().join("target/deeper/grandchild.txt"));
|
||||
std::os::unix::fs::symlink(tree.path().join("target"), tree.path().join("link"))
|
||||
.expect("directory symlink should be created");
|
||||
let request = WalkRequest::new(tree.path())
|
||||
.follow_links(FollowLinks::Always)
|
||||
.filter(WalkFilter::files_only());
|
||||
|
||||
assert_eq!(
|
||||
sorted_parallel_candidates(&request),
|
||||
sorted_serial_candidates(&request),
|
||||
"follow-links traversal should expose the same candidate set through parallel API and \
|
||||
serial collection"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unordered_and_path_collection_return_equal_sets_and_path_walk_sorts_each_directory() {
|
||||
let tree = TempTree::new("serial-order");
|
||||
write_file(tree.path().join("βeta.txt"));
|
||||
write_file(tree.path().join("alpha.txt"));
|
||||
write_file(tree.path().join("space name.txt"));
|
||||
write_file(tree.path().join("dir/猫.txt"));
|
||||
write_file(tree.path().join("dir/a.txt"));
|
||||
write_file(tree.path().join("dir/sub/éclair.txt"));
|
||||
write_file(tree.path().join("dir/sub/plain.txt"));
|
||||
#[cfg(unix)]
|
||||
{
|
||||
let _ = write_file_if_supported(
|
||||
tree
|
||||
.path()
|
||||
.join(OsString::from_vec(b"dir/sub/raw-\xFF.txt".to_vec())),
|
||||
);
|
||||
}
|
||||
|
||||
let unordered = collected_path_set(tree.path(), WalkOrder::Unordered);
|
||||
let ordered = collected_path_set(tree.path(), WalkOrder::Path);
|
||||
assert_eq!(
|
||||
unordered, ordered,
|
||||
"serial unordered and path-ordered collection should return the same entry set"
|
||||
);
|
||||
|
||||
let mut visitor = DirectoryOrderVisitor::default();
|
||||
let status = walk_entries(
|
||||
tree.path(),
|
||||
WalkOptions { order: WalkOrder::Path, ..WalkOptions::default() },
|
||||
&mut visitor,
|
||||
|| Ok::<(), Infallible>(()),
|
||||
)
|
||||
.expect("path-ordered walk should succeed");
|
||||
assert_eq!(status, WalkStatus::Complete);
|
||||
visitor.assert_each_directory_sorted();
|
||||
}
|
||||
|
||||
fn collected_path_set(root: &Path, order: WalkOrder) -> Vec<String> {
|
||||
let mut paths = WalkRequest::new(root)
|
||||
.order(order)
|
||||
.collect()
|
||||
.expect("collection should succeed")
|
||||
.entries
|
||||
.into_iter()
|
||||
.map(|entry| entry.path)
|
||||
.collect::<Vec<_>>();
|
||||
paths.sort_unstable();
|
||||
paths
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct DirectoryOrderVisitor {
|
||||
children_by_parent: BTreeMap<String, Vec<Vec<u8>>>,
|
||||
}
|
||||
|
||||
impl DirectoryOrderVisitor {
|
||||
fn assert_each_directory_sorted(&self) {
|
||||
for (parent, children) in &self.children_by_parent {
|
||||
let mut sorted = children.clone();
|
||||
sorted.sort_unstable();
|
||||
assert_eq!(
|
||||
children, &sorted,
|
||||
"children of {parent:?} should be visited in lexicographic path order"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl EntryVisitor for DirectoryOrderVisitor {
|
||||
type Error = Infallible;
|
||||
|
||||
fn visit(&mut self, entry: Entry<'_>) -> Result<WalkControl, Self::Error> {
|
||||
let parent = entry
|
||||
.relative
|
||||
.rsplit_once('/')
|
||||
.map_or_else(String::new, |(parent, _)| parent.to_owned());
|
||||
self
|
||||
.children_by_parent
|
||||
.entry(parent)
|
||||
.or_default()
|
||||
.push(sort_key(entry.name));
|
||||
Ok(WalkControl::Continue)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn sort_key(name: &std::ffi::OsStr) -> Vec<u8> {
|
||||
name.as_bytes().to_vec()
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
fn sort_key(name: &std::ffi::OsStr) -> Vec<u8> {
|
||||
name.to_string_lossy().into_owned().into_bytes()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_deep_tree_relative_path_preserves_full_component_chain() {
|
||||
let tree = TempTree::new("deep-tree");
|
||||
let mut dir = tree.path().to_path_buf();
|
||||
let mut components = Vec::new();
|
||||
for depth in 0..40 {
|
||||
let component = format!("a{depth:02}");
|
||||
dir.push(&component);
|
||||
components.push(component);
|
||||
}
|
||||
fs::create_dir_all(&dir).expect("deep directory chain should be created");
|
||||
write_file(dir.join("leaf.txt"));
|
||||
components.push("leaf.txt".to_owned());
|
||||
let expected = components.join("/");
|
||||
let request = WalkRequest::new(tree.path()).filter(WalkFilter::files_only());
|
||||
|
||||
assert_eq!(
|
||||
sorted_parallel_candidates(&request),
|
||||
vec![expected],
|
||||
"parallel path builder should preserve every nested component in the relative file path"
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,206 @@
|
||||
//! Ignored deterministic timing harness for pi-walker.
|
||||
//!
|
||||
//! Run with:
|
||||
//! cargo test --profile ci -p pi-walker --test perf -- --ignored --nocapture
|
||||
//! --test-threads=1
|
||||
|
||||
use std::{
|
||||
fmt::Write as _,
|
||||
fs,
|
||||
hint::black_box,
|
||||
path::{Path, PathBuf},
|
||||
sync::LazyLock,
|
||||
time::{Duration, Instant, SystemTime, UNIX_EPOCH},
|
||||
};
|
||||
|
||||
use pi_walker::{WalkDetail, WalkOrder, WalkRequest};
|
||||
|
||||
const DIRECTORY_FANOUT: [usize; 5] = [5, 5, 5, 4, 2];
|
||||
const CONTENT_FILE_COUNT: usize = 15_000;
|
||||
const NODE_MODULES_PACKAGES: usize = 50;
|
||||
const NODE_MODULES_FILES_PER_PACKAGE: usize = 10;
|
||||
const MEASURED_ITERATIONS: usize = 5;
|
||||
|
||||
static SYNTHETIC_ROOT: LazyLock<PathBuf> = LazyLock::new(build_synthetic_tree);
|
||||
|
||||
#[test]
|
||||
#[ignore = "run with: cargo test --profile ci -p pi-walker --test perf -- --ignored --nocapture \
|
||||
--test-threads=1"]
|
||||
fn perf_walk_candidates_unordered_gitignore() {
|
||||
let root = SYNTHETIC_ROOT.as_path();
|
||||
run_bench("perf_walk_candidates_unordered_gitignore", || {
|
||||
let candidates = WalkRequest::new(root)
|
||||
.hidden(true)
|
||||
.gitignore(true)
|
||||
.skip_git(true)
|
||||
.skip_node_modules(true)
|
||||
.order(WalkOrder::Unordered)
|
||||
.collect_file_candidates()
|
||||
.expect("collect unordered gitignore candidates");
|
||||
let count = candidates.len();
|
||||
assert!(count > 14_000, "expected a full candidate set, got {count}");
|
||||
count
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[ignore = "run with: cargo test --profile ci -p pi-walker --test perf -- --ignored --nocapture \
|
||||
--test-threads=1"]
|
||||
fn perf_walk_candidates_path_order_no_gitignore() {
|
||||
let root = SYNTHETIC_ROOT.as_path();
|
||||
run_bench("perf_walk_candidates_path_order_no_gitignore", || {
|
||||
let candidates = WalkRequest::new(root)
|
||||
.hidden(true)
|
||||
.gitignore(false)
|
||||
.skip_git(true)
|
||||
.skip_node_modules(true)
|
||||
.order(WalkOrder::Path)
|
||||
.collect_file_candidates()
|
||||
.expect("collect path-ordered candidates without gitignore");
|
||||
let count = candidates.len();
|
||||
assert!(count > 15_000, "expected unignored candidates, got {count}");
|
||||
count
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[ignore = "run with: cargo test --profile ci -p pi-walker --test perf -- --ignored --nocapture \
|
||||
--test-threads=1"]
|
||||
fn perf_walk_collect_full_detail() {
|
||||
let root = SYNTHETIC_ROOT.as_path();
|
||||
run_bench("perf_walk_collect_full_detail", || {
|
||||
let outcome = WalkRequest::new(root)
|
||||
.hidden(true)
|
||||
.gitignore(true)
|
||||
.skip_git(true)
|
||||
.skip_node_modules(true)
|
||||
.order(WalkOrder::Unordered)
|
||||
.detail(WalkDetail::Full)
|
||||
.collect()
|
||||
.expect("collect full-detail entries");
|
||||
let count = outcome.entries.len();
|
||||
assert!(count > 15_000, "expected full-detail entries, got {count}");
|
||||
count
|
||||
});
|
||||
}
|
||||
|
||||
fn run_bench(mut name: &str, mut run: impl FnMut() -> usize) {
|
||||
black_box(run());
|
||||
|
||||
let mut timings = [Duration::ZERO; MEASURED_ITERATIONS];
|
||||
for timing in &mut timings {
|
||||
let started = Instant::now();
|
||||
let observed = run();
|
||||
let elapsed = started.elapsed();
|
||||
black_box(observed);
|
||||
*timing = elapsed;
|
||||
}
|
||||
|
||||
timings.sort_unstable();
|
||||
let median_ms = timings[MEASURED_ITERATIONS / 2].as_secs_f64() * 1_000.0;
|
||||
name = black_box(name);
|
||||
println!("BENCH {name}: {median_ms:.3} ms");
|
||||
}
|
||||
|
||||
fn build_synthetic_tree() -> PathBuf {
|
||||
let root = unique_temp_root("pi-walker-perf");
|
||||
fs::create_dir_all(&root).expect("create synthetic root");
|
||||
fs::create_dir_all(root.join(".git")).expect("create repo marker");
|
||||
|
||||
let directories = create_directory_layout(&root);
|
||||
create_gitignores(&directories);
|
||||
create_content_files(&directories);
|
||||
create_node_modules(&root);
|
||||
|
||||
root
|
||||
}
|
||||
|
||||
fn unique_temp_root(prefix: &str) -> PathBuf {
|
||||
let timestamp = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.expect("system time is after UNIX_EPOCH")
|
||||
.as_nanos();
|
||||
let pid = std::process::id();
|
||||
std::env::temp_dir().join(format!("{prefix}-{pid}-{timestamp}"))
|
||||
}
|
||||
|
||||
fn create_directory_layout(root: &Path) -> Vec<PathBuf> {
|
||||
let mut directories = Vec::with_capacity(1_700);
|
||||
directories.push(root.to_path_buf());
|
||||
|
||||
let mut level = vec![root.to_path_buf()];
|
||||
for (depth, fanout) in DIRECTORY_FANOUT.into_iter().enumerate() {
|
||||
let mut next_level = Vec::with_capacity(level.len() * fanout);
|
||||
for (parent_index, parent) in level.iter().enumerate() {
|
||||
for child in 0..fanout {
|
||||
let directory = parent.join(format!("d{depth:02}-{parent_index:04}-{child:02}"));
|
||||
fs::create_dir_all(&directory).expect("create synthetic directory");
|
||||
directories.push(directory.clone());
|
||||
next_level.push(directory);
|
||||
}
|
||||
}
|
||||
level = next_level;
|
||||
}
|
||||
|
||||
directories
|
||||
}
|
||||
|
||||
fn create_gitignores(directories: &[PathBuf]) {
|
||||
for (directory_id, directory) in directories.iter().enumerate() {
|
||||
if directory_id.is_multiple_of(10) {
|
||||
let pattern = format!("/ignored-{directory_id:04}-*.txt\n");
|
||||
fs::write(directory.join(".gitignore"), pattern).expect("write synthetic gitignore");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn create_content_files(directories: &[PathBuf]) {
|
||||
for file_index in 0..CONTENT_FILE_COUNT {
|
||||
let directory_id = file_index % directories.len();
|
||||
let local_index = file_index / directories.len();
|
||||
let file_name = if directory_id.is_multiple_of(10) && local_index == 0 {
|
||||
format!("ignored-{directory_id:04}-{local_index:03}.txt")
|
||||
} else {
|
||||
format!("file-{directory_id:04}-{local_index:03}.txt")
|
||||
};
|
||||
let path = directories[directory_id].join(file_name);
|
||||
fs::write(path, synthetic_content(file_index)).expect("write synthetic content file");
|
||||
}
|
||||
}
|
||||
|
||||
fn create_node_modules(root: &Path) {
|
||||
let node_modules = root.join("node_modules");
|
||||
for package in 0..NODE_MODULES_PACKAGES {
|
||||
let package_dir = node_modules.join(format!("pkg-{package:02}"));
|
||||
fs::create_dir_all(&package_dir).expect("create synthetic node_modules package");
|
||||
for file in 0..NODE_MODULES_FILES_PER_PACKAGE {
|
||||
let content_index = CONTENT_FILE_COUNT + package * NODE_MODULES_FILES_PER_PACKAGE + file;
|
||||
let path = package_dir.join(format!("file-{file:02}.js"));
|
||||
fs::write(path, synthetic_content(content_index))
|
||||
.expect("write synthetic node_modules file");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn synthetic_content(file_index: usize) -> String {
|
||||
let target_len = 512 + (file_index * 73) % 3_488;
|
||||
let has_common_token = (file_index * 37) % 100 < 60;
|
||||
let has_rare_token = file_index.is_multiple_of(100);
|
||||
let mut content = String::with_capacity(target_len + 96);
|
||||
|
||||
if has_common_token {
|
||||
writeln!(content, "common token needle in file {file_index:05}").expect("write to String");
|
||||
} else {
|
||||
writeln!(content, "ordinary haystack line in file {file_index:05}").expect("write to String");
|
||||
}
|
||||
if has_rare_token {
|
||||
writeln!(content, "rare token NEEDLE_RARE in file {file_index:05}").expect("write to String");
|
||||
}
|
||||
|
||||
let filler = format!("line {file_index:05} deterministic pi walker payload text\n");
|
||||
while content.len() < target_len {
|
||||
content.push_str(&filler);
|
||||
}
|
||||
|
||||
content
|
||||
}
|
||||
@@ -2,6 +2,12 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Changed
|
||||
|
||||
- Rewrote native `grep` directory search to stream while the tree is walked: a work-stealing parallel traversal feeds searchers directly, and content-mode match budgets now terminate the walk itself instead of only the search. Limited searches keep deterministic path-ordered first pages at every budget size via windowed commits, with oversized files still deferred behind normal-sized results.
|
||||
- Faster filesystem walker: gitignore/ignore state is now derived from each directory's own listing instead of up to five per-directory stat probes, per-entry allocations were eliminated through pooled directory scratch buffers and reusable path builders, and a new parallel unordered file-candidate walk API backs full-scan grep.
|
||||
- Concurrent `grep` calls are no longer serialized against each other, searchers are reused per worker instead of rebuilt per file, and non-multiline patterns opt into grep-regex's line-terminator fast path with a compatibility fallback.
|
||||
|
||||
## [16.3.0] - 2026-07-02
|
||||
|
||||
### Added
|
||||
|
||||
Reference in New Issue
Block a user