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    /// The digest of the identity of the answer [`Self::changed_report`] last handed out
167    /// ([`RepaintDigest`]).
168    presented: Option<u128>,
169}
170
171/// A 128-bit FNV-1a digest of a repaint's identity, written section by section.
172///
173/// A digest rather than the identity itself: the identity is the whole rendered answer,
174/// and a session keeping it would hold a second copy of a `--view full --limit all`
175/// report for its whole life. FNV-1a is the hash the engine's fingerprints already use,
176/// at its 128-bit width, so two identities a session compares collide with a chance no
177/// repaint rule has to consider, and it needs no dependency. Each section's length is
178/// mixed in where it ends, so no two splits of the same bytes into sections digest alike.
179struct RepaintDigest {
180    hash: u128,
181    section: u64,
182}
183
184impl RepaintDigest {
185    const OFFSET_BASIS: u128 = 0x6c62_272e_07bb_0142_62b8_2175_6295_c58d;
186    const PRIME: u128 = 0x0000_0000_0100_0000_0000_0000_0000_013b;
187
188    const fn new() -> Self {
189        Self { hash: Self::OFFSET_BASIS, section: 0 }
190    }
191
192    fn mix(&mut self, bytes: &[u8]) {
193        for byte in bytes {
194            self.hash ^= u128::from(*byte);
195            self.hash = self.hash.wrapping_mul(Self::PRIME);
196        }
197    }
198
199    /// Close the section written so far.
200    fn end_section(&mut self) {
201        let length = std::mem::take(&mut self.section);
202        self.mix(&length.to_le_bytes());
203    }
204
205    fn finish(mut self) -> u128 {
206        self.end_section();
207        self.hash
208    }
209}
210
211impl std::io::Write for RepaintDigest {
212    fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
213        self.mix(bytes);
214        self.section = self.section.wrapping_add(u64::try_from(bytes.len()).unwrap_or(u64::MAX));
215        Ok(bytes.len())
216    }
217
218    fn flush(&mut self) -> std::io::Result<()> {
219        Ok(())
220    }
221}
222
223impl Session {
224    /// Open a tree and bind observation under the shared execution plan.
225    ///
226    /// Startup persistence is joined before binding the session. A save failure is
227    /// returned by the first `persist_due` call, so it does not discard a valid live
228    /// answer. Filesystem and observation failures still fail startup.
229    pub fn start(request: Request, delivery: Delivery) -> Result<Self> {
230        Self::start_observed(request, delivery, None)
231    }
232
233    /// [`Self::start`], reporting the initial scan through `progress` as it runs.
234    ///
235    /// The same session as [`Self::start`]; the handle observes the start and changes
236    /// nothing about it. A start is two passes: the open (cold, or a load and
237    /// revalidation) with its save, if it writes one, joined under
238    /// [`ProgressPhase::Saving`](crate::ProgressPhase), and then the revalidation that
239    /// closes the gap between that walk and the bound watcher. The second pass begins
240    /// again at [`ProgressPhase::Revalidating`](crate::ProgressPhase) with the walk
241    /// counters restarted, so they end at the tree's totals once, not twice. Once this
242    /// returns, the handle observes one more thing: the first answer's build, asked for
243    /// through [`Self::changed_report_with_progress`]. The session's repaints are the
244    /// progress from there.
245    pub fn start_with_progress(
246        request: Request,
247        delivery: Delivery,
248        progress: &crate::Progress,
249    ) -> Result<Self> {
250        Self::start_observed(request, delivery, Some(progress))
251    }
252
253    fn start_observed(
254        request: Request,
255        mut delivery: Delivery,
256        progress: Option<&crate::Progress>,
257    ) -> Result<Self> {
258        delivery.watch.get_or_insert_with(WatchDelivery::default);
259        let plan =
260            crate::plan(&request, &delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
261        let (index, report, pending, _diagnostics) =
262            crate::execute(&plan, &request.basis, false, progress)?;
263        let startup_save_error = pending.join().err();
264        let index = std::sync::Arc::into_inner(index)
265            .expect("the joined writer released the only other reference");
266        let mut session = Self::new_observed(
267            IndexHandle::new(index),
268            request,
269            &delivery,
270            WatchConfig::default(),
271            progress,
272        )?;
273        session.persistence.pending |= startup_save_error.is_some() || !report.is_complete();
274        session.startup_save_error = startup_save_error;
275        Ok(session)
276    }
277
278    /// Persist pending changes when the watch delivery's interval has elapsed.
279    ///
280    /// Call this after batches and idle timeouts. Only a completed write clears pending
281    /// changes; a refused or failed write is retried on a later interval. `now` is a
282    /// monotonic caller-supplied clock so wall-clock corrections cannot postpone saves.
283    pub fn persist_due(&mut self, now: Instant) -> SaveOutcome {
284        if let Some(error) = self.startup_save_error.take() {
285            self.persistence.last_attempt = now;
286            return SaveOutcome::Failed(error);
287        }
288        let interval = self.plan.delivery().watch.expect("watch plan").interval;
289        self.persistence.persist_due(now, interval, || {
290            if !self.plan.persists() || self.plan.delivery().cache_path.is_none() {
291                return Ok(false);
292            }
293            let index = self.index.snapshot()?;
294            crate::persist_index(&index, &self.plan)
295        })
296    }
297
298    /// Start watching an already-opened index, answering `request` as the tree changes.
299    ///
300    /// `request` carries its own `now`, fixed when it was built: a watch answers one
301    /// request as the tree changes, and a relative time window that slid under it would
302    /// make two repaints answer two different questions.
303    ///
304    /// `delivery` is the one its caller opened the index under. It is taken rather than
305    /// composed here because the cache policy is part of it: a session built against a
306    /// fabricated `Delivery` read `cache: Auto` whatever the caller had asked for, so
307    /// [`RequestError::WatchCacheOnly`](crate::query::RequestError::WatchCacheOnly) could
308    /// not fire inside the engine at all and the rule held only at the two front doors
309    /// (fdu-i18y). Starting a session is what a watch *is*, so the delivery is read as a
310    /// watch whether or not the caller remembered to say so.
311    ///
312    /// Refusals are in the order every route publishes: what no delivery can carry first,
313    /// then what this index was taken under, then what it holds. The request-level rule
314    /// speaks first, so a library caller watching a depth-2 request is told the watch
315    /// cannot narrow its scope rather than that the index has another scope.
316    ///
317    /// # Errors
318    ///
319    /// [`Error::InvalidRequest`] when a watch cannot deliver the request -- a narrowed scan
320    /// scope, content analysis nothing re-reads, a snapshot nothing verified -- or when
321    /// this index cannot answer it, including a selection by ignored state over an index
322    /// that observed no control state.
323    /// [`Error::ScanScopeMismatch`] when the index was not taken under the request's scope.
324    pub fn new(
325        index: IndexHandle,
326        request: Request,
327        delivery: &Delivery,
328        watch: WatchConfig,
329    ) -> Result<Self> {
330        Self::new_observed(index, request, delivery, watch, None)
331    }
332
333    fn new_observed(
334        index: IndexHandle,
335        request: Request,
336        delivery: &Delivery,
337        watch: WatchConfig,
338        progress: Option<&crate::Progress>,
339    ) -> Result<Self> {
340        let root = index.root_path()?;
341        let scan = request.basis.scope.scan_config(delivery);
342        // What no delivery can carry, before anything stored is read and before the
343        // backend is bound: this is the rule each surface used to keep for itself, so a
344        // library caller could watch what `--watch` has always refused.
345        let delivery =
346            Delivery { watch: Some(delivery.watch.unwrap_or_default()), ..delivery.clone() };
347        crate::plan(&request, &delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
348        crate::validate_basis_root(&root, &request.basis)?;
349        // Reject an out-of-scope watch before the backend is bound, so a rejected run
350        // never leaves a watcher registered on the tree.
351        scan.validate_for_scope(index.scope()?)?;
352        // The scope check above proved the index was taken under exactly this scan's
353        // identity, control tier included, so what remains is what this index holds.
354        let held = Basis {
355            root: root.clone(),
356            scope: scan.clone().into(),
357            content: index.read_with(crate::Index::content_set)?,
358        };
359        request.validate_read(&held).map_err(Error::InvalidRequest)?;
360        // Bind observation before closing the gap from the scan that produced `index`.
361        // The full reconciliation catches a mutation that completed before registration;
362        // the capture drain applies every hint observed while that pass ran.
363        let watcher = Watcher::new(&root, watch)?;
364        Self::finish_initial_handoff(index, request, &delivery, watcher, scan, progress)
365    }
366
367    /// Finish the two-part initial handoff after observation has been bound.
368    ///
369    /// Kept separate so the scripted watcher exercises the same reconciliation, drain, and
370    /// acceptance boundary as an OS watcher. Once this returns, later partial observations are
371    /// valid live state; `accept_partial` governs only the coherent state handed to the caller.
372    ///
373    /// `progress` observes the handoff's own walk and drain only. The session keeps
374    /// `scan` without it, so nothing it reconciles later reports through a handle whose
375    /// poller has long since stopped.
376    fn finish_initial_handoff(
377        index: IndexHandle,
378        request: Request,
379        delivery: &Delivery,
380        watcher: Watcher,
381        scan: ScanConfig,
382        progress: Option<&crate::Progress>,
383    ) -> Result<Self> {
384        if let Some(progress) = progress {
385            progress.begin_pass(crate::ProgressPhase::Revalidating);
386        }
387        let observed = ScanConfig { progress: progress.cloned(), ..scan.clone() };
388        let mut dirty = false;
389        let reconciliation = crate::scan::reconcile_handle(&index, &observed, &mut |commit| {
390            dirty |= !commit.changes.is_empty();
391        })?;
392        if !reconciliation.scan.is_complete() && !delivery.accept_partial {
393            return Err(Error::ObservationHandoffIncomplete);
394        }
395        dirty |= drain_initial_capture(&watcher, &index, &observed)?;
396        if !delivery.accept_partial
397            && !index.read_with(|index| crate::query::TreeStatus::of(index, &request).complete)?
398        {
399            return Err(Error::ObservationHandoffIncomplete);
400        }
401        let plan =
402            crate::plan(&request, delivery, crate::Route::Watch).map_err(Error::InvalidRequest)?;
403        Ok(Self {
404            index,
405            watcher,
406            scan,
407            request,
408            plan,
409            persistence: Persistence { pending: dirty, last_attempt: Instant::now() },
410            startup_save_error: None,
411            presented: None,
412        })
413    }
414
415    /// The request this session answers.
416    pub fn request(&self) -> &Request {
417        &self.request
418    }
419
420    /// The query this session answers.
421    pub fn query(&self) -> &Query {
422        &self.request.query
423    }
424
425    /// Render the current answer.
426    ///
427    /// The same `report` a one-shot run produces, from the same index, which is what
428    /// makes "watch is the same query repeated" true rather than aspirational.
429    pub fn report(&self, generated_at: std::time::SystemTime) -> Result<Report> {
430        let index = self.index.snapshot()?;
431        report(&index, &self.request, generated_at)
432    }
433
434    /// The current answer, unless a reader of `format` would see nothing new in it.
435    ///
436    /// A tree changes more often than its answer does. A touch that leaves a file's
437    /// size alone, or a write to an entry the selection leaves out, moves the index and
438    /// marks its batch dirty, and a caller that repaints on every dirty batch then prints
439    /// the rows it printed a moment ago under a new timestamp (fdu-wb5n). The answer's
440    /// identity is what `format` renders of it with its generation instant held fixed,
441    /// together with its tree status, its source and freshness, and the diagnostics a
442    /// frontend writes beside it: its notes and tips
443    /// ([`diagnostic_lines`](crate::report_format::diagnostic_lines)) and its warnings
444    /// ([`report_warnings`](crate::report_format::report_warnings)). So a tree that shows
445    /// sizes repaints when a size moves, a listing that shows dates repaints when a date
446    /// does, machine output repaints when any field it carries does, and a retained
447    /// observation gap, a coverage change, a freshness change, or a new note such as a
448    /// `.gitignore` refused for its limits repaints whether or not `format` shows it,
449    /// while the instant a repaint is generated at never counts on its own. A caller that
450    /// prints no notes, as a quiet one does, may repaint once more than it had to, never
451    /// once fewer. A session keeps a 128-bit digest of the identity, not the identity.
452    /// The change records of [`Self::next_batch`] are never deduplicated; only this
453    /// repaint is. The first call always answers, and [`Self::report`] always answers.
454    ///
455    /// # Errors
456    ///
457    /// As [`Self::report`], and [`Error::Io`] when `format` cannot render the answer.
458    pub fn changed_report(
459        &mut self,
460        generated_at: std::time::SystemTime,
461        format: crate::report_format::Format,
462        options: crate::report_format::RenderOptions,
463    ) -> Result<Option<Report>> {
464        use std::io::Write as _;
465
466        let mut report = self.report(generated_at)?;
467        let pinned = std::time::SystemTime::UNIX_EPOCH;
468        let stamped = std::mem::replace(&mut report.provenance.generated_at, pinned);
469        let mut identity = RepaintDigest::new();
470        let rendered =
471            crate::report_format::write_with_options(&report, format, options, &mut identity)
472                .and_then(|()| {
473                    identity.end_section();
474                    write!(
475                        identity,
476                        "{:?}\n{:?}\n{:?}",
477                        report.status, report.provenance.source, report.provenance.freshness
478                    )?;
479                    identity.end_section();
480                    let notes = crate::report_format::diagnostic_lines(&report).into_lines();
481                    for line in notes.iter().chain(&crate::report_format::report_warnings(&report))
482                    {
483                        writeln!(identity, "{line}")?;
484                    }
485                    Ok(())
486                });
487        report.provenance.generated_at = stamped;
488        rendered.map_err(|error| Error::io(&self.request.basis.root, error))?;
489        let identity = identity.finish();
490        if self.presented == Some(identity) {
491            return Ok(None);
492        }
493        self.presented = Some(identity);
494        Ok(Some(report))
495    }
496
497    /// [`Self::changed_report`], reporting the answer's construction through `progress`.
498    ///
499    /// The same answer as [`Self::changed_report`]; the handle observes the build and
500    /// changes nothing about it. The build enters
501    /// [`ProgressPhase::Summarizing`](crate::ProgressPhase), as a one-shot report's does,
502    /// so a caller that drew the start through [`Self::start_with_progress`] can keep
503    /// drawing until the first answer exists: a heavy view over a large tree takes
504    /// seconds to build, and a line stopped when the start returned said nothing about
505    /// them (fdu-wku3). Later repaints are the progress from there and take
506    /// [`Self::changed_report`].
507    ///
508    /// # Errors
509    ///
510    /// As [`Self::changed_report`].
511    pub fn changed_report_with_progress(
512        &mut self,
513        generated_at: std::time::SystemTime,
514        format: crate::report_format::Format,
515        options: crate::report_format::RenderOptions,
516        progress: &crate::Progress,
517    ) -> Result<Option<Report>> {
518        progress.enter(crate::ProgressPhase::Summarizing);
519        self.changed_report(generated_at, format, options)
520    }
521
522    /// A consistent copy of the current index.
523    ///
524    /// Used to persist a live session without holding a lock across the write.
525    pub fn index_snapshot(&self) -> Result<crate::Index> {
526        self.index.snapshot()
527    }
528
529    /// Wait for the next batch of changes, up to `timeout`.
530    ///
531    /// Returns `None` when nothing arrived in the window, which is the idle case and
532    /// costs no filesystem work.
533    ///
534    /// Takes `&mut self` because consuming from the event queue is a mutation: two
535    /// callers draining one session would each see an arbitrary half of the stream.
536    pub fn next_batch(&mut self, timeout: Duration) -> Result<Option<Batch>> {
537        let mut commits: Vec<Commit> = Vec::new();
538        let outcome =
539            self.watcher.apply_next(&self.index, &self.scan, timeout, &mut |commit: &Commit| {
540                commits.push(commit.clone());
541            });
542
543        self.persistence.pending |= commits.iter().any(|commit| !commit.changes.is_empty());
544        let Some(_report) = outcome? else {
545            return Ok(None);
546        };
547
548        let mut batch = Batch {
549            changes: Vec::new(),
550            dirty: commits.iter().any(|commit| !commit.changes.is_empty()),
551        };
552        let facts = self.batch_facts(&commits)?;
553        for commit in &commits {
554            for effective in &commit.changes {
555                if let Some(change) = self.change_for(effective, commit.clock.0, &facts) {
556                    batch.changes.push(change);
557                }
558            }
559        }
560        Ok(Some(batch))
561    }
562
563    /// What one batch needs from the index, read once under one lock rather than per
564    /// change.
565    ///
566    /// Two things: the ignore classification of every entry the batch touched, which each
567    /// record carries and an ignored-state selection filters on, and the retained facts of
568    /// every reclassified entry, which is what lets a rule edit that moves an entry into
569    /// the selection be streamed as the upsert a consumer needs to draw the row.
570    fn batch_facts(&self, commits: &[Commit]) -> Result<BatchFacts> {
571        let mut touched: Vec<&PathBuf> = Vec::new();
572        let mut removed: Vec<(&PathBuf, crate::EntryKind)> = Vec::new();
573        let mut reclassified: Vec<&PathBuf> = Vec::new();
574        for effective in commits.iter().flat_map(|commit| &commit.changes) {
575            match effective {
576                EffectiveChange::Inserted { path, .. } | EffectiveChange::Updated { path, .. } => {
577                    touched.push(path);
578                }
579                EffectiveChange::Reclassified { path, .. } => {
580                    reclassified.push(path);
581                }
582                EffectiveChange::Removed { path, kind, .. } => removed.push((path, *kind)),
583                EffectiveChange::Invalidated { .. }
584                | EffectiveChange::ControlUpdated { .. }
585                | EffectiveChange::ControlRefusalUpdated { .. } => {}
586            }
587        }
588
589        self.index.read_with(|index| {
590            let observed = index.observes_controls();
591            let entries = reclassified
592                .into_iter()
593                .filter_map(|path| {
594                    let id = index.lookup(path)?;
595                    let attrs = index.attrs_of(id)?;
596                    Some((
597                        path.clone(),
598                        EntryFacts {
599                            kind: index.kind_of(id)?,
600                            bytes: attrs.size,
601                            allocated: attrs.allocated,
602                            mtime_ns: attrs.mtime_ns,
603                        },
604                    ))
605                })
606                .collect();
607            BatchFacts {
608                ignored: observed.then(|| {
609                    let mut ignored = touched
610                        .into_iter()
611                        .filter_map(|path| match index.is_ignored(path) {
612                            Ok(Some(ignored)) => Some((path.clone(), ignored)),
613                            // Gone from the index, or the index reads no rules; either
614                            // way there is nothing to say about it.
615                            Ok(None) | Err(_) => None,
616                        })
617                        .collect::<BTreeMap<_, _>>();
618                    for (path, kind) in removed {
619                        if index.control_classification_known(path) {
620                            ignored.insert(
621                                path.clone(),
622                                index.control_table().is_ignored(path, kind.is_dir()),
623                            );
624                        }
625                    }
626                    ignored
627                }),
628                reclassified: entries,
629            }
630        })
631    }
632
633    /// Translate one exact effective change into the legacy change view.
634    fn change_for(
635        &self,
636        effective: &EffectiveChange,
637        clock: u64,
638        facts: &BatchFacts,
639    ) -> Option<Change> {
640        match effective {
641            EffectiveChange::Inserted { path, kind, attrs } => {
642                let name = path.file_name()?.to_string_lossy().into_owned();
643                let candidate = crate::query::Candidate {
644                    relative: path,
645                    name: &name,
646                    kind: *kind,
647                    bytes: attrs.size,
648                    allocated: attrs.allocated,
649                    mtime_ns: attrs.mtime_ns,
650                    ignored: facts.is_ignored(path).unwrap_or(false),
651                };
652                self.selection().admits(&candidate).then(|| Change {
653                    path: path.clone(),
654                    kind: ChangeKind::Upsert,
655                    entry_kind: Some(*kind),
656                    bytes: Some(attrs.size),
657                    allocated: Some(attrs.allocated),
658                    mtime_ns: Some(attrs.mtime_ns),
659                    ignored: facts.is_ignored(path),
660                    clock,
661                })
662            }
663            EffectiveChange::Updated { path, kind, previous: _, current } => {
664                let name = path.file_name()?.to_string_lossy().into_owned();
665                let ignored = facts.is_ignored(path).unwrap_or(false);
666                let candidate = crate::query::Candidate {
667                    relative: path,
668                    name: &name,
669                    kind: *kind,
670                    bytes: current.size,
671                    allocated: current.allocated,
672                    mtime_ns: current.mtime_ns,
673                    ignored,
674                };
675                if self.selection().admits(&candidate) {
676                    Some(Change {
677                        path: path.clone(),
678                        kind: ChangeKind::Upsert,
679                        entry_kind: Some(*kind),
680                        bytes: Some(current.size),
681                        allocated: Some(current.allocated),
682                        mtime_ns: Some(current.mtime_ns),
683                        ignored: facts.is_ignored(path),
684                        clock,
685                    })
686                } else if self.admits_by_path(path, &name) {
687                    Some(Change {
688                        path: path.clone(),
689                        kind: ChangeKind::Remove,
690                        entry_kind: None,
691                        bytes: None,
692                        allocated: None,
693                        mtime_ns: None,
694                        ignored: None,
695                        clock,
696                    })
697                } else {
698                    None
699                }
700            }
701            // A removal carries no attributes to filter on, so only the path-shaped parts
702            // of a selection can apply. Filtering it out entirely on a size, time, or
703            // ignored-state bound would hide the disappearance of something the caller was
704            // watching. Its last known kind lets current control state classify it
705            // when the governing rules are known.
706            EffectiveChange::Removed { path, .. } => {
707                let name = path.file_name()?.to_string_lossy().into_owned();
708                self.admits_by_path(path, &name).then(|| Change {
709                    path: path.clone(),
710                    kind: ChangeKind::Remove,
711                    entry_kind: None,
712                    bytes: None,
713                    allocated: None,
714                    mtime_ns: None,
715                    ignored: facts.is_ignored(path),
716                    clock,
717                })
718            }
719            // Escalations are never filtered: they say the consumer's view may have gaps,
720            // and that is true regardless of what the selection asked for.
721            EffectiveChange::Invalidated { path, .. } => Some(Change {
722                path: path.clone(),
723                kind: ChangeKind::Invalidate,
724                entry_kind: None,
725                bytes: None,
726                allocated: None,
727                mtime_ns: None,
728                ignored: None,
729                clock,
730            }),
731            // A rule edit changes what an ignored-state selection contains without
732            // anything on disk changing for the entry, so the entry set the flag promises
733            // is maintained here rather than left to the aggregates: a row that left is
734            // removed and a row that arrived is upserted with the facts to draw it.
735            EffectiveChange::Reclassified { path, previous_ignored, current_ignored } => {
736                let name = path.file_name()?.to_string_lossy().into_owned();
737                // Absent in two cases, each meaning there is nothing to emit: a selection
738                // that admits both partitions, whose membership no edit can change, and an
739                // entry that left the index after the commit, whose removal is already in
740                // this batch.
741                let entry = facts.reclassified.get(path)?;
742                let admits = |ignored: bool| {
743                    self.selection().admits(&crate::query::Candidate {
744                        relative: path,
745                        name: &name,
746                        kind: entry.kind,
747                        bytes: entry.bytes,
748                        allocated: entry.allocated,
749                        mtime_ns: entry.mtime_ns,
750                        ignored,
751                    })
752                };
753                match (admits(*previous_ignored), admits(*current_ignored)) {
754                    (true, false) => Some(Change {
755                        path: path.clone(),
756                        kind: ChangeKind::Remove,
757                        entry_kind: None,
758                        bytes: None,
759                        allocated: None,
760                        mtime_ns: None,
761                        ignored: Some(*current_ignored),
762                        clock,
763                    }),
764                    (_, true) => Some(Change {
765                        path: path.clone(),
766                        kind: ChangeKind::Upsert,
767                        entry_kind: Some(entry.kind),
768                        bytes: Some(entry.bytes),
769                        allocated: Some(entry.allocated),
770                        mtime_ns: Some(entry.mtime_ns),
771                        ignored: Some(*current_ignored),
772                        clock,
773                    }),
774                    _ => None,
775                }
776            }
777            // The legacy watch surface repaints the complete query when `dirty` is true,
778            // so it needs no second row-change vocabulary for control-file effects.
779            // Opened-root consumers read these exact commit variants directly.
780            EffectiveChange::ControlUpdated { .. }
781            | EffectiveChange::ControlRefusalUpdated { .. } => None,
782        }
783    }
784
785    /// Whether the path-shaped parts of the selection admit a path.
786    fn admits_by_path(&self, path: &std::path::Path, name: &str) -> bool {
787        let selection = self.selection();
788        if selection.exclude.iter().any(|pattern| pattern.matches(path, name)) {
789            return false;
790        }
791        selection.include.is_empty()
792            || selection.include.iter().any(|pattern| pattern.matches(path, name))
793    }
794
795    fn selection(&self) -> &Selection {
796        &self.request.query.selection
797    }
798}
799
800fn drain_initial_capture(
801    watcher: &Watcher,
802    index: &IndexHandle,
803    scan: &ScanConfig,
804) -> Result<bool> {
805    let mut dirty = false;
806    for _ in 0..2 {
807        watcher.flush_capture()?;
808        let mut drained = false;
809        for _ in 0..=watcher.capture_backlog_bound() {
810            if watcher
811                .apply_next(index, scan, Duration::ZERO, &mut |commit| {
812                    dirty |= !commit.changes.is_empty();
813                })?
814                .is_none()
815            {
816                drained = true;
817                break;
818            }
819        }
820        if !drained {
821            return Err(Error::ObservationHandoffIncomplete);
822        }
823    }
824    Ok(dirty)
825}
826
827#[cfg(test)]
828mod tests {
829    use super::*;
830
831    /// The watch loop's save throttle, as a table over every state that reaches it.
832    ///
833    /// Two of the three defects review found on this branch were transitions in here, and
834    /// the second was introduced by fixing the first. End-to-end tests could not catch
835    /// either: they observe whether a file changed on disk, which cannot distinguish "not
836    /// due yet" from "due and skipped", nor a cleared flag from a retained one.
837    #[test]
838    fn a_retained_session_rejects_a_request_for_another_root_before_binding() {
839        let a = tempfile::tempdir().expect("root a");
840        let b = tempfile::tempdir().expect("root b");
841        let basis = Basis {
842            root: a.path().into(),
843            scope: crate::query::Scope::default(),
844            content: crate::content::AnalysisSet::NONE,
845        };
846        let delivery = Delivery::new(crate::CachePolicy::Off, None);
847        let (index, _) = crate::open(&basis, &delivery).expect("open a");
848        let handle = IndexHandle::new(index);
849        let before = handle.clock().expect("clock");
850        let request = Request::new(
851            Basis { root: b.path().into(), ..basis },
852            Query::default(),
853            std::time::SystemTime::now(),
854        );
855        assert!(matches!(
856            Session::new(handle.clone(), request, &delivery, WatchConfig::default()),
857            Err(Error::InvalidRequest(crate::query::RequestError::RootMismatch { .. }))
858        ));
859        assert_eq!(handle.clock().expect("clock"), before);
860    }
861
862    #[test]
863    fn a_save_is_due_only_when_a_change_is_pending_and_the_throttle_has_elapsed() {
864        let interval = Duration::from_secs(1);
865        let cases = [
866            // (pending, since last save, due, what this case is)
867            (true, Duration::from_secs(2), true, "pending and past the interval"),
868            (true, interval, true, "pending, exactly at the interval: inclusive"),
869            // The R5 case. Not due *now* -- and the flag stays set, which is the half that
870            // was missing: the idle path saves it once the interval passes.
871            (true, Duration::from_millis(1), false, "pending but throttled"),
872            (false, Duration::from_secs(60), false, "nothing pending, however long it has been"),
873            (false, Duration::ZERO, false, "nothing pending and just saved"),
874        ];
875
876        for (pending, since, want, case) in cases {
877            assert_eq!(save_is_due(pending, since, interval), want, "{case}");
878        }
879    }
880
881    /// A throttled change must survive every outcome except a completed write.
882    #[test]
883    fn only_a_completed_write_clears_the_pending_change() {
884        // The R7 case is Skipped and Failed: clearing the flag for either means the idle
885        // path never retries, so on a quiet tree the change is never persisted at all.
886        assert!(!pending_after(&SaveOutcome::Written), "a completed write persists the change");
887        assert!(
888            pending_after(&SaveOutcome::Skipped),
889            "a skipped save wrote nothing, so the change is still owed to disk",
890        );
891        assert!(
892            pending_after(&SaveOutcome::Failed(Error::Snapshot("failed".into()))),
893            "a failed save must be retried, not forgotten"
894        );
895    }
896
897    /// The sequence that defeated persistence in its most common shape.
898    #[test]
899    fn a_burst_then_a_quiet_tree_still_persists() {
900        let interval = Duration::from_secs(1);
901
902        // A change arrives too soon after the last save, so nothing is written yet.
903        let mut pending = true;
904        assert!(!save_is_due(pending, Duration::from_millis(50), interval));
905        assert!(pending, "the throttle must not consume the change");
906
907        // The tree goes quiet: no further batches will ever arrive. The idle path is the
908        // only remaining caller, and once the interval passes the save must happen.
909        assert!(save_is_due(pending, Duration::from_secs(3), interval));
910
911        // A skip at that point keeps it pending for the next idle tick rather than
912        // silently dropping the session's work.
913        pending = pending_after(&SaveOutcome::Skipped);
914        assert!(pending);
915        pending = pending_after(&SaveOutcome::Written);
916        assert!(!pending, "once written, the loop stops rewriting an unchanged index");
917    }
918
919    #[test]
920    fn skips_and_failures_retry_only_after_another_interval() {
921        let start = Instant::now();
922        let interval = Duration::from_secs(2);
923        let mut persistence = Persistence { pending: true, last_attempt: start };
924        assert!(matches!(
925            persistence.persist_due(start + interval / 2, interval, || panic!("throttled")),
926            SaveOutcome::Skipped
927        ));
928        assert!(matches!(
929            persistence.persist_due(start + interval, interval, || Ok(false)),
930            SaveOutcome::Skipped
931        ));
932        assert!(persistence.pending);
933        assert!(matches!(
934            persistence.persist_due(start + interval, interval, || panic!("skip was throttled")),
935            SaveOutcome::Skipped
936        ));
937        assert!(matches!(
938            persistence.persist_due(start + interval * 2, interval, || {
939                Err(Error::Snapshot("disk unavailable".into()))
940            }),
941            SaveOutcome::Failed(_)
942        ));
943        assert!(persistence.pending);
944        assert!(matches!(
945            persistence
946                .persist_due(start + interval * 2, interval, || panic!("failure was throttled")),
947            SaveOutcome::Skipped
948        ));
949        assert!(matches!(
950            persistence.persist_due(start + interval * 3, interval, || Ok(true)),
951            SaveOutcome::Written
952        ));
953        assert!(!persistence.pending);
954        assert!(matches!(
955            persistence.persist_due(start + interval * 4, interval, || panic!("already persisted")),
956            SaveOutcome::Skipped
957        ));
958    }
959
960    #[test]
961    fn handoff_changes_are_persisted_after_the_tree_goes_quiet() {
962        let root = tempfile::tempdir().expect("root");
963        let cache = tempfile::tempdir().expect("cache");
964        let cache_path = cache.path().join("snapshot");
965        let scan = ScanConfig::default();
966        std::fs::write(root.path().join("before.txt"), b"before").expect("before");
967        let (index, _) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
968        let request = Request::new(
969            Basis {
970                root: root.path().to_path_buf(),
971                scope: scan.clone().into(),
972                content: crate::content::AnalysisSet::NONE,
973            },
974            Query::default(),
975            std::time::SystemTime::now(),
976        );
977        let interval = Duration::from_secs(2);
978        let delivery = Delivery {
979            stale_ok: false,
980            cache: crate::CachePolicy::Auto,
981            cache_path: Some(cache_path.clone()),
982            accept_partial: false,
983            watch: Some(WatchDelivery { interval }),
984            workers: crate::query::Workers::default(),
985            batch_size: ScanConfig::default().batch_size,
986            order: crate::scan::ScanOrder::default(),
987        };
988        let script = tempfile::NamedTempFile::new().expect("script");
989        let (watcher, _sender) =
990            Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
991        std::fs::write(root.path().join("during-handoff.txt"), b"handoff").expect("handoff change");
992        let mut session = Session::finish_initial_handoff(
993            IndexHandle::new(index),
994            request,
995            &delivery,
996            watcher,
997            scan,
998            None,
999        )
1000        .expect("handoff");
1001        let started = session.persistence.last_attempt;
1002        assert!(session.persistence.pending, "handoff changes need persistence too");
1003        assert!(matches!(session.persist_due(started), SaveOutcome::Skipped));
1004        assert!(!cache_path.exists(), "throttle delays the write");
1005        assert!(matches!(session.persist_due(started + interval), SaveOutcome::Written));
1006        let restored = crate::snapshot::load(&cache_path).expect("load").expect("saved");
1007        assert!(matches!(
1008            restored.path_state(std::path::Path::new("during-handoff.txt")),
1009            crate::PathState::Present { .. }
1010        ));
1011        assert!(matches!(session.persist_due(started + interval * 2), SaveOutcome::Skipped));
1012    }
1013
1014    #[test]
1015    fn startup_save_failure_keeps_the_session_live_and_retries() {
1016        let root = tempfile::tempdir().expect("root");
1017        let cache = tempfile::tempdir().expect("cache");
1018        // A session reads before it writes, so the fixture must read as no snapshot yet
1019        // refuse the write: a parent that cannot become a directory. On Unix that is a
1020        // symbolic link to a directory that does not exist yet, since a path through a
1021        // regular file fails to open rather than reading as absent; Windows spells a
1022        // path through a regular file as not found, so a file serves there.
1023        let parent = cache.path().join("blocked");
1024        #[cfg(unix)]
1025        std::os::unix::fs::symlink(cache.path().join("missing"), &parent)
1026            .expect("dangling parent blocks the write");
1027        #[cfg(not(unix))]
1028        std::fs::write(&parent, b"").expect("file parent blocks the write");
1029        let restore = || {
1030            #[cfg(unix)]
1031            std::fs::create_dir(cache.path().join("missing")).expect("restore the parent");
1032            #[cfg(not(unix))]
1033            std::fs::remove_file(&parent).expect("restore the parent");
1034        };
1035        let cache_path = parent.join("snapshot");
1036        std::fs::write(root.path().join("file.txt"), b"content").expect("file");
1037        let interval = Duration::from_secs(2);
1038        let request = Request::new(
1039            Basis {
1040                root: root.path().to_path_buf(),
1041                scope: ScanConfig::default().into(),
1042                content: crate::content::AnalysisSet::NONE,
1043            },
1044            Query::default(),
1045            std::time::SystemTime::now(),
1046        );
1047        let delivery = Delivery {
1048            stale_ok: false,
1049            cache: crate::CachePolicy::On,
1050            cache_path: Some(cache_path.clone()),
1051            accept_partial: false,
1052            watch: Some(WatchDelivery { interval }),
1053            workers: crate::query::Workers::default(),
1054            batch_size: ScanConfig::default().batch_size,
1055            order: crate::scan::ScanOrder::default(),
1056        };
1057        let mut session = Session::start(request, delivery).expect("save failure is nonfatal");
1058        assert!(session.report(std::time::SystemTime::now()).expect("live report").status.complete);
1059        let now = Instant::now();
1060        assert!(matches!(session.persist_due(now), SaveOutcome::Failed(_)));
1061        restore();
1062        assert!(matches!(session.persist_due(now), SaveOutcome::Skipped));
1063        assert!(matches!(session.persist_due(now + interval), SaveOutcome::Written));
1064        assert!(crate::snapshot::load(&cache_path).expect("read snapshot").is_some());
1065    }
1066
1067    #[test]
1068    fn an_update_that_leaves_attribute_selection_emits_remove() {
1069        let root = tempfile::tempdir().expect("tempdir");
1070        std::fs::write(root.path().join("file.txt"), b"12345678").expect("fixture");
1071        let scan = ScanConfig::default();
1072        let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1073        assert!(report.is_complete());
1074        let request = Request::new(
1075            Basis {
1076                root: root.path().to_path_buf(),
1077                scope: scan.clone().into(),
1078                content: crate::content::AnalysisSet::NONE,
1079            },
1080            Query {
1081                selection: Selection { min_size: Some(4), ..Selection::default() },
1082                ..Query::default()
1083            },
1084            std::time::SystemTime::now(),
1085        );
1086        let delivery = Delivery {
1087            stale_ok: false,
1088            cache: crate::CachePolicy::Off,
1089            cache_path: None,
1090            accept_partial: false,
1091            watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1092            workers: crate::query::Workers::default(),
1093            batch_size: ScanConfig::default().batch_size,
1094            order: crate::scan::ScanOrder::default(),
1095        };
1096        let session =
1097            Session::new(IndexHandle::new(index), request, &delivery, WatchConfig::default())
1098                .expect("session");
1099        let path = PathBuf::from("file.txt");
1100        let change = session
1101            .change_for(
1102                &EffectiveChange::Updated {
1103                    path: path.clone(),
1104                    kind: EntryKind::File,
1105                    previous: crate::Attrs { size: 8, allocated: 8, ..crate::Attrs::default() },
1106                    current: crate::Attrs { size: 1, allocated: 1, ..crate::Attrs::default() },
1107                },
1108                1,
1109                &BatchFacts {
1110                    ignored: Some(BTreeMap::from([(path, false)])),
1111                    reclassified: BTreeMap::new(),
1112                },
1113            )
1114            .expect("membership transition");
1115        assert_eq!(change.kind, ChangeKind::Remove);
1116    }
1117
1118    #[test]
1119    fn initial_handoff_drains_a_sticky_overflow_after_a_full_intent_queue() {
1120        let root = tempfile::tempdir().expect("root");
1121        let script = tempfile::NamedTempFile::new().expect("script");
1122        let scan = ScanConfig::default();
1123        let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1124        assert!(report.is_complete());
1125        let handle = IndexHandle::new(index);
1126        let config = WatchConfig {
1127            settle: Duration::from_millis(1),
1128            max_hold: Duration::from_millis(2),
1129            event_capacity: 8,
1130            batch_path_capacity: 1,
1131            intent_capacity: 1,
1132            ..WatchConfig::default()
1133        };
1134        let (watcher, sender) =
1135            Watcher::scripted(root.path(), config, script.path()).expect("scripted watcher");
1136        std::fs::write(root.path().join("a.txt"), b"a").expect("a");
1137        sender.send("create\ta.txt\n").expect("first event");
1138        watcher.flush_capture().expect("first barrier fills the intent queue");
1139        std::fs::write(root.path().join("b.txt"), b"b").expect("b");
1140        sender.send("create\tb.txt\n").expect("second event");
1141        watcher.flush_capture().expect("second barrier retains sticky overflow");
1142
1143        drain_initial_capture(&watcher, &handle, &scan).expect("bounded handoff");
1144
1145        assert!(
1146            handle.snapshot().expect("snapshot").lookup(std::path::Path::new("a.txt")).is_some()
1147        );
1148        assert!(
1149            handle.snapshot().expect("snapshot").lookup(std::path::Path::new("b.txt")).is_some()
1150        );
1151        assert!(
1152            watcher
1153                .apply_next(&handle, &scan, Duration::ZERO, &mut |_| {})
1154                .expect("proof poll")
1155                .is_none(),
1156            "no queued or sticky pre-handoff work remains"
1157        );
1158    }
1159
1160    #[test]
1161    fn initial_handoff_enforces_partial_acceptance_after_its_reconciliation() {
1162        let root = tempfile::tempdir().expect("root");
1163        std::fs::write(root.path().join("kept.txt"), b"kept").expect("fixture");
1164        let scan = ScanConfig::default();
1165        let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1166        assert!(report.is_complete());
1167        let request = Request::new(
1168            Basis {
1169                root: root.path().to_path_buf(),
1170                scope: scan.clone().into(),
1171                content: crate::content::AnalysisSet::NONE,
1172            },
1173            Query::default(),
1174            std::time::SystemTime::now(),
1175        );
1176        let _fault = crate::scan::install_walk_hook(root.path(), |_| {
1177            Some(std::io::Error::new(
1178                std::io::ErrorKind::PermissionDenied,
1179                "deterministic handoff refusal",
1180            ))
1181        });
1182        let delivery = Delivery {
1183            stale_ok: false,
1184            cache: crate::CachePolicy::Off,
1185            cache_path: None,
1186            accept_partial: false,
1187            watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1188            workers: crate::query::Workers::default(),
1189            batch_size: ScanConfig::default().batch_size,
1190            order: crate::scan::ScanOrder::default(),
1191        };
1192
1193        let Err(error) = Session::new(
1194            IndexHandle::new(index.clone()),
1195            request.clone(),
1196            &delivery,
1197            WatchConfig::default(),
1198        ) else {
1199            panic!("a partial handoff is refused");
1200        };
1201        assert!(matches!(error, Error::ObservationHandoffIncomplete));
1202
1203        let accepted = Session::new(
1204            IndexHandle::new(index),
1205            request,
1206            &Delivery { accept_partial: true, ..delivery },
1207            WatchConfig::default(),
1208        )
1209        .expect("the caller explicitly accepts a partial handoff");
1210        assert!(
1211            !accepted.report(std::time::SystemTime::now()).expect("partial report").status.complete
1212        );
1213    }
1214
1215    #[test]
1216    fn initial_handoff_rechecks_partial_acceptance_after_draining_capture() {
1217        use std::sync::Arc;
1218        use std::sync::atomic::{AtomicUsize, Ordering};
1219
1220        let root = tempfile::tempdir().expect("root");
1221        std::fs::write(root.path().join("kept.txt"), b"kept").expect("fixture");
1222        let scan = ScanConfig::default();
1223        let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1224        assert!(report.is_complete());
1225        let request = Request::new(
1226            Basis {
1227                root: root.path().to_path_buf(),
1228                scope: scan.clone().into(),
1229                content: crate::content::AnalysisSet::NONE,
1230            },
1231            Query::default(),
1232            std::time::SystemTime::now(),
1233        );
1234        let delivery = Delivery {
1235            stale_ok: false,
1236            cache: crate::CachePolicy::Off,
1237            cache_path: None,
1238            accept_partial: false,
1239            watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1240            workers: crate::query::Workers::default(),
1241            batch_size: ScanConfig::default().batch_size,
1242            order: crate::scan::ScanOrder::default(),
1243        };
1244        let script = tempfile::NamedTempFile::new().expect("script");
1245        let (watcher, sender) =
1246            Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1247        sender.send("rescan\t.\n").expect("queue initial gap");
1248
1249        // The startup reconciliation succeeds. The scripted overflow then reaches the same
1250        // tree during the handoff drain, where its reconciliation fails and must be refused.
1251        let attempts = Arc::new(AtomicUsize::new(0));
1252        let hook_attempts = Arc::clone(&attempts);
1253        let _fault = crate::scan::install_walk_hook(root.path(), move |_| {
1254            (hook_attempts.fetch_add(1, Ordering::SeqCst) > 0).then(|| {
1255                std::io::Error::new(
1256                    std::io::ErrorKind::PermissionDenied,
1257                    "deterministic drain-only refusal",
1258                )
1259            })
1260        });
1261
1262        let Err(error) = Session::finish_initial_handoff(
1263            IndexHandle::new(index),
1264            request,
1265            &delivery,
1266            watcher,
1267            scan,
1268            None,
1269        ) else {
1270            panic!("a partial state created while draining is refused");
1271        };
1272        assert!(matches!(error, Error::ObservationHandoffIncomplete));
1273        assert!(attempts.load(Ordering::SeqCst) > 1, "the drain ran after startup reconciliation");
1274    }
1275
1276    /// The handoff pass restarts the walk counters and enters `Revalidating`, whatever
1277    /// the first pass left in the handle: with a scripted watcher that reports nothing,
1278    /// the counts afterwards are exactly one walk of the tree. Identical event streams
1279    /// also make the observed and unobserved reports exactly comparable.
1280    #[test]
1281    fn the_handoff_pass_restarts_the_counts() {
1282        let root = tempfile::tempdir().expect("root");
1283        let mut bytes = 0;
1284        for directory in 0..2 {
1285            let dir = root.path().join(format!("d{directory}"));
1286            std::fs::create_dir(&dir).expect("directory");
1287            for file in 0..3 {
1288                let size = directory * 3 + file + 1;
1289                std::fs::write(dir.join(format!("f{file}.txt")), vec![b'.'; size]).expect("file");
1290                bytes += size as u64;
1291            }
1292        }
1293        let scan = ScanConfig::default();
1294        let (index, report) = crate::scan::scan_into_index(root.path(), &scan).expect("scan");
1295        assert!(report.is_complete());
1296        let request = Request::new(
1297            Basis {
1298                root: root.path().to_path_buf(),
1299                scope: scan.clone().into(),
1300                content: crate::content::AnalysisSet::NONE,
1301            },
1302            Query { views: vec![crate::query::ViewSpec::Summary], ..Query::default() },
1303            std::time::SystemTime::now(),
1304        );
1305        let delivery = Delivery {
1306            stale_ok: false,
1307            cache: crate::CachePolicy::Off,
1308            cache_path: None,
1309            accept_partial: false,
1310            watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1311            workers: crate::query::Workers::default(),
1312            batch_size: ScanConfig::default().batch_size,
1313            order: crate::scan::ScanOrder::default(),
1314        };
1315        let script = tempfile::NamedTempFile::new().expect("script");
1316        let (watcher, _sender) =
1317            Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1318        let progress = crate::Progress::new();
1319        progress.add_walked(100, 100, 100, 100);
1320        progress.enter(crate::ProgressPhase::Saving);
1321
1322        let session = Session::finish_initial_handoff(
1323            IndexHandle::new(index.clone()),
1324            request.clone(),
1325            &delivery,
1326            watcher,
1327            scan.clone(),
1328            Some(&progress),
1329        )
1330        .expect("handoff");
1331
1332        let snapshot = progress.snapshot();
1333        assert_eq!(snapshot.phase, crate::ProgressPhase::Revalidating);
1334        assert_eq!(
1335            (snapshot.directories, snapshot.files, snapshot.bytes),
1336            (3, 6, bytes),
1337            "one walk of the root and its two directories, the first pass not added in"
1338        );
1339        assert_eq!(snapshot.allocated, report.allocated_walked, "allocated restarts with them");
1340
1341        let (plain_watcher, _plain_sender) =
1342            Watcher::scripted(root.path(), WatchConfig::default(), script.path()).expect("watcher");
1343        let plain = Session::finish_initial_handoff(
1344            IndexHandle::new(index),
1345            request,
1346            &delivery,
1347            plain_watcher,
1348            scan,
1349            None,
1350        )
1351        .expect("plain handoff");
1352        let generated_at = std::time::SystemTime::now();
1353        let observed_report = session.report(generated_at).expect("observed report");
1354        let mut plain_report = plain.report(generated_at).expect("plain report");
1355        plain_report.provenance = observed_report.provenance.clone();
1356        let json = |report: &Report| {
1357            crate::report_format::render(report, crate::report_format::Format::Json, false)
1358                .expect("render")
1359        };
1360        assert_eq!(json(&plain_report), json(&observed_report));
1361        assert_eq!(progress.snapshot(), snapshot, "the second handoff and reads were not observed");
1362    }
1363
1364    /// A record says what the index can be asked, and nothing more.
1365    ///
1366    /// The three answers are distinct and a consumer acts on each differently: a bit, "no
1367    /// rules were read", and "there is no such entry". A set of the ignored paths collapsed
1368    /// the last two into `false`, so a batch that created and removed one file in the same
1369    /// window upserted it as unignored before removing it -- a classification claim about
1370    /// an entry that never survived the batch.
1371    #[test]
1372    fn a_record_claims_a_classification_only_for_an_entry_the_index_still_holds() {
1373        let unobserved = BatchFacts { ignored: None, reclassified: BTreeMap::new() };
1374        assert_eq!(unobserved.is_ignored(std::path::Path::new("any.txt")), None);
1375
1376        let observed = BatchFacts {
1377            ignored: Some(BTreeMap::from([
1378                (PathBuf::from("build/out.bin"), true),
1379                (PathBuf::from("src/main.rs"), false),
1380            ])),
1381            reclassified: BTreeMap::new(),
1382        };
1383        assert_eq!(observed.is_ignored(std::path::Path::new("build/out.bin")), Some(true));
1384        assert_eq!(observed.is_ignored(std::path::Path::new("src/main.rs")), Some(false));
1385        assert_eq!(
1386            observed.is_ignored(std::path::Path::new("gone.tmp")),
1387            None,
1388            "an entry the batch removed is in neither partition, not in the unignored one"
1389        );
1390    }
1391
1392    /// A start is the open and then the handoff revalidation, a second pass whose
1393    /// counts restart, which keeps the line moving on a large tree after the save,
1394    /// where a frozen count would look like a hang. The session it returns is the one
1395    /// [`Session::start`] returns, and it reports nothing further through the handle
1396    /// once started. The exact reset and report equivalence are pinned by the scripted
1397    /// test above; independent native watchers may replay different creation hints and
1398    /// record different legitimate setup-gap diagnostics, so only lower bounds hold.
1399    #[test]
1400    fn a_started_session_reports_its_second_pass_and_then_nothing() {
1401        let root = tempfile::tempdir().expect("root");
1402        let cache = tempfile::tempdir().expect("cache");
1403        let mut bytes = 0;
1404        for directory in 0..4 {
1405            let dir = root.path().join(format!("d{directory}"));
1406            std::fs::create_dir(&dir).expect("directory");
1407            for file in 0..3 {
1408                let size = directory * 3 + file + 1;
1409                std::fs::write(dir.join(format!("f{file}.txt")), vec![b'.'; size]).expect("file");
1410                bytes += size as u64;
1411            }
1412        }
1413        let request = || {
1414            Request::new(
1415                Basis {
1416                    root: root.path().to_path_buf(),
1417                    scope: ScanConfig::default().into(),
1418                    content: crate::content::AnalysisSet::NONE,
1419                },
1420                Query::default(),
1421                std::time::UNIX_EPOCH,
1422            )
1423        };
1424        let delivery = Delivery {
1425            stale_ok: false,
1426            cache: crate::CachePolicy::Auto,
1427            cache_path: Some(cache.path().join("snapshot")),
1428            accept_partial: false,
1429            watch: Some(WatchDelivery { interval: Duration::from_secs(2) }),
1430            workers: crate::query::Workers::default(),
1431            batch_size: 4,
1432            order: crate::scan::ScanOrder::default(),
1433        };
1434
1435        let progress = crate::Progress::new();
1436        let session = Session::start_with_progress(request(), delivery.clone(), &progress)
1437            .expect("observed start");
1438        let after_start = progress.snapshot();
1439        assert_eq!(
1440            after_start.phase,
1441            crate::ProgressPhase::Revalidating,
1442            "the handoff revalidation follows the joined save"
1443        );
1444        // The closing pass restarts the counters, so they show its walk alone. A backend
1445        // that reports a file created just before the watch began can make that pass
1446        // read more, never less.
1447        assert!(after_start.directories >= 5, "the root and four children: {after_start:?}");
1448        assert!(after_start.files >= 12, "{after_start:?}");
1449        assert!(after_start.bytes >= bytes, "{after_start:?}");
1450        assert_eq!(after_start.analysis, None);
1451
1452        let plain = Session::start(request(), delivery).expect("plain start");
1453        let generated_at = std::time::SystemTime::now();
1454        let observed_report = session.report(generated_at).expect("observed report");
1455        let mut plain_report = plain.report(generated_at).expect("plain report");
1456        // A native watcher registered a moment after a directory was created may still
1457        // report it, and the engine records that truthfully as a setup-race observation
1458        // gap beside the facts it re-verified (fdu-21ns). That diagnostic is the only
1459        // way the watched answer may differ from the plain one; anything else is a real
1460        // difference, and the facts both report are compared with it set aside.
1461        //
1462        // An incomplete answer is allowed only when it retains at least one issue, every
1463        // one a setup-race gap, and none omitted: an incomplete answer with no issue, or
1464        // with any other, is a regression. The facts it re-verified are as fresh as the
1465        // plain answer's, and with no issue its coverage is the plain answer's too; both
1466        // are compared before the status and provenance are set aside.
1467        let setup_race = format!("{:?}", crate::InvalidateReason::WatchSetupRace);
1468        let errors = &observed_report.status.errors;
1469        let setup_race_only = !errors.is_empty()
1470            && errors.iter().all(|issue| {
1471                issue.kind == crate::IssueKind::ObservationGap
1472                    && issue.message.ends_with(&setup_race)
1473            });
1474        assert!(
1475            observed_report.status.complete || setup_race_only,
1476            "the watched answer is complete or carries only setup-race gaps: {:?}",
1477            observed_report.status
1478        );
1479        assert_eq!(observed_report.status.errors_omitted, 0, "{:?}", observed_report.status);
1480        assert!(plain_report.status.complete, "{:?}", plain_report.status);
1481        assert_eq!(
1482            observed_report.provenance.freshness, plain_report.provenance.freshness,
1483            "the watched answer is as fresh as the plain one"
1484        );
1485        if errors.is_empty() {
1486            assert_eq!(
1487                observed_report.status.coverage, plain_report.status.coverage,
1488                "with no retained issue, the coverage is the plain answer's"
1489            );
1490        }
1491        plain_report.provenance = observed_report.provenance.clone();
1492        plain_report.status = observed_report.status.clone();
1493        let json = |report: &Report| {
1494            crate::report_format::render(report, crate::report_format::Format::Json, false)
1495                .expect("render")
1496        };
1497        assert_eq!(json(&plain_report), json(&observed_report), "the same facts either way");
1498        assert_eq!(progress.snapshot(), after_start, "the second start was not observed");
1499    }
1500
1501    /// The first answer's build is the last thing a start's progress covers (fdu-wku3):
1502    /// asked for through the handle, it enters `Summarizing`, answers as the plain call
1503    /// does, and is the answer later repaints are measured against.
1504    #[test]
1505    fn the_first_answer_is_built_under_the_summarizing_phase() {
1506        let root = tempfile::tempdir().expect("root");
1507        std::fs::write(root.path().join("a.txt"), b"alpha").expect("file");
1508        std::fs::create_dir(root.path().join("d")).expect("directory");
1509        std::fs::write(root.path().join("d/b.txt"), b"beta").expect("nested file");
1510        let (mut session, _sender) = scripted_session(root.path(), Query::default());
1511        let format = crate::report_format::Format::Json;
1512        let options = crate::report_format::RenderOptions { color: false, bar_size: 0 };
1513        let generated_at = std::time::SystemTime::now();
1514        let plain_answer = session.report(generated_at).expect("plain answer");
1515
1516        let progress = crate::Progress::new();
1517        progress.enter(crate::ProgressPhase::Revalidating);
1518        let observed = session
1519            .changed_report_with_progress(generated_at, format, options, &progress)
1520            .expect("observed first answer")
1521            .expect("a session's first answer is always given");
1522        assert_eq!(
1523            progress.snapshot().phase,
1524            crate::ProgressPhase::Summarizing,
1525            "the build is reported as the answer's construction"
1526        );
1527
1528        let json =
1529            |report: &Report| crate::report_format::render(report, format, false).expect("render");
1530        assert_eq!(json(&observed), json(&plain_answer), "the same answer either way");
1531        assert!(
1532            session.changed_report(generated_at, format, options).expect("repaint").is_none(),
1533            "nothing changed since the observed build, so nothing repaints"
1534        );
1535    }
1536
1537    /// A scripted session over `root` answering `query`, and the sender that scripts its
1538    /// events.
1539    fn scripted_session(
1540        root: &std::path::Path,
1541        query: Query,
1542    ) -> (Session, crate::watch::ScriptedSender) {
1543        scripted_session_under(root, query, ScanConfig::default())
1544    }
1545
1546    /// [`scripted_session`] under `scan`, such as one with tighter control limits.
1547    fn scripted_session_under(
1548        root: &std::path::Path,
1549        query: Query,
1550        scan: ScanConfig,
1551    ) -> (Session, crate::watch::ScriptedSender) {
1552        let (index, report) = crate::scan::scan_into_index(root, &scan).expect("scan");
1553        assert!(report.is_complete());
1554        let request = Request::new(
1555            Basis {
1556                root: root.to_path_buf(),
1557                scope: scan.clone().into(),
1558                content: crate::content::AnalysisSet::NONE,
1559            },
1560            query,
1561            std::time::SystemTime::now(),
1562        );
1563        let delivery = Delivery {
1564            stale_ok: false,
1565            cache: crate::CachePolicy::Off,
1566            cache_path: None,
1567            accept_partial: false,
1568            watch: Some(WatchDelivery { interval: Duration::from_millis(50) }),
1569            workers: crate::query::Workers::default(),
1570            batch_size: ScanConfig::default().batch_size,
1571            order: crate::scan::ScanOrder::default(),
1572        };
1573        let script = tempfile::NamedTempFile::new().expect("script");
1574        let (watcher, sender) =
1575            Watcher::scripted(root, WatchConfig::default(), script.path()).expect("watcher");
1576        let session = Session::finish_initial_handoff(
1577            IndexHandle::new(index),
1578            request,
1579            &delivery,
1580            watcher,
1581            scan,
1582            None,
1583        )
1584        .expect("handoff");
1585        (session, sender)
1586    }
1587
1588    /// Move a file's modification time forward without touching its bytes.
1589    fn touch(path: &std::path::Path) {
1590        let file = std::fs::File::options().write(true).open(path).expect("open for touch");
1591        let modified = file.metadata().expect("metadata").modified().expect("mtime");
1592        file.set_modified(modified + Duration::from_secs(5)).expect("set mtime");
1593    }
1594
1595    /// A tree changes more often than its answer (fdu-wb5n): an idle tree, a touch that
1596    /// leaves a size alone, and a change to an entry the selection leaves out all leave
1597    /// the aggregate a reader sees as it was, and only a visible change repaints it.
1598    #[test]
1599    fn an_aggregate_repaint_is_skipped_when_a_reader_would_see_no_change() {
1600        use crate::report_format::{Format, RenderOptions};
1601
1602        let root = tempfile::tempdir().expect("root");
1603        let big = root.path().join("big.txt");
1604        std::fs::write(&big, vec![b'x'; 4096]).expect("big");
1605        std::fs::write(root.path().join("small.txt"), b"small").expect("small");
1606        let query = Query {
1607            views: vec![crate::query::ViewSpec::Summary],
1608            selection: Selection {
1609                size: crate::query::SizeMetric::Apparent,
1610                min_size: Some(4096),
1611                ..Selection::default()
1612            },
1613            ..Query::default()
1614        };
1615        let (mut session, sender) = scripted_session(root.path(), query);
1616        let now = std::time::SystemTime::now;
1617        let text = |session: &mut Session| {
1618            session.changed_report(now(), Format::Text, RenderOptions::default()).expect("report")
1619        };
1620        let summary = |report: &Report| match report.sections.first() {
1621            Some(crate::query::Section::Summary(row)) => (row.files, row.bytes),
1622            other => panic!("expected a summary, got {other:?}"),
1623        };
1624
1625        let first = text(&mut session).expect("the first answer is always given");
1626        assert_eq!(summary(&first), (1, 4096));
1627
1628        // Idle: nothing arrived, nothing to say.
1629        assert!(session.next_batch(Duration::from_millis(50)).expect("idle").is_none());
1630        assert!(text(&mut session).is_none(), "an idle tree repaints nothing");
1631
1632        // A touch moves the index, so its batch is dirty, but no size a reader of the
1633        // summary sees has moved.
1634        touch(&big);
1635        sender.send("modify\tbig.txt\n").expect("script a touch");
1636        let batch = session.next_batch(Duration::from_secs(10)).expect("touch").expect("observed");
1637        assert!(batch.dirty, "the index records the new modification time");
1638        assert!(text(&mut session).is_none(), "a touch that changes no size repaints nothing");
1639
1640        // A change to an entry the selection leaves out is a change to the tree, and
1641        // still not a change to the answer.
1642        std::fs::write(root.path().join("small.txt"), b"still small").expect("grow small");
1643        sender.send("modify\tsmall.txt\n").expect("script a filtered change");
1644        let batch = session.next_batch(Duration::from_secs(10)).expect("small").expect("observed");
1645        assert!(batch.dirty);
1646        assert!(text(&mut session).is_none(), "a filtered change repaints nothing");
1647
1648        // A size the summary shows moves: repaint, with the new answer.
1649        std::fs::write(&big, vec![b'x'; 8192]).expect("grow big");
1650        sender.send("modify\tbig.txt\n").expect("script a visible change");
1651        let batch = session.next_batch(Duration::from_secs(10)).expect("big").expect("observed");
1652        assert!(batch.dirty);
1653        let repainted = text(&mut session).expect("a visible change repaints");
1654        assert_eq!(summary(&repainted), (1, 8192));
1655        assert!(text(&mut session).is_none(), "and only once");
1656
1657        // The plain report always answers, and never counts as a repaint.
1658        assert_eq!(summary(&session.report(now()).expect("report")), (1, 8192));
1659        assert!(text(&mut session).is_none());
1660    }
1661
1662    /// The identity follows the format a reader sees: machine output carries the newest
1663    /// modification time, so the touch that text ignores repaints JSON.
1664    #[test]
1665    fn a_repaint_identity_is_what_the_format_renders() {
1666        use crate::report_format::{Format, RenderOptions};
1667
1668        let root = tempfile::tempdir().expect("root");
1669        let big = root.path().join("big.txt");
1670        std::fs::write(&big, vec![b'x'; 4096]).expect("big");
1671        let query = Query { views: vec![crate::query::ViewSpec::Summary], ..Query::default() };
1672        let (mut session, sender) = scripted_session(root.path(), query);
1673        let json = |session: &mut Session| {
1674            session
1675                .changed_report(
1676                    std::time::SystemTime::now(),
1677                    Format::Json,
1678                    RenderOptions::default(),
1679                )
1680                .expect("report")
1681        };
1682        assert!(json(&mut session).is_some());
1683        assert!(json(&mut session).is_none(), "a second generation instant alone is no change");
1684
1685        touch(&big);
1686        sender.send("modify\tbig.txt\n").expect("script a touch");
1687        assert!(session.next_batch(Duration::from_secs(10)).expect("touch").is_some());
1688        assert!(json(&mut session).is_some(), "JSON shows the newest modification time");
1689    }
1690
1691    /// The repaint digest is FNV-1a at 128 bits, and a section boundary is part of what
1692    /// it digests.
1693    #[test]
1694    fn a_repaint_digest_is_fnv1a_128_framed_by_section() {
1695        use std::io::Write as _;
1696
1697        let mut raw = RepaintDigest::new();
1698        raw.mix(b"a");
1699        assert_eq!(raw.hash, 0xd228_cb69_6f1a_8caf_7891_2b70_4e4a_8964, "the published vector");
1700        let digest = |sections: &[&[u8]]| {
1701            let mut digest = RepaintDigest::new();
1702            for (at, section) in sections.iter().enumerate() {
1703                if at > 0 {
1704                    digest.end_section();
1705                }
1706                digest.write_all(section).expect("digest");
1707            }
1708            digest.finish()
1709        };
1710        assert_eq!(digest(&[b"ab", b"c"]), digest(&[b"ab", b"c"]));
1711        assert_ne!(digest(&[b"ab", b"c"]), digest(&[b"a", b"bc"]));
1712        assert_ne!(digest(&[b"abc", b""]), digest(&[b"", b"abc"]));
1713    }
1714
1715    /// A reader of a text report reads its diagnostics too (R164-4): a `.gitignore` edited
1716    /// past the line limit adds a refusal note while the tree, its sizes, and its status
1717    /// stay as they were, and that note alone repaints.
1718    #[test]
1719    fn a_diagnostic_that_appears_on_an_unchanged_tree_repaints() {
1720        use crate::report_format::{Format, RenderOptions};
1721
1722        let root = tempfile::tempdir().expect("root");
1723        std::fs::write(root.path().join("big.txt"), vec![b'x'; 4096]).expect("big");
1724        let control = root.path().join(".gitignore");
1725        // Ten bytes either way, so no size a reader sees moves; only the longest line does.
1726        std::fs::write(&control, b"a\nb\nc\nd\ne\n").expect("short lines");
1727        let scan = ScanConfig {
1728            control_limits: crate::control::ControlLimits {
1729                line_limit: Some(4),
1730                ..crate::control::ControlLimits::default()
1731            },
1732            ..ScanConfig::default()
1733        };
1734        let (mut session, sender) = scripted_session_under(root.path(), Query::default(), scan);
1735        let text = |session: &mut Session| {
1736            session
1737                .changed_report(
1738                    std::time::SystemTime::now(),
1739                    Format::Text,
1740                    RenderOptions::default(),
1741                )
1742                .expect("report")
1743        };
1744        let first = text(&mut session).expect("the first answer is always given");
1745        let notes = |report: &Report| crate::report_format::diagnostic_lines(report).into_lines();
1746        assert!(
1747            !notes(&first).iter().any(|line| line.contains("line limit")),
1748            "{:?}",
1749            notes(&first)
1750        );
1751
1752        std::fs::write(&control, b"abcdefghi\n").expect("one long line");
1753        sender.send("modify\t.gitignore\n").expect("script the edit");
1754        let batch = session.next_batch(Duration::from_secs(10)).expect("edit").expect("observed");
1755        assert!(batch.dirty);
1756        let repainted = text(&mut session).expect("a new refusal note repaints");
1757        assert_eq!(
1758            format!("{:?}", repainted.status),
1759            format!("{:?}", first.status),
1760            "the status alone would not have repainted"
1761        );
1762        assert!(
1763            notes(&repainted).iter().any(|line| line.contains("line limit")),
1764            "{:?}",
1765            notes(&repainted)
1766        );
1767        assert!(text(&mut session).is_none(), "and only once");
1768    }
1769
1770    /// An invalidation is never deduplicated as a change record, and the answer it leaves
1771    /// repaints exactly when a reader would see its status or its rows differ.
1772    #[test]
1773    fn an_invalidation_keeps_its_change_record_and_repaints_by_the_answer() {
1774        use crate::report_format::{Format, RenderOptions};
1775
1776        let root = tempfile::tempdir().expect("root");
1777        std::fs::write(root.path().join("a.txt"), b"aaaa").expect("a");
1778        let query = Query { views: vec![crate::query::ViewSpec::Summary], ..Query::default() };
1779        let (mut session, sender) = scripted_session(root.path(), query);
1780        let text = |session: &mut Session| {
1781            session
1782                .changed_report(
1783                    std::time::SystemTime::now(),
1784                    Format::Text,
1785                    RenderOptions::default(),
1786                )
1787                .expect("report")
1788        };
1789        let first = text(&mut session).expect("first answer");
1790
1791        sender.send("rescan\t.\n").expect("script an invalidation");
1792        let batch =
1793            session.next_batch(Duration::from_secs(10)).expect("invalidation").expect("observed");
1794        assert!(
1795            batch.changes.iter().any(|change| change.kind == ChangeKind::Invalidate),
1796            "the invalidation reaches the change stream: {:?}",
1797            batch.changes
1798        );
1799        // The closed loop re-verified the tree, which is as it was; whether the answer
1800        // repaints is decided by its status, which the identity carries explicitly.
1801        let after = session.report(std::time::SystemTime::now()).expect("report");
1802        let repainted = text(&mut session);
1803        assert_eq!(
1804            repainted.is_some(),
1805            format!("{:?}", after.status) != format!("{:?}", first.status),
1806            "a repaint follows a status change and nothing else: {:?}",
1807            after.status
1808        );
1809
1810        // The same invalidation again leaves the same status, so no second repaint.
1811        sender.send("rescan\t.\n").expect("script another invalidation");
1812        let batch =
1813            session.next_batch(Duration::from_secs(10)).expect("invalidation").expect("observed");
1814        assert!(batch.changes.iter().any(|change| change.kind == ChangeKind::Invalidate));
1815        assert!(text(&mut session).is_none());
1816    }
1817}