Skip to main content

fdu_core/query/
query_status.rs

1//! Completeness and provenance for one coherent report snapshot.
2
3use std::path::PathBuf;
4use std::time::{Duration, SystemTime};
5
6use crate::{Coverage, CoverageReason, Freshness, Index, Issue, Source};
7
8use super::{ReportSource, Request};
9
10/// Completeness of the facts used to answer one report.
11#[derive(Clone, Debug)]
12pub struct TreeStatus {
13    /// Whether every requested structural and content fact is represented.
14    pub complete: bool,
15    /// Coverage of the least complete requested tier.
16    pub coverage: Coverage,
17    /// Bounded, path-ordered operational failure details.
18    pub errors: Vec<Issue>,
19    /// Further failure details omitted by the shared retention bound.
20    pub errors_omitted: u64,
21}
22
23impl TreeStatus {
24    /// Derive status from the same immutable index snapshot used to build report rows.
25    pub fn of(index: &Index, request: &Request) -> Self {
26        let state = index.state();
27        let mut details: Vec<(PathBuf, Issue)> = Vec::with_capacity(crate::MAX_RETAINED_ISSUES);
28        let mut detail_count = 0_u64;
29        for issue in index.issues() {
30            detail_count = detail_count.saturating_add(1);
31            retain_first_detail(
32                &mut details,
33                (issue.path.clone().unwrap_or_default(), issue.clone()),
34            );
35        }
36        let wanted = index.content_identity(request.basis.content);
37        let admitted = index.content().and_then(|content| content.admit(&wanted));
38        let mut content_failures = 0_u64;
39        if request.basis.content.is_enabled() {
40            if let Some(content) = admitted {
41                for (path, analysis) in content.records() {
42                    let Some(reason) = analysis.operational_failure() else {
43                        continue;
44                    };
45                    content_failures = content_failures.saturating_add(1);
46                    detail_count = detail_count.saturating_add(1);
47                    let detail = analysis.error.clone().unwrap_or_else(|| match reason {
48                        crate::content::CoverageReason::IoError => {
49                            "content analysis could not read the file".to_string()
50                        }
51                        crate::content::CoverageReason::ChangedDuringRead => {
52                            "file changed during content analysis".to_string()
53                        }
54                        _ => unreachable!("operational_failure returns only operational reasons"),
55                    });
56                    retain_first_detail(
57                        &mut details,
58                        (path.to_path_buf(), Issue::provider_failure(Some(path), detail)),
59                    );
60                }
61            }
62        }
63        let content_tier_partial = request.basis.content.is_enabled()
64            && admitted
65                .and_then(crate::stored_state::ContentProjection::state)
66                .is_some_and(|tier| tier.freshness == Freshness::Partial);
67        let content_pending = index.content_has_pending(request.basis.content);
68        if content_tier_partial && content_failures == 0 {
69            detail_count = detail_count.saturating_add(1);
70            retain_first_detail(
71                &mut details,
72                (
73                    PathBuf::new(),
74                    Issue::provider_failure(
75                        None,
76                        "content analysis results became stale before they could be retained"
77                            .to_string(),
78                    ),
79                ),
80            );
81        }
82        let retained = u64::try_from(details.len()).unwrap_or(u64::MAX);
83        let errors = details.into_iter().map(|(_, issue)| issue).collect();
84        let complete = state.coverage == Coverage::Complete
85            && content_failures == 0
86            && !content_tier_partial
87            && !content_pending;
88        Self {
89            complete,
90            coverage: if content_failures == 0 && !content_tier_partial && !content_pending {
91                state.coverage
92            } else {
93                Coverage::Partial(CoverageReason::Failed)
94            },
95            errors,
96            errors_omitted: state
97                .issues
98                .omitted
99                .saturating_add(detail_count.saturating_sub(retained)),
100        }
101    }
102
103    /// Derive status from a one-shot walk, collapsing repeated causes first.
104    ///
105    /// Workers can meet one unreadable path more than once; the index retains each cause
106    /// once, so this path must too or the two routes disagree about `errors_omitted`.
107    pub(crate) fn of_walk(root: &std::path::Path, scan: &mut crate::ScanReport) -> Self {
108        crate::scan::normalize_walk_errors(root, &mut scan.errors);
109        let complete = scan.is_complete();
110        let mut details = Vec::with_capacity(crate::MAX_RETAINED_ISSUES);
111        let mut count = 0_u64;
112        for error in &scan.errors {
113            count = count.saturating_add(1);
114            let issue = crate::Issue::from_error_under(root, error);
115            retain_first_detail(&mut details, (issue.path.clone().unwrap_or_default(), issue));
116        }
117        let retained = u64::try_from(details.len()).unwrap_or(u64::MAX);
118        let errors = details.into_iter().map(|(_, issue)| issue).collect();
119        Self {
120            complete,
121            coverage: if complete {
122                Coverage::Complete
123            } else {
124                Coverage::Partial(CoverageReason::Inaccessible)
125            },
126            errors,
127            errors_omitted: count.saturating_sub(retained),
128        }
129    }
130}
131
132fn retain_first_detail(details: &mut Vec<(PathBuf, Issue)>, detail: (PathBuf, Issue)) {
133    let position = details
134        .binary_search_by(|current| {
135            current.0.cmp(&detail.0).then_with(|| current.1.message.cmp(&detail.1.message))
136        })
137        .unwrap_or_else(|position| position);
138    if position >= crate::MAX_RETAINED_ISSUES {
139        return;
140    }
141    details.insert(position, detail);
142    if details.len() > crate::MAX_RETAINED_ISSUES {
143        details.pop();
144    }
145}
146
147/// Source and currency of one retained tier.
148#[derive(Clone, Copy, PartialEq, Eq, Debug)]
149pub struct TierState {
150    /// Weakest source represented by this tier.
151    pub source: Source,
152    /// Current trust state of this tier.
153    pub freshness: Freshness,
154    /// When this tier's facts were observed, or `None` when the stored format cannot say.
155    pub observed_at_ns: Option<i64>,
156}
157
158/// Provenance of each tier contributing to a report.
159#[derive(Clone, Copy, PartialEq, Eq, Debug)]
160pub struct TierProvenance {
161    /// Retained filesystem-entry tier.
162    pub entries: TierState,
163    /// Sparse content-analysis tier when requested and present.
164    pub content: Option<TierState>,
165}
166
167/// How and when the coherent answer was produced.
168#[derive(Clone, Debug)]
169pub struct ReportProvenance {
170    /// Weakest source among tiers contributing to the report.
171    pub source: ReportSource,
172    /// Least fresh tier contributing to the report.
173    pub freshness: Freshness,
174    /// Conservative start watermark of the entry verification pass.
175    pub scan_started_at: Option<SystemTime>,
176    /// Caller-supplied instant at which this answer was generated.
177    pub generated_at: SystemTime,
178    /// Per-tier source, currency, and observation time.
179    pub tiers: TierProvenance,
180}
181
182impl ReportProvenance {
183    /// Derive report provenance for the requested tiers from one immutable index snapshot.
184    pub fn of(
185        index: &Index,
186        content_requested: crate::content::AnalysisSet,
187        generated_at: SystemTime,
188    ) -> Self {
189        let state = index.state();
190        let entries = TierState {
191            source: state.source,
192            freshness: state.freshness,
193            observed_at_ns: Some(index.writing_pass_started_at_ns()),
194        };
195        let scan_started_at = u64::try_from(index.writing_pass_started_at_ns())
196            .ok()
197            .map(|nanos| SystemTime::UNIX_EPOCH + Duration::from_nanos(nanos));
198        let content_pending = index.content_has_pending(content_requested);
199        let wanted = index.content_identity(content_requested);
200        let content = if content_requested.is_enabled() {
201            index.content().and_then(|content| content.admit(&wanted)).and_then(|content| {
202                content.state().map(|state| TierState {
203                    source: state.source,
204                    freshness: if content_pending { Freshness::Partial } else { state.freshness },
205                    observed_at_ns: state.observed_at_ns,
206                })
207            })
208        } else {
209            None
210        };
211        let source = content.map_or(entries.source, |state| entries.source.max(state.source));
212        let freshness = content
213            .map_or(entries.freshness, |state| least_fresh(entries.freshness, state.freshness));
214        Self {
215            source: report_source(source),
216            freshness,
217            scan_started_at,
218            generated_at,
219            tiers: TierProvenance { entries, content },
220        }
221    }
222
223    pub(crate) fn of_walk(started: SystemTime, generated_at: SystemTime, complete: bool) -> Self {
224        let observed_at_ns = crate::query::system_time_to_nanos(started);
225        let freshness = if complete { Freshness::Fresh } else { Freshness::Partial };
226        let entries = TierState { source: Source::Scanned, freshness, observed_at_ns };
227        Self {
228            source: ReportSource::ColdScan,
229            freshness,
230            scan_started_at: Some(started),
231            generated_at,
232            tiers: TierProvenance { entries, content: None },
233        }
234    }
235}
236
237fn least_fresh(left: Freshness, right: Freshness) -> Freshness {
238    let rank = |freshness| match freshness {
239        Freshness::Fresh => 0,
240        Freshness::Reconciling => 1,
241        Freshness::Stale => 2,
242        Freshness::Partial => 3,
243    };
244    if rank(left) >= rank(right) { left } else { right }
245}
246
247fn report_source(source: Source) -> ReportSource {
248    match source {
249        Source::Scanned => ReportSource::ColdScan,
250        Source::Revalidated | Source::JournalScoped => ReportSource::WarmRevalidate,
251        Source::Cached => ReportSource::CacheOnly,
252    }
253}
254
255#[cfg(test)]
256mod tests {
257    use std::fs;
258    use std::time::UNIX_EPOCH;
259
260    use super::{ReportProvenance, TreeStatus};
261    use crate::content::{AnalysisRequest, AnalysisSet, analyze_index};
262    use crate::query::{Basis, Query, Request};
263    use crate::scan::ScanConfig;
264    use crate::{Freshness, Index, Source};
265
266    #[test]
267    fn partial_content_tier_keeps_status_incomplete_without_a_failure_record() {
268        let mut index = Index::new("/unused");
269        let analysis = AnalysisSet::NONE.with_lines();
270        index.prepare_content_analysis(AnalysisRequest {
271            profile: analysis,
272            ..AnalysisRequest::default()
273        });
274        index.set_content_tier_state(Source::Scanned, Freshness::Partial, Some(1));
275        let mut basis = Basis::held_by(&index);
276        basis.content = analysis;
277        let request = Request::new(basis, Query::default(), UNIX_EPOCH);
278
279        let status = TreeStatus::of(&index, &request);
280
281        assert!(!status.complete);
282        assert_eq!(status.coverage, crate::Coverage::Partial(crate::CoverageReason::Failed));
283        assert_eq!(
284            status.errors[0].message,
285            "content analysis results became stale before they could be retained"
286        );
287    }
288
289    #[test]
290    fn metadata_only_provenance_excludes_an_unrequested_partial_content_tier() {
291        let mut index = Index::new("/unused");
292        index.prepare_content_analysis(AnalysisRequest {
293            profile: AnalysisSet::NONE.with_lines(),
294            ..AnalysisRequest::default()
295        });
296        index.set_content_tier_state(Source::Cached, Freshness::Partial, None);
297
298        let provenance = ReportProvenance::of(&index, AnalysisSet::NONE, UNIX_EPOCH);
299
300        assert_eq!(provenance.tiers.content, None);
301        assert_eq!(provenance.freshness, Freshness::Fresh);
302        assert_eq!(provenance.source, crate::query::ReportSource::ColdScan);
303    }
304
305    #[test]
306    fn repeated_content_read_failure_stays_incomplete_until_a_verified_recovery() {
307        let root = tempfile::tempdir().expect("root");
308        let path = root.path().join("failed.txt");
309        fs::write(&path, b"hello").expect("fixture");
310        let scan = ScanConfig::default();
311        let (mut index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
312        assert!(report.is_complete());
313        fs::remove_file(&path).expect("make the retained candidate unreadable");
314        let analysis = AnalysisRequest { profile: AnalysisSet::NONE.with_lines(), workers: 1 };
315        let first = analyze_index(&mut index, analysis);
316        assert_eq!(first.lines.io_errors, 1);
317        let request = Request::new(Basis::held_by(&index), Query::default(), UNIX_EPOCH);
318
319        let first_status = TreeStatus::of(&index, &request);
320        assert!(!first_status.complete);
321        assert_eq!(first_status.errors.len(), 1);
322        assert_eq!(
323            first_status.errors[0].path.as_deref(),
324            Some(std::path::Path::new("failed.txt"))
325        );
326
327        let repeated = analyze_index(&mut index, analysis);
328        assert_eq!(repeated.lines.io_errors, 1, "a failed record must be retried");
329        assert!(!TreeStatus::of(&index, &request).complete);
330
331        fs::write(&path, b"recovered").expect("recover file");
332        crate::scan::reconcile(&mut index, &scan, &mut |_| {}).expect("reconcile recovery");
333        let recovered = analyze_index(&mut index, analysis);
334        assert_eq!(recovered.lines.analyzed, 1);
335        let status = TreeStatus::of(&index, &request);
336        assert!(status.complete, "a successful reread clears the operational failure");
337        assert!(status.errors.is_empty());
338    }
339
340    /// A walk can meet one unreadable directory from several workers. The summary fast
341    /// path must collapse those repeats as the index does, or `--view summary` and
342    /// `--view tree` disagree about how many errors were omitted (R113-4).
343    #[test]
344    fn walk_status_counts_each_unreadable_path_once() {
345        let root = std::path::Path::new("/root");
346        let denied = |number: usize| crate::Error::Io {
347            path: root.join(format!("denied-{number:02}")),
348            source: std::io::Error::from(std::io::ErrorKind::PermissionDenied),
349        };
350        let mut scan = crate::ScanReport::default();
351        for number in (0..40).rev() {
352            scan.errors.push(denied(number));
353            scan.errors.push(denied(number));
354        }
355
356        let status = TreeStatus::of_walk(root, &mut scan);
357
358        assert!(!status.complete);
359        assert_eq!(status.errors.len(), 40);
360        assert_eq!(status.errors_omitted, 0);
361        assert_eq!(status.errors[0].path.as_deref(), Some(std::path::Path::new("denied-00")));
362    }
363
364    #[test]
365    fn content_failures_retain_the_first_paths_and_count_every_omission() {
366        let root = tempfile::tempdir().expect("root");
367        for number in (0..66).rev() {
368            fs::write(root.path().join(format!("file-{number:02}.txt")), b"x").expect("fixture");
369        }
370        let (mut index, report) =
371            crate::scan::scan_into_index(root.path(), &ScanConfig::default()).expect("scan");
372        assert!(report.is_complete());
373        for number in 0..66 {
374            fs::remove_file(root.path().join(format!("file-{number:02}.txt")))
375                .expect("make candidate unreadable");
376        }
377        analyze_index(
378            &mut index,
379            AnalysisRequest { profile: AnalysisSet::NONE.with_lines(), workers: 2 },
380        );
381        let request = Request::new(Basis::held_by(&index), Query::default(), UNIX_EPOCH);
382
383        let status = TreeStatus::of(&index, &request);
384
385        assert!(!status.complete);
386        assert_eq!(status.errors.len(), crate::MAX_RETAINED_ISSUES);
387        assert_eq!(status.errors_omitted, 2);
388        let retained_paths: Vec<_> =
389            status.errors.iter().map(|issue| issue.path.as_deref().expect("path")).collect();
390        assert_eq!(retained_paths[0], std::path::Path::new("file-00.txt"));
391        assert_eq!(retained_paths[63], std::path::Path::new("file-63.txt"));
392        assert!(retained_paths.windows(2).all(|pair| pair[0] < pair[1]));
393        assert!(!retained_paths.contains(&std::path::Path::new("file-64.txt")));
394        assert!(!retained_paths.contains(&std::path::Path::new("file-65.txt")));
395        assert_eq!(
396            index
397                .content()
398                .expect("content")
399                .records()
400                .filter(|(_, record)| { record.operational_failure().is_some() })
401                .count(),
402            66,
403            "the diagnostic bound must not discard retained failure state"
404        );
405    }
406}