lsp: Serialize LSP notifications on background threads (#39403)
This should reduce hiccups when opening large files. Release Notes: - N/A
This commit is contained in:
+59
-29
@@ -80,11 +80,14 @@ pub struct LanguageServerBinaryOptions {
|
||||
pub pre_release: bool,
|
||||
}
|
||||
|
||||
struct NotificationSerializer(Box<dyn FnOnce() -> String + Send + Sync>);
|
||||
|
||||
/// A running language server process.
|
||||
pub struct LanguageServer {
|
||||
server_id: LanguageServerId,
|
||||
next_id: AtomicI32,
|
||||
outbound_tx: channel::Sender<String>,
|
||||
notification_tx: channel::Sender<NotificationSerializer>,
|
||||
name: LanguageServerName,
|
||||
process_name: Arc<str>,
|
||||
binary: LanguageServerBinary,
|
||||
@@ -477,9 +480,24 @@ impl LanguageServer {
|
||||
}
|
||||
.into();
|
||||
|
||||
let (notification_tx, notification_rx) = channel::unbounded::<NotificationSerializer>();
|
||||
cx.background_spawn({
|
||||
let outbound_tx = outbound_tx.clone();
|
||||
async move {
|
||||
while let Ok(serializer) = notification_rx.recv().await {
|
||||
let serialized = (serializer.0)();
|
||||
let Ok(_) = outbound_tx.send(serialized).await else {
|
||||
return;
|
||||
};
|
||||
}
|
||||
outbound_tx.close();
|
||||
}
|
||||
})
|
||||
.detach();
|
||||
Self {
|
||||
server_id,
|
||||
notification_handlers,
|
||||
notification_tx,
|
||||
response_handlers,
|
||||
io_handlers,
|
||||
name: server_name,
|
||||
@@ -906,7 +924,7 @@ impl LanguageServer {
|
||||
self.capabilities = RwLock::new(response.capabilities);
|
||||
self.configuration = configuration;
|
||||
|
||||
self.notify::<notification::Initialized>(&InitializedParams {})?;
|
||||
self.notify::<notification::Initialized>(InitializedParams {})?;
|
||||
Ok(Arc::new(self))
|
||||
})
|
||||
}
|
||||
@@ -918,11 +936,13 @@ impl LanguageServer {
|
||||
let next_id = AtomicI32::new(self.next_id.load(SeqCst));
|
||||
let outbound_tx = self.outbound_tx.clone();
|
||||
let executor = self.executor.clone();
|
||||
let notification_serializers = self.notification_tx.clone();
|
||||
let mut output_done = self.output_done_rx.lock().take().unwrap();
|
||||
let shutdown_request = Self::request_internal::<request::Shutdown>(
|
||||
&next_id,
|
||||
&response_handlers,
|
||||
&outbound_tx,
|
||||
¬ification_serializers,
|
||||
&executor,
|
||||
(),
|
||||
);
|
||||
@@ -956,8 +976,8 @@ impl LanguageServer {
|
||||
}
|
||||
|
||||
response_handlers.lock().take();
|
||||
Self::notify_internal::<notification::Exit>(&outbound_tx, &()).ok();
|
||||
outbound_tx.close();
|
||||
Self::notify_internal::<notification::Exit>(¬ification_serializers, ()).ok();
|
||||
notification_serializers.close();
|
||||
output_done.recv().await;
|
||||
server.lock().take().map(|mut child| child.kill());
|
||||
drop(tasks);
|
||||
@@ -1179,6 +1199,7 @@ impl LanguageServer {
|
||||
&self.next_id,
|
||||
&self.response_handlers,
|
||||
&self.outbound_tx,
|
||||
&self.notification_tx,
|
||||
&self.executor,
|
||||
params,
|
||||
)
|
||||
@@ -1200,6 +1221,7 @@ impl LanguageServer {
|
||||
&self.next_id,
|
||||
&self.response_handlers,
|
||||
&self.outbound_tx,
|
||||
&self.notification_tx,
|
||||
&self.executor,
|
||||
timer,
|
||||
params,
|
||||
@@ -1210,6 +1232,7 @@ impl LanguageServer {
|
||||
next_id: &AtomicI32,
|
||||
response_handlers: &Mutex<Option<HashMap<RequestId, ResponseHandler>>>,
|
||||
outbound_tx: &channel::Sender<String>,
|
||||
notification_serializers: &channel::Sender<NotificationSerializer>,
|
||||
executor: &BackgroundExecutor,
|
||||
timer: U,
|
||||
params: T::Params,
|
||||
@@ -1261,7 +1284,7 @@ impl LanguageServer {
|
||||
.try_send(message)
|
||||
.context("failed to write to language server's stdin");
|
||||
|
||||
let outbound_tx = outbound_tx.downgrade();
|
||||
let notification_serializers = notification_serializers.downgrade();
|
||||
let started = Instant::now();
|
||||
LspRequest::new(id, async move {
|
||||
if let Err(e) = handle_response {
|
||||
@@ -1272,10 +1295,10 @@ impl LanguageServer {
|
||||
}
|
||||
|
||||
let cancel_on_drop = util::defer(move || {
|
||||
if let Some(outbound_tx) = outbound_tx.upgrade() {
|
||||
if let Some(notification_serializers) = notification_serializers.upgrade() {
|
||||
Self::notify_internal::<notification::Cancel>(
|
||||
&outbound_tx,
|
||||
&CancelParams {
|
||||
¬ification_serializers,
|
||||
CancelParams {
|
||||
id: NumberOrString::Number(id),
|
||||
},
|
||||
)
|
||||
@@ -1310,6 +1333,7 @@ impl LanguageServer {
|
||||
next_id: &AtomicI32,
|
||||
response_handlers: &Mutex<Option<HashMap<RequestId, ResponseHandler>>>,
|
||||
outbound_tx: &channel::Sender<String>,
|
||||
notification_serializers: &channel::Sender<NotificationSerializer>,
|
||||
executor: &BackgroundExecutor,
|
||||
params: T::Params,
|
||||
) -> impl LspRequestFuture<T::Result> + use<T>
|
||||
@@ -1321,6 +1345,7 @@ impl LanguageServer {
|
||||
next_id,
|
||||
response_handlers,
|
||||
outbound_tx,
|
||||
notification_serializers,
|
||||
executor,
|
||||
Self::default_request_timer(executor.clone()),
|
||||
params,
|
||||
@@ -1336,21 +1361,25 @@ impl LanguageServer {
|
||||
/// Sends a RPC notification to the language server.
|
||||
///
|
||||
/// [LSP Specification](https://microsoft.github.io/language-server-protocol/specifications/lsp/3.17/specification/#notificationMessage)
|
||||
pub fn notify<T: notification::Notification>(&self, params: &T::Params) -> Result<()> {
|
||||
Self::notify_internal::<T>(&self.outbound_tx, params)
|
||||
pub fn notify<T: notification::Notification>(&self, params: T::Params) -> Result<()> {
|
||||
let outbound = self.notification_tx.clone();
|
||||
Self::notify_internal::<T>(&outbound, params)
|
||||
}
|
||||
|
||||
fn notify_internal<T: notification::Notification>(
|
||||
outbound_tx: &channel::Sender<String>,
|
||||
params: &T::Params,
|
||||
outbound_tx: &channel::Sender<NotificationSerializer>,
|
||||
params: T::Params,
|
||||
) -> Result<()> {
|
||||
let message = serde_json::to_string(&Notification {
|
||||
jsonrpc: JSON_RPC_VERSION,
|
||||
method: T::METHOD,
|
||||
params,
|
||||
})
|
||||
.unwrap();
|
||||
outbound_tx.try_send(message)?;
|
||||
let serializer = NotificationSerializer(Box::new(move || {
|
||||
serde_json::to_string(&Notification {
|
||||
jsonrpc: JSON_RPC_VERSION,
|
||||
method: T::METHOD,
|
||||
params,
|
||||
})
|
||||
.unwrap()
|
||||
}));
|
||||
|
||||
outbound_tx.send_blocking(serializer)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -1385,7 +1414,7 @@ impl LanguageServer {
|
||||
removed: vec![],
|
||||
},
|
||||
};
|
||||
self.notify::<DidChangeWorkspaceFolders>(¶ms).ok();
|
||||
self.notify::<DidChangeWorkspaceFolders>(params).ok();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1419,7 +1448,7 @@ impl LanguageServer {
|
||||
}],
|
||||
},
|
||||
};
|
||||
self.notify::<DidChangeWorkspaceFolders>(¶ms).ok();
|
||||
self.notify::<DidChangeWorkspaceFolders>(params).ok();
|
||||
}
|
||||
}
|
||||
pub fn set_workspace_folders(&self, folders: BTreeSet<Uri>) {
|
||||
@@ -1451,7 +1480,7 @@ impl LanguageServer {
|
||||
let params = DidChangeWorkspaceFoldersParams {
|
||||
event: WorkspaceFoldersChangeEvent { added, removed },
|
||||
};
|
||||
self.notify::<DidChangeWorkspaceFolders>(¶ms).ok();
|
||||
self.notify::<DidChangeWorkspaceFolders>(params).ok();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1469,14 +1498,14 @@ impl LanguageServer {
|
||||
version: i32,
|
||||
initial_text: String,
|
||||
) {
|
||||
self.notify::<notification::DidOpenTextDocument>(&DidOpenTextDocumentParams {
|
||||
self.notify::<notification::DidOpenTextDocument>(DidOpenTextDocumentParams {
|
||||
text_document: TextDocumentItem::new(uri, language_id, version, initial_text),
|
||||
})
|
||||
.ok();
|
||||
}
|
||||
|
||||
pub fn unregister_buffer(&self, uri: Uri) {
|
||||
self.notify::<notification::DidCloseTextDocument>(&DidCloseTextDocumentParams {
|
||||
self.notify::<notification::DidCloseTextDocument>(DidCloseTextDocumentParams {
|
||||
text_document: TextDocumentIdentifier::new(uri),
|
||||
})
|
||||
.ok();
|
||||
@@ -1692,7 +1721,7 @@ impl LanguageServer {
|
||||
#[cfg(any(test, feature = "test-support"))]
|
||||
impl FakeLanguageServer {
|
||||
/// See [`LanguageServer::notify`].
|
||||
pub fn notify<T: notification::Notification>(&self, params: &T::Params) {
|
||||
pub fn notify<T: notification::Notification>(&self, params: T::Params) {
|
||||
self.server.notify::<T>(params).ok();
|
||||
}
|
||||
|
||||
@@ -1801,7 +1830,7 @@ impl FakeLanguageServer {
|
||||
.await
|
||||
.into_response()
|
||||
.unwrap();
|
||||
self.notify::<notification::Progress>(&ProgressParams {
|
||||
self.notify::<notification::Progress>(ProgressParams {
|
||||
token: NumberOrString::String(token),
|
||||
value: ProgressParamsValue::WorkDone(WorkDoneProgress::Begin(progress)),
|
||||
});
|
||||
@@ -1809,7 +1838,7 @@ impl FakeLanguageServer {
|
||||
|
||||
/// Simulate that the server has completed work and notifies about that with the specified token.
|
||||
pub fn end_progress(&self, token: impl Into<String>) {
|
||||
self.notify::<notification::Progress>(&ProgressParams {
|
||||
self.notify::<notification::Progress>(ProgressParams {
|
||||
token: NumberOrString::String(token.into()),
|
||||
value: ProgressParamsValue::WorkDone(WorkDoneProgress::End(Default::default())),
|
||||
});
|
||||
@@ -1868,7 +1897,7 @@ mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
server
|
||||
.notify::<notification::DidOpenTextDocument>(&DidOpenTextDocumentParams {
|
||||
.notify::<notification::DidOpenTextDocument>(DidOpenTextDocumentParams {
|
||||
text_document: TextDocumentItem::new(
|
||||
Uri::from_str("file://a/b").unwrap(),
|
||||
"rust".to_string(),
|
||||
@@ -1886,11 +1915,11 @@ mod tests {
|
||||
"file://a/b"
|
||||
);
|
||||
|
||||
fake.notify::<notification::ShowMessage>(&ShowMessageParams {
|
||||
fake.notify::<notification::ShowMessage>(ShowMessageParams {
|
||||
typ: MessageType::ERROR,
|
||||
message: "ok".to_string(),
|
||||
});
|
||||
fake.notify::<notification::PublishDiagnostics>(&PublishDiagnosticsParams {
|
||||
fake.notify::<notification::PublishDiagnostics>(PublishDiagnosticsParams {
|
||||
uri: Uri::from_str("file://b/c").unwrap(),
|
||||
version: Some(5),
|
||||
diagnostics: vec![],
|
||||
@@ -1904,6 +1933,7 @@ mod tests {
|
||||
fake.set_request_handler::<request::Shutdown, _, _>(|_, _| async move { Ok(()) });
|
||||
|
||||
drop(server);
|
||||
cx.run_until_parked();
|
||||
fake.receive_notification::<notification::Exit>().await;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user