use crate::{
cancellation::AgentCancellation,
config::LspSettings,
lsp::{
client::LspClient,
diagnostics::format_injected_diagnostics,
sync::{LanguageRoute, route_for_path},
},
};
use std::{
collections::BTreeMap,
path::{Path, PathBuf},
sync::{Arc, Mutex, MutexGuard, TryLockError, mpsc},
thread::{self, JoinHandle},
time::{Duration, Instant},
};
const MAX_UNEXPECTED_FAILURES_PER_SESSION: u8 = 3;
const LOCK_POLL_INTERVAL: Duration = Duration::from_millis(5);
fn lock_until<'a, T>(
mutex: &'a Mutex<T>,
deadline: Instant,
cancellation: &AgentCancellation,
) -> Option<MutexGuard<'a, T>> {
loop {
if cancellation.is_canceled() || Instant::now() >= deadline {
return None;
}
match mutex.try_lock() {
Ok(guard) => return Some(guard),
Err(TryLockError::Poisoned(_)) => return None,
Err(TryLockError::WouldBlock) => {
thread::sleep(
deadline
.saturating_duration_since(Instant::now())
.min(LOCK_POLL_INTERVAL),
);
}
}
}
}
pub(crate) struct LspManager {
settings: LspSettings,
workspace_root: PathBuf,
servers: Arc<Mutex<BTreeMap<String, Arc<Mutex<ManagedServer>>>>>,
watchdog_stop: Mutex<Option<mpsc::Sender<()>>>,
watchdog: Mutex<Option<JoinHandle<()>>>,
}
struct ManagedServer {
route: LanguageRoute,
client: Option<LspClient>,
unexpected_failures: u8,
disabled_for_session: bool,
last_used: Instant,
}
pub(crate) struct EditDiagnosticsRequest<'a> {
pub(crate) path: &'a Path,
pub(crate) content: String,
}
struct WatchdogAccess {
servers: Arc<Mutex<BTreeMap<String, Arc<Mutex<ManagedServer>>>>>,
idle_after: Duration,
}
impl WatchdogAccess {
fn shutdown_idle(&self) {
let now = Instant::now();
let entries = self
.servers
.try_lock()
.ok()
.map(|servers| servers.values().cloned().collect::<Vec<_>>())
.unwrap_or_default();
let clients = entries
.into_iter()
.filter_map(|entry| {
let mut server = entry.try_lock().ok()?;
if server.client.is_some()
&& now.duration_since(server.last_used) >= self.idle_after
{
server.client.take()
} else {
None
}
})
.collect::<Vec<_>>();
for client in clients {
client.shutdown();
}
}
}
impl LspManager {
pub(crate) fn new(settings: LspSettings, workspace_root: PathBuf) -> Self {
let idle_after = Duration::from_secs(settings.idle_shutdown_minutes.saturating_mul(60));
let tick = idle_after
.min(Duration::from_secs(1))
.max(Duration::from_millis(10));
Self::new_with_timing_impl(settings, workspace_root, idle_after, tick)
}
fn new_with_timing_impl(
settings: LspSettings,
workspace_root: PathBuf,
idle_after: Duration,
tick: Duration,
) -> Self {
let workspace_root = workspace_root.canonicalize().unwrap_or(workspace_root);
let (stop_tx, stop_rx) = mpsc::channel();
let servers = Arc::new(Mutex::new(BTreeMap::new()));
let access = Arc::new(WatchdogAccess {
servers: Arc::clone(&servers),
idle_after,
});
let handle = thread::spawn(move || {
loop {
match stop_rx.recv_timeout(tick) {
Ok(()) | Err(mpsc::RecvTimeoutError::Disconnected) => break,
Err(mpsc::RecvTimeoutError::Timeout) => access.shutdown_idle(),
}
}
});
Self {
settings,
workspace_root,
servers,
watchdog_stop: Mutex::new(Some(stop_tx)),
watchdog: Mutex::new(Some(handle)),
}
}
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, deadline, cancellation)?;
let mut server = lock_until(&server, deadline, cancellation)?;
server.last_used = Instant::now();
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(_) => {
let failed = Self::take_failed_client(&mut server);
drop(server);
if let Some(client) = failed {
client.shutdown();
}
return None;
}
};
let diagnostics = client.wait_for_fresh_diagnostics(
request.path,
version,
started,
deadline,
cancellation,
);
if diagnostics.is_none() && client.is_closed() {
let failed = Self::take_failed_client(&mut server);
drop(server);
if let Some(client) = failed {
client.shutdown();
}
return None;
}
let diagnostics = diagnostics?;
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 shutdown_all(&self) {
let clients = self
.servers
.lock()
.ok()
.map(|servers| {
servers
.values()
.filter_map(|entry| entry.lock().ok()?.client.take())
.collect::<Vec<_>>()
})
.unwrap_or_default();
for client in clients {
client.shutdown();
}
}
fn server_entry(
&self,
route: LanguageRoute,
deadline: Instant,
cancellation: &AgentCancellation,
) -> Option<Arc<Mutex<ManagedServer>>> {
let mut servers = lock_until(&self.servers, deadline, cancellation)?;
Some(
servers
.entry(route.server_id.clone())
.or_insert_with(|| {
Arc::new(Mutex::new(ManagedServer {
route,
client: None,
unexpected_failures: 0,
disabled_for_session: false,
last_used: Instant::now(),
}))
})
.clone(),
)
}
fn client_for<'a>(
&self,
server: &'a mut ManagedServer,
deadline: Instant,
cancellation: &AgentCancellation,
) -> anyhow::Result<&'a mut LspClient> {
cancellation.check()?;
if server.disabled_for_session {
anyhow::bail!("LSP server disabled for session");
}
if server.client.is_none() {
if server.unexpected_failures >= MAX_UNEXPECTED_FAILURES_PER_SESSION {
server.disabled_for_session = true;
anyhow::bail!("LSP server unexpected failure 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);
}
Err(error) => {
server.unexpected_failures = server.unexpected_failures.saturating_add(1);
if server.unexpected_failures >= MAX_UNEXPECTED_FAILURES_PER_SESSION {
server.disabled_for_session = true;
}
return Err(error);
}
}
}
server
.client
.as_mut()
.ok_or_else(|| anyhow::anyhow!("LSP client unavailable"))
}
fn take_failed_client(server: &mut ManagedServer) -> Option<LspClient> {
let client = server.client.take();
if client.is_some() {
server.unexpected_failures = server.unexpected_failures.saturating_add(1);
}
if server.unexpected_failures >= MAX_UNEXPECTED_FAILURES_PER_SESSION {
server.disabled_for_session = true;
}
client
}
fn is_inside_workspace(&self, path: &Path) -> bool {
path.canonicalize()
.map(|path| path.starts_with(&self.workspace_root))
.unwrap_or(false)
}
}
impl Drop for LspManager {
fn drop(&mut self) {
if let Ok(mut stop) = self.watchdog_stop.lock() {
stop.take();
}
if let Ok(mut watchdog) = self.watchdog.lock()
&& let Some(handle) = watchdog.take()
{
let _ = handle.join();
}
self.shutdown_all();
}
}