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.
11//! - When a directory is created, backends that watch per directory register the new
12//!   watch *after* the fact — anything created inside that window produces no event.
13//! - Kernel queues overflow. inotify's `Q_OVERFLOW`, `FSEvents`' `MustScanSubDirs`, and
14//!   Windows buffer overruns all mean "your view is now incomplete". They surface here
15//!   as `Flag::Rescan`, and dropping that signal is precisely how an event-driven index
16//!   silently diverges from the filesystem — which is what the `watchfiles` layer
17//!   metabrowser runs today does, mapping notify's rich event model down to
18//!   `(change, path)` and letting the rescan flag fall through a match arm.
19//!
20//! So this layer never forwards an event. It coalesces, then **verifies by stat**, and
21//! emits only [`Op::Upsert`] with a fresh fingerprint, [`Op::Remove`], or —
22//! when it genuinely cannot describe the change — [`Op::InvalidateSubtree`], which the
23//! scan layer resolves back into precise committed changes.
24//!
25//! Building on notify rather than on raw platform APIs is deliberate: its six backends
26//! and its overflow signaling are proven, and the information loss that motivates this
27//! module all happens in layers *above* it.
28
29use std::collections::BTreeMap;
30use std::panic::{AssertUnwindSafe, catch_unwind};
31use std::path::{Path, PathBuf};
32use std::sync::Arc;
33use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
34use std::sync::mpsc::{Receiver, RecvTimeoutError, SyncSender, TrySendError, sync_channel};
35use std::thread::JoinHandle;
36use std::time::{Duration, Instant};
37
38use notify::{EventKind, RecommendedWatcher, RecursiveMode, Watcher as NotifyWatcher};
39
40#[cfg(test)]
41mod scripted_events;
42
43use crate::engine_contract::{
44    Commit, Error, InvalidateReason, Observation, ObservationOp, Op, Result,
45};
46use crate::scan;
47use crate::{ApplyOutcome, IndexHandle, ScanConfig};
48
49/// Optimistic re-verification retries before the applying driver guarantees progress
50/// through conservative root invalidation and reconciliation.
51const MAX_OPTIMISTIC_APPLY_ATTEMPTS: usize = 3;
52
53const WORKER_RUNNING: u8 = 0;
54const WORKER_STOPPED: u8 = 1;
55const WORKER_PANICKED: u8 = 2;
56const MAX_EVENT_CAPACITY: usize = 64 * 1024;
57const MAX_BATCH_PATH_CAPACITY: usize = 64 * 1024;
58const MAX_BUFFERED_INTENT_PATHS: usize = 1024 * 1024;
59
60/// Tuning for event coalescing.
61#[derive(Clone, Copy, Debug)]
62pub struct WatchConfig {
63    /// How long the event stream must be quiet before a batch is emitted.
64    pub settle: Duration,
65    /// Longest a batch may be held open while events keep arriving. Without a ceiling, a
66    /// continuously busy tree would never produce a delta at all.
67    pub max_hold: Duration,
68    /// Whether a newly created directory triggers a re-list of its contents.
69    ///
70    /// On by default, and it should stay on for inotify and kqueue: those backends
71    /// install a directory's watch only after the create event arrives, so files created
72    /// in between are never reported by anything.
73    pub relist_new_dirs: bool,
74    /// Maximum raw backend events queued before overload collapses to one root
75    /// invalidation. Backend callback threads never block on this queue.
76    pub event_capacity: usize,
77    /// Maximum distinct paths retained in one coalesced intent.
78    pub batch_path_capacity: usize,
79    /// Maximum coalesced intents awaiting a consumer.
80    pub intent_capacity: usize,
81}
82
83impl Default for WatchConfig {
84    fn default() -> Self {
85        Self {
86            // The 50 ms step / 1.6 s ceiling pairing is watchfiles' batching loop, which
87            // is the part of that stack worth keeping.
88            settle: Duration::from_millis(50),
89            max_hold: Duration::from_millis(1600),
90            relist_new_dirs: true,
91            event_capacity: 4096,
92            batch_path_capacity: 4096,
93            intent_capacity: 16,
94        }
95    }
96}
97
98impl WatchConfig {
99    fn validate(self) -> Result<()> {
100        let buffered_paths = self.batch_path_capacity.checked_mul(self.intent_capacity);
101        if self.settle.is_zero()
102            || self.max_hold.is_zero()
103            || self.max_hold < self.settle
104            || self.event_capacity == 0
105            || self.batch_path_capacity == 0
106            || self.intent_capacity == 0
107            || self.event_capacity > MAX_EVENT_CAPACITY
108            || self.batch_path_capacity > MAX_BATCH_PATH_CAPACITY
109            || buffered_paths.is_none_or(|paths| paths > MAX_BUFFERED_INTENT_PATHS)
110        {
111            return Err(Error::UnsupportedScanConfig(
112                "watch durations and capacities exceed the supported nonzero bounds, or max_hold is less than settle",
113            ));
114        }
115        Ok(())
116    }
117}
118
119/// What a coalesced path still needs before it can become an observation.
120#[derive(Clone, Copy, PartialEq, Eq, Debug)]
121enum Pending {
122    /// Stat it and decide. Covers creates, writes, removes, and every rename shape:
123    /// letting the stat decide is what makes the same code correct on all backends.
124    Verify {
125        /// Preserve whether a create event occurred while this path was coalesced. Only
126        /// a newly created directory has the watch-registration race that needs a relist.
127        relist_if_dir: bool,
128    },
129    /// The producer already knows it cannot describe this precisely.
130    Escalate(InvalidateReason),
131}
132
133#[derive(Debug, Default)]
134struct CoalescedIntent {
135    pending: BTreeMap<PathBuf, Pending>,
136}
137
138enum RawMessage {
139    Event(notify::Result<notify::Event>),
140    Flush(SyncSender<()>),
141    Stop,
142}
143
144/// A live watcher over one tree.
145///
146/// The watcher is movable to one consuming thread. It is intentionally not shareable:
147/// its private standard-library receiver enforces one ordered consumer for coalesced
148/// intents. Use a separate [`IndexHandle`] to serve concurrent readers and writers.
149/// Dropping the watcher stops the OS watch and shuts the worker thread down.
150pub struct Watcher {
151    root: PathBuf,
152    config: WatchConfig,
153    /// `Option` only so [`Drop`] can release it before joining the worker.
154    inner: Option<RecommendedWatcher>,
155    intents: Receiver<CoalescedIntent>,
156    control: Option<SyncSender<RawMessage>>,
157    cancelled: Arc<AtomicBool>,
158    worker_status: Arc<AtomicU8>,
159    worker: Option<JoinHandle<()>>,
160}
161
162#[cfg(test)]
163pub(crate) struct ScriptedSender {
164    root: PathBuf,
165    raw: SyncSender<RawMessage>,
166    overflowed: Arc<AtomicBool>,
167}
168
169#[cfg(test)]
170impl ScriptedSender {
171    pub(crate) fn send(&self, source: &str) -> Result<()> {
172        let events =
173            scripted_events::parse_script(source, &self.root).map_err(Error::WatchScript)?;
174        for event in events {
175            enqueue_raw(&self.raw, &self.overflowed, event);
176        }
177        Ok(())
178    }
179}
180
181/// Effects of one watch observation and any reconciliation it requested.
182#[derive(Debug)]
183pub struct WatchApplyReport {
184    /// Effect of the verified watch intent itself.
185    pub apply: ApplyOutcome,
186    /// Effect of closing any invalidation loop opened by the intent.
187    pub reconciliation: scan::ReconcileReport,
188}
189
190impl Watcher {
191    /// Start watching `root` recursively.
192    pub fn new(root: &Path, config: WatchConfig) -> Result<Self> {
193        config.validate()?;
194        let root = root.canonicalize().map_err(|e| Error::io(root, e))?;
195
196        let (raw_tx, raw_rx) = sync_channel::<RawMessage>(config.event_capacity);
197        let (intent_tx, intent_rx) = sync_channel::<CoalescedIntent>(config.intent_capacity);
198        let control_tx = raw_tx.clone();
199        let overflowed = Arc::new(AtomicBool::new(false));
200        let callback_overflowed = Arc::clone(&overflowed);
201
202        let mut inner = notify::recommended_watcher(move |result| {
203            enqueue_raw(&raw_tx, &callback_overflowed, result);
204        })
205        .map_err(|error| notify_error(&root, error))?;
206        inner.watch(&root, RecursiveMode::Recursive).map_err(|error| notify_error(&root, error))?;
207
208        let worker_root = root.clone();
209        let cancelled = Arc::new(AtomicBool::new(false));
210        let worker_cancelled = Arc::clone(&cancelled);
211        let worker_overflowed = Arc::clone(&overflowed);
212        let worker_status = Arc::new(AtomicU8::new(WORKER_RUNNING));
213        let tracked_status = Arc::clone(&worker_status);
214        let worker = std::thread::Builder::new()
215            .name("fdu-watch".into())
216            .spawn(move || {
217                let _counter_guard = crate::counters::thread_flush_guard();
218                run_tracked_worker(&tracked_status, || {
219                    run_worker(
220                        &worker_root,
221                        config,
222                        &raw_rx,
223                        &intent_tx,
224                        &worker_overflowed,
225                        &worker_cancelled,
226                    );
227                });
228            })
229            .map_err(|e| Error::io(&root, e))?;
230
231        Ok(Self {
232            root,
233            config,
234            inner: Some(inner),
235            intents: intent_rx,
236            control: Some(control_tx),
237            cancelled,
238            worker_status,
239            worker: Some(worker),
240        })
241    }
242
243    #[cfg(test)]
244    pub(crate) fn scripted(
245        root: &Path,
246        config: WatchConfig,
247        events: &Path,
248    ) -> Result<(Self, ScriptedSender)> {
249        config.validate()?;
250        let root = root.canonicalize().map_err(|error| Error::io(root, error))?;
251        let scripted = scripted_events::read_script(events, &root).map_err(Error::WatchScript)?;
252        let (raw_tx, raw_rx) = sync_channel::<RawMessage>(config.event_capacity);
253        let (intent_tx, intent_rx) = sync_channel::<CoalescedIntent>(config.intent_capacity);
254        let control_tx = raw_tx.clone();
255        let overflowed = Arc::new(AtomicBool::new(false));
256        for event in scripted {
257            enqueue_raw(&raw_tx, &overflowed, event);
258        }
259
260        let worker_root = root.clone();
261        let cancelled = Arc::new(AtomicBool::new(false));
262        let worker_cancelled = Arc::clone(&cancelled);
263        let worker_overflowed = Arc::clone(&overflowed);
264        let worker_status = Arc::new(AtomicU8::new(WORKER_RUNNING));
265        let tracked_status = Arc::clone(&worker_status);
266        let worker = std::thread::Builder::new()
267            .name("fdu-scripted-watch".into())
268            .spawn(move || {
269                let _counter_guard = crate::counters::thread_flush_guard();
270                run_tracked_worker(&tracked_status, || {
271                    run_worker(
272                        &worker_root,
273                        config,
274                        &raw_rx,
275                        &intent_tx,
276                        &worker_overflowed,
277                        &worker_cancelled,
278                    );
279                });
280            })
281            .map_err(|error| Error::io(&root, error))?;
282
283        let sender = ScriptedSender {
284            root: root.clone(),
285            raw: control_tx.clone(),
286            overflowed: Arc::clone(&overflowed),
287        };
288        Ok((
289            Self {
290                root,
291                config,
292                inner: None,
293                intents: intent_rx,
294                control: Some(control_tx),
295                cancelled,
296                worker_status,
297                worker: Some(worker),
298            },
299            sender,
300        ))
301    }
302
303    /// Block for and verify the next coalesced intent, up to `timeout`.
304    ///
305    /// Filesystem calls happen synchronously on this consuming thread and may have
306    /// ordinary filesystem latency. A timeout, a stopped worker, and a panicked worker
307    /// are distinct outcomes. The returned observation contains relative paths but no
308    /// root identity; use [`Self::apply_next`] for the supported root-checked applying
309    /// driver rather than applying it to an arbitrary index.
310    pub fn next_observation(&self, timeout: Duration) -> Result<Option<Observation>> {
311        let Some(intent) = self.next_intent(timeout)? else {
312            return Ok(None);
313        };
314        Ok(Some(verify_intent(&self.root, self.config, &intent, &ScanConfig::default())))
315    }
316
317    fn next_intent(&self, timeout: Duration) -> Result<Option<CoalescedIntent>> {
318        match self.intents.recv_timeout(timeout) {
319            Ok(intent) => Ok(Some(intent)),
320            Err(RecvTimeoutError::Timeout) => Ok(None),
321            Err(RecvTimeoutError::Disconnected) => match self.worker_status.load(Ordering::Acquire)
322            {
323                WORKER_PANICKED => Err(Error::WatchWorkerPanicked),
324                _ => Err(Error::WatchStopped),
325            },
326        }
327    }
328
329    /// Advance capture through every raw hint accepted before this call.
330    ///
331    /// A full intent queue may retain their loss as one internal sticky overflow marker;
332    /// the next barrier advances that marker after the consumer drains queue capacity.
333    pub(crate) fn flush_capture(&self) -> Result<()> {
334        let Some(control) = self.control.as_ref() else {
335            return Err(Error::WatchStopped);
336        };
337        let (acknowledge, acknowledged) = sync_channel(0);
338        control.send(RawMessage::Flush(acknowledge)).map_err(|_| Error::WatchStopped)?;
339        acknowledged.recv().map_err(|_| Error::WatchStopped)
340    }
341
342    /// Maximum intents that can represent the raw hints preceding a flush barrier.
343    ///
344    /// The intent queue supplies the ordinary bound. One additional root invalidation
345    /// may be held as `sticky_overflow` when that queue filled, so the handoff derives
346    /// its drain work from the configured capacity rather than imposing another
347    /// unrelated limit.
348    pub(crate) const fn capture_backlog_bound(&self) -> usize {
349        self.config.intent_capacity.saturating_add(1)
350    }
351
352    /// Apply one verified watch observation and close any invalidation loop it opens.
353    ///
354    /// Restricted `max_depth` and `one_filesystem` scopes are rejected until the watch
355    /// adapter can filter raw backend events against those boundaries.
356    pub fn apply_next(
357        &self,
358        index: &IndexHandle,
359        scan_config: &ScanConfig,
360        timeout: Duration,
361        sink: &mut dyn FnMut(&Commit),
362    ) -> Result<Option<WatchApplyReport>> {
363        scan_config.validate_for_watch_scope(index.scope()?)?;
364        let indexed_root = index.root_path()?;
365        if indexed_root != self.root {
366            return Err(Error::WatchRootMismatch {
367                watched: self.root.clone(),
368                indexed: indexed_root,
369            });
370        }
371        let Some(intent) = self.next_intent(timeout)? else {
372            return Ok(None);
373        };
374        apply_intent(index, &self.root, self.config, &intent, scan_config, sink).map(Some)
375    }
376
377    /// Apply one intent through the opened-root lifecycle and exact resource boundary.
378    pub(crate) fn apply_next_controlled(
379        &self,
380        index: &IndexHandle,
381        scan_config: &ScanConfig,
382        timeout: Duration,
383        control: &dyn scan::ReconcileControl,
384        sink: &mut dyn FnMut(&Commit),
385    ) -> Result<Option<WatchApplyReport>> {
386        scan_config.validate_for_watch_scope(index.scope()?)?;
387        let indexed_root = index.root_path()?;
388        if indexed_root != self.root {
389            return Err(Error::WatchRootMismatch {
390                watched: self.root.clone(),
391                indexed: indexed_root,
392            });
393        }
394        let Some(intent) = self.next_intent(timeout)? else {
395            return Ok(None);
396        };
397        apply_intent_controlled(index, &self.root, self.config, &intent, scan_config, control, sink)
398            .map(Some)
399    }
400}
401
402fn apply_intent(
403    index: &IndexHandle,
404    root: &Path,
405    watch_config: WatchConfig,
406    intent: &CoalescedIntent,
407    scan_config: &ScanConfig,
408    sink: &mut dyn FnMut(&Commit),
409) -> Result<WatchApplyReport> {
410    let mut verifier =
411        |_: &Path, _: &Observation| Ok(verify_intent(root, watch_config, intent, scan_config));
412    let apply = apply_reverified_with(index, &Observation::default(), scan_config, &mut verifier)?;
413    if let Some(commit) = apply.commit.as_ref() {
414        sink(commit);
415    }
416    let reconciliation = scan::reconcile_pending_handle(index, scan_config, sink)?;
417    retain_unreadable(index, root, &reconciliation, sink)?;
418    Ok(WatchApplyReport { apply, reconciliation })
419}
420
421/// Retain why a reconciliation could not read part of the tree.
422///
423/// The shared reconcile settles an unreadable subtree rather than walking it again on
424/// every later event, so this report is the only time its error is seen. Committing the
425/// causes keeps the partial freshness it left explainable, as the opened root does, and
426/// names each path relative to the root so a later clean walk of it can drop the issue.
427fn retain_unreadable(
428    index: &IndexHandle,
429    root: &Path,
430    reconciliation: &scan::ReconcileReport,
431    sink: &mut dyn FnMut(&Commit),
432) -> Result<()> {
433    if reconciliation.scan.errors.is_empty() {
434        return Ok(());
435    }
436    let retained = reconciliation.scan.errors.len().min(crate::MAX_RETAINED_ISSUES);
437    let issues = reconciliation.scan.errors[..retained]
438        .iter()
439        .map(|error| crate::Issue::from_error_under(root, error))
440        .collect();
441    let omitted = u64::try_from(reconciliation.scan.errors.len() - retained).unwrap_or(u64::MAX);
442    let outcome =
443        index.transition_observation(crate::index::ObservationTransition::Unreadable {
444            issues,
445            omitted,
446        })?;
447    if let Some(commit) = outcome.commit.as_ref() {
448        sink(commit);
449    }
450    Ok(())
451}
452
453fn apply_intent_controlled(
454    index: &IndexHandle,
455    root: &Path,
456    watch_config: WatchConfig,
457    intent: &CoalescedIntent,
458    scan_config: &ScanConfig,
459    control: &dyn scan::ReconcileControl,
460    sink: &mut dyn FnMut(&Commit),
461) -> Result<WatchApplyReport> {
462    let mut verifier =
463        |_: &Path, _: &Observation| Ok(verify_intent(root, watch_config, intent, scan_config));
464    let apply = apply_reverified_with_control(
465        index,
466        &Observation::default(),
467        scan_config,
468        control,
469        &mut verifier,
470    )?;
471    if let Some(commit) = apply.commit.as_ref() {
472        sink(commit);
473    }
474    let reconciliation =
475        scan::reconcile_pending_handle_controlled(index, scan_config, control, sink)?;
476    Ok(WatchApplyReport { apply, reconciliation })
477}
478
479/// Test the unrooted observation driver without making it a public apply capability.
480///
481/// Production callers use [`Watcher::apply_next`], which proves that the watcher and
482/// index have the same root before consuming an intent. An [`Observation`] intentionally
483/// remains a generic producer batch and does not claim a filesystem root identity.
484#[cfg(test)]
485fn apply_observation(
486    index: &IndexHandle,
487    observation: &Observation,
488    scan_config: &ScanConfig,
489    sink: &mut dyn FnMut(&Commit),
490) -> Result<WatchApplyReport> {
491    scan_config.validate_for_watch_scope(index.scope()?)?;
492    let apply = apply_reverified(index, observation, scan_config)?;
493    if let Some(commit) = apply.commit.as_ref() {
494        sink(commit);
495    }
496    let reconciliation = scan::reconcile_pending_handle(index, scan_config, sink)?;
497    Ok(WatchApplyReport { apply, reconciliation })
498}
499
500/// Re-stat a queued watch sample against a clock-stable index boundary before applying
501/// it. Filesystem verification always runs outside the index lock. The filesystem itself
502/// cannot be locked by this process: a sample is valid at its `stat` linearization point,
503/// and any mutation after that point remains a later backend event. Queue loss or
504/// ambiguity becomes an invalidation and reconciliation, rather than a claim that the
505/// disk stayed frozen between `stat` and the in-memory commit. Sustained competing index
506/// writes conservatively invalidate the root without blocking readers on filesystem I/O.
507#[cfg(test)]
508fn apply_reverified(
509    index: &IndexHandle,
510    observation: &Observation,
511    scan_config: &ScanConfig,
512) -> Result<ApplyOutcome> {
513    let mut verifier = |root: &Path, observation: &Observation| {
514        reverify_observation(root, observation, scan_config)
515    };
516    apply_reverified_with(index, observation, scan_config, &mut verifier)
517}
518
519fn apply_reverified_with(
520    index: &IndexHandle,
521    observation: &Observation,
522    scan_config: &ScanConfig,
523    verifier: &mut impl FnMut(&Path, &Observation) -> Result<Observation>,
524) -> Result<ApplyOutcome> {
525    let (root, scope, _) = index.watch_boundary()?;
526    scan_config.validate_for_watch_scope(scope)?;
527    for _ in 0..MAX_OPTIMISTIC_APPLY_ATTEMPTS {
528        let clock = index.clock()?;
529        let candidate = verifier(&root, observation)?;
530        let candidate = escalate_unknown_ancestry(index, candidate)?;
531        if let Some(outcome) = index.apply_if_clock(clock, &candidate)? {
532            return Ok(outcome);
533        }
534    }
535
536    index.invalidate_root(InvalidateReason::WatchContention)
537}
538
539fn apply_reverified_with_control(
540    index: &IndexHandle,
541    observation: &Observation,
542    scan_config: &ScanConfig,
543    control: &dyn scan::ReconcileControl,
544    verifier: &mut impl FnMut(&Path, &Observation) -> Result<Observation>,
545) -> Result<ApplyOutcome> {
546    let (root, scope, _) = index.watch_boundary()?;
547    scan_config.validate_for_watch_scope(scope)?;
548    for _ in 0..MAX_OPTIMISTIC_APPLY_ATTEMPTS {
549        control.check_active()?;
550        let clock = index.clock()?;
551        let candidate = verifier(&root, observation)?;
552        let candidate = escalate_unknown_ancestry(index, candidate)?;
553        control.before_conditional_commit()?;
554        if let Some(outcome) =
555            index.apply_opened_if_clock(clock, &candidate, control.max_files())?
556        {
557            return Ok(outcome);
558        }
559    }
560
561    control.before_conditional_commit()?;
562    let outcome = index.apply_opened(
563        &Observation::new(vec![Op::InvalidateSubtree {
564            path: PathBuf::new(),
565            reason: InvalidateReason::WatchContention,
566        }]),
567        control.max_files(),
568    )?;
569    Ok(outcome)
570}
571
572/// Replace unverifiable child facts with bounded reconciliation hints.
573fn escalate_unknown_ancestry(index: &IndexHandle, candidate: Observation) -> Result<Observation> {
574    let unknown = index.unknown_ancestry(&candidate)?;
575    if unknown.is_empty() {
576        return Ok(candidate);
577    }
578
579    let mut roots: Vec<PathBuf> = unknown.into_iter().map(|(_, root)| root).collect();
580    roots.sort_by(|left, right| {
581        left.components().count().cmp(&right.components().count()).then_with(|| left.cmp(right))
582    });
583    roots.dedup();
584    let mut covering: Vec<PathBuf> = Vec::with_capacity(roots.len());
585    for root in roots {
586        if !covering.iter().any(|ancestor| root.starts_with(ancestor)) {
587            covering.push(root);
588        }
589    }
590
591    let mut ops: Vec<ObservationOp> = candidate
592        .ops
593        .into_iter()
594        .filter(|observed| !covering.iter().any(|root| observed.op.path().starts_with(root)))
595        .collect();
596    ops.extend(covering.into_iter().map(|path| {
597        ObservationOp::unconditional(Op::InvalidateSubtree {
598            path,
599            reason: InvalidateReason::UnknownAncestry,
600        })
601    }));
602    Ok(Observation::from_ops(ops))
603}
604
605#[cfg(test)]
606fn reverify_observation(
607    root: &Path,
608    observation: &Observation,
609    scan_config: &ScanConfig,
610) -> Result<Observation> {
611    let mut ops = Vec::with_capacity(observation.len().saturating_mul(2));
612    for observed in &observation.ops {
613        let relative = scan::normalize_subtree(observed.op.path())?;
614        match &observed.op {
615            Op::InvalidateSubtree { reason, .. } => {
616                ops.push(Op::InvalidateSubtree { path: relative, reason: *reason });
617            }
618            Op::Upsert { .. } | Op::Remove { .. } => {
619                let absolute = root.join(&relative);
620                match std::fs::symlink_metadata(&absolute) {
621                    Ok(metadata) => {
622                        let (kind, attrs) = scan::observe(&absolute, &metadata)
623                            .map_err(|source| Error::io(&absolute, source))?;
624                        match crate::admission::decide_path(
625                            &relative,
626                            kind,
627                            scan_config.hidden(),
628                            scan_config.exclude_special,
629                        ) {
630                            crate::admission::Disposition::Retain => {
631                                ops.push(Op::Upsert { path: relative.clone(), kind, attrs });
632                                if let Some(control) =
633                                    scan::read_control_op(scan_config, root, &relative, kind)?
634                                {
635                                    ops.push(control);
636                                }
637                            }
638                            crate::admission::Disposition::ControlOnly => {
639                                ops.push(Op::Remove { path: relative.clone() });
640                                if let Some(control) =
641                                    scan::read_control_op(scan_config, root, &relative, kind)?
642                                {
643                                    ops.push(control);
644                                }
645                            }
646                            crate::admission::Disposition::Reject => {
647                                ops.push(Op::Remove { path: relative });
648                            }
649                        }
650                    }
651                    Err(error) => ops.push(op_for_stat_error(relative, &error)),
652                }
653            }
654            Op::ControlUpsert { source, .. } => {
655                ops.push(Op::ControlUpsert { path: relative, source: source.clone() });
656            }
657            Op::ControlRemove { .. } => ops.push(Op::ControlRemove { path: relative }),
658        }
659    }
660    Ok(Observation::new(ops))
661}
662
663impl Drop for Watcher {
664    fn drop(&mut self) {
665        self.cancelled.store(true, Ordering::Release);
666        if let Some(control) = self.control.take() {
667            let _ = control.try_send(RawMessage::Stop);
668        }
669        // Stop the backend and release its callback sender before joining. Every worker
670        // send is nonblocking, so a full consumer queue cannot deadlock teardown.
671        self.inner.take();
672        if let Some(worker) = self.worker.take() {
673            let _ = worker.join();
674        }
675    }
676}
677
678fn notify_error(path: &Path, err: notify::Error) -> Error {
679    Error::io(path, std::io::Error::other(err))
680}
681
682fn enqueue_raw(
683    sender: &SyncSender<RawMessage>,
684    overflowed: &AtomicBool,
685    event: notify::Result<notify::Event>,
686) {
687    match sender.try_send(RawMessage::Event(event)) {
688        Ok(()) | Err(TrySendError::Disconnected(_)) => {}
689        Err(TrySendError::Full(_)) => overflowed.store(true, Ordering::Release),
690    }
691}
692
693fn run_tracked_worker(status: &AtomicU8, worker: impl FnOnce()) {
694    let outcome = catch_unwind(AssertUnwindSafe(worker));
695    status.store(if outcome.is_ok() { WORKER_STOPPED } else { WORKER_PANICKED }, Ordering::Release);
696}
697
698fn run_worker(
699    root: &Path,
700    config: WatchConfig,
701    raw: &Receiver<RawMessage>,
702    out: &SyncSender<CoalescedIntent>,
703    overflowed: &AtomicBool,
704    cancelled: &AtomicBool,
705) {
706    let mut pending: BTreeMap<PathBuf, Pending> = BTreeMap::new();
707    let mut batch_started: Option<Instant> = None;
708    let mut sticky_overflow = false;
709
710    loop {
711        if cancelled.load(Ordering::Acquire) {
712            return;
713        }
714        if overflowed.swap(false, Ordering::AcqRel) {
715            collapse_to_overflow(&mut pending);
716            batch_started.get_or_insert_with(Instant::now);
717        }
718        if sticky_overflow {
719            match try_deliver_overflow(out) {
720                Ok(true) => sticky_overflow = false,
721                Ok(false) => {}
722                Err(()) => return,
723            }
724        }
725
726        match raw.recv_timeout(config.settle) {
727            Ok(RawMessage::Event(Ok(event))) => {
728                record(root, &event, &mut pending, config.batch_path_capacity);
729                batch_started.get_or_insert_with(Instant::now);
730            }
731            Ok(RawMessage::Event(Err(err))) => {
732                // notify reports a watch failure. It cannot say what was missed, so the
733                // only honest response is to escalate the whole tree.
734                let _ = err;
735                collapse_to_overflow(&mut pending);
736                batch_started.get_or_insert_with(Instant::now);
737            }
738            Ok(RawMessage::Flush(acknowledge)) => {
739                if overflowed.swap(false, Ordering::AcqRel) {
740                    collapse_to_overflow(&mut pending);
741                }
742                if !pending.is_empty()
743                    && try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err()
744                {
745                    return;
746                }
747                // The consumer may have drained the queue while this worker waited
748                // for the barrier. Publish sticky loss before acknowledging; doing it
749                // at the next loop iteration races the consumer's final empty poll.
750                if sticky_overflow {
751                    match try_deliver_overflow(out) {
752                        Ok(true) => sticky_overflow = false,
753                        Ok(false) => {}
754                        Err(()) => return,
755                    }
756                }
757                batch_started = None;
758                let _ = acknowledge.send(());
759            }
760            Ok(RawMessage::Stop) => return,
761            Err(RecvTimeoutError::Timeout) => {
762                // Quiet for a full step: the batch has settled.
763                if !pending.is_empty()
764                    && try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err()
765                {
766                    return;
767                }
768                batch_started = None;
769                continue;
770            }
771            Err(RecvTimeoutError::Disconnected) => {
772                if !cancelled.load(Ordering::Acquire) && !pending.is_empty() {
773                    let _ = try_deliver_pending(&mut pending, out, &mut sticky_overflow);
774                }
775                return;
776            }
777        }
778
779        // A tree under continuous churn never goes quiet, so cap how long a batch waits.
780        if max_hold_elapsed(batch_started, config.max_hold) {
781            if try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err() {
782                return;
783            }
784            batch_started = None;
785        }
786    }
787}
788
789fn max_hold_elapsed(started: Option<Instant>, max_hold: Duration) -> bool {
790    started.is_some_and(|start| start.elapsed() >= max_hold)
791}
792
793/// Fold one event into the pending set.
794fn record(
795    root: &Path,
796    event: &notify::Event,
797    pending: &mut BTreeMap<PathBuf, Pending>,
798    capacity: usize,
799) {
800    if event.need_rescan() {
801        // The kernel dropped events. Escalate the narrowest subtree the event names, or
802        // the whole root when it names zero/multiple paths or crosses the watch boundary.
803        let target = if event.paths.len() == 1 {
804            relative_to(root, &event.paths[0]).unwrap_or_default()
805        } else {
806            PathBuf::new()
807        };
808        queue_pending(
809            pending,
810            target,
811            Pending::Escalate(InvalidateReason::WatchOverflow),
812            capacity,
813        );
814        return;
815    }
816
817    let rename_mode = match event.kind {
818        EventKind::Modify(notify::event::ModifyKind::Name(mode)) => Some(mode),
819        _ => None,
820    };
821    let paired_rename = matches!(rename_mode, Some(notify::event::RenameMode::Both))
822        && event.paths.len() == 2
823        && event.paths.iter().all(|path| relative_to(root, path).is_some());
824    if rename_mode.is_some() && !paired_rename {
825        // A one-sided rename gives no safe bound on where its counterpart lives. A full
826        // reconciliation is more expensive than guessing a parent, but it cannot leave
827        // the old name behind or miss a moved-in subtree.
828        queue_pending(
829            pending,
830            PathBuf::new(),
831            Pending::Escalate(InvalidateReason::UnpairedRename),
832            capacity,
833        );
834    }
835
836    for (position, path) in event.paths.iter().enumerate() {
837        let Some(rel) = relative_to(root, path) else {
838            continue;
839        };
840        if matches!(event.kind, EventKind::Access(_)) {
841            continue; // Reads change nothing this engine records.
842        }
843        let relist_if_dir = matches!(event.kind, EventKind::Create(_))
844            || matches!(rename_mode, Some(notify::event::RenameMode::To))
845            || (paired_rename && position == 1);
846        queue_pending(pending, rel, Pending::Verify { relist_if_dir }, capacity);
847    }
848}
849
850fn queue_pending(
851    pending: &mut BTreeMap<PathBuf, Pending>,
852    path: PathBuf,
853    state: Pending,
854    capacity: usize,
855) {
856    if matches!(pending.get(Path::new("")), Some(Pending::Escalate(_))) {
857        return;
858    }
859    if path.as_os_str().is_empty() && matches!(state, Pending::Escalate(_)) {
860        pending.clear();
861        pending.insert(path, state);
862        return;
863    }
864    if let Some(existing) = pending.get_mut(&path) {
865        match (existing, state) {
866            (Pending::Escalate(_), _) => {}
867            (Pending::Verify { relist_if_dir }, Pending::Verify { relist_if_dir: additional }) => {
868                *relist_if_dir |= additional;
869            }
870            (slot @ Pending::Verify { .. }, Pending::Escalate(reason)) => {
871                *slot = Pending::Escalate(reason);
872            }
873        }
874        return;
875    }
876    if pending.len() >= capacity {
877        collapse_to_overflow(pending);
878    } else {
879        pending.insert(path, state);
880    }
881}
882
883fn collapse_to_overflow(pending: &mut BTreeMap<PathBuf, Pending>) {
884    pending.clear();
885    pending.insert(PathBuf::new(), Pending::Escalate(InvalidateReason::WatchOverflow));
886}
887
888fn try_deliver_pending(
889    pending: &mut BTreeMap<PathBuf, Pending>,
890    out: &SyncSender<CoalescedIntent>,
891    sticky_overflow: &mut bool,
892) -> std::result::Result<(), ()> {
893    if pending.is_empty() {
894        return Ok(());
895    }
896    let intent = CoalescedIntent { pending: std::mem::take(pending) };
897    match out.try_send(intent) {
898        Ok(()) => Ok(()),
899        Err(TrySendError::Full(_)) => {
900            *sticky_overflow = true;
901            Ok(())
902        }
903        Err(TrySendError::Disconnected(_)) => Err(()),
904    }
905}
906
907fn try_deliver_overflow(out: &SyncSender<CoalescedIntent>) -> std::result::Result<bool, ()> {
908    let mut pending = BTreeMap::new();
909    collapse_to_overflow(&mut pending);
910    match out.try_send(CoalescedIntent { pending }) {
911        Ok(()) => Ok(true),
912        Err(TrySendError::Full(_)) => Ok(false),
913        Err(TrySendError::Disconnected(_)) => Err(()),
914    }
915}
916
917/// Verify one bounded intent: stat once per path, never once per backend event.
918fn verify_intent(
919    root: &Path,
920    config: WatchConfig,
921    intent: &CoalescedIntent,
922    scan_config: &ScanConfig,
923) -> Observation {
924    let mut ops = Vec::with_capacity(intent.pending.len());
925
926    for (rel, state) in &intent.pending {
927        match state {
928            Pending::Escalate(reason) => {
929                ops.push(Op::InvalidateSubtree { path: rel.clone(), reason: *reason });
930            }
931            Pending::Verify { relist_if_dir } => {
932                let absolute = root.join(rel);
933                match std::fs::symlink_metadata(&absolute) {
934                    Ok(meta) => {
935                        let Ok((kind, attrs)) = scan::observe(&absolute, &meta) else {
936                            ops.push(Op::InvalidateSubtree {
937                                path: rel.parent().map_or_else(PathBuf::new, Path::to_path_buf),
938                                reason: InvalidateReason::VerificationFailed,
939                            });
940                            continue;
941                        };
942                        let disposition = crate::admission::decide_path(
943                            rel,
944                            kind,
945                            scan_config.hidden(),
946                            scan_config.exclude_special,
947                        );
948                        match disposition {
949                            crate::admission::Disposition::Retain => {
950                                ops.push(Op::Upsert { path: rel.clone(), kind, attrs });
951                                match scan::read_control_op(scan_config, root, rel, kind) {
952                                    Ok(Some(control)) => ops.push(control),
953                                    Ok(None) => {}
954                                    Err(_) => ops.push(Op::InvalidateSubtree {
955                                        path: rel
956                                            .parent()
957                                            .map_or_else(PathBuf::new, Path::to_path_buf),
958                                        reason: InvalidateReason::VerificationFailed,
959                                    }),
960                                }
961                            }
962                            crate::admission::Disposition::ControlOnly => {
963                                match scan::read_control_op(scan_config, root, rel, kind) {
964                                    Ok(Some(control)) => {
965                                        ops.push(Op::Remove { path: rel.clone() });
966                                        ops.push(control);
967                                    }
968                                    Ok(None) => ops.push(Op::Remove { path: rel.clone() }),
969                                    Err(_) => ops.push(Op::InvalidateSubtree {
970                                        path: rel
971                                            .parent()
972                                            .map_or_else(PathBuf::new, Path::to_path_buf),
973                                        reason: InvalidateReason::VerificationFailed,
974                                    }),
975                                }
976                            }
977                            crate::admission::Disposition::Reject => {
978                                ops.push(Op::Remove { path: rel.clone() });
979                            }
980                        }
981                        if disposition == crate::admission::Disposition::Retain
982                            && kind.is_dir()
983                            && *relist_if_dir
984                            && config.relist_new_dirs
985                        {
986                            // The watch for this directory was installed after it was
987                            // created, so anything already inside produced no event.
988                            ops.push(Op::InvalidateSubtree {
989                                path: rel.clone(),
990                                reason: InvalidateReason::WatchSetupRace,
991                            });
992                        }
993                    }
994                    Err(error) => ops.push(op_for_stat_error(rel.clone(), &error)),
995                }
996            }
997        }
998    }
999    Observation::new(ops)
1000}
1001
1002fn op_for_stat_error(path: PathBuf, error: &std::io::Error) -> Op {
1003    match error.kind() {
1004        std::io::ErrorKind::NotFound if path.as_os_str().is_empty() => {
1005            Op::InvalidateSubtree { path, reason: InvalidateReason::VerificationFailed }
1006        }
1007        std::io::ErrorKind::NotFound => Op::Remove { path },
1008        std::io::ErrorKind::NotADirectory => Op::InvalidateSubtree {
1009            path: path.parent().map_or_else(PathBuf::new, Path::to_path_buf),
1010            reason: InvalidateReason::VerificationFailed,
1011        },
1012        _ => Op::InvalidateSubtree { path, reason: InvalidateReason::VerificationFailed },
1013    }
1014}
1015
1016/// Express an absolute path relative to the watch root.
1017///
1018/// Returns `None` for anything outside the root, which should not happen but is not
1019/// worth trusting a backend about.
1020fn relative_to(root: &Path, path: &Path) -> Option<PathBuf> {
1021    path.strip_prefix(root).ok().map(Path::to_path_buf)
1022}
1023
1024#[cfg(test)]
1025mod tests {
1026    use super::*;
1027    use notify::event::{CreateKind, Flag, MetadataKind, ModifyKind, RenameMode};
1028    use std::fs;
1029
1030    fn queued_test_watcher(root: PathBuf) -> (SyncSender<CoalescedIntent>, Watcher) {
1031        let (sender, intents) = sync_channel(1);
1032        let watcher = Watcher {
1033            root,
1034            config: WatchConfig::default(),
1035            inner: None,
1036            intents,
1037            control: None,
1038            cancelled: Arc::new(AtomicBool::new(false)),
1039            worker_status: Arc::new(AtomicU8::new(WORKER_RUNNING)),
1040            worker: None,
1041        };
1042        (sender, watcher)
1043    }
1044
1045    #[test]
1046    fn watcher_can_move_to_its_single_consumer_thread() {
1047        fn assert_send<T: Send>() {}
1048
1049        assert_send::<Watcher>();
1050    }
1051
1052    /// Collect deltas until `want` is satisfied or the deadline passes.
1053    ///
1054    /// Event latency varies by orders of magnitude across backends (inotify is
1055    /// immediate, `FSEvents` batches), so the test waits on a condition rather than
1056    /// sleeping for a fixed guess.
1057    /// Serializes the tests that bind a real filesystem watcher.
1058    ///
1059    /// Real-backend delivery latency depends on how much else is happening on the
1060    /// volume. Roughly two hundred tempdirs churn concurrently across this binary, and
1061    /// `FSEvents` is volume-wide, so a stream bound while that is going on can take far
1062    /// longer to deliver its first event than one bound on a quiet machine.
1063    ///
1064    /// Measured at `origin/main`, before any change here, so this is pre-existing
1065    /// backend behavior rather than something a caller introduced:
1066    ///
1067    /// | libtest threads | Result |
1068    /// | --- | --- |
1069    /// | 1 | 406 passed in 8.39 s |
1070    /// | 4 | 3 failed in 22.76 s |
1071    /// | 10 | 3 failed in 21.87 s |
1072    ///
1073    /// Two things follow, and both are applied. Contention between the real-backend
1074    /// tests themselves is removed by this lock, while the rest of the suite keeps
1075    /// running in parallel — the engine is not implicated, since every other watch test
1076    /// drives the same worker through scripted observations and none of them is
1077    /// affected. Residual variance from the surrounding churn is absorbed by
1078    /// [`REAL_BACKEND_DELIVERY`] rather than by retrying.
1079    static REAL_WATCHER: std::sync::Mutex<()> = std::sync::Mutex::new(());
1080
1081    /// How long a real backend may take to deliver its first event before the test fails.
1082    ///
1083    /// The events do arrive; earlier analysis here claimed they did not, inferred from a
1084    /// suite that finished in about the old deadline, and that inference was wrong. With
1085    /// a deadline long enough not to cut delivery off, the full parallel binary passes
1086    /// five runs out of five and finishes in about 18.6 s — *faster* than the runs that
1087    /// failed, because a dead timeout was the longest thing in those.
1088    ///
1089    /// Sixty seconds is chosen to be far outside the observed distribution rather than
1090    /// tuned to its edge, since a deadline that merely covers today's variance becomes
1091    /// tomorrow's flake on a busier machine. The cost is asymmetric and cheap: this
1092    /// duration is only ever spent when delivery genuinely fails, and a passing run
1093    /// never approaches it.
1094    const REAL_BACKEND_DELIVERY: Duration = Duration::from_secs(60);
1095
1096    /// Hold exclusive access to the real watch backend for the rest of the test.
1097    ///
1098    /// Poisoning is deliberately ignored. The lock guards an OS resource rather than
1099    /// shared data, so a panic in one test leaves nothing for the next to observe, and
1100    /// propagating the poison would turn one real failure into three.
1101    fn real_watcher_guard() -> std::sync::MutexGuard<'static, ()> {
1102        REAL_WATCHER.lock().unwrap_or_else(std::sync::PoisonError::into_inner)
1103    }
1104
1105    /// What a wait for events produced.
1106    enum Waited {
1107        /// The awaited ops arrived, with everything seen before them.
1108        Delivered(Vec<Op>),
1109        /// Nothing whatsoever arrived before the deadline.
1110        Silent,
1111    }
1112
1113    /// Collect ops until `want` is satisfied, separating three outcomes that are not the
1114    /// same failure.
1115    ///
1116    /// `Delivered` — the events arrived and matched. The caller's assertions then run at
1117    /// full strength.
1118    ///
1119    /// Panic — events arrived but never matched. That is a real disagreement about
1120    /// content and must fail. This previously returned the partial list instead, so every
1121    /// caller asserted on it and announced a violated product contract when nothing had
1122    /// arrived at all, sending the next reader after a defect that was not there.
1123    ///
1124    /// `Silent` — *nothing whatsoever* arrived before the deadline. What that means
1125    /// depends on whether the watch was already known to deliver, so the caller decides:
1126    /// [`establish_watch`] may read it as the host's silence, and [`wait_established`]
1127    /// reads it as fdu's.
1128    ///
1129    /// The distinction is observable rather than assumed: a working backend delivers the
1130    /// test's own writes within milliseconds, so an empty list after a full minute means
1131    /// the stream is dead, not slow. On this project's macOS development host a degraded
1132    /// `fseventsd` produces exactly that, while both CI platforms never have.
1133    fn wait_for(
1134        watcher: &Watcher,
1135        deadline: Duration,
1136        mut want: impl FnMut(&[Op]) -> bool,
1137    ) -> Waited {
1138        let start = Instant::now();
1139        let mut seen: Vec<Op> = Vec::new();
1140        while start.elapsed() < deadline {
1141            match watcher.next_observation(Duration::from_millis(200)) {
1142                Ok(Some(observation)) => {
1143                    seen.extend(observation.ops.into_iter().map(|observed| observed.op));
1144                    if want(&seen) {
1145                        return Waited::Delivered(seen);
1146                    }
1147                }
1148                Ok(None) => {}
1149                Err(error) => panic!("watcher stopped while waiting: {error}"),
1150            }
1151        }
1152        if seen.is_empty() {
1153            return Waited::Silent;
1154        }
1155        panic!(
1156            "the backend delivered {} op(s) in {deadline:?} but never the one awaited, so \
1157             this is a disagreement about content rather than a delivery failure: {seen:?}",
1158            seen.len()
1159        );
1160    }
1161
1162    /// Wait on a watch that [`establish_watch`] has already proven live.
1163    ///
1164    /// Silence here is evidence about fdu rather than about the host: the backend delivered
1165    /// the warm-up, so a later write it never reports is a lost event. No opt-out applies,
1166    /// which is why the message offers none.
1167    fn wait_established(
1168        watcher: &Watcher,
1169        deadline: Duration,
1170        want: impl FnMut(&[Op]) -> bool,
1171    ) -> Vec<Op> {
1172        match wait_for(watcher, deadline, want) {
1173            Waited::Delivered(ops) => ops,
1174            Waited::Silent => panic!(
1175                "the watch was established and then delivered nothing in {deadline:?}, so \
1176                 this is a lost event rather than a host precondition; \
1177                 FDU_TEST_ALLOW_NO_NATIVE_WATCH does not apply here"
1178            ),
1179        }
1180    }
1181
1182    /// Write into the root until the watch is provably live, then return.
1183    ///
1184    /// A watcher bound to a directory is not yet watching it. `Watcher::new` returns once
1185    /// registration is *requested*, and anything written before it takes effect produces
1186    /// no event — the engine detects exactly this and answers
1187    /// `InvalidateSubtree { path: "", reason: WatchSetupRace }`, meaning "I missed a
1188    /// window, relist the root".
1189    ///
1190    /// That is the correct answer, and it is why these tests were failing intermittently.
1191    /// A test that writes its subject immediately after `Watcher::new` is racing
1192    /// registration: when it loses, the engine reports the race rather than the file, and
1193    /// an assertion waiting for that file's own upsert rejects a valid reply. The failure
1194    /// looked like flakiness, then like a dead backend, and was neither — it was a real
1195    /// race the test set up for itself, and which the engine reported faithfully.
1196    ///
1197    /// Writing a warm-up file and waiting for *any* event settles it: whichever arrives,
1198    /// the stream is delivering and registration is complete, so a write afterwards
1199    /// cannot fall into the setup window. Everything the test then asserts is about fdu.
1200    ///
1201    /// Returns `false` only for a host explicitly declared unable to deliver native watch
1202    /// events, with `FDU_TEST_ALLOW_NO_NATIVE_WATCH=1`. Without that declaration, silence
1203    /// during the warm-up is an actionable precondition failure rather than a passing test.
1204    fn establish_watch(watcher: &Watcher, dir: &Path) -> bool {
1205        let warmup = dir.join(".fdu-watch-warmup");
1206        fs::write(&warmup, b"warmup").expect("warmup write");
1207        let waited = wait_for(watcher, REAL_BACKEND_DELIVERY, |ops| !ops.is_empty());
1208        let _ = fs::remove_file(&warmup);
1209        match waited {
1210            Waited::Delivered(_) => true,
1211            Waited::Silent => {
1212                if std::env::var_os("FDU_TEST_ALLOW_NO_NATIVE_WATCH").as_deref()
1213                    == Some(std::ffi::OsStr::new("1"))
1214                {
1215                    eprintln!(
1216                        "skipped by FDU_TEST_ALLOW_NO_NATIVE_WATCH=1: the host event service \
1217                         delivered no events to this stream"
1218                    );
1219                    return false;
1220                }
1221                panic!(
1222                    "native watch precondition failed: the host delivered no events in \
1223                     {REAL_BACKEND_DELIVERY:?}; run on a host with event delivery, or \
1224                     explicitly opt out with FDU_TEST_ALLOW_NO_NATIVE_WATCH=1"
1225                );
1226            }
1227        }
1228    }
1229
1230    #[test]
1231    fn created_files_arrive_as_verified_upserts() {
1232        let _serialized = real_watcher_guard();
1233        let dir = tempfile::tempdir().expect("tempdir");
1234        let watcher = Watcher::new(dir.path(), WatchConfig::default()).expect("watcher");
1235        if !establish_watch(&watcher, dir.path()) {
1236            return;
1237        }
1238
1239        fs::write(dir.path().join("hello.txt"), b"hello world").expect("write");
1240
1241        let ops = wait_established(&watcher, REAL_BACKEND_DELIVERY, |ops| {
1242            ops.iter().any(|op| op.path() == Path::new("hello.txt"))
1243        });
1244
1245        let found = ops
1246            .iter()
1247            .find(|op| op.path() == Path::new("hello.txt"))
1248            .expect("an op for the new file");
1249        match found {
1250            Op::Upsert { attrs, kind, .. } => {
1251                assert!(!kind.is_dir());
1252                // The point of verify-then-emit: the delta carries real stat data, which
1253                // no backend put in the event.
1254                assert_eq!(attrs.size, 11);
1255                assert!(attrs.mtime_ns > 0);
1256            }
1257            other => panic!("expected an upsert, got {other:?}"),
1258        }
1259    }
1260
1261    #[test]
1262    fn deleted_files_arrive_as_removes() {
1263        let _serialized = real_watcher_guard();
1264        let dir = tempfile::tempdir().expect("tempdir");
1265        let path = dir.path().join("doomed.txt");
1266        fs::write(&path, b"x").expect("write");
1267
1268        let watcher = Watcher::new(dir.path(), WatchConfig::default()).expect("watcher");
1269        if !establish_watch(&watcher, dir.path()) {
1270            return;
1271        }
1272        fs::remove_file(&path).expect("remove");
1273
1274        let ops = wait_established(&watcher, REAL_BACKEND_DELIVERY, |ops| {
1275            ops.iter()
1276                .any(|op| matches!(op, Op::Remove { path } if path == Path::new("doomed.txt")))
1277        });
1278
1279        assert!(
1280            ops.iter()
1281                .any(|op| matches!(op, Op::Remove { path } if path == Path::new("doomed.txt"))),
1282            "expected a remove, saw {ops:?}"
1283        );
1284    }
1285
1286    // The watch-setup race — a directory created and populated before its watch is
1287    // installed, which must escalate to a relist — is asserted by
1288    // `opened::tests::scripted_directory_creation_closes_the_registration_gap`. That test
1289    // drives the same `WatchSetupRace` invalidation and the same subsequent discovery of
1290    // the child through a scripted observation, so the semantics are pinned on every
1291    // platform without depending on when a backend happens to register.
1292    //
1293    // A real-backend version of it lived here and was removed rather than repaired. It
1294    // could not make the claim it appeared to: the macOS backend watches recursively from
1295    // the root, so there is no per-directory registration window for it to lose, and the
1296    // test spent its full twenty-second deadline waiting for an escalation that platform
1297    // has no reason to emit. Serializing the real-watcher tests fixed the other two and
1298    // left this one failing, which is what showed the difference is in the scenario and
1299    // not in the contention.
1300    //
1301    // What stays here is the minimal real-backend smoke the architecture asks for:
1302    // create and remove actually arrive, and a missing root actually fails.
1303
1304    #[test]
1305    fn watching_a_missing_path_is_an_error() {
1306        let dir = tempfile::tempdir().expect("tempdir");
1307        let missing = dir.path().join("not-there");
1308        assert!(Watcher::new(&missing, WatchConfig::default()).is_err());
1309    }
1310
1311    #[test]
1312    fn paths_outside_the_root_are_ignored() {
1313        let root = Path::new("/a/b");
1314        assert_eq!(relative_to(root, Path::new("/a/b/c/d")), Some(PathBuf::from("c/d")));
1315        assert_eq!(relative_to(root, Path::new("/elsewhere")), None);
1316    }
1317
1318    #[test]
1319    fn verification_errors_distinguish_absence_from_an_invalid_ancestor() {
1320        let path = PathBuf::from("parent/known.txt");
1321        let missing = op_for_stat_error(
1322            path.clone(),
1323            &std::io::Error::new(std::io::ErrorKind::NotFound, "gone"),
1324        );
1325        assert!(matches!(missing, Op::Remove { path: removed } if removed == path));
1326
1327        let not_a_directory = op_for_stat_error(
1328            path.clone(),
1329            &std::io::Error::new(std::io::ErrorKind::NotADirectory, "ancestor is a file"),
1330        );
1331        assert!(matches!(
1332            not_a_directory,
1333            Op::InvalidateSubtree {
1334                path: invalidated,
1335                reason: InvalidateReason::VerificationFailed,
1336            } if invalidated == Path::new("parent")
1337        ));
1338
1339        let denied = op_for_stat_error(
1340            path.clone(),
1341            &std::io::Error::new(std::io::ErrorKind::PermissionDenied, "denied"),
1342        );
1343        assert!(matches!(
1344            denied,
1345            Op::InvalidateSubtree {
1346                path: invalidated,
1347                reason: InvalidateReason::VerificationFailed,
1348            } if invalidated == path
1349        ));
1350    }
1351
1352    #[test]
1353    fn unknown_watch_ancestry_reconciles_from_the_nearest_known_directory() {
1354        let dir = tempfile::tempdir().expect("tempdir");
1355        let (index, _) =
1356            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
1357        let handle = crate::IndexHandle::new(index);
1358        let nested = dir.path().join("new/deep");
1359        fs::create_dir_all(&nested).expect("nested directories");
1360        fs::write(nested.join("file.txt"), b"verified").expect("nested file");
1361
1362        let report = apply_observation(
1363            &handle,
1364            &Observation::new(vec![Op::Upsert {
1365                path: PathBuf::from("new/deep/file.txt"),
1366                kind: crate::EntryKind::File,
1367                attrs: crate::Attrs::default(),
1368            }]),
1369            &crate::ScanConfig::default(),
1370            &mut |_| {},
1371        )
1372        .expect("unknown ancestry schedules reconciliation");
1373
1374        assert_eq!(report.apply.invalidated, 1);
1375        assert!(report.reconciliation.is_complete());
1376        assert_eq!(
1377            handle.kind(Path::new("new")).expect("new directory"),
1378            Some(crate::EntryKind::Dir)
1379        );
1380        assert_eq!(
1381            handle.kind(Path::new("new/deep/file.txt")).expect("nested file"),
1382            Some(crate::EntryKind::File)
1383        );
1384        assert_ne!(
1385            handle.attrs(Path::new("new")).expect("verified parent attrs"),
1386            Some(crate::Attrs::default())
1387        );
1388    }
1389
1390    #[test]
1391    fn create_intent_survives_coalescing_but_metadata_only_does_not_relist() {
1392        let root = Path::new("/watch-root");
1393        let path = root.join("directory");
1394        let mut pending = BTreeMap::new();
1395
1396        record(
1397            root,
1398            &notify::Event::new(EventKind::Create(CreateKind::Folder)).add_path(path.clone()),
1399            &mut pending,
1400            16,
1401        );
1402        record(
1403            root,
1404            &notify::Event::new(EventKind::Modify(ModifyKind::Metadata(MetadataKind::Any)))
1405                .add_path(path),
1406            &mut pending,
1407            16,
1408        );
1409        assert_eq!(
1410            pending.get(Path::new("directory")),
1411            Some(&Pending::Verify { relist_if_dir: true })
1412        );
1413
1414        let mut metadata_only = BTreeMap::new();
1415        record(
1416            root,
1417            &notify::Event::new(EventKind::Modify(ModifyKind::Metadata(MetadataKind::Any)))
1418                .add_path(root.join("existing")),
1419            &mut metadata_only,
1420            16,
1421        );
1422        assert_eq!(
1423            metadata_only.get(Path::new("existing")),
1424            Some(&Pending::Verify { relist_if_dir: false })
1425        );
1426    }
1427
1428    #[test]
1429    fn unpaired_renames_and_ambiguous_rescans_escalate_the_root() {
1430        let root = Path::new("/watch-root");
1431        let mut rename_pending = BTreeMap::new();
1432        record(
1433            root,
1434            &notify::Event::new(EventKind::Modify(ModifyKind::Name(RenameMode::From)))
1435                .add_path(root.join("old")),
1436            &mut rename_pending,
1437            16,
1438        );
1439        assert_eq!(
1440            rename_pending.get(Path::new("")),
1441            Some(&Pending::Escalate(InvalidateReason::UnpairedRename))
1442        );
1443
1444        let mut rescan_pending = BTreeMap::new();
1445        record(
1446            root,
1447            &notify::Event::new(EventKind::Any)
1448                .add_path(root.join("a"))
1449                .add_path(root.join("b"))
1450                .set_flag(Flag::Rescan),
1451            &mut rescan_pending,
1452            16,
1453        );
1454        assert_eq!(
1455            rescan_pending.get(Path::new("")),
1456            Some(&Pending::Escalate(InvalidateReason::WatchOverflow))
1457        );
1458    }
1459
1460    #[test]
1461    fn paired_rename_preserves_the_new_directory_relist_intent() {
1462        let root = Path::new("/watch-root");
1463        let mut pending = BTreeMap::new();
1464        record(
1465            root,
1466            &notify::Event::new(EventKind::Modify(ModifyKind::Name(RenameMode::Both)))
1467                .add_path(root.join("old"))
1468                .add_path(root.join("new")),
1469            &mut pending,
1470            16,
1471        );
1472
1473        assert!(!matches!(
1474            pending.get(Path::new("")),
1475            Some(Pending::Escalate(InvalidateReason::UnpairedRename))
1476        ));
1477        assert_eq!(pending.get(Path::new("old")), Some(&Pending::Verify { relist_if_dir: false }));
1478        assert_eq!(pending.get(Path::new("new")), Some(&Pending::Verify { relist_if_dir: true }));
1479    }
1480
1481    #[test]
1482    fn pending_path_overload_collapses_to_one_root_invalidation() {
1483        let root = Path::new("/watch-root");
1484        let mut pending = BTreeMap::new();
1485        for name in ["one", "two", "three"] {
1486            record(
1487                root,
1488                &notify::Event::new(EventKind::Any).add_path(root.join(name)),
1489                &mut pending,
1490                2,
1491            );
1492        }
1493
1494        assert_eq!(pending.len(), 1);
1495        assert_eq!(
1496            pending.get(Path::new("")),
1497            Some(&Pending::Escalate(InvalidateReason::WatchOverflow))
1498        );
1499    }
1500
1501    #[test]
1502    fn continuous_churn_has_a_deterministic_max_hold_ceiling() {
1503        let past = Instant::now()
1504            .checked_sub(Duration::from_secs(2))
1505            .expect("representable earlier instant");
1506        assert!(max_hold_elapsed(Some(past), Duration::from_secs(1)));
1507        assert!(!max_hold_elapsed(None, Duration::from_secs(1)));
1508    }
1509
1510    #[test]
1511    fn backend_enqueue_is_nonblocking_and_marks_overflow() {
1512        let (sender, receiver) = sync_channel(1);
1513        let overflowed = AtomicBool::new(false);
1514        enqueue_raw(&sender, &overflowed, Ok(notify::Event::new(EventKind::Any)));
1515        enqueue_raw(&sender, &overflowed, Ok(notify::Event::new(EventKind::Any)));
1516
1517        assert!(overflowed.load(Ordering::Acquire));
1518        assert!(matches!(receiver.try_recv(), Ok(RawMessage::Event(Ok(_)))));
1519    }
1520
1521    #[test]
1522    fn full_intent_queue_retains_a_sticky_root_invalidation() {
1523        let (sender, receiver) = sync_channel(1);
1524        sender.try_send(CoalescedIntent::default()).expect("fill output");
1525        let mut pending =
1526            BTreeMap::from([(PathBuf::from("lost.txt"), Pending::Verify { relist_if_dir: false })]);
1527        let mut sticky_overflow = false;
1528
1529        try_deliver_pending(&mut pending, &sender, &mut sticky_overflow).expect("connected");
1530        assert!(pending.is_empty());
1531        assert!(sticky_overflow);
1532
1533        receiver.try_recv().expect("make output capacity");
1534        assert!(try_deliver_overflow(&sender).expect("connected"));
1535        let intent = receiver.try_recv().expect("sticky overflow intent");
1536        let observation = verify_intent(
1537            Path::new("/unused"),
1538            WatchConfig::default(),
1539            &intent,
1540            &ScanConfig::default(),
1541        );
1542        assert!(matches!(
1543            &observation.ops[0].op,
1544            Op::InvalidateSubtree {
1545                path,
1546                reason: InvalidateReason::WatchOverflow,
1547            } if path.as_os_str().is_empty()
1548        ));
1549    }
1550
1551    #[test]
1552    fn cancellation_wakes_and_joins_with_a_full_intent_queue() {
1553        let dir = tempfile::tempdir().expect("tempdir");
1554        let root = dir.path().canonicalize().expect("canonical root");
1555        let config = WatchConfig {
1556            settle: Duration::from_secs(30),
1557            max_hold: Duration::from_secs(30),
1558            event_capacity: 1,
1559            batch_path_capacity: 1,
1560            intent_capacity: 1,
1561            ..WatchConfig::default()
1562        };
1563        let (control, raw) = sync_channel(1);
1564        let (output, intents) = sync_channel(1);
1565        output.try_send(CoalescedIntent::default()).expect("fill intent queue");
1566        let cancelled = Arc::new(AtomicBool::new(false));
1567        let worker_cancelled = Arc::clone(&cancelled);
1568        let status = Arc::new(AtomicU8::new(WORKER_RUNNING));
1569        let tracked_status = Arc::clone(&status);
1570        let overflowed = Arc::new(AtomicBool::new(false));
1571        let worker_overflowed = Arc::clone(&overflowed);
1572        let worker_root = root.clone();
1573        let worker = std::thread::spawn(move || {
1574            run_tracked_worker(&tracked_status, || {
1575                run_worker(
1576                    &worker_root,
1577                    config,
1578                    &raw,
1579                    &output,
1580                    &worker_overflowed,
1581                    &worker_cancelled,
1582                );
1583            });
1584        });
1585        let watcher = Watcher {
1586            root,
1587            config,
1588            inner: None,
1589            intents,
1590            control: Some(control),
1591            cancelled,
1592            worker_status: status,
1593            worker: Some(worker),
1594        };
1595        let (done_tx, done_rx) = sync_channel(1);
1596
1597        std::thread::spawn(move || {
1598            drop(watcher);
1599            done_tx.send(()).expect("report drop");
1600        });
1601
1602        done_rx
1603            .recv_timeout(Duration::from_secs(5))
1604            .expect("watcher drop must wake and join promptly");
1605    }
1606
1607    #[test]
1608    fn coalescing_defers_filesystem_verification_to_the_consumer() {
1609        let dir = tempfile::tempdir().expect("tempdir");
1610        let relative = PathBuf::from("appeared.txt");
1611        let intent = CoalescedIntent {
1612            pending: BTreeMap::from([(relative.clone(), Pending::Verify { relist_if_dir: false })]),
1613        };
1614
1615        fs::write(dir.path().join(&relative), b"current").expect("create after coalescing");
1616        let observation =
1617            verify_intent(dir.path(), WatchConfig::default(), &intent, &ScanConfig::default());
1618
1619        assert!(matches!(
1620            &observation.ops[0].op,
1621            Op::Upsert { path, attrs, .. } if path == &relative && attrs.size == 7
1622        ));
1623    }
1624
1625    #[test]
1626    fn control_verification_emits_exact_source_with_the_entry_fact() {
1627        let dir = tempfile::tempdir().expect("tempdir");
1628        let relative = PathBuf::from(".gitignore");
1629        fs::write(dir.path().join(&relative), b"*.log\n").expect("write control");
1630        let intent = CoalescedIntent {
1631            pending: BTreeMap::from([(relative.clone(), Pending::Verify { relist_if_dir: false })]),
1632        };
1633
1634        let config = ScanConfig { read_controls: true, ..ScanConfig::default() };
1635        let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
1636
1637        assert!(matches!(
1638            &observation.ops[0].op,
1639            Op::Upsert { path, kind: crate::EntryKind::File, .. } if path == &relative
1640        ));
1641        assert!(matches!(
1642            &observation.ops[1].op,
1643            Op::ControlUpsert { path, source } if path == &relative && source == b"*.log\n"
1644        ));
1645    }
1646
1647    /// A watch maintains exactly the control state its scan policy claims.
1648    ///
1649    /// A controls-off scan retains no control table and stamps that into its scope, and a
1650    /// watch with the same policy accepts the scope as its own. Verification that read
1651    /// control files regardless would grow a partial rule set -- only the sources some
1652    /// event happened to touch -- on an index whose scope says it has none, and the next
1653    /// save would persist it under that scope (fdu-ajsu).
1654    #[test]
1655    fn verification_observes_no_control_state_under_a_controls_off_policy() {
1656        let dir = tempfile::tempdir().expect("tempdir");
1657        fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("write control");
1658        fs::write(dir.path().join("debug.log"), b"x").expect("write file");
1659        let config = ScanConfig { read_controls: false, ..ScanConfig::default() };
1660        let (mut index, _) = crate::scan::scan_into_index(dir.path(), &config).expect("scan");
1661        assert!(index.control_table().is_empty());
1662        assert!(matches!(index.controls(), Err(crate::Error::ControlStateNotObserved)));
1663        let control = PathBuf::from(".gitignore");
1664        let intent = CoalescedIntent {
1665            pending: BTreeMap::from([(control.clone(), Pending::Verify { relist_if_dir: false })]),
1666        };
1667
1668        let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
1669        let reverified = reverify_observation(
1670            dir.path(),
1671            &Observation::new(vec![Op::Remove { path: control }]),
1672            &config,
1673        )
1674        .expect("reverify");
1675
1676        for verified in [&observation, &reverified] {
1677            assert!(
1678                !verified.ops.iter().any(|observed| matches!(
1679                    observed.op,
1680                    Op::ControlUpsert { .. } | Op::ControlRemove { .. }
1681                )),
1682                "a controls-off policy observed control state: {:?}",
1683                verified.ops
1684            );
1685        }
1686        index.apply(&observation).expect("apply the verified observation");
1687        assert!(index.control_table().is_empty());
1688        assert_eq!(index.scope(), config.scope());
1689    }
1690
1691    #[cfg(unix)]
1692    #[test]
1693    fn applying_verification_uses_the_index_admission_scope() {
1694        use std::os::unix::net::UnixListener;
1695
1696        let dir = tempfile::tempdir().expect("tempdir");
1697        fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("write control");
1698        fs::write(dir.path().join(".secret"), b"hidden").expect("write hidden");
1699        let _listener = UnixListener::bind(dir.path().join("service.sock")).expect("bind socket");
1700        let intent = CoalescedIntent {
1701            pending: [".gitignore", ".secret", "service.sock"]
1702                .into_iter()
1703                .map(|path| (PathBuf::from(path), Pending::Verify { relist_if_dir: false }))
1704                .collect(),
1705        };
1706        let config = ScanConfig {
1707            hidden: Some(Arc::new(crate::HiddenPolicy::prune_hidden::<[&str; 0], &str>([]))),
1708            exclude_special: true,
1709            read_controls: true,
1710            ..ScanConfig::default()
1711        };
1712
1713        let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
1714
1715        assert!(observation.ops.iter().any(|observed| matches!(
1716            &observed.op,
1717            Op::ControlUpsert { path, source }
1718                if path == Path::new(".gitignore") && source == b"*.log\n"
1719        )));
1720        for path in [".gitignore", ".secret", "service.sock"] {
1721            assert!(observation.ops.iter().any(|observed| matches!(
1722                &observed.op,
1723                Op::Remove { path: removed } if removed == Path::new(path)
1724            )));
1725            assert!(!observation.ops.iter().any(|observed| matches!(
1726                &observed.op,
1727                Op::Upsert { path: retained, .. } if retained == Path::new(path)
1728            )));
1729        }
1730    }
1731
1732    #[test]
1733    fn timeout_stop_and_worker_panic_are_distinct() {
1734        let dir = tempfile::tempdir().expect("tempdir");
1735        let (live_sender, live) =
1736            queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
1737        assert!(live.next_observation(Duration::ZERO).expect("timeout").is_none());
1738        drop(live_sender);
1739
1740        let (stopped_sender, stopped) =
1741            queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
1742        stopped.worker_status.store(WORKER_STOPPED, Ordering::Release);
1743        drop(stopped_sender);
1744        assert!(matches!(stopped.next_observation(Duration::ZERO), Err(Error::WatchStopped)));
1745
1746        let (panicked_sender, panicked) =
1747            queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
1748        panicked.worker_status.store(WORKER_PANICKED, Ordering::Release);
1749        drop(panicked_sender);
1750        assert!(matches!(
1751            panicked.next_observation(Duration::ZERO),
1752            Err(Error::WatchWorkerPanicked)
1753        ));
1754    }
1755
1756    #[test]
1757    fn tracked_worker_records_a_panic() {
1758        let status = AtomicU8::new(WORKER_RUNNING);
1759        run_tracked_worker(&status, || panic!("injected worker panic"));
1760        assert_eq!(status.load(Ordering::Acquire), WORKER_PANICKED);
1761    }
1762
1763    #[test]
1764    fn observation_driver_closes_the_invalidation_loop() {
1765        let dir = tempfile::tempdir().expect("tempdir");
1766        let (index, _) =
1767            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
1768        let handle = crate::IndexHandle::new(index);
1769        fs::write(dir.path().join("raced.txt"), b"raced").expect("write");
1770        let observation = Observation::new(vec![Op::InvalidateSubtree {
1771            path: PathBuf::new(),
1772            reason: InvalidateReason::WatchSetupRace,
1773        }]);
1774
1775        apply_observation(&handle, &observation, &crate::ScanConfig::default(), &mut |_| {})
1776            .expect("apply and reconcile");
1777
1778        assert!(handle.kind(Path::new("raced.txt")).expect("query").is_some());
1779    }
1780
1781    #[test]
1782    fn applying_driver_reverifies_a_queued_sample_after_reconciliation() {
1783        let dir = tempfile::tempdir().expect("tempdir");
1784        let path = dir.path().join("sample.txt");
1785        fs::write(&path, b"old").expect("write old sample");
1786        let (index, _) =
1787            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
1788        let old_attrs = *index.attrs(Path::new("sample.txt")).expect("sample attributes");
1789        let delayed = Observation::new(vec![Op::Upsert {
1790            path: PathBuf::from("sample.txt"),
1791            kind: crate::EntryKind::File,
1792            attrs: old_attrs,
1793        }]);
1794        let handle = crate::IndexHandle::new(index);
1795
1796        fs::write(&path, b"new contents").expect("write current sample");
1797        crate::scan::reconcile_handle(&handle, &crate::ScanConfig::default(), &mut |_| {})
1798            .expect("reconcile newer sample");
1799        let current_size = fs::metadata(&path).expect("sample metadata").len();
1800
1801        apply_observation(&handle, &delayed, &crate::ScanConfig::default(), &mut |_| {})
1802            .expect("apply delayed watch sample");
1803
1804        assert_eq!(
1805            handle.attrs(Path::new("sample.txt")).expect("query").expect("sample remains").size,
1806            current_size
1807        );
1808    }
1809
1810    #[test]
1811    fn blocked_verifier_holds_no_index_lock_and_commits_only_at_current_clock() {
1812        let dir = tempfile::tempdir().expect("tempdir");
1813        let (index, _) =
1814            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
1815        let handle = crate::IndexHandle::new(index);
1816        let queued = Observation::new(vec![Op::Upsert {
1817            path: PathBuf::from("queued.txt"),
1818            kind: crate::EntryKind::File,
1819            attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
1820        }]);
1821        let applying = handle.clone();
1822        let (entered_tx, entered_rx) = sync_channel(1);
1823        let (release_tx, release_rx) = sync_channel(1);
1824        let (done_tx, done_rx) = sync_channel(1);
1825        let apply_thread = std::thread::spawn(move || {
1826            let mut first = true;
1827            let mut verifier = |_: &Path, observation: &Observation| {
1828                if first {
1829                    first = false;
1830                    entered_tx.send(()).expect("signal blocked verifier");
1831                    release_rx.recv().expect("release blocked verifier");
1832                }
1833                Ok(observation.clone())
1834            };
1835            let result = apply_reverified_with(
1836                &applying,
1837                &queued,
1838                &crate::ScanConfig::default(),
1839                &mut verifier,
1840            );
1841            done_tx.send(result).expect("report apply result");
1842        });
1843
1844        entered_rx
1845            .recv_timeout(Duration::from_secs(5))
1846            .expect("verifier must reach the injected block");
1847        let progressing = handle.clone();
1848        let (progress_tx, progress_rx) = sync_channel(1);
1849        let progress_thread = std::thread::spawn(move || {
1850            let total = progressing.total().expect("reader progresses");
1851            let write = progressing.apply(&Observation::new(vec![Op::Upsert {
1852                path: PathBuf::from("competitor.txt"),
1853                kind: crate::EntryKind::File,
1854                attrs: crate::Attrs { size: 3, allocated: 3, ..crate::Attrs::default() },
1855            }]));
1856            progress_tx.send((total, write)).expect("report progress");
1857        });
1858        let (_, competing_write) = progress_rx
1859            .recv_timeout(Duration::from_secs(5))
1860            .expect("reader and writer must progress while verification is blocked");
1861        competing_write.expect("competing write");
1862        release_tx.send(()).expect("release verifier");
1863
1864        let outcome = done_rx
1865            .recv_timeout(Duration::from_secs(5))
1866            .expect("applying driver completes")
1867            .expect("applying driver succeeds");
1868        apply_thread.join().expect("apply thread");
1869        progress_thread.join().expect("progress thread");
1870        assert_eq!(outcome.inserted, 1);
1871        assert!(handle.kind(Path::new("competitor.txt")).expect("query").is_some());
1872        assert!(handle.kind(Path::new("queued.txt")).expect("query").is_some());
1873        assert_eq!(handle.clock().expect("clock"), crate::Clock(2));
1874    }
1875
1876    #[test]
1877    fn exhausted_watch_contention_stays_unfresh_until_reconciliation() {
1878        let dir = tempfile::tempdir().expect("tempdir");
1879        let (index, _) =
1880            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
1881        let handle = crate::IndexHandle::new(index);
1882        let queued = Observation::new(vec![Op::Upsert {
1883            path: PathBuf::from("never-committed.txt"),
1884            kind: crate::EntryKind::File,
1885            attrs: crate::Attrs { size: 1, allocated: 1, ..crate::Attrs::default() },
1886        }]);
1887        let mut attempts = 0_usize;
1888        let mut verifier = |_: &Path, observation: &Observation| {
1889            attempts += 1;
1890            handle
1891                .apply(&Observation::new(vec![Op::Upsert {
1892                    path: PathBuf::from(format!("competitor-{attempts}.txt")),
1893                    kind: crate::EntryKind::File,
1894                    attrs: crate::Attrs {
1895                        size: attempts as u64,
1896                        allocated: attempts as u64,
1897                        ..crate::Attrs::default()
1898                    },
1899                }]))
1900                .expect("force a clock conflict");
1901            Ok(observation.clone())
1902        };
1903
1904        let outcome =
1905            apply_reverified_with(&handle, &queued, &crate::ScanConfig::default(), &mut verifier)
1906                .expect("contention escalates");
1907
1908        assert_eq!(attempts, MAX_OPTIMISTIC_APPLY_ATTEMPTS);
1909        assert_eq!(outcome.invalidated, 1);
1910        assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Stale);
1911        assert!(handle.kind(Path::new("never-committed.txt")).expect("query").is_none());
1912        let pending = handle.take_pending_invalidations().expect("pending invalidation");
1913        assert_eq!(pending, vec![(PathBuf::new(), InvalidateReason::WatchContention)]);
1914        handle.restore_pending_invalidations(pending).expect("restore invalidation");
1915
1916        crate::scan::reconcile_pending_handle(&handle, &crate::ScanConfig::default(), &mut |_| {})
1917            .expect("reconcile contention");
1918        assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Fresh);
1919    }
1920
1921    #[test]
1922    fn verifier_error_mutates_no_shared_state() {
1923        let dir = tempfile::tempdir().expect("tempdir");
1924        let (index, _) =
1925            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
1926        let handle = crate::IndexHandle::new(index);
1927        let before_clock = handle.clock().expect("clock");
1928        let before_total = handle.total().expect("total");
1929        let mut verifier = |_: &Path, _: &Observation| {
1930            Err(Error::io(
1931                PathBuf::from("blocked"),
1932                std::io::Error::new(std::io::ErrorKind::PermissionDenied, "injected"),
1933            ))
1934        };
1935
1936        let error = apply_reverified_with(
1937            &handle,
1938            &Observation::default(),
1939            &crate::ScanConfig::default(),
1940            &mut verifier,
1941        )
1942        .expect_err("verification error");
1943
1944        assert!(matches!(error, Error::Io { .. }));
1945        assert_eq!(handle.clock().expect("clock"), before_clock);
1946        assert_eq!(handle.total().expect("total"), before_total);
1947        assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Fresh);
1948        assert!(handle.take_pending_invalidations().expect("pending").is_empty());
1949    }
1950
1951    #[test]
1952    fn stable_watch_arbitration_verifies_exactly_once() {
1953        let dir = tempfile::tempdir().expect("tempdir");
1954        let (index, _) =
1955            crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
1956        let handle = crate::IndexHandle::new(index);
1957        let mut calls = 0_u8;
1958        let mut verifier = |_: &Path, observation: &Observation| {
1959            calls += 1;
1960            Ok(observation.clone())
1961        };
1962
1963        apply_reverified_with(
1964            &handle,
1965            &Observation::new(vec![Op::InvalidateSubtree {
1966                path: PathBuf::new(),
1967                reason: InvalidateReason::Requested,
1968            }]),
1969            &crate::ScanConfig::default(),
1970            &mut verifier,
1971        )
1972        .expect("stable apply");
1973
1974        assert_eq!(calls, 1);
1975    }
1976
1977    #[test]
1978    fn disappearing_watch_root_escalates_instead_of_removing_the_index_root() {
1979        let error = std::io::Error::new(std::io::ErrorKind::NotFound, "root disappeared");
1980
1981        assert!(matches!(
1982            op_for_stat_error(PathBuf::new(), &error),
1983            Op::InvalidateSubtree {
1984                path,
1985                reason: InvalidateReason::VerificationFailed,
1986            } if path.as_os_str().is_empty()
1987        ));
1988    }
1989
1990    #[test]
1991    fn observation_driver_rejects_scope_mismatch_before_apply() {
1992        let dir = tempfile::tempdir().expect("tempdir");
1993        let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
1994        let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
1995        let handle = crate::IndexHandle::new(index);
1996        let observation = Observation::new(vec![Op::Upsert {
1997            path: PathBuf::from("deep/nested.txt"),
1998            kind: crate::EntryKind::File,
1999            attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2000        }]);
2001
2002        let error =
2003            apply_observation(&handle, &observation, &crate::ScanConfig::default(), &mut |_| {})
2004                .expect_err("mismatched scope must fail");
2005
2006        assert!(matches!(error, Error::ScanScopeMismatch { .. }));
2007        assert!(handle.kind(Path::new("deep/nested.txt")).expect("query").is_none());
2008    }
2009
2010    #[test]
2011    fn observation_driver_rejects_restricted_scopes_until_events_are_filtered() {
2012        let dir = tempfile::tempdir().expect("tempdir");
2013        let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2014        let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2015        let handle = crate::IndexHandle::new(index);
2016        let observation = Observation::new(vec![Op::Upsert {
2017            path: PathBuf::from("deep/nested.txt"),
2018            kind: crate::EntryKind::File,
2019            attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2020        }]);
2021
2022        let error = apply_observation(&handle, &observation, &shallow, &mut |_| {})
2023            .expect_err("unfiltered bounded watch scope must fail");
2024
2025        assert!(matches!(error, Error::UnsupportedScanConfig(_)));
2026        assert!(handle.kind(Path::new("deep/nested.txt")).expect("query").is_none());
2027    }
2028
2029    #[test]
2030    fn apply_next_rejects_restricted_scope_without_consuming_an_observation() {
2031        let dir = tempfile::tempdir().expect("tempdir");
2032        let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2033        let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2034        let handle = crate::IndexHandle::new(index);
2035        let (sender, watcher) =
2036            queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2037        sender.try_send(CoalescedIntent::default()).expect("queue intent");
2038
2039        let error = watcher
2040            .apply_next(&handle, &shallow, Duration::ZERO, &mut |_| {})
2041            .expect_err("restricted scope must fail before receive");
2042
2043        assert!(matches!(error, Error::UnsupportedScanConfig(_)));
2044        assert!(watcher.next_observation(Duration::ZERO).expect("receive").is_some());
2045    }
2046
2047    #[test]
2048    fn apply_next_rejects_a_watcher_for_another_root_without_consuming() {
2049        let indexed = tempfile::tempdir().expect("indexed root");
2050        let watched_root_dir = tempfile::tempdir().expect("watched root");
2051        let (index, _) =
2052            crate::scan::scan_into_index(indexed.path(), &crate::ScanConfig::default())
2053                .expect("scan indexed root");
2054        let handle = crate::IndexHandle::new(index);
2055        let (sender, watcher) = queued_test_watcher(
2056            watched_root_dir.path().canonicalize().expect("canonical watched root"),
2057        );
2058        sender.try_send(CoalescedIntent::default()).expect("queue intent");
2059
2060        let error = watcher
2061            .apply_next(&handle, &crate::ScanConfig::default(), Duration::ZERO, &mut |_| {})
2062            .expect_err("mismatched root must fail");
2063
2064        assert!(matches!(error, Error::WatchRootMismatch { .. }));
2065        assert!(watcher.next_observation(Duration::ZERO).expect("receive").is_some());
2066    }
2067
2068    /// A gap over an unreadable directory is walked once by the per-event driver, which
2069    /// then goes quiet and keeps the cause.
2070    ///
2071    /// `fdu --watch` and the Python watch session drain the invalidation queue after every
2072    /// event through [`Watcher::apply_next`]. While an incomplete reconciliation was queued
2073    /// again, each unrelated event re-walked the same unreadable subtree -- for a root
2074    /// escalation, a full-tree walk per event, for the life of the session.
2075    #[cfg(unix)]
2076    #[test]
2077    fn apply_next_walks_an_unreadable_gap_once_and_retains_its_cause() {
2078        use std::os::unix::fs::PermissionsExt;
2079
2080        fn walks_of(commits: &[Commit], path: &Path) -> usize {
2081            commits
2082                .iter()
2083                .flat_map(|commit| commit.state.iter())
2084                .filter(|transition| {
2085                    matches!(
2086                        transition,
2087                        crate::StateTransition::Freshness { path: marked, current, .. }
2088                            if marked == path && *current == crate::Freshness::Reconciling
2089                    )
2090                })
2091                .count()
2092        }
2093
2094        if !crate::test_support::require_permission_bits() {
2095            return;
2096        }
2097        let dir = tempfile::tempdir().expect("tempdir");
2098        let root = dir.path().canonicalize().expect("canonical root");
2099        let blocked = root.join("blocked");
2100        fs::create_dir(&blocked).expect("blocked");
2101        fs::write(blocked.join("secret"), b"s").expect("fixture");
2102        let (index, _) =
2103            crate::scan::scan_into_index(&root, &crate::ScanConfig::default()).expect("scan");
2104        let handle = crate::IndexHandle::new(index);
2105        let (sender, watcher) = queued_test_watcher(root.clone());
2106        let config = crate::ScanConfig::default();
2107        let mut commits = Vec::new();
2108
2109        fs::set_permissions(&blocked, fs::Permissions::from_mode(0o000)).expect("deny reads");
2110        let mut gap = CoalescedIntent::default();
2111        gap.pending
2112            .insert(PathBuf::from("blocked"), Pending::Escalate(InvalidateReason::WatchOverflow));
2113        sender.try_send(gap).expect("queue the gap");
2114        let first = watcher
2115            .apply_next(&handle, &config, Duration::ZERO, &mut |commit| {
2116                commits.push(commit.clone());
2117            })
2118            .expect("apply the gap")
2119            .expect("an intent was queued");
2120        for name in ["live.txt", "marker.txt"] {
2121            fs::write(root.join(name), name).expect("unrelated mutation");
2122            let mut event = CoalescedIntent::default();
2123            event.pending.insert(PathBuf::from(name), Pending::Verify { relist_if_dir: false });
2124            sender.try_send(event).expect("queue the event");
2125            watcher
2126                .apply_next(&handle, &config, Duration::ZERO, &mut |commit| {
2127                    commits.push(commit.clone());
2128                })
2129                .expect("apply the event")
2130                .expect("an intent was queued");
2131        }
2132
2133        let pending = handle.take_pending_invalidations().expect("pending");
2134        let freshness = handle.freshness_at(Path::new("blocked")).expect("freshness");
2135        let issues = handle.issues().expect("issues");
2136        fs::set_permissions(&blocked, fs::Permissions::from_mode(0o700)).expect("restore reads");
2137        assert!(!first.reconciliation.is_complete(), "the gap must be unreadable");
2138
2139        assert_eq!(
2140            walks_of(&commits, Path::new("blocked")),
2141            1,
2142            "an unreadable subtree must not be re-walked per unrelated event"
2143        );
2144        assert!(pending.is_empty(), "the queue must settle: {pending:?}");
2145        assert_eq!(freshness, crate::Freshness::Partial);
2146        assert!(handle.kind(Path::new("marker.txt")).expect("lookup").is_some());
2147        assert!(
2148            issues.iter().any(|issue| issue.kind == crate::IssueKind::Permission
2149                && issue.path.as_deref() == Some(Path::new("blocked"))),
2150            "{issues:?}"
2151        );
2152    }
2153
2154    #[test]
2155    fn zero_settle_is_rejected_before_starting_a_busy_worker() {
2156        let config = WatchConfig { settle: Duration::ZERO, ..WatchConfig::default() };
2157
2158        assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2159
2160        let config = WatchConfig { intent_capacity: 0, ..WatchConfig::default() };
2161        assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2162
2163        let config = WatchConfig {
2164            batch_path_capacity: MAX_BATCH_PATH_CAPACITY,
2165            intent_capacity: MAX_BUFFERED_INTENT_PATHS / MAX_BATCH_PATH_CAPACITY + 1,
2166            ..WatchConfig::default()
2167        };
2168        assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2169    }
2170}