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 #[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 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}