Skip to main content

fdu_core/
watch.rs

1//! Th OS-native watch layer: turning an unreliable event stream into verified observations.
2//!
3//! This module's whole job is that conversion. Filesystem events are **hints, not
4//! truth**, and the ways they lie are documented per platform:
5//!
6//! - Most events carry no metadata at all, so a producer must stat before it can say
7//!   what an entry now looks like.
8//! - Only inotify pairs the two sides of a rename (via a kernel cookie). `FSEvents` emits
9//!   one path with no mechanism to associate old and new; Windows delivers both sides
10//!   with no cookie; poll-based watching cannot see renames at all. This layer never
11//!   needs the pairing: each named side is verified as its own path, like a create or a
12//!   remove of that name (see `RenameReporting` for what each backend promises).
13//! - When a directory is created, backends that watch per directory register the new
14//!   watch *after* the fact — anything created inside that window produces no event.
15//! - Kernel queues overflow. inotify's `Q_OVERFLOW` and `FSEvents`' `MustScanSubDirs`
16//!   mean "your view is now incomplete" and surface here as `Flag::Rescan`. A Windows
17//!   buffer overrun does not: notify 8.2.0 logs `ERROR_NOTIFY_ENUM_DIR` and drops that
18//!   directory's watch without signaling it. Dropping the rescan signal is precisely how an event-driven index
19//!   silently diverges from the filesystem — which is what the `watchfiles` layer
20//!   metabrowser runs today does, mapping notify's rich event model down to
21//!   `(change, path)` and letting the rescan flag fall through a match arm.
22//!
23//! So this layer never forwards an event. It coalesces, then **verifies by stat**, and
24//! emits only [`Op::Upsert`] with a fresh fingerprint, [`Op::Remove`], or —
25//! when it genuinely cannot describe the change — [`Op::InvalidateSubtree`], which the
26//! scan layer resolves back into precise committed changes.
27//!
28//! Building on notify rather than on raw platform APIs is deliberate: its six backends
29//! and its overflow signaling are proven, and the information loss that motivates this
30//! module all happens in layers *above* it.
31
32use std::collections::BTreeMap;
33use std::panic::{AssertUnwindSafe, catch_unwind};
34use std::path::{Path, PathBuf};
35use std::sync::Arc;
36use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
37use std::sync::mpsc::{Receiver, RecvTimeoutError, SyncSender, TrySendError, sync_channel};
38use std::thread::JoinHandle;
39use std::time::{Duration, Instant};
40
41use notify::{EventKind, RecommendedWatcher, RecursiveMode, Watcher as NotifyWatcher};
42
43#[cfg(test)]
44mod scripted_events;
45
46use crate::engine_contract::{
47    Commit, Error, InvalidateReason, Observation, ObservationOp, Op, Result,
48};
49use crate::scan;
50use crate::{ApplyOutcome, IndexHandle, ScanConfig};
51
52/// Optimistic re-verification retries before the applying driver guarantees progress
53/// through conservative root invalidation and reconciliation.
54const MAX_OPTIMISTIC_APPLY_ATTEMPTS: usize = 3;
55
56const WORKER_RUNNING: u8 = 0;
57const WORKER_STOPPED: u8 = 1;
58const WORKER_PANICKED: u8 = 2;
59const MAX_EVENT_CAPACITY: usize = 64 * 1024;
60const MAX_BATCH_PATH_CAPACITY: usize = 64 * 1024;
61const MAX_BUFFERED_INTENT_PATHS: usize = 1024 * 1024;
62
63/// Tuning for event coalescing.
64#[derive(Clone, Copy, Debug)]
65pub struct WatchConfig {
66    /// How long the event stream must be quiet before a batch is emitted.
67    pub settle: Duration,
68    /// Longest a batch may be held open while events keep arriving. Without a ceiling, a
69    /// continuously busy tree would never produce a delta at all.
70    pub max_hold: Duration,
71    /// Whether a newly created directory triggers a re-list of its contents.
72    ///
73    /// On by default, and it should stay on for inotify and kqueue: those backends
74    /// install a directory's watch only after the create event arrives, so files created
75    /// in between are never reported by anything.
76    pub relist_new_dirs: bool,
77    /// Maximum raw backend events queued before overload collapses to one root
78    /// invalidation. Backend callback threads never block on this queue.
79    pub event_capacity: usize,
80    /// Maximum distinct paths retained in one coalesced intent.
81    pub batch_path_capacity: usize,
82    /// Maximum coalesced intents awaiting a consumer.
83    pub intent_capacity: usize,
84}
85
86impl Default for WatchConfig {
87    fn default() -> Self {
88        Self {
89            // The 50 ms step / 1.6 s ceiling pairing is watchfiles' batching loop, which
90            // is the part of that stack worth keeping.
91            settle: Duration::from_millis(50),
92            max_hold: Duration::from_millis(1600),
93            relist_new_dirs: true,
94            event_capacity: 4096,
95            batch_path_capacity: 4096,
96            intent_capacity: 16,
97        }
98    }
99}
100
101impl WatchConfig {
102    fn validate(self) -> Result<()> {
103        let buffered_paths = self.batch_path_capacity.checked_mul(self.intent_capacity);
104        if self.settle.is_zero()
105            || self.max_hold.is_zero()
106            || self.max_hold < self.settle
107            || self.event_capacity == 0
108            || self.batch_path_capacity == 0
109            || self.intent_capacity == 0
110            || self.event_capacity > MAX_EVENT_CAPACITY
111            || self.batch_path_capacity > MAX_BATCH_PATH_CAPACITY
112            || buffered_paths.is_none_or(|paths| paths > MAX_BUFFERED_INTENT_PATHS)
113        {
114            return Err(Error::UnsupportedScanConfig(
115                "watch durations and capacities exceed the supported nonzero bounds, or max_hold is less than settle",
116            ));
117        }
118        Ok(())
119    }
120}
121
122/// What a coalesced path still needs before it can become an observation.
123#[derive(Clone, Copy, PartialEq, Eq, Debug)]
124enum Pending {
125    /// Stat it and decide. Covers creates, writes, removes, and every rename shape:
126    /// letting the stat decide is what makes the same code correct on all backends.
127    Verify {
128        /// Preserve whether a create event occurred while this path was coalesced. Only
129        /// a newly created directory has the watch-registration race that needs a relist.
130        relist_if_dir: bool,
131        /// Preserve whether a rename named this path while it was coalesced.
132        ///
133        /// A rename is a removal at one name and an arrival at another, so each side is
134        /// verified exactly like a remove or a create of its own path. Two things differ
135        /// from a create: a directory that arrives by rename brings contents that no
136        /// backend reports, so it is always relisted; and a rename is how a name changes
137        /// only in case or Unicode normalization, which a lookup on an insensitive
138        /// filesystem cannot tell apart, so the name is checked against its parent's
139        /// listing before it is trusted (see [`verify_intent`]).
140        renamed: bool,
141    },
142    /// The producer already knows it cannot describe this precisely.
143    Escalate(InvalidateReason),
144}
145
146/// What a backend's rename events promise about the rename's other side.
147///
148/// Scoping a rename to the paths it names is sound only when every side of it that lies
149/// inside the watched tree is named by some event. That is a property of the backend,
150/// established from notify 8.2's sources rather than assumed:
151///
152/// - `FSEvents` reports each side as its own `ItemRenamed` record naming that side's
153///   path (notify: one `Modify(Name(Any))` per record), and delivers every record under
154///   the watched root; loss is signalled by `MustScanSubDirs`, which arrives here as
155///   `Flag::Rescan`.
156/// - inotify reports `IN_MOVED_FROM` and `IN_MOVED_TO` for each side inside a watched
157///   directory (notify: `From`, `To`, and a cookie-paired `Both`); queue loss is
158///   `IN_Q_OVERFLOW`, again `Flag::Rescan`.
159/// - `ReadDirectoryChangesW` reports the old and new names inside the tree (notify:
160///   `From` and `To`, never paired), and a move across the tree's boundary as a plain
161///   removal or addition. notify 8.2.0 does not signal its buffer overflow.
162/// - kqueue reports only the renamed vnode's old path; the new name surfaces at most as
163///   a write on its new parent directory, which this layer cannot turn into the entry.
164#[derive(Clone, Copy, PartialEq, Eq, Debug)]
165enum RenameReporting {
166    /// Every in-root side of a rename arrives as an event naming it.
167    EachSide,
168    /// A rename may name only one side; the other can be anywhere in the tree.
169    OldSideOnly,
170}
171
172impl RenameReporting {
173    /// The promise of the backend this build's [`RecommendedWatcher`] uses.
174    ///
175    /// Asked at run time rather than decided by `cfg`, because a dependent crate can turn
176    /// on notify's `macos_kqueue` build feature and change the macOS backend under us.
177    fn of_recommended_backend() -> Self {
178        match <RecommendedWatcher as NotifyWatcher>::kind() {
179            notify::WatcherKind::Fsevent
180            | notify::WatcherKind::Inotify
181            | notify::WatcherKind::ReadDirectoryChangesWatcher => Self::EachSide,
182            _ => Self::OldSideOnly,
183        }
184    }
185}
186
187#[derive(Debug, Default)]
188struct CoalescedIntent {
189    pending: BTreeMap<PathBuf, Pending>,
190}
191
192enum RawMessage {
193    Event(notify::Result<notify::Event>),
194    Flush(SyncSender<()>),
195    Stop,
196}
197
198/// A live watcher over one tree.
199///
200/// The watcher is movable to one consuming thread. It is intentionally not shareable:
201/// its private standard-library receiver enforces one ordered consumer for coalesced
202/// intents. Use a separate [`IndexHandle`] to serve concurrent readers and writers.
203/// Dropping the watcher stops the OS watch and shuts the worker thread down.
204pub struct Watcher {
205    root: PathBuf,
206    config: WatchConfig,
207    /// `Option` only so [`Drop`] can release it before joining the worker.
208    inner: Option<RecommendedWatcher>,
209    intents: Receiver<CoalescedIntent>,
210    control: Option<SyncSender<RawMessage>>,
211    cancelled: Arc<AtomicBool>,
212    worker_status: Arc<AtomicU8>,
213    worker: Option<JoinHandle<()>>,
214}
215
216#[cfg(test)]
217pub(crate) struct ScriptedSender {
218    root: PathBuf,
219    raw: SyncSender<RawMessage>,
220    overflowed: Arc<AtomicBool>,
221}
222
223#[cfg(test)]
224impl ScriptedSender {
225    pub(crate) fn send(&self, source: &str) -> Result<()> {
226        let events =
227            scripted_events::parse_script(source, &self.root).map_err(Error::WatchScript)?;
228        for event in events {
229            enqueue_raw(&self.raw, &self.overflowed, event);
230        }
231        Ok(())
232    }
233}
234
235/// Effects of one watch observation and any reconciliation it requested.
236#[derive(Debug)]
237pub struct WatchApplyReport {
238    /// Effect of the verified watch intent itself.
239    pub apply: ApplyOutcome,
240    /// Effect of closing any invalidation loop opened by the intent.
241    pub reconciliation: scan::ReconcileReport,
242}
243
244impl Watcher {
245    /// Start watching `root` recursively.
246    pub fn new(root: &Path, config: WatchConfig) -> Result<Self> {
247        config.validate()?;
248        let root = root.canonicalize().map_err(|e| Error::io(root, e))?;
249
250        let (raw_tx, raw_rx) = sync_channel::<RawMessage>(config.event_capacity);
251        let (intent_tx, intent_rx) = sync_channel::<CoalescedIntent>(config.intent_capacity);
252        let control_tx = raw_tx.clone();
253        let overflowed = Arc::new(AtomicBool::new(false));
254        let callback_overflowed = Arc::clone(&overflowed);
255
256        let mut inner = notify::recommended_watcher(move |result| {
257            enqueue_raw(&raw_tx, &callback_overflowed, result);
258        })
259        .map_err(|error| notify_error(&root, error))?;
260        inner.watch(&root, RecursiveMode::Recursive).map_err(|error| notify_error(&root, error))?;
261
262        let worker_root = root.clone();
263        let cancelled = Arc::new(AtomicBool::new(false));
264        let worker_cancelled = Arc::clone(&cancelled);
265        let worker_overflowed = Arc::clone(&overflowed);
266        let worker_status = Arc::new(AtomicU8::new(WORKER_RUNNING));
267        let tracked_status = Arc::clone(&worker_status);
268        let renames = RenameReporting::of_recommended_backend();
269        let worker = std::thread::Builder::new()
270            .name("fdu-watch".into())
271            .spawn(move || {
272                let _counter_guard = crate::counters::thread_flush_guard();
273                run_tracked_worker(&tracked_status, || {
274                    run_worker(
275                        &worker_root,
276                        config,
277                        renames,
278                        &raw_rx,
279                        &intent_tx,
280                        &worker_overflowed,
281                        &worker_cancelled,
282                    );
283                });
284            })
285            .map_err(|e| Error::io(&root, e))?;
286
287        Ok(Self {
288            root,
289            config,
290            inner: Some(inner),
291            intents: intent_rx,
292            control: Some(control_tx),
293            cancelled,
294            worker_status,
295            worker: Some(worker),
296        })
297    }
298
299    #[cfg(test)]
300    pub(crate) fn scripted(
301        root: &Path,
302        config: WatchConfig,
303        events: &Path,
304    ) -> Result<(Self, ScriptedSender)> {
305        config.validate()?;
306        let root = root.canonicalize().map_err(|error| Error::io(root, error))?;
307        let scripted = scripted_events::read_script(events, &root).map_err(Error::WatchScript)?;
308        let (raw_tx, raw_rx) = sync_channel::<RawMessage>(config.event_capacity);
309        let (intent_tx, intent_rx) = sync_channel::<CoalescedIntent>(config.intent_capacity);
310        let control_tx = raw_tx.clone();
311        let overflowed = Arc::new(AtomicBool::new(false));
312        for event in scripted {
313            enqueue_raw(&raw_tx, &overflowed, event);
314        }
315
316        let worker_root = root.clone();
317        let cancelled = Arc::new(AtomicBool::new(false));
318        let worker_cancelled = Arc::clone(&cancelled);
319        let worker_overflowed = Arc::clone(&overflowed);
320        let worker_status = Arc::new(AtomicU8::new(WORKER_RUNNING));
321        let tracked_status = Arc::clone(&worker_status);
322        // A script replaces the event source only, so it runs under this platform's
323        // production rename policy.
324        let renames = RenameReporting::of_recommended_backend();
325        let worker = std::thread::Builder::new()
326            .name("fdu-scripted-watch".into())
327            .spawn(move || {
328                let _counter_guard = crate::counters::thread_flush_guard();
329                run_tracked_worker(&tracked_status, || {
330                    run_worker(
331                        &worker_root,
332                        config,
333                        renames,
334                        &raw_rx,
335                        &intent_tx,
336                        &worker_overflowed,
337                        &worker_cancelled,
338                    );
339                });
340            })
341            .map_err(|error| Error::io(&root, error))?;
342
343        let sender = ScriptedSender {
344            root: root.clone(),
345            raw: control_tx.clone(),
346            overflowed: Arc::clone(&overflowed),
347        };
348        Ok((
349            Self {
350                root,
351                config,
352                inner: None,
353                intents: intent_rx,
354                control: Some(control_tx),
355                cancelled,
356                worker_status,
357                worker: Some(worker),
358            },
359            sender,
360        ))
361    }
362
363    /// Block for and verify the next coalesced intent, up to `timeout`.
364    ///
365    /// Filesystem calls happen synchronously on this consuming thread and may have
366    /// ordinary filesystem latency. A timeout, a stopped worker, and a panicked worker
367    /// are distinct outcomes. The returned observation contains relative paths but no
368    /// root identity; use [`Self::apply_next`] for the supported root-checked applying
369    /// driver rather than applying it to an arbitrary index.
370    pub fn next_observation(&self, timeout: Duration) -> Result<Option<Observation>> {
371        let Some(intent) = self.next_intent(timeout)? else {
372            return Ok(None);
373        };
374        Ok(Some(verify_intent(&self.root, self.config, &intent, &ScanConfig::default())))
375    }
376
377    fn next_intent(&self, timeout: Duration) -> Result<Option<CoalescedIntent>> {
378        match self.intents.recv_timeout(timeout) {
379            Ok(intent) => Ok(Some(intent)),
380            Err(RecvTimeoutError::Timeout) => Ok(None),
381            Err(RecvTimeoutError::Disconnected) => match self.worker_status.load(Ordering::Acquire)
382            {
383                WORKER_PANICKED => Err(Error::WatchWorkerPanicked),
384                _ => Err(Error::WatchStopped),
385            },
386        }
387    }
388
389    /// Advance capture through every raw hint accepted before this call.
390    ///
391    /// A full intent queue may retain their loss as one internal sticky overflow marker;
392    /// the next barrier advances that marker after the consumer drains queue capacity.
393    pub(crate) fn flush_capture(&self) -> Result<()> {
394        let Some(control) = self.control.as_ref() else {
395            return Err(Error::WatchStopped);
396        };
397        let (acknowledge, acknowledged) = sync_channel(0);
398        control.send(RawMessage::Flush(acknowledge)).map_err(|_| Error::WatchStopped)?;
399        acknowledged.recv().map_err(|_| Error::WatchStopped)
400    }
401
402    /// Maximum intents that can represent the raw hints preceding a flush barrier.
403    ///
404    /// The intent queue supplies the ordinary bound. One additional root invalidation
405    /// may be held as `sticky_overflow` when that queue filled, so the handoff derives
406    /// its drain work from the configured capacity rather than imposing another
407    /// unrelated limit.
408    pub(crate) const fn capture_backlog_bound(&self) -> usize {
409        self.config.intent_capacity.saturating_add(1)
410    }
411
412    /// Apply one verified watch observation and close any invalidation loop it opens.
413    ///
414    /// Restricted `max_depth` and `one_filesystem` scopes are rejected until the watch
415    /// adapter can filter raw backend events against those boundaries.
416    pub fn apply_next(
417        &self,
418        index: &IndexHandle,
419        scan_config: &ScanConfig,
420        timeout: Duration,
421        sink: &mut dyn FnMut(&Commit),
422    ) -> Result<Option<WatchApplyReport>> {
423        scan_config.validate_for_watch_scope(index.scope()?)?;
424        let indexed_root = index.root_path()?;
425        if indexed_root != self.root {
426            return Err(Error::WatchRootMismatch {
427                watched: self.root.clone(),
428                indexed: indexed_root,
429            });
430        }
431        let Some(intent) = self.next_intent(timeout)? else {
432            return Ok(None);
433        };
434        apply_intent(index, &self.root, self.config, &intent, scan_config, sink).map(Some)
435    }
436
437    /// Apply one intent through the opened-root lifecycle and exact resource boundary.
438    pub(crate) fn apply_next_controlled(
439        &self,
440        index: &IndexHandle,
441        scan_config: &ScanConfig,
442        timeout: Duration,
443        control: &dyn scan::ReconcileControl,
444        sink: &mut dyn FnMut(&Commit),
445    ) -> Result<Option<WatchApplyReport>> {
446        scan_config.validate_for_watch_scope(index.scope()?)?;
447        let indexed_root = index.root_path()?;
448        if indexed_root != self.root {
449            return Err(Error::WatchRootMismatch {
450                watched: self.root.clone(),
451                indexed: indexed_root,
452            });
453        }
454        let Some(intent) = self.next_intent(timeout)? else {
455            return Ok(None);
456        };
457        apply_intent_controlled(index, &self.root, self.config, &intent, scan_config, control, sink)
458            .map(Some)
459    }
460}
461
462fn apply_intent(
463    index: &IndexHandle,
464    root: &Path,
465    watch_config: WatchConfig,
466    intent: &CoalescedIntent,
467    scan_config: &ScanConfig,
468    sink: &mut dyn FnMut(&Commit),
469) -> Result<WatchApplyReport> {
470    let mut verifier =
471        |_: &Path, _: &Observation| Ok(verify_intent(root, watch_config, intent, scan_config));
472    let apply = apply_reverified_with(index, &Observation::default(), scan_config, &mut verifier)?;
473    if let Some(commit) = apply.commit.as_ref() {
474        sink(commit);
475    }
476    let reconciliation = scan::reconcile_pending_handle(index, scan_config, sink)?;
477    retain_unreadable(index, root, &reconciliation, sink)?;
478    Ok(WatchApplyReport { apply, reconciliation })
479}
480
481/// Retain why a reconciliation could not read part of the tree.
482///
483/// The shared reconcile settles an unreadable subtree rather than walking it again on
484/// every later event, so this report is the only time its error is seen. Committing the
485/// causes keeps the partial freshness it left explainable, as the opened root does, and
486/// names each path relative to the root so a later clean walk of it can drop the issue.
487fn retain_unreadable(
488    index: &IndexHandle,
489    root: &Path,
490    reconciliation: &scan::ReconcileReport,
491    sink: &mut dyn FnMut(&Commit),
492) -> Result<()> {
493    if reconciliation.scan.errors.is_empty() {
494        return Ok(());
495    }
496    let retained = reconciliation.scan.errors.len().min(crate::MAX_RETAINED_ISSUES);
497    let issues = reconciliation.scan.errors[..retained]
498        .iter()
499        .map(|error| crate::Issue::from_error_under(root, error))
500        .collect();
501    let omitted = u64::try_from(reconciliation.scan.errors.len() - retained).unwrap_or(u64::MAX);
502    let outcome =
503        index.transition_observation(crate::index::ObservationTransition::Unreadable {
504            issues,
505            omitted,
506        })?;
507    if let Some(commit) = outcome.commit.as_ref() {
508        sink(commit);
509    }
510    Ok(())
511}
512
513fn apply_intent_controlled(
514    index: &IndexHandle,
515    root: &Path,
516    watch_config: WatchConfig,
517    intent: &CoalescedIntent,
518    scan_config: &ScanConfig,
519    control: &dyn scan::ReconcileControl,
520    sink: &mut dyn FnMut(&Commit),
521) -> Result<WatchApplyReport> {
522    let mut verifier =
523        |_: &Path, _: &Observation| Ok(verify_intent(root, watch_config, intent, scan_config));
524    let apply = apply_reverified_with_control(
525        index,
526        &Observation::default(),
527        scan_config,
528        control,
529        &mut verifier,
530    )?;
531    if let Some(commit) = apply.commit.as_ref() {
532        sink(commit);
533    }
534    let reconciliation =
535        scan::reconcile_pending_handle_controlled(index, scan_config, control, sink)?;
536    Ok(WatchApplyReport { apply, reconciliation })
537}
538
539/// Test the unrooted observation driver without making it a public apply capability.
540///
541/// Production callers use [`Watcher::apply_next`], which proves that the watcher and
542/// index have the same root before consuming an intent. An [`Observation`] intentionally
543/// remains a generic producer batch and does not claim a filesystem root identity.
544#[cfg(test)]
545fn apply_observation(
546    index: &IndexHandle,
547    observation: &Observation,
548    scan_config: &ScanConfig,
549    sink: &mut dyn FnMut(&Commit),
550) -> Result<WatchApplyReport> {
551    scan_config.validate_for_watch_scope(index.scope()?)?;
552    let apply = apply_reverified(index, observation, scan_config)?;
553    if let Some(commit) = apply.commit.as_ref() {
554        sink(commit);
555    }
556    let reconciliation = scan::reconcile_pending_handle(index, scan_config, sink)?;
557    Ok(WatchApplyReport { apply, reconciliation })
558}
559
560/// Re-stat a queued watch sample against a clock-stable index boundary before applying
561/// it. Filesystem verification always runs outside the index lock. The filesystem itself
562/// cannot be locked by this process: a sample is valid at its `stat` linearization point,
563/// and any mutation after that point remains a later backend event. Queue loss or
564/// ambiguity becomes an invalidation and reconciliation, rather than a claim that the
565/// disk stayed frozen between `stat` and the in-memory commit. Sustained competing index
566/// writes conservatively invalidate the root without blocking readers on filesystem I/O.
567#[cfg(test)]
568fn apply_reverified(
569    index: &IndexHandle,
570    observation: &Observation,
571    scan_config: &ScanConfig,
572) -> Result<ApplyOutcome> {
573    let mut verifier = |root: &Path, observation: &Observation| {
574        reverify_observation(root, observation, scan_config)
575    };
576    apply_reverified_with(index, observation, scan_config, &mut verifier)
577}
578
579fn apply_reverified_with(
580    index: &IndexHandle,
581    observation: &Observation,
582    scan_config: &ScanConfig,
583    verifier: &mut impl FnMut(&Path, &Observation) -> Result<Observation>,
584) -> Result<ApplyOutcome> {
585    let (root, scope, _) = index.watch_boundary()?;
586    scan_config.validate_for_watch_scope(scope)?;
587    for _ in 0..MAX_OPTIMISTIC_APPLY_ATTEMPTS {
588        let clock = index.clock()?;
589        let candidate = verifier(&root, observation)?;
590        let candidate = escalate_unknown_ancestry(index, candidate)?;
591        if let Some(outcome) = index.apply_if_clock(clock, &candidate)? {
592            return Ok(outcome);
593        }
594    }
595
596    index.invalidate_root(InvalidateReason::WatchContention)
597}
598
599fn apply_reverified_with_control(
600    index: &IndexHandle,
601    observation: &Observation,
602    scan_config: &ScanConfig,
603    control: &dyn scan::ReconcileControl,
604    verifier: &mut impl FnMut(&Path, &Observation) -> Result<Observation>,
605) -> Result<ApplyOutcome> {
606    let (root, scope, _) = index.watch_boundary()?;
607    scan_config.validate_for_watch_scope(scope)?;
608    for _ in 0..MAX_OPTIMISTIC_APPLY_ATTEMPTS {
609        control.check_active()?;
610        let clock = index.clock()?;
611        let candidate = verifier(&root, observation)?;
612        let candidate = escalate_unknown_ancestry(index, candidate)?;
613        control.before_conditional_commit()?;
614        if let Some(outcome) =
615            index.apply_opened_if_clock(clock, &candidate, control.max_files())?
616        {
617            return Ok(outcome);
618        }
619    }
620
621    control.before_conditional_commit()?;
622    let outcome = index.apply_opened(
623        &Observation::new(vec![Op::InvalidateSubtree {
624            path: PathBuf::new(),
625            reason: InvalidateReason::WatchContention,
626        }]),
627        control.max_files(),
628    )?;
629    Ok(outcome)
630}
631
632/// Replace unverifiable child facts with bounded reconciliation hints.
633fn escalate_unknown_ancestry(index: &IndexHandle, candidate: Observation) -> Result<Observation> {
634    let unknown = index.unknown_ancestry(&candidate)?;
635    if unknown.is_empty() {
636        return Ok(candidate);
637    }
638
639    let mut roots: Vec<PathBuf> = unknown.into_iter().map(|(_, root)| root).collect();
640    roots.sort_by(|left, right| {
641        left.components().count().cmp(&right.components().count()).then_with(|| left.cmp(right))
642    });
643    roots.dedup();
644    let mut covering: Vec<PathBuf> = Vec::with_capacity(roots.len());
645    for root in roots {
646        if !covering.iter().any(|ancestor| root.starts_with(ancestor)) {
647            covering.push(root);
648        }
649    }
650
651    let mut ops: Vec<ObservationOp> = candidate
652        .ops
653        .into_iter()
654        .filter(|observed| !covering.iter().any(|root| observed.op.path().starts_with(root)))
655        .collect();
656    ops.extend(covering.into_iter().map(|path| {
657        ObservationOp::unconditional(Op::InvalidateSubtree {
658            path,
659            reason: InvalidateReason::UnknownAncestry,
660        })
661    }));
662    Ok(Observation::from_ops(ops))
663}
664
665#[cfg(test)]
666fn reverify_observation(
667    root: &Path,
668    observation: &Observation,
669    scan_config: &ScanConfig,
670) -> Result<Observation> {
671    let mut ops = Vec::with_capacity(observation.len().saturating_mul(2));
672    for observed in &observation.ops {
673        let relative = scan::normalize_subtree(observed.op.path())?;
674        match &observed.op {
675            Op::InvalidateSubtree { reason, .. } => {
676                ops.push(Op::InvalidateSubtree { path: relative, reason: *reason });
677            }
678            Op::Upsert { .. } | Op::Remove { .. } => {
679                let absolute = root.join(&relative);
680                match std::fs::symlink_metadata(&absolute) {
681                    Ok(metadata) => {
682                        let (kind, attrs) = scan::observe(&absolute, &metadata)
683                            .map_err(|source| Error::io(&absolute, source))?;
684                        match crate::admission::decide_path(
685                            &relative,
686                            kind,
687                            scan_config.hidden(),
688                            scan_config.exclude_special,
689                        ) {
690                            crate::admission::Disposition::Retain => {
691                                ops.push(Op::Upsert { path: relative.clone(), kind, attrs });
692                                if let Some(control) =
693                                    scan::read_control_op(scan_config, root, &relative, kind)?
694                                {
695                                    ops.push(control);
696                                }
697                            }
698                            crate::admission::Disposition::ControlOnly => {
699                                ops.push(Op::Remove { path: relative.clone() });
700                                if let Some(control) =
701                                    scan::read_control_op(scan_config, root, &relative, kind)?
702                                {
703                                    ops.push(control);
704                                }
705                            }
706                            crate::admission::Disposition::Reject => {
707                                ops.push(Op::Remove { path: relative });
708                            }
709                        }
710                    }
711                    Err(error) => ops.push(op_for_stat_error(relative, &error)),
712                }
713            }
714            Op::ControlUpsert { source, .. } => {
715                ops.push(Op::ControlUpsert { path: relative, source: source.clone() });
716            }
717            Op::ControlRemove { .. } => ops.push(Op::ControlRemove { path: relative }),
718        }
719    }
720    Ok(Observation::new(ops))
721}
722
723impl Drop for Watcher {
724    fn drop(&mut self) {
725        self.cancelled.store(true, Ordering::Release);
726        if let Some(control) = self.control.take() {
727            let _ = control.try_send(RawMessage::Stop);
728        }
729        // Stop the backend and release its callback sender before joining. Every worker
730        // send is nonblocking, so a full consumer queue cannot deadlock teardown.
731        self.inner.take();
732        if let Some(worker) = self.worker.take() {
733            let _ = worker.join();
734        }
735    }
736}
737
738fn notify_error(path: &Path, err: notify::Error) -> Error {
739    Error::io(path, std::io::Error::other(err))
740}
741
742fn enqueue_raw(
743    sender: &SyncSender<RawMessage>,
744    overflowed: &AtomicBool,
745    event: notify::Result<notify::Event>,
746) {
747    match sender.try_send(RawMessage::Event(event)) {
748        Ok(()) | Err(TrySendError::Disconnected(_)) => {}
749        Err(TrySendError::Full(_)) => overflowed.store(true, Ordering::Release),
750    }
751}
752
753fn run_tracked_worker(status: &AtomicU8, worker: impl FnOnce()) {
754    let outcome = catch_unwind(AssertUnwindSafe(worker));
755    status.store(if outcome.is_ok() { WORKER_STOPPED } else { WORKER_PANICKED }, Ordering::Release);
756}
757
758fn run_worker(
759    root: &Path,
760    config: WatchConfig,
761    renames: RenameReporting,
762    raw: &Receiver<RawMessage>,
763    out: &SyncSender<CoalescedIntent>,
764    overflowed: &AtomicBool,
765    cancelled: &AtomicBool,
766) {
767    let mut pending: BTreeMap<PathBuf, Pending> = BTreeMap::new();
768    let mut batch_started: Option<Instant> = None;
769    let mut sticky_overflow = false;
770
771    loop {
772        if cancelled.load(Ordering::Acquire) {
773            return;
774        }
775        if overflowed.swap(false, Ordering::AcqRel) {
776            collapse_to_overflow(&mut pending);
777            batch_started.get_or_insert_with(Instant::now);
778        }
779        if sticky_overflow {
780            match try_deliver_overflow(out) {
781                Ok(true) => sticky_overflow = false,
782                Ok(false) => {}
783                Err(()) => return,
784            }
785        }
786
787        match raw.recv_timeout(config.settle) {
788            Ok(RawMessage::Event(Ok(event))) => {
789                record(root, &event, &mut pending, config.batch_path_capacity, renames);
790                batch_started.get_or_insert_with(Instant::now);
791            }
792            Ok(RawMessage::Event(Err(err))) => {
793                // notify reports a watch failure. It cannot say what was missed, so the
794                // only honest response is to escalate the whole tree.
795                let _ = err;
796                collapse_to_overflow(&mut pending);
797                batch_started.get_or_insert_with(Instant::now);
798            }
799            Ok(RawMessage::Flush(acknowledge)) => {
800                if overflowed.swap(false, Ordering::AcqRel) {
801                    collapse_to_overflow(&mut pending);
802                }
803                if !pending.is_empty()
804                    && try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err()
805                {
806                    return;
807                }
808                // The consumer may have drained the queue while this worker waited
809                // for the barrier. Publish sticky loss before acknowledging; doing it
810                // at the next loop iteration races the consumer's final empty poll.
811                if sticky_overflow {
812                    match try_deliver_overflow(out) {
813                        Ok(true) => sticky_overflow = false,
814                        Ok(false) => {}
815                        Err(()) => return,
816                    }
817                }
818                batch_started = None;
819                let _ = acknowledge.send(());
820            }
821            Ok(RawMessage::Stop) => return,
822            Err(RecvTimeoutError::Timeout) => {
823                // Quiet for a full step: the batch has settled.
824                if !pending.is_empty()
825                    && try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err()
826                {
827                    return;
828                }
829                batch_started = None;
830                continue;
831            }
832            Err(RecvTimeoutError::Disconnected) => {
833                if !cancelled.load(Ordering::Acquire) && !pending.is_empty() {
834                    let _ = try_deliver_pending(&mut pending, out, &mut sticky_overflow);
835                }
836                return;
837            }
838        }
839
840        // A tree under continuous churn never goes quiet, so cap how long a batch waits.
841        if max_hold_elapsed(batch_started, config.max_hold) {
842            if try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err() {
843                return;
844            }
845            batch_started = None;
846        }
847    }
848}
849
850fn max_hold_elapsed(started: Option<Instant>, max_hold: Duration) -> bool {
851    started.is_some_and(|start| start.elapsed() >= max_hold)
852}
853
854/// Fold one event into the pending set.
855fn record(
856    root: &Path,
857    event: &notify::Event,
858    pending: &mut BTreeMap<PathBuf, Pending>,
859    capacity: usize,
860    reporting: RenameReporting,
861) {
862    if event.need_rescan() {
863        // The kernel dropped events. Escalate the narrowest subtree the event names, or
864        // the whole root when it names zero/multiple paths or crosses the watch boundary.
865        let target = if event.paths.len() == 1 {
866            relative_to(root, &event.paths[0]).unwrap_or_default()
867        } else {
868            PathBuf::new()
869        };
870        queue_pending(
871            pending,
872            target,
873            Pending::Escalate(InvalidateReason::WatchOverflow),
874            capacity,
875        );
876        return;
877    }
878
879    // Every rename shape is handled one named path at a time. Events are hints, not
880    // facts: a rename says only that the name it carries may have gained or lost an
881    // entry, which is exactly what a create or a remove of that name says. So each named
882    // path inside the root is verified like one, and nothing is inferred about where its
883    // counterpart went. That needs no pairing and no guess at a parent:
884    //
885    // - an old name that is gone verifies as a removal of it and its subtree;
886    // - a new name verifies as an upsert, and a directory there is relisted because its
887    //   contents arrived without events (see `verify_intent`);
888    // - a counterpart inside the root is named by its own event on every backend that
889    //   promises `RenameReporting::EachSide`, and is verified the same way;
890    // - a counterpart outside the root changes nothing inside it.
891    //
892    // Scoping to the named path adds no trust in the backend beyond what create and
893    // remove verification already needs: that every name whose entry changed is named
894    // by some event, or else covered by a loss signal, which escalates above.
895    let renamed = matches!(event.kind, EventKind::Modify(notify::event::ModifyKind::Name(_)));
896    if renamed
897        && (reporting == RenameReporting::OldSideOnly
898            || !event.paths.iter().any(|path| relative_to(root, path).is_some()))
899    {
900        // Two renames this layer cannot bound. A backend that may never name the new
901        // side (kqueue) leaves the moved entry anywhere in the tree. A rename naming no
902        // path this root can place is a hint about somewhere the engine cannot locate.
903        // Both reconcile the whole root, which cannot leave an old name behind or miss a
904        // moved-in subtree.
905        queue_pending(
906            pending,
907            PathBuf::new(),
908            Pending::Escalate(InvalidateReason::UnpairedRename),
909            capacity,
910        );
911    }
912
913    for path in &event.paths {
914        let Some(rel) = relative_to(root, path) else {
915            continue; // Outside the root: nothing inside it changed at this name.
916        };
917        if matches!(event.kind, EventKind::Access(_)) {
918            continue; // Reads change nothing this engine records.
919        }
920        if renamed && rel.as_os_str().is_empty() {
921            // The root itself moved (inotify reports its own move as `From`). What now
922            // sits at the root path, if anything, is unrelated to the index, so the
923            // bound is the whole root.
924            queue_pending(
925                pending,
926                PathBuf::new(),
927                Pending::Escalate(InvalidateReason::UnpairedRename),
928                capacity,
929            );
930            continue;
931        }
932        let relist_if_dir = matches!(event.kind, EventKind::Create(_));
933        queue_pending(pending, rel, Pending::Verify { relist_if_dir, renamed }, capacity);
934    }
935}
936
937fn queue_pending(
938    pending: &mut BTreeMap<PathBuf, Pending>,
939    path: PathBuf,
940    state: Pending,
941    capacity: usize,
942) {
943    if matches!(pending.get(Path::new("")), Some(Pending::Escalate(_))) {
944        return;
945    }
946    if path.as_os_str().is_empty() && matches!(state, Pending::Escalate(_)) {
947        pending.clear();
948        pending.insert(path, state);
949        return;
950    }
951    if let Some(existing) = pending.get_mut(&path) {
952        match (existing, state) {
953            (Pending::Escalate(_), _) => {}
954            (
955                Pending::Verify { relist_if_dir, renamed },
956                Pending::Verify { relist_if_dir: relist, renamed: rename },
957            ) => {
958                *relist_if_dir |= relist;
959                *renamed |= rename;
960            }
961            (slot @ Pending::Verify { .. }, Pending::Escalate(reason)) => {
962                *slot = Pending::Escalate(reason);
963            }
964        }
965        return;
966    }
967    if pending.len() >= capacity {
968        collapse_to_overflow(pending);
969    } else {
970        pending.insert(path, state);
971    }
972}
973
974fn collapse_to_overflow(pending: &mut BTreeMap<PathBuf, Pending>) {
975    pending.clear();
976    pending.insert(PathBuf::new(), Pending::Escalate(InvalidateReason::WatchOverflow));
977}
978
979fn try_deliver_pending(
980    pending: &mut BTreeMap<PathBuf, Pending>,
981    out: &SyncSender<CoalescedIntent>,
982    sticky_overflow: &mut bool,
983) -> std::result::Result<(), ()> {
984    if pending.is_empty() {
985        return Ok(());
986    }
987    let intent = CoalescedIntent { pending: std::mem::take(pending) };
988    match out.try_send(intent) {
989        Ok(()) => Ok(()),
990        Err(TrySendError::Full(_)) => {
991            *sticky_overflow = true;
992            Ok(())
993        }
994        Err(TrySendError::Disconnected(_)) => Err(()),
995    }
996}
997
998fn try_deliver_overflow(out: &SyncSender<CoalescedIntent>) -> std::result::Result<bool, ()> {
999    let mut pending = BTreeMap::new();
1000    collapse_to_overflow(&mut pending);
1001    match out.try_send(CoalescedIntent { pending }) {
1002        Ok(()) => Ok(true),
1003        Err(TrySendError::Full(_)) => Ok(false),
1004        Err(TrySendError::Disconnected(_)) => Err(()),
1005    }
1006}
1007
1008/// Verify one bounded intent: stat once per path, never once per backend event.
1009fn verify_intent(
1010    root: &Path,
1011    config: WatchConfig,
1012    intent: &CoalescedIntent,
1013    scan_config: &ScanConfig,
1014) -> Observation {
1015    let mut ops = Vec::with_capacity(intent.pending.len());
1016    let mut listings = ParentListings::default();
1017    // Renamed names whose exact spelling their parent does not list. Nothing exists at
1018    // any path below such a name either, and its parent's reconciliation covers the
1019    // subtree, so their pending descendants are not verified through it. The pending map
1020    // is ordered by component, so a name is always settled before its descendants.
1021    let mut unlisted: Vec<PathBuf> = Vec::new();
1022
1023    for (rel, state) in &intent.pending {
1024        match state {
1025            Pending::Escalate(reason) => {
1026                ops.push(Op::InvalidateSubtree { path: rel.clone(), reason: *reason });
1027            }
1028            Pending::Verify { relist_if_dir, renamed } => {
1029                if unlisted.iter().any(|name| rel.starts_with(name)) {
1030                    continue;
1031                }
1032                let absolute = root.join(rel);
1033                let mut stat = std::fs::symlink_metadata(&absolute);
1034                if *renamed && !rel.as_os_str().is_empty() && stat.is_ok() {
1035                    // A lookup on a case- or normalization-insensitive filesystem (APFS
1036                    // and HFS+ by default, NTFS, casefolded ext4) resolves `Readme` to a
1037                    // stored `README`. The old side of a rename that changed only case
1038                    // therefore stats as present, and upserting it would keep both
1039                    // spellings. The parent's listing holds the stored names, so an exact
1040                    // match proves membership.
1041                    let unlisted_reason = match listings.lists(root, rel) {
1042                        Some(true) => None,
1043                        Some(false) => {
1044                            // The listing was read after the first stat, so a miss is
1045                            // either a stale spelling or a name renamed away since. A
1046                            // second stat tells them apart: a name that is now gone is an
1047                            // ordinary removal, verified below like any other.
1048                            stat = std::fs::symlink_metadata(&absolute);
1049                            stat.is_ok().then_some(InvalidateReason::UnpairedRename)
1050                        }
1051                        // An unlistable parent cannot prove membership either way; its
1052                        // reconciliation retries and reports why.
1053                        None => Some(InvalidateReason::VerificationFailed),
1054                    };
1055                    if let Some(reason) = unlisted_reason {
1056                        // This name is not an entry, and which stored name the index holds
1057                        // for it is unknown, so its parent reconciles.
1058                        let parent = parent_of(rel);
1059                        if !ops.iter().any(
1060                            |op| matches!(op, Op::InvalidateSubtree { path, .. } if *path == parent),
1061                        ) {
1062                            ops.push(Op::InvalidateSubtree { path: parent, reason });
1063                        }
1064                        unlisted.push(rel.clone());
1065                        continue;
1066                    }
1067                }
1068                match stat {
1069                    Ok(meta) => {
1070                        let Ok((kind, attrs)) = scan::observe(&absolute, &meta) else {
1071                            ops.push(Op::InvalidateSubtree {
1072                                path: rel.parent().map_or_else(PathBuf::new, Path::to_path_buf),
1073                                reason: InvalidateReason::VerificationFailed,
1074                            });
1075                            continue;
1076                        };
1077                        let disposition = crate::admission::decide_path(
1078                            rel,
1079                            kind,
1080                            scan_config.hidden(),
1081                            scan_config.exclude_special,
1082                        );
1083                        match disposition {
1084                            crate::admission::Disposition::Retain => {
1085                                ops.push(Op::Upsert { path: rel.clone(), kind, attrs });
1086                                match scan::read_control_op(scan_config, root, rel, kind) {
1087                                    Ok(Some(control)) => ops.push(control),
1088                                    Ok(None) => {}
1089                                    Err(_) => ops.push(Op::InvalidateSubtree {
1090                                        path: rel
1091                                            .parent()
1092                                            .map_or_else(PathBuf::new, Path::to_path_buf),
1093                                        reason: InvalidateReason::VerificationFailed,
1094                                    }),
1095                                }
1096                            }
1097                            crate::admission::Disposition::ControlOnly => {
1098                                match scan::read_control_op(scan_config, root, rel, kind) {
1099                                    Ok(Some(control)) => {
1100                                        ops.push(Op::Remove { path: rel.clone() });
1101                                        ops.push(control);
1102                                    }
1103                                    Ok(None) => ops.push(Op::Remove { path: rel.clone() }),
1104                                    Err(_) => ops.push(Op::InvalidateSubtree {
1105                                        path: rel
1106                                            .parent()
1107                                            .map_or_else(PathBuf::new, Path::to_path_buf),
1108                                        reason: InvalidateReason::VerificationFailed,
1109                                    }),
1110                                }
1111                            }
1112                            crate::admission::Disposition::Reject => {
1113                                ops.push(Op::Remove { path: rel.clone() });
1114                            }
1115                        }
1116                        let retained_dir =
1117                            disposition == crate::admission::Disposition::Retain && kind.is_dir();
1118                        if retained_dir && *renamed {
1119                            // A directory that arrived by rename brought its contents
1120                            // with it, and no backend reports a moved tree's contents.
1121                            // This is not the registration race below, so it does not
1122                            // depend on `relist_new_dirs`: nothing else will ever report
1123                            // these entries. The bound is this directory's subtree.
1124                            ops.push(Op::InvalidateSubtree {
1125                                path: rel.clone(),
1126                                reason: InvalidateReason::UnpairedRename,
1127                            });
1128                        } else if retained_dir && *relist_if_dir && config.relist_new_dirs {
1129                            // The watch for this directory was installed after it was
1130                            // created, so anything already inside produced no event.
1131                            ops.push(Op::InvalidateSubtree {
1132                                path: rel.clone(),
1133                                reason: InvalidateReason::WatchSetupRace,
1134                            });
1135                        }
1136                    }
1137                    Err(error) => {
1138                        ops.push(op_for_stat_error(rel.clone(), &error));
1139                        // The exact name's removal drops its rules in the index. A case
1140                        // variant's says nothing by itself: another spelling may still
1141                        // resolve, or it was never the control, so the canonical path is
1142                        // looked up and answers, a miss removing the rules.
1143                        if error.kind() == std::io::ErrorKind::NotFound
1144                            && crate::control::path_control_spelling(rel)
1145                                == Some(crate::control::ControlSpelling::Variant)
1146                        {
1147                            let control = crate::control::sibling_control_path(rel);
1148                            match scan::read_directory_control_or_removal(
1149                                scan_config,
1150                                root,
1151                                &control,
1152                            ) {
1153                                Ok(Some(observed)) => ops.push(observed),
1154                                Ok(None) => {}
1155                                Err(_) => ops.push(Op::InvalidateSubtree {
1156                                    path: parent_of(rel),
1157                                    reason: InvalidateReason::VerificationFailed,
1158                                }),
1159                            }
1160                        }
1161                    }
1162                }
1163                if scan_config.population != crate::query::IgnoredEntries::Include
1164                    && crate::control::path_control_spelling(rel).is_some()
1165                {
1166                    ops.push(Op::InvalidateSubtree {
1167                        path: rel.parent().map_or_else(PathBuf::new, Path::to_path_buf),
1168                        reason: InvalidateReason::ControlPopulationChanged,
1169                    });
1170                }
1171            }
1172        }
1173    }
1174    Observation::new(ops)
1175}
1176
1177/// The relative parent of a non-root path; the root for a top-level name.
1178fn parent_of(rel: &Path) -> PathBuf {
1179    rel.parent().map_or_else(PathBuf::new, Path::to_path_buf)
1180}
1181
1182/// Stored names of the parent directories one intent's renamed paths live in.
1183///
1184/// Only a renamed name that stats as present is looked up, and a parent is listed once
1185/// per verified intent unless a lookup misses, so the cost is bounded by the directories
1186/// renames touched in the batch rather than by the tree.
1187#[derive(Default)]
1188struct ParentListings(BTreeMap<PathBuf, Option<std::collections::HashSet<std::ffi::OsString>>>);
1189
1190impl ParentListings {
1191    /// Whether the parent of `rel` lists its final component byte for byte, or `None`
1192    /// when the parent cannot be listed completely.
1193    ///
1194    /// A miss against a listing read earlier in this intent is re-read before it is
1195    /// answered: on a busy directory the name may have been created since, and a stale
1196    /// listing must not turn a new entry into a parent reconcile.
1197    fn lists(&mut self, root: &Path, rel: &Path) -> Option<bool> {
1198        let name = rel.file_name()?;
1199        let parent = parent_of(rel);
1200        if let Some(Some(names)) = self.0.get(&parent) {
1201            if names.contains(name) {
1202                return Some(true);
1203            }
1204        }
1205        let names: Option<std::collections::HashSet<_>> = std::fs::read_dir(root.join(&parent))
1206            .and_then(|entries| entries.map(|entry| entry.map(|entry| entry.file_name())).collect())
1207            .ok();
1208        let listed = names.as_ref().map(|names| names.contains(name));
1209        self.0.insert(parent, names);
1210        listed
1211    }
1212}
1213
1214fn op_for_stat_error(path: PathBuf, error: &std::io::Error) -> Op {
1215    match error.kind() {
1216        std::io::ErrorKind::NotFound if path.as_os_str().is_empty() => {
1217            Op::InvalidateSubtree { path, reason: InvalidateReason::VerificationFailed }
1218        }
1219        std::io::ErrorKind::NotFound => Op::Remove { path },
1220        std::io::ErrorKind::NotADirectory => Op::InvalidateSubtree {
1221            path: path.parent().map_or_else(PathBuf::new, Path::to_path_buf),
1222            reason: InvalidateReason::VerificationFailed,
1223        },
1224        _ => Op::InvalidateSubtree { path, reason: InvalidateReason::VerificationFailed },
1225    }
1226}
1227
1228/// Express an absolute path relative to the watch root.
1229///
1230/// Returns `None` for anything outside the root, which should not happen but is not
1231/// worth trusting a backend about.
1232fn relative_to(root: &Path, path: &Path) -> Option<PathBuf> {
1233    path.strip_prefix(root).ok().map(Path::to_path_buf)
1234}
1235
1236#[cfg(test)]
1237mod tests {
1238    use super::*;
1239    use notify::event::{CreateKind, Flag, MetadataKind, ModifyKind, RenameMode};
1240    use std::fs;
1241
1242    fn queued_test_watcher(root: PathBuf) -> (SyncSender<CoalescedIntent>, Watcher) {
1243        let (sender, intents) = sync_channel(1);
1244        let watcher = Watcher {
1245            root,
1246            config: WatchConfig::default(),
1247            inner: None,
1248            intents,
1249            control: None,
1250            cancelled: Arc::new(AtomicBool::new(false)),
1251            worker_status: Arc::new(AtomicU8::new(WORKER_RUNNING)),
1252            worker: None,
1253        };
1254        (sender, watcher)
1255    }
1256
1257    #[test]
1258    fn watcher_can_move_to_its_single_consumer_thread() {
1259        fn assert_send<T: Send>() {}
1260
1261        assert_send::<Watcher>();
1262    }
1263
1264    /// Collect deltas until `want` is satisfied or the deadline passes.
1265    ///
1266    /// Event latency varies by orders of magnitude across backends (inotify is
1267    /// immediate, `FSEvents` batches), so the test waits on a condition rather than
1268    /// sleeping for a fixed guess.
1269    /// Serializes the tests that bind a real filesystem watcher.
1270    ///
1271    /// Real-backend delivery latency depends on how much else is happening on the
1272    /// volume. Roughly two hundred tempdirs churn concurrently across this binary, and
1273    /// `FSEvents` is volume-wide, so a stream bound while that is going on can take far
1274    /// longer to deliver its first event than one bound on a quiet machine.
1275    ///
1276    /// Measured at `origin/main`, before any change here, so this is pre-existing
1277    /// backend behavior rather than something a caller introduced:
1278    ///
1279    /// | libtest threads | Result |
1280    /// | --- | --- |
1281    /// | 1 | 406 passed in 8.39 s |
1282    /// | 4 | 3 failed in 22.76 s |
1283    /// | 10 | 3 failed in 21.87 s |
1284    ///
1285    /// Two things follow, and both are applied. Contention between the real-backend
1286    /// tests themselves is removed by this lock, while the rest of the suite keeps
1287    /// running in parallel — the engine is not implicated, since every other watch test
1288    /// drives the same worker through scripted observations and none of them is
1289    /// affected. Residual variance from the surrounding churn is absorbed by
1290    /// [`REAL_BACKEND_DELIVERY`] rather than by retrying.
1291    static REAL_WATCHER: std::sync::Mutex<()> = std::sync::Mutex::new(());
1292
1293    /// How long a real backend may take to deliver its first event before the test fails.
1294    ///
1295    /// The events do arrive; earlier analysis here claimed they did not, inferred from a
1296    /// suite that finished in about the old deadline, and that inference was wrong. With
1297    /// a deadline long enough not to cut delivery off, the full parallel binary passes
1298    /// five runs out of five and finishes in about 18.6 s — *faster* than the runs that
1299    /// failed, because a dead timeout was the longest thing in those.
1300    ///
1301    /// Sixty seconds is chosen to be far outside the observed distribution rather than
1302    /// tuned to its edge, since a deadline that merely covers today's variance becomes
1303    /// tomorrow's flake on a busier machine. The cost is asymmetric and cheap: this
1304    /// duration is only ever spent when delivery genuinely fails, and a passing run
1305    /// never approaches it.
1306    const REAL_BACKEND_DELIVERY: Duration = Duration::from_secs(60);
1307
1308    /// Hold exclusive access to the real watch backend for the rest of the test.
1309    ///
1310    /// Poisoning is deliberately ignored. The lock guards an OS resource rather than
1311    /// shared data, so a panic in one test leaves nothing for the next to observe, and
1312    /// propagating the poison would turn one real failure into three.
1313    fn real_watcher_guard() -> std::sync::MutexGuard<'static, ()> {
1314        REAL_WATCHER.lock().unwrap_or_else(std::sync::PoisonError::into_inner)
1315    }
1316
1317    /// What a wait for events produced.
1318    enum Waited {
1319        /// The awaited ops arrived, with everything seen before them.
1320        Delivered(Vec<Op>),
1321        /// Nothing whatsoever arrived before the deadline.
1322        Silent,
1323    }
1324
1325    /// Collect ops until `want` is satisfied, separating three outcomes that are not the
1326    /// same failure.
1327    ///
1328    /// `Delivered` — the events arrived and matched. The caller's assertions then run at
1329    /// full strength.
1330    ///
1331    /// Panic — events arrived but never matched. That is a real disagreement about
1332    /// content and must fail. This previously returned the partial list instead, so every
1333    /// caller asserted on it and announced a violated product contract when nothing had
1334    /// arrived at all, sending the next reader after a defect that was not there.
1335    ///
1336    /// `Silent` — *nothing whatsoever* arrived before the deadline. What that means
1337    /// depends on whether the watch was already known to deliver, so the caller decides:
1338    /// [`establish_watch`] may read it as the host's silence, and [`wait_established`]
1339    /// reads it as fdu's.
1340    ///
1341    /// The distinction is observable rather than assumed: a working backend delivers the
1342    /// test's own writes within milliseconds, so an empty list after a full minute means
1343    /// the stream is dead, not slow. On this project's macOS development host a degraded
1344    /// `fseventsd` produces exactly that, while both CI platforms never have.
1345    fn wait_for(
1346        watcher: &Watcher,
1347        deadline: Duration,
1348        mut want: impl FnMut(&[Op]) -> bool,
1349    ) -> Waited {
1350        let start = Instant::now();
1351        let mut seen: Vec<Op> = Vec::new();
1352        while start.elapsed() < deadline {
1353            match watcher.next_observation(Duration::from_millis(200)) {
1354                Ok(Some(observation)) => {
1355                    seen.extend(observation.ops.into_iter().map(|observed| observed.op));
1356                    if want(&seen) {
1357                        return Waited::Delivered(seen);
1358                    }
1359                }
1360                Ok(None) => {}
1361                Err(error) => panic!("watcher stopped while waiting: {error}"),
1362            }
1363        }
1364        if seen.is_empty() {
1365            return Waited::Silent;
1366        }
1367        panic!(
1368            "the backend delivered {} op(s) in {deadline:?} but never the one awaited, so \
1369             this is a disagreement about content rather than a delivery failure: {seen:?}",
1370            seen.len()
1371        );
1372    }
1373
1374    /// Wait on a watch that [`establish_watch`] has already proven live.
1375    ///
1376    /// Silence here is evidence about fdu rather than about the host: the backend delivered
1377    /// the warm-up, so a later write it never reports is a lost event. No opt-out applies,
1378    /// which is why the message offers none.
1379    fn wait_established(
1380        watcher: &Watcher,
1381        deadline: Duration,
1382        want: impl FnMut(&[Op]) -> bool,
1383    ) -> Vec<Op> {
1384        match wait_for(watcher, deadline, want) {
1385            Waited::Delivered(ops) => ops,
1386            Waited::Silent => panic!(
1387                "the watch was established and then delivered nothing in {deadline:?}, so \
1388                 this is a lost event rather than a host precondition; \
1389                 FDU_TEST_ALLOW_NO_NATIVE_WATCH does not apply here"
1390            ),
1391        }
1392    }
1393
1394    /// Write into the root until the watch is provably live, then return.
1395    ///
1396    /// A watcher bound to a directory is not yet watching it. `Watcher::new` returns once
1397    /// registration is *requested*, and anything written before it takes effect produces
1398    /// no event — the engine detects exactly this and answers
1399    /// `InvalidateSubtree { path: "", reason: WatchSetupRace }`, meaning "I missed a
1400    /// window, relist the root".
1401    ///
1402    /// That is the correct answer, and it is why these tests were failing intermittently.
1403    /// A test that writes its subject immediately after `Watcher::new` is racing
1404    /// registration: when it loses, the engine reports the race rather than the file, and
1405    /// an assertion waiting for that file's own upsert rejects a valid reply. The failure
1406    /// looked like flakiness, then like a dead backend, and was neither — it was a real
1407    /// race the test set up for itself, and which the engine reported faithfully.
1408    ///
1409    /// Writing a warm-up file and waiting for *any* event settles it: whichever arrives,
1410    /// the stream is delivering and registration is complete, so a write afterwards
1411    /// cannot fall into the setup window. Everything the test then asserts is about fdu.
1412    ///
1413    /// Returns `false` only for a host explicitly declared unable to deliver native watch
1414    /// events, with `FDU_TEST_ALLOW_NO_NATIVE_WATCH=1`. Without that declaration, silence
1415    /// during the warm-up is an actionable precondition failure rather than a passing test.
1416    fn establish_watch(watcher: &Watcher, dir: &Path) -> bool {
1417        let warmup = dir.join(".fdu-watch-warmup");
1418        fs::write(&warmup, b"warmup").expect("warmup write");
1419        let waited = wait_for(watcher, REAL_BACKEND_DELIVERY, |ops| !ops.is_empty());
1420        let _ = fs::remove_file(&warmup);
1421        match waited {
1422            Waited::Delivered(_) => true,
1423            Waited::Silent => {
1424                if std::env::var_os("FDU_TEST_ALLOW_NO_NATIVE_WATCH").as_deref()
1425                    == Some(std::ffi::OsStr::new("1"))
1426                {
1427                    eprintln!(
1428                        "skipped by FDU_TEST_ALLOW_NO_NATIVE_WATCH=1: the host event service \
1429                         delivered no events to this stream"
1430                    );
1431                    return false;
1432                }
1433                panic!(
1434                    "native watch precondition failed: the host delivered no events in \
1435                     {REAL_BACKEND_DELIVERY:?}; run on a host with event delivery, or \
1436                     explicitly opt out with FDU_TEST_ALLOW_NO_NATIVE_WATCH=1"
1437                );
1438            }
1439        }
1440    }
1441
1442    #[test]
1443    fn created_files_arrive_as_verified_upserts() {
1444        let _serialized = real_watcher_guard();
1445        let dir = tempfile::tempdir().expect("tempdir");
1446        let watcher = Watcher::new(dir.path(), WatchConfig::default()).expect("watcher");
1447        if !establish_watch(&watcher, dir.path()) {
1448            return;
1449        }
1450
1451        fs::write(dir.path().join("hello.txt"), b"hello world").expect("write");
1452
1453        let ops = wait_established(&watcher, REAL_BACKEND_DELIVERY, |ops| {
1454            ops.iter().any(|op| op.path() == Path::new("hello.txt"))
1455        });
1456
1457        let found = ops
1458            .iter()
1459            .find(|op| op.path() == Path::new("hello.txt"))
1460            .expect("an op for the new file");
1461        match found {
1462            Op::Upsert { attrs, kind, .. } => {
1463                assert!(!kind.is_dir());
1464                // The point of verify-then-emit: the delta carries real stat data, which
1465                // no backend put in the event.
1466                assert_eq!(attrs.size, 11);
1467                assert!(attrs.mtime_ns > 0);
1468            }
1469            other => panic!("expected an upsert, got {other:?}"),
1470        }
1471    }
1472
1473    #[test]
1474    fn deleted_files_arrive_as_removes() {
1475        let _serialized = real_watcher_guard();
1476        let dir = tempfile::tempdir().expect("tempdir");
1477        let path = dir.path().join("doomed.txt");
1478        fs::write(&path, b"x").expect("write");
1479
1480        let watcher = Watcher::new(dir.path(), WatchConfig::default()).expect("watcher");
1481        if !establish_watch(&watcher, dir.path()) {
1482            return;
1483        }
1484        fs::remove_file(&path).expect("remove");
1485
1486        let ops = wait_established(&watcher, REAL_BACKEND_DELIVERY, |ops| {
1487            ops.iter()
1488                .any(|op| matches!(op, Op::Remove { path } if path == Path::new("doomed.txt")))
1489        });
1490
1491        assert!(
1492            ops.iter()
1493                .any(|op| matches!(op, Op::Remove { path } if path == Path::new("doomed.txt"))),
1494            "expected a remove, saw {ops:?}"
1495        );
1496    }
1497
1498    // The watch-setup race — a directory created and populated before its watch is
1499    // installed, which must escalate to a relist — is asserted by
1500    // `opened::tests::scripted_directory_creation_closes_the_registration_gap`. That test
1501    // drives the same `WatchSetupRace` invalidation and the same subsequent discovery of
1502    // the child through a scripted observation, so the semantics are pinned on every
1503    // platform without depending on when a backend happens to register.
1504    //
1505    // A real-backend version of it lived here and was removed rather than repaired. It
1506    // could not make the claim it appeared to: the macOS backend watches recursively from
1507    // the root, so there is no per-directory registration window for it to lose, and the
1508    // test spent its full twenty-second deadline waiting for an escalation that platform
1509    // has no reason to emit. Serializing the real-watcher tests fixed the other two and
1510    // left this one failing, which is what showed the difference is in the scenario and
1511    // not in the contention.
1512    //
1513    // What stays here is the minimal real-backend smoke the architecture asks for:
1514    // create and remove actually arrive, and a missing root actually fails.
1515
1516    #[test]
1517    fn watching_a_missing_path_is_an_error() {
1518        let dir = tempfile::tempdir().expect("tempdir");
1519        let missing = dir.path().join("not-there");
1520        assert!(Watcher::new(&missing, WatchConfig::default()).is_err());
1521    }
1522
1523    #[test]
1524    fn paths_outside_the_root_are_ignored() {
1525        let root = Path::new("/a/b");
1526        assert_eq!(relative_to(root, Path::new("/a/b/c/d")), Some(PathBuf::from("c/d")));
1527        assert_eq!(relative_to(root, Path::new("/elsewhere")), None);
1528    }
1529
1530    #[test]
1531    fn verification_errors_distinguish_absence_from_an_invalid_ancestor() {
1532        let path = PathBuf::from("parent/known.txt");
1533        let missing = op_for_stat_error(
1534            path.clone(),
1535            &std::io::Error::new(std::io::ErrorKind::NotFound, "gone"),
1536        );
1537        assert!(matches!(missing, Op::Remove { path: removed } if removed == path));
1538
1539        let not_a_directory = op_for_stat_error(
1540            path.clone(),
1541            &std::io::Error::new(std::io::ErrorKind::NotADirectory, "ancestor is a file"),
1542        );
1543        assert!(matches!(
1544            not_a_directory,
1545            Op::InvalidateSubtree {
1546                path: invalidated,
1547                reason: InvalidateReason::VerificationFailed,
1548            } if invalidated == Path::new("parent")
1549        ));
1550
1551        let denied = op_for_stat_error(
1552            path.clone(),
1553            &std::io::Error::new(std::io::ErrorKind::PermissionDenied, "denied"),
1554        );
1555        assert!(matches!(
1556            denied,
1557            Op::InvalidateSubtree {
1558                path: invalidated,
1559                reason: InvalidateReason::VerificationFailed,
1560            } if invalidated == path
1561        ));
1562    }
1563
1564    #[test]
1565    fn unknown_watch_ancestry_reconciles_from_the_nearest_known_directory() {
1566        let dir = tempfile::tempdir().expect("tempdir");
1567        let (index, _) =
1568            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
1569        let handle = crate::IndexHandle::new(index);
1570        let nested = dir.path().join("new/deep");
1571        fs::create_dir_all(&nested).expect("nested directories");
1572        fs::write(nested.join("file.txt"), b"verified").expect("nested file");
1573
1574        let report = apply_observation(
1575            &handle,
1576            &Observation::new(vec![Op::Upsert {
1577                path: PathBuf::from("new/deep/file.txt"),
1578                kind: crate::EntryKind::File,
1579                attrs: crate::Attrs::default(),
1580            }]),
1581            &crate::ScanConfig::default(),
1582            &mut |_| {},
1583        )
1584        .expect("unknown ancestry schedules reconciliation");
1585
1586        assert_eq!(report.apply.invalidated, 1);
1587        assert!(report.reconciliation.is_complete());
1588        assert_eq!(
1589            handle.kind(Path::new("new")).expect("new directory"),
1590            Some(crate::EntryKind::Dir)
1591        );
1592        assert_eq!(
1593            handle.kind(Path::new("new/deep/file.txt")).expect("nested file"),
1594            Some(crate::EntryKind::File)
1595        );
1596        assert_ne!(
1597            handle.attrs(Path::new("new")).expect("verified parent attrs"),
1598            Some(crate::Attrs::default())
1599        );
1600    }
1601
1602    #[test]
1603    fn create_intent_survives_coalescing_but_metadata_only_does_not_relist() {
1604        let root = Path::new("/watch-root");
1605        let path = root.join("directory");
1606        let mut pending = BTreeMap::new();
1607
1608        record(
1609            root,
1610            &notify::Event::new(EventKind::Create(CreateKind::Folder)).add_path(path.clone()),
1611            &mut pending,
1612            16,
1613            RenameReporting::EachSide,
1614        );
1615        record(
1616            root,
1617            &notify::Event::new(EventKind::Modify(ModifyKind::Metadata(MetadataKind::Any)))
1618                .add_path(path),
1619            &mut pending,
1620            16,
1621            RenameReporting::EachSide,
1622        );
1623        assert_eq!(
1624            pending.get(Path::new("directory")),
1625            Some(&Pending::Verify { relist_if_dir: true, renamed: false })
1626        );
1627
1628        let mut metadata_only = BTreeMap::new();
1629        record(
1630            root,
1631            &notify::Event::new(EventKind::Modify(ModifyKind::Metadata(MetadataKind::Any)))
1632                .add_path(root.join("existing")),
1633            &mut metadata_only,
1634            16,
1635            RenameReporting::EachSide,
1636        );
1637        assert_eq!(
1638            metadata_only.get(Path::new("existing")),
1639            Some(&Pending::Verify { relist_if_dir: false, renamed: false })
1640        );
1641    }
1642
1643    /// Record `events` into a fresh pending set under one backend rename policy.
1644    fn recorded(
1645        root: &Path,
1646        renames: RenameReporting,
1647        events: &[notify::Event],
1648    ) -> BTreeMap<PathBuf, Pending> {
1649        let mut pending = BTreeMap::new();
1650        for event in events {
1651            record(root, event, &mut pending, 16, renames);
1652        }
1653        pending
1654    }
1655
1656    fn rename_event(mode: RenameMode, paths: &[PathBuf]) -> notify::Event {
1657        let mut event = notify::Event::new(EventKind::Modify(ModifyKind::Name(mode)));
1658        event.paths = paths.to_vec();
1659        event
1660    }
1661
1662    /// One `FSEvents` `ItemRenamed` record, as notify 8.2 translates it.
1663    fn fsevents_rename(path: PathBuf) -> notify::Event {
1664        rename_event(RenameMode::Any, &[path])
1665    }
1666
1667    const RENAMED: Pending = Pending::Verify { relist_if_dir: false, renamed: true };
1668
1669    /// Every rename shape a reporting backend delivers verifies its own in-root paths and
1670    /// nothing else: `FSEvents` one side per event, inotify `From`/`To` plus the paired
1671    /// `Both`, and Windows `From`/`To`. A side outside the root contributes nothing.
1672    #[test]
1673    fn each_rename_shape_verifies_its_named_paths_without_escalating() {
1674        let root = Path::new("/watch-root");
1675        let (old, new) = (root.join("dir/old"), root.join("other/new"));
1676        for (backend, events) in [
1677            ("FSEvents", vec![fsevents_rename(old.clone()), fsevents_rename(new.clone())]),
1678            (
1679                "inotify",
1680                vec![
1681                    rename_event(RenameMode::From, std::slice::from_ref(&old)),
1682                    rename_event(RenameMode::To, std::slice::from_ref(&new)),
1683                    rename_event(RenameMode::Both, &[old.clone(), new.clone()]),
1684                ],
1685            ),
1686            (
1687                "Windows",
1688                vec![
1689                    rename_event(RenameMode::From, std::slice::from_ref(&old)),
1690                    rename_event(RenameMode::To, std::slice::from_ref(&new)),
1691                ],
1692            ),
1693        ] {
1694            let pending = recorded(root, RenameReporting::EachSide, &events);
1695            assert_eq!(
1696                pending,
1697                BTreeMap::from([
1698                    (PathBuf::from("dir/old"), RENAMED),
1699                    (PathBuf::from("other/new"), RENAMED),
1700                ]),
1701                "{backend}"
1702            );
1703        }
1704
1705        let outside = Path::new("/elsewhere/file");
1706        for (direction, event) in [
1707            ("move in", rename_event(RenameMode::Both, &[outside.to_path_buf(), new.clone()])),
1708            ("move out", rename_event(RenameMode::Both, &[old.clone(), outside.to_path_buf()])),
1709        ] {
1710            let pending = recorded(root, RenameReporting::EachSide, &[event]);
1711            assert_eq!(pending.len(), 1, "{direction}: {pending:?}");
1712            assert!(!pending.contains_key(Path::new("")), "{direction}: {pending:?}");
1713        }
1714    }
1715
1716    /// `FSEvents` keeps `ItemRenamed` on a path's later records, and notify splits one
1717    /// record into an event per flag. The rename fact merges into the path's pending
1718    /// verification with whatever else arrived, and still costs no root reconcile.
1719    #[test]
1720    fn a_sticky_rename_flag_merges_with_the_same_paths_other_events() {
1721        let root = Path::new("/watch-root");
1722        let path = root.join("state.json");
1723        let pending = recorded(
1724            root,
1725            RenameReporting::EachSide,
1726            &[
1727                notify::Event::new(EventKind::Create(CreateKind::File)).add_path(path.clone()),
1728                fsevents_rename(path.clone()),
1729                notify::Event::new(EventKind::Modify(ModifyKind::Data(
1730                    notify::event::DataChange::Content,
1731                )))
1732                .add_path(path),
1733            ],
1734        );
1735
1736        assert_eq!(
1737            pending,
1738            BTreeMap::from([(
1739                PathBuf::from("state.json"),
1740                Pending::Verify { relist_if_dir: true, renamed: true }
1741            )])
1742        );
1743    }
1744
1745    /// The cases a rename's own path cannot bound keep the whole-root reconcile: a backend
1746    /// that may never name the new side, a rename of the root itself, and a rename that
1747    /// names nowhere this root can place. Loss signals escalate exactly as before.
1748    #[test]
1749    fn unboundable_renames_and_ambiguous_rescans_escalate_the_root() {
1750        let root = Path::new("/watch-root");
1751        let escalated = Some(&Pending::Escalate(InvalidateReason::UnpairedRename));
1752
1753        let kqueue =
1754            recorded(root, RenameReporting::OldSideOnly, &[fsevents_rename(root.join("old"))]);
1755        assert_eq!(kqueue.get(Path::new("")), escalated);
1756
1757        let root_moved = recorded(
1758            root,
1759            RenameReporting::EachSide,
1760            &[rename_event(RenameMode::From, &[root.to_path_buf()])],
1761        );
1762        assert_eq!(root_moved.get(Path::new("")), escalated);
1763        assert_eq!(root_moved.len(), 1, "the root is escalated, never verified: {root_moved:?}");
1764
1765        for paths in [vec![], vec![PathBuf::from("/elsewhere/file")]] {
1766            let unplaced =
1767                recorded(root, RenameReporting::EachSide, &[rename_event(RenameMode::Any, &paths)]);
1768            assert_eq!(unplaced.get(Path::new("")), escalated, "{paths:?}");
1769        }
1770
1771        let rescan = recorded(
1772            root,
1773            RenameReporting::EachSide,
1774            &[notify::Event::new(EventKind::Any)
1775                .add_path(root.join("a"))
1776                .add_path(root.join("b"))
1777                .set_flag(Flag::Rescan)],
1778        );
1779        assert_eq!(
1780            rescan.get(Path::new("")),
1781            Some(&Pending::Escalate(InvalidateReason::WatchOverflow))
1782        );
1783    }
1784
1785    #[test]
1786    fn pending_path_overload_collapses_to_one_root_invalidation() {
1787        let root = Path::new("/watch-root");
1788        let mut pending = BTreeMap::new();
1789        for name in ["one", "two", "three"] {
1790            record(
1791                root,
1792                &notify::Event::new(EventKind::Any).add_path(root.join(name)),
1793                &mut pending,
1794                2,
1795                RenameReporting::EachSide,
1796            );
1797        }
1798
1799        assert_eq!(pending.len(), 1);
1800        assert_eq!(
1801            pending.get(Path::new("")),
1802            Some(&Pending::Escalate(InvalidateReason::WatchOverflow))
1803        );
1804    }
1805
1806    #[test]
1807    fn continuous_churn_has_a_deterministic_max_hold_ceiling() {
1808        let past = Instant::now()
1809            .checked_sub(Duration::from_secs(2))
1810            .expect("representable earlier instant");
1811        assert!(max_hold_elapsed(Some(past), Duration::from_secs(1)));
1812        assert!(!max_hold_elapsed(None, Duration::from_secs(1)));
1813    }
1814
1815    #[test]
1816    fn backend_enqueue_is_nonblocking_and_marks_overflow() {
1817        let (sender, receiver) = sync_channel(1);
1818        let overflowed = AtomicBool::new(false);
1819        enqueue_raw(&sender, &overflowed, Ok(notify::Event::new(EventKind::Any)));
1820        enqueue_raw(&sender, &overflowed, Ok(notify::Event::new(EventKind::Any)));
1821
1822        assert!(overflowed.load(Ordering::Acquire));
1823        assert!(matches!(receiver.try_recv(), Ok(RawMessage::Event(Ok(_)))));
1824    }
1825
1826    #[test]
1827    fn full_intent_queue_retains_a_sticky_root_invalidation() {
1828        let (sender, receiver) = sync_channel(1);
1829        sender.try_send(CoalescedIntent::default()).expect("fill output");
1830        let mut pending = BTreeMap::from([(
1831            PathBuf::from("lost.txt"),
1832            Pending::Verify { relist_if_dir: false, renamed: false },
1833        )]);
1834        let mut sticky_overflow = false;
1835
1836        try_deliver_pending(&mut pending, &sender, &mut sticky_overflow).expect("connected");
1837        assert!(pending.is_empty());
1838        assert!(sticky_overflow);
1839
1840        receiver.try_recv().expect("make output capacity");
1841        assert!(try_deliver_overflow(&sender).expect("connected"));
1842        let intent = receiver.try_recv().expect("sticky overflow intent");
1843        let observation = verify_intent(
1844            Path::new("/unused"),
1845            WatchConfig::default(),
1846            &intent,
1847            &ScanConfig::default(),
1848        );
1849        assert!(matches!(
1850            &observation.ops[0].op,
1851            Op::InvalidateSubtree {
1852                path,
1853                reason: InvalidateReason::WatchOverflow,
1854            } if path.as_os_str().is_empty()
1855        ));
1856    }
1857
1858    #[test]
1859    fn cancellation_wakes_and_joins_with_a_full_intent_queue() {
1860        let dir = tempfile::tempdir().expect("tempdir");
1861        let root = dir.path().canonicalize().expect("canonical root");
1862        let config = WatchConfig {
1863            settle: Duration::from_secs(30),
1864            max_hold: Duration::from_secs(30),
1865            event_capacity: 1,
1866            batch_path_capacity: 1,
1867            intent_capacity: 1,
1868            ..WatchConfig::default()
1869        };
1870        let (control, raw) = sync_channel(1);
1871        let (output, intents) = sync_channel(1);
1872        output.try_send(CoalescedIntent::default()).expect("fill intent queue");
1873        let cancelled = Arc::new(AtomicBool::new(false));
1874        let worker_cancelled = Arc::clone(&cancelled);
1875        let status = Arc::new(AtomicU8::new(WORKER_RUNNING));
1876        let tracked_status = Arc::clone(&status);
1877        let overflowed = Arc::new(AtomicBool::new(false));
1878        let worker_overflowed = Arc::clone(&overflowed);
1879        let worker_root = root.clone();
1880        let worker = std::thread::spawn(move || {
1881            run_tracked_worker(&tracked_status, || {
1882                run_worker(
1883                    &worker_root,
1884                    config,
1885                    RenameReporting::EachSide,
1886                    &raw,
1887                    &output,
1888                    &worker_overflowed,
1889                    &worker_cancelled,
1890                );
1891            });
1892        });
1893        let watcher = Watcher {
1894            root,
1895            config,
1896            inner: None,
1897            intents,
1898            control: Some(control),
1899            cancelled,
1900            worker_status: status,
1901            worker: Some(worker),
1902        };
1903        let (done_tx, done_rx) = sync_channel(1);
1904
1905        std::thread::spawn(move || {
1906            drop(watcher);
1907            done_tx.send(()).expect("report drop");
1908        });
1909
1910        done_rx
1911            .recv_timeout(Duration::from_secs(5))
1912            .expect("watcher drop must wake and join promptly");
1913    }
1914
1915    #[test]
1916    fn coalescing_defers_filesystem_verification_to_the_consumer() {
1917        let dir = tempfile::tempdir().expect("tempdir");
1918        let relative = PathBuf::from("appeared.txt");
1919        let intent = CoalescedIntent {
1920            pending: BTreeMap::from([(
1921                relative.clone(),
1922                Pending::Verify { relist_if_dir: false, renamed: false },
1923            )]),
1924        };
1925
1926        fs::write(dir.path().join(&relative), b"current").expect("create after coalescing");
1927        let observation =
1928            verify_intent(dir.path(), WatchConfig::default(), &intent, &ScanConfig::default());
1929
1930        assert!(matches!(
1931            &observation.ops[0].op,
1932            Op::Upsert { path, attrs, .. } if path == &relative && attrs.size == 7
1933        ));
1934    }
1935
1936    #[test]
1937    fn control_verification_emits_exact_source_with_the_entry_fact() {
1938        let dir = tempfile::tempdir().expect("tempdir");
1939        let relative = PathBuf::from(".gitignore");
1940        fs::write(dir.path().join(&relative), b"*.log\n").expect("write control");
1941        let intent = CoalescedIntent {
1942            pending: BTreeMap::from([(
1943                relative.clone(),
1944                Pending::Verify { relist_if_dir: false, renamed: false },
1945            )]),
1946        };
1947
1948        let config = ScanConfig { read_controls: true, ..ScanConfig::default() };
1949        let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
1950
1951        assert!(matches!(
1952            &observation.ops[0].op,
1953            Op::Upsert { path, kind: crate::EntryKind::File, .. } if path == &relative
1954        ));
1955        assert!(matches!(
1956            &observation.ops[1].op,
1957            Op::ControlUpsert { path, source } if path == &relative && source == b"*.log\n"
1958        ));
1959    }
1960
1961    #[test]
1962    fn only_population_control_event_reconciles_previously_absent_file() {
1963        let dir = tempfile::tempdir().expect("tree");
1964        fs::write(dir.path().join(".gitignore"), b"# no ignored files\n").expect("control");
1965        fs::write(dir.path().join("debug.log"), b"debug").expect("file");
1966        let scan =
1967            ScanConfig { population: crate::query::IgnoredEntries::Only, ..ScanConfig::default() };
1968        let (index, cold) = crate::scan::scan_into_index(dir.path(), &scan).expect("cold");
1969        assert!(cold.is_complete());
1970        assert!(index.lookup(Path::new("debug.log")).is_none());
1971        let handle = crate::IndexHandle::new(index);
1972        fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("rule edit");
1973        let intent = CoalescedIntent {
1974            pending: BTreeMap::from([(
1975                PathBuf::from(".gitignore"),
1976                Pending::Verify { relist_if_dir: false, renamed: false },
1977            )]),
1978        };
1979        let report =
1980            apply_intent(&handle, dir.path(), WatchConfig::default(), &intent, &scan, &mut |_| {})
1981                .expect("watch apply");
1982        assert!(report.reconciliation.is_complete());
1983        assert!(handle.kind(Path::new("debug.log")).expect("lookup").is_some());
1984    }
1985
1986    /// A watch maintains exactly the control state its scan policy claims.
1987    ///
1988    /// A controls-off scan retains no control table and stamps that into its scope, and a
1989    /// watch with the same policy accepts the scope as its own. Verification that read
1990    /// control files regardless would grow a partial rule set -- only the sources some
1991    /// event happened to touch -- on an index whose scope says it has none, and the next
1992    /// save would persist it under that scope (fdu-ajsu).
1993    #[test]
1994    fn verification_observes_no_control_state_under_a_controls_off_policy() {
1995        let dir = tempfile::tempdir().expect("tempdir");
1996        fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("write control");
1997        fs::write(dir.path().join("debug.log"), b"x").expect("write file");
1998        let config = ScanConfig { read_controls: false, ..ScanConfig::default() };
1999        let (mut index, _) = crate::scan::scan_into_index(dir.path(), &config).expect("scan");
2000        assert!(index.control_table().is_empty());
2001        assert!(matches!(index.controls(), Err(crate::Error::ControlStateNotObserved)));
2002        let control = PathBuf::from(".gitignore");
2003        let intent = CoalescedIntent {
2004            pending: BTreeMap::from([(
2005                control.clone(),
2006                Pending::Verify { relist_if_dir: false, renamed: false },
2007            )]),
2008        };
2009
2010        let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
2011        let reverified = reverify_observation(
2012            dir.path(),
2013            &Observation::new(vec![Op::Remove { path: control }]),
2014            &config,
2015        )
2016        .expect("reverify");
2017
2018        for verified in [&observation, &reverified] {
2019            assert!(
2020                !verified.ops.iter().any(|observed| matches!(
2021                    observed.op,
2022                    Op::ControlUpsert { .. } | Op::ControlRemove { .. }
2023                )),
2024                "a controls-off policy observed control state: {:?}",
2025                verified.ops
2026            );
2027        }
2028        index.apply(&observation).expect("apply the verified observation");
2029        assert!(index.control_table().is_empty());
2030        assert_eq!(index.scope(), config.scope());
2031    }
2032
2033    #[cfg(unix)]
2034    #[test]
2035    fn applying_verification_uses_the_index_admission_scope() {
2036        use std::os::unix::net::UnixListener;
2037
2038        let dir = tempfile::tempdir().expect("tempdir");
2039        fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("write control");
2040        fs::write(dir.path().join(".secret"), b"hidden").expect("write hidden");
2041        let _listener = UnixListener::bind(dir.path().join("service.sock")).expect("bind socket");
2042        let intent = CoalescedIntent {
2043            pending: [".gitignore", ".secret", "service.sock"]
2044                .into_iter()
2045                .map(|path| {
2046                    (PathBuf::from(path), Pending::Verify { relist_if_dir: false, renamed: false })
2047                })
2048                .collect(),
2049        };
2050        let config = ScanConfig {
2051            hidden: Some(Arc::new(crate::HiddenPolicy::prune_hidden::<[&str; 0], &str>([]))),
2052            exclude_special: true,
2053            read_controls: true,
2054            ..ScanConfig::default()
2055        };
2056
2057        let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
2058
2059        assert!(observation.ops.iter().any(|observed| matches!(
2060            &observed.op,
2061            Op::ControlUpsert { path, source }
2062                if path == Path::new(".gitignore") && source == b"*.log\n"
2063        )));
2064        for path in [".gitignore", ".secret", "service.sock"] {
2065            assert!(observation.ops.iter().any(|observed| matches!(
2066                &observed.op,
2067                Op::Remove { path: removed } if removed == Path::new(path)
2068            )));
2069            assert!(!observation.ops.iter().any(|observed| matches!(
2070                &observed.op,
2071                Op::Upsert { path: retained, .. } if retained == Path::new(path)
2072            )));
2073        }
2074    }
2075
2076    #[test]
2077    fn timeout_stop_and_worker_panic_are_distinct() {
2078        let dir = tempfile::tempdir().expect("tempdir");
2079        let (live_sender, live) =
2080            queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2081        assert!(live.next_observation(Duration::ZERO).expect("timeout").is_none());
2082        drop(live_sender);
2083
2084        let (stopped_sender, stopped) =
2085            queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2086        stopped.worker_status.store(WORKER_STOPPED, Ordering::Release);
2087        drop(stopped_sender);
2088        assert!(matches!(stopped.next_observation(Duration::ZERO), Err(Error::WatchStopped)));
2089
2090        let (panicked_sender, panicked) =
2091            queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2092        panicked.worker_status.store(WORKER_PANICKED, Ordering::Release);
2093        drop(panicked_sender);
2094        assert!(matches!(
2095            panicked.next_observation(Duration::ZERO),
2096            Err(Error::WatchWorkerPanicked)
2097        ));
2098    }
2099
2100    #[test]
2101    fn tracked_worker_records_a_panic() {
2102        let status = AtomicU8::new(WORKER_RUNNING);
2103        run_tracked_worker(&status, || panic!("injected worker panic"));
2104        assert_eq!(status.load(Ordering::Acquire), WORKER_PANICKED);
2105    }
2106
2107    #[test]
2108    fn observation_driver_closes_the_invalidation_loop() {
2109        let dir = tempfile::tempdir().expect("tempdir");
2110        let (index, _) =
2111            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2112        let handle = crate::IndexHandle::new(index);
2113        fs::write(dir.path().join("raced.txt"), b"raced").expect("write");
2114        let observation = Observation::new(vec![Op::InvalidateSubtree {
2115            path: PathBuf::new(),
2116            reason: InvalidateReason::WatchSetupRace,
2117        }]);
2118
2119        apply_observation(&handle, &observation, &crate::ScanConfig::default(), &mut |_| {})
2120            .expect("apply and reconcile");
2121
2122        assert!(handle.kind(Path::new("raced.txt")).expect("query").is_some());
2123    }
2124
2125    #[test]
2126    fn applying_driver_reverifies_a_queued_sample_after_reconciliation() {
2127        let dir = tempfile::tempdir().expect("tempdir");
2128        let path = dir.path().join("sample.txt");
2129        fs::write(&path, b"old").expect("write old sample");
2130        let (index, _) =
2131            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2132        let old_attrs = *index.attrs(Path::new("sample.txt")).expect("sample attributes");
2133        let delayed = Observation::new(vec![Op::Upsert {
2134            path: PathBuf::from("sample.txt"),
2135            kind: crate::EntryKind::File,
2136            attrs: old_attrs,
2137        }]);
2138        let handle = crate::IndexHandle::new(index);
2139
2140        fs::write(&path, b"new contents").expect("write current sample");
2141        crate::scan::reconcile_handle(&handle, &crate::ScanConfig::default(), &mut |_| {})
2142            .expect("reconcile newer sample");
2143        let current_size = fs::metadata(&path).expect("sample metadata").len();
2144
2145        apply_observation(&handle, &delayed, &crate::ScanConfig::default(), &mut |_| {})
2146            .expect("apply delayed watch sample");
2147
2148        assert_eq!(
2149            handle.attrs(Path::new("sample.txt")).expect("query").expect("sample remains").size,
2150            current_size
2151        );
2152    }
2153
2154    #[test]
2155    fn blocked_verifier_holds_no_index_lock_and_commits_only_at_current_clock() {
2156        let dir = tempfile::tempdir().expect("tempdir");
2157        let (index, _) =
2158            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2159        let handle = crate::IndexHandle::new(index);
2160        let queued = Observation::new(vec![Op::Upsert {
2161            path: PathBuf::from("queued.txt"),
2162            kind: crate::EntryKind::File,
2163            attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2164        }]);
2165        let applying = handle.clone();
2166        let (entered_tx, entered_rx) = sync_channel(1);
2167        let (release_tx, release_rx) = sync_channel(1);
2168        let (done_tx, done_rx) = sync_channel(1);
2169        let apply_thread = std::thread::spawn(move || {
2170            let mut first = true;
2171            let mut verifier = |_: &Path, observation: &Observation| {
2172                if first {
2173                    first = false;
2174                    entered_tx.send(()).expect("signal blocked verifier");
2175                    release_rx.recv().expect("release blocked verifier");
2176                }
2177                Ok(observation.clone())
2178            };
2179            let result = apply_reverified_with(
2180                &applying,
2181                &queued,
2182                &crate::ScanConfig::default(),
2183                &mut verifier,
2184            );
2185            done_tx.send(result).expect("report apply result");
2186        });
2187
2188        entered_rx
2189            .recv_timeout(Duration::from_secs(5))
2190            .expect("verifier must reach the injected block");
2191        let progressing = handle.clone();
2192        let (progress_tx, progress_rx) = sync_channel(1);
2193        let progress_thread = std::thread::spawn(move || {
2194            let total = progressing.total().expect("reader progresses");
2195            let write = progressing.apply(&Observation::new(vec![Op::Upsert {
2196                path: PathBuf::from("competitor.txt"),
2197                kind: crate::EntryKind::File,
2198                attrs: crate::Attrs { size: 3, allocated: 3, ..crate::Attrs::default() },
2199            }]));
2200            progress_tx.send((total, write)).expect("report progress");
2201        });
2202        let (_, competing_write) = progress_rx
2203            .recv_timeout(Duration::from_secs(5))
2204            .expect("reader and writer must progress while verification is blocked");
2205        competing_write.expect("competing write");
2206        release_tx.send(()).expect("release verifier");
2207
2208        let outcome = done_rx
2209            .recv_timeout(Duration::from_secs(5))
2210            .expect("applying driver completes")
2211            .expect("applying driver succeeds");
2212        apply_thread.join().expect("apply thread");
2213        progress_thread.join().expect("progress thread");
2214        assert_eq!(outcome.inserted, 1);
2215        assert!(handle.kind(Path::new("competitor.txt")).expect("query").is_some());
2216        assert!(handle.kind(Path::new("queued.txt")).expect("query").is_some());
2217        assert_eq!(handle.clock().expect("clock"), crate::Clock(2));
2218    }
2219
2220    #[test]
2221    fn exhausted_watch_contention_stays_unfresh_until_reconciliation() {
2222        let dir = tempfile::tempdir().expect("tempdir");
2223        let (index, _) =
2224            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2225        let handle = crate::IndexHandle::new(index);
2226        let queued = Observation::new(vec![Op::Upsert {
2227            path: PathBuf::from("never-committed.txt"),
2228            kind: crate::EntryKind::File,
2229            attrs: crate::Attrs { size: 1, allocated: 1, ..crate::Attrs::default() },
2230        }]);
2231        let mut attempts = 0_usize;
2232        let mut verifier = |_: &Path, observation: &Observation| {
2233            attempts += 1;
2234            handle
2235                .apply(&Observation::new(vec![Op::Upsert {
2236                    path: PathBuf::from(format!("competitor-{attempts}.txt")),
2237                    kind: crate::EntryKind::File,
2238                    attrs: crate::Attrs {
2239                        size: attempts as u64,
2240                        allocated: attempts as u64,
2241                        ..crate::Attrs::default()
2242                    },
2243                }]))
2244                .expect("force a clock conflict");
2245            Ok(observation.clone())
2246        };
2247
2248        let outcome =
2249            apply_reverified_with(&handle, &queued, &crate::ScanConfig::default(), &mut verifier)
2250                .expect("contention escalates");
2251
2252        assert_eq!(attempts, MAX_OPTIMISTIC_APPLY_ATTEMPTS);
2253        assert_eq!(outcome.invalidated, 1);
2254        assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Stale);
2255        assert!(handle.kind(Path::new("never-committed.txt")).expect("query").is_none());
2256        let pending = handle.take_pending_invalidations().expect("pending invalidation");
2257        assert_eq!(pending, vec![(PathBuf::new(), InvalidateReason::WatchContention)]);
2258        handle.restore_pending_invalidations(pending).expect("restore invalidation");
2259
2260        crate::scan::reconcile_pending_handle(&handle, &crate::ScanConfig::default(), &mut |_| {})
2261            .expect("reconcile contention");
2262        assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Fresh);
2263    }
2264
2265    #[test]
2266    fn verifier_error_mutates_no_shared_state() {
2267        let dir = tempfile::tempdir().expect("tempdir");
2268        let (index, _) =
2269            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2270        let handle = crate::IndexHandle::new(index);
2271        let before_clock = handle.clock().expect("clock");
2272        let before_total = handle.total().expect("total");
2273        let mut verifier = |_: &Path, _: &Observation| {
2274            Err(Error::io(
2275                PathBuf::from("blocked"),
2276                std::io::Error::new(std::io::ErrorKind::PermissionDenied, "injected"),
2277            ))
2278        };
2279
2280        let error = apply_reverified_with(
2281            &handle,
2282            &Observation::default(),
2283            &crate::ScanConfig::default(),
2284            &mut verifier,
2285        )
2286        .expect_err("verification error");
2287
2288        assert!(matches!(error, Error::Io { .. }));
2289        assert_eq!(handle.clock().expect("clock"), before_clock);
2290        assert_eq!(handle.total().expect("total"), before_total);
2291        assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Fresh);
2292        assert!(handle.take_pending_invalidations().expect("pending").is_empty());
2293    }
2294
2295    #[test]
2296    fn stable_watch_arbitration_verifies_exactly_once() {
2297        let dir = tempfile::tempdir().expect("tempdir");
2298        let (index, _) =
2299            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2300        let handle = crate::IndexHandle::new(index);
2301        let mut calls = 0_u8;
2302        let mut verifier = |_: &Path, observation: &Observation| {
2303            calls += 1;
2304            Ok(observation.clone())
2305        };
2306
2307        apply_reverified_with(
2308            &handle,
2309            &Observation::new(vec![Op::InvalidateSubtree {
2310                path: PathBuf::new(),
2311                reason: InvalidateReason::Requested,
2312            }]),
2313            &crate::ScanConfig::default(),
2314            &mut verifier,
2315        )
2316        .expect("stable apply");
2317
2318        assert_eq!(calls, 1);
2319    }
2320
2321    #[test]
2322    fn disappearing_watch_root_escalates_instead_of_removing_the_index_root() {
2323        let error = std::io::Error::new(std::io::ErrorKind::NotFound, "root disappeared");
2324
2325        assert!(matches!(
2326            op_for_stat_error(PathBuf::new(), &error),
2327            Op::InvalidateSubtree {
2328                path,
2329                reason: InvalidateReason::VerificationFailed,
2330            } if path.as_os_str().is_empty()
2331        ));
2332    }
2333
2334    #[test]
2335    fn observation_driver_rejects_scope_mismatch_before_apply() {
2336        let dir = tempfile::tempdir().expect("tempdir");
2337        let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2338        let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2339        let handle = crate::IndexHandle::new(index);
2340        let observation = Observation::new(vec![Op::Upsert {
2341            path: PathBuf::from("deep/nested.txt"),
2342            kind: crate::EntryKind::File,
2343            attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2344        }]);
2345
2346        let error =
2347            apply_observation(&handle, &observation, &crate::ScanConfig::default(), &mut |_| {})
2348                .expect_err("mismatched scope must fail");
2349
2350        assert!(matches!(error, Error::ScanScopeMismatch { .. }));
2351        assert!(handle.kind(Path::new("deep/nested.txt")).expect("query").is_none());
2352    }
2353
2354    #[test]
2355    fn observation_driver_rejects_restricted_scopes_until_events_are_filtered() {
2356        let dir = tempfile::tempdir().expect("tempdir");
2357        let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2358        let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2359        let handle = crate::IndexHandle::new(index);
2360        let observation = Observation::new(vec![Op::Upsert {
2361            path: PathBuf::from("deep/nested.txt"),
2362            kind: crate::EntryKind::File,
2363            attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2364        }]);
2365
2366        let error = apply_observation(&handle, &observation, &shallow, &mut |_| {})
2367            .expect_err("unfiltered bounded watch scope must fail");
2368
2369        assert!(matches!(error, Error::UnsupportedScanConfig(_)));
2370        assert!(handle.kind(Path::new("deep/nested.txt")).expect("query").is_none());
2371    }
2372
2373    #[test]
2374    fn apply_next_rejects_restricted_scope_without_consuming_an_observation() {
2375        let dir = tempfile::tempdir().expect("tempdir");
2376        let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2377        let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2378        let handle = crate::IndexHandle::new(index);
2379        let (sender, watcher) =
2380            queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2381        sender.try_send(CoalescedIntent::default()).expect("queue intent");
2382
2383        let error = watcher
2384            .apply_next(&handle, &shallow, Duration::ZERO, &mut |_| {})
2385            .expect_err("restricted scope must fail before receive");
2386
2387        assert!(matches!(error, Error::UnsupportedScanConfig(_)));
2388        assert!(watcher.next_observation(Duration::ZERO).expect("receive").is_some());
2389    }
2390
2391    #[test]
2392    fn apply_next_rejects_a_watcher_for_another_root_without_consuming() {
2393        let indexed = tempfile::tempdir().expect("indexed root");
2394        let watched_root_dir = tempfile::tempdir().expect("watched root");
2395        let (index, _) =
2396            crate::scan::scan_into_index(indexed.path(), &crate::ScanConfig::default())
2397                .expect("scan indexed root");
2398        let handle = crate::IndexHandle::new(index);
2399        let (sender, watcher) = queued_test_watcher(
2400            watched_root_dir.path().canonicalize().expect("canonical watched root"),
2401        );
2402        sender.try_send(CoalescedIntent::default()).expect("queue intent");
2403
2404        let error = watcher
2405            .apply_next(&handle, &crate::ScanConfig::default(), Duration::ZERO, &mut |_| {})
2406            .expect_err("mismatched root must fail");
2407
2408        assert!(matches!(error, Error::WatchRootMismatch { .. }));
2409        assert!(watcher.next_observation(Duration::ZERO).expect("receive").is_some());
2410    }
2411
2412    /// A gap over an unreadable directory is walked once by the per-event driver, which
2413    /// then goes quiet and keeps the cause.
2414    ///
2415    /// `fdu --watch` and the Python watch session drain the invalidation queue after every
2416    /// event through [`Watcher::apply_next`]. While an incomplete reconciliation was queued
2417    /// again, each unrelated event re-walked the same unreadable subtree -- for a root
2418    /// escalation, a full-tree walk per event, for the life of the session.
2419    #[cfg(unix)]
2420    #[test]
2421    fn apply_next_walks_an_unreadable_gap_once_and_retains_its_cause() {
2422        use std::os::unix::fs::PermissionsExt;
2423
2424        fn walks_of(commits: &[Commit], path: &Path) -> usize {
2425            commits
2426                .iter()
2427                .flat_map(|commit| commit.state.iter())
2428                .filter(|transition| {
2429                    matches!(
2430                        transition,
2431                        crate::StateTransition::Freshness { path: marked, current, .. }
2432                            if marked == path && *current == crate::Freshness::Reconciling
2433                    )
2434                })
2435                .count()
2436        }
2437
2438        if !crate::test_support::require_permission_bits() {
2439            return;
2440        }
2441        let dir = tempfile::tempdir().expect("tempdir");
2442        let root = dir.path().canonicalize().expect("canonical root");
2443        let blocked = root.join("blocked");
2444        fs::create_dir(&blocked).expect("blocked");
2445        fs::write(blocked.join("secret"), b"s").expect("fixture");
2446        let (index, _) =
2447            crate::scan::scan_into_index(&root, &crate::ScanConfig::default()).expect("scan");
2448        let handle = crate::IndexHandle::new(index);
2449        let (sender, watcher) = queued_test_watcher(root.clone());
2450        let config = crate::ScanConfig::default();
2451        let mut commits = Vec::new();
2452
2453        fs::set_permissions(&blocked, fs::Permissions::from_mode(0o000)).expect("deny reads");
2454        let mut gap = CoalescedIntent::default();
2455        gap.pending
2456            .insert(PathBuf::from("blocked"), Pending::Escalate(InvalidateReason::WatchOverflow));
2457        sender.try_send(gap).expect("queue the gap");
2458        let first = watcher
2459            .apply_next(&handle, &config, Duration::ZERO, &mut |commit| {
2460                commits.push(commit.clone());
2461            })
2462            .expect("apply the gap")
2463            .expect("an intent was queued");
2464        for name in ["live.txt", "marker.txt"] {
2465            fs::write(root.join(name), name).expect("unrelated mutation");
2466            let mut event = CoalescedIntent::default();
2467            event.pending.insert(
2468                PathBuf::from(name),
2469                Pending::Verify { relist_if_dir: false, renamed: false },
2470            );
2471            sender.try_send(event).expect("queue the event");
2472            watcher
2473                .apply_next(&handle, &config, Duration::ZERO, &mut |commit| {
2474                    commits.push(commit.clone());
2475                })
2476                .expect("apply the event")
2477                .expect("an intent was queued");
2478        }
2479
2480        let pending = handle.take_pending_invalidations().expect("pending");
2481        let freshness = handle.freshness_at(Path::new("blocked")).expect("freshness");
2482        let issues = handle.issues().expect("issues");
2483        fs::set_permissions(&blocked, fs::Permissions::from_mode(0o700)).expect("restore reads");
2484        assert!(!first.reconciliation.is_complete(), "the gap must be unreadable");
2485
2486        assert_eq!(
2487            walks_of(&commits, Path::new("blocked")),
2488            1,
2489            "an unreadable subtree must not be re-walked per unrelated event"
2490        );
2491        assert!(pending.is_empty(), "the queue must settle: {pending:?}");
2492        assert_eq!(freshness, crate::Freshness::Partial);
2493        assert!(handle.kind(Path::new("marker.txt")).expect("lookup").is_some());
2494        assert!(
2495            issues.iter().any(|issue| issue.kind == crate::IssueKind::Permission
2496                && issue.path.as_deref() == Some(Path::new("blocked"))),
2497            "{issues:?}"
2498        );
2499    }
2500
2501    /// A watched tree whose rename events are applied as the production driver would.
2502    struct RenameFixture {
2503        _dir: tempfile::TempDir,
2504        root: PathBuf,
2505        handle: IndexHandle,
2506    }
2507
2508    impl RenameFixture {
2509        /// Build `files` (with their parent directories), then index the tree cold.
2510        fn new(files: &[&str]) -> Self {
2511            let dir = tempfile::tempdir().expect("tempdir");
2512            let root = dir.path().canonicalize().expect("canonical root");
2513            for file in files {
2514                let path = root.join(file);
2515                fs::create_dir_all(path.parent().expect("parent")).expect("parents");
2516                fs::write(&path, file.as_bytes()).expect("fixture file");
2517            }
2518            let (index, _) =
2519                crate::scan::scan_into_index(&root, &ScanConfig::default()).expect("cold scan");
2520            Self { _dir: dir, root, handle: IndexHandle::new(index) }
2521        }
2522
2523        fn path(&self, rel: &str) -> PathBuf {
2524            self.root.join(rel)
2525        }
2526
2527        /// Apply `events` as one coalesced intent from a backend that reports each side,
2528        /// returning every invalidation the intent and its reconciliation committed.
2529        fn apply(&self, events: &[notify::Event]) -> Vec<(PathBuf, InvalidateReason)> {
2530            let intent = CoalescedIntent {
2531                pending: recorded(&self.root, RenameReporting::EachSide, events),
2532            };
2533            let mut commits = Vec::new();
2534            let report = apply_intent(
2535                &self.handle,
2536                &self.root,
2537                WatchConfig::default(),
2538                &intent,
2539                &ScanConfig::default(),
2540                &mut |commit| commits.push(commit.clone()),
2541            )
2542            .expect("apply the intent");
2543            assert!(report.reconciliation.is_complete(), "reconciliation must settle");
2544            commits
2545                .iter()
2546                .flat_map(|commit| commit.changes.iter())
2547                .filter_map(|change| match change {
2548                    crate::EffectiveChange::Invalidated { path, reason } => {
2549                        Some((path.clone(), *reason))
2550                    }
2551                    _ => None,
2552                })
2553                .collect()
2554        }
2555
2556        /// The watched index holds exactly what a cold scan of the tree finds now.
2557        fn assert_converged(&self) {
2558            let (cold, _) =
2559                crate::scan::scan_into_index(&self.root, &ScanConfig::default()).expect("cold");
2560            let watched = self.handle.read_with(entries).expect("read the watched index");
2561            assert_eq!(watched, entries(&cold), "watched index diverged from a cold scan");
2562            assert_eq!(self.handle.freshness().expect("freshness"), crate::Freshness::Fresh);
2563        }
2564    }
2565
2566    /// Every entry with its kind and, for a non-directory, its size. Directory metadata
2567    /// is left out: a directory's own stat changes with its listing, and no backend
2568    /// reports that change for the directory itself.
2569    fn entries(index: &crate::Index) -> BTreeMap<PathBuf, (crate::EntryKind, Option<u64>)> {
2570        fn walk(
2571            index: &crate::Index,
2572            dir: &Path,
2573            out: &mut BTreeMap<PathBuf, (crate::EntryKind, Option<u64>)>,
2574        ) {
2575            let Some(children) = index.children(dir) else {
2576                return;
2577            };
2578            let children: Vec<_> =
2579                children.map(|(name, id)| (dir.join(name), id)).collect::<Vec<_>>();
2580            for (path, id) in children {
2581                let kind = index.kind_of(id).expect("live child");
2582                let size = (!kind.is_dir()).then(|| index.attrs_of(id).expect("attrs").size);
2583                out.insert(path.clone(), (kind, size));
2584                walk(index, &path, out);
2585            }
2586        }
2587        let mut out = BTreeMap::new();
2588        walk(index, Path::new(""), &mut out);
2589        out
2590    }
2591
2592    /// Every entry's ignored classification, and every retained control source by its
2593    /// canonical path.
2594    type Classified = (BTreeMap<PathBuf, Option<bool>>, Vec<(PathBuf, Vec<u8>)>);
2595
2596    fn classified(index: &crate::Index) -> Classified {
2597        let ignored = entries(index)
2598            .into_keys()
2599            .map(|path| {
2600                let ignored = index.is_ignored(&path).expect("observed");
2601                (path, ignored)
2602            })
2603            .collect();
2604        let sources = index
2605            .controls()
2606            .expect("observed")
2607            .sources()
2608            .map(|(path, source)| (path, source.to_vec()))
2609            .collect();
2610        (ignored, sources)
2611    }
2612
2613    impl RenameFixture {
2614        /// The watched index classifies every entry as a cold scan of the tree does now, from
2615        /// the same control sources.
2616        fn assert_classified_as_cold(&self, label: &str) -> Classified {
2617            let (cold, _) =
2618                crate::scan::scan_into_index(&self.root, &ScanConfig::default()).expect("cold");
2619            let watched = self.handle.read_with(classified).expect("read the watched index");
2620            assert_eq!(watched, classified(&cold), "{label}: the watch diverged from a cold scan");
2621            watched
2622        }
2623    }
2624
2625    /// A watch follows a control file spelled `.GITIGNORE` through its creation, an edit,
2626    /// and its removal, classifying as a cold scan does after each: its rules govern
2627    /// exactly where a lookup of `.gitignore` resolves to it (fdu-0w1b). Its removal is
2628    /// the case that cannot be verified by a stat, so the watch looks the directory's
2629    /// control up instead. A case-only rename in either direction converges too, under
2630    /// the host's own name resolution (folded lookups change only the control lookup, and
2631    /// a rename's old spelling must stat as the host resolves it).
2632    #[test]
2633    fn a_watch_follows_a_case_variant_control_file() {
2634        use crate::test_support::CaseLookups;
2635
2636        let probe = tempfile::tempdir().expect("tempdir");
2637        for (lookups, governs) in CaseLookups::on_this_host(probe.path()) {
2638            let tree = RenameFixture::new(&["up/x.tmp", "up/notes.txt"]);
2639            let _lookups = lookups.install(&tree.root);
2640            let variant = tree.path("up/.GITIGNORE");
2641            let x_ignored = |classified: &Classified| classified.0[Path::new("up/x.tmp")];
2642
2643            fs::write(&variant, b"*.tmp\n").expect("create the variant");
2644            tree.apply(&[created(variant.clone())]);
2645            let appeared = tree.assert_classified_as_cold(&format!("{lookups:?}: created"));
2646            assert_eq!(x_ignored(&appeared), Some(governs), "{lookups:?}");
2647
2648            fs::write(&variant, b"*.txt\n").expect("edit the variant");
2649            tree.apply(&[modified(variant.clone())]);
2650            let edited = tree.assert_classified_as_cold(&format!("{lookups:?}: edited"));
2651            assert_eq!(x_ignored(&edited), Some(false), "{lookups:?}");
2652            assert_eq!(edited.0[Path::new("up/notes.txt")], Some(governs), "{lookups:?}");
2653
2654            fs::remove_file(&variant).expect("remove the variant");
2655            tree.apply(&[notify::Event::new(EventKind::Remove(notify::event::RemoveKind::File))
2656                .add_path(variant.clone())]);
2657            let removed = tree.assert_classified_as_cold(&format!("{lookups:?}: removed"));
2658            assert_eq!(removed.1, Vec::new(), "{lookups:?}: no rules remain");
2659
2660            if lookups == CaseLookups::Host {
2661                let exact = tree.path("up/.gitignore");
2662                fs::write(&exact, b"*.tmp\n").expect("create the exact name");
2663                tree.apply(&[created(exact.clone())]);
2664                for (from, to) in [(&exact, &variant), (&variant, &exact)] {
2665                    fs::rename(from, to).expect("case-only rename");
2666                    tree.apply(&[fsevents_rename(from.clone()), fsevents_rename(to.clone())]);
2667                    tree.assert_converged();
2668                    tree.assert_classified_as_cold(&format!("{lookups:?}: renamed to {to:?}"));
2669                }
2670            }
2671        }
2672    }
2673
2674    fn created(path: PathBuf) -> notify::Event {
2675        notify::Event::new(EventKind::Create(CreateKind::Any)).add_path(path)
2676    }
2677
2678    fn modified(path: PathBuf) -> notify::Event {
2679        notify::Event::new(EventKind::Modify(ModifyKind::Data(notify::event::DataChange::Content)))
2680            .add_path(path)
2681    }
2682
2683    /// The atomic-save pattern behind fdu-822y: write a temporary name, rename it over the
2684    /// real one. `FSEvents` names both paths with sticky create and modify flags.
2685    #[test]
2686    fn a_file_renamed_within_its_directory_needs_no_invalidation() {
2687        let tree = RenameFixture::new(&["state/config.json", "state/other.txt"]);
2688        fs::write(tree.path("state/config.json.tmp"), b"new configuration").expect("temp");
2689        fs::rename(tree.path("state/config.json.tmp"), tree.path("state/config.json"))
2690            .expect("rename over");
2691
2692        let invalidations = tree.apply(&[
2693            created(tree.path("state/config.json.tmp")),
2694            modified(tree.path("state/config.json.tmp")),
2695            fsevents_rename(tree.path("state/config.json.tmp")),
2696            fsevents_rename(tree.path("state/config.json")),
2697        ]);
2698
2699        assert_eq!(invalidations, vec![], "a file rename is settled by its own paths");
2700        tree.assert_converged();
2701    }
2702
2703    /// A move between directories arrives as two one-sided events, which may land in
2704    /// different batches. Each side settles on its own, in either order.
2705    #[test]
2706    fn a_file_moved_across_directories_settles_one_side_per_event() {
2707        let tree = RenameFixture::new(&["from/moved.txt", "to/resident.txt"]);
2708        fs::rename(tree.path("from/moved.txt"), tree.path("to/moved.txt")).expect("move");
2709
2710        assert_eq!(tree.apply(&[fsevents_rename(tree.path("to/moved.txt"))]), vec![]);
2711        assert_eq!(tree.apply(&[fsevents_rename(tree.path("from/moved.txt"))]), vec![]);
2712        tree.assert_converged();
2713    }
2714
2715    /// A renamed directory's contents produce no events, so the new name is relisted; the
2716    /// old name's subtree is removed. The reconcile is bounded by the moved directory.
2717    #[test]
2718    fn a_renamed_directory_relists_only_its_own_subtree() {
2719        let tree = RenameFixture::new(&[
2720            "project/old/one.txt",
2721            "project/old/nested/two.txt",
2722            "project/untouched/three.txt",
2723        ]);
2724        fs::rename(tree.path("project/old"), tree.path("project/new")).expect("rename directory");
2725
2726        let invalidations = tree.apply(&[
2727            fsevents_rename(tree.path("project/old")),
2728            fsevents_rename(tree.path("project/new")),
2729        ]);
2730
2731        assert_eq!(
2732            invalidations,
2733            vec![(PathBuf::from("project/new"), InvalidateReason::UnpairedRename)]
2734        );
2735        tree.assert_converged();
2736    }
2737
2738    /// Moves across the watch boundary name one side only, and that side is enough: an
2739    /// arriving tree is relisted, a departing one is removed with its subtree.
2740    #[test]
2741    fn moves_across_the_root_boundary_settle_from_the_inside_side() {
2742        let tree = RenameFixture::new(&["resident.txt", "leaving/a.txt", "leaving/deep/b.txt"]);
2743        let outside = tempfile::tempdir().expect("outside the root");
2744        fs::create_dir_all(outside.path().join("arriving/deep")).expect("outside tree");
2745        fs::write(outside.path().join("arriving/deep/c.txt"), b"arrived").expect("outside file");
2746
2747        fs::rename(outside.path().join("arriving"), tree.path("arrived")).expect("move in");
2748        let invalidations = tree.apply(&[fsevents_rename(tree.path("arrived"))]);
2749        assert_eq!(
2750            invalidations,
2751            vec![(PathBuf::from("arrived"), InvalidateReason::UnpairedRename)]
2752        );
2753        tree.assert_converged();
2754
2755        fs::rename(tree.path("leaving"), outside.path().join("left")).expect("move out");
2756        assert_eq!(tree.apply(&[fsevents_rename(tree.path("leaving"))]), vec![]);
2757        tree.assert_converged();
2758    }
2759
2760    /// A name reused after its entry was renamed away is verified as whatever now sits
2761    /// there. A reused directory name is relisted, which drops the departed children.
2762    #[test]
2763    fn a_name_reused_after_a_rename_verifies_as_its_new_entry() {
2764        let tree = RenameFixture::new(&["log.txt", "cache/entry.bin"]);
2765        fs::rename(tree.path("log.txt"), tree.path("log.1.txt")).expect("rotate");
2766        fs::write(tree.path("log.txt"), b"fresh log, longer than the old one").expect("reuse");
2767        fs::rename(tree.path("cache"), tree.path("cache.old")).expect("retire directory");
2768        fs::create_dir(tree.path("cache")).expect("reuse directory name");
2769
2770        let invalidations = tree.apply(&[
2771            fsevents_rename(tree.path("log.txt")),
2772            created(tree.path("log.txt")),
2773            fsevents_rename(tree.path("log.1.txt")),
2774            fsevents_rename(tree.path("cache")),
2775            created(tree.path("cache")),
2776            fsevents_rename(tree.path("cache.old")),
2777        ]);
2778
2779        assert!(
2780            invalidations.iter().all(|(path, _)| !path.as_os_str().is_empty()),
2781            "no root reconcile: {invalidations:?}"
2782        );
2783        tree.assert_converged();
2784        assert_eq!(tree.handle.kind(Path::new("cache/entry.bin")).expect("lookup"), None);
2785    }
2786
2787    /// `FSEvents` keeps `ItemRenamed` on a file's later records. A plain in-place write
2788    /// then arrives flagged as a rename, and it must cost what a write costs.
2789    #[test]
2790    fn a_sticky_rename_flag_on_a_later_write_is_just_a_write() {
2791        let tree = RenameFixture::new(&["sessions/today.jsonl"]);
2792        fs::write(tree.path("sessions/today.jsonl"), b"appended record after the rename")
2793            .expect("write in place");
2794
2795        let invalidations = tree.apply(&[
2796            fsevents_rename(tree.path("sessions/today.jsonl")),
2797            modified(tree.path("sessions/today.jsonl")),
2798        ]);
2799
2800        assert_eq!(invalidations, vec![]);
2801        tree.assert_converged();
2802    }
2803
2804    /// A rename that changes only case leaves the old spelling resolvable on an insensitive
2805    /// filesystem, so a stat alone would keep both names. The parent listing is the
2806    /// arbiter, and a stale spelling costs a reconcile of its parent, not of the root.
2807    #[test]
2808    fn a_case_only_rename_keeps_one_spelling() {
2809        let tree = RenameFixture::new(&["docs/Readme.md", "docs/Guide/intro.md"]);
2810        let insensitive = crate::test_support::resolves_case_insensitively(&tree.root);
2811        fs::rename(tree.path("docs/Readme.md"), tree.path("docs/README.md")).expect("recase");
2812        fs::rename(tree.path("docs/Guide"), tree.path("docs/guide")).expect("recase directory");
2813
2814        let invalidations = tree.apply(&[
2815            fsevents_rename(tree.path("docs/Readme.md")),
2816            fsevents_rename(tree.path("docs/README.md")),
2817            fsevents_rename(tree.path("docs/Guide")),
2818            fsevents_rename(tree.path("docs/guide")),
2819            // A hint queued under the old spelling before the rename.
2820            modified(tree.path("docs/Guide/intro.md")),
2821        ]);
2822
2823        tree.assert_converged();
2824        assert!(
2825            invalidations.iter().all(|(path, _)| !path.as_os_str().is_empty()),
2826            "no root reconcile: {invalidations:?}"
2827        );
2828        if insensitive {
2829            assert!(
2830                invalidations.contains(&(PathBuf::from("docs"), InvalidateReason::UnpairedRename)),
2831                "the stale spelling reconciles its parent: {invalidations:?}"
2832            );
2833        }
2834    }
2835
2836    /// A listing read earlier in the batch is not the arbiter of a name created since: a
2837    /// miss re-reads the parent, so churn cannot pass for a stale spelling.
2838    #[test]
2839    fn a_listing_miss_rereads_the_parent_before_answering() {
2840        let dir = tempfile::tempdir().expect("tempdir");
2841        let root = dir.path();
2842        fs::write(root.join("first"), b"1").expect("first");
2843        let mut listings = ParentListings::default();
2844
2845        assert_eq!(listings.lists(root, Path::new("first")), Some(true));
2846        fs::write(root.join("created-after-the-listing"), b"2").expect("later entry");
2847        assert_eq!(listings.lists(root, Path::new("created-after-the-listing")), Some(true));
2848        assert_eq!(listings.lists(root, Path::new("never-created")), Some(false));
2849        assert_eq!(listings.lists(root, Path::new("missing-directory/child")), None);
2850    }
2851
2852    /// The real backend on this platform: renames inside the root converge on the tree
2853    /// without ever reconciling the whole root.
2854    #[test]
2855    fn native_renames_converge_without_reconciling_the_root() {
2856        let _serialized = real_watcher_guard();
2857        let dir = tempfile::tempdir().expect("tempdir");
2858        let root = dir.path().canonicalize().expect("canonical root");
2859        for file in ["from/moved.txt", "tree/sub/leaf.txt", "to/resident.txt"] {
2860            fs::create_dir_all(root.join(file).parent().expect("parent")).expect("parents");
2861            fs::write(root.join(file), file.as_bytes()).expect("fixture file");
2862        }
2863        let watcher = Watcher::new(&root, WatchConfig::default()).expect("watcher");
2864        if !establish_watch(&watcher, &root) {
2865            return;
2866        }
2867        let config = ScanConfig::default();
2868        let (index, _) = crate::scan::scan_into_index(&root, &config).expect("scan");
2869        let handle = IndexHandle::new(index);
2870
2871        fs::rename(root.join("from/moved.txt"), root.join("to/moved.txt")).expect("move file");
2872        fs::rename(root.join("tree"), root.join("renamed-tree")).expect("rename directory");
2873
2874        let converged = || {
2875            handle.kind(Path::new("to/moved.txt")).expect("lookup").is_some()
2876                && handle.kind(Path::new("from/moved.txt")).expect("lookup").is_none()
2877                && handle.kind(Path::new("renamed-tree/sub/leaf.txt")).expect("lookup").is_some()
2878                && handle.kind(Path::new("tree")).expect("lookup").is_none()
2879        };
2880        let mut commits = Vec::new();
2881        let start = Instant::now();
2882        while !converged() && start.elapsed() < REAL_BACKEND_DELIVERY {
2883            watcher
2884                .apply_next(&handle, &config, Duration::from_millis(200), &mut |commit| {
2885                    commits.push(commit.clone());
2886                })
2887                .expect("apply");
2888        }
2889
2890        assert!(converged(), "the renames never converged: {commits:?}");
2891        let root_reconciles = commits
2892            .iter()
2893            .flat_map(|commit| commit.changes.iter())
2894            .filter(|change| {
2895                matches!(change, crate::EffectiveChange::Invalidated { path, .. }
2896                    if path.as_os_str().is_empty())
2897            })
2898            .count();
2899        assert_eq!(root_reconciles, 0, "a rename reconciled the whole root: {commits:?}");
2900        let (cold, _) = crate::scan::scan_into_index(&root, &config).expect("cold");
2901        assert_eq!(handle.read_with(entries).expect("read"), entries(&cold));
2902    }
2903
2904    #[test]
2905    fn zero_settle_is_rejected_before_starting_a_busy_worker() {
2906        let config = WatchConfig { settle: Duration::ZERO, ..WatchConfig::default() };
2907
2908        assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2909
2910        let config = WatchConfig { intent_capacity: 0, ..WatchConfig::default() };
2911        assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2912
2913        let config = WatchConfig {
2914            batch_path_capacity: MAX_BATCH_PATH_CAPACITY,
2915            intent_capacity: MAX_BUFFERED_INTENT_PATHS / MAX_BATCH_PATH_CAPACITY + 1,
2916            ..WatchConfig::default()
2917        };
2918        assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2919    }
2920}