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