lattice_lsp/actor.rs
1//! Per-server actor (DESIGN.md §5.4 + §5.7). One tokio task owns
2//! the wire-side state for one (workspace, server-id) pair:
3//!
4//! - The pending-request table (RequestId → oneshot to the
5//! editor-side caller).
6//! - The negotiated [`Capabilities`].
7//! - A monotonic JSON-RPC request-id counter.
8//!
9//! Two helper tasks fan I/O in and out:
10//!
11//! - **read_loop** -- reads [`Message`]s from `LspReader` and
12//! pushes them into an inbound channel.
13//! - **write_loop** -- receives [`Message`]s on an outbound
14//! channel and writes them through `LspWriter`.
15//!
16//! ## Why three tasks
17//!
18//! A single-task design works for a typewriter-pace editor but
19//! collapses under burst loads (a server emitting hundreds of
20//! `$/progress` notifications during indexing while the editor
21//! is also sending `didChange` per keystroke). Splitting reads
22//! and writes onto separate tasks lets the OS schedule them
23//! across cores; the actor task itself stays cheap (no I/O).
24//!
25//! ## Lifecycle
26//!
27//! [`spawn`] runs the initialize handshake before returning a
28//! [`ServerHandle`]. Failure during handshake yields
29//! [`LspError::HandshakeFailed`] and tears the child process
30//! down via `kill_on_drop`.
31//!
32//! [`ServerHandle::shutdown`] runs the LSP shutdown protocol:
33//! `shutdown` request → `exit` notification → wait for child
34//! exit. The actor task exits cleanly and all pending
35//! requests resolve with [`LspError::Cancelled`].
36
37use std::collections::HashMap;
38use std::str::FromStr;
39use std::sync::Arc;
40
41use arc_swap::ArcSwap;
42use lsp_types::{
43 ClientInfo, InitializeParams, InitializeResult, InitializedParams, Uri, WorkspaceFolder,
44};
45use serde_json::Value;
46use tokio::io::{AsyncBufRead, AsyncWrite};
47use tokio::process::{Child, ChildStderr};
48use tokio::select;
49use tokio::sync::{Mutex, mpsc, oneshot};
50
51use crate::capabilities::{self, Capabilities};
52use crate::codec::{LspReader, LspWriter};
53use crate::config::ServerConfig;
54use crate::diagnostics::{DiagnosticEvent, DiagnosticsBus};
55use crate::error::{LspError, LspResult};
56use crate::jsonrpc::{Message, Notification, Request, RequestId, Response};
57use crate::logging::{LogLevel, LogSource, LspLogger};
58use crate::pending::{InvocationId, Pending};
59use crate::transport::ChildTransport;
60
61/// Editor-facing handle to one running language-server actor.
62///
63/// Cheap to clone (`Arc` internally); the editor passes one
64/// around per buffer / pane that talks to this server. Dropping
65/// the last clone closes the mailbox -- the actor sees the
66/// channel close and runs the LSP shutdown sequence on its way
67/// out, so no leak.
68#[derive(Clone)]
69pub struct ServerHandle {
70 inner: Arc<HandleInner>,
71}
72
73struct HandleInner {
74 cmd_tx: mpsc::UnboundedSender<ActorCmd>,
75 /// Negotiated capabilities, published as an
76 /// `ArcSwap<Capabilities>` (4.4.n). The handshake snapshot
77 /// goes in at handshake-time; subsequent
78 /// `client/registerCapability` and
79 /// `client/unregisterCapability` notifications publish new
80 /// snapshots in place. Readers
81 /// ([`ServerHandle::capabilities`]) load a fresh `Arc` per
82 /// call -- lock-free and as cheap as an atomic load.
83 capabilities: Arc<ArcSwap<Capabilities>>,
84 /// Server id, for logs / telemetry.
85 server_id: String,
86 /// Workspace root this actor was spawned against (B'.2).
87 /// Pairs with `server_id` as the canonical
88 /// `(server_id, workspace)` instance key so multi-instance
89 /// setups stay distinct in the per-instance log rings.
90 workspace_root: Arc<std::path::Path>,
91 /// Diagnostics broadcast bus -- subscribers (App, plugins,
92 /// future picker) receive every `publishDiagnostics` from
93 /// this server.
94 diagnostics: DiagnosticsBus,
95 /// LSP-subsystem logger. Cloned from the App's shared
96 /// logger; per-instance records land in the
97 /// `*lsp:<server>:<workspace>*` ring.
98 logger: LspLogger,
99}
100
101impl std::fmt::Debug for ServerHandle {
102 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
103 f.debug_struct("ServerHandle")
104 .field("server_id", &self.inner.server_id)
105 .finish_non_exhaustive()
106 }
107}
108
109impl ServerHandle {
110 /// Snapshot of the current negotiated capabilities. Each
111 /// call returns a fresh `Arc` that includes any dynamic
112 /// registrations the server has issued since handshake
113 /// (4.4.n). Callers that need the capability set to be
114 /// stable for a multi-step decision should bind the
115 /// snapshot to a local variable; subsequent
116 /// `capabilities()` calls see whatever the actor has
117 /// published in the interim.
118 pub fn capabilities(&self) -> Arc<Capabilities> {
119 self.inner.capabilities.load_full()
120 }
121
122 /// Server's stable id (e.g. `"rust"`). Useful for logs and
123 /// for the supervisor to look up `ServerConfig`.
124 pub fn server_id(&self) -> &str {
125 &self.inner.server_id
126 }
127
128 /// Workspace root this actor was spawned against. Pairs with
129 /// `server_id()` as the canonical instance identity (B'.2);
130 /// cheap clone (Arc bump).
131 pub fn workspace_root(&self) -> Arc<std::path::Path> {
132 Arc::clone(&self.inner.workspace_root)
133 }
134
135 /// Build the `(server_id, workspace)` instance key for use
136 /// with [`LspLogger::log`] and the per-instance log ring.
137 pub fn instance(&self) -> crate::logging::InstanceKey {
138 crate::logging::InstanceKey::new(
139 Arc::<str>::from(self.inner.server_id.as_str()),
140 Arc::clone(&self.inner.workspace_root),
141 )
142 }
143
144 /// Subscribe to this server's diagnostics broadcast. Each
145 /// subscriber receives every `publishDiagnostics` event
146 /// after the call -- prior events are not replayed
147 /// (a freshly opened pane re-issues the URIs it cares about
148 /// to the server, which republishes diagnostics for those).
149 ///
150 /// The returned `Receiver` is the standard
151 /// `tokio::sync::broadcast::Receiver`. A lagging consumer
152 /// drops oldest first; reconcile by tracking the latest
153 /// `version` per URI and ignoring events older than the
154 /// editor's view of the doc.
155 pub fn subscribe_diagnostics(&self) -> tokio::sync::broadcast::Receiver<DiagnosticEvent> {
156 self.inner.diagnostics.subscribe()
157 }
158
159 /// True iff at least one subscriber is currently listening.
160 /// Used by tests; production code doesn't need this.
161 pub fn diagnostics_subscriber_count(&self) -> usize {
162 self.inner.diagnostics.receiver_count()
163 }
164
165 /// Borrow the logger this actor emits through. The App
166 /// holds the same `LspLogger` (cloned) and uses it for
167 /// supervisor-side records.
168 pub fn logger(&self) -> &LspLogger {
169 &self.inner.logger
170 }
171
172 /// Send a typed JSON-RPC request and return a [`Pending`]
173 /// resolving to the deserialized response.
174 ///
175 /// `R` must match the server's response shape for `method`;
176 /// a mismatch surfaces as [`LspError::ResponseDecode`].
177 pub fn request<P, R>(&self, method: &str, params: P) -> Pending<R>
178 where
179 P: serde::Serialize,
180 R: serde::de::DeserializeOwned + Send + 'static,
181 {
182 let params_json = match serde_json::to_value(params) {
183 Ok(v) => Some(v),
184 Err(e) => return Pending::ready_err(LspError::ResponseDecode(e)),
185 };
186 let (reply_tx, reply_rx) = oneshot::channel::<LspResult<Value>>();
187 let cmd = ActorCmd::Request {
188 method: method.to_string(),
189 params: params_json,
190 reply: reply_tx,
191 };
192 if self.inner.cmd_tx.send(cmd).is_err() {
193 return Pending::ready_err(LspError::ActorGone);
194 }
195 // Adapt Value → R inside a small relay so the public
196 // API returns Pending<R> rather than Pending<Value>.
197 let id = InvocationId::next();
198 let (tx, rx) = oneshot::channel::<LspResult<R>>();
199 tokio::spawn(async move {
200 let result = match reply_rx.await {
201 Ok(Ok(v)) => match serde_json::from_value::<R>(v) {
202 Ok(r) => Ok(r),
203 Err(e) => Err(LspError::ResponseDecode(e)),
204 },
205 Ok(Err(e)) => Err(e),
206 Err(_) => Err(LspError::ResponseDropped),
207 };
208 let _ = tx.send(result);
209 });
210 Pending::new(id, rx)
211 }
212
213 /// Same as [`Self::request`] but with cooperative cancellation
214 /// driven by a [`lattice_protocol::CancellationToken`]. While
215 /// the relay task is awaiting the response from the actor, it
216 /// also polls the token; if the token flips before the response
217 /// arrives, the relay resolves with [`LspError::Cancelled`].
218 ///
219 /// **Local-only cancellation today.** The server may keep
220 /// computing -- we just drop its result if it arrives stale.
221 /// `$/cancelRequest` over the wire is a Phase 4.2 polish item
222 /// (requires plumbing the JSON-RPC id back from the actor for
223 /// server-side cancel; not on a hot path).
224 ///
225 /// Used by every Phase 4.2 navigation feature
226 /// ([`Self::hover`] / [`Self::goto_definition`] / ...) so
227 /// stale popups don't appear after the user moves on.
228 pub fn request_with_cancel<P, R>(
229 &self,
230 method: &str,
231 params: P,
232 token: lattice_protocol::CancellationToken,
233 ) -> Pending<R>
234 where
235 P: serde::Serialize,
236 R: serde::de::DeserializeOwned + Send + 'static,
237 {
238 let params_json = match serde_json::to_value(params) {
239 Ok(v) => Some(v),
240 Err(e) => return Pending::ready_err(LspError::ResponseDecode(e)),
241 };
242 let (reply_tx, mut reply_rx) = oneshot::channel::<LspResult<Value>>();
243 let cmd = ActorCmd::Request {
244 method: method.to_string(),
245 params: params_json,
246 reply: reply_tx,
247 };
248 if self.inner.cmd_tx.send(cmd).is_err() {
249 return Pending::ready_err(LspError::ActorGone);
250 }
251 let id = InvocationId::next();
252 let (tx, rx) = oneshot::channel::<LspResult<R>>();
253 tokio::spawn(async move {
254 // 10ms poll cadence balances responsiveness (a typical
255 // human-perceptible delay is >50ms) against wakeup
256 // overhead. The token is an `Arc<AtomicBool>`; the poll
257 // is one Acquire load.
258 let result = loop {
259 tokio::select! {
260 biased;
261 v = &mut reply_rx => {
262 break match v {
263 Ok(Ok(v)) => match serde_json::from_value::<R>(v) {
264 Ok(r) => Ok(r),
265 Err(e) => Err(LspError::ResponseDecode(e)),
266 },
267 Ok(Err(e)) => Err(e),
268 Err(_) => Err(LspError::ResponseDropped),
269 };
270 }
271 _ = tokio::time::sleep(std::time::Duration::from_millis(10)) => {
272 if token.is_cancelled() {
273 break Err(LspError::Cancelled);
274 }
275 }
276 }
277 };
278 let _ = tx.send(result);
279 });
280 Pending::new(id, rx)
281 }
282
283 /// Fire a JSON-RPC notification (no response expected).
284 pub fn notify<P: serde::Serialize>(&self, method: &str, params: P) -> LspResult<()> {
285 let params_json = serde_json::to_value(params).map_err(LspError::ResponseDecode)?;
286 let cmd = ActorCmd::Notify {
287 method: method.to_string(),
288 params: Some(params_json),
289 };
290 self.inner.cmd_tx.send(cmd).map_err(|_| LspError::ActorGone)
291 }
292
293 /// 4.4.b: send `$/setTrace { value }` to the server (LSP
294 /// §3.18). `Off` silences trace records; `Messages` ships
295 /// the wire shapes; `Verbose` ships shapes + parameter
296 /// contents. The server replies with `$/logTrace`
297 /// notifications which the host routes into the
298 /// `*lsp:<server>:trace*` ring.
299 pub fn set_trace(&self, value: lsp_types::TraceValue) -> LspResult<()> {
300 self.notify("$/setTrace", lsp_types::SetTraceParams { value })
301 }
302
303 /// Cancel an in-flight server-side `$/progress` operation
304 /// (LSP §3.16 `window/workDoneProgress/cancel`). The server
305 /// is asked to wind down the work tied to `token`; whether
306 /// it complies is server-specific. The host treats the
307 /// cancel as best-effort — the modeline keeps the entry
308 /// until an `end` progress notification arrives.
309 pub fn cancel_progress(&self, token: &str) -> LspResult<()> {
310 // Wire shape: { "token": <number | string> }. We always
311 // serialise as a string here; that's the canonical
312 // representation the host uses to key its accumulator,
313 // and servers accept either form.
314 self.notify(
315 "window/workDoneProgress/cancel",
316 serde_json::json!({ "token": token }),
317 )
318 }
319
320 /// Cancel a pending request by JSON-RPC numeric id (LSP
321 /// `$/cancelRequest`). The actor sends the cancel
322 /// notification and resolves the matching pending oneshot
323 /// with [`LspError::Cancelled`].
324 ///
325 /// Note: the JSON-RPC id is server-internal -- callers
326 /// usually don't have it. The higher-level cancellation
327 /// path uses [`lattice_runtime::CancellationToken`] which
328 /// the editor binds to a request via `request_with_cancel`
329 /// (added in 4.2 alongside the navigation features).
330 pub fn cancel(&self, jsonrpc_id: i64) -> LspResult<()> {
331 let cmd = ActorCmd::Cancel { id: jsonrpc_id };
332 self.inner.cmd_tx.send(cmd).map_err(|_| LspError::ActorGone)
333 }
334
335 /// Open a document in the actor's own DocSync mirror and emit
336 /// `textDocument/didOpen`. Single-writer through the actor:
337 /// the supervisor mutex is no longer in this path, so the UI
338 /// thread cannot stall behind a flush. FIFO ordering with
339 /// subsequent `record_edit` calls is preserved by the cmd
340 /// channel.
341 pub fn open_doc(
342 &self,
343 uri: lsp_types::Uri,
344 language_id: impl Into<String>,
345 text: impl Into<String>,
346 ) -> LspResult<()> {
347 let cmd = ActorCmd::OpenDoc {
348 uri,
349 language_id: language_id.into(),
350 text: text.into(),
351 };
352 self.inner.cmd_tx.send(cmd).map_err(|_| LspError::ActorGone)
353 }
354
355 /// Record an edit against the actor's mirror. Coalesced with
356 /// other edits and flushed after the actor's debounce window.
357 /// Drop-free: the cmd channel is unbounded, so the publisher
358 /// (typically the per-server fan-in task) cannot lose work.
359 pub fn record_edit(
360 &self,
361 uri: lsp_types::Uri,
362 edit: lattice_protocol::edit::Edit,
363 ) -> LspResult<()> {
364 let cmd = ActorCmd::RecordEdit { uri, edit };
365 self.inner.cmd_tx.send(cmd).map_err(|_| LspError::ActorGone)
366 }
367
368 /// Force a flush of the pending change queue for one URI.
369 /// Useful before a synchronous request that depends on the
370 /// server having seen the latest text (hover/definition right
371 /// after typing).
372 pub fn flush(&self, uri: lsp_types::Uri) -> LspResult<()> {
373 let cmd = ActorCmd::Flush { uri };
374 self.inner.cmd_tx.send(cmd).map_err(|_| LspError::ActorGone)
375 }
376
377 /// Force a flush of every URI tracked by this actor.
378 pub fn flush_all(&self) -> LspResult<()> {
379 let cmd = ActorCmd::FlushAll;
380 self.inner.cmd_tx.send(cmd).map_err(|_| LspError::ActorGone)
381 }
382
383 /// Close a document: emit any final `textDocument/didChange`
384 /// then `textDocument/didClose`, all from inside the actor.
385 pub fn close_doc(&self, uri: lsp_types::Uri) -> LspResult<()> {
386 let cmd = ActorCmd::CloseDoc { uri };
387 self.inner.cmd_tx.send(cmd).map_err(|_| LspError::ActorGone)
388 }
389
390 /// Run the LSP shutdown sequence: `shutdown` request → `exit`
391 /// notification → wait for child exit. After this resolves
392 /// the actor task is gone; subsequent requests yield
393 /// [`LspError::ActorGone`].
394 pub fn shutdown(&self) -> Pending<()> {
395 let id = InvocationId::next();
396 let (tx, rx) = oneshot::channel::<LspResult<()>>();
397 let cmd = ActorCmd::Shutdown { reply: tx };
398 if self.inner.cmd_tx.send(cmd).is_err() {
399 return Pending::ready_err(LspError::ActorGone);
400 }
401 Pending::new(id, rx)
402 }
403}
404
405/// Internal actor commands -- not part of the public API.
406enum ActorCmd {
407 Request {
408 method: String,
409 params: Option<Value>,
410 reply: oneshot::Sender<LspResult<Value>>,
411 },
412 Notify {
413 method: String,
414 params: Option<Value>,
415 },
416 Cancel {
417 id: i64,
418 },
419 Shutdown {
420 reply: oneshot::Sender<LspResult<()>>,
421 },
422 /// Bring a buffer under the actor's DocSync management. The
423 /// actor builds the `didOpen` payload via `DocSync::open` +
424 /// ships it. v1 is fire-and-forget -- if the channel push
425 /// fails the caller already lost the connection. Phase 4.x
426 /// per-actor edit-path refactor.
427 OpenDoc {
428 uri: lsp_types::Uri,
429 language_id: String,
430 text: String,
431 },
432 /// Apply one committed edit to the actor's DocSync mirror +
433 /// queue the `didChange` event. The per-actor debounce
434 /// timer drives the eventual flush; rapid edits coalesce.
435 /// Phase 4.x.
436 RecordEdit {
437 uri: lsp_types::Uri,
438 edit: lattice_protocol::edit::Edit,
439 },
440 /// Eagerly drain queued change events for `uri` and ship a
441 /// `didChange`. Used by will-save hooks etc. that need a
442 /// coherent server-side view RIGHT NOW. Phase 4.x.
443 Flush {
444 uri: lsp_types::Uri,
445 },
446 /// Same as `Flush` but for every URI the actor tracks.
447 /// Used at editor shutdown. Phase 4.x.
448 FlushAll,
449 /// Drop a buffer's mirror; ship the optional final
450 /// `didChange` followed by `didClose`. Phase 4.x.
451 CloseDoc {
452 uri: lsp_types::Uri,
453 },
454}
455
456/// Spawn a language server from a [`ServerConfig`].
457///
458/// Performs the initialize handshake before returning. Failures
459/// at any handshake step (spawn, framing, decode, server error
460/// response, missing required capability) surface as
461/// [`LspError`]. `logger` is the shared subsystem logger -- in
462/// production every server uses the App's clone; tests pass a
463/// fresh `LspLogger::with_defaults()`.
464pub async fn spawn(
465 config: ServerConfig,
466 workspace_root: std::path::PathBuf,
467 logger: LspLogger,
468 apply_edit_bus: Option<crate::apply_edit::ApplyEditBus>,
469 configuration_bus: Option<crate::configuration::ConfigurationBus>,
470 show_document_bus: Option<crate::show_document::ShowDocumentBus>,
471 show_message_request_bus: Option<crate::show_message_request::ShowMessageRequestBus>,
472 event_bus: Option<Arc<lattice_runtime::EventBus>>,
473) -> LspResult<ServerHandle> {
474 let transport = ChildTransport::spawn(&config.binary, &config.args, Some(&workspace_root))
475 .await
476 .map_err(LspError::Transport)?;
477 let (reader, writer, stderr, child) = transport.split();
478 spawn_with_io(
479 config,
480 workspace_root,
481 reader,
482 writer,
483 stderr,
484 Some(child),
485 logger,
486 apply_edit_bus,
487 configuration_bus,
488 show_document_bus,
489 show_message_request_bus,
490 event_bus,
491 )
492 .await
493}
494
495/// Spawn the actor against pre-existing `LspReader` / `LspWriter`
496/// halves. Used by tests (mock server over a duplex pipe) and
497/// by future embedded transports (TCP, named pipe).
498///
499/// `child` is `None` for in-process tests (the duplex partner is
500/// a tokio task) and `Some(_)` for the real child-process path.
501/// When None, the shutdown sequence skips the child-exit wait.
502#[allow(clippy::too_many_arguments)]
503pub async fn spawn_with_io<R, W>(
504 config: ServerConfig,
505 workspace_root: std::path::PathBuf,
506 reader: LspReader<R>,
507 writer: LspWriter<W>,
508 stderr: Option<ChildStderr>,
509 child: Option<Child>,
510 logger: LspLogger,
511 apply_edit_bus: Option<crate::apply_edit::ApplyEditBus>,
512 configuration_bus: Option<crate::configuration::ConfigurationBus>,
513 show_document_bus: Option<crate::show_document::ShowDocumentBus>,
514 show_message_request_bus: Option<crate::show_message_request::ShowMessageRequestBus>,
515 event_bus: Option<Arc<lattice_runtime::EventBus>>,
516) -> LspResult<ServerHandle>
517where
518 R: AsyncBufRead + Unpin + Send + 'static,
519 W: AsyncWrite + Unpin + Send + 'static,
520{
521 let (cmd_tx, cmd_rx) = mpsc::unbounded_channel::<ActorCmd>();
522 // 4.4.n: handshake hands back the published `ArcSwap` cell
523 // (not a plain `Arc<Capabilities>`) so the host's handle and
524 // the actor task share one publication point. After
525 // handshake the actor swaps in fresh snapshots whenever
526 // `client/(un)registerCapability` modifies the dynamic
527 // registry; readers (`ServerHandle::capabilities()`) load
528 // through the same cell.
529 let (handshake_tx, handshake_rx) = oneshot::channel::<LspResult<Arc<ArcSwap<Capabilities>>>>();
530
531 let server_id = config.id.clone();
532 let init_options = config.initialization_options.clone();
533 let workspace_folder_uri = uri_from_path(&workspace_root);
534 let workspace_name = workspace_root
535 .file_name()
536 .map(|s| s.to_string_lossy().into_owned())
537 .unwrap_or_else(|| "workspace".to_string());
538 // B'.2: the actor's logger calls carry the canonical
539 // `(server_id, workspace)` pair so multi-instance setups
540 // (two `rust-analyzer`s on different workspaces) stay
541 // distinct in the per-instance log rings.
542 let instance = crate::logging::InstanceKey::new(
543 Arc::<str>::from(server_id.as_str()),
544 Arc::<std::path::Path>::from(workspace_root.as_path()),
545 );
546 // One bus per actor, shared with the read_loop's notification
547 // dispatcher.
548 let diagnostics = DiagnosticsBus::new();
549
550 // Subsystem-wide event: server spawn / handshake start.
551 logger.log(
552 None,
553 LogLevel::Info,
554 LogSource::Client,
555 format!(
556 "spawning LSP actor for server {:?} workspace {}",
557 server_id,
558 workspace_root.display()
559 ),
560 );
561
562 tokio::spawn(actor_main(
563 reader,
564 writer,
565 child,
566 stderr,
567 cmd_rx,
568 handshake_tx,
569 server_id.clone(),
570 instance.clone(),
571 workspace_folder_uri,
572 workspace_name,
573 init_options,
574 diagnostics.clone(),
575 logger.clone(),
576 apply_edit_bus,
577 configuration_bus,
578 show_document_bus,
579 show_message_request_bus,
580 event_bus,
581 ));
582
583 let capabilities = handshake_rx
584 .await
585 .map_err(|_| LspError::HandshakeFailed("actor died before handshake".into()))??;
586
587 logger.log(
588 Some(&instance),
589 LogLevel::Info,
590 LogSource::Client,
591 "handshake complete; server attached",
592 );
593 let _server_id_arc: Arc<str> = Arc::clone(&instance.server_id);
594 let workspace_root_arc: Arc<std::path::Path> = Arc::clone(&instance.workspace);
595
596 Ok(ServerHandle {
597 inner: Arc::new(HandleInner {
598 cmd_tx,
599 // 4.4.n: ArcSwap shared with the actor task so
600 // register / unregister updates are observable
601 // from readers without restarting the actor.
602 capabilities,
603 server_id,
604 workspace_root: workspace_root_arc,
605 diagnostics,
606 logger,
607 }),
608 })
609}
610
611/// Convert a filesystem path to a `file://` URI in the form LSP
612/// expects. lsp-types 0.97 dropped the `url` crate; we
613/// percent-encode manually for the small set of bytes that
614/// matter in a path (space → `%20`, etc.). Servers in practice
615/// tolerate plain `file:///<path>` without aggressive encoding.
616/// Inverse of [`uri_from_path`]. Strips the `file://` scheme +
617/// percent-decodes the small set of bytes the encoder rewrites.
618/// Returns the path string (caller decides whether to coerce to
619/// `PathBuf`); `None` for non-`file://` URIs.
620pub fn uri_to_path(uri: &Uri) -> Option<std::path::PathBuf> {
621 let s = uri.as_str();
622 let stripped = s.strip_prefix("file://")?;
623 // Percent-decode the small set the encoder writes. More
624 // exotic encodings (CJK paths etc) round-trip via the
625 // identity branch in uri_from_path.
626 let mut out = String::with_capacity(stripped.len());
627 let mut chars = stripped.chars().peekable();
628 while let Some(c) = chars.next() {
629 if c == '%' {
630 let h1 = chars.next();
631 let h2 = chars.next();
632 if let (Some(h1), Some(h2)) = (h1, h2)
633 && let (Some(h1), Some(h2)) = (h1.to_digit(16), h2.to_digit(16))
634 {
635 let byte = (h1 * 16 + h2) as u8;
636 out.push(byte as char);
637 continue;
638 }
639 // Malformed; keep the literal `%` and any consumed chars.
640 out.push('%');
641 if let Some(h1) = h1 {
642 out.push(h1);
643 }
644 if let Some(h2) = h2 {
645 out.push(h2);
646 }
647 } else {
648 out.push(c);
649 }
650 }
651 Some(std::path::PathBuf::from(without_drive_slash(
652 out,
653 cfg!(windows),
654 )))
655}
656
657/// `/C:/Users/me` → `C:/Users/me`.
658///
659/// [`uri_from_path`] writes a Windows path as `file:///C:/Users/me` — three
660/// slashes, the third opening the path — so stripping `file://` leaves a
661/// leading `/` in front of the drive letter, and `/C:/Users/me` is not a path
662/// Windows can open. Every location a server sent back (a definition, a
663/// reference, a diagnostic's file) resolved to a file that did not exist.
664///
665/// Windows only: elsewhere `/C:/x` is a legitimate, if unlikely, absolute
666/// path. Pure, so the Windows half is tested on every platform.
667fn without_drive_slash(path: String, windows: bool) -> String {
668 let b = path.as_bytes();
669 if windows && b.len() >= 3 && b[0] == b'/' && b[1].is_ascii_alphabetic() && b[2] == b':' {
670 path[1..].to_string()
671 } else {
672 path
673 }
674}
675
676pub fn uri_from_path(p: &std::path::Path) -> Uri {
677 // Promote relative paths to absolute *before* building the URI.
678 // `file:///crates/lattice-core/src/buffer.rs` is interpreted by
679 // every LSP server as `/crates/lattice-core/...` (root-rooted),
680 // not as "the file the editor opened from a relative arg" --
681 // rust-analyzer then can't find the file inside its workspace,
682 // returns null hovers / definitions, and emits notify-watcher
683 // warnings about paths that don't exist on disk. The fix is to
684 // canonicalise to an absolute path here so every URI we send is
685 // wire-correct regardless of how the user invoked the editor
686 // (relative path on the cli, etc.).
687 //
688 // `std::path::absolute` does NOT do I/O and does NOT resolve
689 // symlinks -- both are deliberate. We just want
690 // `cwd().join(p)`-shaped output. If absolute() fails (very rare
691 // -- malformed path on Windows mostly), fall back to the
692 // original; the URI may then still be wrong but the server's
693 // existing failure mode (null reply + warning) is no worse than
694 // before this fix.
695 let absolute = std::path::absolute(p)
696 .ok()
697 .unwrap_or_else(|| p.to_path_buf());
698 let display = absolute.to_string_lossy();
699 // Normalise Windows backslashes to forward slashes so the
700 // URI is well-formed across platforms. Drive letters
701 // (`C:\`) remain in `<C:/path/...>` form, which is what LSP
702 // servers expect.
703 let normalised = display.replace('\\', "/");
704 let raw = if normalised.starts_with('/') {
705 format!("file://{normalised}")
706 } else {
707 format!("file:///{normalised}")
708 };
709 // Percent-encode the small set of byte classes that fluent-uri
710 // rejects (spaces, control chars). Anything else passes through.
711 let mut encoded = String::with_capacity(raw.len());
712 for c in raw.chars() {
713 match c {
714 ' ' => encoded.push_str("%20"),
715 '"' => encoded.push_str("%22"),
716 '<' => encoded.push_str("%3C"),
717 '>' => encoded.push_str("%3E"),
718 '|' => encoded.push_str("%7C"),
719 other => encoded.push(other),
720 }
721 }
722 Uri::from_str(&encoded).unwrap_or_else(|_| {
723 // If even that fails, fall back to a synthetic
724 // host-only URI so the actor doesn't panic.
725 Uri::from_str("file:///").expect("file:/// is a valid URI")
726 })
727}
728
729#[allow(clippy::too_many_arguments)]
730async fn actor_main<R, W>(
731 reader: LspReader<R>,
732 writer: LspWriter<W>,
733 mut child: Option<Child>,
734 stderr: Option<ChildStderr>,
735 mut cmd_rx: mpsc::UnboundedReceiver<ActorCmd>,
736 handshake_tx: oneshot::Sender<LspResult<Arc<ArcSwap<Capabilities>>>>,
737 server_id: String,
738 instance: crate::logging::InstanceKey,
739 workspace_folder_uri: Uri,
740 workspace_name: String,
741 init_options: Option<Value>,
742 diagnostics: DiagnosticsBus,
743 logger: LspLogger,
744 apply_edit_bus: Option<crate::apply_edit::ApplyEditBus>,
745 configuration_bus: Option<crate::configuration::ConfigurationBus>,
746 show_document_bus: Option<crate::show_document::ShowDocumentBus>,
747 show_message_request_bus: Option<crate::show_message_request::ShowMessageRequestBus>,
748 event_bus: Option<Arc<lattice_runtime::EventBus>>,
749) where
750 R: AsyncBufRead + Unpin + Send + 'static,
751 W: AsyncWrite + Unpin + Send + 'static,
752{
753 let server_id_arc: Arc<str> = Arc::clone(&instance.server_id);
754 // Drain stderr through the logger -- each line lands in the
755 // per-instance `*lsp:<server>:<workspace>*` ring at Warn
756 // (server stderr is the canonical "something's up" signal).
757 if let Some(stderr) = stderr {
758 tokio::spawn(stderr_drain(stderr, instance.clone(), logger.clone()));
759 }
760
761 // Spawn write_loop. Mutex around the writer because
762 // `write_loop` also writes the initialize request (via the
763 // outbound channel) before the loop starts processing
764 // mailbox commands. After spawn, the actor only writes
765 // through the channel.
766 let writer = Arc::new(Mutex::new(writer));
767 let (out_tx, out_rx) = mpsc::unbounded_channel::<Message>();
768 tokio::spawn(write_loop(
769 Arc::clone(&writer),
770 out_rx,
771 instance.clone(),
772 logger.clone(),
773 ));
774
775 // Spawn read_loop.
776 let (in_tx, mut in_rx) = mpsc::unbounded_channel::<Message>();
777 tokio::spawn(read_loop(reader, in_tx, instance.clone(), logger.clone()));
778
779 // Handshake.
780 let mut next_id: u64 = 1;
781 let init_id = RequestId::from_u64(next_id);
782 next_id += 1;
783 let init_params = build_initialize_params(workspace_folder_uri, workspace_name, init_options);
784 let init_value = match serde_json::to_value(&init_params) {
785 Ok(v) => v,
786 Err(e) => {
787 let _ = handshake_tx.send(Err(LspError::HandshakeFailed(format!(
788 "could not serialize initialize params: {e}"
789 ))));
790 return;
791 }
792 };
793 let req = Request::new(init_id.clone(), "initialize", Some(init_value));
794 if out_tx.send(Message::Request(req)).is_err() {
795 let _ = handshake_tx.send(Err(LspError::HandshakeFailed(
796 "write_loop closed before initialize".into(),
797 )));
798 return;
799 }
800
801 // Wait for the matching initialize response. While waiting,
802 // the server may send window/logMessage or $/progress -- log
803 // them and keep waiting.
804 let init_result = loop {
805 match in_rx.recv().await {
806 Some(Message::Response(r)) if r.id == init_id => break r,
807 Some(other) => {
808 handle_pre_handshake_message(
809 &instance,
810 other,
811 &diagnostics,
812 &logger,
813 event_bus.as_ref(),
814 );
815 }
816 None => {
817 let _ = handshake_tx.send(Err(LspError::HandshakeFailed(
818 "server stream closed before initialize response".into(),
819 )));
820 return;
821 }
822 }
823 };
824
825 let init_value = match init_result.error {
826 Some(err) => {
827 let _ = handshake_tx.send(Err(LspError::HandshakeFailed(format!(
828 "server rejected initialize: {} ({})",
829 err.message, err.code
830 ))));
831 return;
832 }
833 None => init_result.result.unwrap_or(Value::Null),
834 };
835
836 let server_caps = match serde_json::from_value::<InitializeResult>(init_value) {
837 Ok(r) => r.capabilities,
838 Err(e) => {
839 let _ = handshake_tx.send(Err(LspError::HandshakeFailed(format!(
840 "could not deserialize InitializeResult: {e}"
841 ))));
842 return;
843 }
844 };
845
846 let caps = Capabilities::from_initialize(capabilities::client_capabilities(), server_caps);
847 // 4.4.n: publish the handshake snapshot into the shared
848 // ArcSwap cell. `caps_cell` is the publication point both
849 // the host's `ServerHandle::capabilities()` and the
850 // actor's own dynamic-registration code read / write.
851 let caps_cell: Arc<ArcSwap<Capabilities>> = Arc::new(ArcSwap::new(Arc::clone(&caps)));
852 // Local snapshot used by the actor task between mutations.
853 // After `client/(un)registerCapability` we rebuild + store
854 // a new `Arc<Capabilities>` and refresh this binding so
855 // subsequent reads in this task see the latest state
856 // without a load through the cell.
857 let mut caps: Arc<Capabilities> = caps;
858
859 // Send `initialized` notification per LSP base spec -- the
860 // server is required to wait for this before processing
861 // other requests.
862 let initialized_value = serde_json::to_value(InitializedParams {}).unwrap_or(Value::Null);
863 let _ = out_tx.send(Message::Notification(Notification::new(
864 "initialized",
865 Some(initialized_value),
866 )));
867
868 // Hand the handle back to the caller.
869 if handshake_tx.send(Ok(Arc::clone(&caps_cell))).is_err() {
870 // The caller dropped the spawn future -- nothing else to
871 // do. Run shutdown locally.
872 perform_shutdown(&out_tx, &mut in_rx, &mut next_id, child.as_mut()).await;
873 return;
874 }
875
876 // Main loop.
877 let mut pending: HashMap<RequestId, oneshot::Sender<LspResult<Value>>> = HashMap::new();
878 let mut shutting_down: Option<oneshot::Sender<LspResult<()>>> = None;
879
880 // Per-actor DocSync state (Phase 4.x edit-path refactor).
881 // Owned by the actor so the select! loop is the single
882 // writer to the mirror -- the incremental-sync invariant
883 // is structurally guaranteed, no shared lock to misuse.
884 let mut docsync = crate::sync::DocSync::new();
885 // Per-actor debounce for `textDocument/didChange`. After
886 // every RecordEdit we set `flush_deadline` to ~50ms in the
887 // future; the select! arm guarded by `flush_pending` polls
888 // the sleep, fires when the deadline hits, then clears the
889 // flag. Rapid edits coalesce because each new RecordEdit
890 // resets the deadline.
891 const FLUSH_DEBOUNCE: std::time::Duration = std::time::Duration::from_millis(50);
892 let mut flush_pending = false;
893 let flush_sleep = tokio::time::sleep(std::time::Duration::from_secs(60 * 60));
894 tokio::pin!(flush_sleep);
895
896 'main: loop {
897 select! {
898 // Debounced flush. Only polled when `flush_pending`
899 // is true (a RecordEdit set it); fires once after
900 // FLUSH_DEBOUNCE of idleness past the last edit.
901 _ = &mut flush_sleep, if flush_pending => {
902 flush_pending = false;
903 for (uri, params) in docsync.take_flush_all_payloads(&caps) {
904 let n = Notification::new(
905 "textDocument/didChange",
906 Some(serde_json::to_value(params).unwrap_or(Value::Null)),
907 );
908 if out_tx.send(Message::Notification(n)).is_err() {
909 // write_loop dead -- bail.
910 break 'main;
911 }
912 let _ = uri;
913 }
914 }
915 cmd = cmd_rx.recv() => {
916 match cmd {
917 Some(ActorCmd::Request { method, params, reply }) => {
918 let id = RequestId::from_u64(next_id); next_id += 1;
919 pending.insert(id.clone(), reply);
920 let req = Request::new(id, method, params);
921 if out_tx.send(Message::Request(req)).is_err() {
922 // write_loop dead -- everyone's pending
923 // request resolves with ActorGone.
924 break 'main;
925 }
926 }
927 Some(ActorCmd::Notify { method, params }) => {
928 let n = Notification::new(method, params);
929 let _ = out_tx.send(Message::Notification(n));
930 }
931 Some(ActorCmd::Cancel { id }) => {
932 let cancel = Notification::new(
933 "$/cancelRequest",
934 Some(serde_json::json!({"id": id})),
935 );
936 let _ = out_tx.send(Message::Notification(cancel));
937 // Resolve the matching pending entry locally
938 // so the caller doesn't wait for the server's
939 // ack.
940 let key = RequestId::Number(id);
941 if let Some(reply) = pending.remove(&key) {
942 let _ = reply.send(Err(LspError::Cancelled));
943 }
944 }
945 Some(ActorCmd::OpenDoc { uri, language_id, text }) => {
946 let params = docsync.open(uri, language_id, text);
947 let n = Notification::new(
948 "textDocument/didOpen",
949 Some(serde_json::to_value(params).unwrap_or(Value::Null)),
950 );
951 if out_tx.send(Message::Notification(n)).is_err() {
952 break 'main;
953 }
954 }
955 Some(ActorCmd::RecordEdit { uri, edit }) => {
956 // Single-writer to the mirror -- no locks,
957 // no drops. Invariant: every committed edit
958 // either applies cleanly here or logs a
959 // warning (mirror corruption is impossible
960 // because the actor owns the only mutator).
961 if let Err(e) = docsync.record_edit(&caps, &uri, &edit) {
962 logger.log(
963 Some(&instance),
964 LogLevel::Warn,
965 LogSource::Client,
966 format!(
967 "actor.record_edit on {}: {e}",
968 uri.as_str(),
969 ),
970 );
971 }
972 // Reset the debounce: rapid edits coalesce
973 // into one didChange after the idle window.
974 flush_sleep.as_mut().reset(
975 tokio::time::Instant::now() + FLUSH_DEBOUNCE,
976 );
977 flush_pending = true;
978 }
979 Some(ActorCmd::Flush { uri }) => {
980 if let Some(params) =
981 docsync.take_flush_payload(&caps, &uri)
982 {
983 let n = Notification::new(
984 "textDocument/didChange",
985 Some(serde_json::to_value(params).unwrap_or(Value::Null)),
986 );
987 if out_tx.send(Message::Notification(n)).is_err() {
988 break 'main;
989 }
990 }
991 }
992 Some(ActorCmd::FlushAll) => {
993 for (_uri, params) in
994 docsync.take_flush_all_payloads(&caps)
995 {
996 let n = Notification::new(
997 "textDocument/didChange",
998 Some(serde_json::to_value(params).unwrap_or(Value::Null)),
999 );
1000 if out_tx.send(Message::Notification(n)).is_err() {
1001 break 'main;
1002 }
1003 }
1004 flush_pending = false;
1005 }
1006 Some(ActorCmd::CloseDoc { uri }) => {
1007 if let Some(payloads) = docsync.close(&caps, &uri) {
1008 if let Some(final_changes) = payloads.final_changes {
1009 let n = Notification::new(
1010 "textDocument/didChange",
1011 Some(serde_json::to_value(final_changes)
1012 .unwrap_or(Value::Null)),
1013 );
1014 if out_tx.send(Message::Notification(n)).is_err() {
1015 break 'main;
1016 }
1017 }
1018 let n = Notification::new(
1019 "textDocument/didClose",
1020 Some(serde_json::to_value(payloads.close)
1021 .unwrap_or(Value::Null)),
1022 );
1023 if out_tx.send(Message::Notification(n)).is_err() {
1024 break 'main;
1025 }
1026 }
1027 }
1028 Some(ActorCmd::Shutdown { reply }) => {
1029 shutting_down = Some(reply);
1030 // Send the shutdown request; the response
1031 // arrives back through in_rx and we then
1032 // emit `exit`. We don't break the loop yet
1033 // -- need to await the shutdown response.
1034 // No `next_id += 1` here: we break out of
1035 // the loop immediately after, so the bump
1036 // would be dead.
1037 let id = RequestId::from_u64(next_id);
1038 let (sd_tx, sd_rx) = oneshot::channel::<LspResult<Value>>();
1039 pending.insert(id.clone(), sd_tx);
1040 let _ = out_tx.send(Message::Request(Request::new(
1041 id,
1042 "shutdown",
1043 None,
1044 )));
1045 // Wait briefly for shutdown response, then
1046 // emit exit. Bound at 5s to avoid hanging
1047 // on a misbehaving server.
1048 let timeout = tokio::time::sleep(std::time::Duration::from_secs(5));
1049 tokio::pin!(timeout);
1050 select! {
1051 _ = sd_rx => {},
1052 _ = &mut timeout => {
1053 tracing::warn!(server_id, "shutdown response timed out; sending exit");
1054 }
1055 }
1056 let _ = out_tx.send(Message::Notification(Notification::new(
1057 "exit", None,
1058 )));
1059 if let Some(c) = child.as_mut() {
1060 // Best-effort wait on child exit.
1061 let _ = tokio::time::timeout(
1062 std::time::Duration::from_secs(2),
1063 c.wait(),
1064 )
1065 .await;
1066 }
1067 break 'main;
1068 }
1069 None => {
1070 // Mailbox closed (last ServerHandle dropped).
1071 // Run a graceful shutdown anyway.
1072 perform_shutdown(&out_tx, &mut in_rx, &mut next_id, child.as_mut()).await;
1073 break 'main;
1074 }
1075 }
1076 },
1077 msg = in_rx.recv() => {
1078 match msg {
1079 Some(Message::Response(r)) => {
1080 if let Some(reply) = pending.remove(&r.id) {
1081 let result = match r.error {
1082 Some(e) => Err(LspError::Server(e)),
1083 None => Ok(r.result.unwrap_or(Value::Null)),
1084 };
1085 let _ = reply.send(result);
1086 } else {
1087 tracing::warn!(server_id, ?r.id, "response with unknown id");
1088 }
1089 }
1090 Some(Message::Notification(n)) => {
1091 handle_server_notification(
1092 &instance,
1093 &n,
1094 &diagnostics,
1095 &logger,
1096 event_bus.as_ref(),
1097 );
1098 }
1099 Some(Message::Request(req)) => {
1100 // `workspace/applyEdit` (Phase 4.3) is
1101 // async: the App must apply the edit on
1102 // the UI thread, so we forward the
1103 // request through the apply-edit bus and
1104 // a spawned task awaits the response
1105 // before writing it back to the wire.
1106 // All other server-initiated requests
1107 // resolve synchronously inline.
1108 if req.method == "workspace/applyEdit"
1109 && let Some(bus) = apply_edit_bus.as_ref()
1110 {
1111 let bus = bus.clone();
1112 let instance_clone = instance.clone();
1113 let logger_clone = logger.clone();
1114 let out_tx_clone = out_tx.clone();
1115 let req_id = req.id.clone();
1116 let params = req.params.clone();
1117 tokio::spawn(async move {
1118 let resp = handle_apply_edit_request(
1119 instance_clone,
1120 req_id,
1121 params,
1122 &bus,
1123 &logger_clone,
1124 )
1125 .await;
1126 let _ = out_tx_clone.send(Message::Response(resp));
1127 });
1128 } else if req.method == "workspace/configuration"
1129 && let Some(bus) = configuration_bus.as_ref()
1130 {
1131 // Phase 4.1 follow-up: real values
1132 // (not just `null` per item) require
1133 // the App to walk its cached TOML
1134 // tree. Same async pattern as
1135 // applyEdit -- dispatch via the bus,
1136 // await response, ferry back.
1137 let bus = bus.clone();
1138 let instance_clone = instance.clone();
1139 let logger_clone = logger.clone();
1140 let out_tx_clone = out_tx.clone();
1141 let req_id = req.id.clone();
1142 let params = req.params.clone();
1143 tokio::spawn(async move {
1144 let resp = handle_configuration_request(
1145 instance_clone,
1146 req_id,
1147 params,
1148 &bus,
1149 &logger_clone,
1150 )
1151 .await;
1152 let _ = out_tx_clone.send(Message::Response(resp));
1153 });
1154 } else if req.method == "window/showDocument"
1155 && let Some(bus) = show_document_bus.as_ref()
1156 {
1157 // 4.4.b: server wants the host to open
1158 // a URI (file -> buffer; external ->
1159 // OS handler). The App's drain
1160 // performs the open + writes back via
1161 // the embedded oneshot.
1162 let bus = bus.clone();
1163 let instance_clone = instance.clone();
1164 let logger_clone = logger.clone();
1165 let out_tx_clone = out_tx.clone();
1166 let req_id = req.id.clone();
1167 let params = req.params.clone();
1168 tokio::spawn(async move {
1169 let resp = handle_show_document_request(
1170 instance_clone,
1171 req_id,
1172 params,
1173 &bus,
1174 &logger_clone,
1175 )
1176 .await;
1177 let _ = out_tx_clone.send(Message::Response(resp));
1178 });
1179 } else if req.method == "window/showMessageRequest"
1180 && let Some(bus) = show_message_request_bus.as_ref()
1181 {
1182 // 4.4.b: server-emitted modal action
1183 // request. The App opens an action
1184 // picker; the user's selection
1185 // (or `null` on dismiss) ferries back.
1186 let bus = bus.clone();
1187 let instance_clone = instance.clone();
1188 let logger_clone = logger.clone();
1189 let out_tx_clone = out_tx.clone();
1190 let req_id = req.id.clone();
1191 let params = req.params.clone();
1192 tokio::spawn(async move {
1193 let resp = handle_show_message_request(
1194 instance_clone,
1195 req_id,
1196 params,
1197 &bus,
1198 &logger_clone,
1199 )
1200 .await;
1201 let _ = out_tx_clone.send(Message::Response(resp));
1202 });
1203 } else if req.method == "workspace/inlayHint/refresh" {
1204 // 4.4.g: server-initiated inlay
1205 // hint cache invalidation. Reply
1206 // `null` synchronously per spec;
1207 // publish the typed
1208 // `LspInlayHintRefresh` event so
1209 // the App's drain clears cached
1210 // hints for attached buffers and
1211 // the next render tick re-issues
1212 // `inlayHint`.
1213 let _ = out_tx.send(Message::Response(Response::ok(
1214 req.id.clone(),
1215 Value::Null,
1216 )));
1217 if let Some(bus) = event_bus.as_ref() {
1218 bus.publish_typed(crate::events::LspInlayHintRefresh {
1219 server_id: Arc::clone(&server_id_arc),
1220 });
1221 }
1222 } else if req.method == "workspace/semanticTokens/refresh" {
1223 // 4.4.i: server-initiated semantic
1224 // tokens cache invalidation. Same
1225 // shape as the inlay-hint refresh:
1226 // reply `null` synchronously, publish
1227 // `LspSemanticTokensRefresh` so the
1228 // App's drain drops cached tokens
1229 // (and any stale `result_id`) for
1230 // attached buffers; the next render
1231 // tick re-issues a `full` request to
1232 // rebuild the baseline.
1233 let _ = out_tx.send(Message::Response(Response::ok(
1234 req.id.clone(),
1235 Value::Null,
1236 )));
1237 if let Some(bus) = event_bus.as_ref() {
1238 bus.publish_typed(crate::events::LspSemanticTokensRefresh {
1239 server_id: Arc::clone(&server_id_arc),
1240 });
1241 }
1242 } else if req.method == "workspace/inlineValue/refresh" {
1243 // 4.5.h: server-initiated inline-
1244 // value cache invalidation. The
1245 // renderer trigger is itself
1246 // deferred (no debug-adapter
1247 // integration yet); we still reply
1248 // `null` per spec so the server's
1249 // request resolves, and log so a
1250 // future renderer wire-up can grep
1251 // for the breadcrumb.
1252 let _ = out_tx.send(Message::Response(Response::ok(
1253 req.id.clone(),
1254 Value::Null,
1255 )));
1256 logger.log(
1257 Some(&instance),
1258 LogLevel::Info,
1259 LogSource::Client,
1260 "workspace/inlineValue/refresh accepted (renderer trigger deferred)",
1261 );
1262 } else if req.method == "workspace/codeLens/refresh" {
1263 // 4.5.d: server-initiated code-lens
1264 // cache invalidation. Same shape as
1265 // inlay-hint / semantic-tokens
1266 // refreshes: reply `null` inline +
1267 // publish so the App's drain evicts
1268 // cached lenses for attached buffers
1269 // and the next tick's pump re-issues
1270 // `textDocument/codeLens`.
1271 let _ = out_tx.send(Message::Response(Response::ok(
1272 req.id.clone(),
1273 Value::Null,
1274 )));
1275 if let Some(bus) = event_bus.as_ref() {
1276 bus.publish_typed(crate::events::LspCodeLensRefresh {
1277 server_id: Arc::clone(&server_id_arc),
1278 });
1279 }
1280 } else if req.method == "workspace/diagnostic/refresh" {
1281 // 4.4.j: server-initiated pull-
1282 // diagnostic invalidation. Reply
1283 // `null` inline + publish so the
1284 // App's drain evicts the per-buffer
1285 // `result_id` cache. The next render
1286 // tick re-pulls
1287 // `textDocument/diagnostic` without
1288 // a `previous_result_id`, forcing a
1289 // `Full` report regardless of what
1290 // the server had cached.
1291 let _ = out_tx.send(Message::Response(Response::ok(
1292 req.id.clone(),
1293 Value::Null,
1294 )));
1295 if let Some(bus) = event_bus.as_ref() {
1296 bus.publish_typed(crate::events::LspDiagnosticRefresh {
1297 server_id: Arc::clone(&server_id_arc),
1298 });
1299 }
1300 } else if req.method == "client/registerCapability" {
1301 // 4.4.n: parse the registration batch,
1302 // fold every entry into the dynamic
1303 // registry, publish a new caps snapshot,
1304 // reply `null`. Parse errors degrade to
1305 // a logged warning + still-reply-null --
1306 // throwing the registration away is
1307 // better than failing the request,
1308 // which most servers treat as a fatal
1309 // protocol error.
1310 let parsed = req
1311 .params
1312 .as_ref()
1313 .map(|p| {
1314 serde_json::from_value::<lsp_types::RegistrationParams>(
1315 p.clone(),
1316 )
1317 });
1318 match parsed {
1319 Some(Ok(params)) => {
1320 caps = caps.with_dynamic_mut(|reg| {
1321 for r in params.registrations {
1322 reg.register(
1323 crate::DynamicRegistration {
1324 id: r.id,
1325 method: r.method,
1326 register_options: r
1327 .register_options,
1328 },
1329 );
1330 }
1331 });
1332 caps_cell.store(Arc::clone(&caps));
1333 }
1334 Some(Err(e)) => {
1335 logger.log(
1336 Some(&instance),
1337 LogLevel::Warn,
1338 LogSource::Client,
1339 format!(
1340 "client/registerCapability: malformed params, dropping ({e})"
1341 ),
1342 );
1343 }
1344 None => {
1345 logger.log(
1346 Some(&instance),
1347 LogLevel::Warn,
1348 LogSource::Client,
1349 "client/registerCapability: missing params",
1350 );
1351 }
1352 }
1353 let _ = out_tx.send(Message::Response(Response::ok(
1354 req.id.clone(),
1355 Value::Null,
1356 )));
1357 } else if req.method == "client/unregisterCapability" {
1358 // 4.4.n: evict each entry by id; publish.
1359 // Unknown ids are silently dropped
1360 // (see `DynamicRegistry::unregister`).
1361 let parsed = req
1362 .params
1363 .as_ref()
1364 .map(|p| {
1365 serde_json::from_value::<lsp_types::UnregistrationParams>(
1366 p.clone(),
1367 )
1368 });
1369 match parsed {
1370 Some(Ok(params)) => {
1371 caps = caps.with_dynamic_mut(|reg| {
1372 for u in ¶ms.unregisterations {
1373 reg.unregister(&u.id);
1374 }
1375 });
1376 caps_cell.store(Arc::clone(&caps));
1377 }
1378 Some(Err(e)) => {
1379 logger.log(
1380 Some(&instance),
1381 LogLevel::Warn,
1382 LogSource::Client,
1383 format!(
1384 "client/unregisterCapability: malformed params, dropping ({e})"
1385 ),
1386 );
1387 }
1388 None => {
1389 logger.log(
1390 Some(&instance),
1391 LogLevel::Warn,
1392 LogSource::Client,
1393 "client/unregisterCapability: missing params",
1394 );
1395 }
1396 }
1397 let _ = out_tx.send(Message::Response(Response::ok(
1398 req.id.clone(),
1399 Value::Null,
1400 )));
1401 } else {
1402 let resp = handle_server_request(&instance, &req, &logger);
1403 let _ = out_tx.send(Message::Response(resp));
1404 }
1405 }
1406 None => {
1407 // read_loop ended -- server exited or pipe
1408 // broke. Resolve all pending requests.
1409 break 'main;
1410 }
1411 }
1412 }
1413 }
1414 }
1415
1416 // Drain pending: every outstanding caller resolves with the
1417 // appropriate failure.
1418 for (_id, reply) in pending.drain() {
1419 let _ = reply.send(Err(LspError::ActorGone));
1420 }
1421 let was_clean_shutdown = shutting_down.is_some();
1422 if let Some(reply) = shutting_down {
1423 let _ = reply.send(Ok(()));
1424 }
1425 // 4.4.d: publish actor-exited to the event bus so the
1426 // supervisor can drive auto-restart on unexpected exits.
1427 // Clean exits (the user / shutdown path requested it) are
1428 // tagged distinctly so the supervisor doesn't restart what
1429 // it just asked to die.
1430 if let Some(bus) = event_bus {
1431 let reason = if was_clean_shutdown {
1432 crate::events::LspActorExitReason::Clean
1433 } else {
1434 crate::events::LspActorExitReason::Unexpected
1435 };
1436 bus.publish_typed(crate::events::LspActorExited {
1437 server_id: server_id_arc.clone(),
1438 reason,
1439 });
1440 }
1441}
1442
1443/// Handle a server-initiated notification. `publishDiagnostics`
1444/// fans out to the [`DiagnosticsBus`]; log / show / progress
1445/// notifications land in the per-server log ring.
1446fn handle_server_notification(
1447 instance: &crate::logging::InstanceKey,
1448 n: &Notification,
1449 diagnostics: &DiagnosticsBus,
1450 logger: &LspLogger,
1451 event_bus: Option<&Arc<lattice_runtime::EventBus>>,
1452) {
1453 let server_id = &instance.server_id;
1454 match n.method.as_str() {
1455 "window/logMessage" => {
1456 // LSP severity: 1=Error, 2=Warning, 3=Info, 4=Log/Debug.
1457 let (level, msg) = parse_window_message(&n.params);
1458 logger.log(Some(instance), level, LogSource::LspMessage, msg);
1459 }
1460 "window/showMessage" => {
1461 let (level, msg) = parse_window_message(&n.params);
1462 logger.log(Some(instance), level, LogSource::LspShowMessage, msg);
1463 }
1464 "$/progress" => {
1465 // Server-side work-done progress (LSP §3.16). Parse
1466 // the {token, value} envelope and publish a typed
1467 // `LspProgressUpdate` on the editor bus. The modeline
1468 // (and any plugin subscriber) accumulates by
1469 // (server_id, token).
1470 //
1471 // Logger still gets a Debug breadcrumb so the
1472 // `*lsp:<server>:trace*` ring keeps the raw record.
1473 logger.log(
1474 Some(instance),
1475 LogLevel::Debug,
1476 LogSource::Client,
1477 format!("$/progress: {}", compact_params(&n.params)),
1478 );
1479 if let (Some(bus), Some(update)) =
1480 (event_bus, parse_progress(server_id, n.params.as_ref()))
1481 {
1482 bus.publish_typed(update);
1483 }
1484 }
1485 "experimental/serverStatus" => {
1486 // L2: rust-analyzer readiness notification. `quiescent:
1487 // true` means the server finished indexing and features
1488 // (hover/diagnostics/completion) are reliable; `health`
1489 // reports ok/warning/error. Publish a typed event the
1490 // modeline turns into the ✓/⟳/✗ readiness glyph.
1491 logger.log(
1492 Some(instance),
1493 LogLevel::Debug,
1494 LogSource::Client,
1495 format!("experimental/serverStatus: {}", compact_params(&n.params)),
1496 );
1497 if let (Some(bus), Some(update)) =
1498 (event_bus, parse_server_status(server_id, n.params.as_ref()))
1499 {
1500 bus.publish_typed(update);
1501 }
1502 }
1503 "telemetry/event" => {
1504 // 4.4.a: distinct LogSource so plugin subscribers
1505 // on the typed event bus can filter
1506 // `source == "telemetry"` instead of parsing
1507 // free-form log text. Payload rides as the
1508 // compacted-JSON message tail; subscribers that
1509 // need structured access parse the suffix.
1510 logger.log(
1511 Some(instance),
1512 LogLevel::Debug,
1513 LogSource::Telemetry,
1514 compact_params(&n.params),
1515 );
1516 }
1517 "$/logTrace" => {
1518 // 4.4.b: server-emitted trace record. Shape:
1519 // `{ message: String, verbose: Option<String> }`.
1520 // Append both lines to the trace ring so the
1521 // `*lsp:<server>:trace*` buffer surfaces them; the
1522 // ring drops records by capacity, not by level, so
1523 // a verbose-mode session can produce a lot of data
1524 // -- that's intentional, the user opted in by
1525 // running `:lsp-trace`.
1526 let (message, verbose) = parse_log_trace(n.params.as_ref());
1527 logger.log(Some(instance), LogLevel::Trace, LogSource::Trace, message);
1528 if let Some(verbose) = verbose {
1529 logger.log(
1530 Some(instance),
1531 LogLevel::Trace,
1532 LogSource::Trace,
1533 format!(" {verbose}"),
1534 );
1535 }
1536 }
1537 "textDocument/publishDiagnostics" => {
1538 let params = match n.params.clone() {
1539 Some(v) => v,
1540 None => {
1541 logger.log(
1542 Some(instance),
1543 LogLevel::Warn,
1544 LogSource::Client,
1545 "publishDiagnostics with empty params",
1546 );
1547 return;
1548 }
1549 };
1550 match serde_json::from_value::<lsp_types::PublishDiagnosticsParams>(params) {
1551 Ok(p) => {
1552 let n_diags = p.diagnostics.len();
1553 let uri = p.uri.as_str().to_string();
1554 let event = DiagnosticEvent::from_lsp(Arc::clone(server_id), p);
1555 diagnostics.publish(event);
1556 logger.log(
1557 Some(instance),
1558 LogLevel::Debug,
1559 LogSource::Client,
1560 format!("publishDiagnostics: {n_diags} diag(s) for {uri}"),
1561 );
1562 }
1563 Err(e) => {
1564 logger.log(
1565 Some(instance),
1566 LogLevel::Warn,
1567 LogSource::Client,
1568 format!("publishDiagnostics deserialise failed: {e}"),
1569 );
1570 }
1571 }
1572 }
1573 other => {
1574 logger.log(
1575 Some(instance),
1576 LogLevel::Debug,
1577 LogSource::Client,
1578 format!("unhandled server notification: {other}"),
1579 );
1580 }
1581 }
1582}
1583
1584/// Parse a `$/progress` payload (LSP §3.16) into a typed
1585/// `LspProgressUpdate`. Returns `None` if the envelope is
1586/// missing fields or has the wrong shape — we'd rather drop
1587/// the update than publish a half-filled event.
1588///
1589/// Token can be number or string per spec; we serialise both
1590/// to `String` so the (server_id, token) accumulator key is
1591/// uniform.
1592fn parse_progress(
1593 server_id: &Arc<str>,
1594 params: Option<&Value>,
1595) -> Option<crate::events::LspProgressUpdate> {
1596 use crate::events::{LspProgressKind, LspProgressUpdate};
1597 let p = params?;
1598 let token = match p.get("token")? {
1599 Value::String(s) => s.clone(),
1600 Value::Number(n) => n.to_string(),
1601 _ => return None,
1602 };
1603 let value = p.get("value")?;
1604 let kind_str = value.get("kind")?.as_str()?;
1605 let kind = match kind_str {
1606 "begin" => LspProgressKind::Begin,
1607 "report" => LspProgressKind::Report,
1608 "end" => LspProgressKind::End,
1609 _ => return None,
1610 };
1611 let title = value
1612 .get("title")
1613 .and_then(Value::as_str)
1614 .map(str::to_owned);
1615 let message = value
1616 .get("message")
1617 .and_then(Value::as_str)
1618 .map(str::to_owned);
1619 let percentage = value
1620 .get("percentage")
1621 .and_then(Value::as_u64)
1622 .map(|n| n.min(100) as u32);
1623 let cancellable = value
1624 .get("cancellable")
1625 .and_then(Value::as_bool)
1626 .unwrap_or(false);
1627 Some(LspProgressUpdate {
1628 server_id: Arc::clone(server_id),
1629 token,
1630 kind,
1631 title,
1632 message,
1633 percentage,
1634 cancellable,
1635 })
1636}
1637
1638/// L2: parse an `experimental/serverStatus` params object
1639/// (rust-analyzer). Shape: `{ health: "ok"|"warning"|"error",
1640/// quiescent: bool, message?: string }`. Unknown / missing health
1641/// maps to `Ok` (treat as healthy); missing quiescent defaults to
1642/// `true` (assume ready rather than spin forever on a malformed
1643/// payload). Returns `None` only when params are absent.
1644fn parse_server_status(
1645 server_id: &Arc<str>,
1646 params: Option<&Value>,
1647) -> Option<crate::events::LspServerStatusChanged> {
1648 use crate::events::{LspServerHealth, LspServerStatusChanged};
1649 let p = params?;
1650 let health = match p.get("health").and_then(Value::as_str) {
1651 Some("error") => LspServerHealth::Error,
1652 Some("warning") => LspServerHealth::Warning,
1653 _ => LspServerHealth::Ok,
1654 };
1655 let quiescent = p.get("quiescent").and_then(Value::as_bool).unwrap_or(true);
1656 let message = p.get("message").and_then(Value::as_str).map(str::to_owned);
1657 Some(LspServerStatusChanged {
1658 server_id: Arc::clone(server_id),
1659 quiescent,
1660 health,
1661 message,
1662 })
1663}
1664
1665/// 4.4.b: parse a `$/logTrace` params object. Returns
1666/// `(message, verbose_opt)`. Fallback for malformed shape:
1667/// the compacted-JSON tail as the message and no verbose.
1668fn parse_log_trace(params: Option<&Value>) -> (String, Option<String>) {
1669 let Some(p) = params else {
1670 return ("<empty $/logTrace params>".to_string(), None);
1671 };
1672 let message = p
1673 .get("message")
1674 .and_then(Value::as_str)
1675 .map(str::to_owned)
1676 .unwrap_or_else(|| compact_params(&Some(p.clone())));
1677 let verbose = p.get("verbose").and_then(Value::as_str).map(str::to_owned);
1678 (message, verbose)
1679}
1680
1681/// Pull severity + message out of a `window/logMessage` /
1682/// `window/showMessage` params object. Defaults to Info /
1683/// "<unparseable>" on shape mismatch -- we never drop user
1684/// information silently.
1685fn parse_window_message(params: &Option<Value>) -> (LogLevel, String) {
1686 let Some(v) = params.as_ref() else {
1687 return (LogLevel::Info, "<empty params>".into());
1688 };
1689 let level = match v.get("type").and_then(|t| t.as_i64()) {
1690 Some(1) => LogLevel::Error,
1691 Some(2) => LogLevel::Warn,
1692 Some(3) => LogLevel::Info,
1693 Some(4) => LogLevel::Debug, // Log-class
1694 Some(5) => LogLevel::Debug, // LSP 3.18 Debug
1695 _ => LogLevel::Info,
1696 };
1697 let msg = v
1698 .get("message")
1699 .and_then(|m| m.as_str())
1700 .unwrap_or("<no message>")
1701 .to_string();
1702 (level, msg)
1703}
1704
1705/// Render a JSON value to a single-line compact string. Used
1706/// for the log records so multi-line server payloads don't
1707/// blow the buffer view's per-row layout.
1708fn compact_params(params: &Option<Value>) -> String {
1709 match params {
1710 Some(v) => serde_json::to_string(v).unwrap_or_else(|_| "<unprintable>".into()),
1711 None => String::new(),
1712 }
1713}
1714
1715/// Handle a `workspace/applyEdit` request asynchronously
1716/// (Phase 4.3). Parses the LSP params, dispatches them through
1717/// the apply-edit bus to the App's drain, awaits the App's
1718/// outcome via the embedded oneshot, and converts that into the
1719/// LSP `Response` body. Spec response shape:
1720/// `ApplyWorkspaceEditResponse { applied, failure_reason,
1721/// failed_change }`. We don't track `failed_change` today (the
1722/// per-file apply path is non-atomic), so it stays `None`.
1723///
1724/// Failure modes that surface as `applied: false`:
1725/// - The request params don't deserialize into
1726/// `ApplyWorkspaceEditParams`.
1727/// - The receiver dropped before the App could process the
1728/// edit (App is shutting down).
1729/// - The App reports `applied: false` with its own
1730/// `failure_reason`.
1731async fn handle_apply_edit_request(
1732 instance: crate::logging::InstanceKey,
1733 req_id: RequestId,
1734 params: Option<Value>,
1735 bus: &crate::apply_edit::ApplyEditBus,
1736 logger: &LspLogger,
1737) -> Response {
1738 let server_id = Arc::clone(&instance.server_id);
1739 let parsed: lsp_types::ApplyWorkspaceEditParams = match params {
1740 Some(v) => match serde_json::from_value(v) {
1741 Ok(p) => p,
1742 Err(e) => {
1743 logger.log(
1744 Some(&instance),
1745 LogLevel::Warn,
1746 LogSource::Client,
1747 format!("workspace/applyEdit: malformed params: {e}"),
1748 );
1749 return apply_edit_response(req_id, false, Some(format!("malformed params: {e}")));
1750 }
1751 },
1752 None => {
1753 return apply_edit_response(
1754 req_id,
1755 false,
1756 Some("workspace/applyEdit: missing params".into()),
1757 );
1758 }
1759 };
1760 let (response_tx, response_rx) = oneshot::channel();
1761 let inbound = crate::apply_edit::InboundApplyEdit {
1762 server_id: Arc::clone(&server_id),
1763 workspace: Arc::clone(&instance.workspace),
1764 label: parsed.label,
1765 edit: parsed.edit,
1766 response: response_tx,
1767 };
1768 // BC.8d: the apply-edit bus is now the generic `InboundBus` (host-drained
1769 // variant); `send` wakes the editor so the edit is applied off-keystroke —
1770 // same `Result<(), payload>` shape as the retired `ApplyEditBus::dispatch`.
1771 if bus.send(inbound).is_err() {
1772 return apply_edit_response(
1773 req_id,
1774 false,
1775 Some("client cannot apply edits (no receiver)".into()),
1776 );
1777 }
1778 match response_rx.await {
1779 Ok(outcome) => apply_edit_response(req_id, outcome.applied, outcome.failure_reason),
1780 Err(_) => apply_edit_response(
1781 req_id,
1782 false,
1783 Some("client did not respond before drop".into()),
1784 ),
1785 }
1786}
1787
1788/// Build an `ApplyWorkspaceEditResponse`-shaped LSP `Response`.
1789/// The `failed_change` field stays `None` -- atomic-rollback +
1790/// per-change failure indexing land alongside future
1791/// `apply_workspace_edit_atomic` work.
1792fn apply_edit_response(
1793 req_id: RequestId,
1794 applied: bool,
1795 failure_reason: Option<String>,
1796) -> Response {
1797 let body = lsp_types::ApplyWorkspaceEditResponse {
1798 applied,
1799 failure_reason,
1800 failed_change: None,
1801 };
1802 match serde_json::to_value(body) {
1803 Ok(v) => Response::ok(req_id, v),
1804 Err(e) => Response::err(
1805 req_id,
1806 crate::jsonrpc::ResponseError {
1807 code: crate::jsonrpc::error_codes::INTERNAL_ERROR,
1808 message: format!("encode response: {e}"),
1809 data: None,
1810 },
1811 ),
1812 }
1813}
1814
1815/// Handle a `workspace/configuration` request asynchronously
1816/// (Phase 4.1 follow-up). Parses the LSP `ConfigurationParams`,
1817/// dispatches each item's `section` through the configuration
1818/// bus to the App's drain, and converts the App's
1819/// `Vec<serde_json::Value>` reply into the spec-shaped
1820/// `Vec<Value>` response (one entry per requested item, in
1821/// input order; missing sections come back as `Value::Null`).
1822///
1823/// Failure modes that fall back to `[null, ...]` so the server
1824/// doesn't hang:
1825/// - Malformed params (logged as Warn).
1826/// - Receiver dropped before the App could respond.
1827async fn handle_configuration_request(
1828 instance: crate::logging::InstanceKey,
1829 req_id: RequestId,
1830 params: Option<Value>,
1831 bus: &crate::configuration::ConfigurationBus,
1832 logger: &LspLogger,
1833) -> Response {
1834 let server_id = Arc::clone(&instance.server_id);
1835 let parsed: lsp_types::ConfigurationParams = match params {
1836 Some(v) => match serde_json::from_value(v) {
1837 Ok(p) => p,
1838 Err(e) => {
1839 logger.log(
1840 Some(&instance),
1841 LogLevel::Warn,
1842 LogSource::Client,
1843 format!("workspace/configuration: malformed params: {e}"),
1844 );
1845 // Empty array reply -- the server's caller will
1846 // see "no items" rather than a hang.
1847 return Response::ok(req_id, Value::Array(Vec::new()));
1848 }
1849 },
1850 None => return Response::ok(req_id, Value::Array(Vec::new())),
1851 };
1852 let sections: Vec<String> = parsed
1853 .items
1854 .iter()
1855 .map(|i| i.section.clone().unwrap_or_default())
1856 .collect();
1857 let count = sections.len();
1858 let (response_tx, response_rx) = oneshot::channel();
1859 let inbound = crate::configuration::InboundConfigurationRequest {
1860 server_id: Arc::clone(&server_id),
1861 workspace: Arc::clone(&instance.workspace),
1862 sections,
1863 response: response_tx,
1864 };
1865 // BC.8b: the configuration bus is now the generic `InboundBus`; `send`
1866 // wakes the editor (off-keystroke reply) — same `Result<(), payload>` shape
1867 // as the retired `ConfigurationBus::dispatch`.
1868 if bus.send(inbound).is_err() {
1869 let arr: Vec<Value> = (0..count).map(|_| Value::Null).collect();
1870 return Response::ok(req_id, Value::Array(arr));
1871 }
1872 let values = match response_rx.await {
1873 Ok(v) => v,
1874 Err(_) => (0..count).map(|_| Value::Null).collect(),
1875 };
1876 Response::ok(req_id, Value::Array(values))
1877}
1878
1879/// 4.4.b: `window/showDocument` request handler. Parses the
1880/// LSP shape, dispatches via the bus, awaits the App's
1881/// outcome, and ferries `ShowDocumentResult { success }` back.
1882/// Falls back to `success: false` on malformed params or
1883/// receiver-dropped — the spec lets clients refuse and we
1884/// prefer that over hanging the server.
1885async fn handle_show_document_request(
1886 instance: crate::logging::InstanceKey,
1887 req_id: RequestId,
1888 params: Option<Value>,
1889 bus: &crate::show_document::ShowDocumentBus,
1890 logger: &LspLogger,
1891) -> Response {
1892 let server_id = Arc::clone(&instance.server_id);
1893 let parsed: lsp_types::ShowDocumentParams = match params {
1894 Some(v) => match serde_json::from_value(v) {
1895 Ok(p) => p,
1896 Err(e) => {
1897 logger.log(
1898 Some(&instance),
1899 LogLevel::Warn,
1900 LogSource::Client,
1901 format!("window/showDocument: malformed params: {e}"),
1902 );
1903 return show_document_response(req_id, false);
1904 }
1905 },
1906 None => return show_document_response(req_id, false),
1907 };
1908 let (response_tx, response_rx) = oneshot::channel();
1909 let inbound = crate::show_document::InboundShowDocument {
1910 server_id: Arc::clone(&server_id),
1911 workspace: Arc::clone(&instance.workspace),
1912 uri: parsed.uri,
1913 external: parsed.external.unwrap_or(false),
1914 take_focus: parsed.take_focus.unwrap_or(false),
1915 selection: parsed.selection,
1916 response: response_tx,
1917 };
1918 // BC.8c: the show-document bus is now the generic `InboundBus`; `send`
1919 // wakes the editor (off-keystroke reply) — same `Result<(), payload>`
1920 // shape as the retired `ShowDocumentBus::dispatch`.
1921 if bus.send(inbound).is_err() {
1922 return show_document_response(req_id, false);
1923 }
1924 match response_rx.await {
1925 Ok(outcome) => show_document_response(req_id, outcome.success),
1926 Err(_) => show_document_response(req_id, false),
1927 }
1928}
1929
1930fn show_document_response(req_id: RequestId, success: bool) -> Response {
1931 let body = lsp_types::ShowDocumentResult { success };
1932 match serde_json::to_value(body) {
1933 Ok(v) => Response::ok(req_id, v),
1934 Err(e) => Response::err(
1935 req_id,
1936 crate::jsonrpc::ResponseError {
1937 code: crate::jsonrpc::error_codes::INTERNAL_ERROR,
1938 message: format!("encode response: {e}"),
1939 data: None,
1940 },
1941 ),
1942 }
1943}
1944
1945/// 4.4.b: `window/showMessageRequest` handler. Same shape as
1946/// applyEdit + showDocument: parse, dispatch, await, ferry.
1947/// The reply body is either the selected `MessageActionItem`
1948/// (verbatim) or JSON `null` when the user dismissed without
1949/// picking.
1950async fn handle_show_message_request(
1951 instance: crate::logging::InstanceKey,
1952 req_id: RequestId,
1953 params: Option<Value>,
1954 bus: &crate::show_message_request::ShowMessageRequestBus,
1955 logger: &LspLogger,
1956) -> Response {
1957 let server_id = Arc::clone(&instance.server_id);
1958 let parsed: lsp_types::ShowMessageRequestParams = match params {
1959 Some(v) => match serde_json::from_value(v) {
1960 Ok(p) => p,
1961 Err(e) => {
1962 logger.log(
1963 Some(&instance),
1964 LogLevel::Warn,
1965 LogSource::Client,
1966 format!("window/showMessageRequest: malformed params: {e}"),
1967 );
1968 return Response::ok(req_id, Value::Null);
1969 }
1970 },
1971 None => return Response::ok(req_id, Value::Null),
1972 };
1973 let (response_tx, response_rx) = oneshot::channel();
1974 let inbound = crate::show_message_request::InboundShowMessageRequest {
1975 server_id: Arc::clone(&server_id),
1976 workspace: Arc::clone(&instance.workspace),
1977 level: parsed.typ,
1978 message: parsed.message,
1979 actions: parsed.actions.unwrap_or_default(),
1980 response: response_tx,
1981 };
1982 // BC.8e: the show-message-request bus is now the generic `InboundBus`
1983 // (host-drained variant); `send` wakes the editor so the picker is raised
1984 // off-keystroke — same `Result<(), payload>` shape as the retired
1985 // `ShowMessageRequestBus::dispatch`.
1986 if bus.send(inbound).is_err() {
1987 return Response::ok(req_id, Value::Null);
1988 }
1989 let selected = match response_rx.await {
1990 Ok(outcome) => outcome.selected,
1991 Err(_) => None,
1992 };
1993 match selected {
1994 Some(item) => match serde_json::to_value(item) {
1995 Ok(v) => Response::ok(req_id, v),
1996 Err(e) => Response::err(
1997 req_id,
1998 crate::jsonrpc::ResponseError {
1999 code: crate::jsonrpc::error_codes::INTERNAL_ERROR,
2000 message: format!("encode response: {e}"),
2001 data: None,
2002 },
2003 ),
2004 },
2005 None => Response::ok(req_id, Value::Null),
2006 }
2007}
2008
2009/// Handle a server-initiated request. Default behaviour is "we
2010/// don't implement that yet" (METHOD_NOT_FOUND); per-method
2011/// handlers replace this as features land.
2012fn handle_server_request(
2013 instance: &crate::logging::InstanceKey,
2014 req: &Request,
2015 logger: &LspLogger,
2016) -> Response {
2017 match req.method.as_str() {
2018 // `client/registerCapability` / `client/unregisterCapability`
2019 // are handled inline in `actor_main` (4.4.n) so they can
2020 // mutate the published capability snapshot. Reaching this
2021 // arm means the inline handler is missing a branch --
2022 // log it so we notice in CI.
2023 "client/registerCapability" | "client/unregisterCapability" => {
2024 logger.log(
2025 Some(instance),
2026 LogLevel::Warn,
2027 LogSource::Client,
2028 format!(
2029 "{} reached the inline-fallback handler -- registration dropped (bug)",
2030 req.method
2031 ),
2032 );
2033 Response::ok(req.id.clone(), Value::Null)
2034 }
2035 // Empty configuration -- the §5.12 typed-options layer
2036 // wires real values in later.
2037 "workspace/configuration" => {
2038 // params: {items: [{section: "...", scopeUri: "..."}, ...]}
2039 // Response: array with one entry per requested item.
2040 let n_items = req
2041 .params
2042 .as_ref()
2043 .and_then(|v| v.get("items"))
2044 .and_then(|v| v.as_array())
2045 .map(|a| a.len())
2046 .unwrap_or(0);
2047 let arr: Vec<Value> = (0..n_items).map(|_| Value::Null).collect();
2048 Response::ok(req.id.clone(), Value::Array(arr))
2049 }
2050 "window/workDoneProgress/create" => Response::ok(req.id.clone(), Value::Null),
2051 other => {
2052 logger.log(
2053 Some(instance),
2054 LogLevel::Warn,
2055 LogSource::Client,
2056 format!("server request unhandled: {other}"),
2057 );
2058 Response::err(
2059 req.id.clone(),
2060 crate::jsonrpc::ResponseError {
2061 code: crate::jsonrpc::error_codes::METHOD_NOT_FOUND,
2062 message: format!("client does not implement {other}"),
2063 data: None,
2064 },
2065 )
2066 }
2067 }
2068}
2069
2070/// Pre-handshake message handler: log, ignore. The server might
2071/// emit `window/logMessage` or `$/progress` before responding to
2072/// `initialize`; spec lets it.
2073fn handle_pre_handshake_message(
2074 instance: &crate::logging::InstanceKey,
2075 msg: Message,
2076 diagnostics: &DiagnosticsBus,
2077 logger: &LspLogger,
2078 event_bus: Option<&Arc<lattice_runtime::EventBus>>,
2079) {
2080 match msg {
2081 Message::Notification(n) => {
2082 handle_server_notification(instance, &n, diagnostics, logger, event_bus)
2083 }
2084 Message::Request(r) => {
2085 logger.log(
2086 Some(instance),
2087 LogLevel::Warn,
2088 LogSource::Client,
2089 format!(
2090 "server-initiated {} request before handshake -- ignored",
2091 r.method
2092 ),
2093 );
2094 }
2095 Message::Response(r) => {
2096 logger.log(
2097 Some(instance),
2098 LogLevel::Warn,
2099 LogSource::Client,
2100 format!("stray response before handshake (id {:?})", r.id),
2101 );
2102 }
2103 }
2104}
2105
2106/// Build the `initialize` params with our advertised capabilities.
2107fn build_initialize_params(
2108 workspace_folder_uri: Uri,
2109 workspace_name: String,
2110 initialization_options: Option<Value>,
2111) -> InitializeParams {
2112 let folder = WorkspaceFolder {
2113 uri: workspace_folder_uri.clone(),
2114 name: workspace_name,
2115 };
2116 #[allow(deprecated)]
2117 InitializeParams {
2118 process_id: Some(std::process::id()),
2119 // root_path / root_uri are deprecated but some servers
2120 // still read them. Set the URI for backward compat
2121 // (lsp-types still has the field) and leave root_path as
2122 // None.
2123 root_path: None,
2124 root_uri: Some(workspace_folder_uri.clone()),
2125 initialization_options,
2126 capabilities: capabilities::client_capabilities(),
2127 trace: Some(lsp_types::TraceValue::Off),
2128 workspace_folders: Some(vec![folder]),
2129 client_info: Some(ClientInfo {
2130 name: "lattice".into(),
2131 version: Some(env!("CARGO_PKG_VERSION").into()),
2132 }),
2133 locale: None,
2134 ..Default::default()
2135 }
2136}
2137
2138/// Run the LSP shutdown sequence end-to-end. Used when the
2139/// caller drops their handle without explicit `shutdown`.
2140/// `out_tx` is the channel to the `write_loop`; `child` is the
2141/// process so we can `wait` on its exit.
2142async fn perform_shutdown(
2143 out_tx: &mpsc::UnboundedSender<Message>,
2144 _in_rx: &mut mpsc::UnboundedReceiver<Message>,
2145 next_id: &mut u64,
2146 child: Option<&mut tokio::process::Child>,
2147) {
2148 let id = RequestId::from_u64(*next_id);
2149 *next_id += 1;
2150 let _ = out_tx.send(Message::Request(Request::new(id, "shutdown", None)));
2151 let _ = out_tx.send(Message::Notification(Notification::new("exit", None)));
2152 if let Some(c) = child {
2153 let _ = tokio::time::timeout(std::time::Duration::from_secs(2), c.wait()).await;
2154 }
2155}
2156
2157async fn write_loop<W>(
2158 writer: Arc<Mutex<LspWriter<W>>>,
2159 mut out_rx: mpsc::UnboundedReceiver<Message>,
2160 instance: crate::logging::InstanceKey,
2161 logger: LspLogger,
2162) where
2163 W: AsyncWrite + Unpin + Send,
2164{
2165 while let Some(msg) = out_rx.recv().await {
2166 // Trace interceptor: emit a Trace record before the
2167 // wire write iff trace mode is enabled for this
2168 // instance. `is_tracing` is a single HashSet lookup --
2169 // off path costs almost nothing.
2170 if logger.is_tracing(&instance) {
2171 logger.log(
2172 Some(&instance),
2173 LogLevel::Trace,
2174 LogSource::Trace,
2175 format!("→ {}", trace_render(&msg)),
2176 );
2177 }
2178 let mut w = writer.lock().await;
2179 if let Err(e) = w.write_message(&msg).await {
2180 logger.log(
2181 Some(&instance),
2182 LogLevel::Error,
2183 LogSource::Client,
2184 format!("write_loop terminating: {e}"),
2185 );
2186 break;
2187 }
2188 }
2189}
2190
2191async fn read_loop<R>(
2192 mut reader: LspReader<R>,
2193 in_tx: mpsc::UnboundedSender<Message>,
2194 instance: crate::logging::InstanceKey,
2195 logger: LspLogger,
2196) where
2197 R: AsyncBufRead + Unpin + Send,
2198{
2199 loop {
2200 match reader.read_message().await {
2201 Ok(Some(msg)) => {
2202 if logger.is_tracing(&instance) {
2203 logger.log(
2204 Some(&instance),
2205 LogLevel::Trace,
2206 LogSource::Trace,
2207 format!("← {}", trace_render(&msg)),
2208 );
2209 }
2210 if in_tx.send(msg).is_err() {
2211 // Actor task gone.
2212 return;
2213 }
2214 }
2215 Ok(None) => {
2216 logger.log(
2217 Some(&instance),
2218 LogLevel::Info,
2219 LogSource::Client,
2220 "server closed stdout cleanly",
2221 );
2222 return;
2223 }
2224 Err(e) => {
2225 logger.log(
2226 Some(&instance),
2227 LogLevel::Error,
2228 LogSource::Client,
2229 format!("read_loop terminating: {e}"),
2230 );
2231 return;
2232 }
2233 }
2234 }
2235}
2236
2237async fn stderr_drain(
2238 stderr: ChildStderr,
2239 instance: crate::logging::InstanceKey,
2240 logger: LspLogger,
2241) {
2242 use tokio::io::{AsyncBufReadExt, BufReader};
2243 let mut lines = BufReader::new(stderr).lines();
2244 while let Ok(Some(line)) = lines.next_line().await {
2245 logger.log(Some(&instance), LogLevel::Warn, LogSource::Stderr, line);
2246 }
2247}
2248
2249/// Render a Message as a compact one-line trace string. We use
2250/// the JSON-RPC kind + (for requests/responses) the id +
2251/// method, plus a truncated body. Cheap; runs only when trace
2252/// is on.
2253fn trace_render(msg: &Message) -> String {
2254 const MAX: usize = 240;
2255 let mut s = match msg {
2256 Message::Request(r) => format!("Request id={:?} method={}", r.id, r.method),
2257 Message::Notification(n) => format!("Notification method={}", n.method),
2258 Message::Response(r) => {
2259 if let Some(err) = r.error.as_ref() {
2260 format!("Response id={:?} ERR {} {}", r.id, err.code, err.message)
2261 } else {
2262 format!("Response id={:?} OK", r.id)
2263 }
2264 }
2265 };
2266 if let Ok(body) = msg.to_json() {
2267 let body_str = String::from_utf8_lossy(&body);
2268 if body_str.len() <= MAX - s.len().min(MAX) {
2269 s.push_str(" body=");
2270 s.push_str(&body_str);
2271 } else {
2272 s.push_str(" body=");
2273 s.push_str(&body_str[..MAX.saturating_sub(s.len() + 6)]);
2274 s.push_str("...");
2275 }
2276 }
2277 s
2278}
2279
2280#[cfg(test)]
2281mod uri_tests {
2282 use super::*;
2283 use std::path::PathBuf;
2284
2285 #[test]
2286 fn relative_path_promoted_to_absolute_uri() {
2287 // Bug: a relative path produced `file:///<path>` which the
2288 // server interprets as root-rooted, not as
2289 // "<cwd>/<path>". Result: rust-analyzer can't find the
2290 // file in its workspace and every hover / definition
2291 // returns null. Fix: `uri_from_path` calls
2292 // `std::path::absolute` first.
2293 let cwd = std::env::current_dir().expect("cwd");
2294 let rel = PathBuf::from("crates/lattice-core/src/buffer.rs");
2295 let uri = uri_from_path(&rel);
2296 let uri_str = uri.as_str();
2297 assert!(
2298 uri_str.starts_with("file://"),
2299 "uri must start with file://, got {uri_str:?}"
2300 );
2301 // The absolute prefix (cwd) must appear in the URI.
2302 let cwd_marker = cwd
2303 .to_string_lossy()
2304 .replace('\\', "/")
2305 .trim_start_matches('/')
2306 .to_string();
2307 assert!(
2308 uri_str.contains(&cwd_marker),
2309 "uri should contain absolute cwd; got {uri_str:?} cwd marker {cwd_marker:?}"
2310 );
2311 assert!(
2312 uri_str.ends_with("crates/lattice-core/src/buffer.rs"),
2313 "uri should end with the original relative path; got {uri_str:?}"
2314 );
2315 }
2316
2317 // Unix only: `/tmp/…` has no drive, so on Windows it is not absolute and
2318 // is — correctly — given one (`file:///D:/tmp/…`).
2319 #[cfg(unix)]
2320 #[test]
2321 fn absolute_path_unchanged_in_uri() {
2322 // Already-absolute paths should round-trip without extra
2323 // canonicalisation (no symlink resolution).
2324 let abs = PathBuf::from("/tmp/lattice-test/foo.rs");
2325 let uri = uri_from_path(&abs);
2326 assert_eq!(uri.as_str(), "file:///tmp/lattice-test/foo.rs");
2327 }
2328}
2329
2330#[cfg(test)]
2331mod progress_tests {
2332 use super::*;
2333 use crate::events::LspProgressKind;
2334 use serde_json::json;
2335
2336 fn sid() -> Arc<str> {
2337 Arc::from("rust")
2338 }
2339
2340 #[test]
2341 fn parses_begin_with_full_payload() {
2342 let params = json!({
2343 "token": "build-1",
2344 "value": {
2345 "kind": "begin",
2346 "title": "Building",
2347 "message": "compiling",
2348 "percentage": 12,
2349 "cancellable": true,
2350 }
2351 });
2352 let p = parse_progress(&sid(), Some(¶ms)).expect("begin parses");
2353 assert_eq!(&*p.server_id, "rust");
2354 assert_eq!(p.token, "build-1");
2355 assert_eq!(p.kind, LspProgressKind::Begin);
2356 assert_eq!(p.title.as_deref(), Some("Building"));
2357 assert_eq!(p.message.as_deref(), Some("compiling"));
2358 assert_eq!(p.percentage, Some(12));
2359 assert!(p.cancellable);
2360 }
2361
2362 #[test]
2363 fn parses_numeric_token() {
2364 // Per spec the token can be number or string; we serialise
2365 // either form to a String so the accumulator key stays
2366 // uniform.
2367 let params = json!({
2368 "token": 42,
2369 "value": { "kind": "end" }
2370 });
2371 let p = parse_progress(&sid(), Some(¶ms)).expect("end parses");
2372 assert_eq!(p.token, "42");
2373 assert_eq!(p.kind, LspProgressKind::End);
2374 }
2375
2376 #[test]
2377 fn rejects_unknown_kind() {
2378 let params = json!({
2379 "token": "x",
2380 "value": { "kind": "bogus" }
2381 });
2382 assert!(parse_progress(&sid(), Some(¶ms)).is_none());
2383 }
2384
2385 #[test]
2386 fn rejects_missing_value() {
2387 let params = json!({ "token": "x" });
2388 assert!(parse_progress(&sid(), Some(¶ms)).is_none());
2389 }
2390
2391 #[test]
2392 fn caps_percentage_at_100() {
2393 // Some servers report 0..=100, some over-report briefly;
2394 // we clamp so the modeline can't render `120%`.
2395 let params = json!({
2396 "token": "x",
2397 "value": { "kind": "report", "percentage": 150 }
2398 });
2399 let p = parse_progress(&sid(), Some(¶ms)).expect("report parses");
2400 assert_eq!(p.percentage, Some(100));
2401 }
2402}
2403
2404#[cfg(test)]
2405mod log_trace_tests {
2406 use super::*;
2407 use serde_json::json;
2408
2409 #[test]
2410 fn parses_message_and_optional_verbose() {
2411 let params = json!({
2412 "message": "rpc inbound",
2413 "verbose": "{\"jsonrpc\":\"2.0\",\"id\":1}",
2414 });
2415 let (msg, verbose) = parse_log_trace(Some(¶ms));
2416 assert_eq!(msg, "rpc inbound");
2417 assert_eq!(verbose.as_deref(), Some("{\"jsonrpc\":\"2.0\",\"id\":1}"));
2418 }
2419
2420 #[test]
2421 fn message_required_verbose_optional() {
2422 let params = json!({ "message": "step" });
2423 let (msg, verbose) = parse_log_trace(Some(¶ms));
2424 assert_eq!(msg, "step");
2425 assert!(verbose.is_none());
2426 }
2427
2428 #[test]
2429 fn falls_back_on_missing_message() {
2430 // Spec says message is required; for a malformed
2431 // payload we use the compacted JSON as the message so
2432 // the trace log still shows something useful.
2433 let params = json!({ "other": 1 });
2434 let (msg, _) = parse_log_trace(Some(¶ms));
2435 assert!(!msg.is_empty());
2436 }
2437}
2438
2439#[cfg(test)]
2440mod uri_path_tests {
2441 use super::without_drive_slash;
2442
2443 /// The Windows half of [`super::uri_to_path`], pinned on every platform.
2444 #[test]
2445 fn a_drive_letter_loses_its_leading_slash_on_windows_only() {
2446 assert_eq!(
2447 without_drive_slash("/C:/Users/me/a.rs".to_string(), true),
2448 "C:/Users/me/a.rs"
2449 );
2450 // Not a drive: an ordinary absolute path, and a UNC-ish one.
2451 assert_eq!(
2452 without_drive_slash("/home/me/a.rs".to_string(), true),
2453 "/home/me/a.rs"
2454 );
2455 assert_eq!(without_drive_slash("/1:/x".to_string(), true), "/1:/x");
2456 // Off Windows the same text is a real path and is left alone.
2457 assert_eq!(
2458 without_drive_slash("/C:/Users/me/a.rs".to_string(), false),
2459 "/C:/Users/me/a.rs"
2460 );
2461 }
2462}