fix(natives): bounded sorted glob scans
Kept uncached sort-by-mtime glob traversal bounded to maxResults and emitted onMatch callbacks only for returned matches so broad find scans cannot grow parent memory independently of the limit. Fixes #1761
This commit is contained in:
@@ -403,7 +403,11 @@ fn collect_entries(
|
||||
Ok(entries)
|
||||
}
|
||||
|
||||
fn collect_entry(root: &Path, entry: &ignore::DirEntry, detail: ScanDetail) -> Option<GlobMatch> {
|
||||
pub(crate) fn collect_entry(
|
||||
root: &Path,
|
||||
entry: &ignore::DirEntry,
|
||||
detail: ScanDetail,
|
||||
) -> Option<GlobMatch> {
|
||||
let path = entry.path();
|
||||
let relative = normalize_relative_path(root, path);
|
||||
if relative.is_empty() {
|
||||
|
||||
+138
-25
@@ -14,7 +14,7 @@
|
||||
//! // JS: await native.glob({ pattern: "*.rs", path: "." })
|
||||
//! ```
|
||||
|
||||
use std::path::Path;
|
||||
use std::{cmp::Ordering, collections::BinaryHeap, path::Path};
|
||||
|
||||
use globset::GlobSet;
|
||||
use napi::{
|
||||
@@ -81,6 +81,66 @@ struct GlobConfig {
|
||||
use_cache: bool,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
struct RankedGlobMatch {
|
||||
entry: GlobMatch,
|
||||
}
|
||||
|
||||
impl PartialEq for RankedGlobMatch {
|
||||
fn eq(&self, other: &Self) -> bool {
|
||||
compare_matches_by_rank(&self.entry, &other.entry) == Ordering::Equal
|
||||
}
|
||||
}
|
||||
|
||||
impl Eq for RankedGlobMatch {}
|
||||
|
||||
impl PartialOrd for RankedGlobMatch {
|
||||
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
|
||||
Some(self.cmp(other))
|
||||
}
|
||||
}
|
||||
|
||||
impl Ord for RankedGlobMatch {
|
||||
fn cmp(&self, other: &Self) -> Ordering {
|
||||
if match_is_worse(&self.entry, &other.entry) {
|
||||
Ordering::Greater
|
||||
} else if match_is_worse(&other.entry, &self.entry) {
|
||||
Ordering::Less
|
||||
} else {
|
||||
Ordering::Equal
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn match_mtime(entry: &GlobMatch) -> f64 {
|
||||
entry.mtime.unwrap_or(0.0)
|
||||
}
|
||||
|
||||
fn compare_matches_by_rank(a: &GlobMatch, b: &GlobMatch) -> Ordering {
|
||||
match_mtime(b)
|
||||
.total_cmp(&match_mtime(a))
|
||||
.then_with(|| a.path.cmp(&b.path))
|
||||
}
|
||||
|
||||
fn match_is_worse(a: &GlobMatch, b: &GlobMatch) -> bool {
|
||||
compare_matches_by_rank(a, b) == Ordering::Greater
|
||||
}
|
||||
|
||||
fn push_bounded_match(heap: &mut BinaryHeap<RankedGlobMatch>, entry: GlobMatch, limit: usize) {
|
||||
if heap.len() < limit {
|
||||
heap.push(RankedGlobMatch { entry });
|
||||
return;
|
||||
}
|
||||
|
||||
let Some(worst) = heap.peek() else {
|
||||
return;
|
||||
};
|
||||
if match_is_worse(&worst.entry, &entry) {
|
||||
heap.pop();
|
||||
heap.push(RankedGlobMatch { entry });
|
||||
}
|
||||
}
|
||||
|
||||
fn resolve_symlink_target_type(root: &Path, relative_path: &str) -> Option<FileType> {
|
||||
let target_path = root.join(relative_path);
|
||||
let metadata = std::fs::metadata(target_path).ok()?;
|
||||
@@ -143,7 +203,9 @@ fn filter_entries(
|
||||
};
|
||||
let mut matched_entry = entry.clone();
|
||||
matched_entry.file_type = effective_file_type;
|
||||
if let Some(callback) = on_match {
|
||||
if !config.sort_by_mtime
|
||||
&& let Some(callback) = on_match
|
||||
{
|
||||
callback.call(Ok(matched_entry.clone()), ThreadsafeFunctionCallMode::NonBlocking);
|
||||
}
|
||||
|
||||
@@ -156,6 +218,54 @@ fn filter_entries(
|
||||
Ok(matches)
|
||||
}
|
||||
|
||||
fn collect_sorted_matches_uncached(
|
||||
glob_set: &GlobSet,
|
||||
config: &GlobConfig,
|
||||
ct: &task::CancelToken,
|
||||
) -> Result<Vec<GlobMatch>> {
|
||||
let builder = fs_cache::build_walker(
|
||||
&config.root,
|
||||
config.include_hidden,
|
||||
config.use_gitignore,
|
||||
!config.mentions_node_modules,
|
||||
false,
|
||||
);
|
||||
let mut top_matches = BinaryHeap::with_capacity(config.max_results.min(1024));
|
||||
let mut visited = 0usize;
|
||||
|
||||
for entry in builder.build() {
|
||||
if visited == 0 || visited >= 128 {
|
||||
visited = 0;
|
||||
ct.heartbeat()?;
|
||||
}
|
||||
visited += 1;
|
||||
|
||||
let Ok(entry) = entry else {
|
||||
continue;
|
||||
};
|
||||
let Some(mut matched_entry) =
|
||||
fs_cache::collect_entry(&config.root, &entry, fs_cache::ScanDetail::Full)
|
||||
else {
|
||||
continue;
|
||||
};
|
||||
if fs_cache::should_skip_path(Path::new(&matched_entry.path), config.mentions_node_modules) {
|
||||
continue;
|
||||
}
|
||||
if !glob_set.is_match(&matched_entry.path) {
|
||||
continue;
|
||||
}
|
||||
let Some(effective_file_type) = apply_file_type_filter(&matched_entry, config) else {
|
||||
continue;
|
||||
};
|
||||
matched_entry.file_type = effective_file_type;
|
||||
push_bounded_match(&mut top_matches, matched_entry, config.max_results);
|
||||
}
|
||||
|
||||
let mut matches: Vec<GlobMatch> = top_matches.into_iter().map(|ranked| ranked.entry).collect();
|
||||
matches.sort_by(compare_matches_by_rank);
|
||||
Ok(matches)
|
||||
}
|
||||
|
||||
/// Executes matching/filtering over scanned entries and optionally streams each
|
||||
/// hit.
|
||||
fn run_glob(
|
||||
@@ -180,31 +290,33 @@ fn run_glob(
|
||||
fs_cache::ScanDetail::Minimal
|
||||
},
|
||||
};
|
||||
let mut matches = if config.use_cache {
|
||||
let scan = fs_cache::get_or_scan(&config.root, scan_options, &ct)?;
|
||||
let mut matches = filter_entries(&scan.entries, &glob_set, &config, on_match, &ct)?;
|
||||
// Empty-result recheck: if we got zero matches from a cached scan that's old
|
||||
// enough, force a rescan and try once more before returning empty.
|
||||
if matches.is_empty() && scan.cache_age_ms >= fs_cache::empty_recheck_ms() {
|
||||
let fresh = fs_cache::force_rescan(&config.root, scan_options, true, &ct)?;
|
||||
matches = filter_entries(&fresh, &glob_set, &config, on_match, &ct)?;
|
||||
}
|
||||
matches
|
||||
} else {
|
||||
let fresh = fs_cache::force_rescan(&config.root, scan_options, false, &ct)?;
|
||||
filter_entries(&fresh, &glob_set, &config, on_match, &ct)?
|
||||
};
|
||||
let mut matches =
|
||||
if config.sort_by_mtime && !config.use_cache && config.max_results != usize::MAX {
|
||||
collect_sorted_matches_uncached(&glob_set, &config, &ct)?
|
||||
} else if config.use_cache {
|
||||
let scan = fs_cache::get_or_scan(&config.root, scan_options, &ct)?;
|
||||
let mut matches = filter_entries(&scan.entries, &glob_set, &config, on_match, &ct)?;
|
||||
// Empty-result recheck: if we got zero matches from a cached scan that's old
|
||||
// enough, force a rescan and try once more before returning empty.
|
||||
if matches.is_empty() && scan.cache_age_ms >= fs_cache::empty_recheck_ms() {
|
||||
let fresh = fs_cache::force_rescan(&config.root, scan_options, true, &ct)?;
|
||||
matches = filter_entries(&fresh, &glob_set, &config, on_match, &ct)?;
|
||||
}
|
||||
matches
|
||||
} else {
|
||||
let fresh = fs_cache::force_rescan(&config.root, scan_options, false, &ct)?;
|
||||
filter_entries(&fresh, &glob_set, &config, on_match, &ct)?
|
||||
};
|
||||
|
||||
if config.sort_by_mtime {
|
||||
// Sorting mode: rank by mtime descending, then apply max-results truncation.
|
||||
matches.sort_by(|a, b| {
|
||||
let a_mtime = a.mtime.unwrap_or(0.0);
|
||||
let b_mtime = b.mtime.unwrap_or(0.0);
|
||||
b_mtime
|
||||
.partial_cmp(&a_mtime)
|
||||
.unwrap_or(std::cmp::Ordering::Equal)
|
||||
});
|
||||
matches.sort_by(compare_matches_by_rank);
|
||||
matches.truncate(config.max_results);
|
||||
if let Some(callback) = on_match {
|
||||
for matched_entry in &matches {
|
||||
callback.call(Ok(matched_entry.clone()), ThreadsafeFunctionCallMode::NonBlocking);
|
||||
}
|
||||
}
|
||||
}
|
||||
let total_matches = matches.len().min(u32::MAX as usize) as u32;
|
||||
Ok(GlobResult { matches, total_matches })
|
||||
@@ -215,8 +327,9 @@ fn run_glob(
|
||||
/// Resolves the search root, scans entries, applies glob and optional file-type
|
||||
/// filters, and optionally streams each accepted match through `on_match`.
|
||||
///
|
||||
/// If `sortByMtime` is enabled, all matching entries are collected, sorted by
|
||||
/// descending mtime, then truncated to `maxResults`.
|
||||
/// If `sortByMtime` is enabled with a finite `maxResults`, uncached scans keep
|
||||
/// only the current top results while traversing instead of collecting the full
|
||||
/// tree.
|
||||
///
|
||||
/// # Errors
|
||||
/// Returns an error when the search path cannot be resolved, the path is not a
|
||||
|
||||
@@ -368,8 +368,8 @@ export class FindTool implements AgentTool<typeof findSchema, FindToolDetails> {
|
||||
let timedOut = false;
|
||||
try {
|
||||
const result = await doGlob(useGitignore);
|
||||
// Sort by mtime descending (most recent first) in JS instead of native.
|
||||
// This allows native glob to early-terminate at maxResults.
|
||||
// Native glob returns a bounded mtime-ranked set; keep the JS sort for
|
||||
// deterministic ordering across cached and uncached native paths.
|
||||
result.matches.sort((a, b) => (b.mtime ?? 0) - (a.mtime ?? 0));
|
||||
matches = result.matches;
|
||||
} catch (error) {
|
||||
|
||||
@@ -2,6 +2,10 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed
|
||||
|
||||
- Bounded sorted `glob()` scans to `maxResults` during uncached traversal and capped `onMatch` callbacks to returned matches so broad OMP `find` scans cannot grow parent-process memory independently of the requested limit ([#1761](https://github.com/can1357/oh-my-pi/issues/1761)).
|
||||
|
||||
## [15.7.0] - 2026-05-31
|
||||
|
||||
### Added
|
||||
|
||||
Vendored
+3
-2
@@ -544,8 +544,9 @@ export declare function getWorkProfile(lastSeconds: number): WorkProfile
|
||||
* Resolves the search root, scans entries, applies glob and optional file-type
|
||||
* filters, and optionally streams each accepted match through `on_match`.
|
||||
*
|
||||
* If `sortByMtime` is enabled, all matching entries are collected, sorted by
|
||||
* descending mtime, then truncated to `maxResults`.
|
||||
* If `sortByMtime` is enabled with a finite `maxResults`, uncached scans keep
|
||||
* only the current top results while traversing instead of collecting the full
|
||||
* tree.
|
||||
*
|
||||
* # Errors
|
||||
* Returns an error when the search path cannot be resolved, the path is not a
|
||||
|
||||
@@ -407,6 +407,38 @@ describe("pi-natives", () => {
|
||||
expect(result.matches).toHaveLength(0);
|
||||
});
|
||||
|
||||
it("should bound sorted callbacks to maxResults", async () => {
|
||||
const scopedDir = await fs.mkdtemp(path.join(os.tmpdir(), "natives-glob-limit-"));
|
||||
try {
|
||||
for (let i = 0; i < 40; i++) {
|
||||
await fs.writeFile(path.join(scopedDir, `file-${String(i).padStart(2, "0")}.txt`), `${i}\n`);
|
||||
}
|
||||
|
||||
const streamedPaths: string[] = [];
|
||||
const result = await glob(
|
||||
{
|
||||
pattern: "**/*",
|
||||
path: scopedDir,
|
||||
hidden: true,
|
||||
gitignore: false,
|
||||
sortByMtime: true,
|
||||
maxResults: 5,
|
||||
},
|
||||
(error, match) => {
|
||||
if (error) throw error;
|
||||
if (match?.path) streamedPaths.push(match.path);
|
||||
},
|
||||
);
|
||||
|
||||
await Bun.sleep(10);
|
||||
expect(result.matches).toHaveLength(5);
|
||||
expect(streamedPaths).toHaveLength(5);
|
||||
expect(new Set(streamedPaths)).toEqual(new Set(result.matches.map(match => match.path)));
|
||||
} finally {
|
||||
await fs.rm(scopedDir, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
|
||||
it("should fast-recheck empty cached results when threshold is reached", async () => {
|
||||
const fileName = "cache-empty-recheck-target.txt";
|
||||
const filePath = path.join(testDir, fileName);
|
||||
|
||||
Reference in New Issue
Block a user