1use std::collections::BTreeMap;
33use std::panic::{AssertUnwindSafe, catch_unwind};
34use std::path::{Path, PathBuf};
35use std::sync::Arc;
36use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
37use std::sync::mpsc::{Receiver, RecvTimeoutError, SyncSender, TrySendError, sync_channel};
38use std::thread::JoinHandle;
39use std::time::{Duration, Instant};
40
41use notify::{EventKind, RecommendedWatcher, RecursiveMode, Watcher as NotifyWatcher};
42
43#[cfg(test)]
44mod scripted_events;
45
46use crate::engine_contract::{
47 Commit, Error, InvalidateReason, Observation, ObservationOp, Op, Result,
48};
49use crate::scan;
50use crate::{ApplyOutcome, IndexHandle, ScanConfig};
51
52const MAX_OPTIMISTIC_APPLY_ATTEMPTS: usize = 3;
55
56const WORKER_RUNNING: u8 = 0;
57const WORKER_STOPPED: u8 = 1;
58const WORKER_PANICKED: u8 = 2;
59const MAX_EVENT_CAPACITY: usize = 64 * 1024;
60const MAX_BATCH_PATH_CAPACITY: usize = 64 * 1024;
61const MAX_BUFFERED_INTENT_PATHS: usize = 1024 * 1024;
62
63#[derive(Clone, Copy, Debug)]
65pub struct WatchConfig {
66 pub settle: Duration,
68 pub max_hold: Duration,
71 pub relist_new_dirs: bool,
77 pub event_capacity: usize,
80 pub batch_path_capacity: usize,
82 pub intent_capacity: usize,
84}
85
86impl Default for WatchConfig {
87 fn default() -> Self {
88 Self {
89 settle: Duration::from_millis(50),
92 max_hold: Duration::from_millis(1600),
93 relist_new_dirs: true,
94 event_capacity: 4096,
95 batch_path_capacity: 4096,
96 intent_capacity: 16,
97 }
98 }
99}
100
101impl WatchConfig {
102 fn validate(self) -> Result<()> {
103 let buffered_paths = self.batch_path_capacity.checked_mul(self.intent_capacity);
104 if self.settle.is_zero()
105 || self.max_hold.is_zero()
106 || self.max_hold < self.settle
107 || self.event_capacity == 0
108 || self.batch_path_capacity == 0
109 || self.intent_capacity == 0
110 || self.event_capacity > MAX_EVENT_CAPACITY
111 || self.batch_path_capacity > MAX_BATCH_PATH_CAPACITY
112 || buffered_paths.is_none_or(|paths| paths > MAX_BUFFERED_INTENT_PATHS)
113 {
114 return Err(Error::UnsupportedScanConfig(
115 "watch durations and capacities exceed the supported nonzero bounds, or max_hold is less than settle",
116 ));
117 }
118 Ok(())
119 }
120}
121
122#[derive(Clone, Copy, PartialEq, Eq, Debug)]
124enum Pending {
125 Verify {
128 relist_if_dir: bool,
131 renamed: bool,
141 },
142 Escalate(InvalidateReason),
144}
145
146#[derive(Clone, Copy, PartialEq, Eq, Debug)]
165enum RenameReporting {
166 EachSide,
168 OldSideOnly,
170}
171
172impl RenameReporting {
173 fn of_recommended_backend() -> Self {
178 match <RecommendedWatcher as NotifyWatcher>::kind() {
179 notify::WatcherKind::Fsevent
180 | notify::WatcherKind::Inotify
181 | notify::WatcherKind::ReadDirectoryChangesWatcher => Self::EachSide,
182 _ => Self::OldSideOnly,
183 }
184 }
185}
186
187#[derive(Debug, Default)]
188struct CoalescedIntent {
189 pending: BTreeMap<PathBuf, Pending>,
190}
191
192enum RawMessage {
193 Event(notify::Result<notify::Event>),
194 Flush(SyncSender<()>),
195 Stop,
196}
197
198pub struct Watcher {
205 root: PathBuf,
206 config: WatchConfig,
207 inner: Option<RecommendedWatcher>,
209 intents: Receiver<CoalescedIntent>,
210 control: Option<SyncSender<RawMessage>>,
211 cancelled: Arc<AtomicBool>,
212 worker_status: Arc<AtomicU8>,
213 worker: Option<JoinHandle<()>>,
214}
215
216#[cfg(test)]
217pub(crate) struct ScriptedSender {
218 root: PathBuf,
219 raw: SyncSender<RawMessage>,
220 overflowed: Arc<AtomicBool>,
221}
222
223#[cfg(test)]
224impl ScriptedSender {
225 pub(crate) fn send(&self, source: &str) -> Result<()> {
226 let events =
227 scripted_events::parse_script(source, &self.root).map_err(Error::WatchScript)?;
228 for event in events {
229 enqueue_raw(&self.raw, &self.overflowed, event);
230 }
231 Ok(())
232 }
233}
234
235#[derive(Debug)]
237pub struct WatchApplyReport {
238 pub apply: ApplyOutcome,
240 pub reconciliation: scan::ReconcileReport,
242}
243
244impl Watcher {
245 pub fn new(root: &Path, config: WatchConfig) -> Result<Self> {
247 config.validate()?;
248 let root = root.canonicalize().map_err(|e| Error::io(root, e))?;
249
250 let (raw_tx, raw_rx) = sync_channel::<RawMessage>(config.event_capacity);
251 let (intent_tx, intent_rx) = sync_channel::<CoalescedIntent>(config.intent_capacity);
252 let control_tx = raw_tx.clone();
253 let overflowed = Arc::new(AtomicBool::new(false));
254 let callback_overflowed = Arc::clone(&overflowed);
255
256 let mut inner = notify::recommended_watcher(move |result| {
257 enqueue_raw(&raw_tx, &callback_overflowed, result);
258 })
259 .map_err(|error| notify_error(&root, error))?;
260 inner.watch(&root, RecursiveMode::Recursive).map_err(|error| notify_error(&root, error))?;
261
262 let worker_root = root.clone();
263 let cancelled = Arc::new(AtomicBool::new(false));
264 let worker_cancelled = Arc::clone(&cancelled);
265 let worker_overflowed = Arc::clone(&overflowed);
266 let worker_status = Arc::new(AtomicU8::new(WORKER_RUNNING));
267 let tracked_status = Arc::clone(&worker_status);
268 let renames = RenameReporting::of_recommended_backend();
269 let worker = std::thread::Builder::new()
270 .name("fdu-watch".into())
271 .spawn(move || {
272 let _counter_guard = crate::counters::thread_flush_guard();
273 run_tracked_worker(&tracked_status, || {
274 run_worker(
275 &worker_root,
276 config,
277 renames,
278 &raw_rx,
279 &intent_tx,
280 &worker_overflowed,
281 &worker_cancelled,
282 );
283 });
284 })
285 .map_err(|e| Error::io(&root, e))?;
286
287 Ok(Self {
288 root,
289 config,
290 inner: Some(inner),
291 intents: intent_rx,
292 control: Some(control_tx),
293 cancelled,
294 worker_status,
295 worker: Some(worker),
296 })
297 }
298
299 #[cfg(test)]
300 pub(crate) fn scripted(
301 root: &Path,
302 config: WatchConfig,
303 events: &Path,
304 ) -> Result<(Self, ScriptedSender)> {
305 config.validate()?;
306 let root = root.canonicalize().map_err(|error| Error::io(root, error))?;
307 let scripted = scripted_events::read_script(events, &root).map_err(Error::WatchScript)?;
308 let (raw_tx, raw_rx) = sync_channel::<RawMessage>(config.event_capacity);
309 let (intent_tx, intent_rx) = sync_channel::<CoalescedIntent>(config.intent_capacity);
310 let control_tx = raw_tx.clone();
311 let overflowed = Arc::new(AtomicBool::new(false));
312 for event in scripted {
313 enqueue_raw(&raw_tx, &overflowed, event);
314 }
315
316 let worker_root = root.clone();
317 let cancelled = Arc::new(AtomicBool::new(false));
318 let worker_cancelled = Arc::clone(&cancelled);
319 let worker_overflowed = Arc::clone(&overflowed);
320 let worker_status = Arc::new(AtomicU8::new(WORKER_RUNNING));
321 let tracked_status = Arc::clone(&worker_status);
322 let renames = RenameReporting::of_recommended_backend();
325 let worker = std::thread::Builder::new()
326 .name("fdu-scripted-watch".into())
327 .spawn(move || {
328 let _counter_guard = crate::counters::thread_flush_guard();
329 run_tracked_worker(&tracked_status, || {
330 run_worker(
331 &worker_root,
332 config,
333 renames,
334 &raw_rx,
335 &intent_tx,
336 &worker_overflowed,
337 &worker_cancelled,
338 );
339 });
340 })
341 .map_err(|error| Error::io(&root, error))?;
342
343 let sender = ScriptedSender {
344 root: root.clone(),
345 raw: control_tx.clone(),
346 overflowed: Arc::clone(&overflowed),
347 };
348 Ok((
349 Self {
350 root,
351 config,
352 inner: None,
353 intents: intent_rx,
354 control: Some(control_tx),
355 cancelled,
356 worker_status,
357 worker: Some(worker),
358 },
359 sender,
360 ))
361 }
362
363 pub fn next_observation(&self, timeout: Duration) -> Result<Option<Observation>> {
371 let Some(intent) = self.next_intent(timeout)? else {
372 return Ok(None);
373 };
374 Ok(Some(verify_intent(&self.root, self.config, &intent, &ScanConfig::default())))
375 }
376
377 fn next_intent(&self, timeout: Duration) -> Result<Option<CoalescedIntent>> {
378 match self.intents.recv_timeout(timeout) {
379 Ok(intent) => Ok(Some(intent)),
380 Err(RecvTimeoutError::Timeout) => Ok(None),
381 Err(RecvTimeoutError::Disconnected) => match self.worker_status.load(Ordering::Acquire)
382 {
383 WORKER_PANICKED => Err(Error::WatchWorkerPanicked),
384 _ => Err(Error::WatchStopped),
385 },
386 }
387 }
388
389 pub(crate) fn flush_capture(&self) -> Result<()> {
394 let Some(control) = self.control.as_ref() else {
395 return Err(Error::WatchStopped);
396 };
397 let (acknowledge, acknowledged) = sync_channel(0);
398 control.send(RawMessage::Flush(acknowledge)).map_err(|_| Error::WatchStopped)?;
399 acknowledged.recv().map_err(|_| Error::WatchStopped)
400 }
401
402 pub(crate) const fn capture_backlog_bound(&self) -> usize {
409 self.config.intent_capacity.saturating_add(1)
410 }
411
412 pub fn apply_next(
417 &self,
418 index: &IndexHandle,
419 scan_config: &ScanConfig,
420 timeout: Duration,
421 sink: &mut dyn FnMut(&Commit),
422 ) -> Result<Option<WatchApplyReport>> {
423 scan_config.validate_for_watch_scope(index.scope()?)?;
424 let indexed_root = index.root_path()?;
425 if indexed_root != self.root {
426 return Err(Error::WatchRootMismatch {
427 watched: self.root.clone(),
428 indexed: indexed_root,
429 });
430 }
431 let Some(intent) = self.next_intent(timeout)? else {
432 return Ok(None);
433 };
434 apply_intent(index, &self.root, self.config, &intent, scan_config, sink).map(Some)
435 }
436
437 pub(crate) fn apply_next_controlled(
439 &self,
440 index: &IndexHandle,
441 scan_config: &ScanConfig,
442 timeout: Duration,
443 control: &dyn scan::ReconcileControl,
444 sink: &mut dyn FnMut(&Commit),
445 ) -> Result<Option<WatchApplyReport>> {
446 scan_config.validate_for_watch_scope(index.scope()?)?;
447 let indexed_root = index.root_path()?;
448 if indexed_root != self.root {
449 return Err(Error::WatchRootMismatch {
450 watched: self.root.clone(),
451 indexed: indexed_root,
452 });
453 }
454 let Some(intent) = self.next_intent(timeout)? else {
455 return Ok(None);
456 };
457 apply_intent_controlled(index, &self.root, self.config, &intent, scan_config, control, sink)
458 .map(Some)
459 }
460}
461
462fn apply_intent(
463 index: &IndexHandle,
464 root: &Path,
465 watch_config: WatchConfig,
466 intent: &CoalescedIntent,
467 scan_config: &ScanConfig,
468 sink: &mut dyn FnMut(&Commit),
469) -> Result<WatchApplyReport> {
470 let mut verifier =
471 |_: &Path, _: &Observation| Ok(verify_intent(root, watch_config, intent, scan_config));
472 let apply = apply_reverified_with(index, &Observation::default(), scan_config, &mut verifier)?;
473 if let Some(commit) = apply.commit.as_ref() {
474 sink(commit);
475 }
476 let reconciliation = scan::reconcile_pending_handle(index, scan_config, sink)?;
477 retain_unreadable(index, root, &reconciliation, sink)?;
478 Ok(WatchApplyReport { apply, reconciliation })
479}
480
481fn retain_unreadable(
488 index: &IndexHandle,
489 root: &Path,
490 reconciliation: &scan::ReconcileReport,
491 sink: &mut dyn FnMut(&Commit),
492) -> Result<()> {
493 if reconciliation.scan.errors.is_empty() {
494 return Ok(());
495 }
496 let retained = reconciliation.scan.errors.len().min(crate::MAX_RETAINED_ISSUES);
497 let issues = reconciliation.scan.errors[..retained]
498 .iter()
499 .map(|error| crate::Issue::from_error_under(root, error))
500 .collect();
501 let omitted = u64::try_from(reconciliation.scan.errors.len() - retained).unwrap_or(u64::MAX);
502 let outcome =
503 index.transition_observation(crate::index::ObservationTransition::Unreadable {
504 issues,
505 omitted,
506 })?;
507 if let Some(commit) = outcome.commit.as_ref() {
508 sink(commit);
509 }
510 Ok(())
511}
512
513fn apply_intent_controlled(
514 index: &IndexHandle,
515 root: &Path,
516 watch_config: WatchConfig,
517 intent: &CoalescedIntent,
518 scan_config: &ScanConfig,
519 control: &dyn scan::ReconcileControl,
520 sink: &mut dyn FnMut(&Commit),
521) -> Result<WatchApplyReport> {
522 let mut verifier =
523 |_: &Path, _: &Observation| Ok(verify_intent(root, watch_config, intent, scan_config));
524 let apply = apply_reverified_with_control(
525 index,
526 &Observation::default(),
527 scan_config,
528 control,
529 &mut verifier,
530 )?;
531 if let Some(commit) = apply.commit.as_ref() {
532 sink(commit);
533 }
534 let reconciliation =
535 scan::reconcile_pending_handle_controlled(index, scan_config, control, sink)?;
536 Ok(WatchApplyReport { apply, reconciliation })
537}
538
539#[cfg(test)]
545fn apply_observation(
546 index: &IndexHandle,
547 observation: &Observation,
548 scan_config: &ScanConfig,
549 sink: &mut dyn FnMut(&Commit),
550) -> Result<WatchApplyReport> {
551 scan_config.validate_for_watch_scope(index.scope()?)?;
552 let apply = apply_reverified(index, observation, scan_config)?;
553 if let Some(commit) = apply.commit.as_ref() {
554 sink(commit);
555 }
556 let reconciliation = scan::reconcile_pending_handle(index, scan_config, sink)?;
557 Ok(WatchApplyReport { apply, reconciliation })
558}
559
560#[cfg(test)]
568fn apply_reverified(
569 index: &IndexHandle,
570 observation: &Observation,
571 scan_config: &ScanConfig,
572) -> Result<ApplyOutcome> {
573 let mut verifier = |root: &Path, observation: &Observation| {
574 reverify_observation(root, observation, scan_config)
575 };
576 apply_reverified_with(index, observation, scan_config, &mut verifier)
577}
578
579fn apply_reverified_with(
580 index: &IndexHandle,
581 observation: &Observation,
582 scan_config: &ScanConfig,
583 verifier: &mut impl FnMut(&Path, &Observation) -> Result<Observation>,
584) -> Result<ApplyOutcome> {
585 let (root, scope, _) = index.watch_boundary()?;
586 scan_config.validate_for_watch_scope(scope)?;
587 for _ in 0..MAX_OPTIMISTIC_APPLY_ATTEMPTS {
588 let clock = index.clock()?;
589 let candidate = verifier(&root, observation)?;
590 let candidate = escalate_unknown_ancestry(index, candidate)?;
591 if let Some(outcome) = index.apply_if_clock(clock, &candidate)? {
592 return Ok(outcome);
593 }
594 }
595
596 index.invalidate_root(InvalidateReason::WatchContention)
597}
598
599fn apply_reverified_with_control(
600 index: &IndexHandle,
601 observation: &Observation,
602 scan_config: &ScanConfig,
603 control: &dyn scan::ReconcileControl,
604 verifier: &mut impl FnMut(&Path, &Observation) -> Result<Observation>,
605) -> Result<ApplyOutcome> {
606 let (root, scope, _) = index.watch_boundary()?;
607 scan_config.validate_for_watch_scope(scope)?;
608 for _ in 0..MAX_OPTIMISTIC_APPLY_ATTEMPTS {
609 control.check_active()?;
610 let clock = index.clock()?;
611 let candidate = verifier(&root, observation)?;
612 let candidate = escalate_unknown_ancestry(index, candidate)?;
613 control.before_conditional_commit()?;
614 if let Some(outcome) =
615 index.apply_opened_if_clock(clock, &candidate, control.max_files())?
616 {
617 return Ok(outcome);
618 }
619 }
620
621 control.before_conditional_commit()?;
622 let outcome = index.apply_opened(
623 &Observation::new(vec![Op::InvalidateSubtree {
624 path: PathBuf::new(),
625 reason: InvalidateReason::WatchContention,
626 }]),
627 control.max_files(),
628 )?;
629 Ok(outcome)
630}
631
632fn escalate_unknown_ancestry(index: &IndexHandle, candidate: Observation) -> Result<Observation> {
634 let unknown = index.unknown_ancestry(&candidate)?;
635 if unknown.is_empty() {
636 return Ok(candidate);
637 }
638
639 let mut roots: Vec<PathBuf> = unknown.into_iter().map(|(_, root)| root).collect();
640 roots.sort_by(|left, right| {
641 left.components().count().cmp(&right.components().count()).then_with(|| left.cmp(right))
642 });
643 roots.dedup();
644 let mut covering: Vec<PathBuf> = Vec::with_capacity(roots.len());
645 for root in roots {
646 if !covering.iter().any(|ancestor| root.starts_with(ancestor)) {
647 covering.push(root);
648 }
649 }
650
651 let mut ops: Vec<ObservationOp> = candidate
652 .ops
653 .into_iter()
654 .filter(|observed| !covering.iter().any(|root| observed.op.path().starts_with(root)))
655 .collect();
656 ops.extend(covering.into_iter().map(|path| {
657 ObservationOp::unconditional(Op::InvalidateSubtree {
658 path,
659 reason: InvalidateReason::UnknownAncestry,
660 })
661 }));
662 Ok(Observation::from_ops(ops))
663}
664
665#[cfg(test)]
666fn reverify_observation(
667 root: &Path,
668 observation: &Observation,
669 scan_config: &ScanConfig,
670) -> Result<Observation> {
671 let mut ops = Vec::with_capacity(observation.len().saturating_mul(2));
672 for observed in &observation.ops {
673 let relative = scan::normalize_subtree(observed.op.path())?;
674 match &observed.op {
675 Op::InvalidateSubtree { reason, .. } => {
676 ops.push(Op::InvalidateSubtree { path: relative, reason: *reason });
677 }
678 Op::Upsert { .. } | Op::Remove { .. } => {
679 let absolute = root.join(&relative);
680 match scan::observe_path(&absolute) {
681 Ok((kind, attrs)) => {
682 match crate::admission::decide_path(
683 &relative,
684 kind,
685 scan_config.hidden(),
686 scan_config.exclude_special,
687 ) {
688 crate::admission::Disposition::Retain => {
689 ops.push(Op::Upsert { path: relative.clone(), kind, attrs });
690 if let Some(control) =
691 scan::read_control_op(scan_config, root, &relative, kind)?
692 {
693 ops.push(control);
694 }
695 }
696 crate::admission::Disposition::ControlOnly => {
697 ops.push(Op::Remove { path: relative.clone() });
698 if let Some(control) =
699 scan::read_control_op(scan_config, root, &relative, kind)?
700 {
701 ops.push(control);
702 }
703 }
704 crate::admission::Disposition::Reject => {
705 ops.push(Op::Remove { path: relative });
706 }
707 }
708 }
709 Err(error) => ops.push(op_for_stat_error(relative, &error)),
710 }
711 }
712 Op::ControlUpsert { source, .. } => {
713 ops.push(Op::ControlUpsert { path: relative, source: source.clone() });
714 }
715 Op::ControlRemove { .. } => ops.push(Op::ControlRemove { path: relative }),
716 }
717 }
718 Ok(Observation::new(ops))
719}
720
721impl Drop for Watcher {
722 fn drop(&mut self) {
723 self.cancelled.store(true, Ordering::Release);
724 if let Some(control) = self.control.take() {
725 let _ = control.try_send(RawMessage::Stop);
726 }
727 self.inner.take();
730 if let Some(worker) = self.worker.take() {
731 let _ = worker.join();
732 }
733 }
734}
735
736fn notify_error(path: &Path, err: notify::Error) -> Error {
737 Error::io(path, std::io::Error::other(err))
738}
739
740fn enqueue_raw(
741 sender: &SyncSender<RawMessage>,
742 overflowed: &AtomicBool,
743 event: notify::Result<notify::Event>,
744) {
745 match sender.try_send(RawMessage::Event(event)) {
746 Ok(()) | Err(TrySendError::Disconnected(_)) => {}
747 Err(TrySendError::Full(_)) => overflowed.store(true, Ordering::Release),
748 }
749}
750
751fn run_tracked_worker(status: &AtomicU8, worker: impl FnOnce()) {
752 let outcome = catch_unwind(AssertUnwindSafe(worker));
753 status.store(if outcome.is_ok() { WORKER_STOPPED } else { WORKER_PANICKED }, Ordering::Release);
754}
755
756fn run_worker(
757 root: &Path,
758 config: WatchConfig,
759 renames: RenameReporting,
760 raw: &Receiver<RawMessage>,
761 out: &SyncSender<CoalescedIntent>,
762 overflowed: &AtomicBool,
763 cancelled: &AtomicBool,
764) {
765 let mut pending: BTreeMap<PathBuf, Pending> = BTreeMap::new();
766 let mut batch_started: Option<Instant> = None;
767 let mut sticky_overflow = false;
768
769 loop {
770 if cancelled.load(Ordering::Acquire) {
771 return;
772 }
773 if overflowed.swap(false, Ordering::AcqRel) {
774 collapse_to_overflow(&mut pending);
775 batch_started.get_or_insert_with(Instant::now);
776 }
777 if sticky_overflow {
778 match try_deliver_overflow(out) {
779 Ok(true) => sticky_overflow = false,
780 Ok(false) => {}
781 Err(()) => return,
782 }
783 }
784
785 match raw.recv_timeout(config.settle) {
786 Ok(RawMessage::Event(Ok(event))) => {
787 record(root, &event, &mut pending, config.batch_path_capacity, renames);
788 batch_started.get_or_insert_with(Instant::now);
789 }
790 Ok(RawMessage::Event(Err(err))) => {
791 let _ = err;
794 collapse_to_overflow(&mut pending);
795 batch_started.get_or_insert_with(Instant::now);
796 }
797 Ok(RawMessage::Flush(acknowledge)) => {
798 if overflowed.swap(false, Ordering::AcqRel) {
799 collapse_to_overflow(&mut pending);
800 }
801 if !pending.is_empty()
802 && try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err()
803 {
804 return;
805 }
806 if sticky_overflow {
810 match try_deliver_overflow(out) {
811 Ok(true) => sticky_overflow = false,
812 Ok(false) => {}
813 Err(()) => return,
814 }
815 }
816 batch_started = None;
817 let _ = acknowledge.send(());
818 }
819 Ok(RawMessage::Stop) => return,
820 Err(RecvTimeoutError::Timeout) => {
821 if !pending.is_empty()
823 && try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err()
824 {
825 return;
826 }
827 batch_started = None;
828 continue;
829 }
830 Err(RecvTimeoutError::Disconnected) => {
831 if !cancelled.load(Ordering::Acquire) && !pending.is_empty() {
832 let _ = try_deliver_pending(&mut pending, out, &mut sticky_overflow);
833 }
834 return;
835 }
836 }
837
838 if max_hold_elapsed(batch_started, config.max_hold) {
840 if try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err() {
841 return;
842 }
843 batch_started = None;
844 }
845 }
846}
847
848fn max_hold_elapsed(started: Option<Instant>, max_hold: Duration) -> bool {
849 started.is_some_and(|start| start.elapsed() >= max_hold)
850}
851
852fn record(
854 root: &Path,
855 event: ¬ify::Event,
856 pending: &mut BTreeMap<PathBuf, Pending>,
857 capacity: usize,
858 reporting: RenameReporting,
859) {
860 if event.need_rescan() {
861 let target = if event.paths.len() == 1 {
864 relative_to(root, &event.paths[0]).unwrap_or_default()
865 } else {
866 PathBuf::new()
867 };
868 queue_pending(
869 pending,
870 target,
871 Pending::Escalate(InvalidateReason::WatchOverflow),
872 capacity,
873 );
874 return;
875 }
876
877 let renamed = matches!(event.kind, EventKind::Modify(notify::event::ModifyKind::Name(_)));
894 if renamed
895 && (reporting == RenameReporting::OldSideOnly
896 || !event.paths.iter().any(|path| relative_to(root, path).is_some()))
897 {
898 queue_pending(
904 pending,
905 PathBuf::new(),
906 Pending::Escalate(InvalidateReason::UnpairedRename),
907 capacity,
908 );
909 }
910
911 for path in &event.paths {
912 let Some(rel) = relative_to(root, path) else {
913 continue; };
915 if matches!(event.kind, EventKind::Access(_)) {
916 continue; }
918 if renamed && rel.as_os_str().is_empty() {
919 queue_pending(
923 pending,
924 PathBuf::new(),
925 Pending::Escalate(InvalidateReason::UnpairedRename),
926 capacity,
927 );
928 continue;
929 }
930 let relist_if_dir = matches!(event.kind, EventKind::Create(_));
931 queue_pending(pending, rel, Pending::Verify { relist_if_dir, renamed }, capacity);
932 }
933}
934
935fn queue_pending(
936 pending: &mut BTreeMap<PathBuf, Pending>,
937 path: PathBuf,
938 state: Pending,
939 capacity: usize,
940) {
941 if matches!(pending.get(Path::new("")), Some(Pending::Escalate(_))) {
942 return;
943 }
944 if path.as_os_str().is_empty() && matches!(state, Pending::Escalate(_)) {
945 pending.clear();
946 pending.insert(path, state);
947 return;
948 }
949 if let Some(existing) = pending.get_mut(&path) {
950 match (existing, state) {
951 (Pending::Escalate(_), _) => {}
952 (
953 Pending::Verify { relist_if_dir, renamed },
954 Pending::Verify { relist_if_dir: relist, renamed: rename },
955 ) => {
956 *relist_if_dir |= relist;
957 *renamed |= rename;
958 }
959 (slot @ Pending::Verify { .. }, Pending::Escalate(reason)) => {
960 *slot = Pending::Escalate(reason);
961 }
962 }
963 return;
964 }
965 if pending.len() >= capacity {
966 collapse_to_overflow(pending);
967 } else {
968 pending.insert(path, state);
969 }
970}
971
972fn collapse_to_overflow(pending: &mut BTreeMap<PathBuf, Pending>) {
973 pending.clear();
974 pending.insert(PathBuf::new(), Pending::Escalate(InvalidateReason::WatchOverflow));
975}
976
977fn try_deliver_pending(
978 pending: &mut BTreeMap<PathBuf, Pending>,
979 out: &SyncSender<CoalescedIntent>,
980 sticky_overflow: &mut bool,
981) -> std::result::Result<(), ()> {
982 if pending.is_empty() {
983 return Ok(());
984 }
985 let intent = CoalescedIntent { pending: std::mem::take(pending) };
986 match out.try_send(intent) {
987 Ok(()) => Ok(()),
988 Err(TrySendError::Full(_)) => {
989 *sticky_overflow = true;
990 Ok(())
991 }
992 Err(TrySendError::Disconnected(_)) => Err(()),
993 }
994}
995
996fn try_deliver_overflow(out: &SyncSender<CoalescedIntent>) -> std::result::Result<bool, ()> {
997 let mut pending = BTreeMap::new();
998 collapse_to_overflow(&mut pending);
999 match out.try_send(CoalescedIntent { pending }) {
1000 Ok(()) => Ok(true),
1001 Err(TrySendError::Full(_)) => Ok(false),
1002 Err(TrySendError::Disconnected(_)) => Err(()),
1003 }
1004}
1005
1006fn verify_intent(
1008 root: &Path,
1009 config: WatchConfig,
1010 intent: &CoalescedIntent,
1011 scan_config: &ScanConfig,
1012) -> Observation {
1013 let mut ops = Vec::with_capacity(intent.pending.len());
1014 let mut listings = ParentListings::default();
1015 let mut unlisted: Vec<PathBuf> = Vec::new();
1020
1021 for (rel, state) in &intent.pending {
1022 match state {
1023 Pending::Escalate(reason) => {
1024 ops.push(Op::InvalidateSubtree { path: rel.clone(), reason: *reason });
1025 }
1026 Pending::Verify { relist_if_dir, renamed } => {
1027 if unlisted.iter().any(|name| rel.starts_with(name)) {
1028 continue;
1029 }
1030 let absolute = root.join(rel);
1031 let mut stat = scan::observe_path(&absolute);
1032 if *renamed && !rel.as_os_str().is_empty() && stat.is_ok() {
1033 let unlisted_reason = match listings.lists(root, rel) {
1040 Some(true) => None,
1041 Some(false) => {
1042 stat = scan::observe_path(&absolute);
1047 stat.is_ok().then_some(InvalidateReason::UnpairedRename)
1048 }
1049 None => Some(InvalidateReason::VerificationFailed),
1052 };
1053 if let Some(reason) = unlisted_reason {
1054 let parent = parent_of(rel);
1057 if !ops.iter().any(
1058 |op| matches!(op, Op::InvalidateSubtree { path, .. } if *path == parent),
1059 ) {
1060 ops.push(Op::InvalidateSubtree { path: parent, reason });
1061 }
1062 unlisted.push(rel.clone());
1063 continue;
1064 }
1065 }
1066 match stat {
1067 Ok((kind, attrs)) => {
1068 let disposition = crate::admission::decide_path(
1069 rel,
1070 kind,
1071 scan_config.hidden(),
1072 scan_config.exclude_special,
1073 );
1074 match disposition {
1075 crate::admission::Disposition::Retain => {
1076 ops.push(Op::Upsert { path: rel.clone(), kind, attrs });
1077 match scan::read_control_op(scan_config, root, rel, kind) {
1078 Ok(Some(control)) => ops.push(control),
1079 Ok(None) => {}
1080 Err(_) => ops.push(Op::InvalidateSubtree {
1081 path: rel
1082 .parent()
1083 .map_or_else(PathBuf::new, Path::to_path_buf),
1084 reason: InvalidateReason::VerificationFailed,
1085 }),
1086 }
1087 }
1088 crate::admission::Disposition::ControlOnly => {
1089 match scan::read_control_op(scan_config, root, rel, kind) {
1090 Ok(Some(control)) => {
1091 ops.push(Op::Remove { path: rel.clone() });
1092 ops.push(control);
1093 }
1094 Ok(None) => ops.push(Op::Remove { path: rel.clone() }),
1095 Err(_) => ops.push(Op::InvalidateSubtree {
1096 path: rel
1097 .parent()
1098 .map_or_else(PathBuf::new, Path::to_path_buf),
1099 reason: InvalidateReason::VerificationFailed,
1100 }),
1101 }
1102 }
1103 crate::admission::Disposition::Reject => {
1104 ops.push(Op::Remove { path: rel.clone() });
1105 }
1106 }
1107 let retained_dir =
1108 disposition == crate::admission::Disposition::Retain && kind.is_dir();
1109 if retained_dir && *renamed {
1110 ops.push(Op::InvalidateSubtree {
1116 path: rel.clone(),
1117 reason: InvalidateReason::UnpairedRename,
1118 });
1119 } else if retained_dir && *relist_if_dir && config.relist_new_dirs {
1120 ops.push(Op::InvalidateSubtree {
1123 path: rel.clone(),
1124 reason: InvalidateReason::WatchSetupRace,
1125 });
1126 }
1127 }
1128 Err(error) => {
1129 ops.push(op_for_stat_error(rel.clone(), &error));
1130 if error.kind() == std::io::ErrorKind::NotFound
1135 && crate::control::path_control_spelling(rel)
1136 == Some(crate::control::ControlSpelling::Variant)
1137 {
1138 let control = crate::control::sibling_control_path(rel);
1139 match scan::read_directory_control_or_removal(
1140 scan_config,
1141 root,
1142 &control,
1143 ) {
1144 Ok(Some(observed)) => ops.push(observed),
1145 Ok(None) => {}
1146 Err(_) => ops.push(Op::InvalidateSubtree {
1147 path: parent_of(rel),
1148 reason: InvalidateReason::VerificationFailed,
1149 }),
1150 }
1151 }
1152 }
1153 }
1154 if scan_config.population != crate::query::IgnoredEntries::Include
1155 && crate::control::path_control_spelling(rel).is_some()
1156 {
1157 ops.push(Op::InvalidateSubtree {
1158 path: rel.parent().map_or_else(PathBuf::new, Path::to_path_buf),
1159 reason: InvalidateReason::ControlPopulationChanged,
1160 });
1161 }
1162 }
1163 }
1164 }
1165 Observation::new(ops)
1166}
1167
1168fn parent_of(rel: &Path) -> PathBuf {
1170 rel.parent().map_or_else(PathBuf::new, Path::to_path_buf)
1171}
1172
1173#[derive(Default)]
1179struct ParentListings(BTreeMap<PathBuf, Option<std::collections::HashSet<std::ffi::OsString>>>);
1180
1181impl ParentListings {
1182 fn lists(&mut self, root: &Path, rel: &Path) -> Option<bool> {
1189 let name = rel.file_name()?;
1190 let parent = parent_of(rel);
1191 if let Some(Some(names)) = self.0.get(&parent) {
1192 if names.contains(name) {
1193 return Some(true);
1194 }
1195 }
1196 let names: Option<std::collections::HashSet<_>> = std::fs::read_dir(root.join(&parent))
1197 .and_then(|entries| entries.map(|entry| entry.map(|entry| entry.file_name())).collect())
1198 .ok();
1199 let listed = names.as_ref().map(|names| names.contains(name));
1200 self.0.insert(parent, names);
1201 listed
1202 }
1203}
1204
1205fn op_for_stat_error(path: PathBuf, error: &std::io::Error) -> Op {
1206 match error.kind() {
1207 std::io::ErrorKind::NotFound if path.as_os_str().is_empty() => {
1208 Op::InvalidateSubtree { path, reason: InvalidateReason::VerificationFailed }
1209 }
1210 std::io::ErrorKind::NotFound => Op::Remove { path },
1211 std::io::ErrorKind::NotADirectory => Op::InvalidateSubtree {
1212 path: path.parent().map_or_else(PathBuf::new, Path::to_path_buf),
1213 reason: InvalidateReason::VerificationFailed,
1214 },
1215 _ => Op::InvalidateSubtree { path, reason: InvalidateReason::VerificationFailed },
1216 }
1217}
1218
1219fn relative_to(root: &Path, path: &Path) -> Option<PathBuf> {
1224 path.strip_prefix(root).ok().map(Path::to_path_buf)
1225}
1226
1227#[cfg(test)]
1228mod tests {
1229 use super::*;
1230 use notify::event::{CreateKind, Flag, MetadataKind, ModifyKind, RenameMode};
1231 use std::fs;
1232
1233 fn queued_test_watcher(root: PathBuf) -> (SyncSender<CoalescedIntent>, Watcher) {
1234 let (sender, intents) = sync_channel(1);
1235 let watcher = Watcher {
1236 root,
1237 config: WatchConfig::default(),
1238 inner: None,
1239 intents,
1240 control: None,
1241 cancelled: Arc::new(AtomicBool::new(false)),
1242 worker_status: Arc::new(AtomicU8::new(WORKER_RUNNING)),
1243 worker: None,
1244 };
1245 (sender, watcher)
1246 }
1247
1248 #[test]
1249 fn watcher_can_move_to_its_single_consumer_thread() {
1250 fn assert_send<T: Send>() {}
1251
1252 assert_send::<Watcher>();
1253 }
1254
1255 static REAL_WATCHER: std::sync::Mutex<()> = std::sync::Mutex::new(());
1283
1284 const REAL_BACKEND_DELIVERY: Duration = Duration::from_secs(60);
1298
1299 fn real_watcher_guard() -> std::sync::MutexGuard<'static, ()> {
1305 REAL_WATCHER.lock().unwrap_or_else(std::sync::PoisonError::into_inner)
1306 }
1307
1308 enum Waited {
1310 Delivered(Vec<Op>),
1312 Silent,
1314 }
1315
1316 fn wait_for(
1337 watcher: &Watcher,
1338 deadline: Duration,
1339 mut want: impl FnMut(&[Op]) -> bool,
1340 ) -> Waited {
1341 let start = Instant::now();
1342 let mut seen: Vec<Op> = Vec::new();
1343 while start.elapsed() < deadline {
1344 match watcher.next_observation(Duration::from_millis(200)) {
1345 Ok(Some(observation)) => {
1346 seen.extend(observation.ops.into_iter().map(|observed| observed.op));
1347 if want(&seen) {
1348 return Waited::Delivered(seen);
1349 }
1350 }
1351 Ok(None) => {}
1352 Err(error) => panic!("watcher stopped while waiting: {error}"),
1353 }
1354 }
1355 if seen.is_empty() {
1356 return Waited::Silent;
1357 }
1358 panic!(
1359 "the backend delivered {} op(s) in {deadline:?} but never the one awaited, so \
1360 this is a disagreement about content rather than a delivery failure: {seen:?}",
1361 seen.len()
1362 );
1363 }
1364
1365 fn wait_established(
1371 watcher: &Watcher,
1372 deadline: Duration,
1373 want: impl FnMut(&[Op]) -> bool,
1374 ) -> Vec<Op> {
1375 match wait_for(watcher, deadline, want) {
1376 Waited::Delivered(ops) => ops,
1377 Waited::Silent => panic!(
1378 "the watch was established and then delivered nothing in {deadline:?}, so \
1379 this is a lost event rather than a host precondition; \
1380 FDU_TEST_ALLOW_NO_NATIVE_WATCH does not apply here"
1381 ),
1382 }
1383 }
1384
1385 fn establish_watch(watcher: &Watcher, dir: &Path) -> bool {
1408 let warmup = dir.join(".fdu-watch-warmup");
1409 fs::write(&warmup, b"warmup").expect("warmup write");
1410 let waited = wait_for(watcher, REAL_BACKEND_DELIVERY, |ops| !ops.is_empty());
1411 let _ = fs::remove_file(&warmup);
1412 match waited {
1413 Waited::Delivered(_) => true,
1414 Waited::Silent => {
1415 if std::env::var_os("FDU_TEST_ALLOW_NO_NATIVE_WATCH").as_deref()
1416 == Some(std::ffi::OsStr::new("1"))
1417 {
1418 eprintln!(
1419 "skipped by FDU_TEST_ALLOW_NO_NATIVE_WATCH=1: the host event service \
1420 delivered no events to this stream"
1421 );
1422 return false;
1423 }
1424 panic!(
1425 "native watch precondition failed: the host delivered no events in \
1426 {REAL_BACKEND_DELIVERY:?}; run on a host with event delivery, or \
1427 explicitly opt out with FDU_TEST_ALLOW_NO_NATIVE_WATCH=1"
1428 );
1429 }
1430 }
1431 }
1432
1433 #[test]
1434 fn created_files_arrive_as_verified_upserts() {
1435 let _serialized = real_watcher_guard();
1436 let dir = tempfile::tempdir().expect("tempdir");
1437 let watcher = Watcher::new(dir.path(), WatchConfig::default()).expect("watcher");
1438 if !establish_watch(&watcher, dir.path()) {
1439 return;
1440 }
1441
1442 fs::write(dir.path().join("hello.txt"), b"hello world").expect("write");
1443
1444 let ops = wait_established(&watcher, REAL_BACKEND_DELIVERY, |ops| {
1445 ops.iter().any(|op| op.path() == Path::new("hello.txt"))
1446 });
1447
1448 let found = ops
1449 .iter()
1450 .find(|op| op.path() == Path::new("hello.txt"))
1451 .expect("an op for the new file");
1452 match found {
1453 Op::Upsert { attrs, kind, .. } => {
1454 assert!(!kind.is_dir());
1455 assert_eq!(attrs.size, 11);
1458 assert!(attrs.mtime_ns > 0);
1459 }
1460 other => panic!("expected an upsert, got {other:?}"),
1461 }
1462 }
1463
1464 #[test]
1465 fn deleted_files_arrive_as_removes() {
1466 let _serialized = real_watcher_guard();
1467 let dir = tempfile::tempdir().expect("tempdir");
1468 let path = dir.path().join("doomed.txt");
1469 fs::write(&path, b"x").expect("write");
1470
1471 let watcher = Watcher::new(dir.path(), WatchConfig::default()).expect("watcher");
1472 if !establish_watch(&watcher, dir.path()) {
1473 return;
1474 }
1475 fs::remove_file(&path).expect("remove");
1476
1477 let ops = wait_established(&watcher, REAL_BACKEND_DELIVERY, |ops| {
1478 ops.iter()
1479 .any(|op| matches!(op, Op::Remove { path } if path == Path::new("doomed.txt")))
1480 });
1481
1482 assert!(
1483 ops.iter()
1484 .any(|op| matches!(op, Op::Remove { path } if path == Path::new("doomed.txt"))),
1485 "expected a remove, saw {ops:?}"
1486 );
1487 }
1488
1489 #[test]
1508 fn watching_a_missing_path_is_an_error() {
1509 let dir = tempfile::tempdir().expect("tempdir");
1510 let missing = dir.path().join("not-there");
1511 assert!(Watcher::new(&missing, WatchConfig::default()).is_err());
1512 }
1513
1514 #[test]
1515 fn paths_outside_the_root_are_ignored() {
1516 let root = Path::new("/a/b");
1517 assert_eq!(relative_to(root, Path::new("/a/b/c/d")), Some(PathBuf::from("c/d")));
1518 assert_eq!(relative_to(root, Path::new("/elsewhere")), None);
1519 }
1520
1521 #[test]
1522 fn verification_errors_distinguish_absence_from_an_invalid_ancestor() {
1523 let path = PathBuf::from("parent/known.txt");
1524 let missing = op_for_stat_error(
1525 path.clone(),
1526 &std::io::Error::new(std::io::ErrorKind::NotFound, "gone"),
1527 );
1528 assert!(matches!(missing, Op::Remove { path: removed } if removed == path));
1529
1530 let not_a_directory = op_for_stat_error(
1531 path.clone(),
1532 &std::io::Error::new(std::io::ErrorKind::NotADirectory, "ancestor is a file"),
1533 );
1534 assert!(matches!(
1535 not_a_directory,
1536 Op::InvalidateSubtree {
1537 path: invalidated,
1538 reason: InvalidateReason::VerificationFailed,
1539 } if invalidated == Path::new("parent")
1540 ));
1541
1542 let denied = op_for_stat_error(
1543 path.clone(),
1544 &std::io::Error::new(std::io::ErrorKind::PermissionDenied, "denied"),
1545 );
1546 assert!(matches!(
1547 denied,
1548 Op::InvalidateSubtree {
1549 path: invalidated,
1550 reason: InvalidateReason::VerificationFailed,
1551 } if invalidated == path
1552 ));
1553 }
1554
1555 #[test]
1556 fn unknown_watch_ancestry_reconciles_from_the_nearest_known_directory() {
1557 let dir = tempfile::tempdir().expect("tempdir");
1558 let (index, _) =
1559 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
1560 let handle = crate::IndexHandle::new(index);
1561 let nested = dir.path().join("new/deep");
1562 fs::create_dir_all(&nested).expect("nested directories");
1563 fs::write(nested.join("file.txt"), b"verified").expect("nested file");
1564
1565 let report = apply_observation(
1566 &handle,
1567 &Observation::new(vec![Op::Upsert {
1568 path: PathBuf::from("new/deep/file.txt"),
1569 kind: crate::EntryKind::File,
1570 attrs: crate::Attrs::default(),
1571 }]),
1572 &crate::ScanConfig::default(),
1573 &mut |_| {},
1574 )
1575 .expect("unknown ancestry schedules reconciliation");
1576
1577 assert_eq!(report.apply.invalidated, 1);
1578 assert!(report.reconciliation.is_complete());
1579 assert_eq!(
1580 handle.kind(Path::new("new")).expect("new directory"),
1581 Some(crate::EntryKind::Dir)
1582 );
1583 assert_eq!(
1584 handle.kind(Path::new("new/deep/file.txt")).expect("nested file"),
1585 Some(crate::EntryKind::File)
1586 );
1587 assert_ne!(
1588 handle.attrs(Path::new("new")).expect("verified parent attrs"),
1589 Some(crate::Attrs::default())
1590 );
1591 }
1592
1593 #[test]
1594 fn create_intent_survives_coalescing_but_metadata_only_does_not_relist() {
1595 let root = Path::new("/watch-root");
1596 let path = root.join("directory");
1597 let mut pending = BTreeMap::new();
1598
1599 record(
1600 root,
1601 ¬ify::Event::new(EventKind::Create(CreateKind::Folder)).add_path(path.clone()),
1602 &mut pending,
1603 16,
1604 RenameReporting::EachSide,
1605 );
1606 record(
1607 root,
1608 ¬ify::Event::new(EventKind::Modify(ModifyKind::Metadata(MetadataKind::Any)))
1609 .add_path(path),
1610 &mut pending,
1611 16,
1612 RenameReporting::EachSide,
1613 );
1614 assert_eq!(
1615 pending.get(Path::new("directory")),
1616 Some(&Pending::Verify { relist_if_dir: true, renamed: false })
1617 );
1618
1619 let mut metadata_only = BTreeMap::new();
1620 record(
1621 root,
1622 ¬ify::Event::new(EventKind::Modify(ModifyKind::Metadata(MetadataKind::Any)))
1623 .add_path(root.join("existing")),
1624 &mut metadata_only,
1625 16,
1626 RenameReporting::EachSide,
1627 );
1628 assert_eq!(
1629 metadata_only.get(Path::new("existing")),
1630 Some(&Pending::Verify { relist_if_dir: false, renamed: false })
1631 );
1632 }
1633
1634 fn recorded(
1636 root: &Path,
1637 renames: RenameReporting,
1638 events: &[notify::Event],
1639 ) -> BTreeMap<PathBuf, Pending> {
1640 let mut pending = BTreeMap::new();
1641 for event in events {
1642 record(root, event, &mut pending, 16, renames);
1643 }
1644 pending
1645 }
1646
1647 fn rename_event(mode: RenameMode, paths: &[PathBuf]) -> notify::Event {
1648 let mut event = notify::Event::new(EventKind::Modify(ModifyKind::Name(mode)));
1649 event.paths = paths.to_vec();
1650 event
1651 }
1652
1653 fn fsevents_rename(path: PathBuf) -> notify::Event {
1655 rename_event(RenameMode::Any, &[path])
1656 }
1657
1658 const RENAMED: Pending = Pending::Verify { relist_if_dir: false, renamed: true };
1659
1660 #[test]
1664 fn each_rename_shape_verifies_its_named_paths_without_escalating() {
1665 let root = Path::new("/watch-root");
1666 let (old, new) = (root.join("dir/old"), root.join("other/new"));
1667 for (backend, events) in [
1668 ("FSEvents", vec![fsevents_rename(old.clone()), fsevents_rename(new.clone())]),
1669 (
1670 "inotify",
1671 vec![
1672 rename_event(RenameMode::From, std::slice::from_ref(&old)),
1673 rename_event(RenameMode::To, std::slice::from_ref(&new)),
1674 rename_event(RenameMode::Both, &[old.clone(), new.clone()]),
1675 ],
1676 ),
1677 (
1678 "Windows",
1679 vec![
1680 rename_event(RenameMode::From, std::slice::from_ref(&old)),
1681 rename_event(RenameMode::To, std::slice::from_ref(&new)),
1682 ],
1683 ),
1684 ] {
1685 let pending = recorded(root, RenameReporting::EachSide, &events);
1686 assert_eq!(
1687 pending,
1688 BTreeMap::from([
1689 (PathBuf::from("dir/old"), RENAMED),
1690 (PathBuf::from("other/new"), RENAMED),
1691 ]),
1692 "{backend}"
1693 );
1694 }
1695
1696 let outside = Path::new("/elsewhere/file");
1697 for (direction, event) in [
1698 ("move in", rename_event(RenameMode::Both, &[outside.to_path_buf(), new.clone()])),
1699 ("move out", rename_event(RenameMode::Both, &[old.clone(), outside.to_path_buf()])),
1700 ] {
1701 let pending = recorded(root, RenameReporting::EachSide, &[event]);
1702 assert_eq!(pending.len(), 1, "{direction}: {pending:?}");
1703 assert!(!pending.contains_key(Path::new("")), "{direction}: {pending:?}");
1704 }
1705 }
1706
1707 #[test]
1711 fn a_sticky_rename_flag_merges_with_the_same_paths_other_events() {
1712 let root = Path::new("/watch-root");
1713 let path = root.join("state.json");
1714 let pending = recorded(
1715 root,
1716 RenameReporting::EachSide,
1717 &[
1718 notify::Event::new(EventKind::Create(CreateKind::File)).add_path(path.clone()),
1719 fsevents_rename(path.clone()),
1720 notify::Event::new(EventKind::Modify(ModifyKind::Data(
1721 notify::event::DataChange::Content,
1722 )))
1723 .add_path(path),
1724 ],
1725 );
1726
1727 assert_eq!(
1728 pending,
1729 BTreeMap::from([(
1730 PathBuf::from("state.json"),
1731 Pending::Verify { relist_if_dir: true, renamed: true }
1732 )])
1733 );
1734 }
1735
1736 #[test]
1740 fn unboundable_renames_and_ambiguous_rescans_escalate_the_root() {
1741 let root = Path::new("/watch-root");
1742 let escalated = Some(&Pending::Escalate(InvalidateReason::UnpairedRename));
1743
1744 let kqueue =
1745 recorded(root, RenameReporting::OldSideOnly, &[fsevents_rename(root.join("old"))]);
1746 assert_eq!(kqueue.get(Path::new("")), escalated);
1747
1748 let root_moved = recorded(
1749 root,
1750 RenameReporting::EachSide,
1751 &[rename_event(RenameMode::From, &[root.to_path_buf()])],
1752 );
1753 assert_eq!(root_moved.get(Path::new("")), escalated);
1754 assert_eq!(root_moved.len(), 1, "the root is escalated, never verified: {root_moved:?}");
1755
1756 for paths in [vec![], vec![PathBuf::from("/elsewhere/file")]] {
1757 let unplaced =
1758 recorded(root, RenameReporting::EachSide, &[rename_event(RenameMode::Any, &paths)]);
1759 assert_eq!(unplaced.get(Path::new("")), escalated, "{paths:?}");
1760 }
1761
1762 let rescan = recorded(
1763 root,
1764 RenameReporting::EachSide,
1765 &[notify::Event::new(EventKind::Any)
1766 .add_path(root.join("a"))
1767 .add_path(root.join("b"))
1768 .set_flag(Flag::Rescan)],
1769 );
1770 assert_eq!(
1771 rescan.get(Path::new("")),
1772 Some(&Pending::Escalate(InvalidateReason::WatchOverflow))
1773 );
1774 }
1775
1776 #[test]
1777 fn pending_path_overload_collapses_to_one_root_invalidation() {
1778 let root = Path::new("/watch-root");
1779 let mut pending = BTreeMap::new();
1780 for name in ["one", "two", "three"] {
1781 record(
1782 root,
1783 ¬ify::Event::new(EventKind::Any).add_path(root.join(name)),
1784 &mut pending,
1785 2,
1786 RenameReporting::EachSide,
1787 );
1788 }
1789
1790 assert_eq!(pending.len(), 1);
1791 assert_eq!(
1792 pending.get(Path::new("")),
1793 Some(&Pending::Escalate(InvalidateReason::WatchOverflow))
1794 );
1795 }
1796
1797 #[test]
1798 fn continuous_churn_has_a_deterministic_max_hold_ceiling() {
1799 let past = Instant::now()
1800 .checked_sub(Duration::from_secs(2))
1801 .expect("representable earlier instant");
1802 assert!(max_hold_elapsed(Some(past), Duration::from_secs(1)));
1803 assert!(!max_hold_elapsed(None, Duration::from_secs(1)));
1804 }
1805
1806 #[test]
1807 fn backend_enqueue_is_nonblocking_and_marks_overflow() {
1808 let (sender, receiver) = sync_channel(1);
1809 let overflowed = AtomicBool::new(false);
1810 enqueue_raw(&sender, &overflowed, Ok(notify::Event::new(EventKind::Any)));
1811 enqueue_raw(&sender, &overflowed, Ok(notify::Event::new(EventKind::Any)));
1812
1813 assert!(overflowed.load(Ordering::Acquire));
1814 assert!(matches!(receiver.try_recv(), Ok(RawMessage::Event(Ok(_)))));
1815 }
1816
1817 #[test]
1818 fn full_intent_queue_retains_a_sticky_root_invalidation() {
1819 let (sender, receiver) = sync_channel(1);
1820 sender.try_send(CoalescedIntent::default()).expect("fill output");
1821 let mut pending = BTreeMap::from([(
1822 PathBuf::from("lost.txt"),
1823 Pending::Verify { relist_if_dir: false, renamed: false },
1824 )]);
1825 let mut sticky_overflow = false;
1826
1827 try_deliver_pending(&mut pending, &sender, &mut sticky_overflow).expect("connected");
1828 assert!(pending.is_empty());
1829 assert!(sticky_overflow);
1830
1831 receiver.try_recv().expect("make output capacity");
1832 assert!(try_deliver_overflow(&sender).expect("connected"));
1833 let intent = receiver.try_recv().expect("sticky overflow intent");
1834 let observation = verify_intent(
1835 Path::new("/unused"),
1836 WatchConfig::default(),
1837 &intent,
1838 &ScanConfig::default(),
1839 );
1840 assert!(matches!(
1841 &observation.ops[0].op,
1842 Op::InvalidateSubtree {
1843 path,
1844 reason: InvalidateReason::WatchOverflow,
1845 } if path.as_os_str().is_empty()
1846 ));
1847 }
1848
1849 #[test]
1850 fn cancellation_wakes_and_joins_with_a_full_intent_queue() {
1851 let dir = tempfile::tempdir().expect("tempdir");
1852 let root = dir.path().canonicalize().expect("canonical root");
1853 let config = WatchConfig {
1854 settle: Duration::from_secs(30),
1855 max_hold: Duration::from_secs(30),
1856 event_capacity: 1,
1857 batch_path_capacity: 1,
1858 intent_capacity: 1,
1859 ..WatchConfig::default()
1860 };
1861 let (control, raw) = sync_channel(1);
1862 let (output, intents) = sync_channel(1);
1863 output.try_send(CoalescedIntent::default()).expect("fill intent queue");
1864 let cancelled = Arc::new(AtomicBool::new(false));
1865 let worker_cancelled = Arc::clone(&cancelled);
1866 let status = Arc::new(AtomicU8::new(WORKER_RUNNING));
1867 let tracked_status = Arc::clone(&status);
1868 let overflowed = Arc::new(AtomicBool::new(false));
1869 let worker_overflowed = Arc::clone(&overflowed);
1870 let worker_root = root.clone();
1871 let worker = std::thread::spawn(move || {
1872 run_tracked_worker(&tracked_status, || {
1873 run_worker(
1874 &worker_root,
1875 config,
1876 RenameReporting::EachSide,
1877 &raw,
1878 &output,
1879 &worker_overflowed,
1880 &worker_cancelled,
1881 );
1882 });
1883 });
1884 let watcher = Watcher {
1885 root,
1886 config,
1887 inner: None,
1888 intents,
1889 control: Some(control),
1890 cancelled,
1891 worker_status: status,
1892 worker: Some(worker),
1893 };
1894 let (done_tx, done_rx) = sync_channel(1);
1895
1896 std::thread::spawn(move || {
1897 drop(watcher);
1898 done_tx.send(()).expect("report drop");
1899 });
1900
1901 done_rx
1902 .recv_timeout(Duration::from_secs(5))
1903 .expect("watcher drop must wake and join promptly");
1904 }
1905
1906 #[test]
1907 fn coalescing_defers_filesystem_verification_to_the_consumer() {
1908 let dir = tempfile::tempdir().expect("tempdir");
1909 let relative = PathBuf::from("appeared.txt");
1910 let intent = CoalescedIntent {
1911 pending: BTreeMap::from([(
1912 relative.clone(),
1913 Pending::Verify { relist_if_dir: false, renamed: false },
1914 )]),
1915 };
1916
1917 fs::write(dir.path().join(&relative), b"current").expect("create after coalescing");
1918 let observation =
1919 verify_intent(dir.path(), WatchConfig::default(), &intent, &ScanConfig::default());
1920
1921 assert!(matches!(
1922 &observation.ops[0].op,
1923 Op::Upsert { path, attrs, .. } if path == &relative && attrs.size == 7
1924 ));
1925 }
1926
1927 #[test]
1928 fn control_verification_emits_exact_source_with_the_entry_fact() {
1929 let dir = tempfile::tempdir().expect("tempdir");
1930 let relative = PathBuf::from(".gitignore");
1931 fs::write(dir.path().join(&relative), b"*.log\n").expect("write control");
1932 let intent = CoalescedIntent {
1933 pending: BTreeMap::from([(
1934 relative.clone(),
1935 Pending::Verify { relist_if_dir: false, renamed: false },
1936 )]),
1937 };
1938
1939 let config = ScanConfig { read_controls: true, ..ScanConfig::default() };
1940 let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
1941
1942 assert!(matches!(
1943 &observation.ops[0].op,
1944 Op::Upsert { path, kind: crate::EntryKind::File, .. } if path == &relative
1945 ));
1946 assert!(matches!(
1947 &observation.ops[1].op,
1948 Op::ControlUpsert { path, source } if path == &relative && source == b"*.log\n"
1949 ));
1950 }
1951
1952 #[test]
1953 fn only_population_control_event_reconciles_previously_absent_file() {
1954 let dir = tempfile::tempdir().expect("tree");
1955 fs::write(dir.path().join(".gitignore"), b"# no ignored files\n").expect("control");
1956 fs::write(dir.path().join("debug.log"), b"debug").expect("file");
1957 let scan =
1958 ScanConfig { population: crate::query::IgnoredEntries::Only, ..ScanConfig::default() };
1959 let (index, cold) = crate::scan::scan_into_index(dir.path(), &scan).expect("cold");
1960 assert!(cold.is_complete());
1961 assert!(index.lookup(Path::new("debug.log")).is_none());
1962 let handle = crate::IndexHandle::new(index);
1963 fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("rule edit");
1964 let intent = CoalescedIntent {
1965 pending: BTreeMap::from([(
1966 PathBuf::from(".gitignore"),
1967 Pending::Verify { relist_if_dir: false, renamed: false },
1968 )]),
1969 };
1970 let report =
1971 apply_intent(&handle, dir.path(), WatchConfig::default(), &intent, &scan, &mut |_| {})
1972 .expect("watch apply");
1973 assert!(report.reconciliation.is_complete());
1974 assert!(handle.kind(Path::new("debug.log")).expect("lookup").is_some());
1975 }
1976
1977 #[test]
1985 fn verification_observes_no_control_state_under_a_controls_off_policy() {
1986 let dir = tempfile::tempdir().expect("tempdir");
1987 fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("write control");
1988 fs::write(dir.path().join("debug.log"), b"x").expect("write file");
1989 let config = ScanConfig { read_controls: false, ..ScanConfig::default() };
1990 let (mut index, _) = crate::scan::scan_into_index(dir.path(), &config).expect("scan");
1991 assert!(index.control_table().is_empty());
1992 assert!(matches!(index.controls(), Err(crate::Error::ControlStateNotObserved)));
1993 let control = PathBuf::from(".gitignore");
1994 let intent = CoalescedIntent {
1995 pending: BTreeMap::from([(
1996 control.clone(),
1997 Pending::Verify { relist_if_dir: false, renamed: false },
1998 )]),
1999 };
2000
2001 let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
2002 let reverified = reverify_observation(
2003 dir.path(),
2004 &Observation::new(vec![Op::Remove { path: control }]),
2005 &config,
2006 )
2007 .expect("reverify");
2008
2009 for verified in [&observation, &reverified] {
2010 assert!(
2011 !verified.ops.iter().any(|observed| matches!(
2012 observed.op,
2013 Op::ControlUpsert { .. } | Op::ControlRemove { .. }
2014 )),
2015 "a controls-off policy observed control state: {:?}",
2016 verified.ops
2017 );
2018 }
2019 index.apply(&observation).expect("apply the verified observation");
2020 assert!(index.control_table().is_empty());
2021 assert_eq!(index.scope(), config.scope());
2022 }
2023
2024 #[cfg(unix)]
2025 #[test]
2026 fn applying_verification_uses_the_index_admission_scope() {
2027 use std::os::unix::net::UnixListener;
2028
2029 let dir = tempfile::tempdir().expect("tempdir");
2030 fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("write control");
2031 fs::write(dir.path().join(".secret"), b"hidden").expect("write hidden");
2032 let _listener = UnixListener::bind(dir.path().join("service.sock")).expect("bind socket");
2033 let intent = CoalescedIntent {
2034 pending: [".gitignore", ".secret", "service.sock"]
2035 .into_iter()
2036 .map(|path| {
2037 (PathBuf::from(path), Pending::Verify { relist_if_dir: false, renamed: false })
2038 })
2039 .collect(),
2040 };
2041 let config = ScanConfig {
2042 hidden: Some(Arc::new(crate::HiddenPolicy::prune_hidden::<[&str; 0], &str>([]))),
2043 exclude_special: true,
2044 read_controls: true,
2045 ..ScanConfig::default()
2046 };
2047
2048 let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
2049
2050 assert!(observation.ops.iter().any(|observed| matches!(
2051 &observed.op,
2052 Op::ControlUpsert { path, source }
2053 if path == Path::new(".gitignore") && source == b"*.log\n"
2054 )));
2055 for path in [".gitignore", ".secret", "service.sock"] {
2056 assert!(observation.ops.iter().any(|observed| matches!(
2057 &observed.op,
2058 Op::Remove { path: removed } if removed == Path::new(path)
2059 )));
2060 assert!(!observation.ops.iter().any(|observed| matches!(
2061 &observed.op,
2062 Op::Upsert { path: retained, .. } if retained == Path::new(path)
2063 )));
2064 }
2065 }
2066
2067 #[test]
2068 fn timeout_stop_and_worker_panic_are_distinct() {
2069 let dir = tempfile::tempdir().expect("tempdir");
2070 let (live_sender, live) =
2071 queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2072 assert!(live.next_observation(Duration::ZERO).expect("timeout").is_none());
2073 drop(live_sender);
2074
2075 let (stopped_sender, stopped) =
2076 queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2077 stopped.worker_status.store(WORKER_STOPPED, Ordering::Release);
2078 drop(stopped_sender);
2079 assert!(matches!(stopped.next_observation(Duration::ZERO), Err(Error::WatchStopped)));
2080
2081 let (panicked_sender, panicked) =
2082 queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2083 panicked.worker_status.store(WORKER_PANICKED, Ordering::Release);
2084 drop(panicked_sender);
2085 assert!(matches!(
2086 panicked.next_observation(Duration::ZERO),
2087 Err(Error::WatchWorkerPanicked)
2088 ));
2089 }
2090
2091 #[test]
2092 fn tracked_worker_records_a_panic() {
2093 let status = AtomicU8::new(WORKER_RUNNING);
2094 run_tracked_worker(&status, || panic!("injected worker panic"));
2095 assert_eq!(status.load(Ordering::Acquire), WORKER_PANICKED);
2096 }
2097
2098 #[test]
2099 fn observation_driver_closes_the_invalidation_loop() {
2100 let dir = tempfile::tempdir().expect("tempdir");
2101 let (index, _) =
2102 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2103 let handle = crate::IndexHandle::new(index);
2104 fs::write(dir.path().join("raced.txt"), b"raced").expect("write");
2105 let observation = Observation::new(vec![Op::InvalidateSubtree {
2106 path: PathBuf::new(),
2107 reason: InvalidateReason::WatchSetupRace,
2108 }]);
2109
2110 apply_observation(&handle, &observation, &crate::ScanConfig::default(), &mut |_| {})
2111 .expect("apply and reconcile");
2112
2113 assert!(handle.kind(Path::new("raced.txt")).expect("query").is_some());
2114 }
2115
2116 #[test]
2117 fn applying_driver_reverifies_a_queued_sample_after_reconciliation() {
2118 let dir = tempfile::tempdir().expect("tempdir");
2119 let path = dir.path().join("sample.txt");
2120 fs::write(&path, b"old").expect("write old sample");
2121 let (index, _) =
2122 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2123 let old_attrs = *index.attrs(Path::new("sample.txt")).expect("sample attributes");
2124 let delayed = Observation::new(vec![Op::Upsert {
2125 path: PathBuf::from("sample.txt"),
2126 kind: crate::EntryKind::File,
2127 attrs: old_attrs,
2128 }]);
2129 let handle = crate::IndexHandle::new(index);
2130
2131 fs::write(&path, b"new contents").expect("write current sample");
2132 crate::scan::reconcile_handle(&handle, &crate::ScanConfig::default(), &mut |_| {})
2133 .expect("reconcile newer sample");
2134 let current_size = fs::metadata(&path).expect("sample metadata").len();
2135
2136 apply_observation(&handle, &delayed, &crate::ScanConfig::default(), &mut |_| {})
2137 .expect("apply delayed watch sample");
2138
2139 assert_eq!(
2140 handle.attrs(Path::new("sample.txt")).expect("query").expect("sample remains").size,
2141 current_size
2142 );
2143 }
2144
2145 #[test]
2146 fn blocked_verifier_holds_no_index_lock_and_commits_only_at_current_clock() {
2147 let dir = tempfile::tempdir().expect("tempdir");
2148 let (index, _) =
2149 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2150 let handle = crate::IndexHandle::new(index);
2151 let queued = Observation::new(vec![Op::Upsert {
2152 path: PathBuf::from("queued.txt"),
2153 kind: crate::EntryKind::File,
2154 attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2155 }]);
2156 let applying = handle.clone();
2157 let (entered_tx, entered_rx) = sync_channel(1);
2158 let (release_tx, release_rx) = sync_channel(1);
2159 let (done_tx, done_rx) = sync_channel(1);
2160 let apply_thread = std::thread::spawn(move || {
2161 let mut first = true;
2162 let mut verifier = |_: &Path, observation: &Observation| {
2163 if first {
2164 first = false;
2165 entered_tx.send(()).expect("signal blocked verifier");
2166 release_rx.recv().expect("release blocked verifier");
2167 }
2168 Ok(observation.clone())
2169 };
2170 let result = apply_reverified_with(
2171 &applying,
2172 &queued,
2173 &crate::ScanConfig::default(),
2174 &mut verifier,
2175 );
2176 done_tx.send(result).expect("report apply result");
2177 });
2178
2179 entered_rx
2180 .recv_timeout(Duration::from_secs(5))
2181 .expect("verifier must reach the injected block");
2182 let progressing = handle.clone();
2183 let (progress_tx, progress_rx) = sync_channel(1);
2184 let progress_thread = std::thread::spawn(move || {
2185 let total = progressing.total().expect("reader progresses");
2186 let write = progressing.apply(&Observation::new(vec![Op::Upsert {
2187 path: PathBuf::from("competitor.txt"),
2188 kind: crate::EntryKind::File,
2189 attrs: crate::Attrs { size: 3, allocated: 3, ..crate::Attrs::default() },
2190 }]));
2191 progress_tx.send((total, write)).expect("report progress");
2192 });
2193 let (_, competing_write) = progress_rx
2194 .recv_timeout(Duration::from_secs(5))
2195 .expect("reader and writer must progress while verification is blocked");
2196 competing_write.expect("competing write");
2197 release_tx.send(()).expect("release verifier");
2198
2199 let outcome = done_rx
2200 .recv_timeout(Duration::from_secs(5))
2201 .expect("applying driver completes")
2202 .expect("applying driver succeeds");
2203 apply_thread.join().expect("apply thread");
2204 progress_thread.join().expect("progress thread");
2205 assert_eq!(outcome.inserted, 1);
2206 assert!(handle.kind(Path::new("competitor.txt")).expect("query").is_some());
2207 assert!(handle.kind(Path::new("queued.txt")).expect("query").is_some());
2208 assert_eq!(handle.clock().expect("clock"), crate::Clock(2));
2209 }
2210
2211 #[test]
2212 fn exhausted_watch_contention_stays_unfresh_until_reconciliation() {
2213 let dir = tempfile::tempdir().expect("tempdir");
2214 let (index, _) =
2215 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2216 let handle = crate::IndexHandle::new(index);
2217 let queued = Observation::new(vec![Op::Upsert {
2218 path: PathBuf::from("never-committed.txt"),
2219 kind: crate::EntryKind::File,
2220 attrs: crate::Attrs { size: 1, allocated: 1, ..crate::Attrs::default() },
2221 }]);
2222 let mut attempts = 0_usize;
2223 let mut verifier = |_: &Path, observation: &Observation| {
2224 attempts += 1;
2225 handle
2226 .apply(&Observation::new(vec![Op::Upsert {
2227 path: PathBuf::from(format!("competitor-{attempts}.txt")),
2228 kind: crate::EntryKind::File,
2229 attrs: crate::Attrs {
2230 size: attempts as u64,
2231 allocated: attempts as u64,
2232 ..crate::Attrs::default()
2233 },
2234 }]))
2235 .expect("force a clock conflict");
2236 Ok(observation.clone())
2237 };
2238
2239 let outcome =
2240 apply_reverified_with(&handle, &queued, &crate::ScanConfig::default(), &mut verifier)
2241 .expect("contention escalates");
2242
2243 assert_eq!(attempts, MAX_OPTIMISTIC_APPLY_ATTEMPTS);
2244 assert_eq!(outcome.invalidated, 1);
2245 assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Stale);
2246 assert!(handle.kind(Path::new("never-committed.txt")).expect("query").is_none());
2247 let pending = handle.take_pending_invalidations().expect("pending invalidation");
2248 assert_eq!(pending, vec![(PathBuf::new(), InvalidateReason::WatchContention)]);
2249 handle.restore_pending_invalidations(pending).expect("restore invalidation");
2250
2251 crate::scan::reconcile_pending_handle(&handle, &crate::ScanConfig::default(), &mut |_| {})
2252 .expect("reconcile contention");
2253 assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Fresh);
2254 }
2255
2256 #[test]
2257 fn verifier_error_mutates_no_shared_state() {
2258 let dir = tempfile::tempdir().expect("tempdir");
2259 let (index, _) =
2260 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2261 let handle = crate::IndexHandle::new(index);
2262 let before_clock = handle.clock().expect("clock");
2263 let before_total = handle.total().expect("total");
2264 let mut verifier = |_: &Path, _: &Observation| {
2265 Err(Error::io(
2266 PathBuf::from("blocked"),
2267 std::io::Error::new(std::io::ErrorKind::PermissionDenied, "injected"),
2268 ))
2269 };
2270
2271 let error = apply_reverified_with(
2272 &handle,
2273 &Observation::default(),
2274 &crate::ScanConfig::default(),
2275 &mut verifier,
2276 )
2277 .expect_err("verification error");
2278
2279 assert!(matches!(error, Error::Io { .. }));
2280 assert_eq!(handle.clock().expect("clock"), before_clock);
2281 assert_eq!(handle.total().expect("total"), before_total);
2282 assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Fresh);
2283 assert!(handle.take_pending_invalidations().expect("pending").is_empty());
2284 }
2285
2286 #[test]
2287 fn stable_watch_arbitration_verifies_exactly_once() {
2288 let dir = tempfile::tempdir().expect("tempdir");
2289 let (index, _) =
2290 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2291 let handle = crate::IndexHandle::new(index);
2292 let mut calls = 0_u8;
2293 let mut verifier = |_: &Path, observation: &Observation| {
2294 calls += 1;
2295 Ok(observation.clone())
2296 };
2297
2298 apply_reverified_with(
2299 &handle,
2300 &Observation::new(vec![Op::InvalidateSubtree {
2301 path: PathBuf::new(),
2302 reason: InvalidateReason::Requested,
2303 }]),
2304 &crate::ScanConfig::default(),
2305 &mut verifier,
2306 )
2307 .expect("stable apply");
2308
2309 assert_eq!(calls, 1);
2310 }
2311
2312 #[test]
2313 fn disappearing_watch_root_escalates_instead_of_removing_the_index_root() {
2314 let error = std::io::Error::new(std::io::ErrorKind::NotFound, "root disappeared");
2315
2316 assert!(matches!(
2317 op_for_stat_error(PathBuf::new(), &error),
2318 Op::InvalidateSubtree {
2319 path,
2320 reason: InvalidateReason::VerificationFailed,
2321 } if path.as_os_str().is_empty()
2322 ));
2323 }
2324
2325 #[test]
2326 fn observation_driver_rejects_scope_mismatch_before_apply() {
2327 let dir = tempfile::tempdir().expect("tempdir");
2328 let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2329 let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2330 let handle = crate::IndexHandle::new(index);
2331 let observation = Observation::new(vec![Op::Upsert {
2332 path: PathBuf::from("deep/nested.txt"),
2333 kind: crate::EntryKind::File,
2334 attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2335 }]);
2336
2337 let error =
2338 apply_observation(&handle, &observation, &crate::ScanConfig::default(), &mut |_| {})
2339 .expect_err("mismatched scope must fail");
2340
2341 assert!(matches!(error, Error::ScanScopeMismatch { .. }));
2342 assert!(handle.kind(Path::new("deep/nested.txt")).expect("query").is_none());
2343 }
2344
2345 #[test]
2346 fn observation_driver_rejects_restricted_scopes_until_events_are_filtered() {
2347 let dir = tempfile::tempdir().expect("tempdir");
2348 let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2349 let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2350 let handle = crate::IndexHandle::new(index);
2351 let observation = Observation::new(vec![Op::Upsert {
2352 path: PathBuf::from("deep/nested.txt"),
2353 kind: crate::EntryKind::File,
2354 attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2355 }]);
2356
2357 let error = apply_observation(&handle, &observation, &shallow, &mut |_| {})
2358 .expect_err("unfiltered bounded watch scope must fail");
2359
2360 assert!(matches!(error, Error::UnsupportedScanConfig(_)));
2361 assert!(handle.kind(Path::new("deep/nested.txt")).expect("query").is_none());
2362 }
2363
2364 #[test]
2365 fn apply_next_rejects_restricted_scope_without_consuming_an_observation() {
2366 let dir = tempfile::tempdir().expect("tempdir");
2367 let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2368 let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2369 let handle = crate::IndexHandle::new(index);
2370 let (sender, watcher) =
2371 queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2372 sender.try_send(CoalescedIntent::default()).expect("queue intent");
2373
2374 let error = watcher
2375 .apply_next(&handle, &shallow, Duration::ZERO, &mut |_| {})
2376 .expect_err("restricted scope must fail before receive");
2377
2378 assert!(matches!(error, Error::UnsupportedScanConfig(_)));
2379 assert!(watcher.next_observation(Duration::ZERO).expect("receive").is_some());
2380 }
2381
2382 #[test]
2383 fn apply_next_rejects_a_watcher_for_another_root_without_consuming() {
2384 let indexed = tempfile::tempdir().expect("indexed root");
2385 let watched_root_dir = tempfile::tempdir().expect("watched root");
2386 let (index, _) =
2387 crate::scan::scan_into_index(indexed.path(), &crate::ScanConfig::default())
2388 .expect("scan indexed root");
2389 let handle = crate::IndexHandle::new(index);
2390 let (sender, watcher) = queued_test_watcher(
2391 watched_root_dir.path().canonicalize().expect("canonical watched root"),
2392 );
2393 sender.try_send(CoalescedIntent::default()).expect("queue intent");
2394
2395 let error = watcher
2396 .apply_next(&handle, &crate::ScanConfig::default(), Duration::ZERO, &mut |_| {})
2397 .expect_err("mismatched root must fail");
2398
2399 assert!(matches!(error, Error::WatchRootMismatch { .. }));
2400 assert!(watcher.next_observation(Duration::ZERO).expect("receive").is_some());
2401 }
2402
2403 #[cfg(unix)]
2411 #[test]
2412 fn apply_next_walks_an_unreadable_gap_once_and_retains_its_cause() {
2413 use std::os::unix::fs::PermissionsExt;
2414
2415 fn walks_of(commits: &[Commit], path: &Path) -> usize {
2416 commits
2417 .iter()
2418 .flat_map(|commit| commit.state.iter())
2419 .filter(|transition| {
2420 matches!(
2421 transition,
2422 crate::StateTransition::Freshness { path: marked, current, .. }
2423 if marked == path && *current == crate::Freshness::Reconciling
2424 )
2425 })
2426 .count()
2427 }
2428
2429 if !crate::test_support::require_permission_bits() {
2430 return;
2431 }
2432 let dir = tempfile::tempdir().expect("tempdir");
2433 let root = dir.path().canonicalize().expect("canonical root");
2434 let blocked = root.join("blocked");
2435 fs::create_dir(&blocked).expect("blocked");
2436 fs::write(blocked.join("secret"), b"s").expect("fixture");
2437 let (index, _) =
2438 crate::scan::scan_into_index(&root, &crate::ScanConfig::default()).expect("scan");
2439 let handle = crate::IndexHandle::new(index);
2440 let (sender, watcher) = queued_test_watcher(root.clone());
2441 let config = crate::ScanConfig::default();
2442 let mut commits = Vec::new();
2443
2444 fs::set_permissions(&blocked, fs::Permissions::from_mode(0o000)).expect("deny reads");
2445 let mut gap = CoalescedIntent::default();
2446 gap.pending
2447 .insert(PathBuf::from("blocked"), Pending::Escalate(InvalidateReason::WatchOverflow));
2448 sender.try_send(gap).expect("queue the gap");
2449 let first = watcher
2450 .apply_next(&handle, &config, Duration::ZERO, &mut |commit| {
2451 commits.push(commit.clone());
2452 })
2453 .expect("apply the gap")
2454 .expect("an intent was queued");
2455 for name in ["live.txt", "marker.txt"] {
2456 fs::write(root.join(name), name).expect("unrelated mutation");
2457 let mut event = CoalescedIntent::default();
2458 event.pending.insert(
2459 PathBuf::from(name),
2460 Pending::Verify { relist_if_dir: false, renamed: false },
2461 );
2462 sender.try_send(event).expect("queue the event");
2463 watcher
2464 .apply_next(&handle, &config, Duration::ZERO, &mut |commit| {
2465 commits.push(commit.clone());
2466 })
2467 .expect("apply the event")
2468 .expect("an intent was queued");
2469 }
2470
2471 let pending = handle.take_pending_invalidations().expect("pending");
2472 let freshness = handle.freshness_at(Path::new("blocked")).expect("freshness");
2473 let issues = handle.issues().expect("issues");
2474 fs::set_permissions(&blocked, fs::Permissions::from_mode(0o700)).expect("restore reads");
2475 assert!(!first.reconciliation.is_complete(), "the gap must be unreadable");
2476
2477 assert_eq!(
2478 walks_of(&commits, Path::new("blocked")),
2479 1,
2480 "an unreadable subtree must not be re-walked per unrelated event"
2481 );
2482 assert!(pending.is_empty(), "the queue must settle: {pending:?}");
2483 assert_eq!(freshness, crate::Freshness::Partial);
2484 assert!(handle.kind(Path::new("marker.txt")).expect("lookup").is_some());
2485 assert!(
2486 issues.iter().any(|issue| issue.kind == crate::IssueKind::Permission
2487 && issue.path.as_deref() == Some(Path::new("blocked"))),
2488 "{issues:?}"
2489 );
2490 }
2491
2492 struct RenameFixture {
2494 _dir: tempfile::TempDir,
2495 root: PathBuf,
2496 handle: IndexHandle,
2497 }
2498
2499 impl RenameFixture {
2500 fn new(files: &[&str]) -> Self {
2502 let dir = tempfile::tempdir().expect("tempdir");
2503 let root = dir.path().canonicalize().expect("canonical root");
2504 for file in files {
2505 let path = root.join(file);
2506 fs::create_dir_all(path.parent().expect("parent")).expect("parents");
2507 fs::write(&path, file.as_bytes()).expect("fixture file");
2508 }
2509 let (index, _) =
2510 crate::scan::scan_into_index(&root, &ScanConfig::default()).expect("cold scan");
2511 Self { _dir: dir, root, handle: IndexHandle::new(index) }
2512 }
2513
2514 fn path(&self, rel: &str) -> PathBuf {
2515 self.root.join(rel)
2516 }
2517
2518 fn apply(&self, events: &[notify::Event]) -> Vec<(PathBuf, InvalidateReason)> {
2521 let intent = CoalescedIntent {
2522 pending: recorded(&self.root, RenameReporting::EachSide, events),
2523 };
2524 let mut commits = Vec::new();
2525 let report = apply_intent(
2526 &self.handle,
2527 &self.root,
2528 WatchConfig::default(),
2529 &intent,
2530 &ScanConfig::default(),
2531 &mut |commit| commits.push(commit.clone()),
2532 )
2533 .expect("apply the intent");
2534 assert!(report.reconciliation.is_complete(), "reconciliation must settle");
2535 commits
2536 .iter()
2537 .flat_map(|commit| commit.changes.iter())
2538 .filter_map(|change| match change {
2539 crate::EffectiveChange::Invalidated { path, reason } => {
2540 Some((path.clone(), *reason))
2541 }
2542 _ => None,
2543 })
2544 .collect()
2545 }
2546
2547 fn assert_converged(&self) {
2549 let (cold, _) =
2550 crate::scan::scan_into_index(&self.root, &ScanConfig::default()).expect("cold");
2551 let watched = self.handle.read_with(entries).expect("read the watched index");
2552 assert_eq!(watched, entries(&cold), "watched index diverged from a cold scan");
2553 assert_eq!(self.handle.freshness().expect("freshness"), crate::Freshness::Fresh);
2554 }
2555 }
2556
2557 fn entries(index: &crate::Index) -> BTreeMap<PathBuf, (crate::EntryKind, Option<u64>)> {
2561 fn walk(
2562 index: &crate::Index,
2563 dir: &Path,
2564 out: &mut BTreeMap<PathBuf, (crate::EntryKind, Option<u64>)>,
2565 ) {
2566 let Some(children) = index.children(dir) else {
2567 return;
2568 };
2569 let children: Vec<_> =
2570 children.map(|(name, id)| (dir.join(name), id)).collect::<Vec<_>>();
2571 for (path, id) in children {
2572 let kind = index.kind_of(id).expect("live child");
2573 let size = (!kind.is_dir()).then(|| index.attrs_of(id).expect("attrs").size);
2574 out.insert(path.clone(), (kind, size));
2575 walk(index, &path, out);
2576 }
2577 }
2578 let mut out = BTreeMap::new();
2579 walk(index, Path::new(""), &mut out);
2580 out
2581 }
2582
2583 type Classified = (BTreeMap<PathBuf, Option<bool>>, Vec<(PathBuf, Vec<u8>)>);
2586
2587 fn classified(index: &crate::Index) -> Classified {
2588 let ignored = entries(index)
2589 .into_keys()
2590 .map(|path| {
2591 let ignored = index.is_ignored(&path).expect("observed");
2592 (path, ignored)
2593 })
2594 .collect();
2595 let sources = index
2596 .controls()
2597 .expect("observed")
2598 .sources()
2599 .map(|(path, source)| (path, source.to_vec()))
2600 .collect();
2601 (ignored, sources)
2602 }
2603
2604 impl RenameFixture {
2605 fn assert_classified_as_cold(&self, label: &str) -> Classified {
2608 let (cold, _) =
2609 crate::scan::scan_into_index(&self.root, &ScanConfig::default()).expect("cold");
2610 let watched = self.handle.read_with(classified).expect("read the watched index");
2611 assert_eq!(watched, classified(&cold), "{label}: the watch diverged from a cold scan");
2612 watched
2613 }
2614 }
2615
2616 #[test]
2624 fn a_watch_follows_a_case_variant_control_file() {
2625 use crate::test_support::CaseLookups;
2626
2627 let probe = tempfile::tempdir().expect("tempdir");
2628 for (lookups, governs) in CaseLookups::on_this_host(probe.path()) {
2629 let tree = RenameFixture::new(&["up/x.tmp", "up/notes.txt"]);
2630 let _lookups = lookups.install(&tree.root);
2631 let variant = tree.path("up/.GITIGNORE");
2632 let x_ignored = |classified: &Classified| classified.0[Path::new("up/x.tmp")];
2633
2634 fs::write(&variant, b"*.tmp\n").expect("create the variant");
2635 tree.apply(&[created(variant.clone())]);
2636 let appeared = tree.assert_classified_as_cold(&format!("{lookups:?}: created"));
2637 assert_eq!(x_ignored(&appeared), Some(governs), "{lookups:?}");
2638
2639 fs::write(&variant, b"*.txt\n").expect("edit the variant");
2640 tree.apply(&[modified(variant.clone())]);
2641 let edited = tree.assert_classified_as_cold(&format!("{lookups:?}: edited"));
2642 assert_eq!(x_ignored(&edited), Some(false), "{lookups:?}");
2643 assert_eq!(edited.0[Path::new("up/notes.txt")], Some(governs), "{lookups:?}");
2644
2645 fs::remove_file(&variant).expect("remove the variant");
2646 tree.apply(&[notify::Event::new(EventKind::Remove(notify::event::RemoveKind::File))
2647 .add_path(variant.clone())]);
2648 let removed = tree.assert_classified_as_cold(&format!("{lookups:?}: removed"));
2649 assert_eq!(removed.1, Vec::new(), "{lookups:?}: no rules remain");
2650
2651 if lookups == CaseLookups::Host {
2652 let exact = tree.path("up/.gitignore");
2653 fs::write(&exact, b"*.tmp\n").expect("create the exact name");
2654 tree.apply(&[created(exact.clone())]);
2655 for (from, to) in [(&exact, &variant), (&variant, &exact)] {
2656 fs::rename(from, to).expect("case-only rename");
2657 tree.apply(&[fsevents_rename(from.clone()), fsevents_rename(to.clone())]);
2658 tree.assert_converged();
2659 tree.assert_classified_as_cold(&format!("{lookups:?}: renamed to {to:?}"));
2660 }
2661 }
2662 }
2663 }
2664
2665 fn created(path: PathBuf) -> notify::Event {
2666 notify::Event::new(EventKind::Create(CreateKind::Any)).add_path(path)
2667 }
2668
2669 fn modified(path: PathBuf) -> notify::Event {
2670 notify::Event::new(EventKind::Modify(ModifyKind::Data(notify::event::DataChange::Content)))
2671 .add_path(path)
2672 }
2673
2674 #[test]
2677 fn a_file_renamed_within_its_directory_needs_no_invalidation() {
2678 let tree = RenameFixture::new(&["state/config.json", "state/other.txt"]);
2679 fs::write(tree.path("state/config.json.tmp"), b"new configuration").expect("temp");
2680 fs::rename(tree.path("state/config.json.tmp"), tree.path("state/config.json"))
2681 .expect("rename over");
2682
2683 let invalidations = tree.apply(&[
2684 created(tree.path("state/config.json.tmp")),
2685 modified(tree.path("state/config.json.tmp")),
2686 fsevents_rename(tree.path("state/config.json.tmp")),
2687 fsevents_rename(tree.path("state/config.json")),
2688 ]);
2689
2690 assert_eq!(invalidations, vec![], "a file rename is settled by its own paths");
2691 tree.assert_converged();
2692 }
2693
2694 #[test]
2697 fn a_file_moved_across_directories_settles_one_side_per_event() {
2698 let tree = RenameFixture::new(&["from/moved.txt", "to/resident.txt"]);
2699 fs::rename(tree.path("from/moved.txt"), tree.path("to/moved.txt")).expect("move");
2700
2701 assert_eq!(tree.apply(&[fsevents_rename(tree.path("to/moved.txt"))]), vec![]);
2702 assert_eq!(tree.apply(&[fsevents_rename(tree.path("from/moved.txt"))]), vec![]);
2703 tree.assert_converged();
2704 }
2705
2706 #[test]
2709 fn a_renamed_directory_relists_only_its_own_subtree() {
2710 let tree = RenameFixture::new(&[
2711 "project/old/one.txt",
2712 "project/old/nested/two.txt",
2713 "project/untouched/three.txt",
2714 ]);
2715 fs::rename(tree.path("project/old"), tree.path("project/new")).expect("rename directory");
2716
2717 let invalidations = tree.apply(&[
2718 fsevents_rename(tree.path("project/old")),
2719 fsevents_rename(tree.path("project/new")),
2720 ]);
2721
2722 assert_eq!(
2723 invalidations,
2724 vec![(PathBuf::from("project/new"), InvalidateReason::UnpairedRename)]
2725 );
2726 tree.assert_converged();
2727 }
2728
2729 #[test]
2732 fn moves_across_the_root_boundary_settle_from_the_inside_side() {
2733 let tree = RenameFixture::new(&["resident.txt", "leaving/a.txt", "leaving/deep/b.txt"]);
2734 let outside = tempfile::tempdir().expect("outside the root");
2735 fs::create_dir_all(outside.path().join("arriving/deep")).expect("outside tree");
2736 fs::write(outside.path().join("arriving/deep/c.txt"), b"arrived").expect("outside file");
2737
2738 fs::rename(outside.path().join("arriving"), tree.path("arrived")).expect("move in");
2739 let invalidations = tree.apply(&[fsevents_rename(tree.path("arrived"))]);
2740 assert_eq!(
2741 invalidations,
2742 vec![(PathBuf::from("arrived"), InvalidateReason::UnpairedRename)]
2743 );
2744 tree.assert_converged();
2745
2746 fs::rename(tree.path("leaving"), outside.path().join("left")).expect("move out");
2747 assert_eq!(tree.apply(&[fsevents_rename(tree.path("leaving"))]), vec![]);
2748 tree.assert_converged();
2749 }
2750
2751 #[test]
2754 fn a_name_reused_after_a_rename_verifies_as_its_new_entry() {
2755 let tree = RenameFixture::new(&["log.txt", "cache/entry.bin"]);
2756 fs::rename(tree.path("log.txt"), tree.path("log.1.txt")).expect("rotate");
2757 fs::write(tree.path("log.txt"), b"fresh log, longer than the old one").expect("reuse");
2758 fs::rename(tree.path("cache"), tree.path("cache.old")).expect("retire directory");
2759 fs::create_dir(tree.path("cache")).expect("reuse directory name");
2760
2761 let invalidations = tree.apply(&[
2762 fsevents_rename(tree.path("log.txt")),
2763 created(tree.path("log.txt")),
2764 fsevents_rename(tree.path("log.1.txt")),
2765 fsevents_rename(tree.path("cache")),
2766 created(tree.path("cache")),
2767 fsevents_rename(tree.path("cache.old")),
2768 ]);
2769
2770 assert!(
2771 invalidations.iter().all(|(path, _)| !path.as_os_str().is_empty()),
2772 "no root reconcile: {invalidations:?}"
2773 );
2774 tree.assert_converged();
2775 assert_eq!(tree.handle.kind(Path::new("cache/entry.bin")).expect("lookup"), None);
2776 }
2777
2778 #[test]
2781 fn a_sticky_rename_flag_on_a_later_write_is_just_a_write() {
2782 let tree = RenameFixture::new(&["sessions/today.jsonl"]);
2783 fs::write(tree.path("sessions/today.jsonl"), b"appended record after the rename")
2784 .expect("write in place");
2785
2786 let invalidations = tree.apply(&[
2787 fsevents_rename(tree.path("sessions/today.jsonl")),
2788 modified(tree.path("sessions/today.jsonl")),
2789 ]);
2790
2791 assert_eq!(invalidations, vec![]);
2792 tree.assert_converged();
2793 }
2794
2795 #[test]
2799 fn a_case_only_rename_keeps_one_spelling() {
2800 let tree = RenameFixture::new(&["docs/Readme.md", "docs/Guide/intro.md"]);
2801 let insensitive = crate::test_support::resolves_case_insensitively(&tree.root);
2802 fs::rename(tree.path("docs/Readme.md"), tree.path("docs/README.md")).expect("recase");
2803 fs::rename(tree.path("docs/Guide"), tree.path("docs/guide")).expect("recase directory");
2804
2805 let invalidations = tree.apply(&[
2806 fsevents_rename(tree.path("docs/Readme.md")),
2807 fsevents_rename(tree.path("docs/README.md")),
2808 fsevents_rename(tree.path("docs/Guide")),
2809 fsevents_rename(tree.path("docs/guide")),
2810 modified(tree.path("docs/Guide/intro.md")),
2812 ]);
2813
2814 tree.assert_converged();
2815 assert!(
2816 invalidations.iter().all(|(path, _)| !path.as_os_str().is_empty()),
2817 "no root reconcile: {invalidations:?}"
2818 );
2819 if insensitive {
2820 assert!(
2821 invalidations.contains(&(PathBuf::from("docs"), InvalidateReason::UnpairedRename)),
2822 "the stale spelling reconciles its parent: {invalidations:?}"
2823 );
2824 }
2825 }
2826
2827 #[test]
2830 fn a_listing_miss_rereads_the_parent_before_answering() {
2831 let dir = tempfile::tempdir().expect("tempdir");
2832 let root = dir.path();
2833 fs::write(root.join("first"), b"1").expect("first");
2834 let mut listings = ParentListings::default();
2835
2836 assert_eq!(listings.lists(root, Path::new("first")), Some(true));
2837 fs::write(root.join("created-after-the-listing"), b"2").expect("later entry");
2838 assert_eq!(listings.lists(root, Path::new("created-after-the-listing")), Some(true));
2839 assert_eq!(listings.lists(root, Path::new("never-created")), Some(false));
2840 assert_eq!(listings.lists(root, Path::new("missing-directory/child")), None);
2841 }
2842
2843 #[test]
2846 fn native_renames_converge_without_reconciling_the_root() {
2847 let _serialized = real_watcher_guard();
2848 let dir = tempfile::tempdir().expect("tempdir");
2849 let root = dir.path().canonicalize().expect("canonical root");
2850 for file in ["from/moved.txt", "tree/sub/leaf.txt", "to/resident.txt"] {
2851 fs::create_dir_all(root.join(file).parent().expect("parent")).expect("parents");
2852 fs::write(root.join(file), file.as_bytes()).expect("fixture file");
2853 }
2854 let watcher = Watcher::new(&root, WatchConfig::default()).expect("watcher");
2855 if !establish_watch(&watcher, &root) {
2856 return;
2857 }
2858 let config = ScanConfig::default();
2859 let (index, _) = crate::scan::scan_into_index(&root, &config).expect("scan");
2860 let handle = IndexHandle::new(index);
2861
2862 fs::rename(root.join("from/moved.txt"), root.join("to/moved.txt")).expect("move file");
2863 fs::rename(root.join("tree"), root.join("renamed-tree")).expect("rename directory");
2864
2865 let converged = || {
2866 handle.kind(Path::new("to/moved.txt")).expect("lookup").is_some()
2867 && handle.kind(Path::new("from/moved.txt")).expect("lookup").is_none()
2868 && handle.kind(Path::new("renamed-tree/sub/leaf.txt")).expect("lookup").is_some()
2869 && handle.kind(Path::new("tree")).expect("lookup").is_none()
2870 };
2871 let mut commits = Vec::new();
2872 let start = Instant::now();
2873 while !converged() && start.elapsed() < REAL_BACKEND_DELIVERY {
2874 watcher
2875 .apply_next(&handle, &config, Duration::from_millis(200), &mut |commit| {
2876 commits.push(commit.clone());
2877 })
2878 .expect("apply");
2879 }
2880
2881 assert!(converged(), "the renames never converged: {commits:?}");
2882 let root_reconciles = commits
2883 .iter()
2884 .flat_map(|commit| commit.changes.iter())
2885 .filter(|change| {
2886 matches!(change, crate::EffectiveChange::Invalidated { path, .. }
2887 if path.as_os_str().is_empty())
2888 })
2889 .count();
2890 assert_eq!(root_reconciles, 0, "a rename reconciled the whole root: {commits:?}");
2891 let (cold, _) = crate::scan::scan_into_index(&root, &config).expect("cold");
2892 assert_eq!(handle.read_with(entries).expect("read"), entries(&cold));
2893 }
2894
2895 #[test]
2896 fn zero_settle_is_rejected_before_starting_a_busy_worker() {
2897 let config = WatchConfig { settle: Duration::ZERO, ..WatchConfig::default() };
2898
2899 assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2900
2901 let config = WatchConfig { intent_capacity: 0, ..WatchConfig::default() };
2902 assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2903
2904 let config = WatchConfig {
2905 batch_path_capacity: MAX_BATCH_PATH_CAPACITY,
2906 intent_capacity: MAX_BUFFERED_INTENT_PATHS / MAX_BATCH_PATH_CAPACITY + 1,
2907 ..WatchConfig::default()
2908 };
2909 assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2910 }
2911}