use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use lsp_types::notification::Progress;
use lsp_types::request::WorkDoneProgressCreate;
use lsp_types::{
ProgressParams, ProgressParamsValue, ProgressToken, WorkDoneProgress, WorkDoneProgressBegin,
WorkDoneProgressCreateParams, WorkDoneProgressEnd, WorkDoneProgressReport,
};
use tokio_util::sync::CancellationToken;
use tracing::warn;
use crate::client::Client;
use crate::error::{ClientError, ProgressError};
const MAX_PROGRESS_TOKEN: u32 = i32::MAX as u32;
pub(crate) struct ActiveProgress {
pub(crate) cancellable: bool,
pub(crate) cancellation: CancellationToken,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ProgressCancel {
Cancelled,
NotCancellable,
NotActive,
}
#[derive(Clone, Default)]
pub(crate) struct ProgressRegistry {
inner: Arc<Mutex<ProgressInner>>,
}
struct ProgressInner {
next_token: u32,
active: HashMap<ProgressToken, ActiveProgress>,
}
impl Default for ProgressInner {
fn default() -> Self {
Self {
next_token: 1,
active: HashMap::new(),
}
}
}
impl ProgressRegistry {
pub(crate) fn allocate(&self) -> Option<ProgressToken> {
let mut inner = self.inner.lock().unwrap();
let mut candidate = inner.next_token;
loop {
if candidate > MAX_PROGRESS_TOKEN {
return None;
}
let token = ProgressToken::Number(candidate as i32);
if !inner.active.contains_key(&token) {
inner.next_token = candidate + 1;
return Some(token);
}
candidate += 1;
}
}
pub(crate) fn register(
&self,
token: ProgressToken,
cancellable: bool,
cancellation: CancellationToken,
) -> bool {
let mut inner = self.inner.lock().unwrap();
if inner.active.contains_key(&token) {
return false;
}
inner.active.insert(
token,
ActiveProgress {
cancellable,
cancellation,
},
);
true
}
pub(crate) fn remove(&self, token: &ProgressToken) -> Option<ActiveProgress> {
self.inner.lock().unwrap().active.remove(token)
}
pub(crate) fn cancel(&self, token: &ProgressToken) -> ProgressCancel {
let inner = self.inner.lock().unwrap();
match inner.active.get(token) {
Some(entry) if entry.cancellable => {
entry.cancellation.cancel();
ProgressCancel::Cancelled
}
Some(_) => ProgressCancel::NotCancellable,
None => ProgressCancel::NotActive,
}
}
pub(crate) fn clear(&self) {
self.inner.lock().unwrap().active.clear();
}
pub(crate) fn is_active(&self, token: &ProgressToken) -> bool {
self.inner.lock().unwrap().active.contains_key(token)
}
#[cfg(test)]
pub(crate) fn active_len(&self) -> usize {
self.inner.lock().unwrap().active.len()
}
#[cfg(test)]
pub(crate) fn set_next_token(&self, token: u32) {
self.inner.lock().unwrap().next_token = token;
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProgressOptions {
pub title: String,
pub cancellable: bool,
pub message: Option<String>,
pub percentage: Option<u32>,
}
impl ProgressOptions {
pub fn new(title: impl Into<String>) -> Self {
Self {
title: title.into(),
cancellable: false,
message: None,
percentage: None,
}
}
#[must_use]
pub fn cancellable(mut self, cancellable: bool) -> Self {
self.cancellable = cancellable;
self
}
#[must_use]
pub fn message(mut self, message: impl Into<String>) -> Self {
self.message = Some(message.into());
self
}
#[must_use]
pub fn percentage(mut self, percentage: u32) -> Self {
self.percentage = Some(percentage);
self
}
}
#[derive(Clone)]
pub(crate) struct SharedProgress {
inner: Arc<SharedInner>,
}
struct SharedInner {
client: Client,
registry: ProgressRegistry,
token: ProgressToken,
cancellable: bool,
cancellation: CancellationToken,
ended: AtomicBool,
}
impl SharedProgress {
fn new(
client: Client,
registry: ProgressRegistry,
token: ProgressToken,
cancellable: bool,
cancellation: CancellationToken,
) -> Self {
Self {
inner: Arc::new(SharedInner {
client,
registry,
token,
cancellable,
cancellation,
ended: AtomicBool::new(false),
}),
}
}
fn is_ended(&self) -> bool {
self.inner.ended.load(Ordering::SeqCst)
}
fn report(&self, message: Option<String>, percentage: Option<u32>) -> Result<(), ClientError> {
if self.is_ended() {
return Err(ProgressError::AlreadyEnded.into());
}
if self.inner.cancellation.is_cancelled() {
return Err(ProgressError::Cancelled.into());
}
if !self.inner.registry.is_active(&self.inner.token) {
return Err(ProgressError::UnknownToken.into());
}
if let Some(percentage) = percentage
&& percentage > 100
{
return Err(ProgressError::InvalidPercentage(percentage).into());
}
self.inner.client.notify_logged::<Progress>(ProgressParams {
token: self.inner.token.clone(),
value: ProgressParamsValue::WorkDone(WorkDoneProgress::Report(
WorkDoneProgressReport {
cancellable: Some(self.inner.cancellable),
message,
percentage,
},
)),
})
}
fn end(&self, message: Option<String>) -> Result<(), ClientError> {
if self.inner.ended.swap(true, Ordering::SeqCst) {
return Err(ProgressError::AlreadyEnded.into());
}
let result = self.inner.client.notify_logged::<Progress>(ProgressParams {
token: self.inner.token.clone(),
value: ProgressParamsValue::WorkDone(WorkDoneProgress::End(WorkDoneProgressEnd {
message,
})),
});
self.inner.registry.remove(&self.inner.token);
result
}
}
pub struct ProgressHandle {
shared: SharedProgress,
}
impl std::fmt::Debug for ProgressHandle {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("ProgressHandle")
.field("token", &self.shared.inner.token)
.finish_non_exhaustive()
}
}
impl ProgressHandle {
pub fn token(&self) -> ProgressToken {
self.shared.inner.token.clone()
}
pub fn cancellation_token(&self) -> CancellationToken {
self.shared.inner.cancellation.clone()
}
pub fn report(
&self,
message: Option<String>,
percentage: Option<u32>,
) -> Result<(), ClientError> {
self.shared.report(message, percentage)
}
pub fn end(self, message: Option<String>) -> Result<(), ClientError> {
self.shared.end(message)
}
}
impl Drop for ProgressHandle {
fn drop(&mut self) {
if !self.shared.is_ended() {
self.shared.inner.registry.remove(&self.shared.inner.token);
warn!(
token = ?self.shared.inner.token,
"progress handle dropped without end; token removed, no end notification sent"
);
}
}
}
impl Client {
pub async fn begin_progress(
&self,
options: ProgressOptions,
) -> Result<ProgressHandle, ClientError> {
if let Some(percentage) = options.percentage
&& percentage > 100
{
return Err(ClientError::InvalidHelperParams(format!(
"work-done begin percentage {percentage} is outside the range 0..=100"
)));
}
let registry = self.progress_registry().clone();
let token = registry.allocate().ok_or(ClientError::IdExhausted)?;
self.request::<WorkDoneProgressCreate>(WorkDoneProgressCreateParams {
token: token.clone(),
})
.await?;
let cancellation = CancellationToken::new();
if !registry.register(token.clone(), options.cancellable, cancellation.clone()) {
registry.remove(&token);
return Err(ProgressError::UnknownToken.into());
}
let begin = ProgressParams {
token: token.clone(),
value: ProgressParamsValue::WorkDone(WorkDoneProgress::Begin(WorkDoneProgressBegin {
title: options.title,
cancellable: Some(options.cancellable),
message: options.message,
percentage: options.percentage,
})),
};
if let Err(error) = self.notify_logged::<Progress>(begin) {
registry.remove(&token);
return Err(error);
}
Ok(ProgressHandle {
shared: SharedProgress::new(
self.clone(),
registry,
token,
options.cancellable,
cancellation,
),
})
}
}
#[cfg(test)]
mod tests {
use bytes::Bytes;
use lsp_types::notification::Notification;
use serde_json::{Value, json};
use tokio::sync::mpsc;
use tracing_subscriber::layer::SubscriberExt;
use super::*;
use crate::client::{OutboundQueue, OutboundRegistry};
use crate::raw::{RawMessage, RequestId};
fn make_client() -> (
Client,
OutboundRegistry,
mpsc::UnboundedReceiver<RawMessage>,
) {
let (outgoing, receiver) = OutboundQueue::new(crate::DEFAULT_OUTBOUND_WARNING_THRESHOLD);
let outbound = OutboundRegistry::default();
let client = Client::new(outgoing, outbound.clone());
(client, outbound, receiver)
}
async fn answer_create_ok(
outbound: &OutboundRegistry,
receiver: &mut mpsc::UnboundedReceiver<RawMessage>,
) -> Value {
let message = receiver.recv().await.expect("create request");
match message {
RawMessage::Request { id, method, params } => {
assert_eq!(method, "window/workDoneProgress/create");
let RequestId::Number(id) = id else {
panic!("numeric request id");
};
let params: Value = serde_json::from_slice(¶ms).unwrap();
assert!(outbound.complete(id as u32, Ok(Bytes::from_static(b"null"))));
params
}
other => panic!("expected create request, got {other:?}"),
}
}
async fn next_notification(
receiver: &mut mpsc::UnboundedReceiver<RawMessage>,
) -> (String, Value) {
match receiver.recv().await.expect("notification") {
RawMessage::Notification { method, params } => (
method.into_owned(),
serde_json::from_slice(¶ms).unwrap(),
),
other => panic!("expected notification, got {other:?}"),
}
}
#[test]
fn options_default_to_not_cancellable_without_message_or_percentage() {
let options = ProgressOptions::new("Indexing");
assert_eq!(options.title, "Indexing");
assert!(!options.cancellable);
assert_eq!(options.message, None);
assert_eq!(options.percentage, None);
let options = ProgressOptions::new("Indexing")
.cancellable(true)
.message("starting")
.percentage(10);
assert!(options.cancellable);
assert_eq!(options.message.as_deref(), Some("starting"));
assert_eq!(options.percentage, Some(10));
}
#[test]
fn tokens_start_at_one_and_increase_monotonically() {
let registry = ProgressRegistry::default();
assert_eq!(registry.allocate(), Some(ProgressToken::Number(1)));
assert_eq!(registry.allocate(), Some(ProgressToken::Number(2)));
assert_eq!(registry.allocate(), Some(ProgressToken::Number(3)));
}
#[test]
fn allocation_skips_active_server_and_client_originated_tokens() {
let registry = ProgressRegistry::default();
assert!(registry.register(ProgressToken::Number(1), false, CancellationToken::new()));
assert!(registry.register(ProgressToken::Number(2), false, CancellationToken::new()));
assert!(registry.register(
ProgressToken::String("client-token".into()),
false,
CancellationToken::new()
));
assert_eq!(registry.allocate(), Some(ProgressToken::Number(3)));
registry.remove(&ProgressToken::Number(1));
assert_eq!(registry.allocate(), Some(ProgressToken::Number(4)));
assert!(!registry.register(ProgressToken::Number(2), true, CancellationToken::new()));
}
#[test]
fn allocation_exhausts_at_the_positive_i32_boundary() {
let registry = ProgressRegistry::default();
registry.set_next_token(MAX_PROGRESS_TOKEN);
assert_eq!(registry.allocate(), Some(ProgressToken::Number(i32::MAX)));
assert_eq!(registry.allocate(), None);
}
#[test]
fn cancel_fires_only_a_matching_active_cancellable_token() {
let registry = ProgressRegistry::default();
let cancellable = CancellationToken::new();
let plain = CancellationToken::new();
assert!(registry.register(ProgressToken::Number(1), true, cancellable.clone()));
assert!(registry.register(ProgressToken::Number(2), false, plain.clone()));
assert_eq!(
registry.cancel(&ProgressToken::Number(2)),
ProgressCancel::NotCancellable
);
assert!(!plain.is_cancelled());
assert_eq!(
registry.cancel(&ProgressToken::Number(3)),
ProgressCancel::NotActive
);
assert_eq!(
registry.cancel(&ProgressToken::String("ended".into())),
ProgressCancel::NotActive
);
assert_eq!(
registry.cancel(&ProgressToken::Number(1)),
ProgressCancel::Cancelled
);
assert!(cancellable.is_cancelled());
assert!(registry.is_active(&ProgressToken::Number(1)));
assert_eq!(
registry.cancel(&ProgressToken::Number(1)),
ProgressCancel::Cancelled
);
}
#[test]
fn clear_removes_every_active_token() {
let registry = ProgressRegistry::default();
let token = registry.allocate().expect("the token space is fresh");
assert!(registry.register(token.clone(), true, CancellationToken::new()));
assert!(registry.register(
ProgressToken::String("client-token".into()),
false,
CancellationToken::new()
));
registry.clear();
assert!(!registry.is_active(&token));
assert!(!registry.is_active(&ProgressToken::String("client-token".into())));
assert_eq!(registry.active_len(), 0);
assert_eq!(registry.allocate(), Some(ProgressToken::Number(2)));
}
#[tokio::test]
async fn begin_progress_happy_path_sends_create_then_one_begin() {
let (client, outbound, mut receiver) = make_client();
let registry = client.progress_registry().clone();
let begin = tokio::spawn({
let client = client.clone();
async move {
client
.begin_progress(
ProgressOptions::new("Indexing")
.cancellable(true)
.message("starting")
.percentage(0),
)
.await
}
});
let params = answer_create_ok(&outbound, &mut receiver).await;
assert_eq!(params, json!({ "token": 1 }));
let handle = begin.await.unwrap().expect("begin succeeds");
assert_eq!(handle.token(), ProgressToken::Number(1));
assert!(!handle.cancellation_token().is_cancelled());
let (method, params) = next_notification(&mut receiver).await;
assert_eq!(method, Progress::METHOD);
assert_eq!(
params,
json!({
"token": 1,
"value": {
"kind": "begin",
"title": "Indexing",
"cancellable": true,
"message": "starting",
"percentage": 0
}
})
);
assert!(registry.is_active(&ProgressToken::Number(1)));
assert!(receiver.try_recv().is_err(), "nothing else is sent");
handle.end(Some("done".into())).unwrap();
}
#[tokio::test]
async fn invalid_begin_percentage_is_rejected_before_any_io() {
let (client, _outbound, mut receiver) = make_client();
let registry = client.progress_registry().clone();
let error = client
.begin_progress(ProgressOptions::new("Indexing").percentage(101))
.await
.unwrap_err();
assert!(matches!(error, ClientError::InvalidHelperParams(_)));
assert!(
receiver.try_recv().is_err(),
"no request or notification sent"
);
assert_eq!(registry.active_len(), 0, "no token registered");
}
#[tokio::test]
async fn create_remote_failure_sends_no_begin_and_registers_no_token() {
let (client, outbound, mut receiver) = make_client();
let registry = client.progress_registry().clone();
let begin = tokio::spawn({
let client = client.clone();
async move {
client
.begin_progress(ProgressOptions::new("Indexing"))
.await
}
});
match receiver.recv().await.expect("create request") {
RawMessage::Request { id, .. } => {
let RequestId::Number(id) = id else {
panic!("numeric id")
};
assert!(outbound.complete(
id as u32,
Err(crate::raw::JsonRpcError {
code: -32803,
message: "Request failed".into(),
data: None,
})
));
}
other => panic!("expected create request, got {other:?}"),
}
let error = begin.await.unwrap().unwrap_err();
assert!(matches!(error, ClientError::Remote(_)));
assert!(receiver.try_recv().is_err(), "no begin notification sent");
assert_eq!(registry.active_len(), 0, "no token left registered");
}
#[tokio::test]
async fn begin_enqueue_failure_removes_the_token() {
let (client, outbound, mut receiver) = make_client();
let registry = client.progress_registry().clone();
let begin = tokio::spawn({
let client = client.clone();
async move {
client
.begin_progress(ProgressOptions::new("Indexing"))
.await
}
});
answer_create_ok(&outbound, &mut receiver).await;
drop(receiver);
let error = begin.await.unwrap().unwrap_err();
assert!(matches!(error, ClientError::OutboundClosed));
assert_eq!(registry.active_len(), 0, "begin failure removed the token");
}
#[tokio::test]
async fn report_sends_the_exact_work_done_report_shape() {
let (client, outbound, mut receiver) = make_client();
let begin = tokio::spawn({
let client = client.clone();
async move {
client
.begin_progress(ProgressOptions::new("Indexing").cancellable(true))
.await
}
});
answer_create_ok(&outbound, &mut receiver).await;
let handle = begin.await.unwrap().unwrap();
next_notification(&mut receiver).await;
handle.report(Some("half".into()), Some(50)).unwrap();
handle.report(None, Some(100)).unwrap();
handle.report(Some("restarted".into()), Some(0)).unwrap();
let (method, params) = next_notification(&mut receiver).await;
assert_eq!(method, Progress::METHOD);
assert_eq!(
params,
json!({
"token": 1,
"value": {
"kind": "report",
"cancellable": true,
"message": "half",
"percentage": 50
}
})
);
let (_, params) = next_notification(&mut receiver).await;
assert_eq!(
params,
json!({ "token": 1, "value": { "kind": "report", "cancellable": true, "percentage": 100 } })
);
let (_, params) = next_notification(&mut receiver).await;
assert_eq!(
params,
json!({
"token": 1,
"value": { "kind": "report", "cancellable": true, "message": "restarted", "percentage": 0 }
})
);
handle.end(None).unwrap();
}
#[tokio::test]
async fn invalid_report_percentage_sends_nothing() {
let (client, outbound, mut receiver) = make_client();
let begin = tokio::spawn({
let client = client.clone();
async move {
client
.begin_progress(ProgressOptions::new("Indexing"))
.await
}
});
answer_create_ok(&outbound, &mut receiver).await;
let handle = begin.await.unwrap().unwrap();
next_notification(&mut receiver).await;
let error = handle.report(None, Some(101)).unwrap_err();
assert!(matches!(
error,
ClientError::Progress(ProgressError::InvalidPercentage(101))
));
assert!(receiver.try_recv().is_err(), "invalid report sent nothing");
handle.end(None).unwrap();
}
#[tokio::test]
async fn cancelled_handle_reports_nothing_but_can_still_end() {
let (client, outbound, mut receiver) = make_client();
let begin = tokio::spawn({
let client = client.clone();
async move {
client
.begin_progress(ProgressOptions::new("Indexing").cancellable(true))
.await
}
});
answer_create_ok(&outbound, &mut receiver).await;
let handle = begin.await.unwrap().unwrap();
next_notification(&mut receiver).await;
handle.cancellation_token().cancel();
let error = handle.report(None, Some(10)).unwrap_err();
assert!(matches!(
error,
ClientError::Progress(ProgressError::Cancelled)
));
assert!(
receiver.try_recv().is_err(),
"cancelled report sent nothing"
);
handle.end(Some("stopped".into())).unwrap();
let (method, params) = next_notification(&mut receiver).await;
assert_eq!(method, Progress::METHOD);
assert_eq!(
params,
json!({ "token": 1, "value": { "kind": "end", "message": "stopped" } })
);
}
#[tokio::test]
async fn end_consumes_the_handle_and_removes_the_token() {
let (client, outbound, mut receiver) = make_client();
let registry = client.progress_registry().clone();
let begin = tokio::spawn({
let client = client.clone();
async move {
client
.begin_progress(ProgressOptions::new("Indexing"))
.await
}
});
answer_create_ok(&outbound, &mut receiver).await;
let handle = begin.await.unwrap().unwrap();
next_notification(&mut receiver).await; assert!(registry.is_active(&ProgressToken::Number(1)));
handle.end(Some("done".into())).unwrap();
let (method, params) = next_notification(&mut receiver).await;
assert_eq!(method, Progress::METHOD);
assert_eq!(
params,
json!({ "token": 1, "value": { "kind": "end", "message": "done" } })
);
assert!(!registry.is_active(&ProgressToken::Number(1)));
assert_eq!(registry.active_len(), 0);
}
#[tokio::test]
async fn end_enqueue_failure_still_removes_the_token() {
let (client, outbound, mut receiver) = make_client();
let registry = client.progress_registry().clone();
let begin = tokio::spawn({
let client = client.clone();
async move {
client
.begin_progress(ProgressOptions::new("Indexing"))
.await
}
});
answer_create_ok(&outbound, &mut receiver).await;
let handle = begin.await.unwrap().unwrap();
next_notification(&mut receiver).await;
drop(receiver);
let error = handle.end(None).unwrap_err();
assert!(matches!(error, ClientError::OutboundClosed));
assert_eq!(registry.active_len(), 0, "failed end removed the token");
}
#[tokio::test]
async fn repeated_shared_state_operations_fail_without_sending() {
let (client, outbound, mut receiver) = make_client();
let registry = client.progress_registry().clone();
let begin = tokio::spawn({
let client = client.clone();
async move {
client
.begin_progress(ProgressOptions::new("Indexing"))
.await
}
});
answer_create_ok(&outbound, &mut receiver).await;
let handle = begin.await.unwrap().unwrap();
next_notification(&mut receiver).await;
let shared = handle.shared.clone();
registry.remove(&ProgressToken::Number(1));
let error = shared.report(None, None).unwrap_err();
assert!(matches!(
error,
ClientError::Progress(ProgressError::UnknownToken)
));
assert!(registry.register(ProgressToken::Number(1), false, CancellationToken::new()));
shared.end(None).unwrap();
let (method, _) = next_notification(&mut receiver).await;
assert_eq!(method, Progress::METHOD);
let error = shared.end(None).unwrap_err();
assert!(matches!(
error,
ClientError::Progress(ProgressError::AlreadyEnded)
));
let error = shared.report(None, None).unwrap_err();
assert!(matches!(
error,
ClientError::Progress(ProgressError::AlreadyEnded)
));
assert!(
receiver.try_recv().is_err(),
"repeated operations sent nothing"
);
drop(handle);
assert_eq!(registry.active_len(), 0);
}
#[tokio::test]
async fn dropping_an_active_handle_removes_the_token_without_io() {
let (client, outbound, mut receiver) = make_client();
let registry = client.progress_registry().clone();
let begin = tokio::spawn({
let client = client.clone();
async move {
client
.begin_progress(ProgressOptions::new("Indexing"))
.await
}
});
answer_create_ok(&outbound, &mut receiver).await;
let handle = begin.await.unwrap().unwrap();
next_notification(&mut receiver).await; assert!(registry.is_active(&ProgressToken::Number(1)));
let events = crate::test_util::EventCapture::new();
let subscriber = tracing_subscriber::registry().with(events.clone());
tracing::subscriber::with_default(subscriber, || drop(handle));
assert_eq!(registry.active_len(), 0, "drop removed the token");
assert!(
receiver.try_recv().is_err(),
"drop performed no I/O and sent no implicit end"
);
assert!(
events.contains("progress handle dropped without end"),
"drop logged a warning, got {:?}",
events.messages()
);
}
#[tokio::test]
async fn concurrent_handles_get_distinct_tokens_and_lifecycles() {
let (client, outbound, mut receiver) = make_client();
let begin_a = tokio::spawn({
let client = client.clone();
async move { client.begin_progress(ProgressOptions::new("A")).await }
});
let begin_b = tokio::spawn({
let client = client.clone();
async move { client.begin_progress(ProgressOptions::new("B")).await }
});
let params_a = answer_create_ok(&outbound, &mut receiver).await;
let params_b = answer_create_ok(&outbound, &mut receiver).await;
assert_eq!(params_a, json!({ "token": 1 }));
assert_eq!(params_b, json!({ "token": 2 }));
let handle_a = begin_a.await.unwrap().unwrap();
let handle_b = begin_b.await.unwrap().unwrap();
assert_eq!(handle_a.token(), ProgressToken::Number(1));
assert_eq!(handle_b.token(), ProgressToken::Number(2));
handle_b.report(None, Some(5)).unwrap();
handle_a.end(None).unwrap();
assert!(
!client
.progress_registry()
.is_active(&ProgressToken::Number(1))
);
assert!(
client
.progress_registry()
.is_active(&ProgressToken::Number(2))
);
handle_b.end(None).unwrap();
}
#[tokio::test]
async fn abandoning_the_begin_future_leaves_no_token_behind() {
let (client, outbound, mut receiver) = make_client();
let registry = client.progress_registry().clone();
let mut begin = Box::pin(client.begin_progress(ProgressOptions::new("Indexing")));
tokio::select! {
_ = &mut begin => panic!("begin cannot complete without a response"),
message = receiver.recv() => {
let message = message.expect("create request");
assert!(matches!(message, RawMessage::Request { .. }));
}
}
assert_eq!(outbound.pending_len(), 1);
drop(begin);
assert_eq!(outbound.pending_len(), 0);
let (method, params) = next_notification(&mut receiver).await;
assert_eq!(method, lsp_types::notification::Cancel::METHOD);
assert_eq!(params, json!({ "id": 1 }));
assert_eq!(
registry.active_len(),
0,
"abandoned begin registered no token"
);
assert!(
receiver.try_recv().is_err(),
"no begin notification followed"
);
}
}