Skip to main content

lattice_host/
lsp_watcher.rs

1//! LSP file-watcher subsystem (Phase 5.8.AF.5 / Slice 1).
2//!
3//! Owns a `notify::RecommendedWatcher` and the per-server LSP
4//! subscription map. Lives entirely on a dedicated tokio task on
5//! the LSP runtime — paramount goal #4 (CLAUDE.md): nothing that
6//! does I/O, classification, or LSP fan-out runs on the renderer's
7//! per-tick loop.
8//!
9//! ## Shape
10//!
11//! - `LspFileWatcherHandle` — what `Editor` holds. Just a
12//!   `cmd_tx`. Cheap to construct, cheap to send through, never
13//!   blocks.
14//! - `spawn_lsp_file_watcher_task` — spawns the task on
15//!   `lsp_runtime`. The task owns the watcher, the event rx, and
16//!   a `HashMap<server_id, CachedSubscription>`. It `select!`s
17//!   between commands from `Editor` and `notify` events.
18//! - Editor pushes a `SyncSubscriptions` command whenever the
19//!   actor roster or per-server caps change. The task installs/
20//!   tears down recursive watches and replaces its in-memory
21//!   subscription map atomically.
22//! - When notify fires an event, the task classifies it, matches
23//!   it against every server's `WatcherSubscriptions`, and fans
24//!   out one `workspace/didChangeWatchedFiles` per interested
25//!   server via the cloned `LspSupervisorHandle`. The supervisor
26//!   handle's `did_change_watched_files` notification is
27//!   non-blocking (channel send on the per-server actor).
28//!
29//! Constructing the `notify` watcher and calling `watcher.watch`/
30//! `watcher.unwatch` are sync APIs; they're invoked from inside
31//! the task via [`tokio::task::block_in_place`] so the LSP
32//! runtime (multi-thread) keeps polling other futures while a
33//! large recursive walk is in flight.
34//!
35//! ## What's NOT here (deferred to Slice 2)
36//!
37//! - `.gitignore` / `.ignore` aware filtering at the notify
38//!   callback so events for `target/`, `.git/`, `node_modules/`,
39//!   etc. never enter the channel.
40//! - Per-server registered-glob filtering at the source so we
41//!   don't pay the inotify wakeup cost for irrelevant paths.
42//!
43//! Both are pure-additive improvements that bolt onto this
44//! task's existing structure.
45
46use std::collections::{HashMap, HashSet};
47use std::path::{Path, PathBuf};
48use std::sync::Arc;
49
50use arc_swap::ArcSwap;
51use ignore::gitignore::{Gitignore, GitignoreBuilder};
52use lattice_lsp::lsp_types::FileChangeType;
53use notify::{
54    Event as NotifyEvent, EventKind as NotifyEventKind, RecommendedWatcher, RecursiveMode, Watcher,
55    event::{CreateKind, ModifyKind, RemoveKind},
56};
57use tokio::sync::mpsc;
58
59use lattice_lsp::WatcherSubscriptions;
60
61/// Default exclude globs applied to every workspace root,
62/// independent of `.gitignore`. These cover directories that
63/// modern build tools / package managers churn through
64/// constantly and that no LSP server has any business hearing
65/// about. The list intentionally mirrors VSCode's
66/// `files.watcherExclude` defaults plus a Rust addition.
67///
68/// Phase 5.8.AF.5 Slice 2.
69const DEFAULT_EXCLUDE_GLOBS: &[&str] = &[
70    "**/target/**",
71    "**/.git/**",
72    "**/node_modules/**",
73    "**/.cache/**",
74    "**/.next/**",
75    "**/dist/**",
76    "**/build/**",
77    "**/.venv/**",
78    "**/__pycache__/**",
79    "**/.idea/**",
80    "**/.vscode/**",
81    "**/.DS_Store",
82];
83
84/// Matcher consulted in the notify callback. A path is rejected
85/// when ANY of the per-root `.gitignore` matchers OR the baked
86/// default `GlobSet` says "ignore." Built once per `sync_roots`
87/// call by the watcher task and published into the
88/// `Arc<ArcSwap<...>>` shared with the notify callback.
89pub struct WatcherIgnoreMatcher {
90    /// One per workspace root. Each `Gitignore` covers its root's
91    /// `.gitignore` + `.ignore` (and any nested ones the
92    /// `GitignoreBuilder::add` call picked up).
93    per_root: Vec<(PathBuf, Gitignore)>,
94    /// User-independent baked defaults. Matched against the path
95    /// directly — works whether or not the path lives under a
96    /// known workspace root.
97    defaults: globset::GlobSet,
98}
99
100impl Default for WatcherIgnoreMatcher {
101    fn default() -> Self {
102        Self::build(&HashSet::new())
103    }
104}
105
106impl WatcherIgnoreMatcher {
107    /// Build a matcher covering every root in `roots` plus the
108    /// baked defaults. Reads each root's `.gitignore` + `.ignore`
109    /// from disk; missing files are skipped silently.
110    pub fn build(roots: &HashSet<PathBuf>) -> Self {
111        let mut per_root: Vec<(PathBuf, Gitignore)> = Vec::with_capacity(roots.len());
112        for root in roots {
113            let mut builder = GitignoreBuilder::new(root);
114            // `add` returns `Option<Error>` — missing file is None,
115            // parse error is logged but we keep building.
116            for filename in [".gitignore", ".ignore"] {
117                let p = root.join(filename);
118                if p.exists()
119                    && let Some(e) = builder.add(&p)
120                {
121                    tracing::warn!(
122                        path = %p.display(),
123                        error = %e,
124                        "lsp_watcher: ignore file parse failed (skipping)"
125                    );
126                }
127            }
128            // global gitignore (e.g. ~/.gitignore_global). Empty
129            // result if the user has none configured.
130            let _ = builder.add_line(None, "");
131            match builder.build() {
132                Ok(gi) => per_root.push((root.clone(), gi)),
133                Err(e) => tracing::warn!(
134                    root = %root.display(),
135                    error = %e,
136                    "lsp_watcher: Gitignore build failed (root will only use baked defaults)"
137                ),
138            }
139        }
140        // EF.1: one shared parse-skip-build policy (warn + skip a
141        // bad pattern, empty-set fallback on build failure) instead
142        // of a hand-rolled loop here.
143        let defaults = lattice_runtime::compile_glob_set(DEFAULT_EXCLUDE_GLOBS.iter().copied());
144        Self { per_root, defaults }
145    }
146
147    /// True when the path should be filtered out (never enter
148    /// the channel). Consults the baked defaults first (cheap
149    /// `GlobSet` test) then any matching per-root `.gitignore`.
150    pub fn is_ignored(&self, path: &Path) -> bool {
151        if self.defaults.is_match(path) {
152            return true;
153        }
154        for (root, gi) in &self.per_root {
155            if path.starts_with(root) {
156                // `is_dir = false` is the safe over-approximation:
157                // a file gitignore rule (e.g. `foo.log`) still
158                // matches regardless of whether the path is a
159                // file or dir, and we don't pay the syscall to
160                // resolve the path's true kind.
161                if gi
162                    .matched_path_or_any_parents(path, /* is_dir = */ false)
163                    .is_ignore()
164                {
165                    return true;
166                }
167            }
168        }
169        false
170    }
171
172    /// True when *every* path in the event is ignored. Notify
173    /// events can carry multiple paths (rename = (from, to));
174    /// we filter conservatively — if any path is still
175    /// interesting, the event flows.
176    pub fn event_fully_ignored(&self, event: &NotifyEvent) -> bool {
177        !event.paths.is_empty() && event.paths.iter().all(|p| self.is_ignored(p))
178    }
179}
180
181/// One server's subscription snapshot + the fingerprint used to
182/// detect changes. `Editor` computes the fingerprint, stores its
183/// own `server_id → fingerprint` map, and only sends a
184/// [`WatcherCommand::SyncSubscriptions`] when at least one server's
185/// fingerprint flipped (or the actor roster changed).
186#[derive(Clone)]
187pub struct CachedSubscription {
188    pub fingerprint: u64,
189    pub subs: WatcherSubscriptions,
190}
191
192impl std::fmt::Debug for CachedSubscription {
193    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
194        f.debug_struct("CachedSubscription")
195            .field("fingerprint", &self.fingerprint)
196            .finish_non_exhaustive()
197    }
198}
199
200/// Editor → watcher-task control commands.
201pub enum WatcherCommand {
202    /// Replace the watched-root set + per-server subscription map
203    /// atomically. The task diffs the roots against its current
204    /// set, installs/tears down notify watches, and swaps the
205    /// subscription map.
206    SyncSubscriptions {
207        target_roots: HashSet<PathBuf>,
208        subscriptions: HashMap<String, CachedSubscription>,
209    },
210    /// Tear down everything and exit the task loop. Sent on
211    /// editor shutdown; the task also exits when `cmd_rx` is
212    /// closed (handle dropped).
213    Shutdown,
214}
215
216/// Handle held by `Editor`. Owning this is cheap; sending through
217/// it is non-blocking. Drop the handle to tear down the task
218/// (the watcher itself drops + every inotify watch is released).
219pub struct LspFileWatcherHandle {
220    cmd_tx: mpsc::UnboundedSender<WatcherCommand>,
221}
222
223impl std::fmt::Debug for LspFileWatcherHandle {
224    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
225        f.debug_struct("LspFileWatcherHandle")
226            .finish_non_exhaustive()
227    }
228}
229
230impl LspFileWatcherHandle {
231    /// Send a new target-root + subscription snapshot to the
232    /// task. Best-effort: a dropped task (shouldn't happen in
233    /// production) silently no-ops.
234    pub fn sync(
235        &self,
236        target_roots: HashSet<PathBuf>,
237        subscriptions: HashMap<String, CachedSubscription>,
238    ) {
239        let _ = self.cmd_tx.send(WatcherCommand::SyncSubscriptions {
240            target_roots,
241            subscriptions,
242        });
243    }
244
245    /// Tell the task to exit. Idempotent.
246    pub fn shutdown(&self) {
247        let _ = self.cmd_tx.send(WatcherCommand::Shutdown);
248    }
249}
250
251/// Spawn the watcher task on the LSP runtime. Returns the handle
252/// `Editor` should keep. Returns `Err` only if `notify` itself
253/// fails to construct the OS watcher (out of inotify slots, etc.).
254pub fn spawn_lsp_file_watcher_task(
255    supervisor: lattice_lsp::LspSupervisorHandle,
256    logger: lattice_lsp::LspLogger,
257) -> Result<LspFileWatcherHandle, notify::Error> {
258    let (cmd_tx, cmd_rx) = mpsc::unbounded_channel::<WatcherCommand>();
259    let (event_tx, event_rx) = mpsc::unbounded_channel::<NotifyEvent>();
260    // Slice 2: shared matcher consulted by the notify callback
261    // BEFORE the channel push. Initialised empty (matches only the
262    // baked default globs); the task replaces it inside
263    // `sync_roots` once the workspace roots are known.
264    let ignore_matcher: Arc<ArcSwap<WatcherIgnoreMatcher>> =
265        Arc::new(ArcSwap::from_pointee(WatcherIgnoreMatcher::default()));
266    let matcher_for_callback = ignore_matcher.clone();
267    // Build the OS watcher up front. Its callback is invoked from
268    // notify's own worker thread; we fan only events whose every
269    // path is interesting into our tokio channel.
270    let watcher = RecommendedWatcher::new(
271        move |res: notify::Result<NotifyEvent>| {
272            if let Ok(ev) = res {
273                // Wait-free load — ArcSwap is the right primitive
274                // here, the matcher swaps rarely (sync_roots) and
275                // reads happen at notify-callback rate.
276                let matcher = matcher_for_callback.load();
277                if matcher.event_fully_ignored(&ev) {
278                    return;
279                }
280                let _ = event_tx.send(ev);
281            }
282        },
283        notify::Config::default(),
284    )?;
285
286    lattice_runtime::runtime::spawn_on_lsp_runtime(async move {
287        let mut state = WatcherTaskState {
288            watcher,
289            watched_roots: HashSet::new(),
290            by_server: HashMap::new(),
291            supervisor,
292            logger,
293            ignore_matcher,
294        };
295        state.run(cmd_rx, event_rx).await;
296    });
297
298    Ok(LspFileWatcherHandle { cmd_tx })
299}
300
301/// In-task state. Owned exclusively by the spawned task — no
302/// locks, no shared mutability. Editor talks to it only through
303/// [`WatcherCommand`]s.
304struct WatcherTaskState {
305    watcher: RecommendedWatcher,
306    watched_roots: HashSet<PathBuf>,
307    by_server: HashMap<String, CachedSubscription>,
308    supervisor: lattice_lsp::LspSupervisorHandle,
309    logger: lattice_lsp::LspLogger,
310    /// Slice 2: shared with the notify callback. Rebuilt and
311    /// published via `ArcSwap::store` whenever `sync_roots`
312    /// changes the watched-root set; the callback reads it
313    /// wait-free.
314    ignore_matcher: Arc<ArcSwap<WatcherIgnoreMatcher>>,
315}
316
317impl WatcherTaskState {
318    async fn run(
319        &mut self,
320        mut cmd_rx: mpsc::UnboundedReceiver<WatcherCommand>,
321        mut event_rx: mpsc::UnboundedReceiver<NotifyEvent>,
322    ) {
323        loop {
324            tokio::select! {
325                cmd = cmd_rx.recv() => match cmd {
326                    Some(WatcherCommand::SyncSubscriptions { target_roots, subscriptions }) => {
327                        self.sync_roots(&target_roots);
328                        self.by_server = subscriptions;
329                    }
330                    Some(WatcherCommand::Shutdown) | None => break,
331                },
332                ev = event_rx.recv() => match ev {
333                    Some(event) => self.dispatch_event(event),
334                    None => break,
335                },
336            }
337        }
338    }
339
340    /// Diff `target` against `self.watched_roots`. Add inotify
341    /// watches for new ones, remove for stale. Both `notify`
342    /// calls are sync; wrap in `block_in_place` so a large
343    /// recursive walk on a giant workspace doesn't stall the LSP
344    /// runtime's reactor for other tasks.
345    ///
346    /// Slice 2: rebuild + publish the ignore matcher from the
347    /// new root set BEFORE installing watches. The matcher is
348    /// the gate the notify callback consults; publishing first
349    /// guarantees no flood-window where a freshly-installed
350    /// watch fires events that the matcher hasn't been updated
351    /// to filter.
352    fn sync_roots(&mut self, target: &HashSet<PathBuf>) {
353        // Rebuild + publish ignore matcher first so the notify
354        // callback filters newly-watched paths from the moment
355        // the first event lands.
356        let new_matcher = WatcherIgnoreMatcher::build(target);
357        self.ignore_matcher.store(Arc::new(new_matcher));
358        let stale: Vec<PathBuf> = self
359            .watched_roots
360            .iter()
361            .filter(|p| !target.contains(*p))
362            .cloned()
363            .collect();
364        for p in stale {
365            tracing::info!(path = %p.display(), "lsp_watcher: unwatching stale root");
366            let _ = tokio::task::block_in_place(|| self.watcher.unwatch(&p));
367            self.watched_roots.remove(&p);
368        }
369        let new: Vec<PathBuf> = target
370            .iter()
371            .filter(|p| !self.watched_roots.contains(*p))
372            .cloned()
373            .collect();
374        for p in new {
375            // notify's `Recursive` watch walks the tree
376            // synchronously on Linux. `block_in_place` flags this
377            // to the multi-thread runtime so other tasks on
378            // sibling workers keep making progress.
379            tracing::info!(
380                path = %p.display(),
381                "lsp_watcher: installing recursive watch (may block on large trees)"
382            );
383            let started = std::time::Instant::now();
384            let result =
385                tokio::task::block_in_place(|| self.watcher.watch(&p, RecursiveMode::Recursive));
386            match result {
387                Ok(()) => {
388                    tracing::info!(
389                        path = %p.display(),
390                        elapsed_ms = started.elapsed().as_millis(),
391                        "lsp_watcher: recursive watch installed"
392                    );
393                    self.watched_roots.insert(p);
394                }
395                Err(e) => {
396                    tracing::info!(
397                        path = %p.display(),
398                        elapsed_ms = started.elapsed().as_millis(),
399                        error = %e,
400                        "lsp_watcher: recursive watch failed"
401                    );
402                    self.logger.log(
403                        None,
404                        lattice_lsp::LogLevel::Warn,
405                        lattice_lsp::LogSource::Client,
406                        format!("file-watcher watch {} failed: {e}", p.display()),
407                    );
408                }
409            }
410        }
411    }
412
413    /// Translate a `notify::Event` into per-server batched
414    /// `workspace/didChangeWatchedFiles` notifications and fan
415    /// them out. Notification-only on the LSP side; the actor's
416    /// `did_change_watched_files` is a channel-send.
417    fn dispatch_event(&mut self, event: NotifyEvent) {
418        let Some(kind) = classify(&event) else {
419            return;
420        };
421        // (path, kind) tuples. One event may cover several paths
422        // (notify::Event::paths is a Vec).
423        let classified: Vec<(&Path, FileChangeType)> =
424            event.paths.iter().map(|p| (p.as_path(), kind)).collect();
425        if classified.is_empty() {
426            return;
427        }
428        // Per-server: match the classified paths against the
429        // server's compiled subscription set; build the batch.
430        let mut per_server: HashMap<String, Vec<lattice_lsp::lsp_types::FileEvent>> =
431            HashMap::new();
432        for (server_id, cached) in &self.by_server {
433            if cached.subs.is_empty() {
434                continue;
435            }
436            let mut batch: Vec<lattice_lsp::lsp_types::FileEvent> = Vec::new();
437            for (path, change) in &classified {
438                let hits = cached.subs.matches(path, *change);
439                if hits.is_empty() {
440                    continue;
441                }
442                let uri = lattice_lsp::actor::uri_from_path(path);
443                batch.push(lattice_lsp::lsp_types::FileEvent::new(uri, *change));
444            }
445            if !batch.is_empty() {
446                per_server.insert(server_id.clone(), batch);
447            }
448        }
449        if per_server.is_empty() {
450            return;
451        }
452        // Fan out. The supervisor handle's `running_actors` is a
453        // wait-free ArcSwap load; the per-handle notify is a
454        // channel send.
455        for (_key, handle) in self.supervisor.running_actors() {
456            let server_id = handle.server_id().to_string();
457            let Some(batch) = per_server.remove(&server_id) else {
458                continue;
459            };
460            let params = lattice_lsp::lsp_types::DidChangeWatchedFilesParams { changes: batch };
461            if let Err(e) = handle.did_change_watched_files(params) {
462                let instance = handle.instance();
463                self.logger.log(
464                    Some(&instance),
465                    lattice_lsp::LogLevel::Warn,
466                    lattice_lsp::LogSource::Client,
467                    format!("workspace/didChangeWatchedFiles fan-out failed: {e}"),
468                );
469            }
470        }
471    }
472}
473
474/// Translate one `notify::Event` into an LSP `FileChangeType`.
475pub fn classify(event: &NotifyEvent) -> Option<FileChangeType> {
476    match event.kind {
477        NotifyEventKind::Create(CreateKind::File)
478        | NotifyEventKind::Create(CreateKind::Folder)
479        | NotifyEventKind::Create(CreateKind::Any)
480        | NotifyEventKind::Create(CreateKind::Other) => Some(FileChangeType::CREATED),
481        NotifyEventKind::Modify(ModifyKind::Data(_))
482        | NotifyEventKind::Modify(ModifyKind::Metadata(_))
483        | NotifyEventKind::Modify(ModifyKind::Any)
484        | NotifyEventKind::Modify(ModifyKind::Other)
485        | NotifyEventKind::Modify(ModifyKind::Name(_)) => Some(FileChangeType::CHANGED),
486        NotifyEventKind::Remove(RemoveKind::File)
487        | NotifyEventKind::Remove(RemoveKind::Folder)
488        | NotifyEventKind::Remove(RemoveKind::Any)
489        | NotifyEventKind::Remove(RemoveKind::Other) => Some(FileChangeType::DELETED),
490        _ => None,
491    }
492}
493
494#[cfg(test)]
495mod tests {
496    use super::*;
497    use std::path::PathBuf;
498
499    /// Baked defaults catch the common cases — `target/`, `.git/`,
500    /// `node_modules/` — without any `.gitignore` on disk.
501    /// This is the path that fires for >99% of the fs-event
502    /// flood in the original repro.
503    #[test]
504    fn default_globs_ignore_target_git_node_modules() {
505        let matcher = WatcherIgnoreMatcher::default();
506        for p in [
507            "/foo/bar/target/debug/foo.rlib",
508            "/foo/bar/.git/index",
509            "/foo/bar/node_modules/lodash/index.js",
510            "/foo/bar/.cache/build",
511            "/foo/bar/.next/server/page.js",
512            "/foo/bar/dist/main.js",
513        ] {
514            assert!(
515                matcher.is_ignored(&PathBuf::from(p)),
516                "expected ignored: {p}"
517            );
518        }
519    }
520
521    /// Real source files survive the default filter.
522    #[test]
523    fn default_globs_pass_source_files() {
524        let matcher = WatcherIgnoreMatcher::default();
525        for p in [
526            "/foo/bar/src/main.rs",
527            "/foo/bar/Cargo.toml",
528            "/foo/bar/crates/x/src/lib.rs",
529            "/foo/bar/README.md",
530        ] {
531            assert!(
532                !matcher.is_ignored(&PathBuf::from(p)),
533                "expected not ignored: {p}"
534            );
535        }
536    }
537
538    /// Multi-path events (e.g. notify rename = (from, to)) are
539    /// only suppressed when EVERY path is ignored — preserves
540    /// `target/foo.rs → src/foo.rs` flows where the destination
541    /// is interesting.
542    #[test]
543    fn event_fully_ignored_requires_every_path_ignored() {
544        let matcher = WatcherIgnoreMatcher::default();
545        let mixed = NotifyEvent {
546            kind: NotifyEventKind::Modify(ModifyKind::Any),
547            paths: vec![
548                PathBuf::from("/foo/bar/target/old.rlib"),
549                PathBuf::from("/foo/bar/src/main.rs"),
550            ],
551            attrs: notify::event::EventAttributes::new(),
552        };
553        assert!(
554            !matcher.event_fully_ignored(&mixed),
555            "mixed event must not be suppressed -- src/main.rs is live"
556        );
557
558        let all_ignored = NotifyEvent {
559            kind: NotifyEventKind::Modify(ModifyKind::Any),
560            paths: vec![
561                PathBuf::from("/foo/bar/target/a.rlib"),
562                PathBuf::from("/foo/bar/.git/HEAD"),
563            ],
564            attrs: notify::event::EventAttributes::new(),
565        };
566        assert!(
567            matcher.event_fully_ignored(&all_ignored),
568            "all-paths-ignored event must be suppressed"
569        );
570
571        // Empty paths: spec-vague edge case; preserve the
572        // event rather than dropping (matches notify's intent).
573        let empty = NotifyEvent {
574            kind: NotifyEventKind::Modify(ModifyKind::Any),
575            paths: vec![],
576            attrs: notify::event::EventAttributes::new(),
577        };
578        assert!(!matcher.event_fully_ignored(&empty));
579    }
580
581    /// `.gitignore` files at the workspace root are loaded and
582    /// applied to paths under that root. Verifies the
583    /// `GitignoreBuilder::add` + `matched_path_or_any_parents`
584    /// path actually wires up.
585    #[test]
586    fn per_root_gitignore_filters_workspace_paths() {
587        // Use a unique temp dir per test to avoid cross-pollination.
588        let dir = std::env::temp_dir().join(format!(
589            "lattice-watcher-test-{}-{}",
590            std::process::id(),
591            std::time::SystemTime::now()
592                .duration_since(std::time::UNIX_EPOCH)
593                .unwrap()
594                .as_nanos()
595        ));
596        std::fs::create_dir_all(&dir).expect("create temp dir");
597        std::fs::write(dir.join(".gitignore"), "secret/\n*.log\n").expect("write .gitignore");
598        let mut roots = HashSet::new();
599        roots.insert(dir.clone());
600        let matcher = WatcherIgnoreMatcher::build(&roots);
601
602        assert!(
603            matcher.is_ignored(&dir.join("secret").join("creds.txt")),
604            "secret/ should be ignored per .gitignore"
605        );
606        assert!(
607            matcher.is_ignored(&dir.join("debug.log")),
608            "*.log should be ignored per .gitignore"
609        );
610        assert!(
611            !matcher.is_ignored(&dir.join("src").join("main.rs")),
612            "src/main.rs should pass"
613        );
614
615        // Cleanup.
616        let _ = std::fs::remove_dir_all(&dir);
617    }
618}