Skip to main content

weavatrix_scan/scanner/
stream.rs

1use super::{Scanner, discovery::discover_repository_with_options};
2use crate::content::inspect_files;
3use crate::error::Result;
4use crate::scan_finalize::{RevisionBuilder, sort_report_evidence};
5use crate::scan_stream::{ScanSink, ScanSinkControl, ScanStreamReport};
6
7impl Scanner {
8    /// Inspects and emits the deterministic manifest under synchronous
9    /// backpressure without retaining selected file records.
10    ///
11    /// The sink is invoked in normalized relative-path order. Returning
12    /// [`ScanSinkControl::Stop`] stops inspection and marks the stream summary
13    /// incomplete.
14    ///
15    /// # Errors
16    ///
17    /// Returns the same errors as [`Self::scan`].
18    pub fn scan_into<S>(self, mut sink: S) -> Result<ScanStreamReport>
19    where
20        S: ScanSink,
21    {
22        let (mut report, runtime) =
23            discover_repository_with_options(&self.root, &self.options, &self.runtime)?;
24        sort_report_evidence(&mut report);
25        let files = std::mem::take(&mut report.files);
26        let mut revision = RevisionBuilder::new(&report);
27        let mut selected = 0_u64;
28        let mut emitted = 0_u64;
29        let mut stopped = false;
30        for file in files {
31            let inspected = inspect_files(vec![file], &self.options, runtime.started, None)?;
32            report.skipped.extend(inspected.skipped);
33            if !inspected.warnings.is_empty() {
34                report.complete = false;
35                report.warnings.extend(inspected.warnings);
36            }
37            report.cache.reused_hashes = report
38                .cache
39                .reused_hashes
40                .saturating_add(inspected.cache.reused_hashes);
41            report.cache.content_reads = report
42                .cache
43                .content_reads
44                .saturating_add(inspected.cache.content_reads);
45            report.cache.fingerprint_reads = report
46                .cache
47                .fingerprint_reads
48                .saturating_add(inspected.cache.fingerprint_reads);
49            if let Some(reason) = inspected.termination {
50                report.terminate(reason);
51            }
52            for file in inspected.files {
53                revision.push(&file);
54                selected = selected.saturating_add(1);
55                emitted = emitted.saturating_add(1);
56                if sink.on_file(&file) == ScanSinkControl::Stop {
57                    stopped = true;
58                    report.complete = false;
59                    break;
60                }
61            }
62            if stopped || report.termination.is_some() {
63                break;
64            }
65        }
66        sort_report_evidence(&mut report);
67        report.revision = revision.finish(&report);
68        report.finish_recording();
69        Ok(ScanStreamReport {
70            root: report.root,
71            selected,
72            emitted,
73            stopped,
74            skipped: report.skipped,
75            warnings: report.warnings,
76            ignore_sources: report.ignore_sources,
77            revision: report.revision,
78            complete: report.complete,
79            termination: report.termination,
80            portable: report.portable,
81            cache: report.cache,
82        })
83    }
84}