Skip to main content

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 &params.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(&params)).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(&params)).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(&params)).is_none());
2383    }
2384
2385    #[test]
2386    fn rejects_missing_value() {
2387        let params = json!({ "token": "x" });
2388        assert!(parse_progress(&sid(), Some(&params)).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(&params)).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(&params));
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(&params));
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(&params));
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}