use std::collections::BTreeMap;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
use std::sync::mpsc::{Receiver, RecvTimeoutError, SyncSender, TrySendError, sync_channel};
use std::thread::JoinHandle;
use std::time::{Duration, Instant};
use notify::{EventKind, RecommendedWatcher, RecursiveMode, Watcher as NotifyWatcher};
#[cfg(test)]
mod scripted_events;
use crate::engine_contract::{
Commit, Error, InvalidateReason, Observation, ObservationOp, Op, Result,
};
use crate::scan;
use crate::{ApplyOutcome, IndexHandle, ScanConfig};
const MAX_OPTIMISTIC_APPLY_ATTEMPTS: usize = 3;
const WORKER_RUNNING: u8 = 0;
const WORKER_STOPPED: u8 = 1;
const WORKER_PANICKED: u8 = 2;
const MAX_EVENT_CAPACITY: usize = 64 * 1024;
const MAX_BATCH_PATH_CAPACITY: usize = 64 * 1024;
const MAX_BUFFERED_INTENT_PATHS: usize = 1024 * 1024;
#[derive(Clone, Copy, Debug)]
pub struct WatchConfig {
pub settle: Duration,
pub max_hold: Duration,
pub relist_new_dirs: bool,
pub event_capacity: usize,
pub batch_path_capacity: usize,
pub intent_capacity: usize,
}
impl Default for WatchConfig {
fn default() -> Self {
Self {
settle: Duration::from_millis(50),
max_hold: Duration::from_millis(1600),
relist_new_dirs: true,
event_capacity: 4096,
batch_path_capacity: 4096,
intent_capacity: 16,
}
}
}
impl WatchConfig {
fn validate(self) -> Result<()> {
let buffered_paths = self.batch_path_capacity.checked_mul(self.intent_capacity);
if self.settle.is_zero()
|| self.max_hold.is_zero()
|| self.max_hold < self.settle
|| self.event_capacity == 0
|| self.batch_path_capacity == 0
|| self.intent_capacity == 0
|| self.event_capacity > MAX_EVENT_CAPACITY
|| self.batch_path_capacity > MAX_BATCH_PATH_CAPACITY
|| buffered_paths.is_none_or(|paths| paths > MAX_BUFFERED_INTENT_PATHS)
{
return Err(Error::UnsupportedScanConfig(
"watch durations and capacities exceed the supported nonzero bounds, or max_hold is less than settle",
));
}
Ok(())
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
enum Pending {
Verify {
relist_if_dir: bool,
renamed: bool,
},
Escalate(InvalidateReason),
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
enum RenameReporting {
EachSide,
OldSideOnly,
}
impl RenameReporting {
fn of_recommended_backend() -> Self {
match <RecommendedWatcher as NotifyWatcher>::kind() {
notify::WatcherKind::Fsevent
| notify::WatcherKind::Inotify
| notify::WatcherKind::ReadDirectoryChangesWatcher => Self::EachSide,
_ => Self::OldSideOnly,
}
}
}
#[derive(Debug, Default)]
struct CoalescedIntent {
pending: BTreeMap<PathBuf, Pending>,
}
enum RawMessage {
Event(notify::Result<notify::Event>),
Flush(SyncSender<()>),
Stop,
}
pub struct Watcher {
root: PathBuf,
config: WatchConfig,
inner: Option<RecommendedWatcher>,
intents: Receiver<CoalescedIntent>,
control: Option<SyncSender<RawMessage>>,
cancelled: Arc<AtomicBool>,
worker_status: Arc<AtomicU8>,
worker: Option<JoinHandle<()>>,
}
#[cfg(test)]
pub(crate) struct ScriptedSender {
root: PathBuf,
raw: SyncSender<RawMessage>,
overflowed: Arc<AtomicBool>,
}
#[cfg(test)]
impl ScriptedSender {
pub(crate) fn send(&self, source: &str) -> Result<()> {
let events =
scripted_events::parse_script(source, &self.root).map_err(Error::WatchScript)?;
for event in events {
enqueue_raw(&self.raw, &self.overflowed, event);
}
Ok(())
}
}
#[derive(Debug)]
pub struct WatchApplyReport {
pub apply: ApplyOutcome,
pub reconciliation: scan::ReconcileReport,
}
impl Watcher {
pub fn new(root: &Path, config: WatchConfig) -> Result<Self> {
config.validate()?;
let root = root.canonicalize().map_err(|e| Error::io(root, e))?;
let (raw_tx, raw_rx) = sync_channel::<RawMessage>(config.event_capacity);
let (intent_tx, intent_rx) = sync_channel::<CoalescedIntent>(config.intent_capacity);
let control_tx = raw_tx.clone();
let overflowed = Arc::new(AtomicBool::new(false));
let callback_overflowed = Arc::clone(&overflowed);
let mut inner = notify::recommended_watcher(move |result| {
enqueue_raw(&raw_tx, &callback_overflowed, result);
})
.map_err(|error| notify_error(&root, error))?;
inner.watch(&root, RecursiveMode::Recursive).map_err(|error| notify_error(&root, error))?;
let worker_root = root.clone();
let cancelled = Arc::new(AtomicBool::new(false));
let worker_cancelled = Arc::clone(&cancelled);
let worker_overflowed = Arc::clone(&overflowed);
let worker_status = Arc::new(AtomicU8::new(WORKER_RUNNING));
let tracked_status = Arc::clone(&worker_status);
let renames = RenameReporting::of_recommended_backend();
let worker = std::thread::Builder::new()
.name("fdu-watch".into())
.spawn(move || {
let _counter_guard = crate::counters::thread_flush_guard();
run_tracked_worker(&tracked_status, || {
run_worker(
&worker_root,
config,
renames,
&raw_rx,
&intent_tx,
&worker_overflowed,
&worker_cancelled,
);
});
})
.map_err(|e| Error::io(&root, e))?;
Ok(Self {
root,
config,
inner: Some(inner),
intents: intent_rx,
control: Some(control_tx),
cancelled,
worker_status,
worker: Some(worker),
})
}
#[cfg(test)]
pub(crate) fn scripted(
root: &Path,
config: WatchConfig,
events: &Path,
) -> Result<(Self, ScriptedSender)> {
config.validate()?;
let root = root.canonicalize().map_err(|error| Error::io(root, error))?;
let scripted = scripted_events::read_script(events, &root).map_err(Error::WatchScript)?;
let (raw_tx, raw_rx) = sync_channel::<RawMessage>(config.event_capacity);
let (intent_tx, intent_rx) = sync_channel::<CoalescedIntent>(config.intent_capacity);
let control_tx = raw_tx.clone();
let overflowed = Arc::new(AtomicBool::new(false));
for event in scripted {
enqueue_raw(&raw_tx, &overflowed, event);
}
let worker_root = root.clone();
let cancelled = Arc::new(AtomicBool::new(false));
let worker_cancelled = Arc::clone(&cancelled);
let worker_overflowed = Arc::clone(&overflowed);
let worker_status = Arc::new(AtomicU8::new(WORKER_RUNNING));
let tracked_status = Arc::clone(&worker_status);
let renames = RenameReporting::of_recommended_backend();
let worker = std::thread::Builder::new()
.name("fdu-scripted-watch".into())
.spawn(move || {
let _counter_guard = crate::counters::thread_flush_guard();
run_tracked_worker(&tracked_status, || {
run_worker(
&worker_root,
config,
renames,
&raw_rx,
&intent_tx,
&worker_overflowed,
&worker_cancelled,
);
});
})
.map_err(|error| Error::io(&root, error))?;
let sender = ScriptedSender {
root: root.clone(),
raw: control_tx.clone(),
overflowed: Arc::clone(&overflowed),
};
Ok((
Self {
root,
config,
inner: None,
intents: intent_rx,
control: Some(control_tx),
cancelled,
worker_status,
worker: Some(worker),
},
sender,
))
}
pub fn next_observation(&self, timeout: Duration) -> Result<Option<Observation>> {
let Some(intent) = self.next_intent(timeout)? else {
return Ok(None);
};
Ok(Some(verify_intent(&self.root, self.config, &intent, &ScanConfig::default())))
}
fn next_intent(&self, timeout: Duration) -> Result<Option<CoalescedIntent>> {
match self.intents.recv_timeout(timeout) {
Ok(intent) => Ok(Some(intent)),
Err(RecvTimeoutError::Timeout) => Ok(None),
Err(RecvTimeoutError::Disconnected) => match self.worker_status.load(Ordering::Acquire)
{
WORKER_PANICKED => Err(Error::WatchWorkerPanicked),
_ => Err(Error::WatchStopped),
},
}
}
pub(crate) fn flush_capture(&self) -> Result<()> {
let Some(control) = self.control.as_ref() else {
return Err(Error::WatchStopped);
};
let (acknowledge, acknowledged) = sync_channel(0);
control.send(RawMessage::Flush(acknowledge)).map_err(|_| Error::WatchStopped)?;
acknowledged.recv().map_err(|_| Error::WatchStopped)
}
pub(crate) const fn capture_backlog_bound(&self) -> usize {
self.config.intent_capacity.saturating_add(1)
}
pub fn apply_next(
&self,
index: &IndexHandle,
scan_config: &ScanConfig,
timeout: Duration,
sink: &mut dyn FnMut(&Commit),
) -> Result<Option<WatchApplyReport>> {
scan_config.validate_for_watch_scope(index.scope()?)?;
let indexed_root = index.root_path()?;
if indexed_root != self.root {
return Err(Error::WatchRootMismatch {
watched: self.root.clone(),
indexed: indexed_root,
});
}
let Some(intent) = self.next_intent(timeout)? else {
return Ok(None);
};
apply_intent(index, &self.root, self.config, &intent, scan_config, sink).map(Some)
}
pub(crate) fn apply_next_controlled(
&self,
index: &IndexHandle,
scan_config: &ScanConfig,
timeout: Duration,
control: &dyn scan::ReconcileControl,
sink: &mut dyn FnMut(&Commit),
) -> Result<Option<WatchApplyReport>> {
scan_config.validate_for_watch_scope(index.scope()?)?;
let indexed_root = index.root_path()?;
if indexed_root != self.root {
return Err(Error::WatchRootMismatch {
watched: self.root.clone(),
indexed: indexed_root,
});
}
let Some(intent) = self.next_intent(timeout)? else {
return Ok(None);
};
apply_intent_controlled(index, &self.root, self.config, &intent, scan_config, control, sink)
.map(Some)
}
}
fn apply_intent(
index: &IndexHandle,
root: &Path,
watch_config: WatchConfig,
intent: &CoalescedIntent,
scan_config: &ScanConfig,
sink: &mut dyn FnMut(&Commit),
) -> Result<WatchApplyReport> {
let mut verifier =
|_: &Path, _: &Observation| Ok(verify_intent(root, watch_config, intent, scan_config));
let apply = apply_reverified_with(index, &Observation::default(), scan_config, &mut verifier)?;
if let Some(commit) = apply.commit.as_ref() {
sink(commit);
}
let reconciliation = scan::reconcile_pending_handle(index, scan_config, sink)?;
retain_unreadable(index, root, &reconciliation, sink)?;
Ok(WatchApplyReport { apply, reconciliation })
}
fn retain_unreadable(
index: &IndexHandle,
root: &Path,
reconciliation: &scan::ReconcileReport,
sink: &mut dyn FnMut(&Commit),
) -> Result<()> {
if reconciliation.scan.errors.is_empty() {
return Ok(());
}
let retained = reconciliation.scan.errors.len().min(crate::MAX_RETAINED_ISSUES);
let issues = reconciliation.scan.errors[..retained]
.iter()
.map(|error| crate::Issue::from_error_under(root, error))
.collect();
let omitted = u64::try_from(reconciliation.scan.errors.len() - retained).unwrap_or(u64::MAX);
let outcome =
index.transition_observation(crate::index::ObservationTransition::Unreadable {
issues,
omitted,
})?;
if let Some(commit) = outcome.commit.as_ref() {
sink(commit);
}
Ok(())
}
fn apply_intent_controlled(
index: &IndexHandle,
root: &Path,
watch_config: WatchConfig,
intent: &CoalescedIntent,
scan_config: &ScanConfig,
control: &dyn scan::ReconcileControl,
sink: &mut dyn FnMut(&Commit),
) -> Result<WatchApplyReport> {
let mut verifier =
|_: &Path, _: &Observation| Ok(verify_intent(root, watch_config, intent, scan_config));
let apply = apply_reverified_with_control(
index,
&Observation::default(),
scan_config,
control,
&mut verifier,
)?;
if let Some(commit) = apply.commit.as_ref() {
sink(commit);
}
let reconciliation =
scan::reconcile_pending_handle_controlled(index, scan_config, control, sink)?;
Ok(WatchApplyReport { apply, reconciliation })
}
#[cfg(test)]
fn apply_observation(
index: &IndexHandle,
observation: &Observation,
scan_config: &ScanConfig,
sink: &mut dyn FnMut(&Commit),
) -> Result<WatchApplyReport> {
scan_config.validate_for_watch_scope(index.scope()?)?;
let apply = apply_reverified(index, observation, scan_config)?;
if let Some(commit) = apply.commit.as_ref() {
sink(commit);
}
let reconciliation = scan::reconcile_pending_handle(index, scan_config, sink)?;
Ok(WatchApplyReport { apply, reconciliation })
}
#[cfg(test)]
fn apply_reverified(
index: &IndexHandle,
observation: &Observation,
scan_config: &ScanConfig,
) -> Result<ApplyOutcome> {
let mut verifier = |root: &Path, observation: &Observation| {
reverify_observation(root, observation, scan_config)
};
apply_reverified_with(index, observation, scan_config, &mut verifier)
}
fn apply_reverified_with(
index: &IndexHandle,
observation: &Observation,
scan_config: &ScanConfig,
verifier: &mut impl FnMut(&Path, &Observation) -> Result<Observation>,
) -> Result<ApplyOutcome> {
let (root, scope, _) = index.watch_boundary()?;
scan_config.validate_for_watch_scope(scope)?;
for _ in 0..MAX_OPTIMISTIC_APPLY_ATTEMPTS {
let clock = index.clock()?;
let candidate = verifier(&root, observation)?;
let candidate = escalate_unknown_ancestry(index, candidate)?;
if let Some(outcome) = index.apply_if_clock(clock, &candidate)? {
return Ok(outcome);
}
}
index.invalidate_root(InvalidateReason::WatchContention)
}
fn apply_reverified_with_control(
index: &IndexHandle,
observation: &Observation,
scan_config: &ScanConfig,
control: &dyn scan::ReconcileControl,
verifier: &mut impl FnMut(&Path, &Observation) -> Result<Observation>,
) -> Result<ApplyOutcome> {
let (root, scope, _) = index.watch_boundary()?;
scan_config.validate_for_watch_scope(scope)?;
for _ in 0..MAX_OPTIMISTIC_APPLY_ATTEMPTS {
control.check_active()?;
let clock = index.clock()?;
let candidate = verifier(&root, observation)?;
let candidate = escalate_unknown_ancestry(index, candidate)?;
control.before_conditional_commit()?;
if let Some(outcome) =
index.apply_opened_if_clock(clock, &candidate, control.max_files())?
{
return Ok(outcome);
}
}
control.before_conditional_commit()?;
let outcome = index.apply_opened(
&Observation::new(vec![Op::InvalidateSubtree {
path: PathBuf::new(),
reason: InvalidateReason::WatchContention,
}]),
control.max_files(),
)?;
Ok(outcome)
}
fn escalate_unknown_ancestry(index: &IndexHandle, candidate: Observation) -> Result<Observation> {
let unknown = index.unknown_ancestry(&candidate)?;
if unknown.is_empty() {
return Ok(candidate);
}
let mut roots: Vec<PathBuf> = unknown.into_iter().map(|(_, root)| root).collect();
roots.sort_by(|left, right| {
left.components().count().cmp(&right.components().count()).then_with(|| left.cmp(right))
});
roots.dedup();
let mut covering: Vec<PathBuf> = Vec::with_capacity(roots.len());
for root in roots {
if !covering.iter().any(|ancestor| root.starts_with(ancestor)) {
covering.push(root);
}
}
let mut ops: Vec<ObservationOp> = candidate
.ops
.into_iter()
.filter(|observed| !covering.iter().any(|root| observed.op.path().starts_with(root)))
.collect();
ops.extend(covering.into_iter().map(|path| {
ObservationOp::unconditional(Op::InvalidateSubtree {
path,
reason: InvalidateReason::UnknownAncestry,
})
}));
Ok(Observation::from_ops(ops))
}
#[cfg(test)]
fn reverify_observation(
root: &Path,
observation: &Observation,
scan_config: &ScanConfig,
) -> Result<Observation> {
let mut ops = Vec::with_capacity(observation.len().saturating_mul(2));
for observed in &observation.ops {
let relative = scan::normalize_subtree(observed.op.path())?;
match &observed.op {
Op::InvalidateSubtree { reason, .. } => {
ops.push(Op::InvalidateSubtree { path: relative, reason: *reason });
}
Op::Upsert { .. } | Op::Remove { .. } => {
let absolute = root.join(&relative);
match scan::observe_path(&absolute) {
Ok((kind, attrs)) => {
match crate::admission::decide_path(
&relative,
kind,
scan_config.hidden(),
scan_config.exclude_special,
) {
crate::admission::Disposition::Retain => {
ops.push(Op::Upsert { path: relative.clone(), kind, attrs });
if let Some(control) =
scan::read_control_op(scan_config, root, &relative, kind)?
{
ops.push(control);
}
}
crate::admission::Disposition::ControlOnly => {
ops.push(Op::Remove { path: relative.clone() });
if let Some(control) =
scan::read_control_op(scan_config, root, &relative, kind)?
{
ops.push(control);
}
}
crate::admission::Disposition::Reject => {
ops.push(Op::Remove { path: relative });
}
}
}
Err(error) => ops.push(op_for_stat_error(relative, &error)),
}
}
Op::ControlUpsert { source, .. } => {
ops.push(Op::ControlUpsert { path: relative, source: source.clone() });
}
Op::ControlRemove { .. } => ops.push(Op::ControlRemove { path: relative }),
}
}
Ok(Observation::new(ops))
}
impl Drop for Watcher {
fn drop(&mut self) {
self.cancelled.store(true, Ordering::Release);
if let Some(control) = self.control.take() {
let _ = control.try_send(RawMessage::Stop);
}
self.inner.take();
if let Some(worker) = self.worker.take() {
let _ = worker.join();
}
}
}
fn notify_error(path: &Path, err: notify::Error) -> Error {
Error::io(path, std::io::Error::other(err))
}
fn enqueue_raw(
sender: &SyncSender<RawMessage>,
overflowed: &AtomicBool,
event: notify::Result<notify::Event>,
) {
match sender.try_send(RawMessage::Event(event)) {
Ok(()) | Err(TrySendError::Disconnected(_)) => {}
Err(TrySendError::Full(_)) => overflowed.store(true, Ordering::Release),
}
}
fn run_tracked_worker(status: &AtomicU8, worker: impl FnOnce()) {
let outcome = catch_unwind(AssertUnwindSafe(worker));
status.store(if outcome.is_ok() { WORKER_STOPPED } else { WORKER_PANICKED }, Ordering::Release);
}
fn run_worker(
root: &Path,
config: WatchConfig,
renames: RenameReporting,
raw: &Receiver<RawMessage>,
out: &SyncSender<CoalescedIntent>,
overflowed: &AtomicBool,
cancelled: &AtomicBool,
) {
let mut pending: BTreeMap<PathBuf, Pending> = BTreeMap::new();
let mut batch_started: Option<Instant> = None;
let mut sticky_overflow = false;
loop {
if cancelled.load(Ordering::Acquire) {
return;
}
if overflowed.swap(false, Ordering::AcqRel) {
collapse_to_overflow(&mut pending);
batch_started.get_or_insert_with(Instant::now);
}
if sticky_overflow {
match try_deliver_overflow(out) {
Ok(true) => sticky_overflow = false,
Ok(false) => {}
Err(()) => return,
}
}
match raw.recv_timeout(config.settle) {
Ok(RawMessage::Event(Ok(event))) => {
record(root, &event, &mut pending, config.batch_path_capacity, renames);
batch_started.get_or_insert_with(Instant::now);
}
Ok(RawMessage::Event(Err(err))) => {
let _ = err;
collapse_to_overflow(&mut pending);
batch_started.get_or_insert_with(Instant::now);
}
Ok(RawMessage::Flush(acknowledge)) => {
if overflowed.swap(false, Ordering::AcqRel) {
collapse_to_overflow(&mut pending);
}
if !pending.is_empty()
&& try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err()
{
return;
}
if sticky_overflow {
match try_deliver_overflow(out) {
Ok(true) => sticky_overflow = false,
Ok(false) => {}
Err(()) => return,
}
}
batch_started = None;
let _ = acknowledge.send(());
}
Ok(RawMessage::Stop) => return,
Err(RecvTimeoutError::Timeout) => {
if !pending.is_empty()
&& try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err()
{
return;
}
batch_started = None;
continue;
}
Err(RecvTimeoutError::Disconnected) => {
if !cancelled.load(Ordering::Acquire) && !pending.is_empty() {
let _ = try_deliver_pending(&mut pending, out, &mut sticky_overflow);
}
return;
}
}
if max_hold_elapsed(batch_started, config.max_hold) {
if try_deliver_pending(&mut pending, out, &mut sticky_overflow).is_err() {
return;
}
batch_started = None;
}
}
}
fn max_hold_elapsed(started: Option<Instant>, max_hold: Duration) -> bool {
started.is_some_and(|start| start.elapsed() >= max_hold)
}
fn record(
root: &Path,
event: ¬ify::Event,
pending: &mut BTreeMap<PathBuf, Pending>,
capacity: usize,
reporting: RenameReporting,
) {
if event.need_rescan() {
let target = if event.paths.len() == 1 {
relative_to(root, &event.paths[0]).unwrap_or_default()
} else {
PathBuf::new()
};
queue_pending(
pending,
target,
Pending::Escalate(InvalidateReason::WatchOverflow),
capacity,
);
return;
}
let renamed = matches!(event.kind, EventKind::Modify(notify::event::ModifyKind::Name(_)));
if renamed
&& (reporting == RenameReporting::OldSideOnly
|| !event.paths.iter().any(|path| relative_to(root, path).is_some()))
{
queue_pending(
pending,
PathBuf::new(),
Pending::Escalate(InvalidateReason::UnpairedRename),
capacity,
);
}
for path in &event.paths {
let Some(rel) = relative_to(root, path) else {
continue; };
if matches!(event.kind, EventKind::Access(_)) {
continue; }
if renamed && rel.as_os_str().is_empty() {
queue_pending(
pending,
PathBuf::new(),
Pending::Escalate(InvalidateReason::UnpairedRename),
capacity,
);
continue;
}
let relist_if_dir = matches!(event.kind, EventKind::Create(_));
queue_pending(pending, rel, Pending::Verify { relist_if_dir, renamed }, capacity);
}
}
fn queue_pending(
pending: &mut BTreeMap<PathBuf, Pending>,
path: PathBuf,
state: Pending,
capacity: usize,
) {
if matches!(pending.get(Path::new("")), Some(Pending::Escalate(_))) {
return;
}
if path.as_os_str().is_empty() && matches!(state, Pending::Escalate(_)) {
pending.clear();
pending.insert(path, state);
return;
}
if let Some(existing) = pending.get_mut(&path) {
match (existing, state) {
(Pending::Escalate(_), _) => {}
(
Pending::Verify { relist_if_dir, renamed },
Pending::Verify { relist_if_dir: relist, renamed: rename },
) => {
*relist_if_dir |= relist;
*renamed |= rename;
}
(slot @ Pending::Verify { .. }, Pending::Escalate(reason)) => {
*slot = Pending::Escalate(reason);
}
}
return;
}
if pending.len() >= capacity {
collapse_to_overflow(pending);
} else {
pending.insert(path, state);
}
}
fn collapse_to_overflow(pending: &mut BTreeMap<PathBuf, Pending>) {
pending.clear();
pending.insert(PathBuf::new(), Pending::Escalate(InvalidateReason::WatchOverflow));
}
fn try_deliver_pending(
pending: &mut BTreeMap<PathBuf, Pending>,
out: &SyncSender<CoalescedIntent>,
sticky_overflow: &mut bool,
) -> std::result::Result<(), ()> {
if pending.is_empty() {
return Ok(());
}
let intent = CoalescedIntent { pending: std::mem::take(pending) };
match out.try_send(intent) {
Ok(()) => Ok(()),
Err(TrySendError::Full(_)) => {
*sticky_overflow = true;
Ok(())
}
Err(TrySendError::Disconnected(_)) => Err(()),
}
}
fn try_deliver_overflow(out: &SyncSender<CoalescedIntent>) -> std::result::Result<bool, ()> {
let mut pending = BTreeMap::new();
collapse_to_overflow(&mut pending);
match out.try_send(CoalescedIntent { pending }) {
Ok(()) => Ok(true),
Err(TrySendError::Full(_)) => Ok(false),
Err(TrySendError::Disconnected(_)) => Err(()),
}
}
fn verify_intent(
root: &Path,
config: WatchConfig,
intent: &CoalescedIntent,
scan_config: &ScanConfig,
) -> Observation {
let mut ops = Vec::with_capacity(intent.pending.len());
let mut listings = ParentListings::default();
let mut unlisted: Vec<PathBuf> = Vec::new();
for (rel, state) in &intent.pending {
match state {
Pending::Escalate(reason) => {
ops.push(Op::InvalidateSubtree { path: rel.clone(), reason: *reason });
}
Pending::Verify { relist_if_dir, renamed } => {
if unlisted.iter().any(|name| rel.starts_with(name)) {
continue;
}
let absolute = root.join(rel);
let mut stat = scan::observe_path(&absolute);
if *renamed && !rel.as_os_str().is_empty() && stat.is_ok() {
let unlisted_reason = match listings.lists(root, rel) {
Some(true) => None,
Some(false) => {
stat = scan::observe_path(&absolute);
stat.is_ok().then_some(InvalidateReason::UnpairedRename)
}
None => Some(InvalidateReason::VerificationFailed),
};
if let Some(reason) = unlisted_reason {
let parent = parent_of(rel);
if !ops.iter().any(
|op| matches!(op, Op::InvalidateSubtree { path, .. } if *path == parent),
) {
ops.push(Op::InvalidateSubtree { path: parent, reason });
}
unlisted.push(rel.clone());
continue;
}
}
match stat {
Ok((kind, attrs)) => {
let disposition = crate::admission::decide_path(
rel,
kind,
scan_config.hidden(),
scan_config.exclude_special,
);
match disposition {
crate::admission::Disposition::Retain => {
ops.push(Op::Upsert { path: rel.clone(), kind, attrs });
match scan::read_control_op(scan_config, root, rel, kind) {
Ok(Some(control)) => ops.push(control),
Ok(None) => {}
Err(_) => ops.push(Op::InvalidateSubtree {
path: rel
.parent()
.map_or_else(PathBuf::new, Path::to_path_buf),
reason: InvalidateReason::VerificationFailed,
}),
}
}
crate::admission::Disposition::ControlOnly => {
match scan::read_control_op(scan_config, root, rel, kind) {
Ok(Some(control)) => {
ops.push(Op::Remove { path: rel.clone() });
ops.push(control);
}
Ok(None) => ops.push(Op::Remove { path: rel.clone() }),
Err(_) => ops.push(Op::InvalidateSubtree {
path: rel
.parent()
.map_or_else(PathBuf::new, Path::to_path_buf),
reason: InvalidateReason::VerificationFailed,
}),
}
}
crate::admission::Disposition::Reject => {
ops.push(Op::Remove { path: rel.clone() });
}
}
let retained_dir =
disposition == crate::admission::Disposition::Retain && kind.is_dir();
if retained_dir && *renamed {
ops.push(Op::InvalidateSubtree {
path: rel.clone(),
reason: InvalidateReason::UnpairedRename,
});
} else if retained_dir && *relist_if_dir && config.relist_new_dirs {
ops.push(Op::InvalidateSubtree {
path: rel.clone(),
reason: InvalidateReason::WatchSetupRace,
});
}
}
Err(error) => {
ops.push(op_for_stat_error(rel.clone(), &error));
if error.kind() == std::io::ErrorKind::NotFound
&& crate::control::path_control_spelling(rel)
== Some(crate::control::ControlSpelling::Variant)
{
let control = crate::control::sibling_control_path(rel);
match scan::read_directory_control_or_removal(
scan_config,
root,
&control,
) {
Ok(Some(observed)) => ops.push(observed),
Ok(None) => {}
Err(_) => ops.push(Op::InvalidateSubtree {
path: parent_of(rel),
reason: InvalidateReason::VerificationFailed,
}),
}
}
}
}
if scan_config.population != crate::query::IgnoredEntries::Include
&& crate::control::path_control_spelling(rel).is_some()
{
ops.push(Op::InvalidateSubtree {
path: rel.parent().map_or_else(PathBuf::new, Path::to_path_buf),
reason: InvalidateReason::ControlPopulationChanged,
});
}
}
}
}
Observation::new(ops)
}
fn parent_of(rel: &Path) -> PathBuf {
rel.parent().map_or_else(PathBuf::new, Path::to_path_buf)
}
#[derive(Default)]
struct ParentListings(BTreeMap<PathBuf, Option<std::collections::HashSet<std::ffi::OsString>>>);
impl ParentListings {
fn lists(&mut self, root: &Path, rel: &Path) -> Option<bool> {
let name = rel.file_name()?;
let parent = parent_of(rel);
if let Some(Some(names)) = self.0.get(&parent) {
if names.contains(name) {
return Some(true);
}
}
let names: Option<std::collections::HashSet<_>> = std::fs::read_dir(root.join(&parent))
.and_then(|entries| entries.map(|entry| entry.map(|entry| entry.file_name())).collect())
.ok();
let listed = names.as_ref().map(|names| names.contains(name));
self.0.insert(parent, names);
listed
}
}
fn op_for_stat_error(path: PathBuf, error: &std::io::Error) -> Op {
match error.kind() {
std::io::ErrorKind::NotFound if path.as_os_str().is_empty() => {
Op::InvalidateSubtree { path, reason: InvalidateReason::VerificationFailed }
}
std::io::ErrorKind::NotFound => Op::Remove { path },
std::io::ErrorKind::NotADirectory => Op::InvalidateSubtree {
path: path.parent().map_or_else(PathBuf::new, Path::to_path_buf),
reason: InvalidateReason::VerificationFailed,
},
_ => Op::InvalidateSubtree { path, reason: InvalidateReason::VerificationFailed },
}
}
fn relative_to(root: &Path, path: &Path) -> Option<PathBuf> {
path.strip_prefix(root).ok().map(Path::to_path_buf)
}
#[cfg(test)]
mod tests {
use super::*;
use notify::event::{CreateKind, Flag, MetadataKind, ModifyKind, RenameMode};
use std::fs;
fn queued_test_watcher(root: PathBuf) -> (SyncSender<CoalescedIntent>, Watcher) {
let (sender, intents) = sync_channel(1);
let watcher = Watcher {
root,
config: WatchConfig::default(),
inner: None,
intents,
control: None,
cancelled: Arc::new(AtomicBool::new(false)),
worker_status: Arc::new(AtomicU8::new(WORKER_RUNNING)),
worker: None,
};
(sender, watcher)
}
#[test]
fn watcher_can_move_to_its_single_consumer_thread() {
fn assert_send<T: Send>() {}
assert_send::<Watcher>();
}
static REAL_WATCHER: std::sync::Mutex<()> = std::sync::Mutex::new(());
const REAL_BACKEND_DELIVERY: Duration = Duration::from_secs(60);
fn real_watcher_guard() -> std::sync::MutexGuard<'static, ()> {
REAL_WATCHER.lock().unwrap_or_else(std::sync::PoisonError::into_inner)
}
enum Waited {
Delivered(Vec<Op>),
Silent,
}
fn wait_for(
watcher: &Watcher,
deadline: Duration,
mut want: impl FnMut(&[Op]) -> bool,
) -> Waited {
let start = Instant::now();
let mut seen: Vec<Op> = Vec::new();
while start.elapsed() < deadline {
match watcher.next_observation(Duration::from_millis(200)) {
Ok(Some(observation)) => {
seen.extend(observation.ops.into_iter().map(|observed| observed.op));
if want(&seen) {
return Waited::Delivered(seen);
}
}
Ok(None) => {}
Err(error) => panic!("watcher stopped while waiting: {error}"),
}
}
if seen.is_empty() {
return Waited::Silent;
}
panic!(
"the backend delivered {} op(s) in {deadline:?} but never the one awaited, so \
this is a disagreement about content rather than a delivery failure: {seen:?}",
seen.len()
);
}
fn wait_established(
watcher: &Watcher,
deadline: Duration,
want: impl FnMut(&[Op]) -> bool,
) -> Vec<Op> {
match wait_for(watcher, deadline, want) {
Waited::Delivered(ops) => ops,
Waited::Silent => panic!(
"the watch was established and then delivered nothing in {deadline:?}, so \
this is a lost event rather than a host precondition; \
FDU_TEST_ALLOW_NO_NATIVE_WATCH does not apply here"
),
}
}
fn establish_watch(watcher: &Watcher, dir: &Path) -> bool {
let warmup = dir.join(".fdu-watch-warmup");
fs::write(&warmup, b"warmup").expect("warmup write");
let waited = wait_for(watcher, REAL_BACKEND_DELIVERY, |ops| !ops.is_empty());
let _ = fs::remove_file(&warmup);
match waited {
Waited::Delivered(_) => true,
Waited::Silent => {
if std::env::var_os("FDU_TEST_ALLOW_NO_NATIVE_WATCH").as_deref()
== Some(std::ffi::OsStr::new("1"))
{
eprintln!(
"skipped by FDU_TEST_ALLOW_NO_NATIVE_WATCH=1: the host event service \
delivered no events to this stream"
);
return false;
}
panic!(
"native watch precondition failed: the host delivered no events in \
{REAL_BACKEND_DELIVERY:?}; run on a host with event delivery, or \
explicitly opt out with FDU_TEST_ALLOW_NO_NATIVE_WATCH=1"
);
}
}
}
#[test]
fn created_files_arrive_as_verified_upserts() {
let _serialized = real_watcher_guard();
let dir = tempfile::tempdir().expect("tempdir");
let watcher = Watcher::new(dir.path(), WatchConfig::default()).expect("watcher");
if !establish_watch(&watcher, dir.path()) {
return;
}
fs::write(dir.path().join("hello.txt"), b"hello world").expect("write");
let ops = wait_established(&watcher, REAL_BACKEND_DELIVERY, |ops| {
ops.iter().any(|op| op.path() == Path::new("hello.txt"))
});
let found = ops
.iter()
.find(|op| op.path() == Path::new("hello.txt"))
.expect("an op for the new file");
match found {
Op::Upsert { attrs, kind, .. } => {
assert!(!kind.is_dir());
assert_eq!(attrs.size, 11);
assert!(attrs.mtime_ns > 0);
}
other => panic!("expected an upsert, got {other:?}"),
}
}
#[test]
fn deleted_files_arrive_as_removes() {
let _serialized = real_watcher_guard();
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("doomed.txt");
fs::write(&path, b"x").expect("write");
let watcher = Watcher::new(dir.path(), WatchConfig::default()).expect("watcher");
if !establish_watch(&watcher, dir.path()) {
return;
}
fs::remove_file(&path).expect("remove");
let ops = wait_established(&watcher, REAL_BACKEND_DELIVERY, |ops| {
ops.iter()
.any(|op| matches!(op, Op::Remove { path } if path == Path::new("doomed.txt")))
});
assert!(
ops.iter()
.any(|op| matches!(op, Op::Remove { path } if path == Path::new("doomed.txt"))),
"expected a remove, saw {ops:?}"
);
}
#[test]
fn watching_a_missing_path_is_an_error() {
let dir = tempfile::tempdir().expect("tempdir");
let missing = dir.path().join("not-there");
assert!(Watcher::new(&missing, WatchConfig::default()).is_err());
}
#[test]
fn paths_outside_the_root_are_ignored() {
let root = Path::new("/a/b");
assert_eq!(relative_to(root, Path::new("/a/b/c/d")), Some(PathBuf::from("c/d")));
assert_eq!(relative_to(root, Path::new("/elsewhere")), None);
}
#[test]
fn verification_errors_distinguish_absence_from_an_invalid_ancestor() {
let path = PathBuf::from("parent/known.txt");
let missing = op_for_stat_error(
path.clone(),
&std::io::Error::new(std::io::ErrorKind::NotFound, "gone"),
);
assert!(matches!(missing, Op::Remove { path: removed } if removed == path));
let not_a_directory = op_for_stat_error(
path.clone(),
&std::io::Error::new(std::io::ErrorKind::NotADirectory, "ancestor is a file"),
);
assert!(matches!(
not_a_directory,
Op::InvalidateSubtree {
path: invalidated,
reason: InvalidateReason::VerificationFailed,
} if invalidated == Path::new("parent")
));
let denied = op_for_stat_error(
path.clone(),
&std::io::Error::new(std::io::ErrorKind::PermissionDenied, "denied"),
);
assert!(matches!(
denied,
Op::InvalidateSubtree {
path: invalidated,
reason: InvalidateReason::VerificationFailed,
} if invalidated == path
));
}
#[test]
fn unknown_watch_ancestry_reconciles_from_the_nearest_known_directory() {
let dir = tempfile::tempdir().expect("tempdir");
let (index, _) =
crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
let handle = crate::IndexHandle::new(index);
let nested = dir.path().join("new/deep");
fs::create_dir_all(&nested).expect("nested directories");
fs::write(nested.join("file.txt"), b"verified").expect("nested file");
let report = apply_observation(
&handle,
&Observation::new(vec![Op::Upsert {
path: PathBuf::from("new/deep/file.txt"),
kind: crate::EntryKind::File,
attrs: crate::Attrs::default(),
}]),
&crate::ScanConfig::default(),
&mut |_| {},
)
.expect("unknown ancestry schedules reconciliation");
assert_eq!(report.apply.invalidated, 1);
assert!(report.reconciliation.is_complete());
assert_eq!(
handle.kind(Path::new("new")).expect("new directory"),
Some(crate::EntryKind::Dir)
);
assert_eq!(
handle.kind(Path::new("new/deep/file.txt")).expect("nested file"),
Some(crate::EntryKind::File)
);
assert_ne!(
handle.attrs(Path::new("new")).expect("verified parent attrs"),
Some(crate::Attrs::default())
);
}
#[test]
fn create_intent_survives_coalescing_but_metadata_only_does_not_relist() {
let root = Path::new("/watch-root");
let path = root.join("directory");
let mut pending = BTreeMap::new();
record(
root,
¬ify::Event::new(EventKind::Create(CreateKind::Folder)).add_path(path.clone()),
&mut pending,
16,
RenameReporting::EachSide,
);
record(
root,
¬ify::Event::new(EventKind::Modify(ModifyKind::Metadata(MetadataKind::Any)))
.add_path(path),
&mut pending,
16,
RenameReporting::EachSide,
);
assert_eq!(
pending.get(Path::new("directory")),
Some(&Pending::Verify { relist_if_dir: true, renamed: false })
);
let mut metadata_only = BTreeMap::new();
record(
root,
¬ify::Event::new(EventKind::Modify(ModifyKind::Metadata(MetadataKind::Any)))
.add_path(root.join("existing")),
&mut metadata_only,
16,
RenameReporting::EachSide,
);
assert_eq!(
metadata_only.get(Path::new("existing")),
Some(&Pending::Verify { relist_if_dir: false, renamed: false })
);
}
fn recorded(
root: &Path,
renames: RenameReporting,
events: &[notify::Event],
) -> BTreeMap<PathBuf, Pending> {
let mut pending = BTreeMap::new();
for event in events {
record(root, event, &mut pending, 16, renames);
}
pending
}
fn rename_event(mode: RenameMode, paths: &[PathBuf]) -> notify::Event {
let mut event = notify::Event::new(EventKind::Modify(ModifyKind::Name(mode)));
event.paths = paths.to_vec();
event
}
fn fsevents_rename(path: PathBuf) -> notify::Event {
rename_event(RenameMode::Any, &[path])
}
const RENAMED: Pending = Pending::Verify { relist_if_dir: false, renamed: true };
#[test]
fn each_rename_shape_verifies_its_named_paths_without_escalating() {
let root = Path::new("/watch-root");
let (old, new) = (root.join("dir/old"), root.join("other/new"));
for (backend, events) in [
("FSEvents", vec![fsevents_rename(old.clone()), fsevents_rename(new.clone())]),
(
"inotify",
vec![
rename_event(RenameMode::From, std::slice::from_ref(&old)),
rename_event(RenameMode::To, std::slice::from_ref(&new)),
rename_event(RenameMode::Both, &[old.clone(), new.clone()]),
],
),
(
"Windows",
vec![
rename_event(RenameMode::From, std::slice::from_ref(&old)),
rename_event(RenameMode::To, std::slice::from_ref(&new)),
],
),
] {
let pending = recorded(root, RenameReporting::EachSide, &events);
assert_eq!(
pending,
BTreeMap::from([
(PathBuf::from("dir/old"), RENAMED),
(PathBuf::from("other/new"), RENAMED),
]),
"{backend}"
);
}
let outside = Path::new("/elsewhere/file");
for (direction, event) in [
("move in", rename_event(RenameMode::Both, &[outside.to_path_buf(), new.clone()])),
("move out", rename_event(RenameMode::Both, &[old.clone(), outside.to_path_buf()])),
] {
let pending = recorded(root, RenameReporting::EachSide, &[event]);
assert_eq!(pending.len(), 1, "{direction}: {pending:?}");
assert!(!pending.contains_key(Path::new("")), "{direction}: {pending:?}");
}
}
#[test]
fn a_sticky_rename_flag_merges_with_the_same_paths_other_events() {
let root = Path::new("/watch-root");
let path = root.join("state.json");
let pending = recorded(
root,
RenameReporting::EachSide,
&[
notify::Event::new(EventKind::Create(CreateKind::File)).add_path(path.clone()),
fsevents_rename(path.clone()),
notify::Event::new(EventKind::Modify(ModifyKind::Data(
notify::event::DataChange::Content,
)))
.add_path(path),
],
);
assert_eq!(
pending,
BTreeMap::from([(
PathBuf::from("state.json"),
Pending::Verify { relist_if_dir: true, renamed: true }
)])
);
}
#[test]
fn unboundable_renames_and_ambiguous_rescans_escalate_the_root() {
let root = Path::new("/watch-root");
let escalated = Some(&Pending::Escalate(InvalidateReason::UnpairedRename));
let kqueue =
recorded(root, RenameReporting::OldSideOnly, &[fsevents_rename(root.join("old"))]);
assert_eq!(kqueue.get(Path::new("")), escalated);
let root_moved = recorded(
root,
RenameReporting::EachSide,
&[rename_event(RenameMode::From, &[root.to_path_buf()])],
);
assert_eq!(root_moved.get(Path::new("")), escalated);
assert_eq!(root_moved.len(), 1, "the root is escalated, never verified: {root_moved:?}");
for paths in [vec![], vec![PathBuf::from("/elsewhere/file")]] {
let unplaced =
recorded(root, RenameReporting::EachSide, &[rename_event(RenameMode::Any, &paths)]);
assert_eq!(unplaced.get(Path::new("")), escalated, "{paths:?}");
}
let rescan = recorded(
root,
RenameReporting::EachSide,
&[notify::Event::new(EventKind::Any)
.add_path(root.join("a"))
.add_path(root.join("b"))
.set_flag(Flag::Rescan)],
);
assert_eq!(
rescan.get(Path::new("")),
Some(&Pending::Escalate(InvalidateReason::WatchOverflow))
);
}
#[test]
fn pending_path_overload_collapses_to_one_root_invalidation() {
let root = Path::new("/watch-root");
let mut pending = BTreeMap::new();
for name in ["one", "two", "three"] {
record(
root,
¬ify::Event::new(EventKind::Any).add_path(root.join(name)),
&mut pending,
2,
RenameReporting::EachSide,
);
}
assert_eq!(pending.len(), 1);
assert_eq!(
pending.get(Path::new("")),
Some(&Pending::Escalate(InvalidateReason::WatchOverflow))
);
}
#[test]
fn continuous_churn_has_a_deterministic_max_hold_ceiling() {
let past = Instant::now()
.checked_sub(Duration::from_secs(2))
.expect("representable earlier instant");
assert!(max_hold_elapsed(Some(past), Duration::from_secs(1)));
assert!(!max_hold_elapsed(None, Duration::from_secs(1)));
}
#[test]
fn backend_enqueue_is_nonblocking_and_marks_overflow() {
let (sender, receiver) = sync_channel(1);
let overflowed = AtomicBool::new(false);
enqueue_raw(&sender, &overflowed, Ok(notify::Event::new(EventKind::Any)));
enqueue_raw(&sender, &overflowed, Ok(notify::Event::new(EventKind::Any)));
assert!(overflowed.load(Ordering::Acquire));
assert!(matches!(receiver.try_recv(), Ok(RawMessage::Event(Ok(_)))));
}
#[test]
fn full_intent_queue_retains_a_sticky_root_invalidation() {
let (sender, receiver) = sync_channel(1);
sender.try_send(CoalescedIntent::default()).expect("fill output");
let mut pending = BTreeMap::from([(
PathBuf::from("lost.txt"),
Pending::Verify { relist_if_dir: false, renamed: false },
)]);
let mut sticky_overflow = false;
try_deliver_pending(&mut pending, &sender, &mut sticky_overflow).expect("connected");
assert!(pending.is_empty());
assert!(sticky_overflow);
receiver.try_recv().expect("make output capacity");
assert!(try_deliver_overflow(&sender).expect("connected"));
let intent = receiver.try_recv().expect("sticky overflow intent");
let observation = verify_intent(
Path::new("/unused"),
WatchConfig::default(),
&intent,
&ScanConfig::default(),
);
assert!(matches!(
&observation.ops[0].op,
Op::InvalidateSubtree {
path,
reason: InvalidateReason::WatchOverflow,
} if path.as_os_str().is_empty()
));
}
#[test]
fn cancellation_wakes_and_joins_with_a_full_intent_queue() {
let dir = tempfile::tempdir().expect("tempdir");
let root = dir.path().canonicalize().expect("canonical root");
let config = WatchConfig {
settle: Duration::from_secs(30),
max_hold: Duration::from_secs(30),
event_capacity: 1,
batch_path_capacity: 1,
intent_capacity: 1,
..WatchConfig::default()
};
let (control, raw) = sync_channel(1);
let (output, intents) = sync_channel(1);
output.try_send(CoalescedIntent::default()).expect("fill intent queue");
let cancelled = Arc::new(AtomicBool::new(false));
let worker_cancelled = Arc::clone(&cancelled);
let status = Arc::new(AtomicU8::new(WORKER_RUNNING));
let tracked_status = Arc::clone(&status);
let overflowed = Arc::new(AtomicBool::new(false));
let worker_overflowed = Arc::clone(&overflowed);
let worker_root = root.clone();
let worker = std::thread::spawn(move || {
run_tracked_worker(&tracked_status, || {
run_worker(
&worker_root,
config,
RenameReporting::EachSide,
&raw,
&output,
&worker_overflowed,
&worker_cancelled,
);
});
});
let watcher = Watcher {
root,
config,
inner: None,
intents,
control: Some(control),
cancelled,
worker_status: status,
worker: Some(worker),
};
let (done_tx, done_rx) = sync_channel(1);
std::thread::spawn(move || {
drop(watcher);
done_tx.send(()).expect("report drop");
});
done_rx
.recv_timeout(Duration::from_secs(5))
.expect("watcher drop must wake and join promptly");
}
#[test]
fn coalescing_defers_filesystem_verification_to_the_consumer() {
let dir = tempfile::tempdir().expect("tempdir");
let relative = PathBuf::from("appeared.txt");
let intent = CoalescedIntent {
pending: BTreeMap::from([(
relative.clone(),
Pending::Verify { relist_if_dir: false, renamed: false },
)]),
};
fs::write(dir.path().join(&relative), b"current").expect("create after coalescing");
let observation =
verify_intent(dir.path(), WatchConfig::default(), &intent, &ScanConfig::default());
assert!(matches!(
&observation.ops[0].op,
Op::Upsert { path, attrs, .. } if path == &relative && attrs.size == 7
));
}
#[test]
fn control_verification_emits_exact_source_with_the_entry_fact() {
let dir = tempfile::tempdir().expect("tempdir");
let relative = PathBuf::from(".gitignore");
fs::write(dir.path().join(&relative), b"*.log\n").expect("write control");
let intent = CoalescedIntent {
pending: BTreeMap::from([(
relative.clone(),
Pending::Verify { relist_if_dir: false, renamed: false },
)]),
};
let config = ScanConfig { read_controls: true, ..ScanConfig::default() };
let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
assert!(matches!(
&observation.ops[0].op,
Op::Upsert { path, kind: crate::EntryKind::File, .. } if path == &relative
));
assert!(matches!(
&observation.ops[1].op,
Op::ControlUpsert { path, source } if path == &relative && source == b"*.log\n"
));
}
#[test]
fn only_population_control_event_reconciles_previously_absent_file() {
let dir = tempfile::tempdir().expect("tree");
fs::write(dir.path().join(".gitignore"), b"# no ignored files\n").expect("control");
fs::write(dir.path().join("debug.log"), b"debug").expect("file");
let scan =
ScanConfig { population: crate::query::IgnoredEntries::Only, ..ScanConfig::default() };
let (index, cold) = crate::scan::scan_into_index(dir.path(), &scan).expect("cold");
assert!(cold.is_complete());
assert!(index.lookup(Path::new("debug.log")).is_none());
let handle = crate::IndexHandle::new(index);
fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("rule edit");
let intent = CoalescedIntent {
pending: BTreeMap::from([(
PathBuf::from(".gitignore"),
Pending::Verify { relist_if_dir: false, renamed: false },
)]),
};
let report =
apply_intent(&handle, dir.path(), WatchConfig::default(), &intent, &scan, &mut |_| {})
.expect("watch apply");
assert!(report.reconciliation.is_complete());
assert!(handle.kind(Path::new("debug.log")).expect("lookup").is_some());
}
#[test]
fn verification_observes_no_control_state_under_a_controls_off_policy() {
let dir = tempfile::tempdir().expect("tempdir");
fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("write control");
fs::write(dir.path().join("debug.log"), b"x").expect("write file");
let config = ScanConfig { read_controls: false, ..ScanConfig::default() };
let (mut index, _) = crate::scan::scan_into_index(dir.path(), &config).expect("scan");
assert!(index.control_table().is_empty());
assert!(matches!(index.controls(), Err(crate::Error::ControlStateNotObserved)));
let control = PathBuf::from(".gitignore");
let intent = CoalescedIntent {
pending: BTreeMap::from([(
control.clone(),
Pending::Verify { relist_if_dir: false, renamed: false },
)]),
};
let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
let reverified = reverify_observation(
dir.path(),
&Observation::new(vec![Op::Remove { path: control }]),
&config,
)
.expect("reverify");
for verified in [&observation, &reverified] {
assert!(
!verified.ops.iter().any(|observed| matches!(
observed.op,
Op::ControlUpsert { .. } | Op::ControlRemove { .. }
)),
"a controls-off policy observed control state: {:?}",
verified.ops
);
}
index.apply(&observation).expect("apply the verified observation");
assert!(index.control_table().is_empty());
assert_eq!(index.scope(), config.scope());
}
#[cfg(unix)]
#[test]
fn applying_verification_uses_the_index_admission_scope() {
use std::os::unix::net::UnixListener;
let dir = tempfile::tempdir().expect("tempdir");
fs::write(dir.path().join(".gitignore"), b"*.log\n").expect("write control");
fs::write(dir.path().join(".secret"), b"hidden").expect("write hidden");
let _listener = UnixListener::bind(dir.path().join("service.sock")).expect("bind socket");
let intent = CoalescedIntent {
pending: [".gitignore", ".secret", "service.sock"]
.into_iter()
.map(|path| {
(PathBuf::from(path), Pending::Verify { relist_if_dir: false, renamed: false })
})
.collect(),
};
let config = ScanConfig {
hidden: Some(Arc::new(crate::HiddenPolicy::prune_hidden::<[&str; 0], &str>([]))),
exclude_special: true,
read_controls: true,
..ScanConfig::default()
};
let observation = verify_intent(dir.path(), WatchConfig::default(), &intent, &config);
assert!(observation.ops.iter().any(|observed| matches!(
&observed.op,
Op::ControlUpsert { path, source }
if path == Path::new(".gitignore") && source == b"*.log\n"
)));
for path in [".gitignore", ".secret", "service.sock"] {
assert!(observation.ops.iter().any(|observed| matches!(
&observed.op,
Op::Remove { path: removed } if removed == Path::new(path)
)));
assert!(!observation.ops.iter().any(|observed| matches!(
&observed.op,
Op::Upsert { path: retained, .. } if retained == Path::new(path)
)));
}
}
#[test]
fn timeout_stop_and_worker_panic_are_distinct() {
let dir = tempfile::tempdir().expect("tempdir");
let (live_sender, live) =
queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
assert!(live.next_observation(Duration::ZERO).expect("timeout").is_none());
drop(live_sender);
let (stopped_sender, stopped) =
queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
stopped.worker_status.store(WORKER_STOPPED, Ordering::Release);
drop(stopped_sender);
assert!(matches!(stopped.next_observation(Duration::ZERO), Err(Error::WatchStopped)));
let (panicked_sender, panicked) =
queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
panicked.worker_status.store(WORKER_PANICKED, Ordering::Release);
drop(panicked_sender);
assert!(matches!(
panicked.next_observation(Duration::ZERO),
Err(Error::WatchWorkerPanicked)
));
}
#[test]
fn tracked_worker_records_a_panic() {
let status = AtomicU8::new(WORKER_RUNNING);
run_tracked_worker(&status, || panic!("injected worker panic"));
assert_eq!(status.load(Ordering::Acquire), WORKER_PANICKED);
}
#[test]
fn observation_driver_closes_the_invalidation_loop() {
let dir = tempfile::tempdir().expect("tempdir");
let (index, _) =
crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
let handle = crate::IndexHandle::new(index);
fs::write(dir.path().join("raced.txt"), b"raced").expect("write");
let observation = Observation::new(vec![Op::InvalidateSubtree {
path: PathBuf::new(),
reason: InvalidateReason::WatchSetupRace,
}]);
apply_observation(&handle, &observation, &crate::ScanConfig::default(), &mut |_| {})
.expect("apply and reconcile");
assert!(handle.kind(Path::new("raced.txt")).expect("query").is_some());
}
#[test]
fn applying_driver_reverifies_a_queued_sample_after_reconciliation() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("sample.txt");
fs::write(&path, b"old").expect("write old sample");
let (index, _) =
crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
let old_attrs = *index.attrs(Path::new("sample.txt")).expect("sample attributes");
let delayed = Observation::new(vec![Op::Upsert {
path: PathBuf::from("sample.txt"),
kind: crate::EntryKind::File,
attrs: old_attrs,
}]);
let handle = crate::IndexHandle::new(index);
fs::write(&path, b"new contents").expect("write current sample");
crate::scan::reconcile_handle(&handle, &crate::ScanConfig::default(), &mut |_| {})
.expect("reconcile newer sample");
let current_size = fs::metadata(&path).expect("sample metadata").len();
apply_observation(&handle, &delayed, &crate::ScanConfig::default(), &mut |_| {})
.expect("apply delayed watch sample");
assert_eq!(
handle.attrs(Path::new("sample.txt")).expect("query").expect("sample remains").size,
current_size
);
}
#[test]
fn blocked_verifier_holds_no_index_lock_and_commits_only_at_current_clock() {
let dir = tempfile::tempdir().expect("tempdir");
let (index, _) =
crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
let handle = crate::IndexHandle::new(index);
let queued = Observation::new(vec![Op::Upsert {
path: PathBuf::from("queued.txt"),
kind: crate::EntryKind::File,
attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
}]);
let applying = handle.clone();
let (entered_tx, entered_rx) = sync_channel(1);
let (release_tx, release_rx) = sync_channel(1);
let (done_tx, done_rx) = sync_channel(1);
let apply_thread = std::thread::spawn(move || {
let mut first = true;
let mut verifier = |_: &Path, observation: &Observation| {
if first {
first = false;
entered_tx.send(()).expect("signal blocked verifier");
release_rx.recv().expect("release blocked verifier");
}
Ok(observation.clone())
};
let result = apply_reverified_with(
&applying,
&queued,
&crate::ScanConfig::default(),
&mut verifier,
);
done_tx.send(result).expect("report apply result");
});
entered_rx
.recv_timeout(Duration::from_secs(5))
.expect("verifier must reach the injected block");
let progressing = handle.clone();
let (progress_tx, progress_rx) = sync_channel(1);
let progress_thread = std::thread::spawn(move || {
let total = progressing.total().expect("reader progresses");
let write = progressing.apply(&Observation::new(vec![Op::Upsert {
path: PathBuf::from("competitor.txt"),
kind: crate::EntryKind::File,
attrs: crate::Attrs { size: 3, allocated: 3, ..crate::Attrs::default() },
}]));
progress_tx.send((total, write)).expect("report progress");
});
let (_, competing_write) = progress_rx
.recv_timeout(Duration::from_secs(5))
.expect("reader and writer must progress while verification is blocked");
competing_write.expect("competing write");
release_tx.send(()).expect("release verifier");
let outcome = done_rx
.recv_timeout(Duration::from_secs(5))
.expect("applying driver completes")
.expect("applying driver succeeds");
apply_thread.join().expect("apply thread");
progress_thread.join().expect("progress thread");
assert_eq!(outcome.inserted, 1);
assert!(handle.kind(Path::new("competitor.txt")).expect("query").is_some());
assert!(handle.kind(Path::new("queued.txt")).expect("query").is_some());
assert_eq!(handle.clock().expect("clock"), crate::Clock(2));
}
#[test]
fn exhausted_watch_contention_stays_unfresh_until_reconciliation() {
let dir = tempfile::tempdir().expect("tempdir");
let (index, _) =
crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
let handle = crate::IndexHandle::new(index);
let queued = Observation::new(vec![Op::Upsert {
path: PathBuf::from("never-committed.txt"),
kind: crate::EntryKind::File,
attrs: crate::Attrs { size: 1, allocated: 1, ..crate::Attrs::default() },
}]);
let mut attempts = 0_usize;
let mut verifier = |_: &Path, observation: &Observation| {
attempts += 1;
handle
.apply(&Observation::new(vec![Op::Upsert {
path: PathBuf::from(format!("competitor-{attempts}.txt")),
kind: crate::EntryKind::File,
attrs: crate::Attrs {
size: attempts as u64,
allocated: attempts as u64,
..crate::Attrs::default()
},
}]))
.expect("force a clock conflict");
Ok(observation.clone())
};
let outcome =
apply_reverified_with(&handle, &queued, &crate::ScanConfig::default(), &mut verifier)
.expect("contention escalates");
assert_eq!(attempts, MAX_OPTIMISTIC_APPLY_ATTEMPTS);
assert_eq!(outcome.invalidated, 1);
assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Stale);
assert!(handle.kind(Path::new("never-committed.txt")).expect("query").is_none());
let pending = handle.take_pending_invalidations().expect("pending invalidation");
assert_eq!(pending, vec![(PathBuf::new(), InvalidateReason::WatchContention)]);
handle.restore_pending_invalidations(pending).expect("restore invalidation");
crate::scan::reconcile_pending_handle(&handle, &crate::ScanConfig::default(), &mut |_| {})
.expect("reconcile contention");
assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Fresh);
}
#[test]
fn verifier_error_mutates_no_shared_state() {
let dir = tempfile::tempdir().expect("tempdir");
let (index, _) =
crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
let handle = crate::IndexHandle::new(index);
let before_clock = handle.clock().expect("clock");
let before_total = handle.total().expect("total");
let mut verifier = |_: &Path, _: &Observation| {
Err(Error::io(
PathBuf::from("blocked"),
std::io::Error::new(std::io::ErrorKind::PermissionDenied, "injected"),
))
};
let error = apply_reverified_with(
&handle,
&Observation::default(),
&crate::ScanConfig::default(),
&mut verifier,
)
.expect_err("verification error");
assert!(matches!(error, Error::Io { .. }));
assert_eq!(handle.clock().expect("clock"), before_clock);
assert_eq!(handle.total().expect("total"), before_total);
assert_eq!(handle.freshness().expect("freshness"), crate::Freshness::Fresh);
assert!(handle.take_pending_invalidations().expect("pending").is_empty());
}
#[test]
fn stable_watch_arbitration_verifies_exactly_once() {
let dir = tempfile::tempdir().expect("tempdir");
let (index, _) =
crate::scan::scan_into_index(dir.path(), &crate::ScanConfig::default()).expect("scan");
let handle = crate::IndexHandle::new(index);
let mut calls = 0_u8;
let mut verifier = |_: &Path, observation: &Observation| {
calls += 1;
Ok(observation.clone())
};
apply_reverified_with(
&handle,
&Observation::new(vec![Op::InvalidateSubtree {
path: PathBuf::new(),
reason: InvalidateReason::Requested,
}]),
&crate::ScanConfig::default(),
&mut verifier,
)
.expect("stable apply");
assert_eq!(calls, 1);
}
#[test]
fn disappearing_watch_root_escalates_instead_of_removing_the_index_root() {
let error = std::io::Error::new(std::io::ErrorKind::NotFound, "root disappeared");
assert!(matches!(
op_for_stat_error(PathBuf::new(), &error),
Op::InvalidateSubtree {
path,
reason: InvalidateReason::VerificationFailed,
} if path.as_os_str().is_empty()
));
}
#[test]
fn observation_driver_rejects_scope_mismatch_before_apply() {
let dir = tempfile::tempdir().expect("tempdir");
let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
let handle = crate::IndexHandle::new(index);
let observation = Observation::new(vec![Op::Upsert {
path: PathBuf::from("deep/nested.txt"),
kind: crate::EntryKind::File,
attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
}]);
let error =
apply_observation(&handle, &observation, &crate::ScanConfig::default(), &mut |_| {})
.expect_err("mismatched scope must fail");
assert!(matches!(error, Error::ScanScopeMismatch { .. }));
assert!(handle.kind(Path::new("deep/nested.txt")).expect("query").is_none());
}
#[test]
fn observation_driver_rejects_restricted_scopes_until_events_are_filtered() {
let dir = tempfile::tempdir().expect("tempdir");
let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
let handle = crate::IndexHandle::new(index);
let observation = Observation::new(vec![Op::Upsert {
path: PathBuf::from("deep/nested.txt"),
kind: crate::EntryKind::File,
attrs: crate::Attrs { size: 5, allocated: 5, ..crate::Attrs::default() },
}]);
let error = apply_observation(&handle, &observation, &shallow, &mut |_| {})
.expect_err("unfiltered bounded watch scope must fail");
assert!(matches!(error, Error::UnsupportedScanConfig(_)));
assert!(handle.kind(Path::new("deep/nested.txt")).expect("query").is_none());
}
#[test]
fn apply_next_rejects_restricted_scope_without_consuming_an_observation() {
let dir = tempfile::tempdir().expect("tempdir");
let shallow = crate::ScanConfig { max_depth: Some(1), ..crate::ScanConfig::default() };
let (index, _) = crate::scan::scan_into_index(dir.path(), &shallow).expect("scan");
let handle = crate::IndexHandle::new(index);
let (sender, watcher) =
queued_test_watcher(dir.path().canonicalize().expect("canonical root"));
sender.try_send(CoalescedIntent::default()).expect("queue intent");
let error = watcher
.apply_next(&handle, &shallow, Duration::ZERO, &mut |_| {})
.expect_err("restricted scope must fail before receive");
assert!(matches!(error, Error::UnsupportedScanConfig(_)));
assert!(watcher.next_observation(Duration::ZERO).expect("receive").is_some());
}
#[test]
fn apply_next_rejects_a_watcher_for_another_root_without_consuming() {
let indexed = tempfile::tempdir().expect("indexed root");
let watched_root_dir = tempfile::tempdir().expect("watched root");
let (index, _) =
crate::scan::scan_into_index(indexed.path(), &crate::ScanConfig::default())
.expect("scan indexed root");
let handle = crate::IndexHandle::new(index);
let (sender, watcher) = queued_test_watcher(
watched_root_dir.path().canonicalize().expect("canonical watched root"),
);
sender.try_send(CoalescedIntent::default()).expect("queue intent");
let error = watcher
.apply_next(&handle, &crate::ScanConfig::default(), Duration::ZERO, &mut |_| {})
.expect_err("mismatched root must fail");
assert!(matches!(error, Error::WatchRootMismatch { .. }));
assert!(watcher.next_observation(Duration::ZERO).expect("receive").is_some());
}
#[cfg(unix)]
#[test]
fn apply_next_walks_an_unreadable_gap_once_and_retains_its_cause() {
use std::os::unix::fs::PermissionsExt;
fn walks_of(commits: &[Commit], path: &Path) -> usize {
commits
.iter()
.flat_map(|commit| commit.state.iter())
.filter(|transition| {
matches!(
transition,
crate::StateTransition::Freshness { path: marked, current, .. }
if marked == path && *current == crate::Freshness::Reconciling
)
})
.count()
}
if !crate::test_support::require_permission_bits() {
return;
}
let dir = tempfile::tempdir().expect("tempdir");
let root = dir.path().canonicalize().expect("canonical root");
let blocked = root.join("blocked");
fs::create_dir(&blocked).expect("blocked");
fs::write(blocked.join("secret"), b"s").expect("fixture");
let (index, _) =
crate::scan::scan_into_index(&root, &crate::ScanConfig::default()).expect("scan");
let handle = crate::IndexHandle::new(index);
let (sender, watcher) = queued_test_watcher(root.clone());
let config = crate::ScanConfig::default();
let mut commits = Vec::new();
fs::set_permissions(&blocked, fs::Permissions::from_mode(0o000)).expect("deny reads");
let mut gap = CoalescedIntent::default();
gap.pending
.insert(PathBuf::from("blocked"), Pending::Escalate(InvalidateReason::WatchOverflow));
sender.try_send(gap).expect("queue the gap");
let first = watcher
.apply_next(&handle, &config, Duration::ZERO, &mut |commit| {
commits.push(commit.clone());
})
.expect("apply the gap")
.expect("an intent was queued");
for name in ["live.txt", "marker.txt"] {
fs::write(root.join(name), name).expect("unrelated mutation");
let mut event = CoalescedIntent::default();
event.pending.insert(
PathBuf::from(name),
Pending::Verify { relist_if_dir: false, renamed: false },
);
sender.try_send(event).expect("queue the event");
watcher
.apply_next(&handle, &config, Duration::ZERO, &mut |commit| {
commits.push(commit.clone());
})
.expect("apply the event")
.expect("an intent was queued");
}
let pending = handle.take_pending_invalidations().expect("pending");
let freshness = handle.freshness_at(Path::new("blocked")).expect("freshness");
let issues = handle.issues().expect("issues");
fs::set_permissions(&blocked, fs::Permissions::from_mode(0o700)).expect("restore reads");
assert!(!first.reconciliation.is_complete(), "the gap must be unreadable");
assert_eq!(
walks_of(&commits, Path::new("blocked")),
1,
"an unreadable subtree must not be re-walked per unrelated event"
);
assert!(pending.is_empty(), "the queue must settle: {pending:?}");
assert_eq!(freshness, crate::Freshness::Partial);
assert!(handle.kind(Path::new("marker.txt")).expect("lookup").is_some());
assert!(
issues.iter().any(|issue| issue.kind == crate::IssueKind::Permission
&& issue.path.as_deref() == Some(Path::new("blocked"))),
"{issues:?}"
);
}
struct RenameFixture {
_dir: tempfile::TempDir,
root: PathBuf,
handle: IndexHandle,
}
impl RenameFixture {
fn new(files: &[&str]) -> Self {
let dir = tempfile::tempdir().expect("tempdir");
let root = dir.path().canonicalize().expect("canonical root");
for file in files {
let path = root.join(file);
fs::create_dir_all(path.parent().expect("parent")).expect("parents");
fs::write(&path, file.as_bytes()).expect("fixture file");
}
let (index, _) =
crate::scan::scan_into_index(&root, &ScanConfig::default()).expect("cold scan");
Self { _dir: dir, root, handle: IndexHandle::new(index) }
}
fn path(&self, rel: &str) -> PathBuf {
self.root.join(rel)
}
fn apply(&self, events: &[notify::Event]) -> Vec<(PathBuf, InvalidateReason)> {
let intent = CoalescedIntent {
pending: recorded(&self.root, RenameReporting::EachSide, events),
};
let mut commits = Vec::new();
let report = apply_intent(
&self.handle,
&self.root,
WatchConfig::default(),
&intent,
&ScanConfig::default(),
&mut |commit| commits.push(commit.clone()),
)
.expect("apply the intent");
assert!(report.reconciliation.is_complete(), "reconciliation must settle");
commits
.iter()
.flat_map(|commit| commit.changes.iter())
.filter_map(|change| match change {
crate::EffectiveChange::Invalidated { path, reason } => {
Some((path.clone(), *reason))
}
_ => None,
})
.collect()
}
fn assert_converged(&self) {
let (cold, _) =
crate::scan::scan_into_index(&self.root, &ScanConfig::default()).expect("cold");
let watched = self.handle.read_with(entries).expect("read the watched index");
assert_eq!(watched, entries(&cold), "watched index diverged from a cold scan");
assert_eq!(self.handle.freshness().expect("freshness"), crate::Freshness::Fresh);
}
}
fn entries(index: &crate::Index) -> BTreeMap<PathBuf, (crate::EntryKind, Option<u64>)> {
fn walk(
index: &crate::Index,
dir: &Path,
out: &mut BTreeMap<PathBuf, (crate::EntryKind, Option<u64>)>,
) {
let Some(children) = index.children(dir) else {
return;
};
let children: Vec<_> =
children.map(|(name, id)| (dir.join(name), id)).collect::<Vec<_>>();
for (path, id) in children {
let kind = index.kind_of(id).expect("live child");
let size = (!kind.is_dir()).then(|| index.attrs_of(id).expect("attrs").size);
out.insert(path.clone(), (kind, size));
walk(index, &path, out);
}
}
let mut out = BTreeMap::new();
walk(index, Path::new(""), &mut out);
out
}
type Classified = (BTreeMap<PathBuf, Option<bool>>, Vec<(PathBuf, Vec<u8>)>);
fn classified(index: &crate::Index) -> Classified {
let ignored = entries(index)
.into_keys()
.map(|path| {
let ignored = index.is_ignored(&path).expect("observed");
(path, ignored)
})
.collect();
let sources = index
.controls()
.expect("observed")
.sources()
.map(|(path, source)| (path, source.to_vec()))
.collect();
(ignored, sources)
}
impl RenameFixture {
fn assert_classified_as_cold(&self, label: &str) -> Classified {
let (cold, _) =
crate::scan::scan_into_index(&self.root, &ScanConfig::default()).expect("cold");
let watched = self.handle.read_with(classified).expect("read the watched index");
assert_eq!(watched, classified(&cold), "{label}: the watch diverged from a cold scan");
watched
}
}
#[test]
fn a_watch_follows_a_case_variant_control_file() {
use crate::test_support::CaseLookups;
let probe = tempfile::tempdir().expect("tempdir");
for (lookups, governs) in CaseLookups::on_this_host(probe.path()) {
let tree = RenameFixture::new(&["up/x.tmp", "up/notes.txt"]);
let _lookups = lookups.install(&tree.root);
let variant = tree.path("up/.GITIGNORE");
let x_ignored = |classified: &Classified| classified.0[Path::new("up/x.tmp")];
fs::write(&variant, b"*.tmp\n").expect("create the variant");
tree.apply(&[created(variant.clone())]);
let appeared = tree.assert_classified_as_cold(&format!("{lookups:?}: created"));
assert_eq!(x_ignored(&appeared), Some(governs), "{lookups:?}");
fs::write(&variant, b"*.txt\n").expect("edit the variant");
tree.apply(&[modified(variant.clone())]);
let edited = tree.assert_classified_as_cold(&format!("{lookups:?}: edited"));
assert_eq!(x_ignored(&edited), Some(false), "{lookups:?}");
assert_eq!(edited.0[Path::new("up/notes.txt")], Some(governs), "{lookups:?}");
fs::remove_file(&variant).expect("remove the variant");
tree.apply(&[notify::Event::new(EventKind::Remove(notify::event::RemoveKind::File))
.add_path(variant.clone())]);
let removed = tree.assert_classified_as_cold(&format!("{lookups:?}: removed"));
assert_eq!(removed.1, Vec::new(), "{lookups:?}: no rules remain");
if lookups == CaseLookups::Host {
let exact = tree.path("up/.gitignore");
fs::write(&exact, b"*.tmp\n").expect("create the exact name");
tree.apply(&[created(exact.clone())]);
for (from, to) in [(&exact, &variant), (&variant, &exact)] {
fs::rename(from, to).expect("case-only rename");
tree.apply(&[fsevents_rename(from.clone()), fsevents_rename(to.clone())]);
tree.assert_converged();
tree.assert_classified_as_cold(&format!("{lookups:?}: renamed to {to:?}"));
}
}
}
}
fn created(path: PathBuf) -> notify::Event {
notify::Event::new(EventKind::Create(CreateKind::Any)).add_path(path)
}
fn modified(path: PathBuf) -> notify::Event {
notify::Event::new(EventKind::Modify(ModifyKind::Data(notify::event::DataChange::Content)))
.add_path(path)
}
#[test]
fn a_file_renamed_within_its_directory_needs_no_invalidation() {
let tree = RenameFixture::new(&["state/config.json", "state/other.txt"]);
fs::write(tree.path("state/config.json.tmp"), b"new configuration").expect("temp");
fs::rename(tree.path("state/config.json.tmp"), tree.path("state/config.json"))
.expect("rename over");
let invalidations = tree.apply(&[
created(tree.path("state/config.json.tmp")),
modified(tree.path("state/config.json.tmp")),
fsevents_rename(tree.path("state/config.json.tmp")),
fsevents_rename(tree.path("state/config.json")),
]);
assert_eq!(invalidations, vec![], "a file rename is settled by its own paths");
tree.assert_converged();
}
#[test]
fn a_file_moved_across_directories_settles_one_side_per_event() {
let tree = RenameFixture::new(&["from/moved.txt", "to/resident.txt"]);
fs::rename(tree.path("from/moved.txt"), tree.path("to/moved.txt")).expect("move");
assert_eq!(tree.apply(&[fsevents_rename(tree.path("to/moved.txt"))]), vec![]);
assert_eq!(tree.apply(&[fsevents_rename(tree.path("from/moved.txt"))]), vec![]);
tree.assert_converged();
}
#[test]
fn a_renamed_directory_relists_only_its_own_subtree() {
let tree = RenameFixture::new(&[
"project/old/one.txt",
"project/old/nested/two.txt",
"project/untouched/three.txt",
]);
fs::rename(tree.path("project/old"), tree.path("project/new")).expect("rename directory");
let invalidations = tree.apply(&[
fsevents_rename(tree.path("project/old")),
fsevents_rename(tree.path("project/new")),
]);
assert_eq!(
invalidations,
vec![(PathBuf::from("project/new"), InvalidateReason::UnpairedRename)]
);
tree.assert_converged();
}
#[test]
fn moves_across_the_root_boundary_settle_from_the_inside_side() {
let tree = RenameFixture::new(&["resident.txt", "leaving/a.txt", "leaving/deep/b.txt"]);
let outside = tempfile::tempdir().expect("outside the root");
fs::create_dir_all(outside.path().join("arriving/deep")).expect("outside tree");
fs::write(outside.path().join("arriving/deep/c.txt"), b"arrived").expect("outside file");
fs::rename(outside.path().join("arriving"), tree.path("arrived")).expect("move in");
let invalidations = tree.apply(&[fsevents_rename(tree.path("arrived"))]);
assert_eq!(
invalidations,
vec![(PathBuf::from("arrived"), InvalidateReason::UnpairedRename)]
);
tree.assert_converged();
fs::rename(tree.path("leaving"), outside.path().join("left")).expect("move out");
assert_eq!(tree.apply(&[fsevents_rename(tree.path("leaving"))]), vec![]);
tree.assert_converged();
}
#[test]
fn a_name_reused_after_a_rename_verifies_as_its_new_entry() {
let tree = RenameFixture::new(&["log.txt", "cache/entry.bin"]);
fs::rename(tree.path("log.txt"), tree.path("log.1.txt")).expect("rotate");
fs::write(tree.path("log.txt"), b"fresh log, longer than the old one").expect("reuse");
fs::rename(tree.path("cache"), tree.path("cache.old")).expect("retire directory");
fs::create_dir(tree.path("cache")).expect("reuse directory name");
let invalidations = tree.apply(&[
fsevents_rename(tree.path("log.txt")),
created(tree.path("log.txt")),
fsevents_rename(tree.path("log.1.txt")),
fsevents_rename(tree.path("cache")),
created(tree.path("cache")),
fsevents_rename(tree.path("cache.old")),
]);
assert!(
invalidations.iter().all(|(path, _)| !path.as_os_str().is_empty()),
"no root reconcile: {invalidations:?}"
);
tree.assert_converged();
assert_eq!(tree.handle.kind(Path::new("cache/entry.bin")).expect("lookup"), None);
}
#[test]
fn a_sticky_rename_flag_on_a_later_write_is_just_a_write() {
let tree = RenameFixture::new(&["sessions/today.jsonl"]);
fs::write(tree.path("sessions/today.jsonl"), b"appended record after the rename")
.expect("write in place");
let invalidations = tree.apply(&[
fsevents_rename(tree.path("sessions/today.jsonl")),
modified(tree.path("sessions/today.jsonl")),
]);
assert_eq!(invalidations, vec![]);
tree.assert_converged();
}
#[test]
fn a_case_only_rename_keeps_one_spelling() {
let tree = RenameFixture::new(&["docs/Readme.md", "docs/Guide/intro.md"]);
let insensitive = crate::test_support::resolves_case_insensitively(&tree.root);
fs::rename(tree.path("docs/Readme.md"), tree.path("docs/README.md")).expect("recase");
fs::rename(tree.path("docs/Guide"), tree.path("docs/guide")).expect("recase directory");
let invalidations = tree.apply(&[
fsevents_rename(tree.path("docs/Readme.md")),
fsevents_rename(tree.path("docs/README.md")),
fsevents_rename(tree.path("docs/Guide")),
fsevents_rename(tree.path("docs/guide")),
modified(tree.path("docs/Guide/intro.md")),
]);
tree.assert_converged();
assert!(
invalidations.iter().all(|(path, _)| !path.as_os_str().is_empty()),
"no root reconcile: {invalidations:?}"
);
if insensitive {
assert!(
invalidations.contains(&(PathBuf::from("docs"), InvalidateReason::UnpairedRename)),
"the stale spelling reconciles its parent: {invalidations:?}"
);
}
}
#[test]
fn a_listing_miss_rereads_the_parent_before_answering() {
let dir = tempfile::tempdir().expect("tempdir");
let root = dir.path();
fs::write(root.join("first"), b"1").expect("first");
let mut listings = ParentListings::default();
assert_eq!(listings.lists(root, Path::new("first")), Some(true));
fs::write(root.join("created-after-the-listing"), b"2").expect("later entry");
assert_eq!(listings.lists(root, Path::new("created-after-the-listing")), Some(true));
assert_eq!(listings.lists(root, Path::new("never-created")), Some(false));
assert_eq!(listings.lists(root, Path::new("missing-directory/child")), None);
}
#[test]
fn native_renames_converge_without_reconciling_the_root() {
let _serialized = real_watcher_guard();
let dir = tempfile::tempdir().expect("tempdir");
let root = dir.path().canonicalize().expect("canonical root");
for file in ["from/moved.txt", "tree/sub/leaf.txt", "to/resident.txt"] {
fs::create_dir_all(root.join(file).parent().expect("parent")).expect("parents");
fs::write(root.join(file), file.as_bytes()).expect("fixture file");
}
let watcher = Watcher::new(&root, WatchConfig::default()).expect("watcher");
if !establish_watch(&watcher, &root) {
return;
}
let config = ScanConfig::default();
let (index, _) = crate::scan::scan_into_index(&root, &config).expect("scan");
let handle = IndexHandle::new(index);
fs::rename(root.join("from/moved.txt"), root.join("to/moved.txt")).expect("move file");
fs::rename(root.join("tree"), root.join("renamed-tree")).expect("rename directory");
let converged = || {
handle.kind(Path::new("to/moved.txt")).expect("lookup").is_some()
&& handle.kind(Path::new("from/moved.txt")).expect("lookup").is_none()
&& handle.kind(Path::new("renamed-tree/sub/leaf.txt")).expect("lookup").is_some()
&& handle.kind(Path::new("tree")).expect("lookup").is_none()
};
let mut commits = Vec::new();
let start = Instant::now();
while !converged() && start.elapsed() < REAL_BACKEND_DELIVERY {
watcher
.apply_next(&handle, &config, Duration::from_millis(200), &mut |commit| {
commits.push(commit.clone());
})
.expect("apply");
}
assert!(converged(), "the renames never converged: {commits:?}");
let root_reconciles = commits
.iter()
.flat_map(|commit| commit.changes.iter())
.filter(|change| {
matches!(change, crate::EffectiveChange::Invalidated { path, .. }
if path.as_os_str().is_empty())
})
.count();
assert_eq!(root_reconciles, 0, "a rename reconciled the whole root: {commits:?}");
let (cold, _) = crate::scan::scan_into_index(&root, &config).expect("cold");
assert_eq!(handle.read_with(entries).expect("read"), entries(&cold));
}
#[test]
fn zero_settle_is_rejected_before_starting_a_busy_worker() {
let config = WatchConfig { settle: Duration::ZERO, ..WatchConfig::default() };
assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
let config = WatchConfig { intent_capacity: 0, ..WatchConfig::default() };
assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
let config = WatchConfig {
batch_path_capacity: MAX_BATCH_PATH_CAPACITY,
intent_capacity: MAX_BUFFERED_INTENT_PATHS / MAX_BATCH_PATH_CAPACITY + 1,
..WatchConfig::default()
};
assert!(matches!(config.validate(), Err(Error::UnsupportedScanConfig(_))));
}
}