use crate::{
agent::cancellation::AgentCancellation,
config::LspSettings,
lsp::{
MAX_REFERENCES, MAX_TOOL_DIAGNOSTICS,
client::LspClient,
diagnostics::{DiagnosticsSummary, format_diagnostics_tool, format_injected_diagnostics},
queries::{format_references, one_based_to_lsp},
sync::{LanguageRoute, read_text_for_sync, route_for_path},
},
output::redact_sensitive_text,
};
use std::{
collections::BTreeMap,
path::{Path, PathBuf},
sync::{Arc, Mutex},
time::{Duration, Instant},
};
const MAX_RESPAWNS_PER_SESSION: u8 = 2;
pub(crate) struct LspManager {
settings: LspSettings,
workspace_root: PathBuf,
servers: Mutex<BTreeMap<String, Arc<Mutex<ManagedServer>>>>,
}
struct ManagedServer {
route: LanguageRoute,
client: Option<LspClient>,
respawns: u8,
disabled_for_session: bool,
last_used: Instant,
notice_emitted: bool,
}
pub(crate) struct EditDiagnosticsRequest<'a> {
pub(crate) path: &'a Path,
pub(crate) content: String,
}
impl LspManager {
pub(crate) fn new(settings: LspSettings, workspace_root: PathBuf) -> Self {
let workspace_root = workspace_root.canonicalize().unwrap_or(workspace_root);
Self {
settings,
workspace_root,
servers: Mutex::new(BTreeMap::new()),
}
}
pub(crate) fn is_enabled(&self) -> bool {
self.settings.enabled
}
pub(crate) fn inject_diagnostics_on_edit(&self) -> bool {
self.settings.inject_diagnostics_on_edit
}
pub(crate) fn edit_budget(&self) -> Duration {
Duration::from_millis(self.settings.diagnostics_wait_ms)
}
pub(crate) fn sync_and_wait_diagnostics(
&self,
request: EditDiagnosticsRequest<'_>,
deadline: Instant,
cancellation: &AgentCancellation,
) -> Option<String> {
if !self.is_enabled() || !self.inject_diagnostics_on_edit() || cancellation.is_canceled() {
return None;
}
if !self.is_inside_workspace(request.path) {
return None;
}
let route = route_for_path(request.path, &self.settings)?;
let server = self.server_entry(route)?;
let mut server = server.lock().ok()?;
let client = self.client_for(&mut server, deadline, cancellation).ok()?;
client.drain_pending_notifications();
let (version, started) = match client.sync_full_document(
request.path,
request.content,
deadline,
cancellation,
) {
Ok(result) => result,
Err(_) => {
Self::drop_failed_client(&mut server);
return None;
}
};
let diagnostics = match client.wait_for_fresh_diagnostics(
request.path,
version,
started,
deadline,
cancellation,
) {
Some(diagnostics) => diagnostics,
None => {
if client.is_closed() {
Self::drop_failed_client(&mut server);
}
return None;
}
};
server.last_used = Instant::now();
let block =
format_injected_diagnostics(&server.route.server_id, request.path, &diagnostics);
(!block.trim().is_empty()).then_some(block)
}
pub(crate) fn diagnostics_tool(
&self,
path: Option<&Path>,
limit: usize,
_cancellation: &AgentCancellation,
) -> anyhow::Result<String> {
if !self.is_enabled() {
return Ok("LSP diagnostics disabled; set lsp.enabled=true in settings".to_string());
}
let limit = limit.min(MAX_TOOL_DIAGNOSTICS);
let Some(path) = path else {
let mut summaries = Vec::new();
if let Ok(servers) = self.servers.lock() {
for server in servers.values() {
if let Ok(mut server) = server.lock()
&& let Some(client) = server.client.as_mut()
{
summaries.extend(client.workspace_diagnostics_summary(limit).files);
}
}
}
summaries.truncate(limit);
let summary = DiagnosticsSummary {
total_files: summaries.len(),
files: summaries,
};
return Ok(format_diagnostics_tool(None, summary));
};
if !self.is_inside_workspace(path) || route_for_path(path, &self.settings).is_none() {
return Ok(redact_sensitive_text(&format!(
"no cached diagnostics for {}; run after edit or use project checks",
path.display()
)));
}
let mut diagnostics = None;
if let Ok(servers) = self.servers.lock()
&& let Some(server) = route_for_path(path, &self.settings)
.and_then(|route| servers.get(&route.server_id).cloned())
&& let Ok(mut server) = server.lock()
&& let Some(client) = server.client.as_mut()
&& client.has_cached_diagnostics_for_path(path)
{
diagnostics = Some(client.diagnostics_for_path(path));
}
let Some(diagnostics) = diagnostics else {
return Ok(format_diagnostics_tool(
Some(path),
DiagnosticsSummary {
files: Vec::new(),
total_files: 0,
},
));
};
Ok(format_injected_diagnostics("cached", path, &diagnostics))
}
pub(crate) fn references_tool(
&self,
path: &Path,
line: u32,
column: u32,
include_declaration: bool,
limit: usize,
cancellation: &AgentCancellation,
) -> anyhow::Result<String> {
if !self.is_enabled() {
return Ok("LSP references disabled; set lsp.enabled=true in settings".to_string());
}
if !self.is_inside_workspace(path) {
anyhow::bail!("path is outside LSP workspace root");
}
let text = read_text_for_sync(path)?;
let position = one_based_to_lsp(line, column, &text)?;
let route = route_for_path(path, &self.settings)
.ok_or_else(|| anyhow::anyhow!("unsupported file type for LSP references"))?;
let deadline = Instant::now() + self.query_timeout();
let server = self
.server_entry(route)
.ok_or_else(|| anyhow::anyhow!("LSP manager lock poisoned"))?;
let mut server = server
.lock()
.map_err(|_| anyhow::anyhow!("LSP server lock poisoned"))?;
let client = self.client_for(&mut server, deadline, cancellation)?;
if let Err(error) = client.sync_full_document(path, text, deadline, cancellation) {
Self::drop_failed_client(&mut server);
return Err(error);
}
let references = match client.references(
path,
position,
include_declaration,
limit.min(MAX_REFERENCES),
self.query_timeout(),
cancellation,
) {
Ok(references) => references,
Err(error) => {
Self::drop_failed_client(&mut server);
return Err(error);
}
};
server.last_used = Instant::now();
Ok(format_references(
&references,
limit.min(MAX_REFERENCES),
&self.workspace_root,
))
}
pub(crate) fn shutdown_idle(&self) {
let idle_after =
Duration::from_secs(self.settings.idle_shutdown_minutes.saturating_mul(60));
let now = Instant::now();
if let Ok(servers) = self.servers.lock() {
for server in servers.values() {
if let Ok(mut server) = server.lock()
&& server
.client
.as_ref()
.is_some_and(|_| now.duration_since(server.last_used) >= idle_after)
&& let Some(client) = server.client.take()
{
client.shutdown();
}
}
}
}
pub(crate) fn shutdown_all(&self) {
if let Ok(servers) = self.servers.lock() {
for server in servers.values() {
if let Ok(mut server) = server.lock()
&& let Some(client) = server.client.take()
{
client.shutdown();
}
}
}
}
fn server_entry(&self, route: LanguageRoute) -> Option<Arc<Mutex<ManagedServer>>> {
let mut servers = self.servers.lock().ok()?;
Some(
servers
.entry(route.server_id.clone())
.or_insert_with(|| {
Arc::new(Mutex::new(ManagedServer {
route,
client: None,
respawns: 0,
disabled_for_session: false,
last_used: Instant::now(),
notice_emitted: false,
}))
})
.clone(),
)
}
fn client_for<'a>(
&self,
server: &'a mut ManagedServer,
deadline: Instant,
_cancellation: &AgentCancellation,
) -> anyhow::Result<&'a mut LspClient> {
if server.disabled_for_session {
anyhow::bail!("LSP server disabled for session");
}
if server.client.is_none() {
if server.respawns >= MAX_RESPAWNS_PER_SESSION {
server.disabled_for_session = true;
server.notice_emitted = true;
anyhow::bail!("LSP server respawn limit reached");
}
let timeout = deadline.saturating_duration_since(Instant::now());
if timeout.is_zero() {
anyhow::bail!("LSP deadline elapsed before spawn");
}
match LspClient::start_with_timeout(
server.route.server_id.clone(),
&server.route.config,
&self.workspace_root,
server.route.language_id.clone(),
timeout,
) {
Ok(client) => {
server.client = Some(client);
server.respawns += 1;
}
Err(error) => {
server.respawns += 1;
if server.respawns >= MAX_RESPAWNS_PER_SESSION {
server.disabled_for_session = true;
server.notice_emitted = true;
}
return Err(error);
}
}
}
server
.client
.as_mut()
.ok_or_else(|| anyhow::anyhow!("LSP client unavailable"))
}
fn drop_failed_client(server: &mut ManagedServer) {
if let Some(client) = server.client.take() {
client.shutdown();
}
if server.respawns >= MAX_RESPAWNS_PER_SESSION {
server.disabled_for_session = true;
server.notice_emitted = true;
}
}
fn is_inside_workspace(&self, path: &Path) -> bool {
path.canonicalize()
.map(|path| path.starts_with(&self.workspace_root))
.unwrap_or(false)
}
fn query_timeout(&self) -> Duration {
Duration::from_millis(self.settings.diagnostics_wait_ms.min(30_000))
}
}
impl Drop for LspManager {
fn drop(&mut self) {
self.shutdown_all();
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::LspServerConfig;
use std::fs;
fn fake_server_script(mode: &str) -> String {
format!(
r#"
import json, sys, time
mode = {mode:?}
root_uri = None
def read_msg():
header = b''
while not header.endswith(b'\r\n\r\n'):
chunk = sys.stdin.buffer.readline()
if not chunk:
return None
header += chunk
length = 0
for line in header.decode().splitlines():
if line.lower().startswith('content-length:'):
length = int(line.split(':', 1)[1].strip())
return json.loads(sys.stdin.buffer.read(length).decode())
def send(value):
body = json.dumps(value, separators=(',', ':')).encode()
sys.stdout.buffer.write(b'Content-Length: ' + str(len(body)).encode() + b'\r\n\r\n' + body)
sys.stdout.buffer.flush()
while True:
msg = read_msg()
if msg is None:
break
method = msg.get('method')
if method == 'initialize':
if mode == 'crash':
sys.exit(7)
if mode == 'crash_after_initialize':
root_uri = msg.get('params', {{}}).get('rootUri')
send({{'jsonrpc':'2.0','id':msg['id'],'result':{{'capabilities':{{'textDocumentSync':1,'referencesProvider':True}}}}}})
continue
root_uri = msg.get('params', {{}}).get('rootUri')
if mode == 'stale_unversioned' and root_uri:
send({{'jsonrpc':'2.0','method':'textDocument/publishDiagnostics','params':{{'uri':root_uri.rstrip('/') + '/lib.rs','diagnostics':[{{'range':{{'start':{{'line':0,'character':0}},'end':{{'line':0,'character':1}}}},'severity':1,'source':'fake','message':'old unversioned'}}]}}}})
send({{'jsonrpc':'2.0','id':msg['id'],'result':{{'capabilities':{{'textDocumentSync':1,'referencesProvider':True}}}}}})
elif method == 'initialized':
if mode == 'crash_after_initialize':
sys.exit(8)
pass
elif method in ('textDocument/didOpen','textDocument/didChange'):
if mode == 'stale_unversioned':
continue
td = msg['params'].get('textDocument', {{}})
uri = td.get('uri')
version = td.get('version')
send({{'jsonrpc':'2.0','method':'textDocument/publishDiagnostics','params':{{'uri':uri,'version':version,'diagnostics':[{{'range':{{'start':{{'line':0,'character':0}},'end':{{'line':0,'character':1}}}},'severity':1,'source':'fake','message':'boom'}}]}}}})
elif method == 'textDocument/references':
uri = msg['params']['textDocument']['uri']
send({{'jsonrpc':'2.0','id':msg['id'],'result':[{{'uri':uri,'range':{{'start':{{'line':1,'character':0}},'end':{{'line':1,'character':1}}}}}}]}})
elif method == 'shutdown':
send({{'jsonrpc':'2.0','id':msg['id'],'result':None}})
elif method == 'exit':
break
"#
)
}
fn settings_with_fake(mode: &str) -> LspSettings {
let mut settings = LspSettings {
enabled: true,
diagnostics_wait_ms: 1_000,
..LspSettings::default()
};
settings.servers.insert(
"rust-analyzer".to_string(),
LspServerConfig {
command: "python3".to_string(),
args: vec!["-u".to_string(), "-c".to_string(), fake_server_script(mode)],
enabled: true,
},
);
settings
}
#[test]
fn disabled_and_unsupported_paths_do_not_spawn() {
let temp = tempfile::TempDir::new().unwrap();
let path = temp.path().join("lib.rs");
fs::write(&path, "fn main() {}\n").unwrap();
let disabled = LspManager::new(LspSettings::default(), temp.path().to_path_buf());
assert!(
disabled
.sync_and_wait_diagnostics(
EditDiagnosticsRequest {
path: &path,
content: "fn main() {}\n".to_string(),
},
Instant::now() + Duration::from_millis(100),
&AgentCancellation::default(),
)
.is_none()
);
assert!(
disabled
.diagnostics_tool(Some(&path), 10, &AgentCancellation::default())
.unwrap()
.contains("disabled")
);
let enabled = LspManager::new(settings_with_fake("normal"), temp.path().to_path_buf());
let txt = temp.path().join("notes.txt");
fs::write(&txt, "notes\n").unwrap();
assert!(
enabled
.sync_and_wait_diagnostics(
EditDiagnosticsRequest {
path: &txt,
content: "notes\n".to_string(),
},
Instant::now() + Duration::from_millis(100),
&AgentCancellation::default(),
)
.is_none()
);
assert!(enabled.servers.lock().unwrap().is_empty());
}
#[test]
fn sync_injects_diagnostics_and_reuses_server() {
let temp = tempfile::TempDir::new().unwrap();
let path = temp.path().join("lib.rs");
fs::write(&path, "fn main() {}\n").unwrap();
let manager = LspManager::new(settings_with_fake("normal"), temp.path().to_path_buf());
let first = manager
.sync_and_wait_diagnostics(
EditDiagnosticsRequest {
path: &path,
content: "fn main() {}\n".to_string(),
},
Instant::now() + Duration::from_secs(2),
&AgentCancellation::default(),
)
.unwrap();
assert!(first.contains("DIAGNOSTICS"));
assert!(first.contains("boom"));
let second = manager
.sync_and_wait_diagnostics(
EditDiagnosticsRequest {
path: &path,
content: "fn main() { }\n".to_string(),
},
Instant::now() + Duration::from_secs(2),
&AgentCancellation::default(),
)
.unwrap();
assert!(second.contains("boom"));
let server = manager.servers.lock().unwrap()["rust-analyzer"].clone();
assert_eq!(server.lock().unwrap().respawns, 1);
}
#[test]
fn stale_unversioned_diagnostics_before_sync_are_not_injected() {
let temp = tempfile::TempDir::new().unwrap();
let path = temp.path().join("lib.rs");
fs::write(&path, "fn main() {}\n").unwrap();
let manager = LspManager::new(
settings_with_fake("stale_unversioned"),
temp.path().to_path_buf(),
);
let block = manager.sync_and_wait_diagnostics(
EditDiagnosticsRequest {
path: &path,
content: "fn main() {}\n".to_string(),
},
Instant::now() + Duration::from_millis(200),
&AgentCancellation::default(),
);
assert!(
block.is_none(),
"stale unversioned diagnostics were injected: {block:?}"
);
}
#[test]
fn diagnostics_before_edit_reports_no_cache_and_references_syncs_first() {
let temp = tempfile::TempDir::new().unwrap();
let path = temp.path().join("lib.rs");
fs::write(&path, "fn main() {}\nlet x = 1;\n").unwrap();
let manager = LspManager::new(settings_with_fake("normal"), temp.path().to_path_buf());
let no_cache = manager
.diagnostics_tool(Some(&path), 10, &AgentCancellation::default())
.unwrap();
assert!(no_cache.contains("no cached diagnostics"));
let refs = manager
.references_tool(&path, 1, 1, false, 10, &AgentCancellation::default())
.unwrap();
assert_eq!(refs, "lib.rs:2:1");
let server = manager.servers.lock().unwrap()["rust-analyzer"].clone();
assert!(server.lock().unwrap().client.is_some());
}
#[test]
fn outside_workspace_does_not_sync_and_spawn_failures_hit_respawn_cap() {
let temp = tempfile::TempDir::new().unwrap();
let outside = tempfile::NamedTempFile::new().unwrap();
let manager = LspManager::new(settings_with_fake("normal"), temp.path().to_path_buf());
assert!(
manager
.sync_and_wait_diagnostics(
EditDiagnosticsRequest {
path: outside.path(),
content: "fn main() {}\n".to_string(),
},
Instant::now() + Duration::from_millis(100),
&AgentCancellation::default(),
)
.is_none()
);
assert!(manager.servers.lock().unwrap().is_empty());
let path = temp.path().join("lib.rs");
fs::write(&path, "fn main() {}\n").unwrap();
let crashing = LspManager::new(settings_with_fake("crash"), temp.path().to_path_buf());
for _ in 0..3 {
assert!(
crashing
.sync_and_wait_diagnostics(
EditDiagnosticsRequest {
path: &path,
content: "fn main() {}\n".to_string(),
},
Instant::now() + Duration::from_secs(1),
&AgentCancellation::default(),
)
.is_none()
);
}
let server = crashing.servers.lock().unwrap()["rust-analyzer"].clone();
let server = server.lock().unwrap();
assert!(server.disabled_for_session);
assert_eq!(server.respawns, MAX_RESPAWNS_PER_SESSION);
}
#[test]
fn post_initialize_crash_drops_stale_client_and_disables_after_respawn_cap() {
let temp = tempfile::TempDir::new().unwrap();
let path = temp.path().join("lib.rs");
fs::write(&path, "fn main() {}\n").unwrap();
let manager = LspManager::new(
settings_with_fake("crash_after_initialize"),
temp.path().to_path_buf(),
);
for _ in 0..3 {
assert!(
manager
.sync_and_wait_diagnostics(
EditDiagnosticsRequest {
path: &path,
content: "fn main() {}\n".to_string(),
},
Instant::now() + Duration::from_secs(1),
&AgentCancellation::default(),
)
.is_none()
);
}
let server = manager.servers.lock().unwrap()["rust-analyzer"].clone();
let server = server.lock().unwrap();
assert!(server.disabled_for_session);
assert!(server.client.is_none());
assert_eq!(server.respawns, MAX_RESPAWNS_PER_SESSION);
}
}