use std::collections::HashMap;
use std::io::{BufReader, Read, Write};
use std::path::Path;
use std::process::{Child, ChildStdin, Command, Stdio};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::mpsc::{Receiver as StdReceiver, SyncSender as StdSyncSender, TrySendError};
use std::sync::{Mutex, MutexGuard};
use serde_json::Value;
use crate::error::{Error, Result};
use crate::tools::{ProcessCleanupMode, ProcessGuard};
const MAX_FRAME_BYTES: usize = 64 * 1024 * 1024;
const NOTIFICATION_QUEUE_CAP: usize = 1024;
const STDERR_TAIL_CAP: usize = 32 * 1024;
#[derive(Debug, Clone, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct RpcErrorObject {
pub code: i64,
pub message: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub data: Option<Value>,
}
#[derive(Debug, Clone, serde::Serialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum TransportError {
Server(RpcErrorObject),
Closed(String),
Io(String),
}
impl TransportError {
#[must_use]
pub const fn code(&self) -> &'static str {
match self {
Self::Server(_) => "LSP_SERVER_ERROR",
Self::Closed(_) => "LSP_TRANSPORT_CLOSED",
Self::Io(_) => "LSP_TRANSPORT_IO",
}
}
#[must_use]
pub fn mcp_code(&self) -> String {
self.code().replace("LSP_", "MCP_")
}
#[must_use]
pub fn message(&self) -> String {
match self {
Self::Server(err) => format!("server error {}: {}", err.code, err.message),
Self::Closed(reason) => format!("transport closed: {reason}"),
Self::Io(reason) => format!("transport I/O error: {reason}"),
}
}
}
#[derive(Debug, Clone)]
pub struct ServerNotification {
pub method: String,
pub params: Value,
}
#[derive(Debug, Clone)]
pub enum EnvPolicy {
InheritAndScrub(&'static [&'static str]),
Allowlist(&'static [&'static str]),
}
pub const MCP_ENV_ALLOWLIST: &[&str] = &[
"PATH",
"HOME",
"LANG",
"LC_ALL",
"LC_CTYPE",
"TMPDIR",
"TEMP",
"TMP",
"NO_COLOR",
"TERM",
"SystemRoot",
"SYSTEMROOT",
"APPDATA",
"LOCALAPPDATA",
"USERPROFILE",
"COMSPEC",
];
type SharedWriter = Mutex<ChildStdin>;
type PendingMap = Mutex<HashMap<u64, StdSyncSender<std::result::Result<Value, TransportError>>>>;
pub type ServerRequestHandler = std::sync::Arc<dyn Fn(&str, &Value) -> Option<Value> + Send + Sync>;
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
#[must_use]
pub fn encode_frame(body: &Value) -> Vec<u8> {
let json = serde_json::to_vec(body).unwrap_or_else(|_| b"null".to_vec());
let mut out = Vec::with_capacity(json.len() + 32);
out.extend_from_slice(format!("Content-Length: {}\r\n\r\n", json.len()).as_bytes());
out.extend_from_slice(&json);
out
}
pub(crate) fn read_frame(reader: &mut BufReader<impl Read>) -> std::io::Result<Option<Value>> {
let mut header: Vec<u8> = Vec::new();
let mut byte = [0u8; 1];
loop {
if reader.read(&mut byte)? == 0 {
if header.is_empty() {
return Ok(None); }
return Err(std::io::Error::new(
std::io::ErrorKind::UnexpectedEof,
"EOF mid-headers",
));
}
header.push(byte[0]);
if header.ends_with(b"\r\n\r\n") {
break;
}
if header.len() > 64 * 1024 {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"headers exceed 64 KiB",
));
}
}
let length = parse_content_length(&header)?;
let mut body = vec![0u8; length];
reader.read_exact(&mut body)?;
let value = serde_json::from_slice(&body).map_err(|err| {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("invalid JSON body: {err}"),
)
})?;
Ok(Some(value))
}
fn parse_content_length(header_bytes: &[u8]) -> std::io::Result<usize> {
let headers = String::from_utf8_lossy(header_bytes);
let mut content_length: Option<usize> = None;
for line in headers.split("\r\n") {
if let Some(value) = line
.split_once(':')
.map(|(k, v)| (k.trim(), v.trim()))
.filter(|(k, _)| k.eq_ignore_ascii_case("content-length"))
.map(|(_, v)| v)
{
content_length = value.parse::<usize>().ok();
}
}
let length = content_length.ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"missing Content-Length header",
)
})?;
if length > MAX_FRAME_BYTES {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("frame body {length} bytes exceeds cap {MAX_FRAME_BYTES}"),
));
}
Ok(length)
}
pub(crate) fn read_frame_with_scratch(
reader: &mut BufReader<impl Read>,
scratch: &mut Vec<u8>,
) -> std::io::Result<Option<Value>> {
let trace = std::env::var_os("PI_DAP_TRACE").is_some();
let mut chunk = [0u8; 8192];
let body_start = loop {
if let Some(pos) = find_subslice(scratch, b"\r\n\r\n") {
break pos + 4;
}
let read = reader.read(&mut chunk)?;
if read == 0 {
return Ok(None); }
scratch.extend_from_slice(&chunk[..read]);
if scratch.len() > 64 * 1024 {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"headers exceed 64 KiB",
));
}
};
if trace {
let headers = String::from_utf8_lossy(&scratch[..body_start]);
eprintln!("[dap-frame] headers: {headers:?}");
}
let length = parse_content_length(&scratch[..body_start])?;
let mut body: Vec<u8> = scratch.split_off(body_start);
while body.len() < length {
let read = reader.read(&mut chunk)?;
if read == 0 {
return Err(std::io::Error::new(
std::io::ErrorKind::UnexpectedEof,
"EOF mid-body",
));
}
body.extend_from_slice(&chunk[..read]);
}
let leftover = body.split_off(length);
*scratch = leftover;
if trace {
eprintln!(
"[dap-frame] body {} bytes: {:?}",
length,
String::from_utf8_lossy(&body[..length.min(120)])
);
}
let value = serde_json::from_slice(&body).map_err(|err| {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("invalid JSON body: {err}"),
)
})?;
Ok(Some(value))
}
fn find_subslice(hay: &[u8], needle: &[u8]) -> Option<usize> {
if needle.is_empty() || hay.len() < needle.len() {
return None;
}
hay.windows(needle.len()).position(|w| w == needle)
}
#[derive(Debug, Default)]
struct TailBuffer {
data: String,
cap: usize,
}
impl TailBuffer {
fn push(&mut self, chunk: &str) {
self.data.push_str(chunk);
if self.data.len() > self.cap {
let keep_from = self.data.len() - self.cap;
let boundary = self.data.ceil_char_boundary(keep_from);
self.data.drain(..boundary);
}
}
}
#[derive(Debug)]
pub(crate) struct PublicTailBuffer {
inner: TailBuffer,
}
impl PublicTailBuffer {
#[must_use]
pub(crate) const fn new() -> Self {
Self {
inner: TailBuffer {
data: String::new(),
cap: 32 * 1024,
},
}
}
pub(crate) fn push(&mut self, chunk: &str) {
self.inner.push(chunk);
}
#[must_use]
pub(crate) fn tail(&self) -> String {
self.inner.data.clone()
}
}
pub async fn await_completion<T>(
rx: StdReceiver<T>,
timeout: std::time::Duration,
on_abandon: impl FnOnce(),
) -> std::result::Result<T, CompletionWaitError> {
let cx = crate::agent_cx::AgentCx::for_current_or_request();
let start = cx
.cx()
.timer_driver()
.map_or_else(asupersync::time::wall_now, |timer| timer.now());
loop {
match rx.try_recv() {
Ok(value) => return Ok(value),
Err(std::sync::mpsc::TryRecvError::Empty) => {}
Err(std::sync::mpsc::TryRecvError::Disconnected) => {
return Err(CompletionWaitError::Closed);
}
}
let now = cx
.cx()
.timer_driver()
.map_or_else(asupersync::time::wall_now, |timer| timer.now());
if std::time::Duration::from_nanos(now.duration_since(start)) >= timeout {
on_abandon();
return Err(CompletionWaitError::Timeout);
}
if cx.checkpoint().is_err() {
on_abandon();
return Err(CompletionWaitError::Cancelled);
}
asupersync::time::sleep(now, std::time::Duration::from_millis(10)).await;
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CompletionWaitError {
Timeout,
Cancelled,
Closed,
}
fn apply_env_policy(cmd: &mut Command, policy: &EnvPolicy, env: &[(String, String)]) {
match policy {
EnvPolicy::InheritAndScrub(scrub) => {
for var in *scrub {
cmd.env_remove(var);
}
}
EnvPolicy::Allowlist(allowlist) => {
cmd.env_clear();
for var in *allowlist {
if let Some(value) = std::env::var_os(var) {
cmd.env(var, value);
}
}
}
}
cmd.envs(env.iter().map(|(k, v)| (k.as_str(), v.as_str())));
}
#[allow(clippy::too_many_arguments)]
fn reader_loop(
stdout: std::process::ChildStdout,
pending: std::sync::Arc<PendingMap>,
alive: std::sync::Arc<std::sync::atomic::AtomicBool>,
writer: std::sync::Arc<SharedWriter>,
stderr_tail: std::sync::Arc<Mutex<TailBuffer>>,
notification_tx: StdSyncSender<ServerNotification>,
dropped: std::sync::Arc<AtomicU64>,
handler: std::sync::Arc<Mutex<Option<ServerRequestHandler>>>,
) -> impl FnOnce() + Send + 'static {
move || {
let mut reader = BufReader::new(stdout);
let mut scratch = Vec::new();
let close_reason = loop {
match read_frame_with_scratch(&mut reader, &mut scratch) {
Ok(Some(message)) => {
handle_message(
&message,
&pending,
&writer,
¬ification_tx,
&dropped,
&stderr_tail,
&handler,
);
}
Ok(None) => break "server closed stdout (EOF)".to_string(),
Err(err) => break format!("frame read error: {err}"),
}
};
alive.store(false, Ordering::SeqCst);
let mut pending = lock(&pending);
for (_, sender) in pending.drain() {
let _ = sender.send(Err(TransportError::Closed(close_reason.clone())));
}
}
}
pub struct JsonRpcClient {
child: Mutex<ProcessGuard>,
writer: std::sync::Arc<SharedWriter>,
pending: std::sync::Arc<PendingMap>,
next_id: AtomicU64,
alive: std::sync::Arc<std::sync::atomic::AtomicBool>,
stderr_tail: std::sync::Arc<Mutex<TailBuffer>>,
notification_rx: Mutex<StdReceiver<ServerNotification>>,
dropped_notifications: std::sync::Arc<AtomicU64>,
server_request_handler: std::sync::Arc<Mutex<Option<ServerRequestHandler>>>,
}
impl JsonRpcClient {
pub fn spawn(
command: &str,
args: &[String],
env: &[(String, String)],
cwd: &Path,
) -> Result<Self> {
Self::spawn_with_policy(
command,
args,
env,
cwd,
&EnvPolicy::InheritAndScrub(&["CARGO_TARGET_DIR"]),
"lsp",
)
}
pub fn spawn_with_policy(
command: &str,
args: &[String],
env: &[(String, String)],
cwd: &Path,
policy: &EnvPolicy,
flavor: &'static str,
) -> Result<Self> {
let mut cmd = Command::new(command);
cmd.args(args)
.current_dir(cwd)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
apply_env_policy(&mut cmd, policy, env);
let mut child: Child = cmd.spawn().map_err(|err| {
Error::tool(
flavor,
format!(
"[{}_SERVER_MISSING] failed to spawn {command:?}: {err}",
flavor.to_ascii_uppercase()
),
)
})?;
let stdin = child
.stdin
.take()
.ok_or_else(|| Error::tool("lsp", "missing child stdin".to_string()))?;
let stdout = child
.stdout
.take()
.ok_or_else(|| Error::tool("lsp", "missing child stdout".to_string()))?;
let stderr = child
.stderr
.take()
.ok_or_else(|| Error::tool("lsp", "missing child stderr".to_string()))?;
let writer = std::sync::Arc::new(Mutex::new(stdin));
let pending: std::sync::Arc<PendingMap> = std::sync::Arc::new(Mutex::new(HashMap::new()));
let alive = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(true));
let stderr_tail = std::sync::Arc::new(Mutex::new(TailBuffer {
data: String::new(),
cap: STDERR_TAIL_CAP,
}));
let (notification_tx, notification_rx) =
std::sync::mpsc::sync_channel::<ServerNotification>(NOTIFICATION_QUEUE_CAP);
let dropped_notifications = std::sync::Arc::new(AtomicU64::new(0));
let server_request_handler: std::sync::Arc<Mutex<Option<ServerRequestHandler>>> =
std::sync::Arc::new(Mutex::new(None));
{
let reader_body = reader_loop(
stdout,
std::sync::Arc::clone(&pending),
std::sync::Arc::clone(&alive),
std::sync::Arc::clone(&writer),
std::sync::Arc::clone(&stderr_tail),
notification_tx,
std::sync::Arc::clone(&dropped_notifications),
std::sync::Arc::clone(&server_request_handler),
);
std::thread::spawn(reader_body); }
{
let stderr_tail = std::sync::Arc::clone(&stderr_tail);
let pump_body = move || {
let mut reader = BufReader::new(stderr);
let mut buf = [0u8; 4096];
loop {
match reader.read(&mut buf) {
Ok(0) | Err(_) => break,
Ok(n) => {
let chunk = String::from_utf8_lossy(&buf[..n]);
lock(&stderr_tail).push(&chunk);
}
}
}
};
std::thread::spawn(pump_body); }
Ok(Self {
child: Mutex::new(ProcessGuard::new(child, ProcessCleanupMode::ChildOnly)),
writer,
pending,
next_id: AtomicU64::new(1),
alive,
stderr_tail,
notification_rx: Mutex::new(notification_rx),
dropped_notifications,
server_request_handler,
})
}
pub fn set_server_request_handler(&self, handler: ServerRequestHandler) {
*lock(&self.server_request_handler) = Some(handler);
}
#[must_use]
pub fn is_alive(&self) -> bool {
self.alive.load(Ordering::SeqCst)
}
pub fn request(
&self,
method: &str,
params: Value,
) -> std::result::Result<
(u64, StdReceiver<std::result::Result<Value, TransportError>>),
TransportError,
> {
if !self.is_alive() {
return Err(TransportError::Closed(
"server transport is not alive".to_string(),
));
}
let id = self.next_id.fetch_add(1, Ordering::SeqCst);
let (tx, rx) = std::sync::mpsc::sync_channel(1);
lock(&self.pending).insert(id, tx);
let mut frame_value = serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"method": method,
});
frame_value["params"] = params;
let frame = encode_frame(&frame_value);
let write_result = {
let mut guard = lock(&self.writer);
guard.write_all(&frame).and_then(|()| guard.flush())
};
if let Err(err) = write_result {
lock(&self.pending).remove(&id);
return Err(TransportError::Io(format!("request write failed: {err}")));
}
Ok((id, rx))
}
pub fn notify(&self, method: &str, params: Value) -> std::result::Result<(), TransportError> {
let mut frame_value = serde_json::json!({
"jsonrpc": "2.0",
"method": method,
});
frame_value["params"] = params;
let frame = encode_frame(&frame_value);
let mut guard = lock(&self.writer);
guard
.write_all(&frame)
.and_then(|()| guard.flush())
.map_err(|err| TransportError::Io(format!("notification write failed: {err}")))
}
pub fn cancel_request(&self, id: u64) {
lock(&self.pending).remove(&id);
let _ = self.notify("$/cancelRequest", serde_json::json!({ "id": id }));
}
pub fn drain_notifications(&self) -> Vec<ServerNotification> {
let mut out = Vec::new();
let rx = lock(&self.notification_rx);
while let Ok(notification) = rx.try_recv() {
out.push(notification);
}
out
}
#[must_use]
pub fn dropped_notifications(&self) -> u64 {
self.dropped_notifications.load(Ordering::SeqCst)
}
#[must_use]
pub fn stderr_tail(&self) -> String {
lock(&self.stderr_tail).data.clone()
}
#[must_use]
pub fn child_exited(&self) -> bool {
let status = lock(&self.child).try_wait_child();
matches!(status, Ok(Some(_)))
}
pub fn shutdown(&self) {
let _ = self.notify("exit", Value::Null);
self.kill();
}
pub fn kill(&self) {
let _ = lock(&self.child).kill();
self.alive.store(false, Ordering::SeqCst);
}
}
impl Drop for JsonRpcClient {
fn drop(&mut self) {
self.kill();
}
}
#[allow(clippy::option_if_let_else)] fn handle_message<W: Write>(
message: &Value,
pending: &PendingMap,
writer: &Mutex<W>,
notification_tx: &StdSyncSender<ServerNotification>,
dropped: &AtomicU64,
stderr_tail: &Mutex<TailBuffer>,
handler: &Mutex<Option<ServerRequestHandler>>,
) {
let id = message.get("id").cloned();
let method = message.get("method").and_then(Value::as_str);
let is_response = method.is_none() && id.is_some();
if is_response {
let Some(numeric_id) = id.and_then(|v| v.as_u64()) else {
return;
};
let sender = lock(pending).remove(&numeric_id);
if let Some(sender) = sender {
let outcome = if let Some(error) = message.get("error") {
Err(TransportError::Server(RpcErrorObject {
code: error.get("code").and_then(Value::as_i64).unwrap_or(-1),
message: error
.get("message")
.and_then(Value::as_str)
.unwrap_or("unknown server error")
.to_string(),
data: error.get("data").cloned(),
}))
} else {
Ok(message.get("result").cloned().unwrap_or(Value::Null))
};
let _ = sender.send(outcome);
}
return;
}
if let (Some(id), Some(method)) = (id, method) {
let custom = lock(handler).as_ref().and_then(|handler| {
handler(
method,
&message.get("params").cloned().unwrap_or(Value::Null),
)
});
let frame = encode_frame(&serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"result": custom.unwrap_or(Value::Null),
}));
{
let mut guard = lock(writer);
let _ = guard.write_all(&frame).and_then(|()| guard.flush());
}
return;
}
if let Some(method) = method {
if method == "window/logMessage" {
let rendered = message
.get("params")
.and_then(|p| p.get("message"))
.and_then(Value::as_str)
.unwrap_or("");
lock(stderr_tail).push(&format!("[logMessage] {rendered}\n"));
}
let notification = ServerNotification {
method: method.to_string(),
params: message.get("params").cloned().unwrap_or(Value::Null),
};
match notification_tx.try_send(notification) {
Ok(()) | Err(TrySendError::Disconnected(_)) => {}
Err(TrySendError::Full(_)) => {
dropped.fetch_add(1, Ordering::SeqCst);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn encode_frame_uses_content_length() {
let body = serde_json::json!({"jsonrpc":"2.0","id":1,"method":"ping"});
let frame = encode_frame(&body);
let text = String::from_utf8(frame).expect("utf8");
let json = serde_json::to_string(&body).expect("json");
assert!(text.starts_with(&format!("Content-Length: {}\r\n\r\n", json.len())));
assert!(text.ends_with(&json));
}
#[test]
fn read_frame_roundtrips() {
let body = serde_json::json!({"jsonrpc":"2.0","id":7,"result":{"ok":true}});
let frame = encode_frame(&body);
let mut reader = BufReader::new(frame.as_slice());
let got = read_frame(&mut reader).expect("read").expect("some");
assert_eq!(got, body);
}
#[test]
fn read_frame_handles_back_to_back_frames() {
let a = serde_json::json!({"a": 1});
let b = serde_json::json!({"b": 2});
let mut bytes = encode_frame(&a);
bytes.extend_from_slice(&encode_frame(&b));
let mut reader = BufReader::new(bytes.as_slice());
assert_eq!(read_frame(&mut reader).expect("read"), Some(a));
assert_eq!(read_frame(&mut reader).expect("read"), Some(b));
assert_eq!(read_frame(&mut reader).expect("read"), None);
}
#[test]
fn read_frame_rejects_missing_length() {
let bytes = b"Content-Type: application/vscode-jsonrpc; charset=utf-8\r\n\r\n{}";
let mut reader = BufReader::new(bytes.as_slice());
assert!(read_frame(&mut reader).is_err());
}
#[test]
fn read_frame_rejects_oversize_length() {
let bytes = format!("Content-Length: {}\r\n\r\n", MAX_FRAME_BYTES + 1);
let mut reader = BufReader::new(bytes.as_bytes());
assert!(read_frame(&mut reader).is_err());
}
#[test]
fn tail_buffer_truncates_to_cap() {
let mut tail = TailBuffer {
data: String::new(),
cap: 16,
};
tail.push("0123456789");
tail.push("abcdefghij");
assert!(tail.data.len() <= 16);
assert!(tail.data.ends_with("abcdefghij"));
}
#[test]
fn handle_message_completes_pending() {
let pending: PendingMap = Mutex::new(HashMap::new());
let (tx, rx) = std::sync::mpsc::sync_channel(1);
lock(&pending).insert(42, tx);
let (notification_tx, _notification_rx) = std::sync::mpsc::sync_channel(8);
let dropped = AtomicU64::new(0);
let stderr_tail = Mutex::new(TailBuffer::default());
let mut cmd = Command::new("cat");
cmd.stdin(Stdio::piped());
let mut child = cmd.spawn().expect("cat");
let writer = Mutex::new(child.stdin.take().expect("stdin"));
handle_message(
&serde_json::json!({"jsonrpc":"2.0","id":42,"result":{"value":1}}),
&pending,
&writer,
¬ification_tx,
&dropped,
&stderr_tail,
&Mutex::new(None),
);
let got = rx.try_recv().expect("completed");
assert_eq!(
got.ok().map(|v| v["value"].clone()),
Some(serde_json::json!(1))
);
let _ = child.kill();
let _ = child.wait();
}
#[test]
fn handle_message_answers_server_requests_with_null() {
let pending: PendingMap = Mutex::new(HashMap::new());
let (notification_tx, _rx) = std::sync::mpsc::sync_channel(8);
let dropped = AtomicU64::new(0);
let stderr_tail = Mutex::new(TailBuffer::default());
let (reader, writer_end) = std::io::pipe().expect("pipe");
let writer = Mutex::new(writer_end);
handle_message(
&serde_json::json!({"jsonrpc":"2.0","id":9,"method":"workspace/configuration","params":{}}),
&pending,
&writer,
¬ification_tx,
&dropped,
&stderr_tail,
&Mutex::new(None),
);
let frame = read_frame(&mut BufReader::new(reader))
.expect("read")
.expect("some");
assert_eq!(frame["id"], 9);
assert!(frame.get("result").is_some());
}
#[test]
fn request_and_notify_writes_do_not_deadlock() {
let temp = tempfile::tempdir().expect("tempdir");
let client = JsonRpcClient::spawn("cat", &[], &[], temp.path()).expect("spawn cat");
let (id, _rx) = client
.request("initialize", serde_json::json!({}))
.expect("request write must succeed");
client
.notify("initialized", serde_json::json!({}))
.expect("notify write must succeed");
client.cancel_request(id);
client.kill();
assert!(!client.is_alive());
}
#[test]
fn handle_message_routes_notifications_and_bounds_queue() {
let pending: PendingMap = Mutex::new(HashMap::new());
let (notification_tx, notification_rx) = std::sync::mpsc::sync_channel(2);
let dropped = AtomicU64::new(0);
let stderr_tail = Mutex::new(TailBuffer::default());
let (_reader, writer_end) = std::io::pipe().expect("pipe");
let writer = Mutex::new(writer_end);
for i in 0..4 {
handle_message(
&serde_json::json!({"jsonrpc":"2.0","method":"textDocument/publishDiagnostics","params":{"i":i}}),
&pending,
&writer,
¬ification_tx,
&dropped,
&stderr_tail,
&Mutex::new(None),
);
}
assert_eq!(dropped.load(Ordering::SeqCst), 2);
assert!(notification_rx.try_recv().is_ok());
assert!(notification_rx.try_recv().is_ok());
assert!(notification_rx.try_recv().is_err());
}
}