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