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() && !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}