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 &mut prepared.matcher,
462 |_path, relative, bytes, version| {
463 selected = Some(content_candidate(relative, bytes, version));
464 },
465 )?;
466 if let Some(file) = selected {
467 if !send_candidate(sender, discovered, file, options, &mut prepared.evidence)? {
468 break;
469 }
470 discovered = discovered.saturating_add(1);
471 }
472 if skip {
473 walker.skip_current_dir();
474 }
475 }
476 Err(error) if options.walk.error_policy == ErrorPolicy::Abort => {
477 return Err(walker_error_into_scan_error(error));
478 }
479 Err(error) => record_walk_error(&error, &prepared.root, &mut prepared.evidence),
480 }
481 }
482 Ok(finish_stream_discovery(prepared, discovered))
483}
484
485fn content_candidate(relative: String, bytes: u64, version: FileVersion) -> CompactScannedFile {
486 CompactScannedFile {
487 relative: relative.into_boxed_str(),
488 bytes,
489 content: Some(Box::new(CompactContentEvidence {
490 content_hash: None,
491 content_fingerprint: None,
492 version,
493 binary_checked: false,
494 })),
495 }
496}
497
498fn send_candidate(
499 sender: &mpsc::SyncSender<(u64, CompactScannedFile)>,
500 sequence: u64,
501 file: CompactScannedFile,
502 options: &ScanOptions,
503 evidence: &mut ScanReport,
504) -> Result<bool> {
505 if options
506 .cancellation
507 .as_ref()
508 .is_some_and(crate::CancellationToken::is_cancelled)
509 {
510 evidence.terminate(crate::ScanTermination::Cancelled);
511 return Ok(false);
512 }
513 if sender.send((sequence, file)).is_ok() {
514 return Ok(true);
515 }
516 if options
517 .cancellation
518 .as_ref()
519 .is_some_and(crate::CancellationToken::is_cancelled)
520 {
521 evidence.terminate(crate::ScanTermination::Cancelled);
522 Ok(false)
523 } else {
524 Err(Error::io(
525 &evidence.root,
526 std::io::Error::new(
527 std::io::ErrorKind::BrokenPipe,
528 "content workers stopped before traversal completed",
529 ),
530 ))
531 }
532}
533
534fn finish_stream_discovery(
535 mut prepared: PreparedDiscovery,
536 discovered: u64,
537) -> (ScanReport, ScanRuntime, u64) {
538 prepared.evidence.ignore_sources = prepared.matcher.sources().to_vec();
539 prepared.evidence.portable = prepared.matcher.portable();
540 if !prepared.matcher.warnings().is_empty() {
541 prepared.evidence.complete = false;
542 prepared
543 .evidence
544 .warnings
545 .extend_from_slice(prepared.matcher.warnings());
546 }
547 (prepared.evidence, prepared.runtime, discovered)
548}
549
550type StreamWorkerOutcome = std::thread::Result<Result<VisitedFiles>>;
551
552fn collect_stream_outcomes(
553 receiver: &mpsc::Receiver<(usize, StreamWorkerOutcome)>,
554 scheduled: usize,
555 root: &std::path::Path,
556) -> Result<Vec<VisitedFiles>> {
557 let mut outcomes = Vec::with_capacity(scheduled);
558 for _ in 0..scheduled {
559 outcomes.push(receiver.recv().map_err(|source| {
560 Error::io(
561 root,
562 std::io::Error::new(std::io::ErrorKind::BrokenPipe, source),
563 )
564 })?);
565 }
566 outcomes.sort_unstable_by_key(|(worker_index, _)| *worker_index);
567 let mut reports = Vec::with_capacity(scheduled);
568 let mut first_error = None;
569 let mut first_panic = None;
570 for (_, outcome) in outcomes {
571 match outcome {
572 Ok(Ok(report)) => reports.push(report),
573 Ok(Err(error)) if first_error.is_none() => first_error = Some(error),
574 Err(panic) if first_panic.is_none() => first_panic = Some(panic),
575 Ok(Err(_)) | Err(_) => {}
576 }
577 }
578 if let Some(panic) = first_panic {
579 std::panic::resume_unwind(panic);
580 }
581 if let Some(error) = first_error {
582 return Err(error);
583 }
584 Ok(reports)
585}
586
587fn finish_content_report(
588 mut evidence: ScanReport,
589 discovered: u64,
590 worker_reports: Vec<VisitedFiles>,
591 mode: ContentVisitMode,
592) -> ContentVisitReport {
593 let mut selected = Vec::new();
594 let mut totals = VisitedFiles::empty(0);
595 for mut worker in worker_reports {
596 selected.append(&mut worker.files);
597 totals.merge(worker);
598 }
599 evidence.skipped.extend(totals.evidence.skipped);
600 evidence.warnings.extend(totals.evidence.warnings);
601 evidence.termination = evidence.termination.or(totals.evidence.termination);
602 evidence.cache = totals.evidence.cache;
603 if totals.visitor_quit {
604 evidence.complete = false;
605 evidence.termination = Some(crate::ScanTermination::Cancelled);
606 }
607 if !evidence.warnings.is_empty() || evidence.termination.is_some() {
608 evidence.complete = false;
609 }
610 sort_evidence(&mut evidence);
611 let revision = if mode == ContentVisitMode::Revision {
612 selected.sort_unstable_by(|left, right| left.1.relative.cmp(&right.1.relative));
613 let files = selected
614 .into_iter()
615 .map(|(_, file)| file)
616 .collect::<Vec<_>>();
617 compact_revision(&evidence, &files)
618 } else {
619 String::new()
620 };
621 let stopped = evidence.termination.is_some();
622 evidence.finish_recording();
623 ContentVisitReport {
624 mode,
625 root: evidence.root,
626 discovered,
627 completed: totals.completed,
628 opened: totals.opened,
629 chunks: totals.chunks,
630 bytes_read: totals.bytes_read,
631 bytes_emitted: totals.bytes_emitted,
632 consumer_skipped: totals.consumer_skipped,
633 stopped,
634 skipped: evidence.skipped,
635 warnings: evidence.warnings,
636 ignore_sources: evidence.ignore_sources,
637 revision,
638 complete: evidence.complete,
639 termination: evidence.termination,
640 portable: evidence.portable,
641 cache: evidence.cache,
642 }
643}
644
645#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
646fn run_workers<Factory, Visitor>(
647 root: PathBuf,
648 files: Vec<CompactScannedFile>,
649 options: ScanOptions,
650 scan_runtime: &ScanRuntime,
651 runtime: &ParallelRuntime,
652 workers: usize,
653 root_index: usize,
654 mode: ContentVisitMode,
655 factory: Factory,
656) -> Result<Vec<VisitedFiles>>
657where
658 Factory: Fn(usize) -> Visitor + Send + Sync + 'static,
659 Visitor: for<'event> FnMut(ContentVisitEvent<'event>) -> ContentVisitControl + Send + 'static,
660{
661 if files.is_empty() {
662 return Ok(Vec::new());
663 }
664 let indexed = files
665 .into_iter()
666 .enumerate()
667 .map(|(sequence, file)| (u64::try_from(sequence).unwrap_or(u64::MAX), file))
668 .collect::<Vec<_>>();
669 if workers <= 1 || runtime.is_worker_thread() {
670 let mut visitor = factory(0);
671 let mut buffer = vec![0_u8; 64 * 1024].into_boxed_slice();
672 return visit_files(
673 indexed,
674 &options,
675 scan_runtime.started,
676 ContentWorkerContext {
677 root: &root,
678 root_index,
679 worker_index: 0,
680 mode,
681 },
682 &mut buffer,
683 &mut visitor,
684 )
685 .map(|report| vec![report]);
686 }
687
688 let chunk_size = indexed.len().div_ceil(workers);
689 let mut indexed = indexed.into_iter();
690 let mut chunks = Vec::with_capacity(workers);
691 loop {
692 let chunk = indexed.by_ref().take(chunk_size).collect::<Vec<_>>();
693 if chunk.is_empty() {
694 break;
695 }
696 chunks.push(chunk);
697 }
698 let root = Arc::new(root);
699 let options = Arc::new(options);
700 let factory = Arc::new(factory);
701 let (sender, receiver) = mpsc::channel();
702 let mut scheduled = 0_usize;
703 let mut schedule_error = None;
704 for (worker_index, chunk) in chunks.into_iter().enumerate() {
705 let worker_root = Arc::clone(&root);
706 let worker_options = Arc::clone(&options);
707 let worker_factory = Arc::clone(&factory);
708 let worker_sender = sender.clone();
709 let started = scan_runtime.started;
710 if let Err(source) = runtime.try_execute(move || {
711 let outcome = catch_unwind(AssertUnwindSafe(|| {
712 let mut visitor = worker_factory(worker_index);
713 let mut buffer = vec![0_u8; 64 * 1024].into_boxed_slice();
714 visit_files(
715 chunk,
716 &worker_options,
717 started,
718 ContentWorkerContext {
719 root: worker_root.as_ref(),
720 root_index,
721 worker_index,
722 mode,
723 },
724 &mut buffer,
725 &mut visitor,
726 )
727 }));
728 let _ = worker_sender.send((worker_index, outcome));
729 }) {
730 options
731 .cancellation
732 .as_ref()
733 .expect("content visit installs cancellation")
734 .cancel();
735 schedule_error = Some(source);
736 break;
737 }
738 scheduled = scheduled.saturating_add(1);
739 }
740 drop(sender);
741
742 let mut outcomes = Vec::with_capacity(scheduled);
743 for _ in 0..scheduled {
744 let (worker_index, outcome) = receiver.recv().map_err(|source| {
745 Error::io(
746 root.as_ref(),
747 std::io::Error::new(std::io::ErrorKind::BrokenPipe, source),
748 )
749 })?;
750 if outcome.as_ref().is_ok_and(std::result::Result::is_err) {
751 options
752 .cancellation
753 .as_ref()
754 .expect("content visit installs cancellation")
755 .cancel();
756 }
757 outcomes.push((worker_index, outcome));
758 }
759 outcomes.sort_unstable_by_key(|(worker_index, _)| *worker_index);
760 if let Some(index) = outcomes.iter().position(|(_, outcome)| outcome.is_err()) {
761 let (_, outcome) = outcomes.swap_remove(index);
762 let Err(panic) = outcome else {
763 unreachable!("panicked worker outcome exists");
764 };
765 std::panic::resume_unwind(panic);
766 }
767 if let Some(source) = schedule_error {
768 return Err(Error::io(root.as_ref(), source));
769 }
770 outcomes
771 .into_iter()
772 .map(|(_, outcome)| outcome.expect("worker panic handled"))
773 .collect()
774}
775
776#[cfg(test)]
777mod tests {
778 use super::*;
779 use std::time::{Duration, SystemTime, UNIX_EPOCH};
780
781 #[test]
782 fn content_visit_is_reentrant_on_its_runtime() {
783 let nonce = SystemTime::now()
784 .duration_since(UNIX_EPOCH)
785 .unwrap()
786 .as_nanos();
787 let root = std::env::temp_dir().join(format!(
788 "weavatrix-content-reentrant-{}-{nonce}",
789 std::process::id()
790 ));
791 std::fs::create_dir_all(&root).unwrap();
792 std::fs::write(root.join("value.rs"), "fn value() {}\n").unwrap();
793 let runtime = ParallelRuntime::dedicated(1).unwrap();
794 let nested_runtime = runtime.clone();
795 let (sender, receiver) = mpsc::channel();
796 runtime
797 .try_execute(move || {
798 let result = Scanner::new(&root)
799 .options(
800 ScanOptions::default()
801 .with_extensions(["rs"])
802 .selected_files_only()
803 .metadata_only(),
804 )
805 .runtime(nested_runtime)
806 .visit_content(|_| |_| ContentVisitControl::Continue)
807 .map(|report| report.completed);
808 let _ = std::fs::remove_dir_all(root);
809 sender.send(result).unwrap();
810 })
811 .unwrap();
812 assert_eq!(
813 receiver
814 .recv_timeout(Duration::from_secs(5))
815 .unwrap()
816 .unwrap(),
817 1
818 );
819 }
820}