use std::{
any::Any,
collections::HashSet,
panic::{AssertUnwindSafe, catch_unwind, resume_unwind},
path::{Path, PathBuf},
sync::{
Arc, Mutex,
atomic::{AtomicBool, AtomicUsize, Ordering},
},
thread,
};
use crossbeam_deque::{Steal, Stealer, Worker};
use super::{
CANCELLATION_STRIDE, CYCLE_KEY_OPERATION, CycleKey, DirectoryBackend, EntryVisitor,
ErrorPolicy, Listing, Verdict, WalkEntry, WalkError, WalkResult, Walker,
classify::{DirectoryTask, EmittedEntry, EntryAction, classify_entry},
own_path, reset_to_directory,
scheduler::{CacheLine, Coordinator, Scheduler, WorkerSlot},
};
pub(super) fn collect<B: DirectoryBackend + Sync>(
walker: Walker,
backend: &B,
visitor: EntryVisitor<'_>,
) -> Result<WalkResult, WalkError> {
let walker = Arc::new(walker);
let shared = Arc::new(Shared::new(Arc::clone(&walker), backend, visitor));
let mut caller = WorkerScratch::new(0);
let caller_slot = shared.coordinator.claim_caller_slot();
let roots = walker.root_tasks(backend);
for root in roots {
shared.coordinator.begin_task();
let _root_task = shared.coordinator.claim_task();
process_directory(&shared, &mut caller, root);
}
caller.flush_into(&shared.scheduler);
let widened =
catch_worker_panic(&shared, || drain_alone(&shared, &mut caller)).unwrap_or(false);
if !widened {
drop(caller_slot);
return finish(shared, std::mem::take(&mut caller.entries));
}
shared.widen();
let spare = (1..walker.threads.max(1))
.map(WorkerScratch::new)
.collect::<Vec<_>>();
let mut stealers = Vec::with_capacity(spare.len() + 1);
stealers.push(caller.queue.stealer());
stealers.extend(spare.iter().map(|worker| worker.queue.stealer()));
let idle = Mutex::new(spare);
let shards = Mutex::new(Vec::new());
let mut entries = thread::scope(|scope| {
let pool = HelperPool {
scope,
shared: &shared,
stealers: &stealers,
idle: &idle,
shards: &shards,
};
pool.grow();
run_worker_catching_panics(pool, &mut caller);
std::mem::take(&mut caller.entries)
});
for shard in lock(&shards).drain(..) {
entries.extend(shard);
}
drop(caller_slot);
finish(shared, entries)
}
#[derive(Clone, Copy)]
struct HelperPool<'scope, 'env> {
scope: &'scope thread::Scope<'scope, 'env>,
shared: &'env Arc<Shared<'env>>,
stealers: &'env [Stealer<DirectoryTask>],
idle: &'env Mutex<Vec<WorkerScratch>>,
shards: &'env Mutex<Vec<Vec<WalkEntry>>>,
}
const HELPER_QUEUE_FLOOR: usize = 8;
const DIRECTORY_WEIGHT: usize = 20;
const HELPER_WORK_FLOOR: usize = 224;
const HELPER_LISTING_FLOOR: usize = 1024;
impl<'scope, 'env> HelperPool<'scope, 'env> {
fn grow(self) {
while !self.shared.should_stop() {
if !self.shared.tree_is_worth_helpers() {
return;
}
let Some(slot) = self
.shared
.coordinator
.claim_worker_slot(self.shared.walker.threads)
else {
return;
};
let Some(scratch) = lock(self.idle).pop() else {
return;
};
if !self.spawn(slot, scratch) {
return;
}
}
}
fn spawn(self, slot: WorkerSlot<'env>, scratch: WorkerScratch) -> bool {
#[cfg(test)]
if should_fail_next_worker_spawn() {
self.shared
.record_startup_error(std::io::Error::other("injected worker start failure"));
lock(self.idle).push(scratch);
return false;
}
let spawn = thread::Builder::new()
.name("ferralk-worker".into())
.spawn_scoped(self.scope, move || {
let _slot = slot;
let mut scratch = scratch;
run_worker_catching_panics(self, &mut scratch);
lock(self.shards).push(std::mem::take(&mut scratch.entries));
});
match spawn {
Ok(_) => true,
Err(source) => {
self.shared.record_startup_error(source);
false
}
}
}
}
fn finish(shared: Arc<Shared>, mut entries: Vec<WalkEntry>) -> Result<WalkResult, WalkError> {
if let Some(payload) = lock(&shared.panic).take() {
resume_unwind(payload);
}
if let Some(error) = lock(&shared.startup_error).take() {
return Err(error);
}
if let Some(error) = lock(&shared.abort_error).take() {
return Err(error);
}
if shared.walker.options.sort {
entries.sort_by(|left, right| left.path.cmp(&right.path));
}
Ok(WalkResult {
entries,
errors: std::mem::take(&mut *lock(&shared.errors)),
cancelled: shared.should_stop(),
})
}
struct Shared<'backend> {
walker: Arc<Walker>,
visitor: EntryVisitor<'backend>,
stopped: AtomicBool,
lone: AtomicBool,
work_seen: CacheLine<AtomicUsize>,
backend: &'backend (dyn DirectoryBackend + Sync),
scheduler: Scheduler<DirectoryTask>,
coordinator: Coordinator,
cancellation: super::CancellationToken,
errors: Mutex<Vec<WalkError>>,
abort_error: Mutex<Option<WalkError>>,
startup_error: Mutex<Option<WalkError>>,
panic: Mutex<Option<Box<dyn Any + Send + 'static>>>,
visited_directories: Mutex<HashSet<CycleKey>>,
}
impl<'backend> Shared<'backend> {
fn new(
walker: Arc<Walker>,
backend: &'backend (dyn DirectoryBackend + Sync),
visitor: EntryVisitor<'backend>,
) -> Self {
let cancellation = walker.cancellation.clone().unwrap_or_default();
Self {
walker,
visitor,
stopped: AtomicBool::new(false),
lone: AtomicBool::new(true),
work_seen: CacheLine(AtomicUsize::new(0)),
backend,
scheduler: Scheduler::new(),
coordinator: Coordinator::new(),
cancellation,
errors: Mutex::new(Vec::new()),
abort_error: Mutex::new(None),
startup_error: Mutex::new(None),
panic: Mutex::new(None),
visited_directories: Mutex::new(HashSet::new()),
}
}
fn schedule(&self, worker: &Worker<DirectoryTask>, task: DirectoryTask) {
self.coordinator.begin_task();
worker.push(task);
if !self.lone.load(Ordering::Relaxed) {
self.coordinator.wake_one_waiter();
}
}
fn widen(&self) {
self.lone.store(false, Ordering::Relaxed);
}
fn tree_is_worth_helpers(&self) -> bool {
let work = self.work_seen.load(Ordering::Acquire);
(work >= HELPER_WORK_FLOOR && self.coordinator.pending() >= HELPER_QUEUE_FLOOR)
|| work >= HELPER_LISTING_FLOOR
}
fn should_stop(&self) -> bool {
self.cancellation.is_cancelled() || self.stopped.load(Ordering::Acquire)
}
fn emit(&self, worker: &mut WorkerScratch, emitted: EmittedEntry) {
let WorkerScratch {
entries,
path,
spare,
..
} = worker;
let entry = emitted.with_path(own_path(spare, path));
match (self.visitor)(&entry) {
Verdict::Keep => entries.push(entry),
Verdict::Skip => *spare = entry.path,
Verdict::Stop => {
*spare = entry.path;
self.stopped.store(true, Ordering::Release);
self.coordinator.wake_waiters();
}
}
}
fn record_error(&self, operation: &'static str, path: PathBuf, source: std::io::Error) {
let error = WalkError::new(operation, path, source);
match self.walker.error_policy {
ErrorPolicy::Abort => {
let mut abort_error = lock(&self.abort_error);
if abort_error.is_none() {
*abort_error = Some(error);
self.cancellation.cancel();
self.coordinator.wake_waiters();
}
}
ErrorPolicy::Skip => {}
ErrorPolicy::Collect => lock(&self.errors).push(error),
}
}
fn record_startup_error(&self, source: std::io::Error) {
let mut startup_error = lock(&self.startup_error);
if startup_error.is_none() {
*startup_error = Some(WalkError::new(
"spawn_worker",
self.walker
.roots()
.next()
.expect("a walk has a root")
.into(),
source,
));
self.cancellation.cancel();
self.coordinator.wake_waiters();
}
}
fn record_panic(&self, payload: Box<dyn Any + Send + 'static>) {
let mut panic = lock(&self.panic);
if panic.is_none() {
*panic = Some(payload);
self.cancellation.cancel();
self.coordinator.wake_waiters();
}
}
}
struct WorkerScratch {
index: usize,
queue: Worker<DirectoryTask>,
entries: Vec<WalkEntry>,
listing: Listing,
path: PathBuf,
spare: PathBuf,
}
impl WorkerScratch {
fn new(index: usize) -> Self {
Self {
index,
queue: Worker::new_fifo(),
entries: Vec::new(),
listing: Listing::default(),
path: PathBuf::new(),
spare: PathBuf::new(),
}
}
fn flush_into(&self, scheduler: &Scheduler<DirectoryTask>) {
while let Some(task) = self.queue.pop() {
scheduler.push(task);
}
}
}
fn run_worker_catching_panics(pool: HelperPool<'_, '_>, worker: &mut WorkerScratch) {
catch_worker_panic(pool.shared, || run_worker(pool, worker));
}
fn catch_worker_panic<T>(shared: &Shared, work: impl FnOnce() -> T) -> Option<T> {
match catch_unwind(AssertUnwindSafe(work)) {
Ok(value) => Some(value),
Err(payload) => {
shared.record_panic(payload);
None
}
}
}
fn drain_alone(shared: &Shared, caller: &mut WorkerScratch) -> bool {
loop {
let has_work = !caller.queue.is_empty() || !shared.scheduler.is_empty();
if has_work && shared.tree_is_worth_helpers() {
return true;
}
let Some(directory) = caller
.queue
.pop()
.or_else(|| shared.scheduler.steal_into(&caller.queue))
else {
return false;
};
let _task = shared.coordinator.claim_task();
if shared.should_stop() {
return false;
}
process_directory(shared, caller, directory);
}
}
fn run_worker(pool: HelperPool<'_, '_>, worker: &mut WorkerScratch) {
let shared = pool.shared;
while let Some(directory) = next_task(shared, worker, pool.stealers) {
let _task = shared.coordinator.claim_task();
if !shared.should_stop() {
process_directory(shared, worker, directory);
pool.grow();
}
#[cfg(test)]
join_worker_rendezvous(shared);
}
}
fn next_task(
shared: &Shared,
worker: &WorkerScratch,
stealers: &[Stealer<DirectoryTask>],
) -> Option<DirectoryTask> {
shared.coordinator.wait_for_task(
|| shared.should_stop(),
|| try_take(shared, worker, stealers),
)
}
fn try_take(
shared: &Shared,
worker: &WorkerScratch,
stealers: &[Stealer<DirectoryTask>],
) -> Option<DirectoryTask> {
worker
.queue
.pop()
.or_else(|| shared.scheduler.steal_into(&worker.queue))
.or_else(|| {
let (head, tail) = stealers.split_at((worker.index + 1).min(stealers.len()));
tail.iter()
.chain(head.iter().take(worker.index))
.find_map(|stealer| {
loop {
match stealer.steal_batch_and_pop(&worker.queue) {
Steal::Success(task) => break Some(task),
Steal::Empty => break None,
Steal::Retry => continue,
}
}
})
})
}
fn process_directory(shared: &Shared, worker: &mut WorkerScratch, task: DirectoryTask) {
let DirectoryTask {
path,
depth,
root,
ignores,
} = task;
#[cfg(test)]
if should_panic_in_directory(&path) {
panic!("injected directory panic");
}
if shared.should_stop() {
return;
}
if shared.walker.options.follow_symlinks && !mark_directory(shared, &path) {
return;
}
if let Err(source) = shared.backend.read_directory(&path, &mut worker.listing) {
shared.record_error("read_dir", path, source);
return;
}
shared.work_seen.fetch_add(
worker.listing.entries().len() + DIRECTORY_WEIGHT,
Ordering::AcqRel,
);
let ignores = ignores.enter(&shared.walker, shared.backend, &path, &worker.listing);
worker.path.clear();
worker.path.push(&path);
for index in 0..worker.listing.entries().len() {
if index.is_multiple_of(CANCELLATION_STRIDE) && shared.should_stop() {
return;
}
worker.path.push(worker.listing.entries()[index].name());
let action = classify_entry(
&shared.walker,
shared.backend,
&worker.path,
&worker.listing.entries()[index],
&ignores,
depth,
root,
);
act(shared, worker, action);
reset_to_directory(&mut worker.path, &path);
}
}
fn mark_directory(shared: &Shared, directory: &Path) -> bool {
match shared.backend.cycle_key(directory) {
Ok(key) => lock(&shared.visited_directories).insert(key),
Err(source) => {
shared.record_error(CYCLE_KEY_OPERATION, directory.to_path_buf(), source);
false
}
}
}
fn act(shared: &Shared, worker: &mut WorkerScratch, action: EntryAction) {
match action {
EntryAction::Skip => {}
EntryAction::Descend(task) => shared.schedule(&worker.queue, task),
EntryAction::DescendAndEmit(entry, task) => {
shared.schedule(&worker.queue, task);
shared.emit(worker, entry);
}
EntryAction::Emit(entry) => shared.emit(worker, entry),
EntryAction::Failed { failure, descend } => {
if let Some(task) = descend {
shared.schedule(&worker.queue, task);
}
shared.record_error(failure.operation, failure.path, failure.source);
}
}
}
fn lock<T>(mutex: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
mutex
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
#[cfg(test)]
std::thread_local! {
static FAIL_NEXT_WORKER_SPAWN: std::cell::Cell<bool> = const { std::cell::Cell::new(false) };
}
#[cfg(test)]
fn fail_next_worker_spawn() {
FAIL_NEXT_WORKER_SPAWN.with(|failure| failure.set(true));
}
#[cfg(test)]
fn should_fail_next_worker_spawn() -> bool {
FAIL_NEXT_WORKER_SPAWN.with(std::cell::Cell::take)
}
#[cfg(test)]
static PANIC_IN_DIRECTORY: Mutex<Option<PathBuf>> = Mutex::new(None);
#[cfg(test)]
fn panic_in_directory(directory: PathBuf) {
*lock(&PANIC_IN_DIRECTORY) = Some(directory);
}
#[cfg(test)]
fn should_panic_in_directory(directory: &std::path::Path) -> bool {
let mut target = lock(&PANIC_IN_DIRECTORY);
if target.as_deref() == Some(directory) {
*target = None;
return true;
}
false
}
#[cfg(test)]
struct WorkerRendezvous {
root: PathBuf,
expected: usize,
threads: HashSet<std::thread::ThreadId>,
released: bool,
}
#[cfg(test)]
static WORKER_RENDEZVOUS: Mutex<Option<WorkerRendezvous>> = Mutex::new(None);
#[cfg(test)]
static WORKER_RENDEZVOUS_GUARD: Mutex<()> = Mutex::new(());
#[cfg(test)]
static WORKER_RENDEZVOUS_WAKE: std::sync::Condvar = std::sync::Condvar::new();
#[cfg(test)]
const WORKER_RENDEZVOUS_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
#[cfg(test)]
const WORKER_RENDEZVOUS_POLL: std::time::Duration = std::time::Duration::from_millis(10);
#[cfg(test)]
fn expect_worker_threads(root: PathBuf, expected: usize) {
*lock(&WORKER_RENDEZVOUS) = Some(WorkerRendezvous {
root,
expected,
threads: HashSet::new(),
released: false,
});
}
#[cfg(test)]
fn observed_worker_threads() -> usize {
lock(&WORKER_RENDEZVOUS)
.take()
.map_or(0, |rendezvous| rendezvous.threads.len())
}
#[cfg(test)]
fn join_worker_rendezvous(shared: &Shared) {
if !shared.tree_is_worth_helpers() {
return;
}
let mut state = lock(&WORKER_RENDEZVOUS);
match state.as_mut() {
Some(rendezvous)
if shared.walker.roots().next() == Some(rendezvous.root.as_path())
&& !rendezvous.released =>
{
rendezvous.threads.insert(std::thread::current().id());
}
_ => return,
}
let deadline = std::time::Instant::now() + WORKER_RENDEZVOUS_TIMEOUT;
loop {
let Some(rendezvous) = state.as_mut() else {
return;
};
let waited_long_enough = std::time::Instant::now() >= deadline;
if rendezvous.released || rendezvous.threads.len() >= rendezvous.expected {
rendezvous.released = true;
WORKER_RENDEZVOUS_WAKE.notify_all();
return;
}
if waited_long_enough {
rendezvous.released = true;
WORKER_RENDEZVOUS_WAKE.notify_all();
return;
}
let (resumed, _) = WORKER_RENDEZVOUS_WAKE
.wait_timeout(state, WORKER_RENDEZVOUS_POLL)
.unwrap_or_else(|poisoned| poisoned.into_inner());
state = resumed;
}
}
#[cfg(test)]
mod tests {
use std::{
collections::HashSet,
fs,
panic::{AssertUnwindSafe, catch_unwind},
path::{Path, PathBuf},
sync::{
Arc, Mutex,
atomic::{AtomicUsize, Ordering},
mpsc,
},
thread,
time::{Duration, SystemTime, UNIX_EPOCH},
};
use super::{
Shared, WORKER_RENDEZVOUS_GUARD, Walker, catch_worker_panic, expect_worker_threads,
fail_next_worker_spawn, finish, lock, observed_worker_threads, panic_in_directory,
};
use crate::{CancellationToken, WalkOptions, WalkResult};
const RESUME_TIMEOUT: Duration = Duration::from_secs(30);
static NEXT_ROOT: AtomicUsize = AtomicUsize::new(0);
fn unique_root(label: &str) -> PathBuf {
std::env::temp_dir().join(format!(
"ferralk-{label}-{}-{}-{}",
std::process::id(),
SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("system clock is after epoch")
.as_nanos(),
NEXT_ROOT.fetch_add(1, Ordering::Relaxed)
))
}
fn create_wide_fixture(root: &Path) {
for branch in 0..10 {
for nested in 0..4 {
let directory = root
.join(format!("branch-{branch}"))
.join(format!("nested-{nested}"));
fs::create_dir_all(&directory).expect("create fixture directory");
for file in 0..4 {
fs::write(directory.join(format!("file-{file}.txt")), b"fixture")
.expect("write fixture file");
}
}
}
}
fn floor_trips_for(dirs: usize, per_dir: usize) -> bool {
let walker = Arc::new(Walker::new("."));
let shared = Shared::new(
Arc::clone(&walker),
&crate::SystemBackend,
&crate::keep_every_entry,
);
let mut queued = Vec::new();
for _ in 0..dirs {
shared.coordinator.begin_task();
queued.push(shared.coordinator.claim_task());
}
shared
.work_seen
.store(dirs + super::DIRECTORY_WEIGHT, Ordering::Release);
if shared.tree_is_worth_helpers() {
return true;
}
while queued.pop().is_some() {
shared
.work_seen
.fetch_add(per_dir + super::DIRECTORY_WEIGHT, Ordering::AcqRel);
if shared.tree_is_worth_helpers() {
return true;
}
}
false
}
#[test]
fn the_floor_divides_where_the_sweep_says_it_should() {
for (dirs, per_dir, expected, note) in [
(9, 16, false, "pooling loses 10%"),
(17, 16, true, "pooling wins 8.5%"),
(4, 16, false, "pooling loses 38%"),
(8, 16, false, "pooling loses 10%"),
(14, 16, true, "pooling wins 15%"),
(10, 4, false, "pooling loses 15%"),
(14, 4, false, "pooling loses 4%"),
(20, 4, true, "pooling wins 11%"),
(12, 1, false, "pooling loses 8%"),
(24, 1, true, "pooling wins 15%"),
(48, 1, true, "pooling wins 15%"),
] {
assert_eq!(
floor_trips_for(dirs, per_dir),
expected,
"{dirs} directories of {per_dir} files: {note}"
);
}
for (dirs, expected, note) in [
(16, false, "inside the noise"),
(24, true, "pooling wins 9%, and used to be refused"),
(28, true, "pooling wins 26%, and used to be refused"),
(31, true, "pooling wins 5%, and used to be refused"),
(40, true, "pooling wins 30%"),
(96, true, "pooling wins 21%"),
] {
assert_eq!(
floor_trips_for(dirs, 0),
expected,
"{dirs} empty directories: {note}"
);
}
}
#[test]
fn the_floor_weighs_the_tree_and_the_queue_together() {
let walker = Arc::new(Walker::new("."));
let shared = Shared::new(
Arc::clone(&walker),
&crate::SystemBackend,
&crate::keep_every_entry,
);
shared
.work_seen
.store(28 + super::DIRECTORY_WEIGHT, Ordering::Release);
for _ in 0..16 {
shared.coordinator.begin_task();
}
assert!(
!shared.tree_is_worth_helpers(),
"a wide but trivial tree must not start a pool"
);
shared
.work_seen
.store(super::HELPER_WORK_FLOOR, Ordering::Release);
assert!(shared.tree_is_worth_helpers());
let shared = Shared::new(
Arc::clone(&walker),
&crate::SystemBackend,
&crate::keep_every_entry,
);
shared
.work_seen
.store(super::HELPER_LISTING_FLOOR, Ordering::Release);
shared.coordinator.begin_task();
assert!(
shared.tree_is_worth_helpers(),
"a single huge directory is worth helpers even with an empty queue"
);
}
#[test]
fn worker_panic_cancels_siblings_and_resumes_on_the_caller() {
let cancellation = CancellationToken::default();
let walker = Arc::new(Walker::new(".").cancellation(cancellation.clone()));
let shared = Arc::new(Shared::new(
walker,
&crate::SystemBackend,
&crate::keep_every_entry,
));
catch_worker_panic(&shared, || panic!("injected worker panic"));
assert!(cancellation.is_cancelled());
assert!(catch_unwind(AssertUnwindSafe(|| finish(shared, Vec::new()))).is_err());
}
fn walked_paths(result: &WalkResult) -> Vec<PathBuf> {
result
.entries()
.iter()
.map(|entry| entry.path().to_path_buf())
.collect()
}
#[test]
fn a_single_root_subdirectory_still_uses_the_configured_threads() {
let _rendezvous = lock(&WORKER_RENDEZVOUS_GUARD);
let root = unique_root("single-subtree");
create_wide_fixture(&root.join("only"));
let serial = Walker::new(&root)
.threads(1)
.options(WalkOptions::default().sort(true))
.collect()
.expect("serial walk succeeds");
expect_worker_threads(root.clone(), 4);
let parallel = Walker::new(&root)
.threads(4)
.options(WalkOptions::default().sort(true))
.collect()
.expect("parallel walk succeeds");
let observed = observed_worker_threads();
let _ = fs::remove_dir_all(&root);
assert_eq!(
observed, 4,
"a subtree below a single root child must still use the thread budget"
);
assert_eq!(walked_paths(¶llel), walked_paths(&serial));
assert!(parallel.errors().is_empty());
}
#[test]
fn several_roots_are_walked_by_one_pool() {
let _rendezvous = lock(&WORKER_RENDEZVOUS_GUARD);
let root = unique_root("multi-root-pool");
for name in ["alpha", "beta", "gamma"] {
create_wide_fixture(&root.join(name));
}
let roots = ["alpha", "beta", "gamma"].map(|name| root.join(name));
let serial = Walker::new(&roots[0])
.add_root(&roots[1])
.expect("root")
.add_root(&roots[2])
.expect("root")
.threads(1)
.options(WalkOptions::default().sort(true))
.collect()
.expect("serial walk succeeds");
expect_worker_threads(roots[0].clone(), 4);
let parallel = Walker::new(&roots[0])
.add_root(&roots[1])
.expect("root")
.add_root(&roots[2])
.expect("root")
.threads(4)
.options(WalkOptions::default().sort(true))
.collect()
.expect("parallel walk succeeds");
let observed = observed_worker_threads();
let _ = fs::remove_dir_all(&root);
assert_eq!(
observed, 4,
"three roots must share the walk's one thread budget"
);
assert_eq!(walked_paths(¶llel), walked_paths(&serial));
assert!(parallel.errors().is_empty());
}
#[test]
fn the_floor_counts_across_roots() {
let walker = Arc::new(Walker::new("."));
let shared = Shared::new(
Arc::clone(&walker),
&crate::SystemBackend,
&crate::keep_every_entry,
);
shared
.work_seen
.store(12 + super::DIRECTORY_WEIGHT, Ordering::Release);
for _ in 0..3 {
shared.coordinator.begin_task();
}
assert!(
!shared.tree_is_worth_helpers(),
"three tiny roots are still a tiny walk"
);
shared
.work_seen
.store(super::HELPER_WORK_FLOOR, Ordering::Release);
for _ in 0..super::HELPER_QUEUE_FLOOR {
shared.coordinator.begin_task();
}
assert!(
shared.tree_is_worth_helpers(),
"work summed across roots reaches the floor like work under one"
);
}
#[test]
fn a_visitor_runs_on_every_worker_of_the_walk() {
let _rendezvous = lock(&WORKER_RENDEZVOUS_GUARD);
let root = unique_root("visitor-threads");
create_wide_fixture(&root);
expect_worker_threads(root.clone(), 4);
let seen = Mutex::new(HashSet::new());
let result = Walker::new(&root)
.threads(4)
.visit(|_| {
lock(&seen).insert(thread::current().id());
crate::Verdict::Keep
})
.expect("visited walk succeeds");
let workers = observed_worker_threads();
let visitor_threads = lock(&seen).len();
let _ = fs::remove_dir_all(&root);
assert_eq!(workers, 4, "the walk did not reach its thread budget");
assert!(
visitor_threads > 1,
"the visitor ran on {visitor_threads} thread(s) across {workers} workers"
);
assert!(!result.entries().is_empty());
}
#[test]
fn a_panic_below_the_floor_resumes_on_the_caller_and_cancels() {
let root = unique_root("lean-panic");
let only = root.join("only");
fs::create_dir_all(&only).expect("create fixture directory");
for index in 0..4 {
fs::write(only.join(format!("file-{index}.txt")), b"fixture")
.expect("write fixture file");
}
panic_in_directory(only);
let cancellation = CancellationToken::default();
let walk_cancellation = cancellation.clone();
let outcome = catch_unwind(AssertUnwindSafe(|| {
Walker::new(&root)
.threads(4)
.cancellation(walk_cancellation)
.collect()
}));
let _ = fs::remove_dir_all(&root);
assert!(
outcome.is_err(),
"the injected panic must resume on the caller"
);
assert!(
cancellation.is_cancelled(),
"and must cancel the walk on its way out"
);
}
#[test]
fn worker_panic_during_traversal_resumes_without_hanging_the_walk() {
for round in 0..4 {
let root = unique_root("worker-panic");
create_wide_fixture(&root);
panic_in_directory(root.join("branch-3").join("nested-2"));
let cancellation = CancellationToken::default();
let walk_root = root.clone();
let walk_cancellation = cancellation.clone();
let (sender, receiver) = mpsc::channel();
let runner = thread::Builder::new()
.name("ferralk-panic-regression".into())
.spawn(move || {
let outcome = catch_unwind(AssertUnwindSafe(|| {
Walker::new(&walk_root)
.threads(4)
.cancellation(walk_cancellation)
.collect()
}));
let _ = sender.send(outcome.is_err());
})
.expect("spawn the walking thread");
let resumed = receiver.recv_timeout(RESUME_TIMEOUT);
let _ = fs::remove_dir_all(&root);
assert_eq!(
resumed,
Ok(true),
"round {round}: the injected worker panic must resume on the caller"
);
runner.join().expect("the walking thread joins");
assert!(
cancellation.is_cancelled(),
"round {round}: the panic must cancel the sibling workers"
);
}
}
#[test]
fn worker_start_failure_returns_a_structured_error_and_cancels() {
let root = unique_root("worker-start");
create_wide_fixture(&root);
let cancellation = CancellationToken::default();
fail_next_worker_spawn();
let error = Walker::new(&root)
.threads(4)
.cancellation(cancellation.clone())
.collect()
.expect_err("injected worker start failure is returned");
let _ = fs::remove_dir_all(&root);
assert_eq!(error.operation(), "spawn_worker");
assert_eq!(error.path(), PathBuf::from(&root));
assert!(cancellation.is_cancelled());
}
}