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