supercode_harness/lsp.rs
1//! §2 module 28 `lsp` (COMPOSABLE-HARNESS-DESIGN.md line 478): "D1 LSP
2//! diagnostics in edit path + query tool" — this module ships the
3//! WEAKEST FORM that satisfies D1: server LIFECYCLE for a HANDFUL of
4//! user-configured language servers, and diagnostics surfaced in the
5//! edit/write TOOL RESULT via the shared D-5 write-path seam
6//! ([`crate::tools::WriteObserver`], P5-9/P5-11).
7//!
8//! # Honest, deliberate gaps (design §4.4 "LSP-in-edit at fleet scale")
9//! - **No auto-spawn/auto-download fleet.** opencode auto-provisions ~38
10//! language servers. This module only ever spawns a server the user
11//! EXPLICITLY configured under `[capabilities.lsp.servers.<name>]` — no
12//! network fetch, no bundled binaries, nothing runs that wasn't named in
13//! config.
14//! - **No `/find/symbol` symbol-indexing query tool.** A real symbol index
15//! (workspace/symbol, textDocument/definition, ...) is D-4-sized work
16//! this module does not attempt — shipping a `lsp.query` tool that
17//! silently no-ops would be a worse outcome than not shipping it (build
18//! brief: "if you expose any query surface, it must work or not
19//! exist"), so none is exposed. `[capabilities.lsp]` carries exactly two
20//! real, wired knobs: `enabled` and `servers` (plus the bounds
21//! `max_diagnostics`/`timeout_secs`) — no `query`/`symbols` key is ever
22//! parsed, so there is no declared-but-dead knob for either gap.
23//!
24//! # Wire protocol — hand-rolled, no new dependency
25//! LSP frames a JSON-RPC message behind a tiny HTTP-style header
26//! (`Content-Length: N\r\n\r\n<N bytes of JSON>`). That framing is a dozen
27//! lines over `tokio::io::AsyncBufReadExt`/`AsyncReadExt` — pulling in a
28//! dedicated `lsp-types`/`lsp-server` crate for it would be the heavy,
29//! over-built option for a module scoped to lifecycle + diagnostics over a
30//! handful of servers (no symbol index, no code actions, no incremental
31//! sync deltas — just `initialize`/`initialized`/`didOpen`/`didChange`/
32//! `publishDiagnostics`/`shutdown`/`exit`), so `write_message`/
33//! `read_message` below hand-roll it instead. `cargo deny check` has
34//! nothing new to license-audit as a result.
35//!
36//! # Process lifecycle (no orphaned language servers, incl. grandchildren)
37//! A server is a long-lived child process, spawned lazily
38//! ([`LspManager::diagnostics_after_write`], on the first write to a file
39//! extension it's configured for) and kept alive in [`LspManager`] for
40//! reuse across writes. A real configured server (`rust-analyzer`,
41//! `typescript-language-server`, `gopls`, ...) commonly spawns its OWN
42//! persistent worker subprocesses (a proc-macro/build server, `tsserver`,
43//! `go`, ...) — so `.kill_on_drop(true)`/`Child::start_kill` alone (which
44//! only ever signal the ONE directly-tracked pid) are not enough; this is
45//! the SAME grandchild-orphan class `crate::agent::kill_job_process_group`
46//! was built to close for background shell jobs (P5-6), and the fix here
47//! reuses that exact mechanism: `LspClient::spawn` puts the server in its
48//! OWN process group (`Command::process_group(0)`, unix), and
49//! `LspClient::kill` SIGKILLs the WHOLE group (`kill_process_group`),
50//! not just the leader — `.kill_on_drop(true)` remains as a second,
51//! independent backstop for the leader pid specifically. On non-unix
52//! targets, no portable process-group primitive is wired up (same posture
53//! as `kill_job_process_group`'s own `#[cfg(not(unix))]` arm) — this falls
54//! back to the pre-fix direct-child-only kill, a documented residual, not
55//! silently claimed fixed there.
56//!
57//! [`LspManager::kill_all_sync`] (this module's group-kill, above) is
58//! called from `impl Drop for crate::Agent` — the ONLY production teardown
59//! path today, provable/traceable rather than relying solely on
60//! `kill_on_drop(true)`'s implicit runtime behavior. [`LspManager::shutdown_all`]
61//! (a graceful LSP `shutdown`/`exit` handshake, letting a well-behaved
62//! server reap its own children before this module force-kills the group)
63//! is NOT wired into any automatic path — `Agent::run_loop` runs once PER
64//! TURN, not once per session, so calling it there would tear down and
65//! respawn a reused server every turn, defeating the "kept alive for reuse
66//! across writes" design above; there is no separate session-level
67//! clean-exit hook distinct from `Drop` in this codebase today. It remains
68//! available as public API (exercised directly by this module's own tests)
69//! for a caller that manages its own `Agent` lifecycle and wants to drain
70//! gracefully before dropping it — but nothing calls it automatically, and
71//! that is the honest, current state (not an aspirational claim about a
72//! code path that doesn't exist).
73
74use std::collections::HashMap;
75use std::path::{Path, PathBuf};
76use std::sync::Arc;
77use std::time::Duration;
78
79use serde_json::{json, Value};
80use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader};
81
82use crate::error::{Error, Result};
83
84/// Hardening cap (same rationale as `crate::mcp::MCP_MAX_RESPONSE_BYTES`):
85/// the largest single Content-Length-framed message this client will
86/// buffer before treating the server as hostile/broken and erroring out —
87/// bounds how much memory a misbehaving configured language server can
88/// force this process to allocate for one message.
89pub const LSP_MAX_MESSAGE_BYTES: usize = 16 * 1024 * 1024;
90
91/// Default cap on the number of diagnostics rendered into a single tool
92/// result (bounded-context requirement — build brief: "a flood mustn't
93/// blow context"). Overridable via `[capabilities.lsp] max_diagnostics`.
94pub const DEFAULT_LSP_MAX_DIAGNOSTICS: usize = 20;
95
96/// Default wait for a configured server to publish diagnostics after a
97/// `didOpen`/`didChange` before giving up gracefully (never blocking the
98/// tool call indefinitely). Overridable via `[capabilities.lsp] timeout_secs`.
99pub const DEFAULT_LSP_TIMEOUT_SECS: u64 = 5;
100
101/// One `[capabilities.lsp.servers.<name>]` entry — a user-configured
102/// language server this module is allowed to spawn. `command`/`args` are
103/// config-borne code execution (D-10) — see `crate::configfile::sanitize_for_project`
104/// / `crate::userconfig`'s project-strip, which refuses this table from an
105/// untrusted project layer exactly like `hooks`/`mcp.servers`.
106#[derive(Debug, Clone, PartialEq, Eq)]
107pub struct LspServerSpec {
108 /// The executable to spawn (searched on `PATH` like any `Command::new`).
109 pub command: String,
110 /// Extra arguments passed to `command`.
111 pub args: Vec<String>,
112 /// File extensions (with or without a leading `.`, matched case-
113 /// insensitively) this server handles — a write to a matching path
114 /// lazily spawns (or reuses) this server.
115 pub extensions: Vec<String>,
116}
117
118/// One diagnostic surfaced from a server's `textDocument/publishDiagnostics`
119/// notification — the small subset of the LSP `Diagnostic` shape this
120/// module renders into a tool result (no `code`/`source`/`relatedInformation`
121/// — weakest form).
122#[derive(Debug, Clone, PartialEq, Eq)]
123struct DiagnosticEntry {
124 severity: &'static str,
125 line: u32,
126 character: u32,
127 message: String,
128}
129
130/// Per-connection JSON-RPC-over-stdio state — bundled behind ONE
131/// `tokio::sync::Mutex` (rather than separate locks for stdin/stdout) so a
132/// full request/response (or notify+wait-for-push) cycle runs atomically:
133/// two concurrent writes to files the SAME server handles can never
134/// interleave their reads and steal each other's response/diagnostics.
135struct LspIo {
136 stdin: tokio::process::ChildStdin,
137 stdout: BufReader<tokio::process::ChildStdout>,
138 next_id: i64,
139 /// uri -> last-sent document version (LSP full-text sync: `didOpen`
140 /// sends version 1, every subsequent `didChange` increments it).
141 opened: HashMap<String, i64>,
142 /// Whether `initialize`/`initialized` has completed on this connection.
143 initialized: bool,
144}
145
146/// A live connection to one configured language server — one spawned child
147/// process, kept alive for reuse across writes to files it handles.
148#[derive(Debug)]
149pub struct LspClient {
150 name: String,
151 // Sync `Mutex` (not `tokio::sync::Mutex`): only ever touched via the
152 // synchronous `start_kill`/`try_wait` (never awaited while held) — see
153 // `crate::agent::kill_job_process_group`'s identical precedent, which
154 // is exactly why this can be called from the non-async `impl Drop for
155 // Agent::drop`.
156 child: std::sync::Mutex<tokio::process::Child>,
157 io: tokio::sync::Mutex<LspIo>,
158}
159
160impl std::fmt::Debug for LspIo {
161 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
162 f.debug_struct("LspIo")
163 .field("next_id", &self.next_id)
164 .field("opened", &self.opened.keys().collect::<Vec<_>>())
165 .field("initialized", &self.initialized)
166 .finish()
167 }
168}
169
170/// Write one JSON-RPC value with LSP's `Content-Length` header framing.
171async fn write_message(stdin: &mut tokio::process::ChildStdin, val: &Value) -> Result<()> {
172 let body = serde_json::to_vec(val).map_err(|e| Error::tool("lsp", format!("encode: {e}")))?;
173 let header = format!("Content-Length: {}\r\n\r\n", body.len());
174 stdin
175 .write_all(header.as_bytes())
176 .await
177 .map_err(|e| Error::tool("lsp", format!("write: {e}")))?;
178 stdin
179 .write_all(&body)
180 .await
181 .map_err(|e| Error::tool("lsp", format!("write: {e}")))?;
182 stdin
183 .flush()
184 .await
185 .map_err(|e| Error::tool("lsp", format!("flush: {e}")))?;
186 Ok(())
187}
188
189/// Read one JSON-RPC value framed with LSP's `Content-Length` header —
190/// bounded by [`LSP_MAX_MESSAGE_BYTES`], and errors (rather than hangs) on
191/// EOF (the server exited or closed its stdout).
192async fn read_message(stdout: &mut BufReader<tokio::process::ChildStdout>) -> Result<Value> {
193 let mut content_length: Option<usize> = None;
194 loop {
195 let mut line = String::new();
196 let n = stdout
197 .read_line(&mut line)
198 .await
199 .map_err(|e| Error::tool("lsp", format!("read: {e}")))?;
200 if n == 0 {
201 return Err(Error::tool("lsp", "server closed stdout (eof)"));
202 }
203 let trimmed = line.trim_end_matches(['\r', '\n']);
204 if trimmed.is_empty() {
205 break; // blank line ends the header block
206 }
207 if let Some(v) = trimmed.strip_prefix("Content-Length:") {
208 content_length = v.trim().parse().ok();
209 }
210 // Any other header (e.g. `Content-Type:`) is read and ignored.
211 }
212 let len = content_length
213 .ok_or_else(|| Error::tool("lsp", "message missing Content-Length header"))?;
214 if len > LSP_MAX_MESSAGE_BYTES {
215 return Err(Error::tool(
216 "lsp",
217 format!("message too large ({len} bytes) — refusing to buffer"),
218 ));
219 }
220 let mut buf = vec![0u8; len];
221 stdout
222 .read_exact(&mut buf)
223 .await
224 .map_err(|e| Error::tool("lsp", format!("read body: {e}")))?;
225 serde_json::from_slice(&buf).map_err(|e| Error::tool("lsp", format!("decode: {e}")))
226}
227
228/// `file://` URI for `path` — a minimal, deterministic encoding (no
229/// percent-escaping beyond backslash normalization) sufficient for the
230/// stdio-local servers this module targets; every server this module talks
231/// to is a local child process reading the SAME literal path this process
232/// resolved, so round-trip fidelity (not RFC 3986 completeness) is what
233/// matters.
234fn path_to_uri(path: &Path) -> String {
235 let s = path.to_string_lossy().replace('\\', "/");
236 if let Some(stripped) = s.strip_prefix('/') {
237 format!("file:///{stripped}")
238 } else {
239 format!("file:///{s}")
240 }
241}
242
243/// Best-effort LSP `languageId` for `path`'s extension — covers the common
244/// languages a configured server would plausibly handle; unrecognized
245/// extensions fall back to `"plaintext"` (a server that cares can still use
246/// the extension embedded in the uri).
247fn language_id_for(path: &Path) -> &'static str {
248 match path
249 .extension()
250 .and_then(|e| e.to_str())
251 .unwrap_or_default()
252 .to_ascii_lowercase()
253 .as_str()
254 {
255 "rs" => "rust",
256 "py" => "python",
257 "js" | "mjs" | "cjs" => "javascript",
258 "jsx" => "javascriptreact",
259 "ts" | "mts" | "cts" => "typescript",
260 "tsx" => "typescriptreact",
261 "go" => "go",
262 "rb" => "ruby",
263 "java" => "java",
264 "c" | "h" => "c",
265 "cpp" | "cc" | "cxx" | "hpp" => "cpp",
266 "cs" => "csharp",
267 "json" => "json",
268 "toml" => "toml",
269 "yaml" | "yml" => "yaml",
270 "md" => "markdown",
271 "sh" | "bash" => "shellscript",
272 _ => "plaintext",
273 }
274}
275
276/// SIGKILL an entire process group — the shared grandchild-orphan-fix
277/// primitive, same mechanism as `crate::agent::kill_job_process_group`
278/// (P5-6). `pid` must be a process-group LEADER's pid (i.e. the process was
279/// spawned with `Command::process_group(0)`, making its pgid equal its own
280/// pid) for `-(pid)` to address the whole group rather than just the
281/// leader. Reused by `crate::formatters` for the same reason — see
282/// `crate::formatters::run_formatter`'s timeout arm.
283///
284/// SAFETY: `libc::kill` with a negative pid is `killpg` — it only ever
285/// sends a signal (never dereferences memory), so this is safe regardless
286/// of whether the group is still alive; a group that already exited yields
287/// `ESRCH`, a documented no-op, not an error worth surfacing.
288#[cfg(unix)]
289pub(crate) fn kill_process_group(pid: u32) {
290 unsafe {
291 libc::kill(-(pid as libc::pid_t), libc::SIGKILL);
292 }
293}
294
295fn severity_label(sev: Option<i64>) -> &'static str {
296 match sev {
297 Some(1) => "error",
298 Some(2) => "warning",
299 Some(3) => "information",
300 Some(4) => "hint",
301 _ => "diagnostic",
302 }
303}
304
305/// `true` iff `msg` is a `textDocument/publishDiagnostics` notification for
306/// `uri` specifically (a server may publish for OTHER files it opened as
307/// part of project analysis — those are ignored, never mixed into this
308/// write's result).
309fn is_publish_diagnostics_for(msg: &Value, uri: &str) -> bool {
310 msg.get("method").and_then(|m| m.as_str()) == Some("textDocument/publishDiagnostics")
311 && msg
312 .get("params")
313 .and_then(|p| p.get("uri"))
314 .and_then(|u| u.as_str())
315 == Some(uri)
316}
317
318/// Parse a `publishDiagnostics` notification's `diagnostics` array into
319/// `(true_total_count, bounded_entries)` — `true_total_count` may exceed
320/// `entries.len()` when the server reported more than `cap`, so the caller
321/// can render an honest "N more not shown" instead of silently dropping
322/// them.
323fn diagnostics_from_message(msg: &Value, cap: usize) -> (usize, Vec<DiagnosticEntry>) {
324 let Some(arr) = msg
325 .get("params")
326 .and_then(|p| p.get("diagnostics"))
327 .and_then(|d| d.as_array())
328 else {
329 return (0, Vec::new());
330 };
331 let total = arr.len();
332 let entries = arr
333 .iter()
334 .take(cap)
335 .map(|d| DiagnosticEntry {
336 severity: severity_label(d.get("severity").and_then(|s| s.as_i64())),
337 line: d
338 .get("range")
339 .and_then(|r| r.get("start"))
340 .and_then(|s| s.get("line"))
341 .and_then(|v| v.as_u64())
342 .unwrap_or(0) as u32,
343 character: d
344 .get("range")
345 .and_then(|r| r.get("start"))
346 .and_then(|s| s.get("character"))
347 .and_then(|v| v.as_u64())
348 .unwrap_or(0) as u32,
349 message: d
350 .get("message")
351 .and_then(|m| m.as_str())
352 .unwrap_or_default()
353 .to_string(),
354 })
355 .collect();
356 (total, entries)
357}
358
359/// Max characters kept from a single diagnostic's message text — a
360/// pathological server can't blow context with one enormous message
361/// either, on top of the diagnostics-COUNT cap.
362const MAX_DIAGNOSTIC_MESSAGE_CHARS: usize = 400;
363
364fn format_diagnostics(
365 server: &str,
366 path: &Path,
367 total: usize,
368 entries: &[DiagnosticEntry],
369) -> String {
370 let mut out = format!(
371 "LSP diagnostics ({server}) for {}: {total} issue(s)",
372 path.display()
373 );
374 for d in entries {
375 let mut msg = d.message.clone();
376 if msg.chars().count() > MAX_DIAGNOSTIC_MESSAGE_CHARS {
377 msg = msg.chars().take(MAX_DIAGNOSTIC_MESSAGE_CHARS).collect();
378 msg.push('\u{2026}');
379 }
380 out.push_str(&format!(
381 "\n {}:{}: {}: {}",
382 d.line + 1,
383 d.character + 1,
384 d.severity,
385 msg
386 ));
387 }
388 if total > entries.len() {
389 out.push_str(&format!(
390 "\n ... and {} more diagnostic(s) not shown",
391 total - entries.len()
392 ));
393 }
394 out
395}
396
397impl LspClient {
398 /// Spawn `spec.command` and hold it open, uninitialized (the LSP
399 /// `initialize` handshake happens lazily, under the SAME `io` lock as
400 /// the first `didOpen`, in [`Self::open_or_change_and_diagnose`]).
401 /// `.kill_on_drop(true)` mirrors `crate::mcp::McpClient::connect`'s own
402 /// stdio precedent — reaps the server if this client is ever dropped
403 /// without an explicit [`Self::kill`]/[`Self::shutdown`].
404 ///
405 /// `.process_group(0)` (unix) puts the server in its OWN new process
406 /// group (pgid == its own pid) — see the module doc's "Process
407 /// lifecycle" section and `crate::agent::kill_job_process_group` for
408 /// why: it's what lets [`Self::kill`] SIGKILL the server's WORKER
409 /// grandchildren (proc-macro/build servers, `tsserver`, `go`, ...) too,
410 /// not just this one directly-tracked pid.
411 async fn spawn(name: &str, spec: &LspServerSpec) -> Result<Arc<LspClient>> {
412 let mut cmd = tokio::process::Command::new(&spec.command);
413 cmd.args(&spec.args)
414 .stdin(std::process::Stdio::piped())
415 .stdout(std::process::Stdio::piped())
416 .stderr(std::process::Stdio::null())
417 .kill_on_drop(true);
418 #[cfg(unix)]
419 cmd.process_group(0);
420 let mut child = cmd
421 .spawn()
422 .map_err(|e| Error::tool("lsp", format!("spawn {}: {e}", spec.command)))?;
423 let stdin = child
424 .stdin
425 .take()
426 .ok_or_else(|| Error::tool("lsp", "no stdin"))?;
427 let stdout = BufReader::new(
428 child
429 .stdout
430 .take()
431 .ok_or_else(|| Error::tool("lsp", "no stdout"))?,
432 );
433 Ok(Arc::new(LspClient {
434 name: name.to_string(),
435 child: std::sync::Mutex::new(child),
436 io: tokio::sync::Mutex::new(LspIo {
437 stdin,
438 stdout,
439 next_id: 0,
440 opened: HashMap::new(),
441 initialized: false,
442 }),
443 }))
444 }
445
446 /// The full D1 cycle for one write: ensure `initialize`/`initialized`
447 /// has happened once, send `didOpen` (first touch of this uri) or
448 /// `didChange` (subsequent touches, full-text sync), then read
449 /// messages until this uri's `publishDiagnostics` notification arrives
450 /// or `timeout` elapses. Holds the `io` lock for the whole cycle —
451 /// see [`LspIo`]'s doc comment for why that's required, not just
452 /// convenient.
453 async fn open_or_change_and_diagnose(
454 &self,
455 root: &Path,
456 uri: &str,
457 text: &str,
458 language_id: &str,
459 timeout: Duration,
460 cap: usize,
461 ) -> Result<(usize, Vec<DiagnosticEntry>)> {
462 let mut io = self.io.lock().await;
463
464 if !io.initialized {
465 let id = io.next_id;
466 io.next_id += 1;
467 let req = json!({
468 "jsonrpc": "2.0",
469 "id": id,
470 "method": "initialize",
471 "params": {
472 "processId": std::process::id(),
473 "rootUri": path_to_uri(root),
474 "capabilities": {},
475 }
476 });
477 write_message(&mut io.stdin, &req).await?;
478 let deadline = tokio::time::Instant::now() + timeout;
479 loop {
480 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
481 if remaining.is_zero() {
482 return Err(Error::tool(
483 "lsp",
484 "timed out waiting for initialize response",
485 ));
486 }
487 let msg = tokio::time::timeout(remaining, read_message(&mut io.stdout))
488 .await
489 .map_err(|_| {
490 Error::tool("lsp", "timed out waiting for initialize response")
491 })??;
492 if msg.get("id").and_then(|v| v.as_i64()) == Some(id) {
493 break; // the initialize response — ignore any notification before it
494 }
495 }
496 let notif = json!({"jsonrpc": "2.0", "method": "initialized", "params": {}});
497 write_message(&mut io.stdin, ¬if).await?;
498 io.initialized = true;
499 }
500
501 let msg = match io.opened.get(uri).copied() {
502 Some(version) => {
503 let next = version + 1;
504 io.opened.insert(uri.to_string(), next);
505 json!({
506 "jsonrpc": "2.0",
507 "method": "textDocument/didChange",
508 "params": {
509 "textDocument": {"uri": uri, "version": next},
510 "contentChanges": [{"text": text}],
511 }
512 })
513 }
514 None => {
515 io.opened.insert(uri.to_string(), 1);
516 json!({
517 "jsonrpc": "2.0",
518 "method": "textDocument/didOpen",
519 "params": {
520 "textDocument": {
521 "uri": uri,
522 "languageId": language_id,
523 "version": 1,
524 "text": text,
525 }
526 }
527 })
528 }
529 };
530 write_message(&mut io.stdin, &msg).await?;
531
532 let deadline = tokio::time::Instant::now() + timeout;
533 loop {
534 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
535 if remaining.is_zero() {
536 return Err(Error::tool("lsp", "timed out waiting for diagnostics"));
537 }
538 let msg = tokio::time::timeout(remaining, read_message(&mut io.stdout))
539 .await
540 .map_err(|_| Error::tool("lsp", "timed out waiting for diagnostics"))??;
541 if is_publish_diagnostics_for(&msg, uri) {
542 return Ok(diagnostics_from_message(&msg, cap));
543 }
544 // Some other notification (or a server-initiated request this
545 // weakest-form client doesn't answer) — ignore and keep
546 // waiting for OUR diagnostics, bounded by `deadline` above.
547 }
548 }
549
550 /// Real, synchronous OS-process kill — callable from the non-async
551 /// `impl Drop for Agent` (same reasoning as
552 /// `crate::agent::kill_job_process_group`). SIGKILLs the server's WHOLE
553 /// process group (unix) — see [`Self::spawn`]'s `.process_group(0)` and
554 /// `kill_process_group` — so worker grandchildren die too, not just
555 /// this directly-tracked pid; falls back to the direct-child-only kill
556 /// on non-unix. A no-op if the process already exited.
557 fn kill(&self) {
558 if let Ok(mut child) = self.child.lock() {
559 #[cfg(unix)]
560 if let Some(pid) = child.id() {
561 kill_process_group(pid);
562 }
563 let _ = child.start_kill();
564 let _ = child.try_wait();
565 }
566 }
567
568 /// Graceful shutdown: `shutdown` request, brief wait for its response,
569 /// `exit` notification — then [`Self::kill`] unconditionally as a
570 /// backstop (a server that ignores `exit` must not linger). Bounded to
571 /// 2s total so a hung server can never block session teardown.
572 async fn shutdown(&self) {
573 let graceful = async {
574 let mut io = self.io.lock().await;
575 if !io.initialized {
576 return; // never handshaked — nothing to shut down gracefully
577 }
578 let id = io.next_id;
579 io.next_id += 1;
580 let req = json!({"jsonrpc": "2.0", "id": id, "method": "shutdown", "params": null});
581 if write_message(&mut io.stdin, &req).await.is_ok() {
582 let _ =
583 tokio::time::timeout(Duration::from_millis(500), read_message(&mut io.stdout))
584 .await;
585 }
586 let notif = json!({"jsonrpc": "2.0", "method": "exit", "params": null});
587 let _ = write_message(&mut io.stdin, ¬if).await;
588 };
589 let _ = tokio::time::timeout(Duration::from_secs(2), graceful).await;
590 self.kill();
591 }
592}
593
594/// Session-scoped registry of configured language servers and the live
595/// connections spawned so far — the [`crate::tools::WriteObserver`] this
596/// module installs ([`LspDiagnosticsObserver`]) is a thin wrapper around a
597/// shared `Arc<LspManager>`; `crate::agent::build_tool_context` keeps its
598/// own `Arc` clone too, so `crate::Agent`'s `Drop` impl can reach
599/// [`Self::kill_all_sync`] regardless of how many observer clones exist.
600#[derive(Debug)]
601pub struct LspManager {
602 servers: std::sync::Mutex<HashMap<String, Arc<LspClient>>>,
603 specs: Vec<(String, LspServerSpec)>,
604 root: PathBuf,
605 timeout: Duration,
606 max_diagnostics: usize,
607}
608
609impl LspManager {
610 /// Build a manager over `specs` (name -> server definition), rooted at
611 /// `root` (used as the LSP `rootUri` and the containment floor every
612 /// touched path is checked against via `crate::safe_path::contained`).
613 pub fn new(
614 root: PathBuf,
615 specs: Vec<(String, LspServerSpec)>,
616 timeout: Duration,
617 max_diagnostics: usize,
618 ) -> Self {
619 LspManager {
620 servers: std::sync::Mutex::new(HashMap::new()),
621 specs,
622 root,
623 timeout,
624 max_diagnostics,
625 }
626 }
627
628 fn spec_for_extension(&self, path: &Path) -> Option<(String, LspServerSpec)> {
629 let ext = path
630 .extension()
631 .and_then(|e| e.to_str())?
632 .to_ascii_lowercase();
633 self.specs
634 .iter()
635 .find(|(_, s)| {
636 s.extensions
637 .iter()
638 .any(|e| e.trim_start_matches('.').to_ascii_lowercase() == ext)
639 })
640 .cloned()
641 }
642
643 async fn get_or_spawn(&self, name: &str, spec: &LspServerSpec) -> Result<Arc<LspClient>> {
644 if let Some(existing) = self.servers.lock().unwrap().get(name).cloned() {
645 return Ok(existing);
646 }
647 let client = LspClient::spawn(name, spec).await?;
648 let mut map = self.servers.lock().unwrap();
649 // A concurrent write to another file this SAME server handles may
650 // have raced this spawn and already inserted — keep whichever
651 // landed first (the loser's freshly-spawned child is dropped here,
652 // which reaps it via `kill_on_drop`, never leaked).
653 let winner = map.entry(name.to_string()).or_insert(client).clone();
654 Ok(winner)
655 }
656
657 /// D1: the write-path diagnostics hook — spawns (or reuses) the
658 /// configured server for `path`'s extension, if any, sends
659 /// `didOpen`/`didChange` with the file's CURRENT on-disk content (the
660 /// caller — `LspDiagnosticsObserver::after_write` — only calls this
661 /// once the write has completed, so this always reads the FINAL bytes,
662 /// which for `[capabilities.formatters]` also on means the FORMATTED
663 /// content, not the model's pre-format draft — see the observer
664 /// ordering `crate::agent::build_tool_context` installs), and waits
665 /// (bounded by `self.timeout`) for that file's diagnostics. Never
666 /// fails the caller: every error (no configured server, spawn
667 /// failure, protocol error, timeout) degrades to `None`, logged once
668 /// via `tracing::warn!`.
669 pub async fn diagnostics_after_write(&self, path: &Path) -> Option<String> {
670 if !crate::safe_path::contained(&self.root, path) {
671 return None; // out of this module's scope — never touch outside the project
672 }
673 let (name, spec) = self.spec_for_extension(path)?;
674 let client = match self.get_or_spawn(&name, &spec).await {
675 Ok(c) => c,
676 Err(e) => {
677 tracing::warn!(server = %name, "lsp: failed to spawn/reuse server: {e}");
678 return None;
679 }
680 };
681 let text = match tokio::fs::read_to_string(path).await {
682 Ok(t) => t,
683 Err(_) => return None, // deleted/unreadable — nothing to diagnose
684 };
685 let uri = path_to_uri(path);
686 let language_id = language_id_for(path);
687 match client
688 .open_or_change_and_diagnose(
689 &self.root,
690 &uri,
691 &text,
692 language_id,
693 self.timeout,
694 self.max_diagnostics,
695 )
696 .await
697 {
698 Ok((total, entries)) if total > 0 => {
699 Some(format_diagnostics(&client.name, path, total, &entries))
700 }
701 Ok(_) => None, // clean file — no noise on every write
702 Err(e) => {
703 tracing::warn!(server = %name, "lsp: diagnostics unavailable: {e}");
704 None
705 }
706 }
707 }
708
709 /// Real, synchronous, provable kill of every server this manager has
710 /// spawned so far — called from `impl Drop for crate::Agent` (a
711 /// non-async context, hence the sync signature). A no-op for any
712 /// server already exited.
713 pub fn kill_all_sync(&self) {
714 let drained: Vec<Arc<LspClient>> = self
715 .servers
716 .lock()
717 .map(|mut m| m.drain().map(|(_, c)| c).collect())
718 .unwrap_or_default();
719 for client in drained {
720 client.kill();
721 }
722 }
723
724 /// Graceful async shutdown of every server this manager has spawned —
725 /// `shutdown`/`exit` handshake per server, `kill_all_sync`-equivalent
726 /// backstop applied per-client by `LspClient::shutdown` itself. Use
727 /// this on a clean-exit path that can afford to `.await`; use
728 /// [`Self::kill_all_sync`] from `Drop`.
729 pub async fn shutdown_all(&self) {
730 let drained: Vec<Arc<LspClient>> = self
731 .servers
732 .lock()
733 .map(|mut m| m.drain().map(|(_, c)| c).collect())
734 .unwrap_or_default();
735 for client in drained {
736 client.shutdown().await;
737 }
738 }
739
740 /// Test/observability hook: how many servers are currently live.
741 pub fn running_server_count(&self) -> usize {
742 self.servers.lock().map(|m| m.len()).unwrap_or(0)
743 }
744}
745
746/// The [`crate::tools::WriteObserver`] `[capabilities.lsp]` installs —
747/// `before_write` is a true no-op (LSP has nothing to capture before a
748/// mutation); `after_write` delegates straight to
749/// [`LspManager::diagnostics_after_write`].
750#[derive(Debug)]
751pub struct LspDiagnosticsObserver {
752 manager: Arc<LspManager>,
753}
754
755impl LspDiagnosticsObserver {
756 /// Wrap `manager` as a [`crate::tools::WriteObserver`].
757 pub fn new(manager: Arc<LspManager>) -> Self {
758 LspDiagnosticsObserver { manager }
759 }
760}
761
762#[async_trait::async_trait]
763impl crate::tools::WriteObserver for LspDiagnosticsObserver {
764 async fn before_write(&self, _path: &Path) {}
765 async fn after_write(&self, path: &Path) -> Option<String> {
766 self.manager.diagnostics_after_write(path).await
767 }
768}
769
770/// Build the [`LspManager`] a fresh [`crate::Agent`] should install, given
771/// a resolved [`crate::Config`] — called once, from
772/// `crate::agent::build_tool_context`. `Config::lsp_enabled` is the ONE
773/// gate: `false` (the default) returns `None` WITHOUT spawning anything —
774/// the default-off byte-identity guarantee. `true` with an EMPTY
775/// `Config::lsp_servers` still returns a (harmless, does-nothing) manager,
776/// but warns once — an enabled module with no configured servers is very
777/// likely a config mistake, not silent-by-design.
778pub fn manager_for_config(config: &crate::Config) -> Option<Arc<LspManager>> {
779 if !config.lsp_enabled {
780 return None;
781 }
782 if config.lsp_servers.is_empty() {
783 eprintln!(
784 "warning: [capabilities.lsp] is enabled but no servers are configured under \
785 [capabilities.lsp.servers.<name>] — no language server will ever be launched"
786 );
787 }
788 Some(Arc::new(LspManager::new(
789 config.cwd.clone(),
790 config.lsp_servers.clone(),
791 Duration::from_secs(config.lsp_timeout_secs.max(1)),
792 config.lsp_max_diagnostics.max(1),
793 )))
794}