Skip to main content

lattice_host/
autoread.rs

1//! Autoread — external-change detection + refresh for file-backed
2//! Document buffers (vim's `autoread`).
3//!
4//! See `docs/dev/architecture/autoread.md` for the design and
5//! `docs/dev/operations/slice-plans/autoread.md` (AR.*) for sequencing.
6//!
7//! AR.0 (this file, first slice) lands the **on-disk fingerprint** only —
8//! the seam every later slice gates on. No watcher yet. The fingerprint is
9//! stamped when a buffer loads and after the editor's own `:w`; the live
10//! `notify` watcher (AR.2) compares an incoming filesystem event's post-read
11//! fingerprint against the stored one to (a) suppress the event its own save
12//! produced and (b) skip no-op `touch`es.
13
14use std::collections::{HashMap, HashSet};
15use std::path::{Path, PathBuf};
16
17// SS.1 (2026-08-11): the fingerprint moved down to `lattice-core` so
18// `lattice-multibuffer` can guard its own source writes with the SAME
19// mechanism. Re-exported here so autoread's callers are unchanged.
20pub use lattice_core::on_disk::OnDiskFingerprint;
21
22use notify::{
23    Event as NotifyEvent, EventKind as NotifyEventKind, RecommendedWatcher, RecursiveMode, Watcher,
24};
25use tokio::sync::mpsc;
26
27// ---------------------------------------------------------------------------
28// AR.2 — the `notify` watcher runtime task.
29//
30// A tokio task on the LSP runtime owns a `notify::RecommendedWatcher` and a
31// dir→basenames map. It watches the **parent directories** of open file-backed
32// buffers **non-recursively** (never a tree — that's what keeps cost tied to
33// open buffers, not project size; see `autoread.md` §3), filters events to the
34// watched basenames, and emits `AutoreadChange`s the host drains (AR.4).
35//
36// Deliberately no task-side debounce: the host's fingerprint gate (`stat`
37// pre-gate + content-hash) already coalesces — a burst of events for one save
38// costs a few cheap host-side `stat`s, the first reloads, the rest are no-ops
39// once the stored fingerprint matches disk. Mirrors `lsp_watcher.rs`.
40// ---------------------------------------------------------------------------
41
42/// What kind of external change the watcher detected for a file.
43#[derive(Debug, Clone, Copy, PartialEq, Eq)]
44pub enum AutoreadChangeKind {
45    /// Created, or content/metadata modified — a reload candidate. The host's
46    /// fingerprint gate decides whether it's a real change.
47    Modified,
48    /// Removed or renamed away — the host keeps the buffer and warns (AR.4).
49    Deleted,
50}
51
52/// A detected external change to a watched file. Emitted by the watcher task,
53/// drained by the host, which maps `path` back to a `BufferId`.
54#[derive(Debug, Clone, PartialEq, Eq)]
55pub struct AutoreadChange {
56    pub path: PathBuf,
57    pub kind: AutoreadChangeKind,
58}
59
60/// Classify a `notify` event for autoread. `Create`/`Modify` ⇒ `Modified`;
61/// `Remove` ⇒ `Deleted`; access/other ⇒ ignored. The classification is only a
62/// hint — the host re-`stat`s on receipt, so a rename mis-labelled `Modified`
63/// still resolves correctly (the host finds the file missing and treats it as
64/// a delete).
65pub fn classify_autoread(event: &NotifyEvent) -> Option<AutoreadChangeKind> {
66    match event.kind {
67        NotifyEventKind::Create(_) | NotifyEventKind::Modify(_) => {
68            Some(AutoreadChangeKind::Modified)
69        }
70        NotifyEventKind::Remove(_) => Some(AutoreadChangeKind::Deleted),
71        _ => None,
72    }
73}
74
75/// True when `path` names a watched file: its parent directory is a watch root
76/// AND its file name is in that directory's watched set. Pure — the O(1) filter
77/// the task applies to every event so a busy shared parent directory costs only
78/// a cheap discard, never a spurious change.
79fn path_is_watched(path: &Path, watches: &HashMap<PathBuf, HashSet<String>>) -> bool {
80    let (Some(parent), Some(name)) = (path.parent(), path.file_name().and_then(|n| n.to_str()))
81    else {
82        return false;
83    };
84    watches
85        .get(parent)
86        .is_some_and(|names| names.contains(name))
87}
88
89/// Editor → watcher-task control commands.
90pub enum AutoreadWatcherCommand {
91    /// Atomically replace the watched set: parent-dir → the file names in that
92    /// dir the editor cares about. The task installs a **non-recursive** watch
93    /// per new dir and removes watches for dirs no longer present.
94    Sync {
95        watches: HashMap<PathBuf, HashSet<String>>,
96    },
97    /// Tear down every watch and exit the loop. The task also exits when
98    /// `cmd_rx` closes (handle dropped).
99    Shutdown,
100}
101
102/// Handle held by `Editor`. Cheap to own; sends are non-blocking. Drop it (or
103/// send `Shutdown`) to tear the task down — the watcher drops and every OS
104/// watch is released.
105pub struct AutoreadWatcherHandle {
106    cmd_tx: mpsc::UnboundedSender<AutoreadWatcherCommand>,
107}
108
109impl std::fmt::Debug for AutoreadWatcherHandle {
110    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
111        f.debug_struct("AutoreadWatcherHandle")
112            .finish_non_exhaustive()
113    }
114}
115
116impl AutoreadWatcherHandle {
117    /// Send the desired watched set. Best-effort: a dropped task silently
118    /// no-ops (shouldn't happen in production).
119    pub fn sync(&self, watches: HashMap<PathBuf, HashSet<String>>) {
120        let _ = self.cmd_tx.send(AutoreadWatcherCommand::Sync { watches });
121    }
122
123    /// Tell the task to exit. Idempotent.
124    pub fn shutdown(&self) {
125        let _ = self.cmd_tx.send(AutoreadWatcherCommand::Shutdown);
126    }
127}
128
129/// Spawn the autoread watcher task on the LSP runtime. Returns the handle
130/// `Editor` keeps plus the change receiver the host drains (AR.4). `Err` only
131/// if `notify` itself fails to construct the OS watcher.
132///
133/// `wake` is the editor's `async_landed` notifier. It is a **required
134/// parameter, not an option**: the change channel is drained by
135/// `run_tick_pending`, which in production runs on the next keystroke or on
136/// an `async_landed` wake. Without the wake an external change (a `git
137/// checkout` from magit, a `git pull` in another terminal) sits in the
138/// channel until the user happens to press a key. Threading it through
139/// [`lattice_mode::inbound::make_inbound_raw`] bakes the wake into every
140/// `send`, so it cannot be forgotten (paramount goal #4 — async-correct by
141/// construction, not by discipline).
142///
143/// `make_inbound_raw` rather than `make_inbound` because the per-item work is
144/// irreducibly `&mut Editor` (the reload rewrites buffer contents, cursor and
145/// scroll), which a `FnMut(T) -> Vec<Effect>` handler cannot capture — the
146/// case that constructor documents.
147pub fn spawn_autoread_watcher_task(
148    wake: std::sync::Arc<tokio::sync::Notify>,
149) -> Result<
150    (
151        AutoreadWatcherHandle,
152        mpsc::UnboundedReceiver<AutoreadChange>,
153    ),
154    notify::Error,
155> {
156    let (cmd_tx, cmd_rx) = mpsc::unbounded_channel::<AutoreadWatcherCommand>();
157    let (event_tx, event_rx) = mpsc::unbounded_channel::<NotifyEvent>();
158    let (change_tx, change_rx) = lattice_mode::inbound::make_inbound_raw::<AutoreadChange>(wake);
159    // The callback runs on notify's own worker thread; forward every event to
160    // the tokio channel and filter inside the task (which holds the map).
161    let watcher = RecommendedWatcher::new(
162        move |res: notify::Result<NotifyEvent>| {
163            if let Ok(ev) = res {
164                let _ = event_tx.send(ev);
165            }
166        },
167        notify::Config::default(),
168    )?;
169    lattice_runtime::runtime::spawn_on_lsp_runtime(async move {
170        let mut task = AutoreadWatcherTask {
171            watcher,
172            watched_dirs: HashSet::new(),
173            watches: HashMap::new(),
174            change_tx,
175        };
176        task.run(cmd_rx, event_rx).await;
177    });
178    Ok((AutoreadWatcherHandle { cmd_tx }, change_rx))
179}
180
181/// In-task state, owned exclusively by the spawned task — no locks. The editor
182/// talks to it only through [`AutoreadWatcherCommand`]s.
183struct AutoreadWatcherTask {
184    watcher: RecommendedWatcher,
185    /// Directories with a live OS watch.
186    watched_dirs: HashSet<PathBuf>,
187    /// dir → the basenames in that dir the editor cares about (the event
188    /// filter).
189    watches: HashMap<PathBuf, HashSet<String>>,
190    /// Wake-baked sender: every `send` notifies the editor's `async_landed`,
191    /// so a change reaches the screen without a keystroke.
192    change_tx: lattice_mode::inbound::InboundBus<AutoreadChange>,
193}
194
195impl AutoreadWatcherTask {
196    async fn run(
197        &mut self,
198        mut cmd_rx: mpsc::UnboundedReceiver<AutoreadWatcherCommand>,
199        mut event_rx: mpsc::UnboundedReceiver<NotifyEvent>,
200    ) {
201        loop {
202            tokio::select! {
203                cmd = cmd_rx.recv() => match cmd {
204                    Some(AutoreadWatcherCommand::Sync { watches }) => self.sync(watches),
205                    Some(AutoreadWatcherCommand::Shutdown) | None => break,
206                },
207                ev = event_rx.recv() => match ev {
208                    Some(event) => self.dispatch(event),
209                    None => break,
210                },
211            }
212        }
213    }
214
215    /// Diff `target` against the live watch set: install a non-recursive watch
216    /// per new dir, drop watches for dirs no longer wanted. Both `notify` calls
217    /// are sync, so wrap them in `block_in_place` to avoid stalling the LSP
218    /// runtime's reactor for sibling tasks.
219    fn sync(&mut self, target: HashMap<PathBuf, HashSet<String>>) {
220        // Canonicalize dir keys so they match the paths `notify` reports.
221        // macOS FSEvents resolves symlinks (`/var` → `/private/var`, and the
222        // temp dir lives under one); inotify echoes the path passed to
223        // `watch()`. Watching the canonical dir — and keying `watches` by it —
224        // keeps event-parent lookups consistent on both. A dir that fails to
225        // canonicalize (vanished) is dropped; its buffers fall back to the
226        // host's on-activate `stat` (AR.3).
227        let target: HashMap<PathBuf, HashSet<String>> = target
228            .into_iter()
229            .filter_map(|(dir, names)| std::fs::canonicalize(&dir).ok().map(|c| (c, names)))
230            .collect();
231        let target_dirs: HashSet<&PathBuf> = target.keys().collect();
232        let stale: Vec<PathBuf> = self
233            .watched_dirs
234            .iter()
235            .filter(|d| !target_dirs.contains(*d))
236            .cloned()
237            .collect();
238        for d in stale {
239            let _ = tokio::task::block_in_place(|| self.watcher.unwatch(&d));
240            self.watched_dirs.remove(&d);
241        }
242        let new: Vec<PathBuf> = target
243            .keys()
244            .filter(|d| !self.watched_dirs.contains(*d))
245            .cloned()
246            .collect();
247        for d in new {
248            // NON-recursive: autoread watches individual parent dirs, never a
249            // tree — this is what bounds cost to open buffers, not project size.
250            match tokio::task::block_in_place(|| {
251                self.watcher.watch(&d, RecursiveMode::NonRecursive)
252            }) {
253                Ok(()) => {
254                    self.watched_dirs.insert(d);
255                }
256                Err(e) => {
257                    // A failed watch downgrades that dir's buffers to the host's
258                    // on-activate `stat` fallback (AR.3); log + skip, never
259                    // panic. `debug!` — watch churn can burst.
260                    tracing::debug!(dir = %d.display(), error = %e, "autoread: watch install failed");
261                }
262            }
263        }
264        self.watches = target;
265    }
266
267    /// Filter one event to the watched basenames and emit a change per match.
268    fn dispatch(&self, event: NotifyEvent) {
269        let Some(kind) = classify_autoread(&event) else {
270            return;
271        };
272        for path in &event.paths {
273            if path_is_watched(path, &self.watches) {
274                let _ = self.change_tx.send(AutoreadChange {
275                    path: path.clone(),
276                    kind,
277                });
278            }
279        }
280    }
281}
282
283#[cfg(test)]
284mod tests {
285    #![allow(clippy::unwrap_used)]
286    use super::*;
287
288    fn temp_path(tag: &str) -> PathBuf {
289        use std::sync::atomic::{AtomicU64, Ordering};
290        static N: AtomicU64 = AtomicU64::new(0);
291        std::env::temp_dir().join(format!(
292            "lattice-autoread-{}-{}-{}",
293            tag,
294            std::process::id(),
295            N.fetch_add(1, Ordering::Relaxed)
296        ))
297    }
298
299    #[test]
300    fn classify_maps_create_modify_remove_and_ignores_access() {
301        use notify::EventKind;
302        use notify::event::{AccessKind, CreateKind, ModifyKind, RemoveKind};
303        assert_eq!(
304            classify_autoread(&NotifyEvent::new(EventKind::Create(CreateKind::File))),
305            Some(AutoreadChangeKind::Modified)
306        );
307        assert_eq!(
308            classify_autoread(&NotifyEvent::new(EventKind::Modify(ModifyKind::Any))),
309            Some(AutoreadChangeKind::Modified)
310        );
311        assert_eq!(
312            classify_autoread(&NotifyEvent::new(EventKind::Remove(RemoveKind::File))),
313            Some(AutoreadChangeKind::Deleted)
314        );
315        assert_eq!(
316            classify_autoread(&NotifyEvent::new(EventKind::Access(AccessKind::Read))),
317            None
318        );
319    }
320
321    #[test]
322    fn path_is_watched_matches_dir_and_basename_only() {
323        let mut names = HashSet::new();
324        names.insert("main.rs".to_string());
325        let mut watches = HashMap::new();
326        watches.insert(PathBuf::from("/proj/src"), names);
327
328        assert!(path_is_watched(Path::new("/proj/src/main.rs"), &watches));
329        // Right dir, wrong basename (the O(1) filter that discards a busy
330        // shared directory's other files).
331        assert!(!path_is_watched(Path::new("/proj/src/other.rs"), &watches));
332        // Right basename, wrong (unwatched) dir.
333        assert!(!path_is_watched(Path::new("/proj/other/main.rs"), &watches));
334    }
335
336    #[test]
337    fn watcher_emits_change_on_external_write() {
338        use std::time::Duration;
339        let dir = temp_path("watch-integ");
340        std::fs::create_dir_all(&dir).unwrap();
341        // Canonicalize so the expected path matches what notify reports
342        // (macOS FSEvents resolves the temp dir's symlink); the watcher
343        // canonicalizes its keys the same way.
344        let dir = std::fs::canonicalize(&dir).unwrap();
345        let file = dir.join("watched.txt");
346        std::fs::write(&file, "v1\n").unwrap();
347
348        let (handle, mut rx) =
349            spawn_autoread_watcher_task(std::sync::Arc::new(tokio::sync::Notify::new()))
350                .expect("spawn watcher");
351        let mut names = HashSet::new();
352        names.insert("watched.txt".to_string());
353        let mut watches = HashMap::new();
354        watches.insert(dir.clone(), names);
355        handle.sync(watches);
356
357        // Let the Sync command install the OS watch before writing.
358        std::thread::sleep(Duration::from_millis(300));
359        std::fs::write(&file, "v2-changed\n").unwrap();
360
361        let got = lattice_runtime::block_on(async move {
362            tokio::time::timeout(Duration::from_secs(5), rx.recv()).await
363        });
364        handle.shutdown();
365        std::fs::remove_dir_all(&dir).ok();
366
367        let change = got
368            .expect("watcher emitted a change before the timeout")
369            .expect("change channel open");
370        assert_eq!(change.path, file);
371        assert_eq!(change.kind, AutoreadChangeKind::Modified);
372    }
373
374    /// An external write must **wake the editor**, not merely queue a change.
375    ///
376    /// `drain_autoread_changes` runs inside `run_tick_pending`, which in
377    /// production is reached two ways: the tail of `App::apply` (i.e. the
378    /// next keystroke) and the actor's `async_landed` select arm. Before
379    /// this test the watcher fired neither, so a `git checkout` (or any
380    /// external edit) sat in the channel until the user happened to press
381    /// a key — the "it works, but only after I hit something" failure mode.
382    ///
383    /// Asserting on `rx.recv()` alone does NOT catch that: the change is
384    /// genuinely in the channel either way. The wake is the thing under
385    /// test, so the wake is what this asserts — and it deliberately does
386    /// not touch the receiver, since draining would mask a missing wake.
387    #[test]
388    fn watcher_wakes_the_editor_on_external_write() {
389        use std::time::Duration;
390        let dir = temp_path("watch-wake");
391        std::fs::create_dir_all(&dir).unwrap();
392        let dir = std::fs::canonicalize(&dir).unwrap();
393        let file = dir.join("watched.txt");
394        std::fs::write(&file, "v1\n").unwrap();
395
396        let wake = std::sync::Arc::new(tokio::sync::Notify::new());
397        let (handle, _rx) =
398            spawn_autoread_watcher_task(std::sync::Arc::clone(&wake)).expect("spawn watcher");
399        let mut names = HashSet::new();
400        names.insert("watched.txt".to_string());
401        let mut watches = HashMap::new();
402        watches.insert(dir.clone(), names);
403        handle.sync(watches);
404
405        std::thread::sleep(Duration::from_millis(300));
406        std::fs::write(&file, "v2-changed\n").unwrap();
407
408        let woke = lattice_runtime::block_on(async move {
409            tokio::time::timeout(Duration::from_secs(5), wake.notified()).await
410        });
411        handle.shutdown();
412        std::fs::remove_dir_all(&dir).ok();
413
414        assert!(
415            woke.is_ok(),
416            "external write must wake the editor; without it the reload \
417             waits for the next keystroke",
418        );
419    }
420}