Skip to main content

weavatrix_scan/parallel/
ordered_pull.rs

1use super::ParallelWalker;
2use super::pull::{ParallelWalkIter, PullBatch};
3use crate::control::CancellationToken;
4use crate::runtime::ParallelRuntime;
5use crate::walk_platform::{DirectoryIdentity, FileSystemId};
6use crate::walker::{
7    ErrorPolicy, WalkEntry, WalkError, WalkOperation, WalkOptions, WalkSkipReason, Walker,
8};
9use std::any::Any;
10use std::collections::{HashMap, HashSet, VecDeque};
11use std::io;
12use std::path::PathBuf;
13use std::sync::Arc;
14use std::sync::mpsc::{self, SyncSender, sync_channel};
15
16struct DirectoryTask {
17    id: u64,
18    path: PathBuf,
19    depth: usize,
20    identity: Option<DirectoryIdentity>,
21    ancestors: Arc<HashSet<DirectoryIdentity>>,
22}
23
24struct WorkerResult {
25    id: u64,
26    outcome: Result<DirectoryBatch, Box<dyn Any + Send>>,
27}
28
29struct DirectoryBatch {
30    entries: Vec<Result<WalkEntry, WalkError>>,
31    ancestors: Arc<HashSet<DirectoryIdentity>>,
32}
33
34struct PreparedItem {
35    item: Result<WalkEntry, WalkError>,
36    child: Option<u64>,
37}
38
39struct DirectoryFrame {
40    items: std::vec::IntoIter<PreparedItem>,
41}
42
43struct OrderedScheduler {
44    root: Arc<PathBuf>,
45    root_file_system: Option<FileSystemId>,
46    options: WalkOptions,
47    cancellation: CancellationToken,
48    runtime: ParallelRuntime,
49    limit: usize,
50    next_id: u64,
51    queued: VecDeque<DirectoryTask>,
52    outstanding: usize,
53    ready: HashMap<u64, DirectoryBatch>,
54    result_sender: mpsc::Sender<WorkerResult>,
55    result_receiver: mpsc::Receiver<WorkerResult>,
56    schedule_error: Option<WalkError>,
57}
58
59impl ParallelWalker {
60    /// Starts bounded parallel traversal and yields entries in strict,
61    /// deterministic depth-first order.
62    ///
63    /// Directory reads are prefetched up to `max_open` and the configured
64    /// parallelism. A capacity of zero is normalized to one.
65    ///
66    /// # Panics
67    ///
68    /// Panics if the coordinator thread cannot be created. Use
69    /// [`Self::try_into_iter_ordered_bounded`] for fallible startup.
70    #[must_use]
71    pub fn into_iter_ordered_bounded(self, capacity: usize) -> ParallelWalkIter {
72        self.try_into_iter_ordered_bounded(capacity)
73            .expect("ordered parallel pull coordinator thread can be created")
74    }
75
76    /// Fallible form of [`Self::into_iter_ordered_bounded`].
77    ///
78    /// # Errors
79    ///
80    /// Returns the coordinator thread spawn error.
81    pub fn try_into_iter_ordered_bounded(self, capacity: usize) -> io::Result<ParallelWalkIter> {
82        let capacity = capacity.max(1);
83        let (sender, receiver) = sync_channel(capacity.saturating_sub(1));
84        let cancellation = CancellationToken::new();
85        let coordinator_cancellation = cancellation.clone();
86        let use_serial = self.runtime.is_worker_thread();
87        let coordinator = std::thread::Builder::new()
88            .name("weavatrix-scan-ordered-pull".to_owned())
89            .spawn(move || {
90                if use_serial {
91                    ordered_serial(&self, &coordinator_cancellation, &sender);
92                } else {
93                    ordered_parallel(&self, &coordinator_cancellation, &sender);
94                }
95            })?;
96        Ok(ParallelWalkIter::from_coordinator(
97            receiver,
98            cancellation,
99            coordinator,
100        ))
101    }
102}
103
104fn ordered_serial(
105    walker: &ParallelWalker,
106    cancellation: &CancellationToken,
107    sender: &SyncSender<PullBatch>,
108) {
109    let error_policy = walker.options.error_policy;
110    let skip_stdout = walker.skip_stdout;
111    let mut options = walker.options.normalized();
112    options.error_policy = ErrorPolicy::Continue;
113    let mut walker = match Walker::with_options(&walker.root, options) {
114        Ok(walker) => walker,
115        Err(error) => {
116            let _ = sender.send(vec![Err(error)]);
117            return;
118        }
119    };
120    while !cancellation.is_cancelled() {
121        let Some(item) = walker.next() else {
122            break;
123        };
124        let abort = error_policy == ErrorPolicy::Abort && item.is_err();
125        let visible = match item.as_ref() {
126            Ok(entry) => !super::matches_stdout(entry, skip_stdout),
127            Err(_) => true,
128        };
129        if (visible && sender.send(vec![item]).is_err()) || abort {
130            break;
131        }
132    }
133}
134
135#[allow(clippy::too_many_lines)]
136fn ordered_parallel(
137    walker: &ParallelWalker,
138    cancellation: &CancellationToken,
139    sender: &SyncSender<PullBatch>,
140) {
141    let options = walker.options.normalized();
142    let mut root_options = options;
143    root_options.min_depth = 0;
144    root_options.error_policy = ErrorPolicy::Continue;
145    let mut root_walker = match Walker::with_options(&walker.root, root_options) {
146        Ok(walker) => walker,
147        Err(error) => {
148            let _ = sender.send(vec![Err(error)]);
149            return;
150        }
151    };
152    let root_file_system = root_walker.root_file_system;
153    let root = Arc::clone(&root_walker.root);
154    let root_entry = match root_walker
155        .next()
156        .expect("a validated root yields one entry")
157    {
158        Ok(entry) => entry,
159        Err(error) => {
160            let _ = sender.send(vec![Err(error)]);
161            return;
162        }
163    };
164    let root_identity = root_entry.directory_identity();
165    let root_can_descend = root_entry.is_dir() && root_entry.skip_reason().is_none();
166    if root_entry.depth() >= options.min_depth
167        && !super::matches_stdout(&root_entry, walker.skip_stdout)
168        && sender.send(vec![Ok(root_entry)]).is_err()
169    {
170        return;
171    }
172    if !root_can_descend || cancellation.is_cancelled() {
173        return;
174    }
175
176    let (result_sender, result_receiver) = mpsc::channel();
177    let mut scheduler = OrderedScheduler {
178        root: Arc::clone(&root),
179        root_file_system,
180        options,
181        cancellation: cancellation.clone(),
182        runtime: walker.runtime.clone(),
183        limit: super::requested_workers(&walker.runtime, walker.parallelism, options.max_open),
184        next_id: 1,
185        queued: VecDeque::new(),
186        outstanding: 0,
187        ready: HashMap::new(),
188        result_sender,
189        result_receiver,
190        schedule_error: None,
191    };
192    let mut ancestors = HashSet::new();
193    if let Some(identity) = root_identity {
194        ancestors.insert(identity);
195    }
196    scheduler.queued.push_back(DirectoryTask {
197        id: 0,
198        path: root.as_ref().clone(),
199        depth: 0,
200        identity: root_identity,
201        ancestors: Arc::new(ancestors),
202    });
203    scheduler.refill();
204    let root_batch = match scheduler.wait_for(0) {
205        Ok(Some(batch)) => batch,
206        Ok(None) => return,
207        Err(error) => {
208            let _ = sender.send(vec![Err(error)]);
209            return;
210        }
211    };
212    let mut frames = vec![scheduler.prepare_frame(root_batch)];
213
214    while !cancellation.is_cancelled() {
215        let Some(frame) = frames.last_mut() else {
216            break;
217        };
218        let Some(prepared) = frame.items.next() else {
219            frames.pop();
220            continue;
221        };
222        let child = prepared.child;
223        let visible = prepared
224            .item
225            .as_ref()
226            .map_or(true, |entry| entry.depth() >= options.min_depth);
227        let abort = options.error_policy == ErrorPolicy::Abort && prepared.item.is_err();
228        let visible = visible
229            && match prepared.item.as_ref() {
230                Ok(entry) => !super::matches_stdout(entry, walker.skip_stdout),
231                Err(_) => true,
232            };
233        if visible && sender.send(vec![prepared.item]).is_err() {
234            scheduler.cancel_and_drain();
235            return;
236        }
237        if abort {
238            scheduler.cancel_and_drain();
239            return;
240        }
241        if let Some(child) = child {
242            let batch = match scheduler.wait_for(child) {
243                Ok(Some(batch)) => batch,
244                Ok(None) => return,
245                Err(error) => {
246                    let _ = sender.send(vec![Err(error)]);
247                    return;
248                }
249            };
250            frames.push(scheduler.prepare_frame(batch));
251        }
252    }
253    scheduler.cancel_and_drain();
254}
255
256impl OrderedScheduler {
257    fn refill(&mut self) {
258        while !self.cancellation.is_cancelled() && self.outstanding < self.limit {
259            let Some(task) = self.queued.pop_front() else {
260                break;
261            };
262            let root = Arc::clone(&self.root);
263            let result_sender = self.result_sender.clone();
264            let cancellation = self.cancellation.clone();
265            let root_file_system = self.root_file_system;
266            let options = self.options;
267            let scheduled = self.runtime.try_execute(move || {
268                let id = task.id;
269                let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
270                    read_directory(&root, root_file_system, options, &cancellation, task)
271                }));
272                let _ = result_sender.send(WorkerResult { id, outcome });
273            });
274            match scheduled {
275                Ok(()) => self.outstanding += 1,
276                Err(source) => {
277                    self.cancellation.cancel();
278                    self.queued.clear();
279                    self.schedule_error = Some(WalkError::new(
280                        self.root.as_ref(),
281                        0,
282                        WalkOperation::ScheduleWorker,
283                        source,
284                    ));
285                    break;
286                }
287            }
288        }
289    }
290
291    fn wait_for(&mut self, id: u64) -> Result<Option<DirectoryBatch>, WalkError> {
292        if let Some(batch) = self.ready.remove(&id) {
293            return Ok(Some(batch));
294        }
295        if let Some(error) = self.schedule_error.take() {
296            self.cancel_and_drain();
297            return Err(error);
298        }
299        loop {
300            let Ok(result) = self.result_receiver.recv() else {
301                if let Some(error) = self.schedule_error.take() {
302                    return Err(error);
303                }
304                return Ok(None);
305            };
306            self.outstanding = self.outstanding.saturating_sub(1);
307            match result.outcome {
308                Ok(batch) if result.id == id => {
309                    self.refill();
310                    if let Some(error) = self.schedule_error.take() {
311                        self.cancel_and_drain();
312                        return Err(error);
313                    }
314                    return Ok(Some(batch));
315                }
316                Ok(batch) => {
317                    self.ready.insert(result.id, batch);
318                    self.refill();
319                }
320                Err(payload) => {
321                    self.cancel_and_drain();
322                    std::panic::resume_unwind(payload);
323                }
324            }
325            if self.cancellation.is_cancelled() {
326                self.cancel_and_drain();
327                if let Some(error) = self.schedule_error.take() {
328                    return Err(error);
329                }
330                return Ok(None);
331            }
332        }
333    }
334
335    fn prepare_frame(&mut self, mut batch: DirectoryBatch) -> DirectoryFrame {
336        batch.entries.sort_by(|left, right| {
337            let left_path = left
338                .as_ref()
339                .map_or_else(|error| error.path(), WalkEntry::path);
340            let right_path = right
341                .as_ref()
342                .map_or_else(|error| error.path(), WalkEntry::path);
343            left_path.cmp(right_path)
344        });
345        let mut items = Vec::with_capacity(batch.entries.len());
346        let mut children = Vec::new();
347        for item in batch.entries {
348            let child = item.as_ref().ok().and_then(|entry| {
349                (entry.is_dir() && entry.skip_reason().is_none()).then(|| {
350                    let id = self.next_id;
351                    self.next_id = self.next_id.saturating_add(1);
352                    let identity = entry.directory_identity();
353                    let mut ancestors = batch.ancestors.as_ref().clone();
354                    if let Some(identity) = identity {
355                        ancestors.insert(identity);
356                    }
357                    children.push(DirectoryTask {
358                        id,
359                        path: entry.path().to_path_buf(),
360                        depth: entry.depth(),
361                        identity,
362                        ancestors: Arc::new(ancestors),
363                    });
364                    id
365                })
366            });
367            items.push(PreparedItem { item, child });
368        }
369        for child in children.into_iter().rev() {
370            self.queued.push_front(child);
371        }
372        self.refill();
373        DirectoryFrame {
374            items: items.into_iter(),
375        }
376    }
377
378    fn cancel_and_drain(&mut self) {
379        self.cancellation.cancel();
380        self.queued.clear();
381        while self.outstanding > 0 {
382            if self.result_receiver.recv().is_err() {
383                break;
384            }
385            self.outstanding -= 1;
386        }
387    }
388}
389
390fn read_directory(
391    root: &Arc<PathBuf>,
392    root_file_system: Option<FileSystemId>,
393    options: WalkOptions,
394    cancellation: &CancellationToken,
395    task: DirectoryTask,
396) -> DirectoryBatch {
397    let ancestors = Arc::clone(&task.ancestors);
398    let mut worker_options = options;
399    worker_options.error_policy = ErrorPolicy::Continue;
400    worker_options.min_depth = 0;
401    worker_options.max_open = 1;
402    worker_options.max_depth = Some(
403        options
404            .max_depth
405            .unwrap_or(task.depth.saturating_add(1))
406            .min(task.depth.saturating_add(1)),
407    );
408    let mut walker = Walker::from_known_directory_with_ancestry(
409        root,
410        task.path,
411        task.depth,
412        worker_options,
413        root_file_system,
414        task.identity,
415        task.ancestors.as_ref().clone(),
416    );
417    let mut entries = Vec::new();
418    while !cancellation.is_cancelled() {
419        let Some(item) = walker.next() else {
420            break;
421        };
422        match item {
423            Ok(mut entry) => {
424                if entry.is_dir()
425                    && entry.skip_reason() == Some(WalkSkipReason::MaxDepth)
426                    && options
427                        .max_depth
428                        .is_none_or(|maximum| entry.depth() < maximum)
429                {
430                    entry.clear_depth_skip();
431                }
432                if entry.is_dir() {
433                    walker.skip_current_dir();
434                }
435                entries.push(Ok(entry));
436            }
437            Err(error) => entries.push(Err(error)),
438        }
439    }
440    DirectoryBatch { entries, ancestors }
441}