Skip to main content

fdu_core/
watch_session.rs

1//! A live session: an index, a watcher, and the query they answer together.
2//!
3//! A watch run is the same query as a one-shot run, re-evaluated as changes arrive —
4//! there is no separate watch grammar. This module owns that composition so the CLI loop
5//! and the Python iterator are both thin consumers rather than two implementations of
6//! the same coordination.
7//!
8//! # Detection is event-driven
9//!
10//! Changes arrive from the operating system's own notification backend, never from
11//! polling: `FSEvents` on macOS, inotify on Linux, `ReadDirectoryChangesW` on Windows. An
12//! idle tree costs no filesystem work at all. Events are hints, so each coalesced path is
13//! verified with one fresh stat before it becomes a delta, and the interval a caller
14//! passes throttles aggregate repaints and persistence — it plays no part in detection.
15
16use std::collections::BTreeMap;
17use std::path::PathBuf;
18use std::time::{Duration, Instant};
19
20use crate::engine_contract::{Commit, EffectiveChange, EntryKind, Error, Result};
21use crate::index::IndexHandle;
22use crate::query::{Basis, Delivery, Query, Report, Request, Selection, WatchDelivery, report};
23use crate::scan::ScanConfig;
24use crate::watch::{WatchConfig, Watcher};
25
26/// One effective change, already filtered through the run's selection.
27#[derive(Clone, Debug, PartialEq, Eq)]
28pub struct Change {
29    /// Path relative to the index root.
30    pub path: PathBuf,
31    /// What happened to it.
32    pub kind: ChangeKind,
33    /// What the entry is, when it still exists.
34    pub entry_kind: Option<EntryKind>,
35    /// Apparent bytes, when the entry still exists.
36    pub bytes: Option<u64>,
37    /// Allocated bytes, when the entry still exists.
38    pub allocated: Option<u64>,
39    /// Modification time, when the entry still exists.
40    pub mtime_ns: Option<i64>,
41    /// Whether `.gitignore` rules ignore this entry after the commit.
42    ///
43    /// Set on an upsert, and on a removal a rule edit caused, where the new classification
44    /// is why the row left the selection. `None` on every other record: a session that
45    /// observes no control state never claims a classification an index without rules can
46    /// make, and an ordinary removal or an invalidation has no entry left to classify. A
47    /// report's rows carry the same split, and a stream that could not would be the one
48    /// place a consumer had to guess.
49    pub ignored: Option<bool>,
50    /// The index clock at which this change was committed.
51    pub clock: u64,
52}
53
54/// What happened to a path.
55#[derive(Clone, Copy, PartialEq, Eq, Debug)]
56pub enum ChangeKind {
57    /// The entry appeared or its metadata changed.
58    Upsert,
59    /// The entry is gone.
60    Remove,
61    /// A producer could not describe the change precisely; the subtree was re-scanned.
62    ///
63    /// Surfaced rather than swallowed: this is the signal that a consumer's own view of
64    /// the subtree may have gaps, and dropping it is how an index silently diverges.
65    Invalidate,
66}
67
68/// One batch of applied changes.
69#[derive(Clone, Debug, Default)]
70pub struct Batch {
71    /// Changes the selection admitted, in commit order.
72    pub changes: Vec<Change>,
73    /// Whether anything was applied at all, before selection filtering.
74    ///
75    /// A batch can be non-empty and still yield no changes, when everything it carried
76    /// was filtered out. Aggregate views re-render on this rather than on `changes`,
77    /// because a filtered-out change still moves the totals a tree view reports.
78    pub dirty: bool,
79}
80
81/// What one batch of commits needs from the index, read once under one lock.
82struct BatchFacts {
83    /// How ignore rules classify each touched entry the index still holds once the batch
84    /// applied, or `None` when the index observed no control state and classifies nothing.
85    ///
86    /// A map rather than a set of the ignored: an entry a later commit in the same batch
87    /// removed is in neither partition, and a set could not tell that from unignored.
88    ignored: Option<BTreeMap<PathBuf, bool>>,
89    /// The retained facts of each reclassified entry the selection could move, so a rule
90    /// edit that admits one can stream the upsert that draws it.
91    reclassified: BTreeMap<PathBuf, EntryFacts>,
92}
93
94impl BatchFacts {
95    /// How rules classify a touched entry: `None` when nothing was classified, and when
96    /// the entry is gone from the index, which leaves no entry to classify.
97    fn is_ignored(&self, path: &std::path::Path) -> Option<bool> {
98        self.ignored.as_ref()?.get(path).copied()
99    }
100}
101
102/// One reclassified entry's retained facts.
103#[derive(Clone, Copy)]
104struct EntryFacts {
105    kind: EntryKind,
106    bytes: u64,
107    allocated: u64,
108    mtime_ns: i64,
109}
110
111/// The result of a throttled attempt to persist a live session.
112#[derive(Debug)]
113pub enum SaveOutcome {
114    /// The metadata snapshot reached disk.
115    Written,
116    /// No write was due, or the plan could not yet persist the current state.
117    Skipped,
118    /// Persistence failed; the live session remains usable and will retry.
119    Failed(Error),
120}
121
122struct Persistence {
123    pending: bool,
124    last_attempt: Instant,
125}
126
127impl Persistence {
128    fn persist_due(
129        &mut self,
130        now: Instant,
131        interval: Duration,
132        save: impl FnOnce() -> Result<bool>,
133    ) -> SaveOutcome {
134        if !save_is_due(self.pending, now.saturating_duration_since(self.last_attempt), interval) {
135            return SaveOutcome::Skipped;
136        }
137        let outcome = match save() {
138            Ok(true) => SaveOutcome::Written,
139            Ok(false) => SaveOutcome::Skipped,
140            Err(error) => SaveOutcome::Failed(error),
141        };
142        self.pending = pending_after(&outcome);
143        // Skips and failures are throttled too, while retaining the pending work.
144        self.last_attempt = now;
145        outcome
146    }
147}
148
149fn save_is_due(pending: bool, since_last_save: Duration, interval: Duration) -> bool {
150    pending && since_last_save >= interval
151}
152
153fn pending_after(outcome: &SaveOutcome) -> bool {
154    !matches!(outcome, SaveOutcome::Written)
155}
156
157/// An index paired with a watcher, answering one request continuously.
158pub struct Session {
159    index: IndexHandle,
160    watcher: Watcher,
161    scan: ScanConfig,
162    request: Request,
163    plan: crate::Plan,
164    persistence: Persistence,
165    startup_save_error: Option<Error>,
166}
167
168impl Session {
169    /// Open a tree and bind observation under the shared execution plan.
170    ///
171    /// Startup persistence is joined before binding the session. A save failure is
172    /// returned by the first `persist_due` call, so it does not discard a valid live
173    /// answer. Filesystem and observation failures still fail startup.
174    pub fn start(request: Request, delivery: Delivery) -> Result<Self> {
175        Self::start_observed(request, delivery, None)
176    }
177
178    /// [`Self::start`], reporting the initial scan through `progress` as it runs.
179    ///
180    /// The same session as [`Self::start`]; the handle observes the start and changes
181    /// nothing about it. A start is two passes: the open (cold, or a load and
182    /// revalidation) with its save, if it writes one, joined under
183    /// [`ProgressPhase::Saving`](crate::ProgressPhase), and then the revalidation that
184    /// closes the gap between that walk and the bound watcher. The second pass begins
185    /// again at [`ProgressPhase::Revalidating`](crate::ProgressPhase) with the walk
186    /// counters restarted, so they end at the tree's totals once, not twice. Once this
187    /// returns the session reports nothing further through the handle; its repaints are
188    /// the progress from there.
189    pub fn start_with_progress(
190        request: Request,
191        delivery: Delivery,
192        progress: &crate::Progress,
193    ) -> Result<Self> {
194        Self::start_observed(request, delivery, Some(progress))
195    }
196
197    fn start_observed(
198        request: Request,
199        mut delivery: Delivery,
200        progress: Option<&crate::Progress>,
201    ) -> Result<Self> {
202        delivery.watch.get_or_insert_with(WatchDelivery::default);
203        let plan =
204            crate::plan(&request, &delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
205        let (index, report, pending, _diagnostics) =
206            crate::execute(&plan, &request.basis, false, progress)?;
207        let startup_save_error = pending.join().err();
208        let index = std::sync::Arc::into_inner(index)
209            .expect("the joined writer released the only other reference");
210        let mut session = Self::new_observed(
211            IndexHandle::new(index),
212            request,
213            &delivery,
214            WatchConfig::default(),
215            progress,
216        )?;
217        session.persistence.pending |= startup_save_error.is_some() || !report.is_complete();
218        session.startup_save_error = startup_save_error;
219        Ok(session)
220    }
221
222    /// Persist pending changes when the watch delivery's interval has elapsed.
223    ///
224    /// Call this after batches and idle timeouts. Only a completed write clears pending
225    /// changes; a refused or failed write is retried on a later interval. `now` is a
226    /// monotonic caller-supplied clock so wall-clock corrections cannot postpone saves.
227    pub fn persist_due(&mut self, now: Instant) -> SaveOutcome {
228        if let Some(error) = self.startup_save_error.take() {
229            self.persistence.last_attempt = now;
230            return SaveOutcome::Failed(error);
231        }
232        let interval = self.plan.delivery().watch.expect("watch plan").interval;
233        self.persistence.persist_due(now, interval, || {
234            if !self.plan.persists() || self.plan.delivery().cache_path.is_none() {
235                return Ok(false);
236            }
237            let index = self.index.snapshot()?;
238            crate::persist_index(&index, &self.plan)
239        })
240    }
241
242    /// Start watching an already-opened index, answering `request` as the tree changes.
243    ///
244    /// `request` carries its own `now`, fixed when it was built: a watch answers one
245    /// request as the tree changes, and a relative time window that slid under it would
246    /// make two repaints answer two different questions.
247    ///
248    /// `delivery` is the one its caller opened the index under. It is taken rather than
249    /// composed here because the cache policy is part of it: a session built against a
250    /// fabricated `Delivery` read `cache: Auto` whatever the caller had asked for, so
251    /// [`RequestError::WatchCacheOnly`](crate::query::RequestError::WatchCacheOnly) could
252    /// not fire inside the engine at all and the rule held only at the two front doors
253    /// (fdu-i18y). Starting a session is what a watch *is*, so the delivery is read as a
254    /// watch whether or not the caller remembered to say so.
255    ///
256    /// Refusals are in the order every route publishes: what no delivery can carry first,
257    /// then what this index was taken under, then what it holds. The request-level rule
258    /// speaks first, so a library caller watching a depth-2 request is told the watch
259    /// cannot narrow its scope rather than that the index has another scope.
260    ///
261    /// # Errors
262    ///
263    /// [`Error::InvalidRequest`] when a watch cannot deliver the request -- a narrowed scan
264    /// scope, content analysis nothing re-reads, a snapshot nothing verified -- or when
265    /// this index cannot answer it, including a selection by ignored state over an index
266    /// that observed no control state.
267    /// [`Error::ScanScopeMismatch`] when the index was not taken under the request's scope.
268    pub fn new(
269        index: IndexHandle,
270        request: Request,
271        delivery: &Delivery,
272        watch: WatchConfig,
273    ) -> Result<Self> {
274        Self::new_observed(index, request, delivery, watch, None)
275    }
276
277    fn new_observed(
278        index: IndexHandle,
279        request: Request,
280        delivery: &Delivery,
281        watch: WatchConfig,
282        progress: Option<&crate::Progress>,
283    ) -> Result<Self> {
284        let root = index.root_path()?;
285        let scan = request.basis.scope.scan_config(delivery);
286        // What no delivery can carry, before anything stored is read and before the
287        // backend is bound: this is the rule each surface used to keep for itself, so a
288        // library caller could watch what `--watch` has always refused.
289        let delivery =
290            Delivery { watch: Some(delivery.watch.unwrap_or_default()), ..delivery.clone() };
291        crate::plan(&request, &delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
292        crate::validate_basis_root(&root, &request.basis)?;
293        // Reject an out-of-scope watch before the backend is bound, so a rejected run
294        // never leaves a watcher registered on the tree.
295        scan.validate_for_scope(index.scope()?)?;
296        // The scope check above proved the index was taken under exactly this scan's
297        // identity, control tier included, so what remains is what this index holds.
298        let held = Basis {
299            root: root.clone(),
300            scope: scan.clone().into(),
301            content: index.read_with(crate::Index::content_set)?,
302        };
303        request.validate_read(&held).map_err(Error::InvalidRequest)?;
304        // Bind observation before closing the gap from the scan that produced `index`.
305        // The full reconciliation catches a mutation that completed before registration;
306        // the capture drain applies every hint observed while that pass ran.
307        let watcher = Watcher::new(&root, watch)?;
308        Self::finish_initial_handoff(index, request, &delivery, watcher, scan, progress)
309    }
310
311    /// Finish the two-part initial handoff after observation has been bound.
312    ///
313    /// Kept separate so the scripted watcher exercises the same reconciliation, drain, and
314    /// acceptance boundary as an OS watcher. Once this returns, later partial observations are
315    /// valid live state; `accept_partial` governs only the coherent state handed to the caller.
316    ///
317    /// `progress` observes the handoff's own walk and drain only. The session keeps
318    /// `scan` without it, so nothing it reconciles later reports through a handle whose
319    /// poller has long since stopped.
320    fn finish_initial_handoff(
321        index: IndexHandle,
322        request: Request,
323        delivery: &Delivery,
324        watcher: Watcher,
325        scan: ScanConfig,
326        progress: Option<&crate::Progress>,
327    ) -> Result<Self> {
328        if let Some(progress) = progress {
329            progress.begin_pass(crate::ProgressPhase::Revalidating);
330        }
331        let observed = ScanConfig { progress: progress.cloned(), ..scan.clone() };
332        let mut dirty = false;
333        let reconciliation = crate::scan::reconcile_handle(&index, &observed, &mut |commit| {
334            dirty |= !commit.changes.is_empty();
335        })?;
336        if !reconciliation.scan.is_complete() && !delivery.accept_partial {
337            return Err(Error::ObservationHandoffIncomplete);
338        }
339        dirty |= drain_initial_capture(&watcher, &index, &observed)?;
340        if !delivery.accept_partial
341            && !index.read_with(|index| crate::query::TreeStatus::of(index, &request).complete)?
342        {
343            return Err(Error::ObservationHandoffIncomplete);
344        }
345        let plan =
346            crate::plan(&request, delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
347        Ok(Self {
348            index,
349            watcher,
350            scan,
351            request,
352            plan,
353            persistence: Persistence { pending: dirty, last_attempt: Instant::now() },
354            startup_save_error: None,
355        })
356    }
357
358    /// The request this session answers.
359    pub fn request(&self) -> &Request {
360        &self.request
361    }
362
363    /// The query this session answers.
364    pub fn query(&self) -> &Query {
365        &self.request.query
366    }
367
368    /// Render the current answer.
369    ///
370    /// The same `report` a one-shot run produces, from the same index, which is what
371    /// makes "watch is the same query repeated" true rather than aspirational.
372    pub fn report(&self, generated_at: std::time::SystemTime) -> Result<Report> {
373        let index = self.index.snapshot()?;
374        report(&index, &self.request, generated_at)
375    }
376
377    /// A consistent copy of the current index.
378    ///
379    /// Used to persist a live session without holding a lock across the write.
380    pub fn index_snapshot(&self) -> Result<crate::Index> {
381        self.index.snapshot()
382    }
383
384    /// Wait for the next batch of changes, up to `timeout`.
385    ///
386    /// Returns `None` when nothing arrived in the window, which is the idle case and
387    /// costs no filesystem work.
388    ///
389    /// Takes `&mut self` because consuming from the event queue is a mutation: two
390    /// callers draining one session would each see an arbitrary half of the stream.
391    pub fn next_batch(&mut self, timeout: Duration) -> Result<Option<Batch>> {
392        let mut commits: Vec<Commit> = Vec::new();
393        let outcome =
394            self.watcher.apply_next(&self.index, &self.scan, timeout, &mut |commit: &Commit| {
395                commits.push(commit.clone());
396            });
397
398        self.persistence.pending |= commits.iter().any(|commit| !commit.changes.is_empty());
399        let Some(_report) = outcome? else {
400            return Ok(None);
401        };
402
403        let mut batch = Batch {
404            changes: Vec::new(),
405            dirty: commits.iter().any(|commit| !commit.changes.is_empty()),
406        };
407        let facts = self.batch_facts(&commits)?;
408        for commit in &commits {
409            for effective in &commit.changes {
410                if let Some(change) = self.change_for(effective, commit.clock.0, &facts) {
411                    batch.changes.push(change);
412                }
413            }
414        }
415        Ok(Some(batch))
416    }
417
418    /// What one batch needs from the index, read once under one lock rather than per
419    /// change.
420    ///
421    /// Two things: the ignore classification of every entry the batch touched, which each
422    /// record carries and an ignored-state selection filters on, and the retained facts of
423    /// every reclassified entry, which is what lets a rule edit that moves an entry into
424    /// the selection be streamed as the upsert a consumer needs to draw the row.
425    fn batch_facts(&self, commits: &[Commit]) -> Result<BatchFacts> {
426        let mut touched: Vec<&PathBuf> = Vec::new();
427        let mut removed: Vec<(&PathBuf, crate::EntryKind)> = Vec::new();
428        let mut reclassified: Vec<&PathBuf> = Vec::new();
429        for effective in commits.iter().flat_map(|commit| &commit.changes) {
430            match effective {
431                EffectiveChange::Inserted { path, .. } | EffectiveChange::Updated { path, .. } => {
432                    touched.push(path);
433                }
434                EffectiveChange::Reclassified { path, .. } => {
435                    reclassified.push(path);
436                }
437                EffectiveChange::Removed { path, kind, .. } => removed.push((path, *kind)),
438                EffectiveChange::Invalidated { .. }
439                | EffectiveChange::ControlUpdated { .. }
440                | EffectiveChange::ControlRefusalUpdated { .. } => {}
441            }
442        }
443
444        self.index.read_with(|index| {
445            let observed = index.observes_controls();
446            let entries = reclassified
447                .into_iter()
448                .filter_map(|path| {
449                    let id = index.lookup(path)?;
450                    let attrs = index.attrs_of(id)?;
451                    Some((
452                        path.clone(),
453                        EntryFacts {
454                            kind: index.kind_of(id)?,
455                            bytes: attrs.size,
456                            allocated: attrs.allocated,
457                            mtime_ns: attrs.mtime_ns,
458                        },
459                    ))
460                })
461                .collect();
462            BatchFacts {
463                ignored: observed.then(|| {
464                    let mut ignored = touched
465                        .into_iter()
466                        .filter_map(|path| match index.is_ignored(path) {
467                            Ok(Some(ignored)) => Some((path.clone(), ignored)),
468                            // Gone from the index, or the index reads no rules; either
469                            // way there is nothing to say about it.
470                            Ok(None) | Err(_) => None,
471                        })
472                        .collect::<BTreeMap<_, _>>();
473                    for (path, kind) in removed {
474                        if index.control_classification_known(path) {
475                            ignored.insert(
476                                path.clone(),
477                                index.control_table().is_ignored(path, kind.is_dir()),
478                            );
479                        }
480                    }
481                    ignored
482                }),
483                reclassified: entries,
484            }
485        })
486    }
487
488    /// Translate one exact effective change into the legacy change view.
489    fn change_for(
490        &self,
491        effective: &EffectiveChange,
492        clock: u64,
493        facts: &BatchFacts,
494    ) -> Option<Change> {
495        match effective {
496            EffectiveChange::Inserted { path, kind, attrs } => {
497                let name = path.file_name()?.to_string_lossy().into_owned();
498                let candidate = crate::query::Candidate {
499                    relative: path,
500                    name: &name,
501                    kind: *kind,
502                    bytes: attrs.size,
503                    allocated: attrs.allocated,
504                    mtime_ns: attrs.mtime_ns,
505                    ignored: facts.is_ignored(path).unwrap_or(false),
506                };
507                self.selection().admits(&candidate).then(|| Change {
508                    path: path.clone(),
509                    kind: ChangeKind::Upsert,
510                    entry_kind: Some(*kind),
511                    bytes: Some(attrs.size),
512                    allocated: Some(attrs.allocated),
513                    mtime_ns: Some(attrs.mtime_ns),
514                    ignored: facts.is_ignored(path),
515                    clock,
516                })
517            }
518            EffectiveChange::Updated { path, kind, previous: _, current } => {
519                let name = path.file_name()?.to_string_lossy().into_owned();
520                let ignored = facts.is_ignored(path).unwrap_or(false);
521                let candidate = crate::query::Candidate {
522                    relative: path,
523                    name: &name,
524                    kind: *kind,
525                    bytes: current.size,
526                    allocated: current.allocated,
527                    mtime_ns: current.mtime_ns,
528                    ignored,
529                };
530                if self.selection().admits(&candidate) {
531                    Some(Change {
532                        path: path.clone(),
533                        kind: ChangeKind::Upsert,
534                        entry_kind: Some(*kind),
535                        bytes: Some(current.size),
536                        allocated: Some(current.allocated),
537                        mtime_ns: Some(current.mtime_ns),
538                        ignored: facts.is_ignored(path),
539                        clock,
540                    })
541                } else if self.admits_by_path(path, &name) {
542                    Some(Change {
543                        path: path.clone(),
544                        kind: ChangeKind::Remove,
545                        entry_kind: None,
546                        bytes: None,
547                        allocated: None,
548                        mtime_ns: None,
549                        ignored: None,
550                        clock,
551                    })
552                } else {
553                    None
554                }
555            }
556            // A removal carries no attributes to filter on, so only the path-shaped parts
557            // of a selection can apply. Filtering it out entirely on a size, time, or
558            // ignored-state bound would hide the disappearance of something the caller was
559            // watching. Its last known kind lets current control state classify it
560            // when the governing rules are known.
561            EffectiveChange::Removed { path, .. } => {
562                let name = path.file_name()?.to_string_lossy().into_owned();
563                self.admits_by_path(path, &name).then(|| Change {
564                    path: path.clone(),
565                    kind: ChangeKind::Remove,
566                    entry_kind: None,
567                    bytes: None,
568                    allocated: None,
569                    mtime_ns: None,
570                    ignored: facts.is_ignored(path),
571                    clock,
572                })
573            }
574            // Escalations are never filtered: they say the consumer's view may have gaps,
575            // and that is true regardless of what the selection asked for.
576            EffectiveChange::Invalidated { path, .. } => Some(Change {
577                path: path.clone(),
578                kind: ChangeKind::Invalidate,
579                entry_kind: None,
580                bytes: None,
581                allocated: None,
582                mtime_ns: None,
583                ignored: None,
584                clock,
585            }),
586            // A rule edit changes what an ignored-state selection contains without
587            // anything on disk changing for the entry, so the entry set the flag promises
588            // is maintained here rather than left to the aggregates: a row that left is
589            // removed and a row that arrived is upserted with the facts to draw it.
590            EffectiveChange::Reclassified { path, previous_ignored, current_ignored } => {
591                let name = path.file_name()?.to_string_lossy().into_owned();
592                // Absent in two cases, each meaning there is nothing to emit: a selection
593                // that admits both partitions, whose membership no edit can change, and an
594                // entry that left the index after the commit, whose removal is already in
595                // this batch.
596                let entry = facts.reclassified.get(path)?;
597                let admits = |ignored: bool| {
598                    self.selection().admits(&crate::query::Candidate {
599                        relative: path,
600                        name: &name,
601                        kind: entry.kind,
602                        bytes: entry.bytes,
603                        allocated: entry.allocated,
604                        mtime_ns: entry.mtime_ns,
605                        ignored,
606                    })
607                };
608                match (admits(*previous_ignored), admits(*current_ignored)) {
609                    (true, false) => Some(Change {
610                        path: path.clone(),
611                        kind: ChangeKind::Remove,
612                        entry_kind: None,
613                        bytes: None,
614                        allocated: None,
615                        mtime_ns: None,
616                        ignored: Some(*current_ignored),
617                        clock,
618                    }),
619                    (_, true) => Some(Change {
620                        path: path.clone(),
621                        kind: ChangeKind::Upsert,
622                        entry_kind: Some(entry.kind),
623                        bytes: Some(entry.bytes),
624                        allocated: Some(entry.allocated),
625                        mtime_ns: Some(entry.mtime_ns),
626                        ignored: Some(*current_ignored),
627                        clock,
628                    }),
629                    _ => None,
630                }
631            }
632            // The legacy watch surface repaints the complete query when `dirty` is true,
633            // so it needs no second row-change vocabulary for control-file effects.
634            // Opened-root consumers read these exact commit variants directly.
635            EffectiveChange::ControlUpdated { .. }
636            | EffectiveChange::ControlRefusalUpdated { .. } => None,
637        }
638    }
639
640    /// Whether the path-shaped parts of the selection admit a path.
641    fn admits_by_path(&self, path: &std::path::Path, name: &str) -> bool {
642        let selection = self.selection();
643        if selection.exclude.iter().any(|pattern| pattern.matches(path, name)) {
644            return false;
645        }
646        selection.include.is_empty()
647            || selection.include.iter().any(|pattern| pattern.matches(path, name))
648    }
649
650    fn selection(&self) -> &Selection {
651        &self.request.query.selection
652    }
653}
654
655fn drain_initial_capture(
656    watcher: &Watcher,
657    index: &IndexHandle,
658    scan: &ScanConfig,
659) -> Result<bool> {
660    let mut dirty = false;
661    for _ in 0..2 {
662        watcher.flush_capture()?;
663        let mut drained = false;
664        for _ in 0..=watcher.capture_backlog_bound() {
665            if watcher
666                .apply_next(index, scan, Duration::ZERO, &mut |commit| {
667                    dirty |= !commit.changes.is_empty();
668                })?
669                .is_none()
670            {
671                drained = true;
672                break;
673            }
674        }
675        if !drained {
676            return Err(Error::ObservationHandoffIncomplete);
677        }
678    }
679    Ok(dirty)
680}
681
682#[cfg(test)]
683mod tests {
684    use super::*;
685
686    /// The watch loop's save throttle, as a table over every state that reaches it.
687    ///
688    /// Two of the three defects review found on this branch were transitions in here, and
689    /// the second was introduced by fixing the first. End-to-end tests could not catch
690    /// either: they observe whether a file changed on disk, which cannot distinguish "not
691    /// due yet" from "due and skipped", nor a cleared flag from a retained one.
692    #[test]
693    fn a_retained_session_rejects_a_request_for_another_root_before_binding() {
694        let a = tempfile::tempdir().expect("root a");
695        let b = tempfile::tempdir().expect("root b");
696        let basis = Basis {
697            root: a.path().into(),
698            scope: crate::query::Scope::default(),
699            content: crate::content::AnalysisSet::NONE,
700        };
701        let delivery = Delivery::new(crate::CachePolicy::Off, None);
702        let (index, _) = crate::open(&basis, &delivery).expect("open a");
703        let handle = IndexHandle::new(index);
704        let before = handle.clock().expect("clock");
705        let request = Request::new(
706            Basis { root: b.path().into(), ..basis },
707            Query::default(),
708            std::time::SystemTime::now(),
709        );
710        assert!(matches!(
711            Session::new(handle.clone(), request, &delivery, WatchConfig::default()),
712            Err(Error::InvalidRequest(crate::query::RequestError::RootMismatch { .. }))
713        ));
714        assert_eq!(handle.clock().expect("clock"), before);
715    }
716
717    #[test]
718    fn a_save_is_due_only_when_a_change_is_pending_and_the_throttle_has_elapsed() {
719        let interval = Duration::from_secs(1);
720        let cases = [
721            // (pending, since last save, due, what this case is)
722            (true, Duration::from_secs(2), true, "pending and past the interval"),
723            (true, interval, true, "pending, exactly at the interval: inclusive"),
724            // The R5 case. Not due *now* -- and the flag stays set, which is the half that
725            // was missing: the idle path saves it once the interval passes.
726            (true, Duration::from_millis(1), false, "pending but throttled"),
727            (false, Duration::from_secs(60), false, "nothing pending, however long it has been"),
728            (false, Duration::ZERO, false, "nothing pending and just saved"),
729        ];
730
731        for (pending, since, want, case) in cases {
732            assert_eq!(save_is_due(pending, since, interval), want, "{case}");
733        }
734    }
735
736    /// A throttled change must survive every outcome except a completed write.
737    #[test]
738    fn only_a_completed_write_clears_the_pending_change() {
739        // The R7 case is Skipped and Failed: clearing the flag for either means the idle
740        // path never retries, so on a quiet tree the change is never persisted at all.
741        assert!(!pending_after(&SaveOutcome::Written), "a completed write persists the change");
742        assert!(
743            pending_after(&SaveOutcome::Skipped),
744            "a skipped save wrote nothing, so the change is still owed to disk",
745        );
746        assert!(
747            pending_after(&SaveOutcome::Failed(Error::Snapshot("failed".into()))),
748            "a failed save must be retried, not forgotten"
749        );
750    }
751
752    /// The sequence that defeated persistence in its most common shape.
753    #[test]
754    fn a_burst_then_a_quiet_tree_still_persists() {
755        let interval = Duration::from_secs(1);
756
757        // A change arrives too soon after the last save, so nothing is written yet.
758        let mut pending = true;
759        assert!(!save_is_due(pending, Duration::from_millis(50), interval));
760        assert!(pending, "the throttle must not consume the change");
761
762        // The tree goes quiet: no further batches will ever arrive. The idle path is the
763        // only remaining caller, and once the interval passes the save must happen.
764        assert!(save_is_due(pending, Duration::from_secs(3), interval));
765
766        // A skip at that point keeps it pending for the next idle tick rather than
767        // silently dropping the session's work.
768        pending = pending_after(&SaveOutcome::Skipped);
769        assert!(pending);
770        pending = pending_after(&SaveOutcome::Written);
771        assert!(!pending, "once written, the loop stops rewriting an unchanged index");
772    }
773
774    #[test]
775    fn skips_and_failures_retry_only_after_another_interval() {
776        let start = Instant::now();
777        let interval = Duration::from_secs(2);
778        let mut persistence = Persistence { pending: true, last_attempt: start };
779        assert!(matches!(
780            persistence.persist_due(start + interval / 2, interval, || panic!("throttled")),
781            SaveOutcome::Skipped
782        ));
783        assert!(matches!(
784            persistence.persist_due(start + interval, interval, || Ok(false)),
785            SaveOutcome::Skipped
786        ));
787        assert!(persistence.pending);
788        assert!(matches!(
789            persistence.persist_due(start + interval, interval, || panic!("skip was throttled")),
790            SaveOutcome::Skipped
791        ));
792        assert!(matches!(
793            persistence.persist_due(start + interval * 2, interval, || {
794                Err(Error::Snapshot("disk unavailable".into()))
795            }),
796            SaveOutcome::Failed(_)
797        ));
798        assert!(persistence.pending);
799        assert!(matches!(
800            persistence
801                .persist_due(start + interval * 2, interval, || panic!("failure was throttled")),
802            SaveOutcome::Skipped
803        ));
804        assert!(matches!(
805            persistence.persist_due(start + interval * 3, interval, || Ok(true)),
806            SaveOutcome::Written
807        ));
808        assert!(!persistence.pending);
809        assert!(matches!(
810            persistence.persist_due(start + interval * 4, interval, || panic!("already persisted")),
811            SaveOutcome::Skipped
812        ));
813    }
814
815    #[test]
816    fn handoff_changes_are_persisted_after_the_tree_goes_quiet() {
817        let root = tempfile::tempdir().expect("root");
818        let cache = tempfile::tempdir().expect("cache");
819        let cache_path = cache.path().join("snapshot");
820        let scan = ScanConfig::default();
821        std::fs::write(root.path().join("before.txt"), b"before").expect("before");
822        let (index, _) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
823        let request = Request::new(
824            Basis {
825                root: root.path().to_path_buf(),
826                scope: scan.clone().into(),
827                content: crate::content::AnalysisSet::NONE,
828            },
829            Query::default(),
830            std::time::SystemTime::now(),
831        );
832        let interval = Duration::from_secs(2);
833        let delivery = Delivery {
834            stale_ok: false,
835            cache: crate::CachePolicy::Auto,
836            cache_path: Some(cache_path.clone()),
837            accept_partial: false,
838            watch: Some(WatchDelivery { interval }),
839            workers: crate::query::Workers::default(),
840            batch_size: ScanConfig::default().batch_size,
841            order: crate::scan::ScanOrder::default(),
842        };
843        let script = tempfile::NamedTempFile::new().expect("script");
844        let (watcher, _sender) =
845            Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
846        std::fs::write(root.path().join("during-handoff.txt"), b"handoff").expect("handoff change");
847        let mut session = Session::finish_initial_handoff(
848            IndexHandle::new(index),
849            request,
850            &delivery,
851            watcher,
852            scan,
853            None,
854        )
855        .expect("handoff");
856        let started = session.persistence.last_attempt;
857        assert!(session.persistence.pending, "handoff changes need persistence too");
858        assert!(matches!(session.persist_due(started), SaveOutcome::Skipped));
859        assert!(!cache_path.exists(), "throttle delays the write");
860        assert!(matches!(session.persist_due(started + interval), SaveOutcome::Written));
861        let restored = crate::snapshot::load(&cache_path).expect("load").expect("saved");
862        assert!(matches!(
863            restored.path_state(std::path::Path::new("during-handoff.txt")),
864            crate::PathState::Present { .. }
865        ));
866        assert!(matches!(session.persist_due(started + interval * 2), SaveOutcome::Skipped));
867    }
868
869    #[test]
870    fn startup_save_failure_keeps_the_session_live_and_retries() {
871        let root = tempfile::tempdir().expect("root");
872        let cache = tempfile::tempdir().expect("cache");
873        // A session reads before it writes, so the fixture must read as no snapshot yet
874        // refuse the write: a parent that cannot become a directory. On Unix that is a
875        // symbolic link to a directory that does not exist yet, since a path through a
876        // regular file fails to open rather than reading as absent; Windows spells a
877        // path through a regular file as not found, so a file serves there.
878        let parent = cache.path().join("blocked");
879        #[cfg(unix)]
880        std::os::unix::fs::symlink(cache.path().join("missing"), &parent)
881            .expect("dangling parent blocks the write");
882        #[cfg(not(unix))]
883        std::fs::write(&parent, b"").expect("file parent blocks the write");
884        let restore = || {
885            #[cfg(unix)]
886            std::fs::create_dir(cache.path().join("missing")).expect("restore the parent");
887            #[cfg(not(unix))]
888            std::fs::remove_file(&parent).expect("restore the parent");
889        };
890        let cache_path = parent.join("snapshot");
891        std::fs::write(root.path().join("file.txt"), b"content").expect("file");
892        let interval = Duration::from_secs(2);
893        let request = Request::new(
894            Basis {
895                root: root.path().to_path_buf(),
896                scope: ScanConfig::default().into(),
897                content: crate::content::AnalysisSet::NONE,
898            },
899            Query::default(),
900            std::time::SystemTime::now(),
901        );
902        let delivery = Delivery {
903            stale_ok: false,
904            cache: crate::CachePolicy::On,
905            cache_path: Some(cache_path.clone()),
906            accept_partial: false,
907            watch: Some(WatchDelivery { interval }),
908            workers: crate::query::Workers::default(),
909            batch_size: ScanConfig::default().batch_size,
910            order: crate::scan::ScanOrder::default(),
911        };
912        let mut session = Session::start(request, delivery).expect("save failure is nonfatal");
913        assert!(session.report(std::time::SystemTime::now()).expect("live report").status.complete);
914        let now = Instant::now();
915        assert!(matches!(session.persist_due(now), SaveOutcome::Failed(_)));
916        restore();
917        assert!(matches!(session.persist_due(now), SaveOutcome::Skipped));
918        assert!(matches!(session.persist_due(now + interval), SaveOutcome::Written));
919        assert!(crate::snapshot::load(&cache_path).expect("read snapshot").is_some());
920    }
921
922    #[test]
923    fn an_update_that_leaves_attribute_selection_emits_remove() {
924        let root = tempfile::tempdir().expect("tempdir");
925        std::fs::write(root.path().join("file.txt"), b"12345678").expect("fixture");
926        let scan = ScanConfig::default();
927        let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
928        assert!(report.is_complete());
929        let request = Request::new(
930            Basis {
931                root: root.path().to_path_buf(),
932                scope: scan.clone().into(),
933                content: crate::content::AnalysisSet::NONE,
934            },
935            Query {
936                selection: Selection { min_size: Some(4), ..Selection::default() },
937                ..Query::default()
938            },
939            std::time::SystemTime::now(),
940        );
941        let delivery = Delivery {
942            stale_ok: false,
943            cache: crate::CachePolicy::Off,
944            cache_path: None,
945            accept_partial: false,
946            watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
947            workers: crate::query::Workers::default(),
948            batch_size: ScanConfig::default().batch_size,
949            order: crate::scan::ScanOrder::default(),
950        };
951        let session =
952            Session::new(IndexHandle::new(index), request, &delivery, WatchConfig::default())
953                .expect("session");
954        let path = PathBuf::from("file.txt");
955        let change = session
956            .change_for(
957                &EffectiveChange::Updated {
958                    path: path.clone(),
959                    kind: EntryKind::File,
960                    previous: crate::Attrs { size: 8, allocated: 8, ..crate::Attrs::default() },
961                    current: crate::Attrs { size: 1, allocated: 1, ..crate::Attrs::default() },
962                },
963                1,
964                &BatchFacts {
965                    ignored: Some(BTreeMap::from([(path, false)])),
966                    reclassified: BTreeMap::new(),
967                },
968            )
969            .expect("membership transition");
970        assert_eq!(change.kind, ChangeKind::Remove);
971    }
972
973    #[test]
974    fn initial_handoff_drains_a_sticky_overflow_after_a_full_intent_queue() {
975        let root = tempfile::tempdir().expect("root");
976        let script = tempfile::NamedTempFile::new().expect("script");
977        let scan = ScanConfig::default();
978        let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
979        assert!(report.is_complete());
980        let handle = IndexHandle::new(index);
981        let config = WatchConfig {
982            settle: Duration::from_millis(1),
983            max_hold: Duration::from_millis(2),
984            event_capacity: 8,
985            batch_path_capacity: 1,
986            intent_capacity: 1,
987            ..WatchConfig::default()
988        };
989        let (watcher, sender) =
990            Watcher::scripted(root.path(), config, script.path()).expect("scripted watcher");
991        std::fs::write(root.path().join("a.txt"), b"a").expect("a");
992        sender.send("create\ta.txt\n").expect("first event");
993        watcher.flush_capture().expect("first barrier fills the intent queue");
994        std::fs::write(root.path().join("b.txt"), b"b").expect("b");
995        sender.send("create\tb.txt\n").expect("second event");
996        watcher.flush_capture().expect("second barrier retains sticky overflow");
997
998        drain_initial_capture(&watcher, &handle, &scan).expect("bounded handoff");
999
1000        assert!(
1001            handle.snapshot().expect("snapshot").lookup(std::path::Path::new("a.txt")).is_some()
1002        );
1003        assert!(
1004            handle.snapshot().expect("snapshot").lookup(std::path::Path::new("b.txt")).is_some()
1005        );
1006        assert!(
1007            watcher
1008                .apply_next(&handle, &scan, Duration::ZERO, &mut |_| {})
1009                .expect("proof poll")
1010                .is_none(),
1011            "no queued or sticky pre-handoff work remains"
1012        );
1013    }
1014
1015    #[test]
1016    fn initial_handoff_enforces_partial_acceptance_after_its_reconciliation() {
1017        let root = tempfile::tempdir().expect("root");
1018        std::fs::write(root.path().join("kept.txt"), b"kept").expect("fixture");
1019        let scan = ScanConfig::default();
1020        let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1021        assert!(report.is_complete());
1022        let request = Request::new(
1023            Basis {
1024                root: root.path().to_path_buf(),
1025                scope: scan.clone().into(),
1026                content: crate::content::AnalysisSet::NONE,
1027            },
1028            Query::default(),
1029            std::time::SystemTime::now(),
1030        );
1031        let _fault = crate::scan::install_walk_hook(root.path(), |_| {
1032            Some(std::io::Error::new(
1033                std::io::ErrorKind::PermissionDenied,
1034                "deterministic handoff refusal",
1035            ))
1036        });
1037        let delivery = Delivery {
1038            stale_ok: false,
1039            cache: crate::CachePolicy::Off,
1040            cache_path: None,
1041            accept_partial: false,
1042            watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1043            workers: crate::query::Workers::default(),
1044            batch_size: ScanConfig::default().batch_size,
1045            order: crate::scan::ScanOrder::default(),
1046        };
1047
1048        let Err(error) = Session::new(
1049            IndexHandle::new(index.clone()),
1050            request.clone(),
1051            &delivery,
1052            WatchConfig::default(),
1053        ) else {
1054            panic!("a partial handoff is refused");
1055        };
1056        assert!(matches!(error, Error::ObservationHandoffIncomplete));
1057
1058        let accepted = Session::new(
1059            IndexHandle::new(index),
1060            request,
1061            &Delivery { accept_partial: true, ..delivery },
1062            WatchConfig::default(),
1063        )
1064        .expect("the caller explicitly accepts a partial handoff");
1065        assert!(
1066            !accepted.report(std::time::SystemTime::now()).expect("partial report").status.complete
1067        );
1068    }
1069
1070    #[test]
1071    fn initial_handoff_rechecks_partial_acceptance_after_draining_capture() {
1072        use std::sync::Arc;
1073        use std::sync::atomic::{AtomicUsize, Ordering};
1074
1075        let root = tempfile::tempdir().expect("root");
1076        std::fs::write(root.path().join("kept.txt"), b"kept").expect("fixture");
1077        let scan = ScanConfig::default();
1078        let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1079        assert!(report.is_complete());
1080        let request = Request::new(
1081            Basis {
1082                root: root.path().to_path_buf(),
1083                scope: scan.clone().into(),
1084                content: crate::content::AnalysisSet::NONE,
1085            },
1086            Query::default(),
1087            std::time::SystemTime::now(),
1088        );
1089        let delivery = Delivery {
1090            stale_ok: false,
1091            cache: crate::CachePolicy::Off,
1092            cache_path: None,
1093            accept_partial: false,
1094            watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1095            workers: crate::query::Workers::default(),
1096            batch_size: ScanConfig::default().batch_size,
1097            order: crate::scan::ScanOrder::default(),
1098        };
1099        let script = tempfile::NamedTempFile::new().expect("script");
1100        let (watcher, sender) =
1101            Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1102        sender.send("rescan\t.\n").expect("queue initial gap");
1103
1104        // The startup reconciliation succeeds. The scripted overflow then reaches the same
1105        // tree during the handoff drain, where its reconciliation fails and must be refused.
1106        let attempts = Arc::new(AtomicUsize::new(0));
1107        let hook_attempts = Arc::clone(&attempts);
1108        let _fault = crate::scan::install_walk_hook(root.path(), move |_| {
1109            (hook_attempts.fetch_add(1, Ordering::SeqCst) > 0).then(|| {
1110                std::io::Error::new(
1111                    std::io::ErrorKind::PermissionDenied,
1112                    "deterministic drain-only refusal",
1113                )
1114            })
1115        });
1116
1117        let Err(error) = Session::finish_initial_handoff(
1118            IndexHandle::new(index),
1119            request,
1120            &delivery,
1121            watcher,
1122            scan,
1123            None,
1124        ) else {
1125            panic!("a partial state created while draining is refused");
1126        };
1127        assert!(matches!(error, Error::ObservationHandoffIncomplete));
1128        assert!(attempts.load(Ordering::SeqCst) > 1, "the drain ran after startup reconciliation");
1129    }
1130
1131    /// The handoff pass restarts the walk counters and enters `Revalidating`, whatever
1132    /// the first pass left in the handle: with a scripted watcher that reports nothing,
1133    /// the counts afterwards are exactly one walk of the tree. Identical event streams
1134    /// also make the observed and unobserved reports exactly comparable.
1135    #[test]
1136    fn the_handoff_pass_restarts_the_counts() {
1137        let root = tempfile::tempdir().expect("root");
1138        let mut bytes = 0;
1139        for directory in 0..2 {
1140            let dir = root.path().join(format!("d{directory}"));
1141            std::fs::create_dir(&dir).expect("directory");
1142            for file in 0..3 {
1143                let size = directory * 3 + file + 1;
1144                std::fs::write(dir.join(format!("f{file}.txt")), vec![b'.'; size]).expect("file");
1145                bytes += size as u64;
1146            }
1147        }
1148        let scan = ScanConfig::default();
1149        let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1150        assert!(report.is_complete());
1151        let request = Request::new(
1152            Basis {
1153                root: root.path().to_path_buf(),
1154                scope: scan.clone().into(),
1155                content: crate::content::AnalysisSet::NONE,
1156            },
1157            Query { views: vec![crate::query::ViewSpec::Summary], ..Query::default() },
1158            std::time::SystemTime::now(),
1159        );
1160        let delivery = Delivery {
1161            stale_ok: false,
1162            cache: crate::CachePolicy::Off,
1163            cache_path: None,
1164            accept_partial: false,
1165            watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1166            workers: crate::query::Workers::default(),
1167            batch_size: ScanConfig::default().batch_size,
1168            order: crate::scan::ScanOrder::default(),
1169        };
1170        let script = tempfile::NamedTempFile::new().expect("script");
1171        let (watcher, _sender) =
1172            Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1173        let progress = crate::Progress::new();
1174        progress.add_walked(100, 100, 100, 100);
1175        progress.enter(crate::ProgressPhase::Saving);
1176
1177        let session = Session::finish_initial_handoff(
1178            IndexHandle::new(index.clone()),
1179            request.clone(),
1180            &delivery,
1181            watcher,
1182            scan.clone(),
1183            Some(&progress),
1184        )
1185        .expect("handoff");
1186
1187        let snapshot = progress.snapshot();
1188        assert_eq!(snapshot.phase, crate::ProgressPhase::Revalidating);
1189        assert_eq!(
1190            (snapshot.directories, snapshot.files, snapshot.bytes),
1191            (3, 6, bytes),
1192            "one walk of the root and its two directories, the first pass not added in"
1193        );
1194        assert_eq!(snapshot.allocated, report.allocated_walked, "allocated restarts with them");
1195
1196        let (plain_watcher, _plain_sender) =
1197            Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1198        let plain = Session::finish_initial_handoff(
1199            IndexHandle::new(index),
1200            request,
1201            &delivery,
1202            plain_watcher,
1203            scan,
1204            None,
1205        )
1206        .expect("plain handoff");
1207        let generated_at = std::time::SystemTime::now();
1208        let observed_report = session.report(generated_at).expect("observed report");
1209        let mut plain_report = plain.report(generated_at).expect("plain report");
1210        plain_report.provenance = observed_report.provenance.clone();
1211        let json = |report: &Report| {
1212            crate::report_format::render(report, crate::report_format::Format::Json, false)
1213                .expect("render")
1214        };
1215        assert_eq!(json(&plain_report), json(&observed_report));
1216        assert_eq!(progress.snapshot(), snapshot, "the second handoff and reads were not observed");
1217    }
1218
1219    /// A record says what the index can be asked, and nothing more.
1220    ///
1221    /// The three answers are distinct and a consumer acts on each differently: a bit, "no
1222    /// rules were read", and "there is no such entry". A set of the ignored paths collapsed
1223    /// the last two into `false`, so a batch that created and removed one file in the same
1224    /// window upserted it as unignored before removing it -- a classification claim about
1225    /// an entry that never survived the batch.
1226    #[test]
1227    fn a_record_claims_a_classification_only_for_an_entry_the_index_still_holds() {
1228        let unobserved = BatchFacts { ignored: None, reclassified: BTreeMap::new() };
1229        assert_eq!(unobserved.is_ignored(std::path::Path::new("any.txt")), None);
1230
1231        let observed = BatchFacts {
1232            ignored: Some(BTreeMap::from([
1233                (PathBuf::from("build/out.bin"), true),
1234                (PathBuf::from("src/main.rs"), false),
1235            ])),
1236            reclassified: BTreeMap::new(),
1237        };
1238        assert_eq!(observed.is_ignored(std::path::Path::new("build/out.bin")), Some(true));
1239        assert_eq!(observed.is_ignored(std::path::Path::new("src/main.rs")), Some(false));
1240        assert_eq!(
1241            observed.is_ignored(std::path::Path::new("gone.tmp")),
1242            None,
1243            "an entry the batch removed is in neither partition, not in the unignored one"
1244        );
1245    }
1246
1247    /// A start is the open and then the handoff revalidation, a second pass whose
1248    /// counts restart, which keeps the line moving on a large tree after the save,
1249    /// where a frozen count would look like a hang. The session it returns is the one
1250    /// [`Session::start`] returns, and it reports nothing further through the handle
1251    /// once started. The exact reset and report equivalence are pinned by the scripted
1252    /// test above; independent native watchers may replay different creation hints and
1253    /// record different legitimate setup-gap diagnostics, so only lower bounds hold.
1254    #[test]
1255    fn a_started_session_reports_its_second_pass_and_then_nothing() {
1256        let root = tempfile::tempdir().expect("root");
1257        let cache = tempfile::tempdir().expect("cache");
1258        let mut bytes = 0;
1259        for directory in 0..4 {
1260            let dir = root.path().join(format!("d{directory}"));
1261            std::fs::create_dir(&dir).expect("directory");
1262            for file in 0..3 {
1263                let size = directory * 3 + file + 1;
1264                std::fs::write(dir.join(format!("f{file}.txt")), vec![b'.'; size]).expect("file");
1265                bytes += size as u64;
1266            }
1267        }
1268        let request = || {
1269            Request::new(
1270                Basis {
1271                    root: root.path().to_path_buf(),
1272                    scope: ScanConfig::default().into(),
1273                    content: crate::content::AnalysisSet::NONE,
1274                },
1275                Query::default(),
1276                std::time::UNIX_EPOCH,
1277            )
1278        };
1279        let delivery = Delivery {
1280            stale_ok: false,
1281            cache: crate::CachePolicy::Auto,
1282            cache_path: Some(cache.path().join("snapshot")),
1283            accept_partial: false,
1284            watch: Some(WatchDelivery { interval: Duration::from_secs(2) }),
1285            workers: crate::query::Workers::default(),
1286            batch_size: 4,
1287            order: crate::scan::ScanOrder::default(),
1288        };
1289
1290        let progress = crate::Progress::new();
1291        let session = Session::start_with_progress(request(), delivery.clone(), &progress)
1292            .expect("observed start");
1293        let after_start = progress.snapshot();
1294        assert_eq!(
1295            after_start.phase,
1296            crate::ProgressPhase::Revalidating,
1297            "the handoff revalidation follows the joined save"
1298        );
1299        // The closing pass restarts the counters, so they show its walk alone. A backend
1300        // that reports a file created just before the watch began can make that pass
1301        // read more, never less.
1302        assert!(after_start.directories >= 5, "the root and four children: {after_start:?}");
1303        assert!(after_start.files >= 12, "{after_start:?}");
1304        assert!(after_start.bytes >= bytes, "{after_start:?}");
1305        assert_eq!(after_start.analysis, None);
1306
1307        let plain = Session::start(request(), delivery).expect("plain start");
1308        let generated_at = std::time::SystemTime::now();
1309        let observed_report = session.report(generated_at).expect("observed report");
1310        let plain_report = plain.report(generated_at).expect("plain report");
1311        assert!(observed_report.status.complete);
1312        assert!(plain_report.status.complete);
1313        assert_eq!(progress.snapshot(), after_start, "the second start was not observed");
1314    }
1315}