Skip to main content

weavatrix_scan/scanner/
content_visit.rs

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    /// Visits selected file bytes once with bounded parallelism.
22    ///
23    /// Ignore rules, file limits, path safety, binary detection, content
24    /// hashing, and revision evidence use the same scanner configuration. A
25    /// worker-local visitor is created by `factory`; callback order is
26    /// intentionally concurrent. Every event carries a monotonic work
27    /// sequence plus its root and normalized relative path; use the paths for
28    /// deterministic cross-run ordering.
29    ///
30    /// `ContentVisitControl::SkipFile` stops delivering chunks for the current
31    /// file. The scanner still finishes reading when hashing or binary
32    /// detection requires complete evidence. `ContentVisitControl::Quit`
33    /// cooperatively cancels every worker.
34    ///
35    /// # Errors
36    ///
37    /// Returns root, traversal, content I/O, or worker-submission failures
38    /// according to the configured error policy.
39    ///
40    /// # Panics
41    ///
42    /// Propagates a panic from the factory or visitor after active workers
43    /// observe cancellation.
44    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    /// Visits selected bytes without retaining a selected-file manifest or
54    /// computing a revision.
55    ///
56    /// Typed skip evidence is still retained when `EvidenceMode::Complete` is
57    /// configured. Use `selected_files_only()` as well for constant-memory
58    /// summary reporting.
59    ///
60    /// # Errors
61    ///
62    /// Returns the same errors as [`Self::visit_content`].
63    ///
64    /// # Panics
65    ///
66    /// Propagates callback panics like [`Self::visit_content`].
67    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    /// Visits only safe changed-file paths from a watcher plan.
80    ///
81    /// No directory traversal occurs. Plans that can affect directory
82    /// structure or file-selection rules return
83    /// [`ChangedContentVisitOutcome::FullRescanRequired`] before invoking the
84    /// factory. Removed paths are returned separately.
85    ///
86    /// # Errors
87    ///
88    /// Returns root, matcher, content I/O, or worker-submission failures
89    /// according to the configured error policy.
90    ///
91    /// # Panics
92    ///
93    /// Propagates a panic from the factory or visitor after active workers
94    /// observe cancellation.
95    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    /// Visits only changed-file bytes without retaining their compact manifest
109    /// or computing a subset revision.
110    ///
111    /// # Errors
112    ///
113    /// Returns the same errors as [`Self::visit_changed_content`].
114    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}