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