git.lucas.co / cce-vault
notes vault library and CLI (Obsidian-compatible)
git clone https://git.lucas.co/cce-vault.git

src/watch.rs (4.9K)

  1 //! A recursive watcher over the vault that hands back debounced batches
  2 //! of changed paths, ready for [`Index::apply_changes`].
  3 //!
  4 //! The debounce is cce-files' shape: wait for a first event, then keep
  5 //! collecting until 150 ms pass without another. A save from an editor is
  6 //! several events (temp file, rename, attribute change) and a sync client
  7 //! landing a folder is hundreds; either arrives as one batch. A stream
  8 //! that never goes quiet (a long sync) is still delivered once a second.
  9 //!
 10 //! The callback runs on the watcher's own thread. A cce-ui app forwards
 11 //! the batch to its event loop (a calloop `Sender`) and applies it there;
 12 //! the index is not shared across threads.
 13 //!
 14 //! [`Index::apply_changes`]: crate::Index::apply_changes
 15 
 16 use std::path::{Component, Path, PathBuf};
 17 use std::sync::mpsc;
 18 use std::time::Duration;
 19 
 20 use notify::{EventKind, RecommendedWatcher, RecursiveMode, Watcher};
 21 
 22 const QUIET: Duration = Duration::from_millis(150);
 23 const MAX_WAIT: Duration = Duration::from_secs(1);
 24 
 25 pub struct VaultWatcher {
 26     // Dropping the watcher closes the channel, which ends the thread.
 27     _watcher: RecommendedWatcher,
 28 }
 29 
 30 impl VaultWatcher {
 31     /// Watch `root` recursively. `on_change` receives each batch of
 32     /// absolute paths, deduplicated and sorted, with anything in a hidden
 33     /// folder or hidden file (`.obsidian/`, `.trash/`, this crate's own
 34     /// `.name.cce-tmp` files) already dropped.
 35     pub fn spawn<F>(root: &Path, mut on_change: F) -> notify::Result<VaultWatcher>
 36     where
 37         F: FnMut(Vec<PathBuf>) + Send + 'static,
 38     {
 39         let root = root.canonicalize().map_err(notify::Error::io)?;
 40         let (tx, rx) = mpsc::channel::<PathBuf>();
 41         let mut watcher = notify::recommended_watcher(move |res: notify::Result<notify::Event>| {
 42             match res {
 43                 Ok(event) if !matches!(event.kind, EventKind::Access(_)) => {
 44                     for p in event.paths {
 45                         let _ = tx.send(p);
 46                     }
 47                 }
 48                 Ok(_) => {}
 49                 Err(e) => log::warn!("vault watcher: {e}"),
 50             }
 51         })?;
 52         watcher.watch(&root, RecursiveMode::Recursive)?;
 53 
 54         let filter_root = root.clone();
 55         std::thread::Builder::new()
 56             .name("cce-vault-watch".into())
 57             .spawn(move || {
 58                 while let Ok(first) = rx.recv() {
 59                     let mut batch = vec![first];
 60                     let started = std::time::Instant::now();
 61                     while let Some(left) = MAX_WAIT.checked_sub(started.elapsed()) {
 62                         match rx.recv_timeout(QUIET.min(left)) {
 63                             Ok(p) => batch.push(p),
 64                             Err(_) => break,
 65                         }
 66                     }
 67                     batch.retain(|p| visible(&filter_root, p));
 68                     batch.sort();
 69                     batch.dedup();
 70                     if !batch.is_empty() {
 71                         on_change(batch);
 72                     }
 73                 }
 74             })
 75             .map_err(notify::Error::io)?;
 76         Ok(VaultWatcher { _watcher: watcher })
 77     }
 78 }
 79 
 80 fn visible(root: &Path, path: &Path) -> bool {
 81     match path.strip_prefix(root) {
 82         Ok(rel) => rel.components().all(|c| match c {
 83             Component::Normal(s) => !s.to_string_lossy().starts_with('.'),
 84             _ => true,
 85         }),
 86         Err(_) => false,
 87     }
 88 }
 89 
 90 #[cfg(test)]
 91 mod tests {
 92     use super::*;
 93     use std::sync::{Arc, Mutex};
 94     use std::time::Instant;
 95 
 96     #[test]
 97     fn batches_arrive_debounced_and_filtered() {
 98         let dir = tempfile::tempdir().unwrap();
 99         let root = dir.path().canonicalize().unwrap();
100         let seen: Arc<Mutex<Vec<Vec<PathBuf>>>> = Arc::default();
101         let sink = seen.clone();
102         let _w = VaultWatcher::spawn(&root, move |b| sink.lock().unwrap().push(b)).unwrap();
103 
104         std::fs::create_dir_all(root.join(".obsidian")).unwrap();
105         std::fs::write(root.join(".obsidian/app.json"), "{}").unwrap();
106         crate::write::atomic_write(&root.join("sub/Note.md"), b"hello").unwrap();
107         std::fs::write(root.join("Other.md"), "x").unwrap();
108 
109         let deadline = Instant::now() + Duration::from_secs(5);
110         loop {
111             let all: Vec<PathBuf> = seen.lock().unwrap().iter().flatten().cloned().collect();
112             // A file written into a brand-new folder can land before the
113             // watch on that folder exists; the folder's own event covers it
114             // (apply_changes walks a folder it is handed).
115             let sub = all.contains(&root.join("sub/Note.md")) || all.contains(&root.join("sub"));
116             if all.contains(&root.join("Other.md")) && sub {
117                 assert!(all.iter().all(|p| visible(&root, p)), "{all:?}");
118                 break;
119             }
120             assert!(Instant::now() < deadline, "no batch within 5 s: {all:?}");
121             std::thread::sleep(Duration::from_millis(20));
122         }
123     }
124 }