1use super::Scanner;
2use super::compact::{apply_total_bytes_limit, compact_revision, discover_compact, sort_evidence};
3use super::entry::{process_entry_with, record_walk_error, walker_error_into_scan_error};
4use crate::config::ScanOptions;
5use crate::content::{ContentWorkerContext, VisitedFiles, visit_files};
6use crate::content_visit::{
7 ChangedContentVisitOutcome, ChangedContentVisitReport, ContentVisitControl, ContentVisitEvent,
8 ContentVisitMode, ContentVisitReport,
9};
10use crate::error::{Error, Result};
11use crate::ignore::RepositoryMatcher;
12use crate::report::{CompactContentEvidence, CompactScannedFile, FileVersion, ScanReport};
13use crate::runtime::ParallelRuntime;
14use crate::scan_limits::ScanRuntime;
15use crate::walker::{ErrorPolicy, Walker};
16use std::panic::{AssertUnwindSafe, catch_unwind};
17use std::path::PathBuf;
18use std::sync::{Arc, Mutex, mpsc};
19
20impl Scanner {
21 pub fn visit_content<Factory, Visitor>(self, factory: Factory) -> Result<ContentVisitReport>
45 where
46 Factory: Fn(usize) -> Visitor + Send + Sync + 'static,
47 Visitor:
48 for<'event> FnMut(ContentVisitEvent<'event>) -> ContentVisitControl + Send + 'static,
49 {
50 self.visit_content_with_root(0, factory)
51 }
52
53 pub fn visit_content_streaming<Factory, Visitor>(
68 self,
69 factory: Factory,
70 ) -> Result<ContentVisitReport>
71 where
72 Factory: Fn(usize) -> Visitor + Send + Sync + 'static,
73 Visitor:
74 for<'event> FnMut(ContentVisitEvent<'event>) -> ContentVisitControl + Send + 'static,
75 {
76 self.visit_content_with_root_mode(0, ContentVisitMode::Streaming, factory)
77 }
78
79 pub fn visit_changed_content<Factory, Visitor>(
96 self,
97 plan: &crate::WatchPlan,
98 factory: Factory,
99 ) -> Result<ChangedContentVisitOutcome>
100 where
101 Factory: Fn(usize) -> Visitor + Send + Sync + 'static,
102 Visitor:
103 for<'event> FnMut(ContentVisitEvent<'event>) -> ContentVisitControl + Send + 'static,
104 {
105 visit_changed_content_plan(self, plan, ContentVisitMode::Revision, factory)
106 }
107
108 pub fn visit_changed_content_streaming<Factory, Visitor>(
115 self,
116 plan: &crate::WatchPlan,
117 factory: Factory,
118 ) -> Result<ChangedContentVisitOutcome>
119 where
120 Factory: Fn(usize) -> Visitor + Send + Sync + 'static,
121 Visitor:
122 for<'event> FnMut(ContentVisitEvent<'event>) -> ContentVisitControl + Send + 'static,
123 {
124 visit_changed_content_plan(self, plan, ContentVisitMode::Streaming, factory)
125 }
126
127 pub(crate) fn visit_content_with_root<Factory, Visitor>(
128 self,
129 root_index: usize,
130 factory: Factory,
131 ) -> Result<ContentVisitReport>
132 where
133 Factory: Fn(usize) -> Visitor + Send + Sync + 'static,
134 Visitor:
135 for<'event> FnMut(ContentVisitEvent<'event>) -> ContentVisitControl + Send + 'static,
136 {
137 self.visit_content_with_root_mode(root_index, ContentVisitMode::Revision, factory)
138 }
139
140 pub(crate) fn visit_content_with_root_mode<Factory, Visitor>(
141 self,
142 root_index: usize,
143 mode: ContentVisitMode,
144 factory: Factory,
145 ) -> Result<ContentVisitReport>
146 where
147 Factory: Fn(usize) -> Visitor + Send + Sync + 'static,
148 Visitor:
149 for<'event> FnMut(ContentVisitEvent<'event>) -> ContentVisitControl + Send + 'static,
150 {
151 if self.options.limits.max_total_bytes.is_none()
152 && self.options.content_discovery == crate::ContentDiscoveryMode::Streaming
153 && !self.runtime.is_worker_thread()
154 {
155 return visit_content_direct(self, root_index, mode, factory);
156 }
157 let mut discovery_options = self.options.clone();
158 discovery_options.detect_binary_files = true;
159 let (mut evidence, mut files, scan_runtime) =
160 discover_compact(&self.root, &discovery_options, &self.runtime)?;
161 files.sort_unstable_by(|left, right| left.relative.cmp(&right.relative));
162 apply_total_bytes_limit(&mut evidence, &mut files, &self.options);
163 let discovered = u64::try_from(files.len()).unwrap_or(u64::MAX);
164
165 let cancellation = self.options.cancellation.clone().unwrap_or_default();
166 let mut visit_options = self.options.clone();
167 visit_options.cancellation = Some(cancellation.clone());
168 let workers = visit_options
169 .content_visit_worker_count(files.len())
170 .min(self.runtime.parallelism())
171 .max(1);
172 let worker_reports = run_workers(
173 evidence.root.clone(),
174 files,
175 visit_options,
176 &scan_runtime,
177 &self.runtime,
178 workers,
179 root_index,
180 mode,
181 factory,
182 )?;
183
184 Ok(finish_content_report(
185 evidence,
186 discovered,
187 worker_reports,
188 mode,
189 ))
190 }
191}
192
193fn visit_changed_content_plan<Factory, Visitor>(
194 scanner: Scanner,
195 plan: &crate::WatchPlan,
196 mode: ContentVisitMode,
197 factory: Factory,
198) -> Result<ChangedContentVisitOutcome>
199where
200 Factory: Fn(usize) -> Visitor + Send + Sync + 'static,
201 Visitor: for<'event> FnMut(ContentVisitEvent<'event>) -> ContentVisitControl + Send + 'static,
202{
203 if plan.full_rescan
204 || plan
205 .invalidated()
206 .any(|relative| !super::watch_update::is_safe_relative(relative))
207 {
208 return Ok(ChangedContentVisitOutcome::FullRescanRequired);
209 }
210
211 let mut options = scanner.options;
212 let cancellation = options.cancellation.clone().unwrap_or_default();
213 options.cancellation = Some(cancellation);
214 let mut prepared = prepare_discovery(&scanner.root, &options)?;
215 let mut changed = plan.changed.clone();
216 changed.sort_unstable();
217 changed.dedup();
218 let mut files = Vec::with_capacity(changed.len());
219 for relative in changed {
220 if let Some(reason) = prepared.runtime.before_next(&options) {
221 prepared.evidence.terminate(reason);
222 break;
223 }
224 prepared.runtime.record_entry();
225 match super::watch_update::changed_candidate(
226 &prepared.root,
227 &relative,
228 &options,
229 &mut prepared.matcher,
230 &mut prepared.evidence,
231 )? {
232 super::watch_update::ChangedPath::Candidate(file) => {
233 let file = *file;
234 files.push(content_candidate(file.relative, file.bytes, file.version));
235 }
236 super::watch_update::ChangedPath::MissingOrSkipped => {}
237 super::watch_update::ChangedPath::NeedsFullScan => {
238 return Ok(ChangedContentVisitOutcome::FullRescanRequired);
239 }
240 }
241 }
242 files.sort_unstable_by(|left, right| left.relative.cmp(&right.relative));
243 let selected = u64::try_from(files.len()).unwrap_or(u64::MAX);
244 let (mut evidence, runtime, _) = finish_stream_discovery(prepared, selected);
245 apply_total_bytes_limit(&mut evidence, &mut files, &options);
246 let discovered = u64::try_from(files.len()).unwrap_or(u64::MAX);
247 let workers = options
248 .content_visit_worker_count(files.len())
249 .min(scanner.runtime.parallelism())
250 .max(1);
251 let worker_reports = run_workers(
252 evidence.root.clone(),
253 files,
254 options,
255 &runtime,
256 &scanner.runtime,
257 workers,
258 0,
259 mode,
260 factory,
261 )?;
262 let mut removed = plan.removed.clone();
263 removed.sort_unstable();
264 removed.dedup();
265 Ok(ChangedContentVisitOutcome::Visited(Box::new(
266 ChangedContentVisitReport {
267 content: finish_content_report(evidence, discovered, worker_reports, mode),
268 removed,
269 },
270 )))
271}
272
273struct PreparedDiscovery {
274 root: PathBuf,
275 evidence: ScanReport,
276 matcher: RepositoryMatcher,
277 runtime: ScanRuntime,
278}
279
280#[allow(clippy::too_many_lines)]
281fn visit_content_direct<Factory, Visitor>(
282 scanner: Scanner,
283 root_index: usize,
284 mode: ContentVisitMode,
285 factory: Factory,
286) -> Result<ContentVisitReport>
287where
288 Factory: Fn(usize) -> Visitor + Send + Sync + 'static,
289 Visitor: for<'event> FnMut(ContentVisitEvent<'event>) -> ContentVisitControl + Send + 'static,
290{
291 let cancellation = scanner.options.cancellation.clone().unwrap_or_default();
292 let mut options = scanner.options;
293 options.cancellation = Some(cancellation.clone());
294 let mut discovery_options = options.clone();
295 discovery_options.detect_binary_files = true;
296 let prepared = prepare_discovery(&scanner.root, &discovery_options)?;
297 let root = Arc::new(prepared.root.clone());
298 let started = prepared.runtime.started;
299 let workers = options
300 .content_visit_worker_count(usize::MAX)
301 .min(scanner.runtime.parallelism())
302 .max(1);
303 let (sender, receiver) = mpsc::sync_channel(workers.saturating_mul(64).max(1));
304 let receiver = Arc::new(Mutex::new(receiver));
305 let factory = Arc::new(factory);
306 let worker_options = Arc::new(options.clone());
307 let (outcome_sender, outcome_receiver) = mpsc::channel();
308 let mut scheduled = 0_usize;
309 let mut schedule_error = None;
310 for worker_index in 0..workers {
311 let worker_receiver = Arc::clone(&receiver);
312 let worker_factory = Arc::clone(&factory);
313 let worker_options = Arc::clone(&worker_options);
314 let worker_root = Arc::clone(&root);
315 let worker_cancellation = cancellation.clone();
316 let worker_outcome = outcome_sender.clone();
317 if let Err(source) = scanner.runtime.try_execute(move || {
318 let outcome = catch_unwind(AssertUnwindSafe(|| {
319 let mut visitor = worker_factory(worker_index);
320 let mut buffer = vec![0_u8; 64 * 1024].into_boxed_slice();
321 let mut aggregate = VisitedFiles::empty(0);
322 loop {
323 let batch = {
324 let receiver = worker_receiver
325 .lock()
326 .unwrap_or_else(std::sync::PoisonError::into_inner);
327 let Ok(first) = receiver.recv() else {
328 break;
329 };
330 let mut batch = Vec::with_capacity(32);
331 batch.push(first);
332 while batch.len() < 32 {
333 match receiver.try_recv() {
334 Ok(work) => batch.push(work),
335 Err(
336 mpsc::TryRecvError::Empty | mpsc::TryRecvError::Disconnected,
337 ) => break,
338 }
339 }
340 batch
341 };
342 let visited = visit_files(
343 batch,
344 &worker_options,
345 started,
346 ContentWorkerContext {
347 root: worker_root.as_ref(),
348 root_index,
349 worker_index,
350 mode,
351 },
352 &mut buffer,
353 &mut visitor,
354 )?;
355 let stop = visited.visitor_quit || visited.evidence.termination.is_some();
356 aggregate.merge(visited);
357 if stop {
358 break;
359 }
360 }
361 Ok(aggregate)
362 }));
363 if outcome.as_ref().is_err() || outcome.as_ref().is_ok_and(std::result::Result::is_err)
364 {
365 worker_cancellation.cancel();
366 }
367 let _ = worker_outcome.send((worker_index, outcome));
368 }) {
369 schedule_error = Some(source);
370 cancellation.cancel();
371 break;
372 }
373 scheduled = scheduled.saturating_add(1);
374 }
375 drop(outcome_sender);
376 drop(receiver);
377 if schedule_error.is_none() {
378 let discovery = stream_discover_serial(prepared, &discovery_options, &sender);
379 if discovery.is_err() {
380 cancellation.cancel();
381 }
382 drop(sender);
383 let worker_reports = collect_stream_outcomes(&outcome_receiver, scheduled, root.as_ref())?;
384 let (mut evidence, scan_runtime, discovered) = discovery?;
385 if evidence.termination.is_none()
386 && let Some(reason) = scan_runtime.external_termination(&options)
387 {
388 evidence.terminate(reason);
389 }
390 return Ok(finish_content_report(
391 evidence,
392 discovered,
393 worker_reports,
394 mode,
395 ));
396 }
397
398 drop(sender);
399 let worker_result = collect_stream_outcomes(&outcome_receiver, scheduled, root.as_ref());
400 let _ = worker_result?;
401 Err(Error::io(
402 root.as_ref(),
403 schedule_error.expect("content worker scheduling failed"),
404 ))
405}
406
407fn prepare_discovery(root: &std::path::Path, options: &ScanOptions) -> Result<PreparedDiscovery> {
408 if options.walk.root_symlink_policy == crate::RootSymlinkPolicy::Reject {
409 let metadata = std::fs::symlink_metadata(root).map_err(|source| Error::io(root, source))?;
410 if metadata.file_type().is_symlink() {
411 return Err(Error::io(
412 root,
413 std::io::Error::new(
414 std::io::ErrorKind::InvalidInput,
415 "root symlink rejected by policy",
416 ),
417 ));
418 }
419 }
420 let canonical = root
421 .canonicalize()
422 .map_err(|source| Error::io(root, source))?;
423 if !canonical.is_dir() {
424 return Err(Error::InvalidRoot(canonical));
425 }
426 Ok(PreparedDiscovery {
427 evidence: ScanReport::new(
428 canonical.clone(),
429 options.evidence == crate::EvidenceMode::Complete,
430 ),
431 matcher: RepositoryMatcher::with_options(&canonical, options)?,
432 runtime: ScanRuntime::new(),
433 root: canonical,
434 })
435}
436
437fn stream_discover_serial(
438 mut prepared: PreparedDiscovery,
439 options: &ScanOptions,
440 sender: &mpsc::SyncSender<(u64, CompactScannedFile)>,
441) -> Result<(ScanReport, ScanRuntime, u64)> {
442 let mut walker = Walker::with_options(&prepared.root, options.walk_options())
443 .map_err(walker_error_into_scan_error)?;
444 let mut discovered = 0_u64;
445 loop {
446 if let Some(reason) = prepared.runtime.before_next(options) {
447 prepared.evidence.terminate(reason);
448 break;
449 }
450 let Some(item) = walker.next() else {
451 break;
452 };
453 prepared.runtime.record_entry();
454 match item {
455 Ok(entry) => {
456 let mut selected = None;
457 let skip = process_entry_with(
458 &entry,
459 options,
460 &mut prepared.evidence,
461 &prepared.matcher,
462 None,
463 |_path, relative, bytes, version| {
464 selected = Some(content_candidate(relative, bytes, version));
465 },
466 )?;
467 if let Some(file) = selected {
468 if !send_candidate(sender, discovered, file, options, &mut prepared.evidence)? {
469 break;
470 }
471 discovered = discovered.saturating_add(1);
472 }
473 if skip {
474 walker.skip_current_dir();
475 } else if entry.is_dir() {
476 prepared.matcher.prepare_directory(entry.path())?;
477 }
478 }
479 Err(error) if options.walk.error_policy == ErrorPolicy::Abort => {
480 return Err(walker_error_into_scan_error(error));
481 }
482 Err(error) => record_walk_error(&error, &prepared.root, &mut prepared.evidence),
483 }
484 }
485 Ok(finish_stream_discovery(prepared, discovered))
486}
487
488fn content_candidate(relative: String, bytes: u64, version: FileVersion) -> CompactScannedFile {
489 CompactScannedFile {
490 relative: relative.into_boxed_str(),
491 bytes,
492 content: Some(Box::new(CompactContentEvidence {
493 content_hash: None,
494 content_fingerprint: None,
495 version,
496 binary_checked: false,
497 })),
498 }
499}
500
501fn send_candidate(
502 sender: &mpsc::SyncSender<(u64, CompactScannedFile)>,
503 sequence: u64,
504 file: CompactScannedFile,
505 options: &ScanOptions,
506 evidence: &mut ScanReport,
507) -> Result<bool> {
508 if options
509 .cancellation
510 .as_ref()
511 .is_some_and(crate::CancellationToken::is_cancelled)
512 {
513 evidence.terminate(crate::ScanTermination::Cancelled);
514 return Ok(false);
515 }
516 if sender.send((sequence, file)).is_ok() {
517 return Ok(true);
518 }
519 if options
520 .cancellation
521 .as_ref()
522 .is_some_and(crate::CancellationToken::is_cancelled)
523 {
524 evidence.terminate(crate::ScanTermination::Cancelled);
525 Ok(false)
526 } else {
527 Err(Error::io(
528 &evidence.root,
529 std::io::Error::new(
530 std::io::ErrorKind::BrokenPipe,
531 "content workers stopped before traversal completed",
532 ),
533 ))
534 }
535}
536
537fn finish_stream_discovery(
538 mut prepared: PreparedDiscovery,
539 discovered: u64,
540) -> (ScanReport, ScanRuntime, u64) {
541 prepared.evidence.ignore_sources = prepared.matcher.sources().to_vec();
542 prepared.evidence.portable = prepared.matcher.portable();
543 if !prepared.matcher.warnings().is_empty() {
544 prepared.evidence.complete = false;
545 prepared
546 .evidence
547 .warnings
548 .extend_from_slice(prepared.matcher.warnings());
549 }
550 (prepared.evidence, prepared.runtime, discovered)
551}
552
553type StreamWorkerOutcome = std::thread::Result<Result<VisitedFiles>>;
554
555fn collect_stream_outcomes(
556 receiver: &mpsc::Receiver<(usize, StreamWorkerOutcome)>,
557 scheduled: usize,
558 root: &std::path::Path,
559) -> Result<Vec<VisitedFiles>> {
560 let mut outcomes = Vec::with_capacity(scheduled);
561 for _ in 0..scheduled {
562 outcomes.push(receiver.recv().map_err(|source| {
563 Error::io(
564 root,
565 std::io::Error::new(std::io::ErrorKind::BrokenPipe, source),
566 )
567 })?);
568 }
569 outcomes.sort_unstable_by_key(|(worker_index, _)| *worker_index);
570 let mut reports = Vec::with_capacity(scheduled);
571 let mut first_error = None;
572 let mut first_panic = None;
573 for (_, outcome) in outcomes {
574 match outcome {
575 Ok(Ok(report)) => reports.push(report),
576 Ok(Err(error)) if first_error.is_none() => first_error = Some(error),
577 Err(panic) if first_panic.is_none() => first_panic = Some(panic),
578 Ok(Err(_)) | Err(_) => {}
579 }
580 }
581 if let Some(panic) = first_panic {
582 std::panic::resume_unwind(panic);
583 }
584 if let Some(error) = first_error {
585 return Err(error);
586 }
587 Ok(reports)
588}
589
590fn finish_content_report(
591 mut evidence: ScanReport,
592 discovered: u64,
593 worker_reports: Vec<VisitedFiles>,
594 mode: ContentVisitMode,
595) -> ContentVisitReport {
596 let mut selected = Vec::new();
597 let mut totals = VisitedFiles::empty(0);
598 for mut worker in worker_reports {
599 selected.append(&mut worker.files);
600 totals.merge(worker);
601 }
602 evidence.skipped.extend(totals.evidence.skipped);
603 evidence.warnings.extend(totals.evidence.warnings);
604 evidence.termination = evidence.termination.or(totals.evidence.termination);
605 evidence.cache = totals.evidence.cache;
606 if totals.visitor_quit {
607 evidence.complete = false;
608 evidence.termination = Some(crate::ScanTermination::Cancelled);
609 }
610 if !evidence.warnings.is_empty() || evidence.termination.is_some() {
611 evidence.complete = false;
612 }
613 sort_evidence(&mut evidence);
614 let revision = if mode == ContentVisitMode::Revision {
615 selected.sort_unstable_by(|left, right| left.1.relative.cmp(&right.1.relative));
616 let files = selected
617 .into_iter()
618 .map(|(_, file)| file)
619 .collect::<Vec<_>>();
620 compact_revision(&evidence, &files)
621 } else {
622 String::new()
623 };
624 let stopped = evidence.termination.is_some();
625 evidence.finish_recording();
626 ContentVisitReport {
627 mode,
628 root: evidence.root,
629 discovered,
630 completed: totals.completed,
631 opened: totals.opened,
632 chunks: totals.chunks,
633 bytes_read: totals.bytes_read,
634 bytes_emitted: totals.bytes_emitted,
635 consumer_skipped: totals.consumer_skipped,
636 stopped,
637 skipped: evidence.skipped,
638 warnings: evidence.warnings,
639 ignore_sources: evidence.ignore_sources,
640 revision,
641 complete: evidence.complete,
642 termination: evidence.termination,
643 portable: evidence.portable,
644 cache: evidence.cache,
645 }
646}
647
648#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
649fn run_workers<Factory, Visitor>(
650 root: PathBuf,
651 files: Vec<CompactScannedFile>,
652 options: ScanOptions,
653 scan_runtime: &ScanRuntime,
654 runtime: &ParallelRuntime,
655 workers: usize,
656 root_index: usize,
657 mode: ContentVisitMode,
658 factory: Factory,
659) -> Result<Vec<VisitedFiles>>
660where
661 Factory: Fn(usize) -> Visitor + Send + Sync + 'static,
662 Visitor: for<'event> FnMut(ContentVisitEvent<'event>) -> ContentVisitControl + Send + 'static,
663{
664 if files.is_empty() {
665 return Ok(Vec::new());
666 }
667 let indexed = files
668 .into_iter()
669 .enumerate()
670 .map(|(sequence, file)| (u64::try_from(sequence).unwrap_or(u64::MAX), file))
671 .collect::<Vec<_>>();
672 if workers <= 1 || runtime.is_worker_thread() {
673 let mut visitor = factory(0);
674 let mut buffer = vec![0_u8; 64 * 1024].into_boxed_slice();
675 return visit_files(
676 indexed,
677 &options,
678 scan_runtime.started,
679 ContentWorkerContext {
680 root: &root,
681 root_index,
682 worker_index: 0,
683 mode,
684 },
685 &mut buffer,
686 &mut visitor,
687 )
688 .map(|report| vec![report]);
689 }
690
691 let chunk_size = indexed.len().div_ceil(workers);
692 let mut indexed = indexed.into_iter();
693 let mut chunks = Vec::with_capacity(workers);
694 loop {
695 let chunk = indexed.by_ref().take(chunk_size).collect::<Vec<_>>();
696 if chunk.is_empty() {
697 break;
698 }
699 chunks.push(chunk);
700 }
701 let root = Arc::new(root);
702 let options = Arc::new(options);
703 let factory = Arc::new(factory);
704 let (sender, receiver) = mpsc::channel();
705 let mut scheduled = 0_usize;
706 let mut schedule_error = None;
707 for (worker_index, chunk) in chunks.into_iter().enumerate() {
708 let worker_root = Arc::clone(&root);
709 let worker_options = Arc::clone(&options);
710 let worker_factory = Arc::clone(&factory);
711 let worker_sender = sender.clone();
712 let started = scan_runtime.started;
713 if let Err(source) = runtime.try_execute(move || {
714 let outcome = catch_unwind(AssertUnwindSafe(|| {
715 let mut visitor = worker_factory(worker_index);
716 let mut buffer = vec![0_u8; 64 * 1024].into_boxed_slice();
717 visit_files(
718 chunk,
719 &worker_options,
720 started,
721 ContentWorkerContext {
722 root: worker_root.as_ref(),
723 root_index,
724 worker_index,
725 mode,
726 },
727 &mut buffer,
728 &mut visitor,
729 )
730 }));
731 let _ = worker_sender.send((worker_index, outcome));
732 }) {
733 options
734 .cancellation
735 .as_ref()
736 .expect("content visit installs cancellation")
737 .cancel();
738 schedule_error = Some(source);
739 break;
740 }
741 scheduled = scheduled.saturating_add(1);
742 }
743 drop(sender);
744
745 let mut outcomes = Vec::with_capacity(scheduled);
746 for _ in 0..scheduled {
747 let (worker_index, outcome) = receiver.recv().map_err(|source| {
748 Error::io(
749 root.as_ref(),
750 std::io::Error::new(std::io::ErrorKind::BrokenPipe, source),
751 )
752 })?;
753 if outcome.as_ref().is_ok_and(std::result::Result::is_err) {
754 options
755 .cancellation
756 .as_ref()
757 .expect("content visit installs cancellation")
758 .cancel();
759 }
760 outcomes.push((worker_index, outcome));
761 }
762 outcomes.sort_unstable_by_key(|(worker_index, _)| *worker_index);
763 if let Some(index) = outcomes.iter().position(|(_, outcome)| outcome.is_err()) {
764 let (_, outcome) = outcomes.swap_remove(index);
765 let Err(panic) = outcome else {
766 unreachable!("panicked worker outcome exists");
767 };
768 std::panic::resume_unwind(panic);
769 }
770 if let Some(source) = schedule_error {
771 return Err(Error::io(root.as_ref(), source));
772 }
773 outcomes
774 .into_iter()
775 .map(|(_, outcome)| outcome.expect("worker panic handled"))
776 .collect()
777}
778
779#[cfg(test)]
780mod tests {
781 use super::*;
782 use std::time::{Duration, SystemTime, UNIX_EPOCH};
783
784 #[test]
785 fn content_visit_is_reentrant_on_its_runtime() {
786 let nonce = SystemTime::now()
787 .duration_since(UNIX_EPOCH)
788 .unwrap()
789 .as_nanos();
790 let root = std::env::temp_dir().join(format!(
791 "weavatrix-content-reentrant-{}-{nonce}",
792 std::process::id()
793 ));
794 std::fs::create_dir_all(&root).unwrap();
795 std::fs::write(root.join("value.rs"), "fn value() {}\n").unwrap();
796 let runtime = ParallelRuntime::dedicated(1).unwrap();
797 let nested_runtime = runtime.clone();
798 let (sender, receiver) = mpsc::channel();
799 runtime
800 .try_execute(move || {
801 let result = Scanner::new(&root)
802 .options(
803 ScanOptions::default()
804 .with_extensions(["rs"])
805 .selected_files_only()
806 .metadata_only(),
807 )
808 .runtime(nested_runtime)
809 .visit_content(|_| |_| ContentVisitControl::Continue)
810 .map(|report| report.completed);
811 let _ = std::fs::remove_dir_all(root);
812 sender.send(result).unwrap();
813 })
814 .unwrap();
815 assert_eq!(
816 receiver
817 .recv_timeout(Duration::from_secs(5))
818 .unwrap()
819 .unwrap(),
820 1
821 );
822 }
823}