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}