use std::time::Duration;
use lsp_types::{
HoverParams as LspHoverParams, PartialResultParams, ReferenceContext, ReferenceParams,
TextDocumentIdentifier, TextDocumentPositionParams, WorkDoneProgressParams,
};
use tokio::time::Instant;
use super::Translator;
use super::dto::{
DefinitionResult, HoverResult, Location, LocationsResult, Position, ReferencesResult,
};
use super::encoding_ctx::EncodingCtx;
use super::routing::{Capability, IndexingGate};
use crate::bridge::IndexingState;
use crate::bridge::indexing::{
DEFAULT_INDEXING_READY_TIMEOUT_SECS, INDEXING_STALENESS_BOUND, PROGRESS_LATCH_IDLE,
PROGRESS_SETTLE,
};
use crate::config::{ServerId, ToolKind};
use crate::error::{Error, Result};
pub(super) const INDEXING_READY_TIMEOUT: Duration =
Duration::from_secs(DEFAULT_INDEXING_READY_TIMEOUT_SECS);
const INDEXING_POLL_INTERVAL: Duration = Duration::from_millis(100);
const _: () = assert!(
INDEXING_STALENESS_BOUND.as_nanos() > INDEXING_READY_TIMEOUT.as_nanos(),
"INDEXING_STALENESS_BOUND must be greater than INDEXING_READY_TIMEOUT"
);
const _: () = assert!(
PROGRESS_SETTLE.as_nanos() < PROGRESS_LATCH_IDLE.as_nanos(),
"PROGRESS_SETTLE must be less than PROGRESS_LATCH_IDLE"
);
const _: () = assert!(
PROGRESS_SETTLE.as_nanos() < INDEXING_READY_TIMEOUT.as_nanos(),
"PROGRESS_SETTLE must be less than INDEXING_READY_TIMEOUT"
);
fn definition_to_locations(definition: lsp_types::Definition) -> Vec<lsp_types::Location> {
match definition {
lsp_types::Definition::Location(loc) => vec![loc],
lsp_types::Definition::LocationList(locs) => locs,
}
}
fn definition_link_to_location(link: lsp_types::DefinitionLink) -> lsp_types::Location {
lsp_types::Location {
uri: link.target_uri,
range: link.target_selection_range,
}
}
pub(super) const MAX_NORMALIZED_LOCATIONS: usize = 10_000;
async fn lsp_locations_to_mcp(
mut locs: Vec<lsp_types::Location>,
ctx: &EncodingCtx,
) -> NormalizedLocations {
let truncated = locs.len() > MAX_NORMALIZED_LOCATIONS;
if truncated {
tracing::warn!(
reported = locs.len(),
cap = MAX_NORMALIZED_LOCATIONS,
"LSP response location count exceeds MAX_NORMALIZED_LOCATIONS; truncating"
);
}
locs.truncate(MAX_NORMALIZED_LOCATIONS);
let mut locations = Vec::with_capacity(locs.len());
for loc in locs {
locations.push(Location {
uri: loc.uri.to_string(),
range: ctx.normalize_range(&loc.uri, loc.range).await,
out_of_workspace: ctx.is_out_of_workspace(&loc.uri),
});
}
NormalizedLocations {
locations,
truncated,
positions_degraded: ctx.positions_degraded(),
}
}
struct NormalizedLocations {
locations: Vec<Location>,
truncated: bool,
positions_degraded: bool,
}
enum GotoKind {
Definition(lsp_types::Definition),
DefinitionLinkList(Vec<lsp_types::DefinitionLink>),
}
trait GotoResponse {
fn into_kind(self) -> GotoKind;
}
impl GotoResponse for lsp_types::DefinitionResponse {
fn into_kind(self) -> GotoKind {
match self {
Self::Definition(def) => GotoKind::Definition(def),
Self::DefinitionLinkList(links) => GotoKind::DefinitionLinkList(links),
}
}
}
impl GotoResponse for lsp_types::ImplementationResponse {
fn into_kind(self) -> GotoKind {
match self {
Self::Definition(def) => GotoKind::Definition(def),
Self::DefinitionLinkList(links) => GotoKind::DefinitionLinkList(links),
}
}
}
impl GotoResponse for lsp_types::TypeDefinitionResponse {
fn into_kind(self) -> GotoKind {
match self {
Self::Definition(def) => GotoKind::Definition(def),
Self::DefinitionLinkList(links) => GotoKind::DefinitionLinkList(links),
}
}
}
async fn goto_response_to_locations<R: GotoResponse>(
response: Option<R>,
ctx: &EncodingCtx,
) -> NormalizedLocations {
let lsp_locs = match response.map(GotoResponse::into_kind) {
Some(GotoKind::Definition(def)) => definition_to_locations(def),
Some(GotoKind::DefinitionLinkList(links)) => {
links.into_iter().map(definition_link_to_location).collect()
}
None => vec![],
};
lsp_locations_to_mcp(lsp_locs, ctx).await
}
trait GotoParams: Sized {
fn from_position(text_document_position_params: TextDocumentPositionParams) -> Self;
}
impl GotoParams for lsp_types::DefinitionParams {
fn from_position(text_document_position_params: TextDocumentPositionParams) -> Self {
Self {
text_document_position_params,
work_done_progress_params: WorkDoneProgressParams::default(),
partial_result_params: PartialResultParams::default(),
}
}
}
impl GotoParams for lsp_types::ImplementationParams {
fn from_position(text_document_position_params: TextDocumentPositionParams) -> Self {
Self {
text_document_position_params,
work_done_progress_params: WorkDoneProgressParams::default(),
partial_result_params: PartialResultParams::default(),
}
}
}
impl GotoParams for lsp_types::TypeDefinitionParams {
fn from_position(text_document_position_params: TextDocumentPositionParams) -> Self {
Self {
text_document_position_params,
work_done_progress_params: WorkDoneProgressParams::default(),
partial_result_params: PartialResultParams::default(),
}
}
}
#[allow(deprecated)]
fn extract_hover_contents(contents: lsp_types::Contents) -> String {
match contents {
lsp_types::Contents::MarkedString(marked_string) => marked_string_to_string(marked_string),
lsp_types::Contents::MarkedStringList(marked_strings) => marked_strings
.into_iter()
.map(marked_string_to_string)
.collect::<Vec<_>>()
.join("\n\n"),
lsp_types::Contents::MarkupContent(markup) => markup.value,
}
}
#[allow(deprecated)]
fn marked_string_to_string(marked: lsp_types::MarkedString) -> String {
match marked {
lsp_types::MarkedString::String(s) => s,
lsp_types::MarkedString::MarkedStringWithLanguage(ls) => {
format!("```{}\n{}\n```", ls.language, ls.value)
}
}
}
impl Translator {
pub(super) async fn wait_for_indexing_ready(&self, server_id: &ServerId) -> Result<()> {
self.wait_for_indexing_ready_with(
server_id,
self.indexing_ready_timeout,
INDEXING_POLL_INTERVAL,
)
.await
}
async fn wait_for_indexing_ready_with(
&self,
server_id: &ServerId,
timeout: Duration,
poll_interval: Duration,
) -> Result<()> {
let Some(cache) = self.notification_cache.as_ref() else {
return Ok(());
};
let start = Instant::now();
let deadline = start + timeout;
loop {
let state = cache.lock().await.indexing_state(server_id);
if state != IndexingState::Loading {
return Ok(());
}
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return Err(Error::WorkspaceIndexing {
server_id: server_id.clone(),
elapsed_secs: start.elapsed().as_secs(),
});
}
tokio::time::sleep(poll_interval.min(remaining)).await;
}
}
pub async fn handle_hover(&self, file_path: String, position: Position) -> Result<HoverResult> {
let Position { line, character } = position;
let (server_id, client, uri) = self
.prepare_gated_document(
&file_path,
ToolKind::Hover,
Capability::Hover,
IndexingGate::Required,
)
.await?;
let ctx = self.encoding_ctx(&server_id);
let lsp_position = ctx.to_lsp(&uri, line, character).await;
let response_uri = uri.clone();
let params = LspHoverParams {
text_document_position_params: TextDocumentPositionParams {
text_document: TextDocumentIdentifier { uri },
position: lsp_position,
},
work_done_progress_params: WorkDoneProgressParams::default(),
};
let response = client
.request_typed::<lsp_types::HoverRequest>(params, client.request_timeout())
.await?;
let result = match response {
Some(hover) => {
let contents = extract_hover_contents(hover.contents);
let range = match hover.range {
Some(r) => Some(ctx.normalize_range(&response_uri, r).await),
None => None,
};
HoverResult {
contents,
range,
positions_degraded: ctx.positions_degraded(),
}
}
None => HoverResult {
contents: "No hover information available".to_string(),
range: None,
positions_degraded: ctx.positions_degraded(),
},
};
Ok(result)
}
async fn handle_goto<R, T>(
&self,
file_path: &str,
position: Position,
tool: ToolKind,
capability: Capability,
) -> Result<NormalizedLocations>
where
R: lsp_types::Request<Result = Option<T>>,
R::Params: GotoParams,
T: GotoResponse,
{
let Position { line, character } = position;
let (server_id, client, uri) = self
.prepare_gated_document(file_path, tool, capability, IndexingGate::Required)
.await?;
let ctx = self.encoding_ctx(&server_id);
let lsp_position = ctx.to_lsp(&uri, line, character).await;
let params = R::Params::from_position(TextDocumentPositionParams {
text_document: TextDocumentIdentifier { uri },
position: lsp_position,
});
let response = client
.request_typed::<R>(params, client.request_timeout())
.await?;
Ok(goto_response_to_locations(response, &ctx).await)
}
pub async fn handle_definition(
&self,
file_path: String,
position: Position,
) -> Result<DefinitionResult> {
let NormalizedLocations {
locations,
truncated,
positions_degraded,
} = self
.handle_goto::<lsp_types::DefinitionRequest, _>(
&file_path,
position,
ToolKind::Definition,
Capability::Definition,
)
.await?;
Ok(DefinitionResult {
locations,
truncated,
positions_degraded,
})
}
pub async fn handle_references(
&self,
file_path: String,
position: Position,
include_declaration: bool,
) -> Result<ReferencesResult> {
let Position { line, character } = position;
let (server_id, client, uri) = self
.prepare_gated_document(
&file_path,
ToolKind::References,
Capability::References,
IndexingGate::Required,
)
.await?;
let ctx = self.encoding_ctx(&server_id);
let lsp_position = ctx.to_lsp(&uri, line, character).await;
let params = ReferenceParams {
text_document_position_params: TextDocumentPositionParams {
text_document: TextDocumentIdentifier { uri },
position: lsp_position,
},
work_done_progress_params: WorkDoneProgressParams::default(),
partial_result_params: PartialResultParams::default(),
context: ReferenceContext {
include_declaration,
},
};
let response = client
.request_typed::<lsp_types::ReferencesRequest>(params, client.request_timeout())
.await?;
let locations = response.unwrap_or_default();
let NormalizedLocations {
locations,
truncated,
positions_degraded,
} = lsp_locations_to_mcp(locations, &ctx).await;
Ok(ReferencesResult {
locations,
truncated,
positions_degraded,
})
}
pub async fn handle_implementation(
&self,
file_path: String,
position: Position,
) -> Result<LocationsResult> {
let NormalizedLocations {
locations,
truncated,
positions_degraded,
} = self
.handle_goto::<lsp_types::ImplementationRequest, _>(
&file_path,
position,
ToolKind::Implementation,
Capability::Implementation,
)
.await?;
Ok(LocationsResult {
locations,
truncated,
positions_degraded,
})
}
pub async fn handle_type_definition(
&self,
file_path: String,
position: Position,
) -> Result<LocationsResult> {
let NormalizedLocations {
locations,
truncated,
positions_degraded,
} = self
.handle_goto::<lsp_types::TypeDefinitionRequest, _>(
&file_path,
position,
ToolKind::TypeDefinition,
Capability::TypeDefinition,
)
.await?;
Ok(LocationsResult {
locations,
truncated,
positions_degraded,
})
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, deprecated)]
mod tests {
use std::fs;
use std::sync::Arc;
use std::time::Duration;
use tempfile::TempDir;
use tokio::io::BufReader;
use tokio::sync::Mutex;
use tokio::time::timeout;
use url::Url;
use super::*;
use crate::bridge::encoding::PositionEncoding;
use crate::bridge::translator::testing::*;
use crate::bridge::{NotificationCache, lock_std, path_to_uri};
use crate::config::ServerId;
#[tokio::test]
async fn test_wait_for_indexing_ready_without_cache_is_noop() {
let translator = Translator::new();
let server_id = ServerId::from("rust");
translator
.wait_for_indexing_ready(&server_id)
.await
.unwrap();
}
#[tokio::test]
async fn test_wait_for_indexing_ready_unknown_state_is_noop() {
let translator = Translator::new()
.with_notification_cache(Arc::new(Mutex::new(NotificationCache::new())));
let server_id = ServerId::from("rust");
translator
.wait_for_indexing_ready(&server_id)
.await
.unwrap();
}
#[tokio::test]
async fn test_wait_for_indexing_ready_ready_state_is_noop() {
let cache = Arc::new(Mutex::new(NotificationCache::new()));
let server_id = ServerId::from("rust");
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": true})),
);
let translator = Translator::new().with_notification_cache(cache);
translator
.wait_for_indexing_ready(&server_id)
.await
.unwrap();
}
#[tokio::test(start_paused = true)]
async fn test_wait_for_indexing_ready_uses_configured_timeout_override() {
let cache = Arc::new(Mutex::new(NotificationCache::new()));
let server_id = ServerId::from("rust");
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
let translator = Translator::new()
.with_notification_cache(cache)
.with_indexing_ready_timeout(Duration::from_secs(5));
let start = Instant::now();
let err = translator
.wait_for_indexing_ready(&server_id)
.await
.unwrap_err();
assert!(matches!(err, Error::WorkspaceIndexing { elapsed_secs, .. } if elapsed_secs == 5));
assert_eq!(start.elapsed(), Duration::from_secs(5));
}
#[tokio::test]
async fn test_wait_for_indexing_ready_loading_times_out() {
let cache = Arc::new(Mutex::new(NotificationCache::new()));
let server_id = ServerId::from("rust");
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
let translator = Translator::new().with_notification_cache(cache);
let err = translator
.wait_for_indexing_ready_with(
&server_id,
Duration::from_millis(50),
Duration::from_millis(10),
)
.await
.unwrap_err();
assert!(matches!(
err,
Error::WorkspaceIndexing { server_id: id, .. } if id == ServerId::from("rust")
));
}
#[tokio::test]
async fn test_wait_for_indexing_ready_returns_ok_once_signaled_ready() {
let cache = Arc::new(Mutex::new(NotificationCache::new()));
let server_id = ServerId::from("rust");
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
let translator = Translator::new().with_notification_cache(Arc::clone(&cache));
let waiter = {
let server_id = server_id.clone();
tokio::spawn(async move {
translator
.wait_for_indexing_ready_with(
&server_id,
Duration::from_secs(5),
Duration::from_millis(10),
)
.await
})
};
tokio::time::sleep(Duration::from_millis(30)).await;
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": true})),
);
timeout(Duration::from_secs(1), waiter)
.await
.expect("waiter task timed out")
.expect("waiter task panicked")
.expect("expected Ok once quiescent");
}
#[tokio::test(start_paused = true)]
async fn test_handle_hover_returns_workspace_indexing_error_when_loading() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
hover_provider: Some(lsp_types::HoverProvider::Bool(true)),
..Default::default()
};
let (translator, _server) = translator_with_capabilities(&dir, &server_id, caps);
let cache = Arc::new(Mutex::new(NotificationCache::new()));
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
let translator = translator.with_notification_cache(cache);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let err = translator
.handle_hover(path.to_string_lossy().to_string(), pos(1, 1))
.await
.unwrap_err();
assert!(matches!(
err,
Error::WorkspaceIndexing { server_id: id, elapsed_secs: 30 } if id == server_id
));
}
#[tokio::test]
async fn test_handle_hover_dispatches_when_indexing_ready() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
hover_provider: Some(lsp_types::HoverProvider::Bool(true)),
..Default::default()
};
let (translator, mut server) = translator_with_capabilities(&dir, &server_id, caps);
let cache = Arc::new(Mutex::new(NotificationCache::new()));
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": true})),
);
let translator = Arc::new(translator.with_notification_cache(cache));
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let handle = {
let translator = Arc::clone(&translator);
let path = path.to_string_lossy().to_string();
tokio::spawn(async move { translator.handle_hover(path, pos(1, 1)).await })
};
let mut wire = BufReader::new(&mut server.write_stdout);
let opened = read_framed_message(&mut wire).await;
assert_eq!(opened["method"], "textDocument/didOpen");
let request = read_framed_message(&mut wire).await;
assert_eq!(request["method"], "textDocument/hover");
write_response(
&mut server.read_half_stdin,
&request["id"],
serde_json::json!({
"contents": {"kind": "markdown", "value": "hover text"}
}),
)
.await;
let result = handle.await.unwrap().unwrap();
assert_eq!(result.contents, "hover text");
}
#[tokio::test(start_paused = true)]
async fn test_handle_definition_returns_workspace_indexing_error_when_loading() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
definition_provider: Some(lsp_types::DefinitionProvider::Bool(true)),
..Default::default()
};
let (translator, _server) = translator_with_capabilities(&dir, &server_id, caps);
let cache = Arc::new(Mutex::new(NotificationCache::new()));
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
let translator = translator.with_notification_cache(cache);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let err = translator
.handle_definition(path.to_string_lossy().to_string(), pos(1, 1))
.await
.unwrap_err();
assert!(matches!(
err,
Error::WorkspaceIndexing { server_id: id, .. } if id == server_id
));
}
#[tokio::test]
async fn test_handle_definition_dispatches_when_indexing_ready() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
definition_provider: Some(lsp_types::DefinitionProvider::Bool(true)),
..Default::default()
};
let (translator, mut server) = translator_with_capabilities(&dir, &server_id, caps);
let cache = Arc::new(Mutex::new(NotificationCache::new()));
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": true})),
);
let translator = Arc::new(translator.with_notification_cache(cache));
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let handle = {
let translator = Arc::clone(&translator);
let path = path.to_string_lossy().to_string();
tokio::spawn(async move { translator.handle_definition(path, pos(1, 1)).await })
};
let mut wire = BufReader::new(&mut server.write_stdout);
let opened = read_framed_message(&mut wire).await;
assert_eq!(opened["method"], "textDocument/didOpen");
let request = read_framed_message(&mut wire).await;
assert_eq!(request["method"], "textDocument/definition");
write_response(
&mut server.read_half_stdin,
&request["id"],
serde_json::Value::Null,
)
.await;
let result = handle.await.unwrap().unwrap();
assert!(result.locations.is_empty());
}
#[tokio::test(start_paused = true)]
async fn test_handle_references_returns_workspace_indexing_error_when_loading() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
references_provider: Some(lsp_types::ReferencesProvider::Bool(true)),
..Default::default()
};
let (translator, _server) = translator_with_capabilities(&dir, &server_id, caps);
let cache = Arc::new(Mutex::new(NotificationCache::new()));
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
let translator = translator.with_notification_cache(cache);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let err = translator
.handle_references(path.to_string_lossy().to_string(), pos(1, 1), true)
.await
.unwrap_err();
assert!(matches!(
err,
Error::WorkspaceIndexing { server_id: id, .. } if id == server_id
));
}
#[tokio::test]
async fn test_handle_references_dispatches_when_indexing_ready() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
references_provider: Some(lsp_types::ReferencesProvider::Bool(true)),
..Default::default()
};
let (translator, mut server) = translator_with_capabilities(&dir, &server_id, caps);
let cache = Arc::new(Mutex::new(NotificationCache::new()));
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": true})),
);
let translator = Arc::new(translator.with_notification_cache(cache));
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let handle = {
let translator = Arc::clone(&translator);
let path = path.to_string_lossy().to_string();
tokio::spawn(async move { translator.handle_references(path, pos(1, 1), true).await })
};
let mut wire = BufReader::new(&mut server.write_stdout);
let opened = read_framed_message(&mut wire).await;
assert_eq!(opened["method"], "textDocument/didOpen");
let request = read_framed_message(&mut wire).await;
assert_eq!(request["method"], "textDocument/references");
write_response(
&mut server.read_half_stdin,
&request["id"],
serde_json::Value::Null,
)
.await;
let result = handle.await.unwrap().unwrap();
assert!(result.locations.is_empty());
}
#[tokio::test(start_paused = true)]
async fn test_handle_implementation_returns_workspace_indexing_error_when_loading() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
implementation_provider: Some(lsp_types::ImplementationProvider::Bool(true)),
..Default::default()
};
let (translator, _server) = translator_with_capabilities(&dir, &server_id, caps);
let cache = Arc::new(Mutex::new(NotificationCache::new()));
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
let translator = translator.with_notification_cache(cache);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let err = translator
.handle_implementation(path.to_string_lossy().to_string(), pos(1, 1))
.await
.unwrap_err();
assert!(matches!(
err,
Error::WorkspaceIndexing { server_id: id, .. } if id == server_id
));
}
#[tokio::test(start_paused = true)]
async fn test_handle_type_definition_returns_workspace_indexing_error_when_loading() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
type_definition_provider: Some(lsp_types::TypeDefinitionProvider::Bool(true)),
..Default::default()
};
let (translator, _server) = translator_with_capabilities(&dir, &server_id, caps);
let cache = Arc::new(Mutex::new(NotificationCache::new()));
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
let translator = translator.with_notification_cache(cache);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let err = translator
.handle_type_definition(path.to_string_lossy().to_string(), pos(1, 1))
.await
.unwrap_err();
assert!(matches!(
err,
Error::WorkspaceIndexing { server_id: id, .. } if id == server_id
));
}
#[tokio::test]
async fn test_wait_for_indexing_ready_timeout_does_not_mutate_shared_state() {
let cache = Arc::new(Mutex::new(NotificationCache::new()));
let server_id = ServerId::from("rust");
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
let translator = Translator::new().with_notification_cache(Arc::clone(&cache));
translator
.wait_for_indexing_ready_with(
&server_id,
Duration::from_millis(30),
Duration::from_millis(10),
)
.await
.unwrap_err();
assert_eq!(
cache.lock().await.indexing_state(&server_id),
IndexingState::Loading,
"a timed-out wait must not touch the shared entry -- it is still fresh, so it must \
still read as Loading for any other caller"
);
}
#[tokio::test]
async fn test_wait_for_indexing_ready_one_callers_timeout_does_not_release_another() {
let cache = Arc::new(Mutex::new(NotificationCache::new()));
let server_id = ServerId::from("rust");
cache.lock().await.observe_indexing_signal(
&server_id,
"experimental/serverStatus",
Some(&serde_json::json!({"quiescent": false})),
);
let translator = Arc::new(Translator::new().with_notification_cache(Arc::clone(&cache)));
let short = {
let translator = Arc::clone(&translator);
let server_id = server_id.clone();
tokio::spawn(async move {
translator
.wait_for_indexing_ready_with(
&server_id,
Duration::from_millis(80),
Duration::from_millis(10),
)
.await
})
};
let long = {
let translator = Arc::clone(&translator);
let server_id = server_id.clone();
tokio::spawn(async move {
translator
.wait_for_indexing_ready_with(
&server_id,
Duration::from_secs(30),
Duration::from_millis(10),
)
.await
})
};
let short_result = short.await.unwrap();
assert!(
matches!(short_result, Err(Error::WorkspaceIndexing { .. })),
"the short-timeout waiter must time out on its own schedule, got {short_result:?}"
);
tokio::time::sleep(Duration::from_millis(150)).await;
assert!(
!long.is_finished(),
"a concurrent caller's short timeout must never resolve another caller's \
independent wait early"
);
long.abort();
}
#[test]
fn test_extract_hover_contents_string() {
let marked_string = lsp_types::MarkedString::String("Test hover".to_string());
let contents = lsp_types::Contents::MarkedString(marked_string);
let result = extract_hover_contents(contents);
assert_eq!(result, "Test hover");
}
#[test]
fn test_extract_hover_contents_language_string() {
let marked_string = lsp_types::MarkedString::MarkedStringWithLanguage(
lsp_types::MarkedStringWithLanguage {
language: "rust".to_string(),
value: "fn main() {}".to_string(),
},
);
let contents = lsp_types::Contents::MarkedString(marked_string);
let result = extract_hover_contents(contents);
assert_eq!(result, "```rust\nfn main() {}\n```");
}
#[test]
fn test_extract_hover_contents_markup() {
let markup = lsp_types::MarkupContent {
kind: lsp_types::MarkupKind::Markdown,
value: "# Documentation".to_string(),
};
let contents = lsp_types::Contents::MarkupContent(markup);
let result = extract_hover_contents(contents);
assert_eq!(result, "# Documentation");
}
#[tokio::test]
async fn test_handle_definition_flattens_single_location() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
definition_provider: Some(lsp_types::DefinitionProvider::Bool(true)),
..Default::default()
};
let (translator, mut server) = translator_with_capabilities(&dir, &server_id, caps);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let target_path = dir.path().join("target.rs");
fs::write(&target_path, "fn target() {}").unwrap();
let target_uri = Url::from_file_path(&target_path).unwrap().to_string();
let translator = Arc::new(translator);
let handle = {
let translator = Arc::clone(&translator);
let path = path.to_string_lossy().to_string();
tokio::spawn(async move {
translator
.handle_definition(
path,
Position {
line: 1,
character: 1,
},
)
.await
})
};
let mut wire = BufReader::new(&mut server.write_stdout);
let opened = read_framed_message(&mut wire).await;
assert_eq!(opened["method"], "textDocument/didOpen");
let request = read_framed_message(&mut wire).await;
assert_eq!(request["method"], "textDocument/definition");
write_response(
&mut server.read_half_stdin,
&request["id"],
serde_json::json!({
"uri": target_uri,
"range": {
"start": {"line": 0, "character": 0},
"end": {"line": 0, "character": 6}
}
}),
)
.await;
let result = timeout(Duration::from_secs(2), handle)
.await
.expect("handle_definition should not hang")
.unwrap()
.unwrap();
assert_eq!(result.locations.len(), 1);
assert_eq!(result.locations[0].uri, target_uri);
assert!(
!result.locations[0].out_of_workspace,
"a definition location inside the workspace root must not be marked out_of_workspace"
);
}
#[tokio::test]
async fn test_handle_definition_does_not_filter_out_of_workspace_location() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
definition_provider: Some(lsp_types::DefinitionProvider::Bool(true)),
..Default::default()
};
let (translator, mut server) = translator_with_capabilities(&dir, &server_id, caps);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let outside_uri = "file:///outside/workspace/stdlib.rs";
let translator = Arc::new(translator);
let handle = {
let translator = Arc::clone(&translator);
let path = path.to_string_lossy().to_string();
tokio::spawn(async move {
translator
.handle_definition(
path,
Position {
line: 1,
character: 1,
},
)
.await
})
};
let mut wire = BufReader::new(&mut server.write_stdout);
let opened = read_framed_message(&mut wire).await;
assert_eq!(opened["method"], "textDocument/didOpen");
let request = read_framed_message(&mut wire).await;
assert_eq!(request["method"], "textDocument/definition");
write_response(
&mut server.read_half_stdin,
&request["id"],
serde_json::json!({
"uri": outside_uri,
"range": {
"start": {"line": 0, "character": 0},
"end": {"line": 0, "character": 6}
}
}),
)
.await;
let result = timeout(Duration::from_secs(2), handle)
.await
.expect("handle_definition should not hang")
.unwrap()
.unwrap();
assert_eq!(
result.locations.len(),
1,
"an out-of-workspace definition location (e.g. stdlib/a dependency) must be \
returned, not dropped"
);
assert_eq!(result.locations[0].uri, outside_uri);
assert!(
result.locations[0].out_of_workspace,
"a definition location outside every workspace root must be marked out_of_workspace"
);
}
#[tokio::test]
async fn test_handle_references_does_not_filter_out_of_workspace_location() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
references_provider: Some(lsp_types::ReferencesProvider::Bool(true)),
..Default::default()
};
let (translator, mut server) = translator_with_capabilities(&dir, &server_id, caps);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let inside_path = dir.path().join("inside.rs");
fs::write(&inside_path, "fn used() {}").unwrap();
let inside_uri = Url::from_file_path(&inside_path).unwrap().to_string();
let outside_uri = "file:///outside/workspace/stdlib.rs";
let translator = Arc::new(translator);
let handle = {
let translator = Arc::clone(&translator);
let path = path.to_string_lossy().to_string();
tokio::spawn(async move { translator.handle_references(path, pos(1, 1), true).await })
};
let mut wire = BufReader::new(&mut server.write_stdout);
let opened = read_framed_message(&mut wire).await;
assert_eq!(opened["method"], "textDocument/didOpen");
let request = read_framed_message(&mut wire).await;
assert_eq!(request["method"], "textDocument/references");
write_response(
&mut server.read_half_stdin,
&request["id"],
serde_json::json!([
{
"uri": inside_uri,
"range": {
"start": {"line": 0, "character": 0},
"end": {"line": 0, "character": 4}
}
},
{
"uri": outside_uri,
"range": {
"start": {"line": 0, "character": 0},
"end": {"line": 0, "character": 4}
}
}
]),
)
.await;
let result = timeout(Duration::from_secs(2), handle)
.await
.expect("handle_references should not hang")
.unwrap()
.unwrap();
assert_eq!(
result.locations.len(),
2,
"both the in-workspace and out-of-workspace locations must survive"
);
assert!(result.locations.iter().any(|l| l.uri == inside_uri));
assert!(result.locations.iter().any(|l| l.uri == outside_uri));
assert!(
!result
.locations
.iter()
.find(|l| l.uri == inside_uri)
.unwrap()
.out_of_workspace,
"an in-workspace reference location must not be marked out_of_workspace"
);
assert!(
result
.locations
.iter()
.find(|l| l.uri == outside_uri)
.unwrap()
.out_of_workspace,
"an out-of-workspace reference location must be marked out_of_workspace"
);
}
#[tokio::test]
async fn test_handle_references_sets_truncated_flag_past_cap() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
references_provider: Some(lsp_types::ReferencesProvider::Bool(true)),
..Default::default()
};
let (translator, mut server) = translator_with_capabilities(&dir, &server_id, caps);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let uri = Url::from_file_path(&path).unwrap().to_string();
let translator = Arc::new(translator);
let handle = {
let translator = Arc::clone(&translator);
let path = path.to_string_lossy().to_string();
tokio::spawn(async move { translator.handle_references(path, pos(1, 1), true).await })
};
let mut wire = BufReader::new(&mut server.write_stdout);
let opened = read_framed_message(&mut wire).await;
assert_eq!(opened["method"], "textDocument/didOpen");
let request = read_framed_message(&mut wire).await;
assert_eq!(request["method"], "textDocument/references");
let locations: Vec<serde_json::Value> = (0..MAX_NORMALIZED_LOCATIONS + 500)
.map(|_| {
serde_json::json!({
"uri": uri,
"range": {
"start": {"line": 0, "character": 0},
"end": {"line": 0, "character": 4}
}
})
})
.collect();
write_response(
&mut server.read_half_stdin,
&request["id"],
serde_json::json!(locations),
)
.await;
let result = timeout(Duration::from_secs(5), handle)
.await
.expect("handle_references should not hang")
.unwrap()
.unwrap();
assert_eq!(result.locations.len(), MAX_NORMALIZED_LOCATIONS);
assert!(
result.truncated,
"a references response past MAX_NORMALIZED_LOCATIONS must set truncated: true"
);
}
#[tokio::test]
async fn test_handle_implementation_flattens_location_list() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
implementation_provider: Some(lsp_types::ImplementationProvider::Bool(true)),
..Default::default()
};
let (translator, mut server) = translator_with_capabilities(&dir, &server_id, caps);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let first_impl_path = dir.path().join("impl_a.rs");
fs::write(&first_impl_path, "struct A;").unwrap();
let first_impl_uri = Url::from_file_path(&first_impl_path).unwrap().to_string();
let second_impl_path = dir.path().join("impl_b.rs");
fs::write(&second_impl_path, "struct B;").unwrap();
let second_impl_uri = Url::from_file_path(&second_impl_path).unwrap().to_string();
let translator = Arc::new(translator);
let handle = {
let translator = Arc::clone(&translator);
let path = path.to_string_lossy().to_string();
tokio::spawn(async move {
translator
.handle_implementation(
path,
Position {
line: 1,
character: 1,
},
)
.await
})
};
let mut wire = BufReader::new(&mut server.write_stdout);
let opened = read_framed_message(&mut wire).await;
assert_eq!(opened["method"], "textDocument/didOpen");
let request = read_framed_message(&mut wire).await;
assert_eq!(request["method"], "textDocument/implementation");
write_response(
&mut server.read_half_stdin,
&request["id"],
serde_json::json!([
{
"uri": first_impl_uri,
"range": {
"start": {"line": 0, "character": 0},
"end": {"line": 0, "character": 9}
}
},
{
"uri": second_impl_uri,
"range": {
"start": {"line": 0, "character": 0},
"end": {"line": 0, "character": 9}
}
}
]),
)
.await;
let result = timeout(Duration::from_secs(2), handle)
.await
.expect("handle_implementation should not hang")
.unwrap()
.unwrap();
assert_eq!(result.locations.len(), 2);
assert_eq!(result.locations[0].uri, first_impl_uri);
assert_eq!(result.locations[1].uri, second_impl_uri);
}
#[tokio::test]
async fn test_handle_type_definition_flattens_definition_link_list() {
let dir = TempDir::new().unwrap();
let server_id = ServerId::from("rust");
let caps = lsp_types::ServerCapabilities {
type_definition_provider: Some(lsp_types::TypeDefinitionProvider::Bool(true)),
..Default::default()
};
let (translator, mut server) = translator_with_capabilities(&dir, &server_id, caps);
let path = dir.path().join("main.rs");
fs::write(&path, "fn main() {}").unwrap();
let target_path = dir.path().join("target_type.rs");
fs::write(&target_path, "struct TargetType;").unwrap();
let target_uri = Url::from_file_path(&target_path).unwrap().to_string();
let translator = Arc::new(translator);
let handle = {
let translator = Arc::clone(&translator);
let path = path.to_string_lossy().to_string();
tokio::spawn(async move {
translator
.handle_type_definition(
path,
Position {
line: 1,
character: 1,
},
)
.await
})
};
let mut wire = BufReader::new(&mut server.write_stdout);
let opened = read_framed_message(&mut wire).await;
assert_eq!(opened["method"], "textDocument/didOpen");
let request = read_framed_message(&mut wire).await;
assert_eq!(request["method"], "textDocument/typeDefinition");
write_response(
&mut server.read_half_stdin,
&request["id"],
serde_json::json!([{
"targetUri": target_uri,
"targetRange": {
"start": {"line": 0, "character": 0},
"end": {"line": 0, "character": 18}
},
"targetSelectionRange": {
"start": {"line": 0, "character": 7},
"end": {"line": 0, "character": 17}
}
}]),
)
.await;
let result = timeout(Duration::from_secs(2), handle)
.await
.expect("handle_type_definition should not hang")
.unwrap()
.unwrap();
assert_eq!(result.locations.len(), 1);
assert_eq!(result.locations[0].uri, target_uri);
assert_eq!(result.locations[0].range.start.character, 8);
}
#[tokio::test]
async fn test_lsp_locations_to_mcp_reads_disk_once_per_distinct_file_line() {
let dir = TempDir::new().unwrap();
let mut uris = Vec::new();
for i in 0..3 {
let path = dir.path().join(format!("file{i}.rs"));
fs::write(&path, "hello").unwrap();
uris.push(path_to_uri(&path).unwrap());
}
let ctx = test_ctx_with(PositionEncoding::Utf8);
let locs: Vec<lsp_types::Location> = (0..300)
.map(|i| lsp_types::Location {
uri: uris[i % 3].clone(),
range: lsp_types::Range {
start: lsp_types::Position {
line: 0,
character: 0,
},
end: lsp_types::Position {
line: 0,
character: 3,
},
},
})
.collect();
let result = lsp_locations_to_mcp(locs, &ctx).await;
assert_eq!(result.locations.len(), 300);
assert!(!result.truncated);
assert!(
result.locations.iter().all(|l| l.range.end.character == 4),
"MCP columns are 1-based, so LSP byte offset 3 in all-ASCII \"hello\" must convert \
to 4"
);
assert_eq!(
lock_std(&ctx.line_cache).entries.len(),
3,
"300 locations across 3 distinct files must populate the cache with exactly 3 \
entries (one disk read per distinct file/line), not one per location"
);
}
#[tokio::test]
async fn test_lsp_locations_to_mcp_truncates_to_max_normalized_locations() {
let ctx = test_ctx();
let uri = test_uri();
let locs: Vec<lsp_types::Location> = (0..MAX_NORMALIZED_LOCATIONS + 500)
.map(|_| lsp_types::Location {
uri: uri.clone(),
range: lsp_types::Range {
start: lsp_types::Position {
line: 0,
character: 0,
},
end: lsp_types::Position {
line: 0,
character: 1,
},
},
})
.collect();
let result = lsp_locations_to_mcp(locs, &ctx).await;
assert_eq!(result.locations.len(), MAX_NORMALIZED_LOCATIONS);
assert!(
result.truncated,
"a response naming more than MAX_NORMALIZED_LOCATIONS must report truncated: true"
);
}
}