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 std::fs::symlink_metadata(&absolute) {
681 Ok(metadata) => {
682 let (kind, attrs) = scan::observe(&absolute, &metadata)
683 .map_err(|source| Error::io(&absolute, source))?;
684 match crate::admission::decide_path(
685 &relative,
686 kind,
687 scan_config.hidden(),
688 scan_config.exclude_special,
689 ) {
690 crate::admission::Disposition::Retain => {
691 ops.push(Op::Upsert { path: relative.clone(), kind, attrs });
692 if let Some(control) =
693 scan::read_control_op(scan_config, root, &relative, kind)?
694 {
695 ops.push(control);
696 }
697 }
698 crate::admission::Disposition::ControlOnly => {
699 ops.push(Op::Remove { path: relative.clone() });
700 if let Some(control) =
701 scan::read_control_op(scan_config, root, &relative, kind)?
702 {
703 ops.push(control);
704 }
705 }
706 crate::admission::Disposition::Reject => {
707 ops.push(Op::Remove { path: relative });
708 }
709 }
710 }
711 Err(error) => ops.push(op_for_stat_error(relative, &error)),
712 }
713 }
714 Op::ControlUpsert { source, .. } => {
715 ops.push(Op::ControlUpsert { path: relative, source: source.clone() });
716 }
717 Op::ControlRemove { .. } => ops.push(Op::ControlRemove { path: relative }),
718 }
719 }
720 Ok(Observation::new(ops))
721}
722
723impl Drop for Watcher {
724 fn drop(&mut self) {
725 self.cancelled.store(true, Ordering::Release);
726 if let Some(control) = self.control.take() {
727 let _ = control.try_send(RawMessage::Stop);
728 }
729 self.inner.take();
732 if let Some(worker) = self.worker.take() {
733 let _ = worker.join();
734 }
735 }
736}
737
738fn notify_error(path: &Path, err: notify::Error) -> Error {
739 Error::io(path, std::io::Error::other(err))
740}
741
742fn enqueue_raw(
743 sender: &SyncSender<RawMessage>,
744 overflowed: &AtomicBool,
745 event: notify::Result<notify::Event>,
746) {
747 match sender.try_send(RawMessage::Event(event)) {
748 Ok(()) | Err(TrySendError::Disconnected(_)) => {}
749 Err(TrySendError::Full(_)) => overflowed.store(true, Ordering::Release),
750 }
751}
752
753fn run_tracked_worker(status: &AtomicU8, worker: impl FnOnce()) {
754 let outcome = catch_unwind(AssertUnwindSafe(worker));
755 status.store(if outcome.is_ok() { WORKER_STOPPED } else { WORKER_PANICKED }, Ordering::Release);
756}
757
758fn run_worker(
759 root: &Path,
760 config: WatchConfig,
761 renames: RenameReporting,
762 raw: &Receiver<RawMessage>,
763 out: &SyncSender<CoalescedIntent>,
764 overflowed: &AtomicBool,
765 cancelled: &AtomicBool,
766) {
767 let mut pending: BTreeMap<PathBuf, Pending> = BTreeMap::new();
768 let mut batch_started: Option<Instant> = None;
769 let mut sticky_overflow = false;
770
771 loop {
772 if cancelled.load(Ordering::Acquire) {
773 return;
774 }
775 if overflowed.swap(false, Ordering::AcqRel) {
776 collapse_to_overflow(&mut pending);
777 batch_started.get_or_insert_with(Instant::now);
778 }
779 if sticky_overflow {
780 match try_deliver_overflow(out) {
781 Ok(true) => sticky_overflow = false,
782 Ok(false) => {}
783 Err(()) => return,
784 }
785 }
786
787 match raw.recv_timeout(config.settle) {
788 Ok(RawMessage::Event(Ok(event))) => {
789 record(root, &event, &mut pending, config.batch_path_capacity, renames);
790 batch_started.get_or_insert_with(Instant::now);
791 }
792 Ok(RawMessage::Event(Err(err))) => {
793 let _ = err;
796 collapse_to_overflow(&mut pending);
797 batch_started.get_or_insert_with(Instant::now);
798 }
799 Ok(RawMessage::Flush(acknowledge)) => {
800 if overflowed.swap(false, Ordering::AcqRel) {
801 collapse_to_overflow(&mut pending);
802 }
803 if !pending.is_empty()
804 && try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err()
805 {
806 return;
807 }
808 if sticky_overflow {
812 match try_deliver_overflow(out) {
813 Ok(true) => sticky_overflow = false,
814 Ok(false) => {}
815 Err(()) => return,
816 }
817 }
818 batch_started = None;
819 let _ = acknowledge.send(());
820 }
821 Ok(RawMessage::Stop) => return,
822 Err(RecvTimeoutError::Timeout) => {
823 if !pending.is_empty()
825 && try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err()
826 {
827 return;
828 }
829 batch_started = None;
830 continue;
831 }
832 Err(RecvTimeoutError::Disconnected) => {
833 if !cancelled.load(Ordering::Acquire) && !pending.is_empty() {
834 let _ = try_deliver_pending(&mut pending, out, &mut sticky_overflow);
835 }
836 return;
837 }
838 }
839
840 if max_hold_elapsed(batch_started, config.max_hold) {
842 if try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err() {
843 return;
844 }
845 batch_started = None;
846 }
847 }
848}
849
850fn max_hold_elapsed(started: Option<Instant>, max_hold: Duration) -> bool {
851 started.is_some_and(|start| start.elapsed() >= max_hold)
852}
853
854fn record(
856 root: &Path,
857 event: ¬ify::Event,
858 pending: &mut BTreeMap<PathBuf, Pending>,
859 capacity: usize,
860 reporting: RenameReporting,
861) {
862 if event.need_rescan() {
863 let target = if event.paths.len() == 1 {
866 relative_to(root, &event.paths[0]).unwrap_or_default()
867 } else {
868 PathBuf::new()
869 };
870 queue_pending(
871 pending,
872 target,
873 Pending::Escalate(InvalidateReason::WatchOverflow),
874 capacity,
875 );
876 return;
877 }
878
879 let renamed = matches!(event.kind, EventKind::Modify(notify::event::ModifyKind::Name(_)));
896 if renamed
897 && (reporting == RenameReporting::OldSideOnly
898 || !event.paths.iter().any(|path| relative_to(root, path).is_some()))
899 {
900 queue_pending(
906 pending,
907 PathBuf::new(),
908 Pending::Escalate(InvalidateReason::UnpairedRename),
909 capacity,
910 );
911 }
912
913 for path in &event.paths {
914 let Some(rel) = relative_to(root, path) else {
915 continue; };
917 if matches!(event.kind, EventKind::Access(_)) {
918 continue; }
920 if renamed && rel.as_os_str().is_empty() {
921 queue_pending(
925 pending,
926 PathBuf::new(),
927 Pending::Escalate(InvalidateReason::UnpairedRename),
928 capacity,
929 );
930 continue;
931 }
932 let relist_if_dir = matches!(event.kind, EventKind::Create(_));
933 queue_pending(pending, rel, Pending::Verify { relist_if_dir, renamed }, capacity);
934 }
935}
936
937fn queue_pending(
938 pending: &mut BTreeMap<PathBuf, Pending>,
939 path: PathBuf,
940 state: Pending,
941 capacity: usize,
942) {
943 if matches!(pending.get(Path::new("")), Some(Pending::Escalate(_))) {
944 return;
945 }
946 if path.as_os_str().is_empty() && matches!(state, Pending::Escalate(_)) {
947 pending.clear();
948 pending.insert(path, state);
949 return;
950 }
951 if let Some(existing) = pending.get_mut(&path) {
952 match (existing, state) {
953 (Pending::Escalate(_), _) => {}
954 (
955 Pending::Verify { relist_if_dir, renamed },
956 Pending::Verify { relist_if_dir: relist, renamed: rename },
957 ) => {
958 *relist_if_dir |= relist;
959 *renamed |= rename;
960 }
961 (slot @ Pending::Verify { .. }, Pending::Escalate(reason)) => {
962 *slot = Pending::Escalate(reason);
963 }
964 }
965 return;
966 }
967 if pending.len() >= capacity {
968 collapse_to_overflow(pending);
969 } else {
970 pending.insert(path, state);
971 }
972}
973
974fn collapse_to_overflow(pending: &mut BTreeMap<PathBuf, Pending>) {
975 pending.clear();
976 pending.insert(PathBuf::new(), Pending::Escalate(InvalidateReason::WatchOverflow));
977}
978
979fn try_deliver_pending(
980 pending: &mut BTreeMap<PathBuf, Pending>,
981 out: &SyncSender<CoalescedIntent>,
982 sticky_overflow: &mut bool,
983) -> std::result::Result<(), ()> {
984 if pending.is_empty() {
985 return Ok(());
986 }
987 let intent = CoalescedIntent { pending: std::mem::take(pending) };
988 match out.try_send(intent) {
989 Ok(()) => Ok(()),
990 Err(TrySendError::Full(_)) => {
991 *sticky_overflow = true;
992 Ok(())
993 }
994 Err(TrySendError::Disconnected(_)) => Err(()),
995 }
996}
997
998fn try_deliver_overflow(out: &SyncSender<CoalescedIntent>) -> std::result::Result<bool, ()> {
999 let mut pending = BTreeMap::new();
1000 collapse_to_overflow(&mut pending);
1001 match out.try_send(CoalescedIntent { pending }) {
1002 Ok(()) => Ok(true),
1003 Err(TrySendError::Full(_)) => Ok(false),
1004 Err(TrySendError::Disconnected(_)) => Err(()),
1005 }
1006}
1007
1008fn verify_intent(
1010 root: &Path,
1011 config: WatchConfig,
1012 intent: &CoalescedIntent,
1013 scan_config: &ScanConfig,
1014) -> Observation {
1015 let mut ops = Vec::with_capacity(intent.pending.len());
1016 let mut listings = ParentListings::default();
1017 let mut unlisted: Vec<PathBuf> = Vec::new();
1022
1023 for (rel, state) in &intent.pending {
1024 match state {
1025 Pending::Escalate(reason) => {
1026 ops.push(Op::InvalidateSubtree { path: rel.clone(), reason: *reason });
1027 }
1028 Pending::Verify { relist_if_dir, renamed } => {
1029 if unlisted.iter().any(|name| rel.starts_with(name)) {
1030 continue;
1031 }
1032 let absolute = root.join(rel);
1033 let mut stat = std::fs::symlink_metadata(&absolute);
1034 if *renamed && !rel.as_os_str().is_empty() && stat.is_ok() {
1035 let unlisted_reason = match listings.lists(root, rel) {
1042 Some(true) => None,
1043 Some(false) => {
1044 stat = std::fs::symlink_metadata(&absolute);
1049 stat.is_ok().then_some(InvalidateReason::UnpairedRename)
1050 }
1051 None => Some(InvalidateReason::VerificationFailed),
1054 };
1055 if let Some(reason) = unlisted_reason {
1056 let parent = parent_of(rel);
1059 if !ops.iter().any(
1060 |op| matches!(op, Op::InvalidateSubtree { path, .. } if *path == parent),
1061 ) {
1062 ops.push(Op::InvalidateSubtree { path: parent, reason });
1063 }
1064 unlisted.push(rel.clone());
1065 continue;
1066 }
1067 }
1068 match stat {
1069 Ok(meta) => {
1070 let Ok((kind, attrs)) = scan::observe(&absolute, &meta) else {
1071 ops.push(Op::InvalidateSubtree {
1072 path: rel.parent().map_or_else(PathBuf::new, Path::to_path_buf),
1073 reason: InvalidateReason::VerificationFailed,
1074 });
1075 continue;
1076 };
1077 let disposition = crate::admission::decide_path(
1078 rel,
1079 kind,
1080 scan_config.hidden(),
1081 scan_config.exclude_special,
1082 );
1083 match disposition {
1084 crate::admission::Disposition::Retain => {
1085 ops.push(Op::Upsert { path: rel.clone(), kind, attrs });
1086 match scan::read_control_op(scan_config, root, rel, kind) {
1087 Ok(Some(control)) => ops.push(control),
1088 Ok(None) => {}
1089 Err(_) => ops.push(Op::InvalidateSubtree {
1090 path: rel
1091 .parent()
1092 .map_or_else(PathBuf::new, Path::to_path_buf),
1093 reason: InvalidateReason::VerificationFailed,
1094 }),
1095 }
1096 }
1097 crate::admission::Disposition::ControlOnly => {
1098 match scan::read_control_op(scan_config, root, rel, kind) {
1099 Ok(Some(control)) => {
1100 ops.push(Op::Remove { path: rel.clone() });
1101 ops.push(control);
1102 }
1103 Ok(None) => ops.push(Op::Remove { path: rel.clone() }),
1104 Err(_) => ops.push(Op::InvalidateSubtree {
1105 path: rel
1106 .parent()
1107 .map_or_else(PathBuf::new, Path::to_path_buf),
1108 reason: InvalidateReason::VerificationFailed,
1109 }),
1110 }
1111 }
1112 crate::admission::Disposition::Reject => {
1113 ops.push(Op::Remove { path: rel.clone() });
1114 }
1115 }
1116 let retained_dir =
1117 disposition == crate::admission::Disposition::Retain && kind.is_dir();
1118 if retained_dir && *renamed {
1119 ops.push(Op::InvalidateSubtree {
1125 path: rel.clone(),
1126 reason: InvalidateReason::UnpairedRename,
1127 });
1128 } else if retained_dir && *relist_if_dir && config.relist_new_dirs {
1129 ops.push(Op::InvalidateSubtree {
1132 path: rel.clone(),
1133 reason: InvalidateReason::WatchSetupRace,
1134 });
1135 }
1136 }
1137 Err(error) => ops.push(op_for_stat_error(rel.clone(), &error)),
1138 }
1139 if scan_config.population != crate::query::IgnoredEntries::Include
1140 && crate::control::is_control_file(rel)
1141 {
1142 ops.push(Op::InvalidateSubtree {
1143 path: rel.parent().map_or_else(PathBuf::new, Path::to_path_buf),
1144 reason: InvalidateReason::ControlPopulationChanged,
1145 });
1146 }
1147 }
1148 }
1149 }
1150 Observation::new(ops)
1151}
1152
1153fn parent_of(rel: &Path) -> PathBuf {
1155 rel.parent().map_or_else(PathBuf::new, Path::to_path_buf)
1156}
1157
1158#[derive(Default)]
1164struct ParentListings(BTreeMap<PathBuf, Option<std::collections::HashSet<std::ffi::OsString>>>);
1165
1166impl ParentListings {
1167 fn lists(&mut self, root: &Path, rel: &Path) -> Option<bool> {
1174 let name = rel.file_name()?;
1175 let parent = parent_of(rel);
1176 if let Some(Some(names)) = self.0.get(&parent) {
1177 if names.contains(name) {
1178 return Some(true);
1179 }
1180 }
1181 let names: Option<std::collections::HashSet<_>> = std::fs::read_dir(root.join(&parent))
1182 .and_then(|entries| entries.map(|entry| entry.map(|entry| entry.file_name())).collect())
1183 .ok();
1184 let listed = names.as_ref().map(|names| names.contains(name));
1185 self.0.insert(parent, names);
1186 listed
1187 }
1188}
1189
1190fn op_for_stat_error(path: PathBuf, error: &std::io::Error) -> Op {
1191 match error.kind() {
1192 std::io::ErrorKind::NotFound if path.as_os_str().is_empty() => {
1193 Op::InvalidateSubtree { path, reason: InvalidateReason::VerificationFailed }
1194 }
1195 std::io::ErrorKind::NotFound => Op::Remove { path },
1196 std::io::ErrorKind::NotADirectory => Op::InvalidateSubtree {
1197 path: path.parent().map_or_else(PathBuf::new, Path::to_path_buf),
1198 reason: InvalidateReason::VerificationFailed,
1199 },
1200 _ => Op::InvalidateSubtree { path, reason: InvalidateReason::VerificationFailed },
1201 }
1202}
1203
1204fn relative_to(root: &Path, path: &Path) -> Option<PathBuf> {
1209 path.strip_prefix(root).ok().map(Path::to_path_buf)
1210}
1211
1212#[cfg(test)]
1213mod tests {
1214 use super::*;
1215 use notify::event::{CreateKind, Flag, MetadataKind, ModifyKind, RenameMode};
1216 use std::fs;
1217
1218 fn queued_test_watcher(root: PathBuf) -> (SyncSender<CoalescedIntent>, Watcher) {
1219 let (sender, intents) = sync_channel(1);
1220 let watcher = Watcher {
1221 root,
1222 config: WatchConfig::default(),
1223 inner: None,
1224 intents,
1225 control: None,
1226 cancelled: Arc::new(AtomicBool::new(false)),
1227 worker_status: Arc::new(AtomicU8::new(WORKER_RUNNING)),
1228 worker: None,
1229 };
1230 (sender, watcher)
1231 }
1232
1233 #[test]
1234 fn watcher_can_move_to_its_single_consumer_thread() {
1235 fn assert_send<T: Send>() {}
1236
1237 assert_send::<Watcher>();
1238 }
1239
1240 static REAL_WATCHER: std::sync::Mutex<()> = std::sync::Mutex::new(());
1268
1269 const REAL_BACKEND_DELIVERY: Duration = Duration::from_secs(60);
1283
1284 fn real_watcher_guard() -> std::sync::MutexGuard<'static, ()> {
1290 REAL_WATCHER.lock().unwrap_or_else(std::sync::PoisonError::into_inner)
1291 }
1292
1293 enum Waited {
1295 Delivered(Vec<Op>),
1297 Silent,
1299 }
1300
1301 fn wait_for(
1322 watcher: &Watcher,
1323 deadline: Duration,
1324 mut want: impl FnMut(&[Op]) -> bool,
1325 ) -> Waited {
1326 let start = Instant::now();
1327 let mut seen: Vec<Op> = Vec::new();
1328 while start.elapsed() < deadline {
1329 match watcher.next_observation(Duration::from_millis(200)) {
1330 Ok(Some(observation)) => {
1331 seen.extend(observation.ops.into_iter().map(|observed| observed.op));
1332 if want(&seen) {
1333 return Waited::Delivered(seen);
1334 }
1335 }
1336 Ok(None) => {}
1337 Err(error) => panic!("watcher stopped while waiting: {error}"),
1338 }
1339 }
1340 if seen.is_empty() {
1341 return Waited::Silent;
1342 }
1343 panic!(
1344 "the backend delivered {} op(s) in {deadline:?} but never the one awaited, so \
1345 this is a disagreement about content rather than a delivery failure: {seen:?}",
1346 seen.len()
1347 );
1348 }
1349
1350 fn wait_established(
1356 watcher: &Watcher,
1357 deadline: Duration,
1358 want: impl FnMut(&[Op]) -> bool,
1359 ) -> Vec<Op> {
1360 match wait_for(watcher, deadline, want) {
1361 Waited::Delivered(ops) => ops,
1362 Waited::Silent => panic!(
1363 "the watch was established and then delivered nothing in {deadline:?}, so \
1364 this is a lost event rather than a host precondition; \
1365 FDU_TEST_ALLOW_NO_NATIVE_WATCH does not apply here"
1366 ),
1367 }
1368 }
1369
1370 fn establish_watch(watcher: &Watcher, dir: &Path) -> bool {
1393 let warmup = dir.join(".fdu-watch-warmup");
1394 fs::write(&warmup, b"warmup").expect("warmup write");
1395 let waited = wait_for(watcher, REAL_BACKEND_DELIVERY, |ops| !ops.is_empty());
1396 let _ = fs::remove_file(&warmup);
1397 match waited {
1398 Waited::Delivered(_) => true,
1399 Waited::Silent => {
1400 if std::env::var_os("FDU_TEST_ALLOW_NO_NATIVE_WATCH").as_deref()
1401 == Some(std::ffi::OsStr::new("1"))
1402 {
1403 eprintln!(
1404 "skipped by FDU_TEST_ALLOW_NO_NATIVE_WATCH=1: the host event service \
1405 delivered no events to this stream"
1406 );
1407 return false;
1408 }
1409 panic!(
1410 "native watch precondition failed: the host delivered no events in \
1411 {REAL_BACKEND_DELIVERY:?}; run on a host with event delivery, or \
1412 explicitly opt out with FDU_TEST_ALLOW_NO_NATIVE_WATCH=1"
1413 );
1414 }
1415 }
1416 }
1417
1418 #[test]
1419 fn created_files_arrive_as_verified_upserts() {
1420 let _serialized = real_watcher_guard();
1421 let dir = tempfile::tempdir().expect("tempdir");
1422 let watcher = Watcher::new(dir.path(), WatchConfig::default()).expect("watcher");
1423 if !establish_watch(&watcher, dir.path()) {
1424 return;
1425 }
1426
1427 fs::write(dir.path().join("hello.txt"), b"hello world").expect("write");
1428
1429 let ops = wait_established(&watcher, REAL_BACKEND_DELIVERY, |ops| {
1430 ops.iter().any(|op| op.path() == Path::new("hello.txt"))
1431 });
1432
1433 let found = ops
1434 .iter()
1435 .find(|op| op.path() == Path::new("hello.txt"))
1436 .expect("an op for the new file");
1437 match found {
1438 Op::Upsert { attrs, kind, .. } => {
1439 assert!(!kind.is_dir());
1440 assert_eq!(attrs.size, 11);
1443 assert!(attrs.mtime_ns > 0);
1444 }
1445 other => panic!("expected an upsert, got {other:?}"),
1446 }
1447 }
1448
1449 #[test]
1450 fn deleted_files_arrive_as_removes() {
1451 let _serialized = real_watcher_guard();
1452 let dir = tempfile::tempdir().expect("tempdir");
1453 let path = dir.path().join("doomed.txt");
1454 fs::write(&path, b"x").expect("write");
1455
1456 let watcher = Watcher::new(dir.path(), WatchConfig::default()).expect("watcher");
1457 if !establish_watch(&watcher, dir.path()) {
1458 return;
1459 }
1460 fs::remove_file(&path).expect("remove");
1461
1462 let ops = wait_established(&watcher, REAL_BACKEND_DELIVERY, |ops| {
1463 ops.iter()
1464 .any(|op| matches!(op, Op::Remove { path } if path == Path::new("doomed.txt")))
1465 });
1466
1467 assert!(
1468 ops.iter()
1469 .any(|op| matches!(op, Op::Remove { path } if path == Path::new("doomed.txt"))),
1470 "expected a remove, saw {ops:?}"
1471 );
1472 }
1473
1474 #[test]
1493 fn watching_a_missing_path_is_an_error() {
1494 let dir = tempfile::tempdir().expect("tempdir");
1495 let missing = dir.path().join("not-there");
1496 assert!(Watcher::new(&missing, WatchConfig::default()).is_err());
1497 }
1498
1499 #[test]
1500 fn paths_outside_the_root_are_ignored() {
1501 let root = Path::new("/a/b");
1502 assert_eq!(relative_to(root, Path::new("/a/b/c/d")), Some(PathBuf::from("c/d")));
1503 assert_eq!(relative_to(root, Path::new("/elsewhere")), None);
1504 }
1505
1506 #[test]
1507 fn verification_errors_distinguish_absence_from_an_invalid_ancestor() {
1508 let path = PathBuf::from("parent/known.txt");
1509 let missing = op_for_stat_error(
1510 path.clone(),
1511 &std::io::Error::new(std::io::ErrorKind::NotFound, "gone"),
1512 );
1513 assert!(matches!(missing, Op::Remove { path: removed } if removed == path));
1514
1515 let not_a_directory = op_for_stat_error(
1516 path.clone(),
1517 &std::io::Error::new(std::io::ErrorKind::NotADirectory, "ancestor is a file"),
1518 );
1519 assert!(matches!(
1520 not_a_directory,
1521 Op::InvalidateSubtree {
1522 path: invalidated,
1523 reason: InvalidateReason::VerificationFailed,
1524 } if invalidated == Path::new("parent")
1525 ));
1526
1527 let denied = op_for_stat_error(
1528 path.clone(),
1529 &std::io::Error::new(std::io::ErrorKind::PermissionDenied, "denied"),
1530 );
1531 assert!(matches!(
1532 denied,
1533 Op::InvalidateSubtree {
1534 path: invalidated,
1535 reason: InvalidateReason::VerificationFailed,
1536 } if invalidated == path
1537 ));
1538 }
1539
1540 #[test]
1541 fn unknown_watch_ancestry_reconciles_from_the_nearest_known_directory() {
1542 let dir = tempfile::tempdir().expect("tempdir");
1543 let (index, _) =
1544 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
1545 let handle = crate::IndexHandle::new(index);
1546 let nested = dir.path().join("new/deep");
1547 fs::create_dir_all(&nested).expect("nested directories");
1548 fs::write(nested.join("file.txt"), b"verified").expect("nested file");
1549
1550 let report = apply_observation(
1551 &handle,
1552 &Observation::new(vec![Op::Upsert {
1553 path: PathBuf::from("new/deep/file.txt"),
1554 kind: crate::EntryKind::File,
1555 attrs: crate::Attrs::default(),
1556 }]),
1557 &crate::ScanConfig::default(),
1558 &mut |_| {},
1559 )
1560 .expect("unknown ancestry schedules reconciliation");
1561
1562 assert_eq!(report.apply.invalidated, 1);
1563 assert!(report.reconciliation.is_complete());
1564 assert_eq!(
1565 handle.kind(Path::new("new")).expect("new directory"),
1566 Some(crate::EntryKind::Dir)
1567 );
1568 assert_eq!(
1569 handle.kind(Path::new("new/deep/file.txt")).expect("nested file"),
1570 Some(crate::EntryKind::File)
1571 );
1572 assert_ne!(
1573 handle.attrs(Path::new("new")).expect("verified parent attrs"),
1574 Some(crate::Attrs::default())
1575 );
1576 }
1577
1578 #[test]
1579 fn create_intent_survives_coalescing_but_metadata_only_does_not_relist() {
1580 let root = Path::new("/watch-root");
1581 let path = root.join("directory");
1582 let mut pending = BTreeMap::new();
1583
1584 record(
1585 root,
1586 ¬ify::Event::new(EventKind::Create(CreateKind::Folder)).add_path(path.clone()),
1587 &mut pending,
1588 16,
1589 RenameReporting::EachSide,
1590 );
1591 record(
1592 root,
1593 ¬ify::Event::new(EventKind::Modify(ModifyKind::Metadata(MetadataKind::Any)))
1594 .add_path(path),
1595 &mut pending,
1596 16,
1597 RenameReporting::EachSide,
1598 );
1599 assert_eq!(
1600 pending.get(Path::new("directory")),
1601 Some(&Pending::Verify { relist_if_dir: true, renamed: false })
1602 );
1603
1604 let mut metadata_only = BTreeMap::new();
1605 record(
1606 root,
1607 ¬ify::Event::new(EventKind::Modify(ModifyKind::Metadata(MetadataKind::Any)))
1608 .add_path(root.join("existing")),
1609 &mut metadata_only,
1610 16,
1611 RenameReporting::EachSide,
1612 );
1613 assert_eq!(
1614 metadata_only.get(Path::new("existing")),
1615 Some(&Pending::Verify { relist_if_dir: false, renamed: false })
1616 );
1617 }
1618
1619 fn recorded(
1621 root: &Path,
1622 renames: RenameReporting,
1623 events: &[notify::Event],
1624 ) -> BTreeMap<PathBuf, Pending> {
1625 let mut pending = BTreeMap::new();
1626 for event in events {
1627 record(root, event, &mut pending, 16, renames);
1628 }
1629 pending
1630 }
1631
1632 fn rename_event(mode: RenameMode, paths: &[PathBuf]) -> notify::Event {
1633 let mut event = notify::Event::new(EventKind::Modify(ModifyKind::Name(mode)));
1634 event.paths = paths.to_vec();
1635 event
1636 }
1637
1638 fn fsevents_rename(path: PathBuf) -> notify::Event {
1640 rename_event(RenameMode::Any, &[path])
1641 }
1642
1643 const RENAMED: Pending = Pending::Verify { relist_if_dir: false, renamed: true };
1644
1645 #[test]
1649 fn each_rename_shape_verifies_its_named_paths_without_escalating() {
1650 let root = Path::new("/watch-root");
1651 let (old, new) = (root.join("dir/old"), root.join("other/new"));
1652 for (backend, events) in [
1653 ("FSEvents", vec![fsevents_rename(old.clone()), fsevents_rename(new.clone())]),
1654 (
1655 "inotify",
1656 vec![
1657 rename_event(RenameMode::From, std::slice::from_ref(&old)),
1658 rename_event(RenameMode::To, std::slice::from_ref(&new)),
1659 rename_event(RenameMode::Both, &[old.clone(), new.clone()]),
1660 ],
1661 ),
1662 (
1663 "Windows",
1664 vec![
1665 rename_event(RenameMode::From, std::slice::from_ref(&old)),
1666 rename_event(RenameMode::To, std::slice::from_ref(&new)),
1667 ],
1668 ),
1669 ] {
1670 let pending = recorded(root, RenameReporting::EachSide, &events);
1671 assert_eq!(
1672 pending,
1673 BTreeMap::from([
1674 (PathBuf::from("dir/old"), RENAMED),
1675 (PathBuf::from("other/new"), RENAMED),
1676 ]),
1677 "{backend}"
1678 );
1679 }
1680
1681 let outside = Path::new("/elsewhere/file");
1682 for (direction, event) in [
1683 ("move in", rename_event(RenameMode::Both, &[outside.to_path_buf(), new.clone()])),
1684 ("move out", rename_event(RenameMode::Both, &[old.clone(), outside.to_path_buf()])),
1685 ] {
1686 let pending = recorded(root, RenameReporting::EachSide, &[event]);
1687 assert_eq!(pending.len(), 1, "{direction}: {pending:?}");
1688 assert!(!pending.contains_key(Path::new("")), "{direction}: {pending:?}");
1689 }
1690 }
1691
1692 #[test]
1696 fn a_sticky_rename_flag_merges_with_the_same_paths_other_events() {
1697 let root = Path::new("/watch-root");
1698 let path = root.join("state.json");
1699 let pending = recorded(
1700 root,
1701 RenameReporting::EachSide,
1702 &[
1703 notify::Event::new(EventKind::Create(CreateKind::File)).add_path(path.clone()),
1704 fsevents_rename(path.clone()),
1705 notify::Event::new(EventKind::Modify(ModifyKind::Data(
1706 notify::event::DataChange::Content,
1707 )))
1708 .add_path(path),
1709 ],
1710 );
1711
1712 assert_eq!(
1713 pending,
1714 BTreeMap::from([(
1715 PathBuf::from("state.json"),
1716 Pending::Verify { relist_if_dir: true, renamed: true }
1717 )])
1718 );
1719 }
1720
1721 #[test]
1725 fn unboundable_renames_and_ambiguous_rescans_escalate_the_root() {
1726 let root = Path::new("/watch-root");
1727 let escalated = Some(&Pending::Escalate(InvalidateReason::UnpairedRename));
1728
1729 let kqueue =
1730 recorded(root, RenameReporting::OldSideOnly, &[fsevents_rename(root.join("old"))]);
1731 assert_eq!(kqueue.get(Path::new("")), escalated);
1732
1733 let root_moved = recorded(
1734 root,
1735 RenameReporting::EachSide,
1736 &[rename_event(RenameMode::From, &[root.to_path_buf()])],
1737 );
1738 assert_eq!(root_moved.get(Path::new("")), escalated);
1739 assert_eq!(root_moved.len(), 1, "the root is escalated, never verified: {root_moved:?}");
1740
1741 for paths in [vec![], vec![PathBuf::from("/elsewhere/file")]] {
1742 let unplaced =
1743 recorded(root, RenameReporting::EachSide, &[rename_event(RenameMode::Any, &paths)]);
1744 assert_eq!(unplaced.get(Path::new("")), escalated, "{paths:?}");
1745 }
1746
1747 let rescan = recorded(
1748 root,
1749 RenameReporting::EachSide,
1750 &[notify::Event::new(EventKind::Any)
1751 .add_path(root.join("a"))
1752 .add_path(root.join("b"))
1753 .set_flag(Flag::Rescan)],
1754 );
1755 assert_eq!(
1756 rescan.get(Path::new("")),
1757 Some(&Pending::Escalate(InvalidateReason::WatchOverflow))
1758 );
1759 }
1760
1761 #[test]
1762 fn pending_path_overload_collapses_to_one_root_invalidation() {
1763 let root = Path::new("/watch-root");
1764 let mut pending = BTreeMap::new();
1765 for name in ["one", "two", "three"] {
1766 record(
1767 root,
1768 ¬ify::Event::new(EventKind::Any).add_path(root.join(name)),
1769 &mut pending,
1770 2,
1771 RenameReporting::EachSide,
1772 );
1773 }
1774
1775 assert_eq!(pending.len(), 1);
1776 assert_eq!(
1777 pending.get(Path::new("")),
1778 Some(&Pending::Escalate(InvalidateReason::WatchOverflow))
1779 );
1780 }
1781
1782 #[test]
1783 fn continuous_churn_has_a_deterministic_max_hold_ceiling() {
1784 let past = Instant::now()
1785 .checked_sub(Duration::from_secs(2))
1786 .expect("representable earlier instant");
1787 assert!(max_hold_elapsed(Some(past), Duration::from_secs(1)));
1788 assert!(!max_hold_elapsed(None, Duration::from_secs(1)));
1789 }
1790
1791 #[test]
1792 fn backend_enqueue_is_nonblocking_and_marks_overflow() {
1793 let (sender, receiver) = sync_channel(1);
1794 let overflowed = AtomicBool::new(false);
1795 enqueue_raw(&sender, &overflowed, Ok(notify::Event::new(EventKind::Any)));
1796 enqueue_raw(&sender, &overflowed, Ok(notify::Event::new(EventKind::Any)));
1797
1798 assert!(overflowed.load(Ordering::Acquire));
1799 assert!(matches!(receiver.try_recv(), Ok(RawMessage::Event(Ok(_)))));
1800 }
1801
1802 #[test]
1803 fn full_intent_queue_retains_a_sticky_root_invalidation() {
1804 let (sender, receiver) = sync_channel(1);
1805 sender.try_send(CoalescedIntent::default()).expect("fill output");
1806 let mut pending = BTreeMap::from([(
1807 PathBuf::from("lost.txt"),
1808 Pending::Verify { relist_if_dir: false, renamed: false },
1809 )]);
1810 let mut sticky_overflow = false;
1811
1812 try_deliver_pending(&mut pending, &sender, &mut sticky_overflow).expect("connected");
1813 assert!(pending.is_empty());
1814 assert!(sticky_overflow);
1815
1816 receiver.try_recv().expect("make output capacity");
1817 assert!(try_deliver_overflow(&sender).expect("connected"));
1818 let intent = receiver.try_recv().expect("sticky overflow intent");
1819 let observation = verify_intent(
1820 Path::new("/unused"),
1821 WatchConfig::default(),
1822 &intent,
1823 &ScanConfig::default(),
1824 );
1825 assert!(matches!(
1826 &observation.ops[0].op,
1827 Op::InvalidateSubtree {
1828 path,
1829 reason: InvalidateReason::WatchOverflow,
1830 } if path.as_os_str().is_empty()
1831 ));
1832 }
1833
1834 #[test]
1835 fn cancellation_wakes_and_joins_with_a_full_intent_queue() {
1836 let dir = tempfile::tempdir().expect("tempdir");
1837 let root = dir.path().canonicalize().expect("canonical root");
1838 let config = WatchConfig {
1839 settle: Duration::from_secs(30),
1840 max_hold: Duration::from_secs(30),
1841 event_capacity: 1,
1842 batch_path_capacity: 1,
1843 intent_capacity: 1,
1844 ..WatchConfig::default()
1845 };
1846 let (control, raw) = sync_channel(1);
1847 let (output, intents) = sync_channel(1);
1848 output.try_send(CoalescedIntent::default()).expect("fill intent queue");
1849 let cancelled = Arc::new(AtomicBool::new(false));
1850 let worker_cancelled = Arc::clone(&cancelled);
1851 let status = Arc::new(AtomicU8::new(WORKER_RUNNING));
1852 let tracked_status = Arc::clone(&status);
1853 let overflowed = Arc::new(AtomicBool::new(false));
1854 let worker_overflowed = Arc::clone(&overflowed);
1855 let worker_root = root.clone();
1856 let worker = std::thread::spawn(move || {
1857 run_tracked_worker(&tracked_status, || {
1858 run_worker(
1859 &worker_root,
1860 config,
1861 RenameReporting::EachSide,
1862 &raw,
1863 &output,
1864 &worker_overflowed,
1865 &worker_cancelled,
1866 );
1867 });
1868 });
1869 let watcher = Watcher {
1870 root,
1871 config,
1872 inner: None,
1873 intents,
1874 control: Some(control),
1875 cancelled,
1876 worker_status: status,
1877 worker: Some(worker),
1878 };
1879 let (done_tx, done_rx) = sync_channel(1);
1880
1881 std::thread::spawn(move || {
1882 drop(watcher);
1883 done_tx.send(()).expect("report drop");
1884 });
1885
1886 done_rx
1887 .recv_timeout(Duration::from_secs(5))
1888 .expect("watcher drop must wake and join promptly");
1889 }
1890
1891 #[test]
1892 fn coalescing_defers_filesystem_verification_to_the_consumer() {
1893 let dir = tempfile::tempdir().expect("tempdir");
1894 let relative = PathBuf::from("appeared.txt");
1895 let intent = CoalescedIntent {
1896 pending: BTreeMap::from([(
1897 relative.clone(),
1898 Pending::Verify { relist_if_dir: false, renamed: false },
1899 )]),
1900 };
1901
1902 fs::write(dir.path().join(&relative), b"current").expect("create after coalescing");
1903 let observation =
1904 verify_intent(dir.path(), WatchConfig::default(), &intent, &ScanConfig::default());
1905
1906 assert!(matches!(
1907 &observation.ops[0].op,
1908 Op::Upsert { path, attrs, .. } if path == &relative && attrs.size == 7
1909 ));
1910 }
1911
1912 #[test]
1913 fn control_verification_emits_exact_source_with_the_entry_fact() {
1914 let dir = tempfile::tempdir().expect("tempdir");
1915 let relative = PathBuf::from(".gitignore");
1916 fs::write(dir.path().join(&relative), b"*.log\n").expect("write control");
1917 let intent = CoalescedIntent {
1918 pending: BTreeMap::from([(
1919 relative.clone(),
1920 Pending::Verify { relist_if_dir: false, renamed: false },
1921 )]),
1922 };
1923
1924 let config = ScanConfig { read_controls: true, ..ScanConfig::default() };
1925 let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
1926
1927 assert!(matches!(
1928 &observation.ops[0].op,
1929 Op::Upsert { path, kind: crate::EntryKind::File, .. } if path == &relative
1930 ));
1931 assert!(matches!(
1932 &observation.ops[1].op,
1933 Op::ControlUpsert { path, source } if path == &relative && source == b"*.log\n"
1934 ));
1935 }
1936
1937 #[test]
1938 fn only_population_control_event_reconciles_previously_absent_file() {
1939 let dir = tempfile::tempdir().expect("tree");
1940 fs::write(dir.path().join(".gitignore"), b"# no ignored files\n").expect("control");
1941 fs::write(dir.path().join("debug.log"), b"debug").expect("file");
1942 let scan =
1943 ScanConfig { population: crate::query::IgnoredEntries::Only, ..ScanConfig::default() };
1944 let (index, cold) = crate::scan::scan_into_index(dir.path(), &scan).expect("cold");
1945 assert!(cold.is_complete());
1946 assert!(index.lookup(Path::new("debug.log")).is_none());
1947 let handle = crate::IndexHandle::new(index);
1948 fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("rule edit");
1949 let intent = CoalescedIntent {
1950 pending: BTreeMap::from([(
1951 PathBuf::from(".gitignore"),
1952 Pending::Verify { relist_if_dir: false, renamed: false },
1953 )]),
1954 };
1955 let report =
1956 apply_intent(&handle, dir.path(), WatchConfig::default(), &intent, &scan, &mut |_| {})
1957 .expect("watch apply");
1958 assert!(report.reconciliation.is_complete());
1959 assert!(handle.kind(Path::new("debug.log")).expect("lookup").is_some());
1960 }
1961
1962 #[test]
1970 fn verification_observes_no_control_state_under_a_controls_off_policy() {
1971 let dir = tempfile::tempdir().expect("tempdir");
1972 fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("write control");
1973 fs::write(dir.path().join("debug.log"), b"x").expect("write file");
1974 let config = ScanConfig { read_controls: false, ..ScanConfig::default() };
1975 let (mut index, _) = crate::scan::scan_into_index(dir.path(), &config).expect("scan");
1976 assert!(index.control_table().is_empty());
1977 assert!(matches!(index.controls(), Err(crate::Error::ControlStateNotObserved)));
1978 let control = PathBuf::from(".gitignore");
1979 let intent = CoalescedIntent {
1980 pending: BTreeMap::from([(
1981 control.clone(),
1982 Pending::Verify { relist_if_dir: false, renamed: false },
1983 )]),
1984 };
1985
1986 let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
1987 let reverified = reverify_observation(
1988 dir.path(),
1989 &Observation::new(vec![Op::Remove { path: control }]),
1990 &config,
1991 )
1992 .expect("reverify");
1993
1994 for verified in [&observation, &reverified] {
1995 assert!(
1996 !verified.ops.iter().any(|observed| matches!(
1997 observed.op,
1998 Op::ControlUpsert { .. } | Op::ControlRemove { .. }
1999 )),
2000 "a controls-off policy observed control state: {:?}",
2001 verified.ops
2002 );
2003 }
2004 index.apply(&observation).expect("apply the verified observation");
2005 assert!(index.control_table().is_empty());
2006 assert_eq!(index.scope(), config.scope());
2007 }
2008
2009 #[cfg(unix)]
2010 #[test]
2011 fn applying_verification_uses_the_index_admission_scope() {
2012 use std::os::unix::net::UnixListener;
2013
2014 let dir = tempfile::tempdir().expect("tempdir");
2015 fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("write control");
2016 fs::write(dir.path().join(".secret"), b"hidden").expect("write hidden");
2017 let _listener = UnixListener::bind(dir.path().join("service.sock")).expect("bind socket");
2018 let intent = CoalescedIntent {
2019 pending: [".gitignore", ".secret", "service.sock"]
2020 .into_iter()
2021 .map(|path| {
2022 (PathBuf::from(path), Pending::Verify { relist_if_dir: false, renamed: false })
2023 })
2024 .collect(),
2025 };
2026 let config = ScanConfig {
2027 hidden: Some(Arc::new(crate::HiddenPolicy::prune_hidden::<[&str; 0], &str>([]))),
2028 exclude_special: true,
2029 read_controls: true,
2030 ..ScanConfig::default()
2031 };
2032
2033 let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
2034
2035 assert!(observation.ops.iter().any(|observed| matches!(
2036 &observed.op,
2037 Op::ControlUpsert { path, source }
2038 if path == Path::new(".gitignore") && source == b"*.log\n"
2039 )));
2040 for path in [".gitignore", ".secret", "service.sock"] {
2041 assert!(observation.ops.iter().any(|observed| matches!(
2042 &observed.op,
2043 Op::Remove { path: removed } if removed == Path::new(path)
2044 )));
2045 assert!(!observation.ops.iter().any(|observed| matches!(
2046 &observed.op,
2047 Op::Upsert { path: retained, .. } if retained == Path::new(path)
2048 )));
2049 }
2050 }
2051
2052 #[test]
2053 fn timeout_stop_and_worker_panic_are_distinct() {
2054 let dir = tempfile::tempdir().expect("tempdir");
2055 let (live_sender, live) =
2056 queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2057 assert!(live.next_observation(Duration::ZERO).expect("timeout").is_none());
2058 drop(live_sender);
2059
2060 let (stopped_sender, stopped) =
2061 queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2062 stopped.worker_status.store(WORKER_STOPPED, Ordering::Release);
2063 drop(stopped_sender);
2064 assert!(matches!(stopped.next_observation(Duration::ZERO), Err(Error::WatchStopped)));
2065
2066 let (panicked_sender, panicked) =
2067 queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2068 panicked.worker_status.store(WORKER_PANICKED, Ordering::Release);
2069 drop(panicked_sender);
2070 assert!(matches!(
2071 panicked.next_observation(Duration::ZERO),
2072 Err(Error::WatchWorkerPanicked)
2073 ));
2074 }
2075
2076 #[test]
2077 fn tracked_worker_records_a_panic() {
2078 let status = AtomicU8::new(WORKER_RUNNING);
2079 run_tracked_worker(&status, || panic!("injected worker panic"));
2080 assert_eq!(status.load(Ordering::Acquire), WORKER_PANICKED);
2081 }
2082
2083 #[test]
2084 fn observation_driver_closes_the_invalidation_loop() {
2085 let dir = tempfile::tempdir().expect("tempdir");
2086 let (index, _) =
2087 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2088 let handle = crate::IndexHandle::new(index);
2089 fs::write(dir.path().join("raced.txt"), b"raced").expect("write");
2090 let observation = Observation::new(vec![Op::InvalidateSubtree {
2091 path: PathBuf::new(),
2092 reason: InvalidateReason::WatchSetupRace,
2093 }]);
2094
2095 apply_observation(&handle, &observation, &crate::ScanConfig::default(), &mut |_| {})
2096 .expect("apply and reconcile");
2097
2098 assert!(handle.kind(Path::new("raced.txt")).expect("query").is_some());
2099 }
2100
2101 #[test]
2102 fn applying_driver_reverifies_a_queued_sample_after_reconciliation() {
2103 let dir = tempfile::tempdir().expect("tempdir");
2104 let path = dir.path().join("sample.txt");
2105 fs::write(&path, b"old").expect("write old sample");
2106 let (index, _) =
2107 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2108 let old_attrs = *index.attrs(Path::new("sample.txt")).expect("sample attributes");
2109 let delayed = Observation::new(vec![Op::Upsert {
2110 path: PathBuf::from("sample.txt"),
2111 kind: crate::EntryKind::File,
2112 attrs: old_attrs,
2113 }]);
2114 let handle = crate::IndexHandle::new(index);
2115
2116 fs::write(&path, b"new contents").expect("write current sample");
2117 crate::scan::reconcile_handle(&handle, &crate::ScanConfig::default(), &mut |_| {})
2118 .expect("reconcile newer sample");
2119 let current_size = fs::metadata(&path).expect("sample metadata").len();
2120
2121 apply_observation(&handle, &delayed, &crate::ScanConfig::default(), &mut |_| {})
2122 .expect("apply delayed watch sample");
2123
2124 assert_eq!(
2125 handle.attrs(Path::new("sample.txt")).expect("query").expect("sample remains").size,
2126 current_size
2127 );
2128 }
2129
2130 #[test]
2131 fn blocked_verifier_holds_no_index_lock_and_commits_only_at_current_clock() {
2132 let dir = tempfile::tempdir().expect("tempdir");
2133 let (index, _) =
2134 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2135 let handle = crate::IndexHandle::new(index);
2136 let queued = Observation::new(vec![Op::Upsert {
2137 path: PathBuf::from("queued.txt"),
2138 kind: crate::EntryKind::File,
2139 attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2140 }]);
2141 let applying = handle.clone();
2142 let (entered_tx, entered_rx) = sync_channel(1);
2143 let (release_tx, release_rx) = sync_channel(1);
2144 let (done_tx, done_rx) = sync_channel(1);
2145 let apply_thread = std::thread::spawn(move || {
2146 let mut first = true;
2147 let mut verifier = |_: &Path, observation: &Observation| {
2148 if first {
2149 first = false;
2150 entered_tx.send(()).expect("signal blocked verifier");
2151 release_rx.recv().expect("release blocked verifier");
2152 }
2153 Ok(observation.clone())
2154 };
2155 let result = apply_reverified_with(
2156 &applying,
2157 &queued,
2158 &crate::ScanConfig::default(),
2159 &mut verifier,
2160 );
2161 done_tx.send(result).expect("report apply result");
2162 });
2163
2164 entered_rx
2165 .recv_timeout(Duration::from_secs(5))
2166 .expect("verifier must reach the injected block");
2167 let progressing = handle.clone();
2168 let (progress_tx, progress_rx) = sync_channel(1);
2169 let progress_thread = std::thread::spawn(move || {
2170 let total = progressing.total().expect("reader progresses");
2171 let write = progressing.apply(&Observation::new(vec![Op::Upsert {
2172 path: PathBuf::from("competitor.txt"),
2173 kind: crate::EntryKind::File,
2174 attrs: crate::Attrs { size: 3, allocated: 3, ..crate::Attrs::default() },
2175 }]));
2176 progress_tx.send((total, write)).expect("report progress");
2177 });
2178 let (_, competing_write) = progress_rx
2179 .recv_timeout(Duration::from_secs(5))
2180 .expect("reader and writer must progress while verification is blocked");
2181 competing_write.expect("competing write");
2182 release_tx.send(()).expect("release verifier");
2183
2184 let outcome = done_rx
2185 .recv_timeout(Duration::from_secs(5))
2186 .expect("applying driver completes")
2187 .expect("applying driver succeeds");
2188 apply_thread.join().expect("apply thread");
2189 progress_thread.join().expect("progress thread");
2190 assert_eq!(outcome.inserted, 1);
2191 assert!(handle.kind(Path::new("competitor.txt")).expect("query").is_some());
2192 assert!(handle.kind(Path::new("queued.txt")).expect("query").is_some());
2193 assert_eq!(handle.clock().expect("clock"), crate::Clock(2));
2194 }
2195
2196 #[test]
2197 fn exhausted_watch_contention_stays_unfresh_until_reconciliation() {
2198 let dir = tempfile::tempdir().expect("tempdir");
2199 let (index, _) =
2200 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2201 let handle = crate::IndexHandle::new(index);
2202 let queued = Observation::new(vec![Op::Upsert {
2203 path: PathBuf::from("never-committed.txt"),
2204 kind: crate::EntryKind::File,
2205 attrs: crate::Attrs { size: 1, allocated: 1, ..crate::Attrs::default() },
2206 }]);
2207 let mut attempts = 0_usize;
2208 let mut verifier = |_: &Path, observation: &Observation| {
2209 attempts += 1;
2210 handle
2211 .apply(&Observation::new(vec![Op::Upsert {
2212 path: PathBuf::from(format!("competitor-{attempts}.txt")),
2213 kind: crate::EntryKind::File,
2214 attrs: crate::Attrs {
2215 size: attempts as u64,
2216 allocated: attempts as u64,
2217 ..crate::Attrs::default()
2218 },
2219 }]))
2220 .expect("force a clock conflict");
2221 Ok(observation.clone())
2222 };
2223
2224 let outcome =
2225 apply_reverified_with(&handle, &queued, &crate::ScanConfig::default(), &mut verifier)
2226 .expect("contention escalates");
2227
2228 assert_eq!(attempts, MAX_OPTIMISTIC_APPLY_ATTEMPTS);
2229 assert_eq!(outcome.invalidated, 1);
2230 assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Stale);
2231 assert!(handle.kind(Path::new("never-committed.txt")).expect("query").is_none());
2232 let pending = handle.take_pending_invalidations().expect("pending invalidation");
2233 assert_eq!(pending, vec![(PathBuf::new(), InvalidateReason::WatchContention)]);
2234 handle.restore_pending_invalidations(pending).expect("restore invalidation");
2235
2236 crate::scan::reconcile_pending_handle(&handle, &crate::ScanConfig::default(), &mut |_| {})
2237 .expect("reconcile contention");
2238 assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Fresh);
2239 }
2240
2241 #[test]
2242 fn verifier_error_mutates_no_shared_state() {
2243 let dir = tempfile::tempdir().expect("tempdir");
2244 let (index, _) =
2245 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2246 let handle = crate::IndexHandle::new(index);
2247 let before_clock = handle.clock().expect("clock");
2248 let before_total = handle.total().expect("total");
2249 let mut verifier = |_: &Path, _: &Observation| {
2250 Err(Error::io(
2251 PathBuf::from("blocked"),
2252 std::io::Error::new(std::io::ErrorKind::PermissionDenied, "injected"),
2253 ))
2254 };
2255
2256 let error = apply_reverified_with(
2257 &handle,
2258 &Observation::default(),
2259 &crate::ScanConfig::default(),
2260 &mut verifier,
2261 )
2262 .expect_err("verification error");
2263
2264 assert!(matches!(error, Error::Io { .. }));
2265 assert_eq!(handle.clock().expect("clock"), before_clock);
2266 assert_eq!(handle.total().expect("total"), before_total);
2267 assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Fresh);
2268 assert!(handle.take_pending_invalidations().expect("pending").is_empty());
2269 }
2270
2271 #[test]
2272 fn stable_watch_arbitration_verifies_exactly_once() {
2273 let dir = tempfile::tempdir().expect("tempdir");
2274 let (index, _) =
2275 crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
2276 let handle = crate::IndexHandle::new(index);
2277 let mut calls = 0_u8;
2278 let mut verifier = |_: &Path, observation: &Observation| {
2279 calls += 1;
2280 Ok(observation.clone())
2281 };
2282
2283 apply_reverified_with(
2284 &handle,
2285 &Observation::new(vec![Op::InvalidateSubtree {
2286 path: PathBuf::new(),
2287 reason: InvalidateReason::Requested,
2288 }]),
2289 &crate::ScanConfig::default(),
2290 &mut verifier,
2291 )
2292 .expect("stable apply");
2293
2294 assert_eq!(calls, 1);
2295 }
2296
2297 #[test]
2298 fn disappearing_watch_root_escalates_instead_of_removing_the_index_root() {
2299 let error = std::io::Error::new(std::io::ErrorKind::NotFound, "root disappeared");
2300
2301 assert!(matches!(
2302 op_for_stat_error(PathBuf::new(), &error),
2303 Op::InvalidateSubtree {
2304 path,
2305 reason: InvalidateReason::VerificationFailed,
2306 } if path.as_os_str().is_empty()
2307 ));
2308 }
2309
2310 #[test]
2311 fn observation_driver_rejects_scope_mismatch_before_apply() {
2312 let dir = tempfile::tempdir().expect("tempdir");
2313 let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2314 let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2315 let handle = crate::IndexHandle::new(index);
2316 let observation = Observation::new(vec![Op::Upsert {
2317 path: PathBuf::from("deep/nested.txt"),
2318 kind: crate::EntryKind::File,
2319 attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2320 }]);
2321
2322 let error =
2323 apply_observation(&handle, &observation, &crate::ScanConfig::default(), &mut |_| {})
2324 .expect_err("mismatched scope must fail");
2325
2326 assert!(matches!(error, Error::ScanScopeMismatch { .. }));
2327 assert!(handle.kind(Path::new("deep/nested.txt")).expect("query").is_none());
2328 }
2329
2330 #[test]
2331 fn observation_driver_rejects_restricted_scopes_until_events_are_filtered() {
2332 let dir = tempfile::tempdir().expect("tempdir");
2333 let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2334 let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2335 let handle = crate::IndexHandle::new(index);
2336 let observation = Observation::new(vec![Op::Upsert {
2337 path: PathBuf::from("deep/nested.txt"),
2338 kind: crate::EntryKind::File,
2339 attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
2340 }]);
2341
2342 let error = apply_observation(&handle, &observation, &shallow, &mut |_| {})
2343 .expect_err("unfiltered bounded watch scope must fail");
2344
2345 assert!(matches!(error, Error::UnsupportedScanConfig(_)));
2346 assert!(handle.kind(Path::new("deep/nested.txt")).expect("query").is_none());
2347 }
2348
2349 #[test]
2350 fn apply_next_rejects_restricted_scope_without_consuming_an_observation() {
2351 let dir = tempfile::tempdir().expect("tempdir");
2352 let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
2353 let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
2354 let handle = crate::IndexHandle::new(index);
2355 let (sender, watcher) =
2356 queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
2357 sender.try_send(CoalescedIntent::default()).expect("queue intent");
2358
2359 let error = watcher
2360 .apply_next(&handle, &shallow, Duration::ZERO, &mut |_| {})
2361 .expect_err("restricted scope must fail before receive");
2362
2363 assert!(matches!(error, Error::UnsupportedScanConfig(_)));
2364 assert!(watcher.next_observation(Duration::ZERO).expect("receive").is_some());
2365 }
2366
2367 #[test]
2368 fn apply_next_rejects_a_watcher_for_another_root_without_consuming() {
2369 let indexed = tempfile::tempdir().expect("indexed root");
2370 let watched_root_dir = tempfile::tempdir().expect("watched root");
2371 let (index, _) =
2372 crate::scan::scan_into_index(indexed.path(), &crate::ScanConfig::default())
2373 .expect("scan indexed root");
2374 let handle = crate::IndexHandle::new(index);
2375 let (sender, watcher) = queued_test_watcher(
2376 watched_root_dir.path().canonicalize().expect("canonical watched root"),
2377 );
2378 sender.try_send(CoalescedIntent::default()).expect("queue intent");
2379
2380 let error = watcher
2381 .apply_next(&handle, &crate::ScanConfig::default(), Duration::ZERO, &mut |_| {})
2382 .expect_err("mismatched root must fail");
2383
2384 assert!(matches!(error, Error::WatchRootMismatch { .. }));
2385 assert!(watcher.next_observation(Duration::ZERO).expect("receive").is_some());
2386 }
2387
2388 #[cfg(unix)]
2396 #[test]
2397 fn apply_next_walks_an_unreadable_gap_once_and_retains_its_cause() {
2398 use std::os::unix::fs::PermissionsExt;
2399
2400 fn walks_of(commits: &[Commit], path: &Path) -> usize {
2401 commits
2402 .iter()
2403 .flat_map(|commit| commit.state.iter())
2404 .filter(|transition| {
2405 matches!(
2406 transition,
2407 crate::StateTransition::Freshness { path: marked, current, .. }
2408 if marked == path && *current == crate::Freshness::Reconciling
2409 )
2410 })
2411 .count()
2412 }
2413
2414 if !crate::test_support::require_permission_bits() {
2415 return;
2416 }
2417 let dir = tempfile::tempdir().expect("tempdir");
2418 let root = dir.path().canonicalize().expect("canonical root");
2419 let blocked = root.join("blocked");
2420 fs::create_dir(&blocked).expect("blocked");
2421 fs::write(blocked.join("secret"), b"s").expect("fixture");
2422 let (index, _) =
2423 crate::scan::scan_into_index(&root, &crate::ScanConfig::default()).expect("scan");
2424 let handle = crate::IndexHandle::new(index);
2425 let (sender, watcher) = queued_test_watcher(root.clone());
2426 let config = crate::ScanConfig::default();
2427 let mut commits = Vec::new();
2428
2429 fs::set_permissions(&blocked, fs::Permissions::from_mode(0o000)).expect("deny reads");
2430 let mut gap = CoalescedIntent::default();
2431 gap.pending
2432 .insert(PathBuf::from("blocked"), Pending::Escalate(InvalidateReason::WatchOverflow));
2433 sender.try_send(gap).expect("queue the gap");
2434 let first = watcher
2435 .apply_next(&handle, &config, Duration::ZERO, &mut |commit| {
2436 commits.push(commit.clone());
2437 })
2438 .expect("apply the gap")
2439 .expect("an intent was queued");
2440 for name in ["live.txt", "marker.txt"] {
2441 fs::write(root.join(name), name).expect("unrelated mutation");
2442 let mut event = CoalescedIntent::default();
2443 event.pending.insert(
2444 PathBuf::from(name),
2445 Pending::Verify { relist_if_dir: false, renamed: false },
2446 );
2447 sender.try_send(event).expect("queue the event");
2448 watcher
2449 .apply_next(&handle, &config, Duration::ZERO, &mut |commit| {
2450 commits.push(commit.clone());
2451 })
2452 .expect("apply the event")
2453 .expect("an intent was queued");
2454 }
2455
2456 let pending = handle.take_pending_invalidations().expect("pending");
2457 let freshness = handle.freshness_at(Path::new("blocked")).expect("freshness");
2458 let issues = handle.issues().expect("issues");
2459 fs::set_permissions(&blocked, fs::Permissions::from_mode(0o700)).expect("restore reads");
2460 assert!(!first.reconciliation.is_complete(), "the gap must be unreadable");
2461
2462 assert_eq!(
2463 walks_of(&commits, Path::new("blocked")),
2464 1,
2465 "an unreadable subtree must not be re-walked per unrelated event"
2466 );
2467 assert!(pending.is_empty(), "the queue must settle: {pending:?}");
2468 assert_eq!(freshness, crate::Freshness::Partial);
2469 assert!(handle.kind(Path::new("marker.txt")).expect("lookup").is_some());
2470 assert!(
2471 issues.iter().any(|issue| issue.kind == crate::IssueKind::Permission
2472 && issue.path.as_deref() == Some(Path::new("blocked"))),
2473 "{issues:?}"
2474 );
2475 }
2476
2477 struct RenameFixture {
2479 _dir: tempfile::TempDir,
2480 root: PathBuf,
2481 handle: IndexHandle,
2482 }
2483
2484 impl RenameFixture {
2485 fn new(files: &[&str]) -> Self {
2487 let dir = tempfile::tempdir().expect("tempdir");
2488 let root = dir.path().canonicalize().expect("canonical root");
2489 for file in files {
2490 let path = root.join(file);
2491 fs::create_dir_all(path.parent().expect("parent")).expect("parents");
2492 fs::write(&path, file.as_bytes()).expect("fixture file");
2493 }
2494 let (index, _) =
2495 crate::scan::scan_into_index(&root, &ScanConfig::default()).expect("cold scan");
2496 Self { _dir: dir, root, handle: IndexHandle::new(index) }
2497 }
2498
2499 fn path(&self, rel: &str) -> PathBuf {
2500 self.root.join(rel)
2501 }
2502
2503 fn apply(&self, events: &[notify::Event]) -> Vec<(PathBuf, InvalidateReason)> {
2506 let intent = CoalescedIntent {
2507 pending: recorded(&self.root, RenameReporting::EachSide, events),
2508 };
2509 let mut commits = Vec::new();
2510 let report = apply_intent(
2511 &self.handle,
2512 &self.root,
2513 WatchConfig::default(),
2514 &intent,
2515 &ScanConfig::default(),
2516 &mut |commit| commits.push(commit.clone()),
2517 )
2518 .expect("apply the intent");
2519 assert!(report.reconciliation.is_complete(), "reconciliation must settle");
2520 commits
2521 .iter()
2522 .flat_map(|commit| commit.changes.iter())
2523 .filter_map(|change| match change {
2524 crate::EffectiveChange::Invalidated { path, reason } => {
2525 Some((path.clone(), *reason))
2526 }
2527 _ => None,
2528 })
2529 .collect()
2530 }
2531
2532 fn assert_converged(&self) {
2534 let (cold, _) =
2535 crate::scan::scan_into_index(&self.root, &ScanConfig::default()).expect("cold");
2536 let watched = self.handle.read_with(entries).expect("read the watched index");
2537 assert_eq!(watched, entries(&cold), "watched index diverged from a cold scan");
2538 assert_eq!(self.handle.freshness().expect("freshness"), crate::Freshness::Fresh);
2539 }
2540 }
2541
2542 fn entries(index: &crate::Index) -> BTreeMap<PathBuf, (crate::EntryKind, Option<u64>)> {
2546 fn walk(
2547 index: &crate::Index,
2548 dir: &Path,
2549 out: &mut BTreeMap<PathBuf, (crate::EntryKind, Option<u64>)>,
2550 ) {
2551 let Some(children) = index.children(dir) else {
2552 return;
2553 };
2554 let children: Vec<_> =
2555 children.map(|(name, id)| (dir.join(name), id)).collect::<Vec<_>>();
2556 for (path, id) in children {
2557 let kind = index.kind_of(id).expect("live child");
2558 let size = (!kind.is_dir()).then(|| index.attrs_of(id).expect("attrs").size);
2559 out.insert(path.clone(), (kind, size));
2560 walk(index, &path, out);
2561 }
2562 }
2563 let mut out = BTreeMap::new();
2564 walk(index, Path::new(""), &mut out);
2565 out
2566 }
2567
2568 fn created(path: PathBuf) -> notify::Event {
2569 notify::Event::new(EventKind::Create(CreateKind::Any)).add_path(path)
2570 }
2571
2572 fn modified(path: PathBuf) -> notify::Event {
2573 notify::Event::new(EventKind::Modify(ModifyKind::Data(notify::event::DataChange::Content)))
2574 .add_path(path)
2575 }
2576
2577 #[test]
2580 fn a_file_renamed_within_its_directory_needs_no_invalidation() {
2581 let tree = RenameFixture::new(&["state/config.json", "state/other.txt"]);
2582 fs::write(tree.path("state/config.json.tmp"), b"new configuration").expect("temp");
2583 fs::rename(tree.path("state/config.json.tmp"), tree.path("state/config.json"))
2584 .expect("rename over");
2585
2586 let invalidations = tree.apply(&[
2587 created(tree.path("state/config.json.tmp")),
2588 modified(tree.path("state/config.json.tmp")),
2589 fsevents_rename(tree.path("state/config.json.tmp")),
2590 fsevents_rename(tree.path("state/config.json")),
2591 ]);
2592
2593 assert_eq!(invalidations, vec![], "a file rename is settled by its own paths");
2594 tree.assert_converged();
2595 }
2596
2597 #[test]
2600 fn a_file_moved_across_directories_settles_one_side_per_event() {
2601 let tree = RenameFixture::new(&["from/moved.txt", "to/resident.txt"]);
2602 fs::rename(tree.path("from/moved.txt"), tree.path("to/moved.txt")).expect("move");
2603
2604 assert_eq!(tree.apply(&[fsevents_rename(tree.path("to/moved.txt"))]), vec![]);
2605 assert_eq!(tree.apply(&[fsevents_rename(tree.path("from/moved.txt"))]), vec![]);
2606 tree.assert_converged();
2607 }
2608
2609 #[test]
2612 fn a_renamed_directory_relists_only_its_own_subtree() {
2613 let tree = RenameFixture::new(&[
2614 "project/old/one.txt",
2615 "project/old/nested/two.txt",
2616 "project/untouched/three.txt",
2617 ]);
2618 fs::rename(tree.path("project/old"), tree.path("project/new")).expect("rename directory");
2619
2620 let invalidations = tree.apply(&[
2621 fsevents_rename(tree.path("project/old")),
2622 fsevents_rename(tree.path("project/new")),
2623 ]);
2624
2625 assert_eq!(
2626 invalidations,
2627 vec![(PathBuf::from("project/new"), InvalidateReason::UnpairedRename)]
2628 );
2629 tree.assert_converged();
2630 }
2631
2632 #[test]
2635 fn moves_across_the_root_boundary_settle_from_the_inside_side() {
2636 let tree = RenameFixture::new(&["resident.txt", "leaving/a.txt", "leaving/deep/b.txt"]);
2637 let outside = tempfile::tempdir().expect("outside the root");
2638 fs::create_dir_all(outside.path().join("arriving/deep")).expect("outside tree");
2639 fs::write(outside.path().join("arriving/deep/c.txt"), b"arrived").expect("outside file");
2640
2641 fs::rename(outside.path().join("arriving"), tree.path("arrived")).expect("move in");
2642 let invalidations = tree.apply(&[fsevents_rename(tree.path("arrived"))]);
2643 assert_eq!(
2644 invalidations,
2645 vec![(PathBuf::from("arrived"), InvalidateReason::UnpairedRename)]
2646 );
2647 tree.assert_converged();
2648
2649 fs::rename(tree.path("leaving"), outside.path().join("left")).expect("move out");
2650 assert_eq!(tree.apply(&[fsevents_rename(tree.path("leaving"))]), vec![]);
2651 tree.assert_converged();
2652 }
2653
2654 #[test]
2657 fn a_name_reused_after_a_rename_verifies_as_its_new_entry() {
2658 let tree = RenameFixture::new(&["log.txt", "cache/entry.bin"]);
2659 fs::rename(tree.path("log.txt"), tree.path("log.1.txt")).expect("rotate");
2660 fs::write(tree.path("log.txt"), b"fresh log, longer than the old one").expect("reuse");
2661 fs::rename(tree.path("cache"), tree.path("cache.old")).expect("retire directory");
2662 fs::create_dir(tree.path("cache")).expect("reuse directory name");
2663
2664 let invalidations = tree.apply(&[
2665 fsevents_rename(tree.path("log.txt")),
2666 created(tree.path("log.txt")),
2667 fsevents_rename(tree.path("log.1.txt")),
2668 fsevents_rename(tree.path("cache")),
2669 created(tree.path("cache")),
2670 fsevents_rename(tree.path("cache.old")),
2671 ]);
2672
2673 assert!(
2674 invalidations.iter().all(|(path, _)| !path.as_os_str().is_empty()),
2675 "no root reconcile: {invalidations:?}"
2676 );
2677 tree.assert_converged();
2678 assert_eq!(tree.handle.kind(Path::new("cache/entry.bin")).expect("lookup"), None);
2679 }
2680
2681 #[test]
2684 fn a_sticky_rename_flag_on_a_later_write_is_just_a_write() {
2685 let tree = RenameFixture::new(&["sessions/today.jsonl"]);
2686 fs::write(tree.path("sessions/today.jsonl"), b"appended record after the rename")
2687 .expect("write in place");
2688
2689 let invalidations = tree.apply(&[
2690 fsevents_rename(tree.path("sessions/today.jsonl")),
2691 modified(tree.path("sessions/today.jsonl")),
2692 ]);
2693
2694 assert_eq!(invalidations, vec![]);
2695 tree.assert_converged();
2696 }
2697
2698 fn case_insensitive(dir: &Path) -> bool {
2700 let probe = dir.join("CaseProbe");
2701 fs::write(&probe, b"probe").expect("case probe");
2702 let insensitive = dir.join("caseprobe").exists();
2703 fs::remove_file(probe).expect("remove case probe");
2704 insensitive
2705 }
2706
2707 #[test]
2711 fn a_case_only_rename_keeps_one_spelling() {
2712 let tree = RenameFixture::new(&["docs/Readme.md", "docs/Guide/intro.md"]);
2713 let insensitive = case_insensitive(&tree.root);
2714 fs::rename(tree.path("docs/Readme.md"), tree.path("docs/README.md")).expect("recase");
2715 fs::rename(tree.path("docs/Guide"), tree.path("docs/guide")).expect("recase directory");
2716
2717 let invalidations = tree.apply(&[
2718 fsevents_rename(tree.path("docs/Readme.md")),
2719 fsevents_rename(tree.path("docs/README.md")),
2720 fsevents_rename(tree.path("docs/Guide")),
2721 fsevents_rename(tree.path("docs/guide")),
2722 modified(tree.path("docs/Guide/intro.md")),
2724 ]);
2725
2726 tree.assert_converged();
2727 assert!(
2728 invalidations.iter().all(|(path, _)| !path.as_os_str().is_empty()),
2729 "no root reconcile: {invalidations:?}"
2730 );
2731 if insensitive {
2732 assert!(
2733 invalidations.contains(&(PathBuf::from("docs"), InvalidateReason::UnpairedRename)),
2734 "the stale spelling reconciles its parent: {invalidations:?}"
2735 );
2736 }
2737 }
2738
2739 #[test]
2742 fn a_listing_miss_rereads_the_parent_before_answering() {
2743 let dir = tempfile::tempdir().expect("tempdir");
2744 let root = dir.path();
2745 fs::write(root.join("first"), b"1").expect("first");
2746 let mut listings = ParentListings::default();
2747
2748 assert_eq!(listings.lists(root, Path::new("first")), Some(true));
2749 fs::write(root.join("created-after-the-listing"), b"2").expect("later entry");
2750 assert_eq!(listings.lists(root, Path::new("created-after-the-listing")), Some(true));
2751 assert_eq!(listings.lists(root, Path::new("never-created")), Some(false));
2752 assert_eq!(listings.lists(root, Path::new("missing-directory/child")), None);
2753 }
2754
2755 #[test]
2758 fn native_renames_converge_without_reconciling_the_root() {
2759 let _serialized = real_watcher_guard();
2760 let dir = tempfile::tempdir().expect("tempdir");
2761 let root = dir.path().canonicalize().expect("canonical root");
2762 for file in ["from/moved.txt", "tree/sub/leaf.txt", "to/resident.txt"] {
2763 fs::create_dir_all(root.join(file).parent().expect("parent")).expect("parents");
2764 fs::write(root.join(file), file.as_bytes()).expect("fixture file");
2765 }
2766 let watcher = Watcher::new(&root, WatchConfig::default()).expect("watcher");
2767 if !establish_watch(&watcher, &root) {
2768 return;
2769 }
2770 let config = ScanConfig::default();
2771 let (index, _) = crate::scan::scan_into_index(&root, &config).expect("scan");
2772 let handle = IndexHandle::new(index);
2773
2774 fs::rename(root.join("from/moved.txt"), root.join("to/moved.txt")).expect("move file");
2775 fs::rename(root.join("tree"), root.join("renamed-tree")).expect("rename directory");
2776
2777 let converged = || {
2778 handle.kind(Path::new("to/moved.txt")).expect("lookup").is_some()
2779 && handle.kind(Path::new("from/moved.txt")).expect("lookup").is_none()
2780 && handle.kind(Path::new("renamed-tree/sub/leaf.txt")).expect("lookup").is_some()
2781 && handle.kind(Path::new("tree")).expect("lookup").is_none()
2782 };
2783 let mut commits = Vec::new();
2784 let start = Instant::now();
2785 while !converged() && start.elapsed() < REAL_BACKEND_DELIVERY {
2786 watcher
2787 .apply_next(&handle, &config, Duration::from_millis(200), &mut |commit| {
2788 commits.push(commit.clone());
2789 })
2790 .expect("apply");
2791 }
2792
2793 assert!(converged(), "the renames never converged: {commits:?}");
2794 let root_reconciles = commits
2795 .iter()
2796 .flat_map(|commit| commit.changes.iter())
2797 .filter(|change| {
2798 matches!(change, crate::EffectiveChange::Invalidated { path, .. }
2799 if path.as_os_str().is_empty())
2800 })
2801 .count();
2802 assert_eq!(root_reconciles, 0, "a rename reconciled the whole root: {commits:?}");
2803 let (cold, _) = crate::scan::scan_into_index(&root, &config).expect("cold");
2804 assert_eq!(handle.read_with(entries).expect("read"), entries(&cold));
2805 }
2806
2807 #[test]
2808 fn zero_settle_is_rejected_before_starting_a_busy_worker() {
2809 let config = WatchConfig { settle: Duration::ZERO, ..WatchConfig::default() };
2810
2811 assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2812
2813 let config = WatchConfig { intent_capacity: 0, ..WatchConfig::default() };
2814 assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2815
2816 let config = WatchConfig {
2817 batch_path_capacity: MAX_BATCH_PATH_CAPACITY,
2818 intent_capacity: MAX_BUFFERED_INTENT_PATHS / MAX_BATCH_PATH_CAPACITY + 1,
2819 ..WatchConfig::default()
2820 };
2821 assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
2822 }
2823}