refactor: migrated concurrency primitives to flume and parking_lot

- Replaced standard and tokio mpsc channels with flume channels across workspace crates to simplify thread synchronization.
- Swapped standard Mutex guards for parking_lot Mutexes to avoid manual lock poisoning handling and improve performance.
- Declared workspace-wide dependencies for flume and parking_lot in root and member Cargo manifests.
This commit is contained in:
can1357
2026-06-30 23:45:42 +02:00
parent daf3aaa508
commit 13b1b5b134
22 changed files with 180 additions and 183 deletions
+1
View File
@@ -13,6 +13,7 @@ workspace = true
dashmap.workspace = true
globset.workspace = true
ignore.workspace = true
parking_lot.workspace = true
rayon.workspace = true
[target.'cfg(unix)'.dependencies]
+6 -16
View File
@@ -4,12 +4,13 @@ use std::{
borrow::Cow,
fmt,
path::{Path, PathBuf},
sync::{Arc, LazyLock, Mutex},
sync::{Arc, LazyLock},
time::{Duration, Instant},
};
use dashmap::DashMap;
use ignore::{ParallelVisitor, ParallelVisitorBuilder, WalkBuilder, WalkState};
use parking_lot::Mutex;
use rayon::{ThreadPool, prelude::*};
use crate::{
@@ -280,7 +281,6 @@ fn build_walker_for_options_inner(
if let Some(pruned_dirs) = &pruned_dirs
&& pruned_dirs
.lock()
.expect("pruned directory lock poisoned")
.iter()
.any(|dir| entry.path().starts_with(dir))
{
@@ -384,11 +384,7 @@ impl<H> Drop for EntryVisitor<'_, H> {
return;
}
let entries = std::mem::take(&mut self.entries);
self
.shared_entries
.lock()
.expect("entry collection lock poisoned")
.push(entries);
self.shared_entries.lock().push(entries);
}
}
@@ -401,7 +397,7 @@ where
if self.visited == 0 || self.visited >= 128 {
self.visited = 0;
if let Err(err) = (self.heartbeat)() {
*self.error.lock().expect("error lock poisoned") = Some(err.to_string());
*self.error.lock() = Some(err.to_string());
return WalkState::Quit;
}
}
@@ -489,18 +485,12 @@ where
heartbeat().map_err(|err| WalkError::Interrupted(err.to_string()))?;
builder.build_parallel().visit(&mut visitor_builder);
let walk_error = error.lock().expect("error lock poisoned").take();
let walk_error = error.lock().take();
if let Some(error) = walk_error {
return Err(WalkError::Interrupted(error));
}
entries.extend(
shared_entries
.lock()
.expect("entry collection lock poisoned")
.drain(..)
.flatten(),
);
entries.extend(shared_entries.lock().drain(..).flatten());
entries.sort_unstable_by(|a, b| a.path.cmp(&b.path));
Ok(EntryScan::Entries(CollectedEntries { entries, cache_age_ms: 0 }))
}
+4 -10
View File
@@ -17,7 +17,7 @@ use std::{
hash::{Hash, Hasher},
io,
path::{Path, PathBuf},
sync::{Arc, Mutex},
sync::Arc,
};
pub use cache::{
@@ -26,6 +26,7 @@ pub use cache::{
parallel_for_each, resolve_search_path, should_parallelize, should_skip_path, walk_workers,
};
use globset::{GlobBuilder, GlobSet, GlobSetBuilder};
use parking_lot::Mutex;
const HEARTBEAT_INTERVAL: usize = 128;
@@ -2307,10 +2308,7 @@ where
WalkControl::Quit => return Ok(WalkStatus::Stopped),
WalkControl::SkipDescend => {
if collected.file_type == FileType::Dir {
pruned_dirs
.lock()
.expect("pruned directory lock poisoned")
.push(entry.path().to_path_buf());
pruned_dirs.lock().push(entry.path().to_path_buf());
}
},
WalkControl::Continue => {},
@@ -2361,11 +2359,7 @@ fn ignore_error_to_io(error: &ignore::Error) -> io::Error {
}
fn is_pruned_path(path: &Path, pruned_dirs: &Arc<Mutex<Vec<PathBuf>>>) -> bool {
pruned_dirs
.lock()
.expect("pruned directory lock poisoned")
.iter()
.any(|dir| path.starts_with(dir))
pruned_dirs.lock().iter().any(|dir| path.starts_with(dir))
}
/// Return whether [`WalkDetail::Full`] provides file sizes without per-entry