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}