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 }