From cb30c981a0026580149ee18b901a41b19d242fe5 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sat, 4 Jul 2026 11:22:04 +0200 Subject: [PATCH] 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. --- Cargo.toml | 2 +- crates/pi-natives/src/grep.rs | 790 ++++++++++++++--- crates/pi-walker/src/cache.rs | 23 +- crates/pi-walker/src/lib.rs | 1268 ++++++++++++++++++++++------ crates/pi-walker/tests/parallel.rs | 380 +++++++++ crates/pi-walker/tests/perf.rs | 206 +++++ packages/natives/CHANGELOG.md | 6 + 7 files changed, 2298 insertions(+), 377 deletions(-) create mode 100644 crates/pi-walker/tests/parallel.rs create mode 100644 crates/pi-walker/tests/perf.rs diff --git a/Cargo.toml b/Cargo.toml index c6e23c9fc..2fba0a270 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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" diff --git a/crates/pi-natives/src/grep.rs b/crates/pi-natives/src/grep.rs index ed4c84e64..c59bf4456 100644 --- a/crates/pi-natives/src/grep.rs +++ b/crates/pi-natives/src/grep.rs @@ -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> = + const { RefCell::new(None) }; +} + +fn with_parallel_grep_searcher( + 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) -> SearchResult { } /// Internal configuration for grep, extracted from options. -struct GrepConfig { - pattern: String, - path: String, - glob: Option, - type_filter: Option, - ignore_case: Option, - multiline: Option, - hidden: Option, - gitignore: Option, - max_count: Option, - offset: Option, - context_before: Option, - context_after: Option, - context: Option, - max_columns: Option, - mode: Option, - max_count_per_file: Option, +pub(crate) struct GrepConfig { + pub(crate) pattern: String, + pub(crate) path: String, + pub(crate) glob: Option, + pub(crate) type_filter: Option, + pub(crate) ignore_case: Option, + pub(crate) multiline: Option, + pub(crate) hidden: Option, + pub(crate) gitignore: Option, + pub(crate) max_count: Option, + pub(crate) offset: Option, + pub(crate) context_before: Option, + pub(crate) context_after: Option, + pub(crate) context: Option, + pub(crate) max_columns: Option, + pub(crate) mode: Option, + pub(crate) max_count_per_file: Option, } // --------------------------------------------------------------------------- @@ -956,10 +958,19 @@ fn build_regex_matcher( ignore_case: bool, multiline: bool, ) -> std::result::Result { - 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> { 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, - parallel_search: bool, -) -> Result, u64, u64)>> { - let requires_path_order = stop_after_matches.is_some() || params.offset != 0; - let Some(candidates) = collect_grep_candidates( +) -> Result<(Vec, 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, + results: &mut Vec, + matcher: &grep_regex::RegexMatcher, + file_params: SearchParams, + state: &PassState, + ct: &task::CancelToken, + stop_after_matches: u64, +) -> Result { + 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, 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, 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, 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>, 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, + } + + #[cfg(unix)] + fn file_snapshots(results: &[super::FileSearchResult]) -> Vec { + 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, + } + + #[cfg(unix)] + fn grep_match_snapshots(matches: &[super::GrepMatch]) -> Vec { + 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 = (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 { + 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() { diff --git a/crates/pi-walker/src/cache.rs b/crates/pi-walker/src/cache.rs index 59cd0bf50..0906e8d6d 100644 --- a/crates/pi-walker/src/cache.rs +++ b/crates/pi-walker/src/cache.rs @@ -100,7 +100,7 @@ pub fn walk_workers() -> usize { } /// Run parallel traversal-adjacent work on the centralized walker pool. -fn with_walk_pool(operation: impl FnOnce() -> R + Send) -> R +pub fn with_walk_pool(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( + 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 diff --git a/crates/pi-walker/src/lib.rs b/crates/pi-walker/src/lib.rs index 3ed5dc557..4e7da9cd6 100644 --- a/crates/pi-walker/src/lib.rs +++ b/crates/pi-walker/src/lib.rs @@ -8,22 +8,31 @@ mod cache; +#[cfg(not(unix))] +use std::ffi::OsString; +#[cfg(unix)] +use std::os::unix::ffi::OsStrExt; use std::{ borrow::Cow, + cell::{Cell, RefCell}, cmp::Ordering, convert::Infallible, - ffi::{OsStr, OsString}, + ffi::OsStr, fmt, hash::{Hash, Hasher}, io, path::{Path, PathBuf}, - sync::Arc, + sync::{ + Arc, Mutex, + atomic::{AtomicBool, Ordering as AtomicOrdering}, + }, }; pub use cache::{ cache_ttl_ms, classify_file_type, contains_component, empty_recheck_ms, invalidate_all, invalidate_path, invalidate_path_string, max_cache_entries, normalize_relative_path, - parallel_for_each, resolve_search_path, should_parallelize, should_skip_path, walk_workers, + parallel_for_each, parallel_for_each_init, resolve_search_path, should_parallelize, + should_skip_path, walk_workers, }; use globset::{GlobBuilder, GlobSet, GlobSetBuilder}; @@ -1043,6 +1052,30 @@ impl WalkRequest { Ok(stats) } + /// Visit accepted regular-file candidates using an unordered parallel walk. + /// + /// This is a files-only API for consumers that own their output ordering. + /// Candidates may be delivered in any order. [`WalkOptions::order`], + /// [`WalkRequest::visit_order`], [`WalkOptions::emit_root`], and + /// [`WalkRequest::limit`] are ignored. Directory-open errors are skipped + /// with grep-style semantics instead of being delivered to visitors. + /// + /// [`ParallelWalkControl::Stop`] sets a shared stop flag; workers check that + /// flag before reading each directory and while processing directory + /// entries, then the method returns [`WalkStatus::Stopped`] after in-flight + /// work winds down. If a sink or heartbeat returns an error, the first + /// error wins and is returned as [`WalkError::Interrupted`]. + pub fn for_each_file_candidate_parallel( + &self, + sink: impl Fn(&FileCandidate) -> std::result::Result + Send + Sync, + heartbeat: impl Fn() -> std::result::Result<(), E> + Send + Sync, + ) -> std::result::Result> + where + E: Send, + { + run_file_candidate_parallel(self, &sink, &heartbeat) + } + fn collect_with_rank_and_limit( &self, rank: Option, @@ -1189,6 +1222,464 @@ where parallel_for_each(candidates, operation) } +/// Execute work for regular-file candidates with per-worker state. +pub fn execute_candidates_init( + candidates: &[FileCandidate], + init: impl Fn() -> S + Send + Sync, + operation: impl Fn(&mut S, &FileCandidate) -> std::result::Result<(), E> + Send + Sync, +) -> std::result::Result<(), E> +where + S: Send, + E: Send, +{ + parallel_for_each_init(candidates, init, operation) +} + +struct SerialCandidateVisitor<'a, S> { + filter: &'a WalkFilter, + sink: &'a S, +} + +impl EntryVisitor for SerialCandidateVisitor<'_, S> +where + S: Fn(&FileCandidate) -> std::result::Result + Sync, +{ + type Error = E; + + fn visit(&mut self, _entry: Entry<'_>) -> std::result::Result { + Ok(WalkControl::Continue) + } + + fn visit_pre_decided( + &mut self, + entry: Entry<'_>, + ) -> std::result::Result { + if entry.file_type != FileType::File { + return Ok(WalkControl::Continue); + } + let candidate = FileCandidate { + path: entry.path.to_path_buf(), + relative: entry.relative.to_string(), + mtime: entry.mtime, + size: entry.size, + }; + match (self.sink)(&candidate)? { + ParallelWalkControl::Continue => Ok(WalkControl::Continue), + ParallelWalkControl::Stop => Ok(WalkControl::Quit), + } + } + + fn decide_pre_descend( + &mut self, + meta: &EntryMeta<'_>, + ) -> std::result::Result { + let is_dir = meta.file_type == FileType::Dir; + Ok(match self.filter.stream_decision(meta) { + WalkDecision::Include => { + PreDescendDecision { emit: true, descend: is_dir, stop: false } + }, + WalkDecision::Skip => { + PreDescendDecision { emit: false, descend: is_dir, stop: false } + }, + WalkDecision::SkipDescend => { + PreDescendDecision { emit: false, descend: false, stop: false } + }, + WalkDecision::Stop => PreDescendDecision { emit: false, descend: false, stop: true }, + }) + } +} + +struct ParallelWalkContext { + root: PathBuf, + options: WalkOptions, + filter: WalkFilter, + matcher: FastIgnore, +} + +struct ParallelWalkShared<'a, E, S, H> { + stop: AtomicBool, + error: Mutex>, + sink: &'a S, + heartbeat: &'a H, +} + +impl<'a, E, S, H> ParallelWalkShared<'a, E, S, H> { + const fn new(sink: &'a S, heartbeat: &'a H) -> Self { + Self { stop: AtomicBool::new(false), error: Mutex::new(None), sink, heartbeat } + } + + fn request_stop(&self) { + self.stop.store(true, AtomicOrdering::Release); + } + + fn should_stop(&self) -> bool { + self.stop.load(AtomicOrdering::Acquire) + } + + fn record_error(&self, error: E) { + let mut slot = match self.error.lock() { + Ok(slot) => slot, + Err(poisoned) => poisoned.into_inner(), + }; + if slot.is_none() { + *slot = Some(error); + } + self.request_stop(); + } + + fn take_error(&self) -> Option { + match self.error.lock() { + Ok(mut slot) => slot.take(), + Err(poisoned) => poisoned.into_inner().take(), + } + } +} + +thread_local! { + static PARALLEL_HEARTBEAT_COUNTER: Cell = const { Cell::new(0) }; + static PARALLEL_SCRATCH_POOL: RefCell> = const { RefCell::new(Vec::new()) }; +} + +fn should_use_parallel_file_candidate_walk(options: WalkOptions) -> bool { + walk_workers() > 1 && options.follow_links == FollowLinks::Never && !options.same_file_system +} + +fn run_file_candidate_parallel( + request: &WalkRequest, + sink: &S, + heartbeat: &H, +) -> std::result::Result> +where + E: Send, + S: Fn(&FileCandidate) -> std::result::Result + Sync, + H: Fn() -> std::result::Result<(), E> + Sync, +{ + let mut options = request.effective_options(); + if options.min_depth > options.max_depth { + return Ok(WalkStatus::Complete); + } + heartbeat().map_err(WalkError::Interrupted)?; + if !should_use_parallel_file_candidate_walk(options) { + return run_file_candidate_serial(request, options, sink, heartbeat); + } + + options.cache = false; + let Some(root_entry) = root_entry(&request.root, options.detail, options.follow_links)? else { + return Ok(WalkStatus::Complete); + }; + let context = ParallelWalkContext { + root: request.root.clone(), + options, + filter: request.filter.clone(), + matcher: FastIgnore::new(options.use_gitignore), + }; + let root_ignore = context.matcher.root_state(&context.root); + let shared = ParallelWalkShared::new(sink, heartbeat); + + if root_entry.file_type == FileType::File && options.min_depth == 0 { + emit_parallel_root_file(&context, &shared, &root_entry); + } + if root_entry.file_type == FileType::Dir && options.max_depth > 0 && !shared.should_stop() { + let root_dir = context.root.clone(); + cache::with_walk_pool(|| { + rayon::scope(|scope| { + scope.spawn(|scope| { + walk_parallel_dir( + scope, + &context, + &shared, + root_dir, + String::new(), + 0, + root_ignore, + false, + ); + }); + }); + }); + } + + if let Some(error) = shared.take_error() { + Err(WalkError::Interrupted(error)) + } else if shared.should_stop() { + Ok(WalkStatus::Stopped) + } else { + Ok(WalkStatus::Complete) + } +} + +fn run_file_candidate_serial( + request: &WalkRequest, + mut options: WalkOptions, + sink: &S, + heartbeat: &H, +) -> std::result::Result> +where + S: Fn(&FileCandidate) -> std::result::Result + Sync, + H: Fn() -> std::result::Result<(), E> + Sync, +{ + options.cache = false; + options.order = WalkOrder::Unordered; + options.contents_first = false; + options.emit_root = true; + options.directory_errors = DirectoryErrorMode::Visit; + let mut visitor = SerialCandidateVisitor { filter: &request.filter, sink }; + walk_entries(&request.root, options, &mut visitor, heartbeat) +} + +fn emit_parallel_root_file( + context: &ParallelWalkContext, + shared: &ParallelWalkShared<'_, E, S, H>, + root_entry: &RootEntry, +) where + E: Send, + S: Fn(&FileCandidate) -> std::result::Result + Sync, + H: Fn() -> std::result::Result<(), E> + Sync, +{ + let meta = EntryMeta { + root: &context.root, + absolute_path: Cow::Borrowed(context.root.as_path()), + relative_path: "", + file_type: FileType::File, + mtime: root_entry.mtime, + size: root_entry.size, + depth: 0, + }; + match context.filter.stream_decision(&meta) { + WalkDecision::Include => { + let candidate = FileCandidate { + path: context.root.clone(), + relative: String::new(), + mtime: root_entry.mtime, + size: root_entry.size, + }; + let _ = emit_parallel_candidate(shared, &candidate); + }, + WalkDecision::Stop => shared.request_stop(), + WalkDecision::Skip | WalkDecision::SkipDescend => {}, + } +} + +fn take_parallel_scratch() -> DirScratch { + let mut scratch = PARALLEL_SCRATCH_POOL + .with(|pool| pool.borrow_mut().pop()) + .unwrap_or_default(); + scratch.clear_listing(); + scratch +} + +fn recycle_parallel_scratch(mut scratch: DirScratch) { + scratch.clear_listing(); + PARALLEL_SCRATCH_POOL.with(|pool| pool.borrow_mut().push(scratch)); +} + +fn parallel_heartbeat(shared: &ParallelWalkShared<'_, E, S, H>) -> bool +where + E: Send, + H: Fn() -> std::result::Result<(), E> + Sync, +{ + if shared.should_stop() { + return false; + } + let should_call = PARALLEL_HEARTBEAT_COUNTER.with(|counter| { + let visited = counter.get(); + if visited == 0 || visited >= HEARTBEAT_INTERVAL { + counter.set(1); + true + } else { + counter.set(visited + 1); + false + } + }); + if !should_call { + return true; + } + match (shared.heartbeat)() { + Ok(()) => true, + Err(error) => { + shared.record_error(error); + false + }, + } +} + +fn emit_parallel_candidate( + shared: &ParallelWalkShared<'_, E, S, H>, + candidate: &FileCandidate, +) -> bool +where + E: Send, + S: Fn(&FileCandidate) -> std::result::Result + Sync, +{ + match (shared.sink)(candidate) { + Ok(ParallelWalkControl::Continue) => true, + Ok(ParallelWalkControl::Stop) => { + shared.request_stop(); + false + }, + Err(error) => { + shared.record_error(error); + false + }, + } +} + +fn walk_parallel_dir<'scope, E, S, H>( + scope: &rayon::Scope<'scope>, + context: &'scope ParallelWalkContext, + shared: &'scope ParallelWalkShared<'_, E, S, H>, + dir: PathBuf, + relative_dir: String, + depth: usize, + ignore_state: Arc, + derive_ignore_from_entries: bool, +) where + E: Send + 'scope, + S: Fn(&FileCandidate) -> std::result::Result + Sync + 'scope, + H: Fn() -> std::result::Result<(), E> + Sync + 'scope, +{ + if shared.should_stop() { + return; + } + let mut scratch = take_parallel_scratch(); + let ignore_entries = match collect_directory_entries( + &dir, + context.options.detail, + &mut scratch, + &context.matcher, + derive_ignore_from_entries, + ) { + Ok(ignore_entries) => ignore_entries, + Err(ReadDirError::Io(_) | ReadDirError::Walk(WalkError::InvalidData { .. })) => { + recycle_parallel_scratch(scratch); + return; + }, + Err(ReadDirError::Walk(WalkError::Interrupted(error))) => { + shared.record_error(error); + recycle_parallel_scratch(scratch); + return; + }, + }; + let dir_ignore = context.matcher.state_from_entries( + &ignore_state, + &dir, + ignore_entries, + derive_ignore_from_entries, + ); + let mut absolute = dir; + let mut relative = relative_dir; + + for index in 0..scratch.entries.len() { + if !parallel_heartbeat(shared) { + break; + } + let entry = &scratch.entries[index]; + let name = scratch.name(entry); + if is_dot_entry(name) { + continue; + } + if !context.options.include_hidden && is_hidden_name(name) { + continue; + } + if (context.options.skip_git && is_git_name(name)) + || (context.options.skip_node_modules && is_node_modules_name(name)) + { + continue; + } + + let name_str = entry_name(name); + if name_str.is_empty() { + continue; + } + let next_depth = depth + 1; + if next_depth > context.options.max_depth { + continue; + } + + let relative_len = relative.len(); + absolute.push(name); + push_relative_name(&mut relative, &name_str); + let is_dir = entry.file_type == FileType::Dir; + if !context.matcher.is_ignored(&dir_ignore, &absolute, is_dir) { + let decision = { + let meta = EntryMeta { + root: &context.root, + absolute_path: Cow::Borrowed(absolute.as_path()), + relative_path: &relative, + file_type: entry.file_type, + mtime: entry.mtime, + size: entry.size, + depth: next_depth, + }; + context.filter.stream_decision(&meta) + }; + match decision { + WalkDecision::Include => { + if entry.file_type == FileType::File && next_depth >= context.options.min_depth { + let candidate = FileCandidate { + path: absolute.clone(), + relative: relative.clone(), + mtime: entry.mtime, + size: entry.size, + }; + if !emit_parallel_candidate(shared, &candidate) { + absolute.pop(); + relative.truncate(relative_len); + break; + } + } + if is_dir && next_depth < context.options.max_depth && !shared.should_stop() { + let child_dir = absolute.clone(); + let child_relative = relative.clone(); + let child_ignore = Arc::clone(&dir_ignore); + scope.spawn(move |scope| { + walk_parallel_dir( + scope, + context, + shared, + child_dir, + child_relative, + next_depth, + child_ignore, + true, + ); + }); + } + }, + WalkDecision::Skip => { + if is_dir && next_depth < context.options.max_depth && !shared.should_stop() { + let child_dir = absolute.clone(); + let child_relative = relative.clone(); + let child_ignore = Arc::clone(&dir_ignore); + scope.spawn(move |scope| { + walk_parallel_dir( + scope, + context, + shared, + child_dir, + child_relative, + next_depth, + child_ignore, + true, + ); + }); + } + }, + WalkDecision::SkipDescend => {}, + WalkDecision::Stop => { + shared.request_stop(); + absolute.pop(); + relative.truncate(relative_len); + break; + }, + } + } + absolute.pop(); + relative.truncate(relative_len); + } + recycle_parallel_scratch(scratch); +} + /// Visitor decision for streaming traversal. #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum WalkControl { @@ -1209,6 +1700,15 @@ pub enum WalkStatus { Stopped, } +/// Control returned by unordered parallel file-candidate sinks. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum ParallelWalkControl { + /// Continue walking and delivering candidates. + Continue, + /// Stop all workers as promptly as possible. + Stop, +} + /// Owned entry returned by [`collect_entries`]. #[derive(Clone, Debug, PartialEq)] pub struct CollectedEntry { @@ -1732,21 +2232,88 @@ struct RawDirEntry<'a> { size: Option, } -struct OwnedDirEntry { +#[cfg(unix)] +struct DirEntryRecord { + name_start: usize, + name_len: usize, + file_type: FileType, + mtime: Option, + size: Option, +} + +#[cfg(not(unix))] +struct DirEntryRecord { name: OsString, file_type: FileType, mtime: Option, size: Option, } -impl RawDirEntry<'_> { - fn into_owned(self) -> OwnedDirEntry { - OwnedDirEntry { - name: self.name.into_owned(), - file_type: self.file_type, - mtime: self.mtime, - size: self.size, - } +#[derive(Default)] +struct DirScratch { + entries: Vec, + #[cfg(unix)] + name_bytes: Vec, + read_buffer: Vec, +} + +impl DirScratch { + fn clear_listing(&mut self) { + self.entries.clear(); + #[cfg(unix)] + self.name_bytes.clear(); + } + + #[cfg(unix)] + fn push(&mut self, entry: RawDirEntry<'_>) { + let name_os: &OsStr = entry.name.as_ref(); + let name = name_os.as_bytes(); + let name_start = self.name_bytes.len(); + self.name_bytes.extend_from_slice(name); + self.entries.push(DirEntryRecord { + name_start, + name_len: name.len(), + file_type: entry.file_type, + mtime: entry.mtime, + size: entry.size, + }); + } + + #[cfg(not(unix))] + fn push(&mut self, entry: RawDirEntry<'_>) { + self.entries.push(DirEntryRecord { + name: entry.name.into_owned(), + file_type: entry.file_type, + mtime: entry.mtime, + size: entry.size, + }); + } + + #[cfg(unix)] + fn sort_by_name(&mut self) { + let names = &self.name_bytes; + self.entries.sort_unstable_by(|left, right| { + let left_name = &names[left.name_start..left.name_start + left.name_len]; + let right_name = &names[right.name_start..right.name_start + right.name_len]; + left_name.cmp(right_name) + }); + } + + #[cfg(not(unix))] + fn sort_by_name(&mut self) { + self + .entries + .sort_unstable_by(|left, right| left.name.cmp(&right.name)); + } + + #[cfg(unix)] + fn name<'a>(&'a self, entry: &DirEntryRecord) -> &'a OsStr { + OsStr::from_bytes(&self.name_bytes[entry.name_start..entry.name_start + entry.name_len]) + } + + #[cfg(not(unix))] + fn name<'a>(&'a self, entry: &'a DirEntryRecord) -> &'a OsStr { + &entry.name } } @@ -1908,6 +2475,9 @@ where root_device, symlink_ancestors: SymlinkAncestorStack::default(), matcher, + absolute_path: root.to_path_buf(), + relative_path: String::new(), + scratch_pool: Vec::new(), visited: 0, heartbeat, }; @@ -1921,6 +2491,9 @@ struct WalkContext<'a, H> { root_device: Option, symlink_ancestors: SymlinkAncestorStack, matcher: FastIgnore, + absolute_path: PathBuf, + relative_path: String, + scratch_pool: Vec, visited: usize, heartbeat: H, } @@ -1992,7 +2565,7 @@ impl WalkContext<'_, H> { let stopped = if root_entry.file_type == FileType::Dir && self.options.max_depth > 0 && decision.descend { - self.walk_dir(root, "", 0, root_ignore, visitor)? + self.walk_dir(0, root_ignore, false, visitor)? } else { false }; @@ -2027,12 +2600,34 @@ impl WalkContext<'_, H> { Ok(WalkStatus::Complete) } + fn take_scratch(&mut self) -> DirScratch { + let mut scratch = self.scratch_pool.pop().unwrap_or_default(); + scratch.clear_listing(); + scratch + } + + fn recycle_scratch(&mut self, mut scratch: DirScratch) { + scratch.clear_listing(); + self.scratch_pool.push(scratch); + } + + fn push_entry_path(&mut self, name: &OsStr, name_str: &str) -> usize { + let relative_len = self.relative_path.len(); + self.absolute_path.push(name); + push_relative_name(&mut self.relative_path, name_str); + relative_len + } + + fn pop_entry_path(&mut self, relative_len: usize) { + self.relative_path.truncate(relative_len); + self.absolute_path.pop(); + } + fn walk_dir( &mut self, - dir: &Path, - relative_dir: &str, depth: usize, ignore_state: &Arc, + derive_ignore_from_entries: bool, visitor: &mut V, ) -> std::result::Result> where @@ -2040,49 +2635,62 @@ impl WalkContext<'_, H> { H: FnMut() -> std::result::Result<(), V::Error>, { if self.options.follow_links == FollowLinks::Always - && let Ok(identity) = directory_identity(dir) + && let Ok(identity) = directory_identity(&self.absolute_path) { self.symlink_ancestors.push(identity); - let result = self.walk_dir_inner(dir, relative_dir, depth, ignore_state, visitor); + let result = self.walk_dir_inner(depth, ignore_state, derive_ignore_from_entries, visitor); self.symlink_ancestors.pop(); return result; } - self.walk_dir_inner(dir, relative_dir, depth, ignore_state, visitor) + self.walk_dir_inner(depth, ignore_state, derive_ignore_from_entries, visitor) } fn walk_dir_inner( &mut self, - dir: &Path, - relative_dir: &str, depth: usize, ignore_state: &Arc, + derive_ignore_from_entries: bool, visitor: &mut V, ) -> std::result::Result> where V: EntryVisitor, H: FnMut() -> std::result::Result<(), V::Error>, { - let mut raw_entries = Vec::new(); - match platform::read_dir_entries(dir, self.options.detail, |entry| { - raw_entries.push(entry.into_owned()); - Ok(ReadDirControl::Continue) - }) { - Ok(_) => {}, - Err(err) => return handle_read_dir_error(dir, err, self.options, visitor), - } + let mut scratch = self.take_scratch(); + let ignore_entries = match collect_directory_entries( + &self.absolute_path, + self.options.detail, + &mut scratch, + &self.matcher, + derive_ignore_from_entries, + ) { + Ok(ignore_entries) => ignore_entries, + Err(err) => { + let dir = self.absolute_path.clone(); + self.recycle_scratch(scratch); + return handle_read_dir_error(&dir, err, self.options, visitor); + }, + }; if self.options.order == WalkOrder::Path { - raw_entries.sort_unstable_by(|a, b| a.name.cmp(&b.name)); + scratch.sort_by_name(); } + let dir_ignore = self.matcher.state_from_entries( + ignore_state, + &self.absolute_path, + ignore_entries, + derive_ignore_from_entries, + ); - for entry in raw_entries { + for index in 0..scratch.entries.len() { if self.visited == 0 || self.visited >= HEARTBEAT_INTERVAL { self.visited = 0; (self.heartbeat)().map_err(WalkError::Interrupted)?; } self.visited += 1; - let name = entry.name.as_ref(); + let entry = &scratch.entries[index]; + let name = scratch.name(entry); if is_dot_entry(name) { continue; } @@ -2099,213 +2707,277 @@ impl WalkContext<'_, H> { if name_str.is_empty() { continue; } - let relative = join_relative_path(relative_dir, &name_str); let next_depth = depth + 1; if next_depth > self.options.max_depth { continue; } - let absolute = dir.join(Path::new(name)); - let mut file_type = entry.file_type; - let mut mtime = entry.mtime; - let mut size = entry.size; - let mut is_dir = entry.file_type == FileType::Dir; - let mut descend = is_dir; - let mut followed_symlink_dir = false; - let followed_metadata = if entry.file_type == FileType::Symlink - && self.options.follow_links == FollowLinks::Always - { - match std::fs::metadata(&absolute) { - Ok(metadata) => Some(metadata), - Err(err) => { - if self.options.directory_errors == DirectoryErrorMode::SkipSkippable - && is_skippable_directory_error(&err) - { - continue; - } - if self.options.directory_errors == DirectoryErrorMode::Visit { - match visitor - .visit_directory_error(DirectoryError { path: &absolute, error: &err }) - .map_err(WalkError::Interrupted)? - { - WalkControl::Quit => return Ok(true), - WalkControl::SkipDescend | WalkControl::Continue => { - continue; - }, - } - } - return Err(WalkError::InvalidData { - path: absolute, - message: err.to_string(), - }); - }, - } - } else { - None - }; - - if let Some(target_metadata) = followed_metadata.as_ref() { - let Some(target_file_type) = file_type_from_metadata(target_metadata) else { - continue; - }; - file_type = target_file_type; - if target_file_type == FileType::Dir { - is_dir = true; - descend = true; - followed_symlink_dir = true; - } else { - is_dir = false; - descend = false; - } - if self.options.detail == WalkDetail::Full { - if target_file_type == FileType::File { - size = Some(target_metadata.len() as f64); - } else { - size = None; - } - mtime = target_metadata - .modified() - .ok() - .and_then(|time| time.duration_since(std::time::UNIX_EPOCH).ok()) - .map(|duration| duration.as_millis() as f64); - } - } - - if self.matcher.is_ignored(ignore_state, &absolute, is_dir) { - continue; - } - - if !is_effective_path_on_root_file_system( - &absolute, + let relative_len = self.push_entry_path(name, &name_str); + let entry_result = self.walk_current_entry( + name, + entry.file_type, + entry.mtime, + entry.size, next_depth, - self.options.follow_links, - self.root_device, - followed_metadata.as_ref(), - ) { - continue; + &dir_ignore, + visitor, + ); + self.pop_entry_path(relative_len); + if entry_result? { + self.recycle_scratch(scratch); + return Ok(true); } + } - if followed_symlink_dir && descend { - match directory_identity(&absolute) { - Ok(target_id) => { - if self.symlink_ancestors.contains(&target_id) { - let loop_err = io::Error::other("filesystem loop detected"); - if self.options.directory_errors == DirectoryErrorMode::Visit { - match visitor - .visit_directory_error(DirectoryError { - path: &absolute, - error: &loop_err, - }) - .map_err(WalkError::Interrupted)? - { - WalkControl::Quit => return Ok(true), - WalkControl::SkipDescend | WalkControl::Continue => { - continue; - }, - } - } else if self.options.directory_errors == DirectoryErrorMode::SkipSkippable { - continue; - } - return Err(WalkError::InvalidData { - path: absolute, - message: "filesystem loop detected".to_string(), - }); - } - }, - Err(err) => { - if self.options.directory_errors == DirectoryErrorMode::SkipSkippable - && is_skippable_directory_error(&err) + self.recycle_scratch(scratch); + Ok(false) + } + + fn walk_current_entry( + &mut self, + name: &OsStr, + entry_file_type: FileType, + entry_mtime: Option, + entry_size: Option, + next_depth: usize, + dir_ignore: &Arc, + visitor: &mut V, + ) -> std::result::Result> + where + V: EntryVisitor, + H: FnMut() -> std::result::Result<(), V::Error>, + { + let mut file_type = entry_file_type; + let mut mtime = entry_mtime; + let mut size = entry_size; + let mut is_dir = entry_file_type == FileType::Dir; + let mut descend = is_dir; + let mut followed_symlink_dir = false; + let followed_metadata = if entry_file_type == FileType::Symlink + && self.options.follow_links == FollowLinks::Always + { + match std::fs::metadata(&self.absolute_path) { + Ok(metadata) => Some(metadata), + Err(err) => { + if self.options.directory_errors == DirectoryErrorMode::SkipSkippable + && is_skippable_directory_error(&err) + { + return Ok(false); + } + if self.options.directory_errors == DirectoryErrorMode::Visit { + match visitor + .visit_directory_error(DirectoryError { + path: &self.absolute_path, + error: &err, + }) + .map_err(WalkError::Interrupted)? { - continue; + WalkControl::Quit => return Ok(true), + WalkControl::SkipDescend | WalkControl::Continue => { + return Ok(false); + }, } + } + return Err(WalkError::InvalidData { + path: self.absolute_path.clone(), + message: err.to_string(), + }); + }, + } + } else { + None + }; + + if let Some(target_metadata) = followed_metadata.as_ref() { + let Some(target_file_type) = file_type_from_metadata(target_metadata) else { + return Ok(false); + }; + file_type = target_file_type; + if target_file_type == FileType::Dir { + is_dir = true; + descend = true; + followed_symlink_dir = true; + } else { + is_dir = false; + descend = false; + } + if self.options.detail == WalkDetail::Full { + if target_file_type == FileType::File { + size = Some(target_metadata.len() as f64); + } else { + size = None; + } + mtime = target_metadata + .modified() + .ok() + .and_then(|time| time.duration_since(std::time::UNIX_EPOCH).ok()) + .map(|duration| duration.as_millis() as f64); + } + } + + if self + .matcher + .is_ignored(dir_ignore, &self.absolute_path, is_dir) + { + return Ok(false); + } + + if !is_effective_path_on_root_file_system( + &self.absolute_path, + next_depth, + self.options.follow_links, + self.root_device, + followed_metadata.as_ref(), + ) { + return Ok(false); + } + + if followed_symlink_dir && descend { + match directory_identity(&self.absolute_path) { + Ok(target_id) => { + if self.symlink_ancestors.contains(&target_id) { + let loop_err = io::Error::other("filesystem loop detected"); if self.options.directory_errors == DirectoryErrorMode::Visit { match visitor - .visit_directory_error(DirectoryError { path: &absolute, error: &err }) + .visit_directory_error(DirectoryError { + path: &self.absolute_path, + error: &loop_err, + }) .map_err(WalkError::Interrupted)? { WalkControl::Quit => return Ok(true), WalkControl::SkipDescend | WalkControl::Continue => { - continue; + return Ok(false); }, } + } else if self.options.directory_errors == DirectoryErrorMode::SkipSkippable { + return Ok(false); } return Err(WalkError::InvalidData { - path: absolute, - message: err.to_string(), + path: self.absolute_path.clone(), + message: "filesystem loop detected".to_string(), }); - }, - } + } + }, + Err(err) => { + if self.options.directory_errors == DirectoryErrorMode::SkipSkippable + && is_skippable_directory_error(&err) + { + return Ok(false); + } + if self.options.directory_errors == DirectoryErrorMode::Visit { + match visitor + .visit_directory_error(DirectoryError { + path: &self.absolute_path, + error: &err, + }) + .map_err(WalkError::Interrupted)? + { + WalkControl::Quit => return Ok(true), + WalkControl::SkipDescend | WalkControl::Continue => { + return Ok(false); + }, + } + } + return Err(WalkError::InvalidData { + path: self.absolute_path.clone(), + message: err.to_string(), + }); + }, } + } + let decision = { let meta = EntryMeta { root: self.root_path, - absolute_path: Cow::Borrowed(&absolute), - relative_path: &relative, + absolute_path: Cow::Borrowed(self.absolute_path.as_path()), + relative_path: &self.relative_path, file_type, mtime, size, depth: next_depth, }; - - let decision = visitor + visitor .decide_pre_descend(&meta) - .map_err(WalkError::Interrupted)?; - if decision.stop { - return Ok(true); + .map_err(WalkError::Interrupted)? + }; + if decision.stop { + return Ok(true); + } + + if !self.options.contents_first && decision.emit && next_depth >= self.options.min_depth { + match visitor + .visit_pre_decided(Entry { + path: self.absolute_path.as_path(), + relative: &self.relative_path, + name, + file_type, + mtime, + size, + depth: next_depth, + }) + .map_err(WalkError::Interrupted)? + { + WalkControl::Quit => return Ok(true), + WalkControl::SkipDescend => return Ok(false), + WalkControl::Continue => {}, } + } - if !self.options.contents_first && decision.emit && next_depth >= self.options.min_depth { - match visitor - .visit_pre_decided(Entry { - path: &absolute, - relative: &relative, - name: entry.name.as_ref(), - file_type, - mtime, - size, - depth: next_depth, - }) - .map_err(WalkError::Interrupted)? - { - WalkControl::Quit => return Ok(true), - WalkControl::SkipDescend => continue, - WalkControl::Continue => {}, - } - } + let child_stopped = if descend && next_depth < self.options.max_depth && decision.descend { + self.walk_dir(next_depth, dir_ignore, true, visitor)? + } else { + false + }; - let child_stopped = if descend && next_depth < self.options.max_depth && decision.descend { - let child_ignore = self.matcher.child_state(ignore_state, &absolute); - self.walk_dir(&absolute, &relative, next_depth, &child_ignore, visitor)? - } else { - false - }; + if child_stopped { + return Ok(true); + } - if child_stopped { - return Ok(true); - } - - if self.options.contents_first && decision.emit && next_depth >= self.options.min_depth { - match visitor - .visit_pre_decided(Entry { - path: &absolute, - relative: &relative, - name: entry.name.as_ref(), - file_type, - mtime, - size, - depth: next_depth, - }) - .map_err(WalkError::Interrupted)? - { - WalkControl::Quit => return Ok(true), - WalkControl::SkipDescend | WalkControl::Continue => {}, - } + if self.options.contents_first && decision.emit && next_depth >= self.options.min_depth { + match visitor + .visit_pre_decided(Entry { + path: self.absolute_path.as_path(), + relative: &self.relative_path, + name, + file_type, + mtime, + size, + depth: next_depth, + }) + .map_err(WalkError::Interrupted)? + { + WalkControl::Quit => return Ok(true), + WalkControl::SkipDescend | WalkControl::Continue => {}, } } Ok(false) } } + +fn collect_directory_entries( + dir: &Path, + detail: WalkDetail, + scratch: &mut DirScratch, + matcher: &FastIgnore, + derive_ignore_from_entries: bool, +) -> std::result::Result> { + scratch.clear_listing(); + let mut ignore_entries = IgnoreEntryNames::default(); + let track_ignore_entries = derive_ignore_from_entries && matcher.use_gitignore; + let mut read_buffer = std::mem::take(&mut scratch.read_buffer); + let result = platform::read_dir_entries(dir, detail, &mut read_buffer, |entry| { + if track_ignore_entries { + ignore_entries.record(entry.name.as_ref(), entry.file_type); + } + scratch.push(entry); + Ok(ReadDirControl::Continue) + }); + scratch.read_buffer = read_buffer; + result?; + Ok(ignore_entries) +} /// Return whether [`WalkDetail::Full`] provides file sizes without per-entry /// metadata syscalls on this platform. pub const fn supports_cheap_size_hints() -> bool { @@ -2318,6 +2990,8 @@ struct IgnoreState { gitignore_matcher: Option, git_exclude_matcher: Option, has_git: bool, + chain_has_matchers: bool, + any_git: bool, } struct FastIgnore { @@ -2325,6 +2999,36 @@ struct FastIgnore { use_gitignore: bool, } +#[derive(Clone, Copy, Default)] +struct IgnoreEntryNames { + ignore_file: bool, + gitignore_file: bool, + git_dir: bool, + repo_marker: bool, +} + +impl IgnoreEntryNames { + fn record(&mut self, name: &OsStr, file_type: FileType) { + if matches!(file_type, FileType::File | FileType::Symlink) { + if name == OsStr::new(".ignore") { + self.ignore_file = true; + } else if name == OsStr::new(".gitignore") { + self.gitignore_file = true; + } + } + if name == OsStr::new(".git") { + self.git_dir = true; + self.repo_marker = true; + } else if name == OsStr::new(".jj") { + self.repo_marker = true; + } + } + + const fn has_relevant(self) -> bool { + self.ignore_file || self.gitignore_file || self.git_dir || self.repo_marker + } +} + fn has_repo_marker(dir: &Path) -> bool { dir.join(".git").exists() || dir.join(".jj").exists() } @@ -2342,16 +3046,66 @@ impl IgnoreState { fn build(dir: &Path, parent: Option>) -> Arc { let has_git = has_repo_marker(dir); let git_exclude = dir.join(".git/info/exclude"); - Arc::new(Self { + Self::new( parent, - ignore_matcher: load_gitignore(dir, &dir.join(".ignore")), - gitignore_matcher: load_gitignore(dir, &dir.join(".gitignore")), - git_exclude_matcher: if has_git { + load_gitignore(dir, &dir.join(".ignore")), + load_gitignore(dir, &dir.join(".gitignore")), + if has_git { load_gitignore(dir, &git_exclude) } else { None }, has_git, + ) + } + + fn build_from_entry_names(dir: &Path, parent: &Arc, names: IgnoreEntryNames) -> Arc { + if !names.has_relevant() { + return Arc::clone(parent); + } + let git_exclude = dir.join(".git/info/exclude"); + Self::new( + Some(Arc::clone(parent)), + if names.ignore_file { + load_gitignore(dir, &dir.join(".ignore")) + } else { + None + }, + if names.gitignore_file { + load_gitignore(dir, &dir.join(".gitignore")) + } else { + None + }, + if names.git_dir { + load_gitignore(dir, &git_exclude) + } else { + None + }, + names.repo_marker, + ) + } + + fn new( + parent: Option>, + ignore_matcher: Option, + gitignore_matcher: Option, + git_exclude_matcher: Option, + has_git: bool, + ) -> Arc { + let parent_has_matchers = parent + .as_ref() + .is_some_and(|parent| parent.chain_has_matchers); + let parent_has_git = parent.as_ref().is_some_and(|parent| parent.any_git); + let has_matchers = + ignore_matcher.is_some() || gitignore_matcher.is_some() || git_exclude_matcher.is_some(); + Arc::new(Self { + parent, + ignore_matcher, + gitignore_matcher, + git_exclude_matcher, + has_git, + chain_has_matchers: has_matchers || parent_has_matchers, + any_git: has_git || parent_has_git, }) } @@ -2393,9 +3147,15 @@ impl FastIgnore { IgnoreState::build(root, IgnoreState::build_parents(root, self.use_gitignore)) } - fn child_state(&self, parent: &Arc, abs_dir: &Path) -> Arc { - if self.use_gitignore { - IgnoreState::build(abs_dir, Some(Arc::clone(parent))) + fn state_from_entries( + &self, + parent: &Arc, + dir: &Path, + names: IgnoreEntryNames, + derive_ignore_from_entries: bool, + ) -> Arc { + if self.use_gitignore && derive_ignore_from_entries { + IgnoreState::build_from_entry_names(dir, parent, names) } else { Arc::clone(parent) } @@ -2406,33 +3166,40 @@ impl FastIgnore { return false; } - let any_git = Self::has_git_state(state); + let any_git = state.any_git; + let global_matcher_applies = any_git && self.global.is_some(); + if !state.chain_has_matchers && !global_matcher_applies { + return false; + } + let mut saw_git = false; let mut ignore_match = ignore::Match::None; let mut gitignore_match = ignore::Match::None; let mut git_exclude_match = ignore::Match::None; - let mut current = Some(state.as_ref()); - while let Some(frame) = current { - if ignore_match.is_none() - && let Some(matcher) = &frame.ignore_matcher - { - ignore_match = matcher.matched(path, is_dir); + if state.chain_has_matchers { + let mut current = Some(state.as_ref()); + while let Some(frame) = current { + if ignore_match.is_none() + && let Some(matcher) = &frame.ignore_matcher + { + ignore_match = matcher.matched(path, is_dir); + } + if gitignore_match.is_none() + && let Some(matcher) = &frame.gitignore_matcher + { + gitignore_match = matcher.matched(path, is_dir); + } + if any_git + && !saw_git + && git_exclude_match.is_none() + && let Some(matcher) = &frame.git_exclude_matcher + { + git_exclude_match = matcher.matched(path, is_dir); + } + saw_git = saw_git || frame.has_git; + current = frame.parent.as_deref(); } - if gitignore_match.is_none() - && let Some(matcher) = &frame.gitignore_matcher - { - gitignore_match = matcher.matched(path, is_dir); - } - if any_git - && !saw_git - && git_exclude_match.is_none() - && let Some(matcher) = &frame.git_exclude_matcher - { - git_exclude_match = matcher.matched(path, is_dir); - } - saw_git = saw_git || frame.has_git; - current = frame.parent.as_deref(); } match ignore_match { ignore::Match::Ignore(_) => return true, @@ -2458,17 +3225,6 @@ impl FastIgnore { } false } - - fn has_git_state(state: &Arc) -> bool { - let mut current = Some(state.as_ref()); - while let Some(frame) = current { - if frame.has_git { - return true; - } - current = frame.parent.as_deref(); - } - false - } } fn handle_read_dir_error( @@ -2522,8 +3278,10 @@ fn is_node_modules_name(name: &OsStr) -> bool { name == OsStr::new("node_modules") } -fn entry_name(name: &OsStr) -> String { - name.to_string_lossy().into_owned() +fn entry_name(name: &OsStr) -> Cow<'_, str> { + name + .to_str() + .map_or_else(|| name.to_string_lossy(), Cow::Borrowed) } #[cfg(unix)] @@ -2545,16 +3303,11 @@ fn is_hidden_name(name: &OsStr) -> bool { .is_some_and(|value| value.as_bytes().first() == Some(&b'.')) } -fn join_relative_path(parent: &str, name: &str) -> String { - if parent.is_empty() { - name.to_string() - } else { - let mut path = String::with_capacity(parent.len() + 1 + name.len()); - path.push_str(parent); - path.push('/'); - path.push_str(name); - path +fn push_relative_name(relative: &mut String, name: &str) { + if !relative.is_empty() { + relative.push('/'); } + relative.push_str(name); } fn mtime_millis(seconds: i64, nanos: i64) -> Option { @@ -2579,6 +3332,11 @@ mod platform { FileType, RawDirEntry, ReadDirControl, ReadDirError, WalkDetail, WalkError, mtime_millis, }; + /// `getattrlistbulk` can return data length in the same batch, but + /// requesting full-detail attributes (size + mtime) measurably slows the + /// bulk scan (~+50% walk time on APFS), which outweighs saving one fstat + /// per opened file. Benchmarked via `perf_walk_collect_full_detail` vs + /// minimal detail. pub const CHEAP_SIZE_HINTS: bool = false; const BUFFER_SIZE: usize = 256 * 1024; @@ -2598,6 +3356,7 @@ mod platform { pub fn read_dir_entries( path: &Path, detail: WalkDetail, + buffer: &mut Vec, mut emit: F, ) -> std::result::Result> where @@ -2618,7 +3377,9 @@ mod platform { attrs.fileattr |= libc::ATTR_FILE_DATALENGTH; } - let mut buffer = vec![0u8; BUFFER_SIZE]; + if buffer.len() != BUFFER_SIZE { + buffer.resize(BUFFER_SIZE, 0); + } loop { // SAFETY: `fd` is an open directory descriptor, `attrs` points to a valid // attrlist for the duration of the call, and `buffer` is writable. @@ -2902,13 +3663,16 @@ mod platform { pub fn read_dir_entries( path: &Path, detail: WalkDetail, + buffer: &mut Vec, mut emit: F, ) -> std::result::Result> where F: FnMut(RawDirEntry<'_>) -> std::result::Result>, { let fd = open_dir(path)?; - let mut buffer = vec![0u8; BUFFER_SIZE]; + if buffer.len() != BUFFER_SIZE { + buffer.resize(BUFFER_SIZE, 0); + } loop { // SAFETY: `fd` is an open directory descriptor and `buffer` is writable. let read = unsafe { @@ -3184,13 +3948,16 @@ mod platform { pub fn read_dir_entries( path: &Path, detail: WalkDetail, + buffer: &mut Vec, mut emit: F, ) -> std::result::Result> where F: FnMut(RawDirEntry<'_>) -> std::result::Result>, { let handle = open_dir(path)?; - let mut buffer = vec![0u8; BUFFER_SIZE]; + if buffer.len() != BUFFER_SIZE { + buffer.resize(BUFFER_SIZE, 0); + } let mut restart = true; loop { @@ -3325,6 +4092,7 @@ mod platform { pub fn read_dir_entries( path: &Path, detail: WalkDetail, + _buffer: &mut Vec, mut emit: F, ) -> std::result::Result> where diff --git a/crates/pi-walker/tests/parallel.rs b/crates/pi-walker/tests/parallel.rs new file mode 100644 index 000000000..607d30916 --- /dev/null +++ b/crates/pi-walker/tests/parallel.rs @@ -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) { + 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) -> 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 { + let mut paths = request + .collect_file_candidates() + .expect("serial candidate collection should succeed") + .into_iter() + .map(|candidate| candidate.relative) + .collect::>(); + paths.sort_unstable(); + paths +} + +fn sorted_parallel_candidates(request: &WalkRequest) -> Vec { + 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 { + let mut paths = WalkRequest::new(root) + .order(order) + .collect() + .expect("collection should succeed") + .entries + .into_iter() + .map(|entry| entry.path) + .collect::>(); + paths.sort_unstable(); + paths +} + +#[derive(Default)] +struct DirectoryOrderVisitor { + children_by_parent: BTreeMap>>, +} + +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 { + 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 { + name.as_bytes().to_vec() +} + +#[cfg(not(unix))] +fn sort_key(name: &std::ffi::OsStr) -> Vec { + 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" + ); +} diff --git a/crates/pi-walker/tests/perf.rs b/crates/pi-walker/tests/perf.rs new file mode 100644 index 000000000..3efebd0d6 --- /dev/null +++ b/crates/pi-walker/tests/perf.rs @@ -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 = 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 { + 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 +} diff --git a/packages/natives/CHANGELOG.md b/packages/natives/CHANGELOG.md index 0cb238263..c0d328048 100644 --- a/packages/natives/CHANGELOG.md +++ b/packages/natives/CHANGELOG.md @@ -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