Skip to main content

lattice_lsp/
supervisor.rs

1//! `LspSupervisor` -- the per-buffer attachment manager
2//! (Phase 4.1.h).
3//!
4//! Glues the wire-side primitives ([`ServerHandle`], `DocSync`,
5//! `DiagnosticsBus`, [`DiagnosticsLayer`], [`LspLogger`]) into
6//! one editor-facing facade. The App holds exactly one
7//! `LspSupervisor`; everything else flows through it.
8//!
9//! ## What the supervisor owns
10//!
11//! - **Config registry.** A `Vec<Arc<ServerConfig>>` -- the
12//!   curated builtins plus user overrides.
13//! - **Per-(workspace, server-id) actors.** One spawned actor
14//!   per pair. Reused across buffers in the same workspace.
15//! - **Per-(workspace, server-id) DocSync.** One per actor.
16//!   Tracks every URI the actor cares about.
17//! - **Per-URI attachments.** Which `(workspace, server-id)`
18//!   actors care about each URI. A buffer can have multiple
19//!   attachments (rust-analyzer + a clippy bridge for `.rs`).
20//! - **Shared logger** (cloned into every actor; cloned for
21//!   subsystem-level events like supervisor decisions).
22//! - **Shared diagnostics layer** (every actor's
23//!   `DiagnosticsBus` is pumped into it via a tokio task
24//!   spawned at actor-creation time).
25//!
26//! ## Why URI as the key (not `BufferId`)
27//!
28//! `lattice-lsp` is below the UI layer in the crate graph; it
29//! has no concept of `BufferId` (which lives in
30//! `lattice-ui-tui`). URIs are the LSP-native identifier and
31//! map 1:1 to file paths the editor cares about. The App
32//! maintains its own `BufferId → Uri` mapping and threads URIs
33//! into the supervisor's API.
34//!
35//! ## Lifecycle
36//!
37//! ```text
38//!   App::open_file(path, text)
39//!     -> supervisor.open_buffer(path, text).await
40//!         -> match_configs(path) -> [ServerConfig...]
41//!         -> for each config:
42//!              ensure_actor(workspace_root, &config).await
43//!              ensure_doc_sync(actor)
44//!              doc_sync.open(uri, language_id, text)
45//!              record (workspace, server_id) in attachments[uri]
46//!         -> return list of attached ServerHandles
47//!
48//!   App::apply_edit(uri, edit)
49//!     -> supervisor.record_edit(uri, edit)
50//!         -> for each attached (workspace, server_id):
51//!              syncs[(workspace, server_id)].record_edit(uri, edit)
52//!
53//!   App::idle_flush(uri)  // 50ms after last edit
54//!     -> supervisor.flush(uri)
55//!         -> for each attached server: doc_sync.flush(uri)
56//!
57//!   App::close_buffer(uri)
58//!     -> supervisor.close_buffer(uri).await
59//!         -> for each attached server: doc_sync.close(uri)
60//!         -> drop the URI from attachments
61//!         -> diagnostics.clear_uri(uri)
62//! ```
63//!
64//! ## Server reuse
65//!
66//! Two `.rs` files in the same Cargo workspace share one
67//! rust-analyzer actor. Two `.rs` files in different workspaces
68//! get two actors (different roots = different indexed views).
69//! Two distinct languages in the same workspace get two actors
70//! (different ids).
71//!
72//! ## What the supervisor does NOT do (yet)
73//!
74//! - Per-feature dispatch (hover / goto-def / references /
75//!   ...). The supervisor's job is attachment and lifecycle;
76//!   the App calls `servers_for(uri)` and walks the list to
77//!   issue feature requests. Per-feature merging (ranking,
78//!   first-non-empty, union-dedupe) lands in 4.2 alongside the
79//!   nav features.
80//! - Server crash recovery. Today's actor detects pipe close
81//!   and resolves pending with `ActorGone`; the supervisor's
82//!   restart-with-backoff logic lands in 4.4.
83//! - Reading `lsp.toml` -- that's the App's responsibility
84//!   (lattice-lsp has no parser dependencies). The App calls
85//!   `add_config()` for each entry.
86
87use std::collections::{HashMap, VecDeque};
88use std::path::{Path, PathBuf};
89use std::sync::Arc;
90use std::time::{Duration, Instant};
91
92use arc_swap::ArcSwap;
93use lsp_types::Uri;
94use tokio::sync::{mpsc, oneshot};
95
96use crate::actor::{self, ServerHandle};
97use crate::config::{ServerConfig, resolve_workspace_root};
98use crate::diagnostics_layer::{DiagnosticsLayer, pump_diagnostics};
99use crate::error::{LspError, LspResult};
100use crate::logging::{InstanceKey, LogLevel, LogSource, LspLogger};
101use lattice_protocol::edit::Edit;
102
103/// Stable key for an actor: workspace root + server id. Public
104/// so [`SupervisorSnapshot`] consumers can pattern-match on the
105/// actor map without re-typing the tuple.
106pub type ActorKey = (PathBuf, String);
107
108/// Build an `InstanceKey` from an `ActorKey` for use with
109/// [`LspLogger::log`] (B'.2). Cheap: two `Arc::from` calls
110/// over the existing `String` / `PathBuf` references.
111fn instance_for(key: &ActorKey) -> InstanceKey {
112    InstanceKey::new(
113        Arc::<str>::from(key.1.as_str()),
114        Arc::<std::path::Path>::from(key.0.as_path()),
115    )
116}
117
118/// 4.4.d: result of a successful [`LspSupervisor::restart_server`]
119/// (or the handle-side equivalent). The App surfaces these to
120/// the user in the echo area ("restarted rust in 1 workspace; 4
121/// buffers replayed") and tests assert on them.
122#[derive(Debug, Clone)]
123pub struct RestartReport {
124    pub server_id: String,
125    /// Every `(workspace, server_id)` pair that re-spawned.
126    pub respawned: Vec<ActorKey>,
127    /// URIs that the replay path issued `didOpen` for so the
128    /// new actor sees the same buffer set.
129    pub replayed_uris: Vec<Uri>,
130}
131
132/// One LSP subsystem per editor instance.
133pub struct LspSupervisor {
134    /// Curated config registry. Order is preference order: when
135    /// multiple configs match a path, all of them attach (we
136    /// don't disambiguate by priority at attachment time --
137    /// priority resolves *feature-dispatch* ties only).
138    configs: Vec<Arc<ServerConfig>>,
139    /// Per-(workspace, server-id) actors, lazily spawned.
140    /// `ServerHandle` is itself `Arc`-wrapped internally so
141    /// cloning is cheap; no need to wrap it again.
142    actors: HashMap<ActorKey, ServerHandle>,
143    /// Per-URI attachments: which actors are tracking this URI.
144    /// Stored as ActorKey so feature dispatch can resolve
145    /// `uri -> [ServerHandle, ...]` in O(handles).
146    ///
147    /// DocSync no longer lives here -- each actor owns its own
148    /// mirror. The supervisor stays out of the edit path entirely
149    /// so the UI thread cannot stall behind a flush.
150    attachments: HashMap<Uri, Vec<ActorKey>>,
151    /// Shared logger. Every actor gets a clone; subsystem
152    /// events (supervisor decisions, attach / detach) emit
153    /// directly through this.
154    logger: LspLogger,
155    /// Shared diagnostics layer. Every actor's
156    /// `DiagnosticsBus` is pumped into this.
157    diagnostics: DiagnosticsLayer,
158    /// Server-initiated `workspace/applyEdit` bus (Phase 4.3).
159    /// `Some` once the App has called
160    /// [`Self::set_apply_edit_bus`] at startup; cloned into
161    /// every actor at spawn so they all forward inbound
162    /// applyEdit requests through the same channel. `None`
163    /// disables the feature -- the actor falls back to a
164    /// METHOD_NOT_FOUND response (matches the pre-4.3
165    /// behaviour).
166    apply_edit_bus: Option<crate::apply_edit::ApplyEditBus>,
167    /// Server-initiated `workspace/configuration` bus (Phase
168    /// 4.1 follow-up). `Some` once the App has called
169    /// [`Self::set_configuration_bus`] at startup; cloned into
170    /// every actor at spawn. `None` falls back to
171    /// `Vec<null>`-per-item replies (the pre-this-commit
172    /// stub).
173    configuration_bus: Option<crate::configuration::ConfigurationBus>,
174    /// Server-initiated `window/showDocument` bus (4.4.b).
175    /// `Some` once the App calls [`Self::set_show_document_bus`];
176    /// `None` leaves the actor falling back to `success: false`
177    /// so the server knows the host can't open URIs.
178    show_document_bus: Option<crate::show_document::ShowDocumentBus>,
179    /// Server-initiated `window/showMessageRequest` bus (4.4.b).
180    /// `Some` once the App calls
181    /// [`Self::set_show_message_request_bus`]; `None` leaves the
182    /// actor replying with JSON `null` (spec-compliant "user
183    /// dismissed without picking").
184    show_message_request_bus: Option<crate::show_message_request::ShowMessageRequestBus>,
185    /// Editor-wide event bus. `Some` after the App calls
186    /// [`Self::set_event_bus`] (early in startup, before any
187    /// buffer opens). When set, every spawned actor gets a
188    /// per-actor [`crate::fan_in`] task subscribed to
189    /// `DocumentChanged` events; that task forwards each
190    /// applied edit straight into the actor's `cmd_tx` --
191    /// keeping the UI thread out of the LSP edit path
192    /// entirely.
193    event_bus: Option<Arc<lattice_runtime::EventBus>>,
194    /// Per-actor fan-in subscription ids. Stored alongside the
195    /// actor so shutdown can call `unsubscribe` and stop the
196    /// bus from holding a dead sender.
197    fan_in_subs: HashMap<ActorKey, lattice_runtime::SubscriptionId>,
198    /// 4.4.d: per-server-id recent restart timestamps. Each
199    /// `:lsp-restart` (and each auto-restart on crash) appends
200    /// the current `Instant`; older entries (> 60s) are pruned
201    /// before the gate check fires. The cap (`MAX_RESTARTS`) is
202    /// per-window so a healthy server that crashes once a day
203    /// keeps restarting forever, but a runaway crash loop
204    /// surfaces as a refusal after the third attempt.
205    restart_history: HashMap<String, VecDeque<Instant>>,
206}
207
208/// 4.4.d: maximum restarts allowed inside [`RESTART_WINDOW`].
209/// Beyond this the supervisor refuses to restart the same server
210/// id until the window slides past. Three is the smallest number
211/// that lets a "spawn → crash → spawn" double-tap recover (one
212/// restart for the user's `:lsp-restart`, one for the post-crash
213/// auto, one for a fast self-correcting re-crash) without
214/// hiding a runaway loop.
215pub(crate) const MAX_RESTARTS: usize = 3;
216/// 4.4.d: rolling window the restart counter respects.
217pub(crate) const RESTART_WINDOW: Duration = Duration::from_secs(60);
218/// 4.4.d: base delay between restarts; multiplied by `2^n` where
219/// `n` is the number of restarts already inside the window.
220/// Capped at [`RESTART_BACKOFF_MAX`].
221pub(crate) const RESTART_BACKOFF_BASE: Duration = Duration::from_millis(250);
222pub(crate) const RESTART_BACKOFF_MAX: Duration = Duration::from_secs(30);
223
224fn compute_restart_backoff(prior_count: usize) -> Duration {
225    if prior_count == 0 {
226        return Duration::ZERO;
227    }
228    let factor = 1u32 << prior_count.min(7); // cap shift at 2^7
229    let scaled = RESTART_BACKOFF_BASE.saturating_mul(factor);
230    scaled.min(RESTART_BACKOFF_MAX)
231}
232
233impl LspSupervisor {
234    /// Construct an empty supervisor with shared logger +
235    /// diagnostics layer. Use `add_config` to populate the
236    /// config registry.
237    pub fn new(logger: LspLogger) -> Self {
238        let diagnostics = DiagnosticsLayer::new(logger.clone());
239        Self {
240            configs: Vec::new(),
241            actors: HashMap::new(),
242            attachments: HashMap::new(),
243            logger,
244            diagnostics,
245            apply_edit_bus: None,
246            configuration_bus: None,
247            show_document_bus: None,
248            show_message_request_bus: None,
249            event_bus: None,
250            fan_in_subs: HashMap::new(),
251            restart_history: HashMap::new(),
252        }
253    }
254
255    /// Install the editor event bus. Must be called once at
256    /// startup before any buffer opens; after the call every
257    /// actor spawned by the supervisor gets a per-actor fan-in
258    /// task (see [`crate::fan_in`]) that turns
259    /// `Event::DocumentChanged` into [`crate::actor::ServerHandle::record_edit`]
260    /// without ever taking the supervisor's mutex.
261    ///
262    /// Calling twice replaces the bus reference; existing actors
263    /// keep their original fan-in subscriptions (we don't
264    /// re-spawn them on bus swap because that would race with
265    /// in-flight events).
266    pub fn set_event_bus(&mut self, bus: Arc<lattice_runtime::EventBus>) {
267        self.event_bus = Some(bus.clone());
268        // Spawn fan-ins for any actors that pre-date the bus
269        // (in practice none -- the App calls this before opening
270        // files -- but the path is here for symmetry with
271        // attach_handle).
272        let pending: Vec<(ActorKey, ServerHandle)> = self
273            .actors
274            .iter()
275            .filter(|(k, _)| !self.fan_in_subs.contains_key(*k))
276            .map(|(k, h)| (k.clone(), h.clone()))
277            .collect();
278        for (key, handle) in pending {
279            let id = crate::fan_in::spawn(handle, bus.clone());
280            self.fan_in_subs.insert(key, id);
281        }
282    }
283
284    /// Install the apply-edit bus (Phase 4.3). The App calls
285    /// this once at startup with the sender side of the channel
286    /// it created; every actor spawned after this point gets a
287    /// clone and forwards inbound `workspace/applyEdit`
288    /// requests through it. Calling twice replaces the bus;
289    /// any actors already spawned keep their original clone
290    /// (we don't track them for retro-fitting in v1).
291    pub fn set_apply_edit_bus(&mut self, bus: crate::apply_edit::ApplyEditBus) {
292        self.apply_edit_bus = Some(bus);
293    }
294
295    /// Install the configuration bus (Phase 4.1 follow-up).
296    /// Same shape as [`Self::set_apply_edit_bus`]: cloned into
297    /// every actor spawned after the call so server-initiated
298    /// `workspace/configuration` requests reach the App's
299    /// drain. `None` falls back to per-item `null` replies.
300    pub fn set_configuration_bus(&mut self, bus: crate::configuration::ConfigurationBus) {
301        self.configuration_bus = Some(bus);
302    }
303
304    /// 4.4.b: install the show-document bus. Cloned into every
305    /// actor spawned (or restarted) after this call. The App
306    /// calls this once at startup with the sender side of the
307    /// channel whose receiver it holds.
308    pub fn set_show_document_bus(&mut self, bus: crate::show_document::ShowDocumentBus) {
309        self.show_document_bus = Some(bus);
310    }
311
312    /// 4.4.b: install the show-message-request bus.
313    pub fn set_show_message_request_bus(
314        &mut self,
315        bus: crate::show_message_request::ShowMessageRequestBus,
316    ) {
317        self.show_message_request_bus = Some(bus);
318    }
319
320    /// Borrow the shared logger (so callers can register their
321    /// own subsystem-level events).
322    pub fn logger(&self) -> &LspLogger {
323        &self.logger
324    }
325
326    /// Borrow the shared diagnostics layer.
327    pub fn diagnostics(&self) -> &DiagnosticsLayer {
328        &self.diagnostics
329    }
330
331    /// Add a server config to the registry. The App calls this
332    /// for every builtin + user-override config at startup.
333    pub fn add_config(&mut self, config: ServerConfig) {
334        self.configs.push(Arc::new(config));
335    }
336
337    /// Set the registry from an iterator (e.g. the curated
338    /// builtins). Replaces any prior contents.
339    pub fn set_configs<I: IntoIterator<Item = ServerConfig>>(&mut self, configs: I) {
340        self.configs = configs.into_iter().map(Arc::new).collect();
341    }
342
343    /// All registered configs (read-only; for `:lsp-status`).
344    pub fn configs(&self) -> &[Arc<ServerConfig>] {
345        &self.configs
346    }
347
348    /// True iff at least one configured server's `file_patterns`
349    /// matches `path`. M.5.2 uses this from the App's
350    /// `MajorEntered` hook to decide whether to auto-activate
351    /// `lsp-mode` on the buffer; if no server cares about the
352    /// path, there's nothing to gate on.
353    pub fn has_server_for_path(&self, path: &Path) -> bool {
354        self.configs
355            .iter()
356            .any(|c| matches_any_pattern(path, &c.file_patterns))
357    }
358
359    /// Every actor currently running. Used by `:lsp-status`.
360    pub fn running_actors(&self) -> Vec<(ActorKey, ServerHandle)> {
361        self.actors
362            .iter()
363            .map(|(k, h)| (k.clone(), h.clone()))
364            .collect()
365    }
366
367    /// Snapshot of every running actor's `ServerHandle`. Used by
368    /// workspace-scoped LSP requests (e.g. `workspace/symbol`)
369    /// that fan out across every server, not just servers
370    /// attached to one buffer.
371    pub fn all_running_handles(&self) -> Vec<ServerHandle> {
372        self.actors.values().cloned().collect()
373    }
374
375    /// Number of buffers currently attached to the actor at
376    /// `key`. Cheap walk over `attachments`; used by
377    /// `:lsp-server-log` to surface per-server buffer counts in
378    /// the picker margin.
379    pub fn buffer_count_for(&self, key: &ActorKey) -> usize {
380        self.attachments
381            .values()
382            .filter(|keys| keys.contains(key))
383            .count()
384    }
385
386    /// Open a buffer. Walks the config registry, spawns
387    /// matching actors as needed, attaches the buffer to each,
388    /// and emits `didOpen` per server. Returns the list of
389    /// attached `ServerHandle`s -- the App stores this so
390    /// feature dispatch knows where to issue requests.
391    ///
392    /// `path` is the buffer's filesystem path; `text` is the
393    /// initial buffer text.
394    pub async fn open_buffer(
395        &mut self,
396        path: PathBuf,
397        text: String,
398    ) -> LspResult<Vec<ServerHandle>> {
399        let uri = crate::actor::uri_from_path(&path);
400        // Already open under another path? If the URI is
401        // already in attachments, surface a no-op rather than
402        // re-issuing didOpen (servers reject duplicate
403        // didOpens).
404        if self.attachments.contains_key(&uri) {
405            return Ok(self.servers_for(&uri));
406        }
407
408        let matches: Vec<Arc<ServerConfig>> = self
409            .configs
410            .iter()
411            .filter(|c| matches_any_pattern(&path, &c.file_patterns))
412            .cloned()
413            .collect();
414
415        if matches.is_empty() {
416            // No server cares about this buffer; not an error.
417            // App still tracks the URI; we just don't store
418            // attachments for it.
419            return Ok(Vec::new());
420        }
421
422        let mut handles: Vec<ServerHandle> = Vec::new();
423        let mut keys: Vec<ActorKey> = Vec::new();
424
425        for config in matches {
426            let workspace =
427                resolve_workspace_root(path.parent().unwrap_or(&path), &config.root_markers);
428            let key: ActorKey = (workspace.clone(), config.id.clone());
429            let handle = match self.actors.get(&key) {
430                Some(existing) => existing.clone(),
431                None => {
432                    self.logger.log(
433                        None,
434                        LogLevel::Info,
435                        LogSource::Client,
436                        format!(
437                            "supervisor: spawning {} for workspace {}",
438                            config.id,
439                            workspace.display()
440                        ),
441                    );
442                    let h = match actor::spawn(
443                        (*config).clone(),
444                        workspace.clone(),
445                        self.logger.clone(),
446                        self.apply_edit_bus.clone(),
447                        self.configuration_bus.clone(),
448                        self.show_document_bus.clone(),
449                        self.show_message_request_bus.clone(),
450                        self.event_bus.clone(),
451                    )
452                    .await
453                    {
454                        Ok(h) => h,
455                        Err(e) => {
456                            self.logger.log(
457                                None,
458                                LogLevel::Warn,
459                                LogSource::Client,
460                                format!(
461                                    "supervisor: spawn failed for {} ({}): {}",
462                                    config.id,
463                                    workspace.display(),
464                                    e
465                                ),
466                            );
467                            // Skip this server; continue with
468                            // others. Server-not-on-PATH is the
469                            // common case and shouldn't sink the
470                            // open.
471                            continue;
472                        }
473                    };
474                    // Spawn the diagnostics pump.
475                    let rx = h.subscribe_diagnostics();
476                    let layer = self.diagnostics.clone();
477                    tokio::spawn(pump_diagnostics(layer, rx));
478                    self.actors.insert(key.clone(), h.clone());
479                    // Spawn the per-actor edit fan-in if the
480                    // editor's event bus is wired up. Stored
481                    // SubscriptionId is unsubscribed on shutdown
482                    // so the bus doesn't keep dead senders.
483                    if let Some(bus) = self.event_bus.clone() {
484                        let sub = crate::fan_in::spawn(h.clone(), bus);
485                        self.fan_in_subs.insert(key.clone(), sub);
486                    }
487                    h
488                }
489            };
490
491            // Drive the DocSync inside the actor: single writer
492            // means no contention with the edit / flush path.
493            handle.open_doc(uri.clone(), config.language_id.clone(), text.clone())?;
494
495            handles.push(handle);
496            keys.push(key);
497        }
498
499        if !keys.is_empty() {
500            self.attachments.insert(uri.clone(), keys);
501            self.logger.log(
502                None,
503                LogLevel::Info,
504                LogSource::Client,
505                format!(
506                    "supervisor: opened {} attached to {} server(s)",
507                    uri.as_str(),
508                    handles.len()
509                ),
510            );
511        }
512        Ok(handles)
513    }
514
515    /// Attach a buffer to a pre-built `ServerHandle`. Used by:
516    ///
517    /// - **Tests** -- the in-process `MockServer` returns a
518    ///   handle that the supervisor wouldn't normally spawn.
519    /// - **Custom transports** (future) -- a TCP / named-pipe
520    ///   server that bypasses `ChildTransport`.
521    ///
522    /// Identical effect to the actor-spawning branch of
523    /// `open_buffer`: registers the actor under
524    /// `(workspace_root, server_id)`, builds a `DocSync`,
525    /// emits `didOpen`, records the attachment, and starts the
526    /// diagnostics pump if it isn't already running.
527    pub fn attach_handle(
528        &mut self,
529        uri: Uri,
530        workspace_root: PathBuf,
531        server_id: String,
532        language_id: String,
533        text: String,
534        handle: ServerHandle,
535    ) -> LspResult<()> {
536        let key: ActorKey = (workspace_root, server_id);
537        // Insert / reuse the actor.
538        let was_new = !self.actors.contains_key(&key);
539        if was_new {
540            self.actors.insert(key.clone(), handle.clone());
541            // Spawn the diagnostics pump for the new actor.
542            let rx = handle.subscribe_diagnostics();
543            let layer = self.diagnostics.clone();
544            tokio::spawn(pump_diagnostics(layer, rx));
545            // Spawn the per-actor edit fan-in if the bus is wired.
546            if let Some(bus) = self.event_bus.clone() {
547                let sub = crate::fan_in::spawn(handle.clone(), bus);
548                self.fan_in_subs.insert(key.clone(), sub);
549            }
550        }
551        // Open inside the actor; single-writer DocSync.
552        handle.open_doc(uri.clone(), language_id, text)?;
553        // Record attachment.
554        self.attachments.entry(uri).or_default().push(key);
555        Ok(())
556    }
557
558    /// Close a buffer. Flushes pending changes per attached
559    /// server, sends `didClose`, and drops the URI from
560    /// attachments + the diagnostics layer.
561    pub fn close_buffer(&mut self, uri: &Uri) -> LspResult<()> {
562        let keys = match self.attachments.remove(uri) {
563            Some(k) => k,
564            None => return Ok(()),
565        };
566        for key in &keys {
567            // The actor owns the DocSync, so it pairs the final
568            // flush + didClose internally and sends them in
569            // order.
570            let Some(handle) = self.actors.get(key) else {
571                continue;
572            };
573            let _ = handle.close_doc(uri.clone());
574        }
575        self.diagnostics.clear_uri(uri);
576        self.logger.log(
577            None,
578            LogLevel::Info,
579            LogSource::Client,
580            format!("supervisor: closed {}", uri.as_str()),
581        );
582        Ok(())
583    }
584
585    /// Forward an edit to every attached actor's mailbox.
586    ///
587    /// **Not on the editor's hot path.** Production edits flow
588    /// through the editor event bus to a per-actor fan-in (see
589    /// [`crate::fan_in`]), which means the UI thread never
590    /// takes the supervisor mutex on a keystroke. This method
591    /// remains for tests + admin tooling that drive the
592    /// supervisor directly without standing up an event bus.
593    pub fn record_edit(&mut self, uri: &Uri, edit: &Edit) -> LspResult<()> {
594        let keys = match self.attachments.get(uri) {
595            Some(k) => k.clone(),
596            None => return Ok(()),
597        };
598        for key in &keys {
599            let Some(handle) = self.actors.get(key) else {
600                continue;
601            };
602            if let Err(e) = handle.record_edit(uri.clone(), edit.clone()) {
603                self.logger.log(
604                    Some(&instance_for(key)),
605                    LogLevel::Warn,
606                    LogSource::Client,
607                    format!("record_edit on {}: {}", uri.as_str(), e),
608                );
609            }
610        }
611        Ok(())
612    }
613
614    /// Force-flush queued changes for `uri` across every
615    /// attached server, bypassing the per-actor debounce. Used
616    /// before synchronous requests that require the server to
617    /// have seen the latest text -- notably `willSaveWaitUntil`
618    /// and `:lsp-flush`.
619    pub fn flush(&mut self, uri: &Uri) -> LspResult<()> {
620        let keys = match self.attachments.get(uri) {
621            Some(k) => k.clone(),
622            None => return Ok(()),
623        };
624        for key in &keys {
625            let Some(handle) = self.actors.get(key) else {
626                continue;
627            };
628            if let Err(e) = handle.flush(uri.clone()) {
629                self.logger.log(
630                    Some(&instance_for(key)),
631                    LogLevel::Warn,
632                    LogSource::Client,
633                    format!("flush on {}: {}", uri.as_str(), e),
634                );
635            }
636        }
637        Ok(())
638    }
639
640    /// Force-flush every open URI's queued changes across every
641    /// attached server. Used at editor shutdown so each server
642    /// sees a coherent final state before `didClose`. Per-edit
643    /// flushing is debounced inside each actor; this is the
644    /// only "drain everything now" affordance.
645    pub fn flush_all(&mut self) -> LspResult<()> {
646        for (key, handle) in &self.actors {
647            if let Err(e) = handle.flush_all() {
648                self.logger.log(
649                    Some(&instance_for(key)),
650                    LogLevel::Warn,
651                    LogSource::Client,
652                    format!("flush_all: {}", e),
653                );
654            }
655        }
656        Ok(())
657    }
658
659    /// Every server attached to `uri`. The App walks this list
660    /// to dispatch features (hover, goto-definition, ...).
661    pub fn servers_for(&self, uri: &Uri) -> Vec<ServerHandle> {
662        let keys = match self.attachments.get(uri) {
663            Some(k) => k,
664            None => return Vec::new(),
665        };
666        keys.iter()
667            .filter_map(|k| self.actors.get(k).cloned())
668            .collect()
669    }
670
671    /// Number of currently-attached buffers.
672    pub fn attached_buffer_count(&self) -> usize {
673        self.attachments.len()
674    }
675
676    /// Number of currently-running actors.
677    pub fn running_actor_count(&self) -> usize {
678        self.actors.len()
679    }
680
681    /// 4.4.d: force-restart every actor with id `server_id`.
682    ///
683    /// Steps:
684    ///   1. Prune restart history outside the private `RESTART_WINDOW`.
685    ///   2. If `history.len() >= MAX_RESTARTS`, refuse (caller
686    ///      surfaces the cooldown to the user).
687    ///   3. Sleep `compute_restart_backoff(history.len())`.
688    ///   4. For every `(workspace, server_id)` actor pair: shut
689    ///      down the existing handle, drop the fan-in
690    ///      subscription, spawn a fresh actor with the same
691    ///      config, restart the diagnostics pump + fan-in.
692    ///   5. Replay `didOpen` for every URI that pointed at the
693    ///      old actor so the new actor sees the same workspace
694    ///      state.
695    ///   6. Push `Instant::now()` to history.
696    ///
697    /// Replayed `didOpen` uses each tracked URI's last-known
698    /// text; today the supervisor doesn't keep that text (each
699    /// actor's `DocSync` owns the mirror), so the replay uses
700    /// the empty string and the App is expected to re-send the
701    /// full buffer via the next `didChange` debounce. Future
702    /// revision can snapshot text into `RestartReport` and let
703    /// the App ferry it back via `record_edit` proactively.
704    ///
705    /// Returns the list of `(workspace, server_id)` keys that
706    /// were re-spawned and the URIs that need re-opening.
707    pub async fn restart_server(&mut self, server_id: &str) -> LspResult<RestartReport> {
708        // Step 1: prune the window.
709        let now = Instant::now();
710        let history = self
711            .restart_history
712            .entry(server_id.to_string())
713            .or_default();
714        while let Some(front) = history.front().copied()
715            && now.duration_since(front) >= RESTART_WINDOW
716        {
717            history.pop_front();
718        }
719
720        // Step 2: gate.
721        if history.len() >= MAX_RESTARTS {
722            let oldest = history.front().copied().unwrap_or(now);
723            let resume = oldest + RESTART_WINDOW;
724            let cooldown = resume.saturating_duration_since(now);
725            return Err(LspError::HandshakeFailed(format!(
726                "{server_id}: too many restarts ({} in the last {}s); next allowed in {}ms",
727                history.len(),
728                RESTART_WINDOW.as_secs(),
729                cooldown.as_millis(),
730            )));
731        }
732
733        // Step 3: backoff before doing any work.
734        let backoff = compute_restart_backoff(history.len());
735        if !backoff.is_zero() {
736            tokio::time::sleep(backoff).await;
737        }
738
739        // Snapshot every (workspace, server_id) match -- usually
740        // one actor, but a single server can attach to multiple
741        // workspaces if the user opens files across different
742        // roots.
743        let targets: Vec<ActorKey> = self
744            .actors
745            .keys()
746            .filter(|(_, id)| id == server_id)
747            .cloned()
748            .collect();
749        if targets.is_empty() {
750            return Err(LspError::HandshakeFailed(format!(
751                "{server_id}: no running actor with that id"
752            )));
753        }
754        let config = self
755            .configs
756            .iter()
757            .find(|c| c.id == server_id)
758            .cloned()
759            .ok_or_else(|| {
760                LspError::HandshakeFailed(format!("{server_id}: no config registered"))
761            })?;
762
763        let mut respawned: Vec<ActorKey> = Vec::new();
764        let mut replayed_uris: Vec<Uri> = Vec::new();
765
766        for key in targets {
767            // Step 4: shut down + replace.
768            let workspace = key.0.clone();
769            let old_handle = match self.actors.remove(&key) {
770                Some(h) => h,
771                None => continue,
772            };
773            if let Some(bus) = self.event_bus.as_ref() {
774                if let Some(sub) = self.fan_in_subs.remove(&key) {
775                    bus.unsubscribe(sub);
776                }
777            } else {
778                self.fan_in_subs.remove(&key);
779            }
780            // Best-effort graceful shutdown. We don't fail the
781            // restart on shutdown errors -- the wire path may
782            // already be dead (the very reason the user asked
783            // for a restart).
784            let _ = old_handle.shutdown().await;
785
786            self.logger.log(
787                None,
788                LogLevel::Info,
789                LogSource::Client,
790                format!(
791                    "supervisor: restarting {} in workspace {}",
792                    server_id,
793                    workspace.display(),
794                ),
795            );
796
797            let new_handle = actor::spawn(
798                (*config).clone(),
799                workspace.clone(),
800                self.logger.clone(),
801                self.apply_edit_bus.clone(),
802                self.configuration_bus.clone(),
803                self.show_document_bus.clone(),
804                self.show_message_request_bus.clone(),
805                self.event_bus.clone(),
806            )
807            .await?;
808
809            let rx = new_handle.subscribe_diagnostics();
810            tokio::spawn(pump_diagnostics(self.diagnostics.clone(), rx));
811            self.actors.insert(key.clone(), new_handle.clone());
812            if let Some(bus) = self.event_bus.clone() {
813                let sub = crate::fan_in::spawn(new_handle.clone(), bus);
814                self.fan_in_subs.insert(key.clone(), sub);
815            }
816
817            // Step 5: replay didOpen for every URI bound to this
818            // key. Empty text -- the next edit (or an explicit
819            // App re-send) repopulates the server's mirror.
820            let uris_for_key: Vec<Uri> = self
821                .attachments
822                .iter()
823                .filter_map(|(uri, keys)| {
824                    if keys.contains(&key) {
825                        Some(uri.clone())
826                    } else {
827                        None
828                    }
829                })
830                .collect();
831            for uri in &uris_for_key {
832                if let Err(e) =
833                    new_handle.open_doc(uri.clone(), config.language_id.clone(), String::new())
834                {
835                    self.logger.log(
836                        None,
837                        LogLevel::Warn,
838                        LogSource::Client,
839                        format!(
840                            "supervisor: restart replay didOpen for {} failed: {}",
841                            uri.as_str(),
842                            e,
843                        ),
844                    );
845                }
846            }
847            replayed_uris.extend(uris_for_key);
848            respawned.push(key);
849        }
850
851        // Step 6: record this restart.
852        self.restart_history
853            .entry(server_id.to_string())
854            .or_default()
855            .push_back(now);
856
857        Ok(RestartReport {
858            server_id: server_id.to_string(),
859            respawned,
860            replayed_uris,
861        })
862    }
863
864    /// Detach every buffer + drop every actor. Used at editor
865    /// exit. Each attached buffer's `didClose` is fired.
866    pub async fn shutdown(&mut self) -> LspResult<()> {
867        // Snapshot URIs first to avoid borrow issues.
868        let uris: Vec<Uri> = self.attachments.keys().cloned().collect();
869        for uri in uris {
870            let _ = self.close_buffer(&uri);
871        }
872        // Drop the per-actor fan-in subscriptions before
873        // dropping the actors -- otherwise the bus would briefly
874        // hold a sender pointing at an actor that is on its way
875        // out.
876        if let Some(bus) = self.event_bus.as_ref() {
877            for (_key, sub) in self.fan_in_subs.drain() {
878                bus.unsubscribe(sub);
879            }
880        } else {
881            self.fan_in_subs.clear();
882        }
883        // Then shut down each actor gracefully.
884        let actors: Vec<ServerHandle> = self.actors.values().cloned().collect();
885        self.actors.clear();
886        for handle in actors {
887            let _ = handle.shutdown().await;
888        }
889        self.logger.log(
890            None,
891            LogLevel::Info,
892            LogSource::Client,
893            "supervisor: shutdown complete",
894        );
895        Ok(())
896    }
897}
898
899impl LspSupervisor {
900    /// Snapshot the current attachment + actor state for the
901    /// handle's `ArcSwap`. Pre-resolves attachments to
902    /// `Vec<ServerHandle>` so `LspSupervisorHandle::servers_for`
903    /// is one map lookup + one Vec clone -- no second walk.
904    pub(crate) fn build_snapshot(&self) -> SupervisorSnapshot {
905        let attachments = self
906            .attachments
907            .iter()
908            .map(|(uri, keys)| {
909                let handles = keys
910                    .iter()
911                    .filter_map(|k| self.actors.get(k).cloned())
912                    .collect::<Vec<_>>();
913                (uri.clone(), handles)
914            })
915            .collect();
916        SupervisorSnapshot {
917            configs: self.configs.clone(),
918            actors: self.actors.clone(),
919            attachments,
920        }
921    }
922
923    /// Consume the configured supervisor and start the supervisor
924    /// task on `runtime_handle`. Returns the
925    /// [`LspSupervisorHandle`] App-side code uses for the rest of
926    /// the editor's lifetime.
927    ///
928    /// After this call, all mutating operations route through the
929    /// returned handle's mailbox; reads come from the wait-free
930    /// `ArcSwap<SupervisorSnapshot>`. The `LspSupervisor` itself
931    /// is owned exclusively by the spawned task -- no `Arc`, no
932    /// `Mutex`, no contention possible with the UI thread.
933    ///
934    /// Configuration calls (`set_event_bus`, `set_apply_edit_bus`,
935    /// `set_configuration_bus`, `add_config`, `set_configs`) must
936    /// happen *before* `spawn` -- the post-spawn handle does not
937    /// expose them. This matches the editor's actual lifecycle:
938    /// the App configures the supervisor at startup before any
939    /// buffer opens.
940    ///
941    /// `runtime_handle` is mandatory: the supervisor's command-
942    /// mailbox semantics only make sense against a live tokio
943    /// task, and a missing runtime is a programming error.
944    /// Production callers (`App::new` →
945    /// `build_lsp_subsystem`) pass the editor's shared LSP
946    /// runtime handle (`runtime::lsp_runtime()`); tests that
947    /// exercise the write path build an ad-hoc runtime via
948    /// `tokio::runtime::Builder` and pass its handle. Earlier
949    /// revisions used `tokio::runtime::Handle::try_current()`
950    /// with a silent fallback that dropped `cmd_rx`; that path
951    /// surfaced as `LspError::ActorGone` on every write whenever
952    /// the caller didn't happen to be inside a tokio context,
953    /// which proved to be a real footgun (e.g. `App::new` runs
954    /// before `runtime::run` has entered any context).
955    pub fn spawn(self, runtime_handle: &tokio::runtime::Handle) -> LspSupervisorHandle {
956        let initial_snapshot = Arc::new(self.build_snapshot());
957        let snapshot_cell = Arc::new(ArcSwap::from(initial_snapshot));
958        let diagnostics = self.diagnostics.clone();
959        let logger = self.logger.clone();
960        let (cmd_tx, cmd_rx) = mpsc::unbounded_channel::<SupervisorCmd>();
961        let snapshot_for_task = snapshot_cell.clone();
962        // 4.4.d: bridge `LspActorExited` events from the editor
963        // bus to a tokio mpsc the supervisor task selects on.
964        // None when `set_event_bus` was never called (tests that
965        // skip event-bus wiring; pre-config sequences); the
966        // supervisor task tolerates a never-firing receiver.
967        let exit_rx = self.event_bus.as_ref().map(|bus| {
968            let (tx, rx) = mpsc::unbounded_channel::<crate::events::LspActorExited>();
969            bus.subscribe_typed(tx);
970            rx
971        });
972        runtime_handle.spawn(supervisor_main(self, cmd_rx, exit_rx, snapshot_for_task));
973        LspSupervisorHandle {
974            snapshot: snapshot_cell,
975            cmd_tx,
976            diagnostics,
977            logger,
978        }
979    }
980}
981
982impl std::fmt::Debug for LspSupervisor {
983    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
984        f.debug_struct("LspSupervisor")
985            .field("configs", &self.configs.len())
986            .field("actors", &self.actors.len())
987            .field("attached_buffers", &self.attachments.len())
988            .finish_non_exhaustive()
989    }
990}
991
992// ---------------------------------------------------------------
993// Snapshot + handle (the public, post-spawn API).
994//
995// The audit (docs/dev/architecture/lsp-architecture.md §11 + post-fix audit) found
996// that 14+ App methods called `try_lock` on
997// `Arc<tokio::sync::Mutex<LspSupervisor>>` from the UI thread
998// (modeline render, Insert-mode trigger probes, etc.) and silently
999// dropped work whenever an async path -- typically `:e <path>`
1000// holding the mutex across the LSP `initialize` handshake -- was
1001// in flight. Same class of bug as the pre-fan-in LSP-edit drop.
1002//
1003// Resolution: split the supervisor's surface in two.
1004//
1005//   - `SupervisorSnapshot` carries everything readers need
1006//     (configs, actors, pre-resolved per-URI attachments). It
1007//     lives in an `ArcSwap` cell; readers do one wait-free
1008//     `load_full()` and clone what they need from the returned
1009//     `Arc`.
1010//   - `LspSupervisorHandle` is what the App holds. Reads return
1011//     directly from the snapshot (lock-free); writes send a
1012//     typed `SupervisorCmd` on the task's mailbox.
1013//   - The supervisor task owns `LspSupervisor` exclusively; it
1014//     processes cmds, mutates state, then publishes a fresh
1015//     snapshot via `ArcSwap::store`. No shared mutex anywhere.
1016//
1017// This makes "UI thread takes no LSP-related lock" a structural
1018// guarantee, not a discipline -- you can't write the contended
1019// pattern even by accident.
1020// ---------------------------------------------------------------
1021
1022/// Wait-free read view of supervisor state.
1023///
1024/// Built by the supervisor's private `build_snapshot` helper after every cmd
1025/// that mutates state; held in [`LspSupervisorHandle::snapshot`]
1026/// inside an `ArcSwap` cell. Readers `load_full()` and clone
1027/// per-field as needed. All fields are public so callers can
1028/// project freely without growing the handle's API.
1029#[derive(Debug, Clone, Default)]
1030pub struct SupervisorSnapshot {
1031    /// Curated config registry. Order is preference order: when
1032    /// multiple configs match a path, all of them attach.
1033    pub configs: Vec<Arc<ServerConfig>>,
1034    /// Per-(workspace, server-id) live actor handles. Cloning
1035    /// is cheap (`ServerHandle` is `Arc<HandleInner>` internally).
1036    pub actors: HashMap<ActorKey, ServerHandle>,
1037    /// Per-URI attachments, pre-resolved to `ServerHandle`s so
1038    /// `servers_for(uri)` is one map lookup + one small clone.
1039    pub attachments: HashMap<Uri, Vec<ServerHandle>>,
1040}
1041
1042/// Commands the [`LspSupervisorHandle`] sends to the supervisor
1043/// task. Every state-mutating operation flows through here so the
1044/// task is the single writer to `LspSupervisor` state.
1045enum SupervisorCmd {
1046    /// Open a buffer + spawn / reuse matching actors.
1047    OpenBuffer {
1048        path: PathBuf,
1049        text: String,
1050        reply: oneshot::Sender<LspResult<Vec<ServerHandle>>>,
1051    },
1052    /// Attach a buffer to a pre-built [`ServerHandle`] (tests +
1053    /// custom transports).
1054    AttachHandle {
1055        uri: Uri,
1056        workspace_root: PathBuf,
1057        server_id: String,
1058        language_id: String,
1059        text: String,
1060        handle: ServerHandle,
1061        reply: oneshot::Sender<LspResult<()>>,
1062    },
1063    /// Detach a buffer + emit `didClose` per server. Reply
1064    /// resolves once the supervisor has processed (used by tests
1065    /// that need ordering).
1066    CloseBuffer {
1067        uri: Uri,
1068        reply: oneshot::Sender<LspResult<()>>,
1069    },
1070    /// Force-flush queued changes for `uri`.
1071    Flush { uri: Uri },
1072    /// Force-flush every URI; reply resolves after all flushes
1073    /// have been issued (used by `shutdown`).
1074    FlushAll { reply: oneshot::Sender<()> },
1075    /// Drive a single edit through the per-actor mailbox. Used by
1076    /// tests + admin tooling; production edits flow through the
1077    /// event-bus fan-in (see `lattice_lsp::fan_in`).
1078    RecordEdit { uri: Uri, edit: Edit },
1079    /// 4.4.d: force-restart every actor with id `server_id`.
1080    /// Replays `didOpen` for every URI that was attached.
1081    Restart {
1082        server_id: String,
1083        reply: oneshot::Sender<LspResult<RestartReport>>,
1084    },
1085    /// Editor exit: close every buffer, drop fan-ins, shut down
1086    /// every actor. Reply resolves after the supervisor task
1087    /// itself is exiting (next iteration drops `cmd_rx`).
1088    Shutdown {
1089        reply: oneshot::Sender<LspResult<()>>,
1090    },
1091}
1092
1093/// Editor-facing handle to the LSP subsystem. Cheap to clone
1094/// (every field is `Arc`-shaped internally); the App holds one
1095/// instance and shares it with helpers that need to read or
1096/// mutate supervisor state.
1097///
1098/// **Reads are wait-free.** [`Self::servers_for`],
1099/// [`Self::running_actors`], [`Self::configs`], and the count
1100/// helpers all go through `ArcSwap::load`. The UI thread can
1101/// call them on the keystroke / render path without ever
1102/// blocking.
1103///
1104/// **Writes are mailbox-routed.** [`Self::open_buffer`],
1105/// [`Self::attach_handle`], [`Self::shutdown`], and
1106/// [`Self::flush_all`] are async (await processing); the
1107/// fire-and-forget variants ([`Self::close_buffer`],
1108/// [`Self::flush`], [`Self::record_edit`]) send the cmd and
1109/// return immediately, with the effect observable on the next
1110/// snapshot publish.
1111#[derive(Clone)]
1112pub struct LspSupervisorHandle {
1113    snapshot: Arc<ArcSwap<SupervisorSnapshot>>,
1114    cmd_tx: mpsc::UnboundedSender<SupervisorCmd>,
1115    diagnostics: DiagnosticsLayer,
1116    logger: LspLogger,
1117}
1118
1119/// Placeholder handle for `Editor::default()` / headless
1120/// test scaffolding. The receiver is dropped immediately, so
1121/// any cmd sent through this handle silently no-ops (the
1122/// caller's reply oneshot drops without being completed,
1123/// which `LspSupervisor::send_with_reply` reports as a
1124/// "supervisor gone" error). Real construction goes through
1125/// `LspSupervisor::spawn`. Production code overwrites this
1126/// in `boot.rs` before any cmd flows.
1127impl Default for LspSupervisorHandle {
1128    fn default() -> Self {
1129        let (cmd_tx, _cmd_rx) = mpsc::unbounded_channel();
1130        let logger = LspLogger::default();
1131        Self {
1132            snapshot: Arc::new(ArcSwap::from(Arc::new(SupervisorSnapshot::default()))),
1133            cmd_tx,
1134            diagnostics: DiagnosticsLayer::default(),
1135            logger,
1136        }
1137    }
1138}
1139
1140impl LspSupervisorHandle {
1141    // ----- read API (wait-free) --------------------------------
1142
1143    /// One wait-free `ArcSwap::load_full` returning the current
1144    /// snapshot. Callers that need many fields project from the
1145    /// returned `Arc`; one-off readers should prefer the
1146    /// projection helpers below to avoid retaining the snapshot
1147    /// across awaits.
1148    pub fn snapshot(&self) -> Arc<SupervisorSnapshot> {
1149        self.snapshot.load_full()
1150    }
1151
1152    /// Every server attached to `uri`. Empty when the URI is not
1153    /// open or no config matched its path. Cheap clone.
1154    pub fn servers_for(&self, uri: &Uri) -> Vec<ServerHandle> {
1155        self.snapshot
1156            .load()
1157            .attachments
1158            .get(uri)
1159            .cloned()
1160            .unwrap_or_default()
1161    }
1162
1163    /// Snapshot of every running actor's `(key, handle)` pair.
1164    pub fn running_actors(&self) -> Vec<(ActorKey, ServerHandle)> {
1165        self.snapshot
1166            .load()
1167            .actors
1168            .iter()
1169            .map(|(k, h)| (k.clone(), h.clone()))
1170            .collect()
1171    }
1172
1173    /// Snapshot of every running actor's `ServerHandle`. Used by
1174    /// workspace-scoped LSP requests (e.g. `workspace/symbol`)
1175    /// that fan out across every server, not just servers
1176    /// attached to one buffer.
1177    pub fn all_running_handles(&self) -> Vec<ServerHandle> {
1178        self.snapshot.load().actors.values().cloned().collect()
1179    }
1180
1181    /// Number of buffers currently attached to the actor at `key`.
1182    /// Used by `:lsp-server-log` for the picker margin.
1183    pub fn buffer_count_for(&self, key: &ActorKey) -> usize {
1184        let snap = self.snapshot.load();
1185        let handle = match snap.actors.get(key) {
1186            Some(h) => h,
1187            None => return 0,
1188        };
1189        snap.attachments
1190            .values()
1191            .filter(|handles| handles.iter().any(|h| h.server_id() == handle.server_id()))
1192            .count()
1193    }
1194
1195    /// Curated config registry (preference-ordered). Cheap clone.
1196    pub fn configs(&self) -> Vec<Arc<ServerConfig>> {
1197        self.snapshot.load().configs.clone()
1198    }
1199
1200    /// True iff at least one configured server's `file_patterns`
1201    /// matches `path`. M.5.2 uses this from the App's
1202    /// `MajorEntered` hook to decide whether to auto-activate
1203    /// `lsp-mode` on the buffer.
1204    pub fn has_server_for_path(&self, path: &Path) -> bool {
1205        self.snapshot
1206            .load()
1207            .configs
1208            .iter()
1209            .any(|c| matches_any_pattern(path, &c.file_patterns))
1210    }
1211
1212    /// Number of currently-attached buffers.
1213    pub fn attached_buffer_count(&self) -> usize {
1214        self.snapshot.load().attachments.len()
1215    }
1216
1217    /// Number of currently-running actors.
1218    pub fn running_actor_count(&self) -> usize {
1219        self.snapshot.load().actors.len()
1220    }
1221
1222    /// Borrow the shared logger. Cloning is cheap (`LspLogger`
1223    /// wraps an `Arc` internally) but most callers just want a
1224    /// borrow to call `.log(...)`.
1225    pub fn logger(&self) -> &LspLogger {
1226        &self.logger
1227    }
1228
1229    /// Borrow the shared diagnostics layer.
1230    pub fn diagnostics(&self) -> &DiagnosticsLayer {
1231        &self.diagnostics
1232    }
1233
1234    // ----- write API (async; await mailbox processing) ---------
1235
1236    /// Open a buffer. Returns the list of attached
1237    /// `ServerHandle`s; empty if no config matched the path.
1238    /// May spawn new actors (each one pays the LSP `initialize`
1239    /// handshake cost). The supervisor task runs the open;
1240    /// awaiting here parks the caller while it does, but does
1241    /// NOT block any *other* read of the supervisor.
1242    pub async fn open_buffer(&self, path: PathBuf, text: String) -> LspResult<Vec<ServerHandle>> {
1243        let (reply, rx) = oneshot::channel();
1244        self.cmd_tx
1245            .send(SupervisorCmd::OpenBuffer { path, text, reply })
1246            .map_err(|_| LspError::ActorGone)?;
1247        rx.await.map_err(|_| LspError::ActorGone)?
1248    }
1249
1250    /// Attach a buffer to a pre-built handle (tests + custom
1251    /// transports).
1252    pub async fn attach_handle(
1253        &self,
1254        uri: Uri,
1255        workspace_root: PathBuf,
1256        server_id: String,
1257        language_id: String,
1258        text: String,
1259        handle: ServerHandle,
1260    ) -> LspResult<()> {
1261        let (reply, rx) = oneshot::channel();
1262        self.cmd_tx
1263            .send(SupervisorCmd::AttachHandle {
1264                uri,
1265                workspace_root,
1266                server_id,
1267                language_id,
1268                text,
1269                handle,
1270                reply,
1271            })
1272            .map_err(|_| LspError::ActorGone)?;
1273        rx.await.map_err(|_| LspError::ActorGone)?
1274    }
1275
1276    // ----- write API (fire-and-forget; effect via snapshot) ----
1277
1278    /// Close a buffer. Fire-and-forget: the supervisor task
1279    /// processes the cmd asynchronously, fires `didClose` per
1280    /// attached server, and publishes a fresh snapshot. Tests
1281    /// that need to assert post-close state should call
1282    /// [`Self::close_buffer_ack`].
1283    pub fn close_buffer(&self, uri: Uri) {
1284        let (reply, _rx) = oneshot::channel();
1285        let _ = self.cmd_tx.send(SupervisorCmd::CloseBuffer { uri, reply });
1286    }
1287
1288    /// Close a buffer + await acknowledgement. Used by tests +
1289    /// shutdown paths that need to observe the next snapshot.
1290    pub async fn close_buffer_ack(&self, uri: Uri) -> LspResult<()> {
1291        let (reply, rx) = oneshot::channel();
1292        self.cmd_tx
1293            .send(SupervisorCmd::CloseBuffer { uri, reply })
1294            .map_err(|_| LspError::ActorGone)?;
1295        rx.await.map_err(|_| LspError::ActorGone)?
1296    }
1297
1298    /// Force-flush queued changes for `uri`. Fire-and-forget.
1299    pub fn flush(&self, uri: Uri) {
1300        let _ = self.cmd_tx.send(SupervisorCmd::Flush { uri });
1301    }
1302
1303    /// Force-flush every URI; awaits the supervisor task's ack.
1304    pub async fn flush_all(&self) -> LspResult<()> {
1305        let (reply, rx) = oneshot::channel();
1306        self.cmd_tx
1307            .send(SupervisorCmd::FlushAll { reply })
1308            .map_err(|_| LspError::ActorGone)?;
1309        rx.await.map_err(|_| LspError::ActorGone)?;
1310        Ok(())
1311    }
1312
1313    /// Drive an edit through the supervisor (test / admin path;
1314    /// production edits ride the event-bus fan-in).
1315    pub fn record_edit(&self, uri: Uri, edit: Edit) {
1316        let _ = self.cmd_tx.send(SupervisorCmd::RecordEdit { uri, edit });
1317    }
1318
1319    /// 4.4.d: force-restart every actor with id `server_id`.
1320    /// Subject to the per-server-id restart-history backoff:
1321    /// more than `MAX_RESTARTS` inside `RESTART_WINDOW`
1322    /// returns an error; the caller surfaces the cooldown
1323    /// message via `set_message`.
1324    pub async fn restart_server(&self, server_id: String) -> LspResult<RestartReport> {
1325        let (reply, rx) = oneshot::channel();
1326        self.cmd_tx
1327            .send(SupervisorCmd::Restart { server_id, reply })
1328            .map_err(|_| LspError::ActorGone)?;
1329        rx.await.map_err(|_| LspError::ActorGone)?
1330    }
1331
1332    /// Editor exit: close every buffer, drop fan-ins, shut down
1333    /// every actor. Awaits the supervisor task's final ack.
1334    pub async fn shutdown(&self) -> LspResult<()> {
1335        let (reply, rx) = oneshot::channel();
1336        self.cmd_tx
1337            .send(SupervisorCmd::Shutdown { reply })
1338            .map_err(|_| LspError::ActorGone)?;
1339        rx.await.map_err(|_| LspError::ActorGone)?
1340    }
1341}
1342
1343impl std::fmt::Debug for LspSupervisorHandle {
1344    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1345        let snap = self.snapshot.load();
1346        f.debug_struct("LspSupervisorHandle")
1347            .field("configs", &snap.configs.len())
1348            .field("actors", &snap.actors.len())
1349            .field("attached_buffers", &snap.attachments.len())
1350            .finish_non_exhaustive()
1351    }
1352}
1353
1354/// The supervisor task. Owns `LspSupervisor` exclusively;
1355/// processes `SupervisorCmd`s in FIFO order; publishes a fresh
1356/// snapshot via `ArcSwap::store` after every state mutation.
1357///
1358/// Exits when `cmd_rx` returns `None` (every handle dropped) or
1359/// after a successful `Shutdown` cmd. The mailbox is unbounded;
1360/// publishers never block on a full queue, and the supervisor
1361/// always drains in order so e.g. `OpenBuffer` -> `CloseBuffer`
1362/// is processed open-first regardless of timing.
1363async fn supervisor_main(
1364    mut state: LspSupervisor,
1365    mut cmd_rx: mpsc::UnboundedReceiver<SupervisorCmd>,
1366    mut exit_rx: Option<mpsc::UnboundedReceiver<crate::events::LspActorExited>>,
1367    snapshot: Arc<ArcSwap<SupervisorSnapshot>>,
1368) {
1369    loop {
1370        let cmd = tokio::select! {
1371            biased;
1372            // Drain exit events first so a crash-storm can't
1373            // get starved behind a long sequence of edits.
1374            Some(exit) = async {
1375                match exit_rx.as_mut() {
1376                    Some(rx) => rx.recv().await,
1377                    None => std::future::pending().await,
1378                }
1379            } => {
1380                handle_actor_exit(&mut state, &snapshot, exit).await;
1381                continue;
1382            }
1383            cmd = cmd_rx.recv() => match cmd {
1384                Some(c) => c,
1385                None => break,
1386            },
1387        };
1388        match cmd {
1389            SupervisorCmd::OpenBuffer { path, text, reply } => {
1390                let result = state.open_buffer(path, text).await;
1391                snapshot.store(Arc::new(state.build_snapshot()));
1392                let _ = reply.send(result);
1393            }
1394            SupervisorCmd::AttachHandle {
1395                uri,
1396                workspace_root,
1397                server_id,
1398                language_id,
1399                text,
1400                handle,
1401                reply,
1402            } => {
1403                let result =
1404                    state.attach_handle(uri, workspace_root, server_id, language_id, text, handle);
1405                snapshot.store(Arc::new(state.build_snapshot()));
1406                let _ = reply.send(result);
1407            }
1408            SupervisorCmd::CloseBuffer { uri, reply } => {
1409                let result = state.close_buffer(&uri);
1410                snapshot.store(Arc::new(state.build_snapshot()));
1411                let _ = reply.send(result);
1412            }
1413            SupervisorCmd::Flush { uri } => {
1414                // `flush` returns Ok unless the URI is unknown;
1415                // either way no state changes, so no snapshot
1416                // republish needed.
1417                let _ = state.flush(&uri);
1418            }
1419            SupervisorCmd::FlushAll { reply } => {
1420                let _ = state.flush_all();
1421                let _ = reply.send(());
1422            }
1423            SupervisorCmd::RecordEdit { uri, edit } => {
1424                let _ = state.record_edit(&uri, &edit);
1425            }
1426            SupervisorCmd::Restart { server_id, reply } => {
1427                let result = state.restart_server(&server_id).await;
1428                snapshot.store(Arc::new(state.build_snapshot()));
1429                let _ = reply.send(result);
1430            }
1431            SupervisorCmd::Shutdown { reply } => {
1432                let result = state.shutdown().await;
1433                snapshot.store(Arc::new(state.build_snapshot()));
1434                let _ = reply.send(result);
1435                // Drain anything still in the queue (likely
1436                // empty), then exit. Closing cmd_rx here would
1437                // be racey -- the runtime is the one dropping
1438                // senders.
1439                break;
1440            }
1441        }
1442    }
1443}
1444
1445/// 4.4.d: handle a typed `LspActorExited`. Clean exits are
1446/// already accounted for (the supervisor itself triggered the
1447/// shutdown via `:lsp-restart` or `shutdown`). Unexpected
1448/// exits drive a restart through the same path the user-
1449/// facing `restart_server` uses, including the backoff gate
1450/// -- a server in a crash loop hits `MAX_RESTARTS` and the
1451/// supervisor stops attempting until the window slides past.
1452async fn handle_actor_exit(
1453    state: &mut LspSupervisor,
1454    snapshot: &Arc<ArcSwap<SupervisorSnapshot>>,
1455    exit: crate::events::LspActorExited,
1456) {
1457    use crate::events::LspActorExitReason;
1458    if matches!(exit.reason, LspActorExitReason::Clean) {
1459        // Routine shutdown -- nothing to do. The
1460        // restart_server / shutdown path already removed the
1461        // actor + republished snapshot.
1462        return;
1463    }
1464    let server_id = exit.server_id.as_ref();
1465    // Drop the dead actor (and its fan-in) from state before
1466    // re-spawning so `restart_server`'s view is consistent.
1467    // The actor's `cmd_tx` is closed already (its task
1468    // returned); a stale entry would mislead `running_actors`
1469    // and the snapshot.
1470    let stale_keys: Vec<ActorKey> = state
1471        .actors
1472        .keys()
1473        .filter(|(_, id)| id == server_id)
1474        .cloned()
1475        .collect();
1476    for key in &stale_keys {
1477        state.actors.remove(key);
1478        if let Some(bus) = state.event_bus.as_ref() {
1479            if let Some(sub) = state.fan_in_subs.remove(key) {
1480                bus.unsubscribe(sub);
1481            }
1482        } else {
1483            state.fan_in_subs.remove(key);
1484        }
1485    }
1486    state.logger.log(
1487        None,
1488        LogLevel::Warn,
1489        LogSource::Client,
1490        format!(
1491            "supervisor: {} exited unexpectedly; attempting auto-restart",
1492            server_id,
1493        ),
1494    );
1495    // Re-stash workspace info: restart_server walks actors,
1496    // but those are gone now. Synthesise the calls by stashing
1497    // workspaces from the stale_keys list -- we can't reuse
1498    // restart_server here because it expects actors[key] to
1499    // still exist. Inline the equivalent spawn-and-replay
1500    // logic without the gate's "no running actor" branch:
1501    // the actor *had* been running, it just died.
1502    let history = state
1503        .restart_history
1504        .entry(server_id.to_string())
1505        .or_default();
1506    let now = Instant::now();
1507    while let Some(front) = history.front().copied()
1508        && now.duration_since(front) >= RESTART_WINDOW
1509    {
1510        history.pop_front();
1511    }
1512    if history.len() >= MAX_RESTARTS {
1513        state.logger.log(
1514            None,
1515            LogLevel::Error,
1516            LogSource::Client,
1517            format!(
1518                "supervisor: {} crashed {} times in {}s; giving up auto-restart \
1519                 (use :lsp-restart {} after the window slides past)",
1520                server_id,
1521                history.len(),
1522                RESTART_WINDOW.as_secs(),
1523                server_id,
1524            ),
1525        );
1526        snapshot.store(Arc::new(state.build_snapshot()));
1527        return;
1528    }
1529    let backoff = compute_restart_backoff(history.len());
1530    if !backoff.is_zero() {
1531        tokio::time::sleep(backoff).await;
1532    }
1533    let config = match state.configs.iter().find(|c| c.id == server_id).cloned() {
1534        Some(c) => c,
1535        None => {
1536            state.logger.log(
1537                None,
1538                LogLevel::Warn,
1539                LogSource::Client,
1540                format!(
1541                    "supervisor: {} exited but its config is gone; cannot restart",
1542                    server_id
1543                ),
1544            );
1545            snapshot.store(Arc::new(state.build_snapshot()));
1546            return;
1547        }
1548    };
1549    for key in stale_keys {
1550        let workspace = key.0.clone();
1551        let new_handle = match actor::spawn(
1552            (*config).clone(),
1553            workspace.clone(),
1554            state.logger.clone(),
1555            state.apply_edit_bus.clone(),
1556            state.configuration_bus.clone(),
1557            state.show_document_bus.clone(),
1558            state.show_message_request_bus.clone(),
1559            state.event_bus.clone(),
1560        )
1561        .await
1562        {
1563            Ok(h) => h,
1564            Err(e) => {
1565                state.logger.log(
1566                    None,
1567                    LogLevel::Warn,
1568                    LogSource::Client,
1569                    format!(
1570                        "supervisor: auto-restart spawn failed for {} ({}): {}",
1571                        server_id,
1572                        workspace.display(),
1573                        e,
1574                    ),
1575                );
1576                continue;
1577            }
1578        };
1579        let rx = new_handle.subscribe_diagnostics();
1580        tokio::spawn(pump_diagnostics(state.diagnostics.clone(), rx));
1581        state.actors.insert(key.clone(), new_handle.clone());
1582        if let Some(bus) = state.event_bus.clone() {
1583            let sub = crate::fan_in::spawn(new_handle.clone(), bus);
1584            state.fan_in_subs.insert(key.clone(), sub);
1585        }
1586        let uris_for_key: Vec<Uri> = state
1587            .attachments
1588            .iter()
1589            .filter_map(|(uri, keys)| {
1590                if keys.contains(&key) {
1591                    Some(uri.clone())
1592                } else {
1593                    None
1594                }
1595            })
1596            .collect();
1597        for uri in uris_for_key {
1598            if let Err(e) =
1599                new_handle.open_doc(uri.clone(), config.language_id.clone(), String::new())
1600            {
1601                state.logger.log(
1602                    None,
1603                    LogLevel::Warn,
1604                    LogSource::Client,
1605                    format!(
1606                        "supervisor: auto-restart replay didOpen for {} failed: {}",
1607                        uri.as_str(),
1608                        e,
1609                    ),
1610                );
1611            }
1612        }
1613    }
1614    state
1615        .restart_history
1616        .entry(server_id.to_string())
1617        .or_default()
1618        .push_back(now);
1619    snapshot.store(Arc::new(state.build_snapshot()));
1620}
1621
1622/// Match `path` against any of `patterns`. Supports `*.<ext>`
1623/// (extension match) and bare basename equality (e.g.
1624/// `Cargo.toml`). More elaborate globs are deferred until a
1625/// real use case appears.
1626pub(crate) fn matches_any_pattern(path: &Path, patterns: &[String]) -> bool {
1627    patterns.iter().any(|p| matches_pattern(path, p))
1628}
1629
1630fn matches_pattern(path: &Path, pattern: &str) -> bool {
1631    if let Some(ext) = pattern.strip_prefix("*.") {
1632        return path
1633            .extension()
1634            .is_some_and(|e| e.eq_ignore_ascii_case(std::ffi::OsStr::new(ext)));
1635    }
1636    // Bare pattern: basename equality.
1637    path.file_name()
1638        .is_some_and(|n| n.eq_ignore_ascii_case(std::ffi::OsStr::new(pattern)))
1639}
1640
1641#[cfg(test)]
1642mod tests {
1643    use super::*;
1644
1645    #[test]
1646    fn matches_pattern_handles_extension_globs() {
1647        let path = std::path::PathBuf::from("/tmp/x.rs");
1648        assert!(matches_pattern(&path, "*.rs"));
1649        assert!(!matches_pattern(&path, "*.py"));
1650        // Case-insensitive on extension.
1651        assert!(matches_pattern(&std::path::PathBuf::from("X.RS"), "*.rs"));
1652    }
1653
1654    #[test]
1655    fn matches_pattern_handles_basename_equality() {
1656        let path = std::path::PathBuf::from("/tmp/Cargo.toml");
1657        assert!(matches_pattern(&path, "Cargo.toml"));
1658        assert!(!matches_pattern(&path, "package.json"));
1659    }
1660
1661    #[test]
1662    fn matches_any_pattern_walks_list() {
1663        let path = std::path::PathBuf::from("/tmp/main.go");
1664        let patterns = vec!["*.rs".to_string(), "*.go".to_string(), "go.mod".to_string()];
1665        assert!(matches_any_pattern(&path, &patterns));
1666    }
1667
1668    /// 4.4.d: the backoff growth doubles per prior restart, but
1669    /// stays bounded so a 10-restart-storm doesn't wedge for
1670    /// hours. The cap is the load-bearing property; the exact
1671    /// curve below it isn't.
1672    #[test]
1673    fn restart_backoff_doubles_then_caps() {
1674        assert_eq!(compute_restart_backoff(0), Duration::ZERO);
1675        assert_eq!(compute_restart_backoff(1), Duration::from_millis(500));
1676        assert_eq!(compute_restart_backoff(2), Duration::from_millis(1000));
1677        // Saturates well before pathological values.
1678        let big = compute_restart_backoff(20);
1679        assert_eq!(big, RESTART_BACKOFF_MAX);
1680    }
1681
1682    #[tokio::test]
1683    async fn restart_no_running_actor_returns_error() {
1684        // No actor with `rust` id is running; the supervisor
1685        // must refuse rather than silently spawn one.
1686        let mut sup = LspSupervisor::new(LspLogger::with_defaults());
1687        sup.add_config(ServerConfig::new("rust", "rust-analyzer", "rust"));
1688        let err = sup.restart_server("rust").await.unwrap_err();
1689        let msg = format!("{err}");
1690        assert!(msg.contains("no running actor"), "got: {msg}",);
1691    }
1692
1693    #[tokio::test]
1694    async fn restart_unknown_config_returns_error() {
1695        // Even with restart history clean, asking to restart a
1696        // server with no registered config must fail rather
1697        // than panicking on the `find` below.
1698        let mut sup = LspSupervisor::new(LspLogger::with_defaults());
1699        let err = sup.restart_server("ghost").await.unwrap_err();
1700        let msg = format!("{err}");
1701        // We hit the "no running actor" branch first (since
1702        // both checks fail without an actor); that's the
1703        // correct user-facing message either way.
1704        assert!(
1705            msg.contains("no running actor") || msg.contains("no config"),
1706            "got: {msg}",
1707        );
1708    }
1709
1710    #[test]
1711    fn supervisor_has_empty_state_at_construction() {
1712        let logger = LspLogger::with_defaults();
1713        let sup = LspSupervisor::new(logger);
1714        assert_eq!(sup.attached_buffer_count(), 0);
1715        assert_eq!(sup.running_actor_count(), 0);
1716        assert_eq!(sup.configs().len(), 0);
1717        assert_eq!(sup.diagnostics().count(), 0);
1718    }
1719
1720    #[test]
1721    fn add_config_registers_in_order() {
1722        let mut sup = LspSupervisor::new(LspLogger::with_defaults());
1723        sup.add_config(ServerConfig::new("rust", "rust-analyzer", "rust"));
1724        sup.add_config(ServerConfig::new("python", "pyright", "python"));
1725        assert_eq!(sup.configs().len(), 2);
1726        assert_eq!(sup.configs()[0].id, "rust");
1727        assert_eq!(sup.configs()[1].id, "python");
1728    }
1729
1730    #[test]
1731    fn set_configs_replaces_registry() {
1732        let mut sup = LspSupervisor::new(LspLogger::with_defaults());
1733        sup.add_config(ServerConfig::new("rust", "rust-analyzer", "rust"));
1734        sup.set_configs([ServerConfig::new("go", "gopls", "go")]);
1735        assert_eq!(sup.configs().len(), 1);
1736        assert_eq!(sup.configs()[0].id, "go");
1737    }
1738}