返回 CodeWhale
client.rs
根目录 / crates / tui / src / lsp / client.rs
1 //! Thin JSON-RPC over stdio client for LSP servers.
2 //!
3 //! We deliberately do **not** depend on `tower-lsp` — it is a server-side
4 //! framework and dragging it in here would add hundreds of unnecessary
5 //! transitive dependencies and slow down `cargo build` for every contributor.
6 //! The LSP wire protocol is small enough that handling it ourselves is a
7 //! self-contained ~400 LOC and lets us keep total control of the spawn
8 //! lifecycle, timeouts, and the async surface.
9 //!
10 //! Architecture:
11 //!
12 //! - [`LspTransport`] is the trait the [`super::LspManager`] talks to. The
13 //! real implementation is [`StdioLspTransport`] (forks an LSP server with
14 //! `tokio::process::Command`); tests use `super::tests::FakeTransport`.
15 //! - [`StdioLspTransport`] runs three tokio tasks: a reader, a writer, and
16 //! the public API. Communication uses tokio mpsc channels.
17 //! - We parse `Content-Length`-framed JSON-RPC and route inbound messages
18 //! either to a per-request response slot (for replies) or to the
19 //! diagnostics queue (for `textDocument/publishDiagnostics` notifications).
20 //!
21 //! The transport is one-shot per file in MVP form: the manager spawns a
22 //! transport on demand for a language and reuses it. We do not implement
23 //! workspace sync beyond didOpen/didChange because the goal is "post-edit
24 //! diagnostics," not full IDE smartness.
25
26 use std::collections::HashMap;
27 use std::path::{Path, PathBuf};
28 use std::process::Stdio;
29 use std::sync::Arc;
30 use std::time::Duration;
31
32 use anyhow::{Context, Result, anyhow};
33 use async_trait::async_trait;
34 use serde_json::{Value, json};
35 use tokio::io::{AsyncRead, AsyncReadExt, AsyncWriteExt};
36 use tokio::process::{Child, Command};
37 use tokio::sync::Mutex as AsyncMutex;
38 use tokio::sync::{mpsc, oneshot};
39 use tokio::time::timeout;
40
41 use super::diagnostics::{Diagnostic, Severity};
42 use crate::utils::spawn_supervised;
43
44 const INITIALIZE_TIMEOUT: Duration = Duration::from_secs(5);
45 const MAX_LSP_HEADER_BYTES: usize = 8 * 1024;
46 const MAX_LSP_FRAME_BYTES: usize = 16 * 1024 * 1024;
47
48 /// A publication retains the server's document version instead of pretending
49 /// that a matching URI alone proves which text was checked.
50 #[derive(Debug)]
51 pub struct DiagnosticPublication {
52 pub items: Vec<Diagnostic>,
53 pub document_version: Option<i64>,
54 pub diagnostic_version: Option<i64>,
55 }
56
57 impl DiagnosticPublication {
58 #[must_use]
59 pub fn freshness(&self) -> &'static str {
60 if self.document_version.is_some() && self.document_version == self.diagnostic_version {
61 "verified"
62 } else {
63 "unverified"
64 }
65 }
66 }
67
68 // Diagnostic-only transports have no document-version proof. Existing callers
69 // may still use their results, but must not claim freshness from an empty list.
70 impl From<Vec<Diagnostic>> for DiagnosticPublication {
71 fn from(items: Vec<Diagnostic>) -> Self {
72 Self {
73 items,
74 document_version: None,
75 diagnostic_version: None,
76 }
77 }
78 }
79
80 /// Source synchronization proof for a semantic reply. Legacy transports keep
81 /// the raw result but cannot assert which document version served it.
82 #[derive(Debug)]
83 pub struct SemanticReply {
84 pub result: Value,
85 pub document_version: Option<i64>,
86 }
87
88 /// Trait the LSP manager talks to. A real LSP server speaks this via stdio;
89 /// tests use an in-process fake.
90 #[async_trait]
91 pub trait LspTransport: Send + Sync {
92 /// Notify the server that a file was opened or its contents updated, then
93 /// wait up to `wait` for a `publishDiagnostics` notification for that
94 /// file. Returns the diagnostics list (possibly empty). Implementations
95 /// must NOT block past `wait`.
96 async fn diagnostics_for(
97 &self,
98 path: &Path,
99 text: &str,
100 wait: Duration,
101 ) -> Result<DiagnosticPublication>;
102
103 /// Send a JSON-RPC request and wait up to `wait` for the reply.
104 ///
105 /// Default returns "unsupported" so diagnostic-only fakes keep working.
106 /// Real transports implement this for go-to-definition, symbols, and
107 /// references without spawning a second server lifecycle.
108 async fn request(&self, _method: &str, _params: Value, _wait: Duration) -> Result<Value> {
109 Err(anyhow!("LSP request not supported by this transport"))
110 }
111
112 /// Synchronize and query one document atomically when the transport can
113 /// prove that ordering. Diagnostic-only/legacy transports stay unverified.
114 async fn request_for_document(
115 &self,
116 path: &Path,
117 text: &str,
118 method: &str,
119 params: Value,
120 wait: Duration,
121 ) -> Result<SemanticReply> {
122 timeout(wait, async {
123 self.ensure_open(path, text).await?;
124 Ok(SemanticReply {
125 result: self.request(method, params, wait).await?,
126 document_version: None,
127 })
128 })
129 .await
130 .map_err(|_| anyhow!("LSP semantic request timed out"))?
131 }
132
133 /// Ensure `path` is open with `text` (didOpen/didChange) so position-based
134 /// requests can target it. Default is a no-op; real transports track opens.
135 async fn ensure_open(&self, _path: &Path, _text: &str) -> Result<()> {
136 Ok(())
137 }
138
139 /// A closed transport is never valid cache evidence. Diagnostic-only
140 /// in-process implementations remain usable until their owner removes them.
141 fn is_alive(&self) -> bool {
142 true
143 }
144
145 /// Best-effort shutdown. Called via `LspManager::shutdown_all`.
146 async fn shutdown(&self);
147 }
148
149 type DiagnosticMessage = (PathBuf, Option<i64>, Vec<Diagnostic>);
150
151 /// Stdio-backed transport. Spawns the LSP server as a child process and
152 /// pipes JSON-RPC over stdin/stdout. Stderr is drained without retaining or
153 /// exposing arbitrary server output.
154 pub struct StdioLspTransport {
155 /// JoinHandle for the running server. Held so the child stays alive for
156 /// the transport's lifetime; consumed during `shutdown`.
157 child: Arc<AsyncMutex<Option<Child>>>,
158 tasks: Vec<tokio::task::JoinHandle<()>>,
159 /// Outgoing message sender to the writer task.
160 tx_outbound: mpsc::Sender<Vec<u8>>,
161 /// Inbound diagnostics queue. We push every `publishDiagnostics`
162 /// notification into here and the public API drains the relevant entries.
163 diagnostics_gate: AsyncMutex<()>,
164 diagnostics_rx: AsyncMutex<mpsc::Receiver<DiagnosticMessage>>,
165 /// Map of in-flight request id -> reply slot for model-facing intelligence
166 /// requests (definition, references, symbols).
167 pending: Arc<AsyncMutex<HashMap<i64, oneshot::Sender<Value>>>>,
168 /// Monotonic request id counter for JSON-RPC request/reply methods.
169 next_id: AsyncMutex<i64>,
170 /// Language id passed in `textDocument/didOpen` (e.g. "rust").
171 language_id: String,
172 /// Track which files we have opened so the second touch sends
173 /// `didChange` instead of `didOpen`.
174 opened: AsyncMutex<HashMap<PathBuf, i64>>,
175 }
176
177 impl StdioLspTransport {
178 /// Spawn `command args…` and run the LSP `initialize` handshake. Returns
179 /// `Err` immediately if the binary is not on PATH or `initialize` fails.
180 pub async fn spawn(
181 command: &str,
182 args: &[String],
183 language_id: &str,
184 workspace: PathBuf,
185 ) -> Result<Self> {
186 Self::spawn_with_timeout(command, args, language_id, workspace, INITIALIZE_TIMEOUT).await
187 }
188
189 async fn spawn_with_timeout(
190 command: &str,
191 args: &[String],
192 language_id: &str,
193 workspace: PathBuf,
194 initialize_wait: Duration,
195 ) -> Result<Self> {
196 let mut cmd = Command::new(command);
197 cmd.args(args);
198 // Language servers execute workspace code (rust-analyzer runs build
199 // scripts and proc-macros, tsserver loads tsconfig plugins), so they
200 // start from the sanitized child environment like `exec_shell`.
201 crate::child_env::apply_to_tokio_command(&mut cmd, std::iter::empty::<(&str, &str)>());
202 cmd.stdin(Stdio::piped());
203 cmd.stdout(Stdio::piped());
204 cmd.stderr(Stdio::piped());
205 cmd.kill_on_drop(true);
206
207 let mut child = cmd
208 .spawn()
209 .with_context(|| format!("failed to spawn LSP server `{command}`"))?;
210
211 let stdin = child
212 .stdin
213 .take()
214 .context("LSP child has no stdin handle")?;
215 let stdout = child
216 .stdout
217 .take()
218 .context("LSP child has no stdout handle")?;
219
220 let mut stderr = child
221 .stderr
222 .take()
223 .context("LSP child has no stderr handle")?;
224 let stderr_task =
225 spawn_supervised("lsp-stderr", std::panic::Location::caller(), async move {
226 // Drain bytes, not lines: even a single unbounded log line must
227 // neither block the child nor accumulate in host memory.
228 let _ = tokio::io::copy(&mut stderr, &mut tokio::io::sink()).await;
229 });
230
231 let (tx_outbound, rx_outbound) = mpsc::channel::<Vec<u8>>(64);
232 let (tx_inbound, rx_inbound) = mpsc::channel::<Value>(64);
233 let (tx_diag, rx_diag) = mpsc::channel::<DiagnosticMessage>(64);
234
235 // Writer task: drain outbound channel, frame with Content-Length, write to stdin.
236 let writer_task = spawn_supervised(
237 "lsp-writer",
238 std::panic::Location::caller(),
239 writer_task(stdin, rx_outbound),
240 );
241 // Reader task: parse Content-Length frames from stdout, push to inbound queue.
242 let reader_task = spawn_supervised(
243 "lsp-reader",
244 std::panic::Location::caller(),
245 reader_task(stdout, tx_inbound),
246 );
247 // Inbound dispatcher: routes notifications to `tx_diag`, replies to a
248 // pending map. We keep the pending map for completeness even though
249 // diagnostics polling itself does not reuse it.
250 let pending: Arc<AsyncMutex<HashMap<i64, oneshot::Sender<Value>>>> =
251 Arc::new(AsyncMutex::new(HashMap::new()));
252 let child = Arc::new(AsyncMutex::new(Some(child)));
253 let dispatcher_child = child.clone();
254 let dispatcher_pending = pending.clone();
255 let dispatcher_task = spawn_supervised(
256 "lsp-dispatcher",
257 std::panic::Location::caller(),
258 async move {
259 dispatcher_task(rx_inbound, tx_diag, dispatcher_pending).await;
260 // EOF, malformed frames, or an overflowing diagnostics queue
261 // terminate the producer too, so its pipes cannot stay stuck.
262 if let Some(mut child) = dispatcher_child.lock().await.take() {
263 let _ = child.start_kill();
264 let _ = child.wait().await;
265 }
266 },
267 );
268
269 let transport = Self {
270 child,
271 tasks: vec![stderr_task, writer_task, reader_task, dispatcher_task],
272 tx_outbound,
273 diagnostics_gate: AsyncMutex::new(()),
274 diagnostics_rx: AsyncMutex::new(rx_diag),
275 pending,
276 next_id: AsyncMutex::new(1),
277 language_id: language_id.to_string(),
278 opened: AsyncMutex::new(HashMap::new()),
279 };
280 let result = transport.request("initialize", json!({
281 "processId": std::process::id(),
282 "rootUri": uri_from_path(&workspace),
283 "capabilities": {
284 "general": { "positionEncodings": ["utf-16"] },
285 "textDocument": {
286 "publishDiagnostics": { "relatedInformation": false, "versionSupport": true }
287 }
288 },
289 "workspaceFolders": [{"uri": uri_from_path(&workspace), "name": "workspace"}]
290 }), initialize_wait).await.context("LSP initialization failed")?;
291 if !result.get("capabilities").is_some_and(Value::is_object) {
292 return Err(anyhow!(
293 "LSP initialize response is missing server capabilities"
294 ));
295 }
296 if result
297 .pointer("/capabilities/positionEncoding")
298 .is_some_and(|encoding| encoding.as_str() != Some("utf-16"))
299 {
300 return Err(anyhow!("LSP server must use UTF-16 positions"));
301 }
302 timeout(
303 initialize_wait,
304 send_message(
305 &transport.tx_outbound,
306 &json!({
307 "jsonrpc": "2.0", "method": "initialized", "params": {}
308 }),
309 ),
310 )
311 .await
312 .map_err(|_| anyhow!("LSP initialized notification timed out"))??;
313 Ok(transport)
314 }
315 }
316
317 impl Drop for StdioLspTransport {
318 fn drop(&mut self) {
319 for task in &self.tasks {
320 task.abort();
321 }
322 if let Ok(mut child) = self.child.try_lock()
323 && let Some(child) = child.as_mut()
324 {
325 let _ = child.start_kill();
326 }
327 // If shutdown/dispatcher currently owns the child, abort releases
328 // its local Child and kill_on_drop remains the final fallback.
329 }
330 }
331
332 impl StdioLspTransport {
333 async fn open_or_change(&self, path: &Path, text: &str) -> Result<(String, i64)> {
334 let path_buf = path.to_path_buf();
335 let uri = uri_from_path(&path_buf);
336 let mut opened = self.opened.lock().await;
337 let is_new = !opened.contains_key(&path_buf);
338 let new_version = opened.get(&path_buf).copied().unwrap_or(0) + 1;
339
340 let payload = if is_new {
341 json!({
342 "jsonrpc": "2.0",
343 "method": "textDocument/didOpen",
344 "params": {
345 "textDocument": {
346 "uri": uri.clone(),
347 "languageId": self.language_id,
348 "version": new_version,
349 "text": text
350 }
351 }
352 })
353 } else {
354 json!({
355 "jsonrpc": "2.0",
356 "method": "textDocument/didChange",
357 "params": {
358 "textDocument": {
359 "uri": uri.clone(),
360 "version": new_version
361 },
362 "contentChanges": [{ "text": text }]
363 }
364 })
365 };
366 send_message(&self.tx_outbound, &payload).await?;
367 opened.insert(path_buf, new_version);
368 Ok((uri, new_version))
369 }
370 }
371
372 #[async_trait]
373 impl LspTransport for StdioLspTransport {
374 fn is_alive(&self) -> bool {
375 // stderr may close independently. The writer, reader and dispatcher
376 // are the protocol lifetime; none may have exited or been aborted.
377 !self.tx_outbound.is_closed() && self.tasks.iter().skip(1).all(|task| !task.is_finished())
378 }
379
380 async fn diagnostics_for(
381 &self,
382 path: &Path,
383 text: &str,
384 wait: Duration,
385 ) -> Result<DiagnosticPublication> {
386 // One receiver cannot serve concurrent polling safely: serialize the
387 // open/version/send/wait transaction, including semantic ensure_open.
388 let deadline = tokio::time::Instant::now() + wait;
389 let _gate = timeout(wait, self.diagnostics_gate.lock())
390 .await
391 .map_err(|_| anyhow!("LSP diagnostics timed out waiting for another document"))?;
392 let path_buf = path.to_path_buf();
393 let (_, version) = timeout(
394 deadline.saturating_duration_since(tokio::time::Instant::now()),
395 self.open_or_change(path, text),
396 )
397 .await
398 .map_err(|_| anyhow!("LSP diagnostics timed out sending document"))??;
399 loop {
400 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
401 if remaining.is_zero() {
402 return Err(anyhow!(
403 "LSP diagnostics timed out before a current publication"
404 ));
405 }
406 let mut rx = self.diagnostics_rx.lock().await;
407 let (file, published_version, items) = match timeout(remaining, rx.recv()).await {
408 Ok(Some(item)) => item,
409 Ok(None) => {
410 return Err(anyhow!(
411 "LSP diagnostics channel closed before publishDiagnostics"
412 ));
413 }
414 Err(_) => {
415 return Err(anyhow!(
416 "LSP diagnostics timed out before a current publication"
417 ));
418 }
419 };
420 if file != path_buf || published_version.is_some_and(|published| published != version) {
421 continue;
422 }
423 return Ok(DiagnosticPublication {
424 items,
425 document_version: Some(version),
426 diagnostic_version: published_version,
427 });
428 }
429 }
430
431 async fn request_for_document(
432 &self,
433 path: &Path,
434 text: &str,
435 method: &str,
436 params: Value,
437 wait: Duration,
438 ) -> Result<SemanticReply> {
439 let deadline = tokio::time::Instant::now() + wait;
440 let _gate = timeout(wait, self.diagnostics_gate.lock())
441 .await
442 .map_err(|_| anyhow!("LSP semantic request timed out waiting for another document"))?;
443 let (_, version) = timeout(
444 deadline.saturating_duration_since(tokio::time::Instant::now()),
445 self.open_or_change(path, text),
446 )
447 .await
448 .map_err(|_| anyhow!("LSP semantic request timed out sending document"))??;
449 let result = self
450 .request(
451 method,
452 params,
453 deadline.saturating_duration_since(tokio::time::Instant::now()),
454 )
455 .await?;
456 Ok(SemanticReply {
457 result,
458 document_version: Some(version),
459 })
460 }
461
462 async fn ensure_open(&self, path: &Path, text: &str) -> Result<()> {
463 let _gate = self.diagnostics_gate.lock().await;
464 self.open_or_change(path, text).await?;
465 Ok(())
466 }
467
468 async fn request(&self, method: &str, params: Value, wait: Duration) -> Result<Value> {
469 let id = {
470 let mut next = self.next_id.lock().await;
471 let id = *next;
472 *next = next.saturating_add(1);
473 id
474 };
475 let (tx, rx) = oneshot::channel();
476 {
477 let mut pending = self.pending.lock().await;
478 pending.insert(id, tx);
479 }
480 let payload = json!({
481 "jsonrpc": "2.0",
482 "id": id,
483 "method": method,
484 "params": params,
485 });
486 // The deadline includes queue backpressure, not just the reply.
487 let response = timeout(wait, async {
488 send_message(&self.tx_outbound, &payload).await?;
489 rx.await.map_err(|_| anyhow!("LSP request channel closed"))
490 })
491 .await;
492 self.pending.lock().await.remove(&id);
493 match response {
494 Ok(Ok(reply)) => {
495 if let Some(error) = reply.get("error") {
496 let message = error
497 .get("message")
498 .and_then(Value::as_str)
499 .unwrap_or("LSP request failed");
500 return Err(anyhow!("{message}"));
501 }
502 reply
503 .get("result")
504 .cloned()
505 .ok_or_else(|| anyhow!("LSP response has no result"))
506 }
507 Ok(Err(error)) => Err(error),
508 Err(_) => Err(anyhow!("LSP request timed out for {method}")),
509 }
510 }
511
512 async fn shutdown(&self) {
513 let mut child = self.child.lock().await;
514 if let Some(mut c) = child.take() {
515 let _ = c.start_kill();
516 let _ = c.wait().await;
517 }
518 for task in &self.tasks {
519 task.abort();
520 }
521 self.pending.lock().await.clear();
522 }
523 }
524
525 /// Send a JSON value as one Content-Length-framed JSON-RPC message.
526 async fn send_message(tx: &mpsc::Sender<Vec<u8>>, value: &Value) -> Result<()> {
527 let body = serde_json::to_vec(value).context("serialize LSP message")?;
528 let header = format!("Content-Length: {}\r\n\r\n", body.len());
529 let mut frame = Vec::with_capacity(header.len() + body.len());
530 frame.extend_from_slice(header.as_bytes());
531 frame.extend_from_slice(&body);
532 tx.send(frame)
533 .await
534 .map_err(|_| anyhow!("LSP outbound channel closed"))?;
535 Ok(())
536 }
537
538 /// Background task that drains the outbound queue and writes each frame to
539 /// the LSP server's stdin. Exits cleanly when the channel closes.
540 async fn writer_task(mut stdin: tokio::process::ChildStdin, mut rx: mpsc::Receiver<Vec<u8>>) {
541 while let Some(frame) = rx.recv().await {
542 if stdin.write_all(&frame).await.is_err() {
543 break;
544 }
545 if stdin.flush().await.is_err() {
546 break;
547 }
548 }
549 }
550
551 /// Background task that parses `Content-Length`-framed JSON-RPC frames from
552 /// the LSP server's stdout. Pushes each parsed JSON value to `tx`. Exits
553 /// when stdout closes or a frame is malformed (we choose to fail closed
554 /// rather than risk hanging).
555 async fn reader_task(mut stdout: impl AsyncRead + Unpin, tx: mpsc::Sender<Value>) {
556 let mut buf: Vec<u8> = Vec::with_capacity(8 * 1024);
557 let mut tmp = [0u8; 4096];
558 loop {
559 let n = match stdout.read(&mut tmp).await {
560 Ok(0) => return,
561 Ok(n) => n,
562 Err(_) => return,
563 };
564 buf.extend_from_slice(&tmp[..n]);
565 loop {
566 let (header_end, content_length) = match parse_header(&buf) {
567 Ok(Some(frame)) => frame,
568 Ok(None) => break,
569 Err(_) => return,
570 };
571 // Both operands are bounded by parse_header.
572 let frame_end = header_end + content_length;
573 if buf.len() < frame_end {
574 break;
575 }
576 let value = match serde_json::from_slice::<Value>(&buf[header_end..frame_end]) {
577 Ok(value) => value,
578 Err(_) => return,
579 };
580 buf.drain(..frame_end);
581 if tx.send(value).await.is_err() {
582 return;
583 }
584 }
585 }
586 }
587
588 /// Distinguish incomplete headers from malformed or oversized frames so a
589 /// broken server cannot cause an indefinitely growing input buffer.
590 fn parse_header(buf: &[u8]) -> Result<Option<(usize, usize)>> {
591 let Some(pos) = buf.windows(4).position(|window| window == b"\r\n\r\n") else {
592 if buf.len() > MAX_LSP_HEADER_BYTES {
593 return Err(anyhow!("LSP header exceeds size limit"));
594 }
595 return Ok(None);
596 };
597 if pos + 4 > MAX_LSP_HEADER_BYTES {
598 return Err(anyhow!("LSP header exceeds size limit"));
599 }
600 let header = std::str::from_utf8(&buf[..pos]).context("invalid LSP header encoding")?;
601 let mut content_length = None;
602 for line in header.split("\r\n") {
603 let (name, value) = line.split_once(':').context("malformed LSP header")?;
604 if name.eq_ignore_ascii_case("Content-Length") {
605 if content_length.is_some() {
606 return Err(anyhow!("duplicate LSP Content-Length"));
607 }
608 let length = value
609 .trim()
610 .parse::<usize>()
611 .context("invalid LSP Content-Length")?;
612 if length == 0 || length > MAX_LSP_FRAME_BYTES {
613 return Err(anyhow!("LSP frame exceeds size limit or is empty"));
614 }
615 content_length = Some(length);
616 }
617 }
618 Ok(Some((
619 pos + 4,
620 content_length.context("missing LSP Content-Length")?,
621 )))
622 }
623
624 /// Background task that consumes inbound JSON values, classifies them as
625 /// notifications/responses, and routes accordingly.
626 async fn dispatcher_task(
627 mut rx: mpsc::Receiver<Value>,
628 tx_diag: mpsc::Sender<DiagnosticMessage>,
629 pending: Arc<AsyncMutex<HashMap<i64, oneshot::Sender<Value>>>>,
630 ) {
631 while let Some(value) = rx.recv().await {
632 // Notifications have a `method` and no `id`.
633 let method = value.get("method").and_then(|v| v.as_str());
634 if method == Some("textDocument/publishDiagnostics") {
635 if let Some(publication) = parse_publish_diagnostics(&value) {
636 // Do not let an unconsumed diagnostics burst prevent reply
637 // delivery or EOF cleanup. Overflow closes this transport's
638 // dispatcher, giving callers an explicit channel error.
639 if tx_diag.try_send(publication).is_err() {
640 break;
641 }
642 }
643 continue;
644 }
645 // Replies have an `id` and a `result` or `error`.
646 if let Some(id) = value.get("id").and_then(|v| v.as_i64()) {
647 let mut map = pending.lock().await;
648 if let Some(slot) = map.remove(&id) {
649 let _ = slot.send(value);
650 }
651 }
652 }
653 // Reader EOF/malformed frames and queue overflow wake every pending
654 // request immediately instead of leaving reply slots until their timeout.
655 pending.lock().await.clear();
656 }
657
658 /// Decode a `textDocument/publishDiagnostics` notification.
659 fn parse_publish_diagnostics(value: &Value) -> Option<DiagnosticMessage> {
660 let params = value.get("params")?;
661 let uri = params.get("uri")?.as_str()?;
662 let path = path_from_uri(uri)?;
663 let version = match params.get("version") {
664 None | Some(Value::Null) => None,
665 Some(value) => Some(value.as_i64()?),
666 };
667 let raw = params.get("diagnostics")?.as_array()?;
668 let mut out = Vec::with_capacity(raw.len());
669 for d in raw {
670 let range = d.get("range")?;
671 let start = range.get("start")?;
672 let line = start.get("line")?.as_u64()? as u32 + 1;
673 let column = start.get("character")?.as_u64()? as u32 + 1;
674 let severity = Severity::from_lsp(d.get("severity").and_then(|v| v.as_i64()))
675 .unwrap_or(Severity::Error);
676 let message = d
677 .get("message")
678 .and_then(|v| v.as_str())
679 .unwrap_or("")
680 .to_string();
681 out.push(Diagnostic {
682 line,
683 column,
684 severity,
685 message,
686 });
687 }
688 Some((path, version, out))
689 }
690
691 /// Encode an absolute filesystem path without following its links.
692 pub(crate) fn uri_from_path(path: &Path) -> String {
693 reqwest::Url::from_file_path(path)
694 .map(|url| url.to_string())
695 .unwrap_or_default()
696 }
697
698 /// File URLs only: no network authority, query or fragment. Decoding happens
699 /// before the workspace no-follow opener validates the target path.
700 pub(super) fn path_from_uri(uri: &str) -> Option<PathBuf> {
701 let url = reqwest::Url::parse(uri).ok()?;
702 if url.scheme() != "file"
703 || url.host_str().is_some_and(|host| !host.is_empty())
704 || !url.username().is_empty()
705 || url.password().is_some()
706 || url.port().is_some()
707 || url.query().is_some()
708 || url.fragment().is_some()
709 {
710 return None;
711 }
712 url.to_file_path().ok()
713 }
714
715 #[cfg(test)]
716 pub(super) mod tests {
717 use super::*;
718
719 fn fixture_path(name: &str) -> PathBuf {
720 std::env::temp_dir()
721 .join("codewhale-lsp-fixture")
722 .join(name)
723 }
724
725 #[test]
726 fn parses_lsp_header() {
727 let frame = b"Content-Length: 5\r\n\r\nhello";
728 let (end, len) = parse_header(frame)
729 .expect("valid header")
730 .expect("header parses");
731 assert_eq!(end, 21);
732 assert_eq!(len, 5);
733 }
734
735 #[test]
736 fn parse_header_returns_none_when_truncated() {
737 let frame = b"Content-Length: 5\r\nMissingTerm";
738 assert!(parse_header(frame).unwrap().is_none());
739 }
740
741 #[test]
742 fn parses_publish_diagnostics_payload() {
743 let payload = json!({
744 "jsonrpc": "2.0",
745 "method": "textDocument/publishDiagnostics",
746 "params": {
747 "uri": uri_from_path(&fixture_path("foo.rs")),
748 "diagnostics": [
749 {
750 "range": {
751 "start": { "line": 11, "character": 7 },
752 "end": { "line": 11, "character": 8 }
753 },
754 "severity": 1,
755 "message": "missing semicolon"
756 }
757 ]
758 }
759 });
760 let (path, version, diags) = parse_publish_diagnostics(&payload).expect("parses");
761 assert_eq!(path, fixture_path("foo.rs"));
762 assert_eq!(version, None);
763 assert_eq!(diags.len(), 1);
764 assert_eq!(diags[0].line, 12);
765 assert_eq!(diags[0].column, 8);
766 assert_eq!(diags[0].severity, Severity::Error);
767 assert_eq!(diags[0].message, "missing semicolon");
768 }
769
770 #[test]
771 fn round_trips_uri_path() {
772 let path = fixture_path("example/foo.rs");
773 let uri = uri_from_path(&path);
774 assert_eq!(path_from_uri(&uri), Some(path));
775 }
776
777 #[tokio::test]
778 async fn closed_diagnostics_channel_is_an_error_not_an_empty_result() {
779 let (tx_outbound, _rx_outbound) = mpsc::channel(1);
780 let (tx_diag, rx_diag) = mpsc::channel(1);
781 drop(tx_diag);
782 let transport = StdioLspTransport {
783 child: Arc::new(AsyncMutex::new(None)),
784 tasks: Vec::new(),
785 tx_outbound,
786 diagnostics_gate: AsyncMutex::new(()),
787 diagnostics_rx: AsyncMutex::new(rx_diag),
788 pending: Arc::new(AsyncMutex::new(HashMap::new())),
789 next_id: AsyncMutex::new(1),
790 language_id: "rust".to_string(),
791 opened: AsyncMutex::new(HashMap::new()),
792 };
793
794 let error = transport
795 .diagnostics_for(
796 &fixture_path("closed-channel.rs"),
797 "fn main() {}\n",
798 Duration::from_millis(10),
799 )
800 .await
801 .expect_err("a closed transport must not look like an empty lint result");
802
803 assert!(
804 error.to_string().contains("diagnostics channel closed"),
805 "unexpected error: {error}"
806 );
807 }
808
809 fn diagnostic_fixture() -> (
810 StdioLspTransport,
811 mpsc::Receiver<Vec<u8>>,
812 mpsc::Sender<DiagnosticMessage>,
813 ) {
814 let (tx_outbound, rx_outbound) = mpsc::channel(8);
815 let (tx_diag, rx_diag) = mpsc::channel(8);
816 (
817 StdioLspTransport {
818 child: Arc::new(AsyncMutex::new(None)),
819 tasks: Vec::new(),
820 tx_outbound,
821 diagnostics_gate: AsyncMutex::new(()),
822 diagnostics_rx: AsyncMutex::new(rx_diag),
823 pending: Arc::new(AsyncMutex::new(HashMap::new())),
824 next_id: AsyncMutex::new(1),
825 language_id: "rust".into(),
826 opened: AsyncMutex::new(HashMap::new()),
827 },
828 rx_outbound,
829 tx_diag,
830 )
831 }
832
833 async fn next_document(rx: &mut mpsc::Receiver<Vec<u8>>) -> Value {
834 let frame = rx.recv().await.unwrap();
835 let (start, _) = parse_header(&frame).unwrap().unwrap();
836 serde_json::from_slice::<Value>(&frame[start..]).unwrap()
837 }
838
839 fn diagnostic(line: u32) -> Diagnostic {
840 Diagnostic {
841 line,
842 column: 1,
843 severity: Severity::Error,
844 message: "fixture".into(),
845 }
846 }
847
848 #[tokio::test]
849 async fn diagnostic_freshness_rejects_old_version_and_accepts_current_empty_publication() {
850 let (transport, mut outbound, diag) = diagnostic_fixture();
851 let server = tokio::spawn(async move {
852 let first = next_document(&mut outbound).await;
853 assert_eq!(first["method"], "textDocument/didOpen");
854 assert_eq!(first["params"]["textDocument"]["version"], 1);
855 let path =
856 path_from_uri(first["params"]["textDocument"]["uri"].as_str().unwrap()).unwrap();
857 diag.send((path.clone(), Some(1), vec![diagnostic(1)]))
858 .await
859 .unwrap();
860 let second = next_document(&mut outbound).await;
861 assert_eq!(second["method"], "textDocument/didChange");
862 assert_eq!(second["params"]["textDocument"]["version"], 2);
863 assert_eq!(second["params"]["contentChanges"][0]["text"], "new 🐋 text");
864 diag.send((path.clone(), Some(1), vec![diagnostic(99)]))
865 .await
866 .unwrap();
867 diag.send((path, Some(2), vec![])).await.unwrap();
868 });
869 let path = &fixture_path("freshness.rs");
870 let first = transport
871 .diagnostics_for(path, "old text", Duration::from_secs(1))
872 .await
873 .unwrap();
874 assert_eq!(first.freshness(), "verified");
875 assert_eq!(first.items[0].line, 1);
876 let second = transport
877 .diagnostics_for(path, "new 🐋 text", Duration::from_secs(1))
878 .await
879 .unwrap();
880 assert_eq!(second.freshness(), "verified");
881 assert_eq!(second.document_version, Some(2));
882 assert_eq!(second.diagnostic_version, Some(2));
883 assert!(
884 second.items.is_empty(),
885 "old error must not apply to the new text"
886 );
887 server.await.unwrap();
888 }
889
890 #[tokio::test]
891 async fn diagnostic_freshness_serializes_concurrent_file_requests() {
892 let (transport, mut outbound, diag) = diagnostic_fixture();
893 let server = tokio::spawn(async move {
894 for _ in 0..2 {
895 let request = next_document(&mut outbound).await;
896 assert!(
897 matches!(outbound.try_recv(), Err(mpsc::error::TryRecvError::Empty)),
898 "second file must not advance before this publication"
899 );
900 let path =
901 path_from_uri(request["params"]["textDocument"]["uri"].as_str().unwrap())
902 .unwrap();
903 let line = if path.ends_with("one.rs") { 1 } else { 2 };
904 diag.send((fixture_path("unrelated.rs"), Some(1), vec![diagnostic(99)]))
905 .await
906 .unwrap();
907 diag.send((path, Some(1), vec![diagnostic(line)]))
908 .await
909 .unwrap();
910 }
911 });
912 let one_path = fixture_path("one.rs");
913 let two_path = fixture_path("two.rs");
914 let (one, two) = tokio::join!(
915 transport.diagnostics_for(&one_path, "one", Duration::from_secs(1)),
916 transport.diagnostics_for(&two_path, "two", Duration::from_secs(1)),
917 );
918 assert_eq!(one.unwrap().items[0].line, 1);
919 assert_eq!(two.unwrap().items[0].line, 2);
920 server.await.unwrap();
921 }
922
923 #[tokio::test]
924 async fn diagnostic_freshness_unversioned_is_unverified_and_silence_is_error() {
925 let (transport, mut outbound, diag) = diagnostic_fixture();
926 let server = tokio::spawn(async move {
927 let request = next_document(&mut outbound).await;
928 let path =
929 path_from_uri(request["params"]["textDocument"]["uri"].as_str().unwrap()).unwrap();
930 diag.send((path, None, vec![])).await.unwrap();
931 // Keep the channels alive while the second request times out.
932 let _ = next_document(&mut outbound).await;
933 tokio::time::sleep(Duration::from_millis(100)).await;
934 });
935 let response = transport
936 .diagnostics_for(
937 &fixture_path("unversioned.rs"),
938 "text",
939 Duration::from_secs(1),
940 )
941 .await
942 .unwrap();
943 assert_eq!(response.freshness(), "unverified");
944 assert_eq!(response.diagnostic_version, None);
945 assert!(
946 transport
947 .diagnostics_for(
948 &fixture_path("silent.rs"),
949 "text",
950 Duration::from_millis(10)
951 )
952 .await
953 .unwrap_err()
954 .to_string()
955 .contains("timed out")
956 );
957 server.abort();
958 }
959
960 #[test]
961 fn diagnostic_freshness_parser_retains_version_and_rejects_malformed_version() {
962 let mut payload = json!({"params":{"uri":uri_from_path(&fixture_path("version.rs")),"version":4,"diagnostics":[]}});
963 assert_eq!(parse_publish_diagnostics(&payload).unwrap().1, Some(4));
964 payload["params"]["version"] = json!("4");
965 assert!(parse_publish_diagnostics(&payload).is_none());
966 }
967 #[cfg(unix)]
968 pub(crate) const STDIO_FIXTURE: &str = r#"
969 import json, os, select, sys, time
970 mode, pid_path = sys.argv[1:]
971 input_stream = os.fdopen(0, 'rb', buffering=0)
972 with open(pid_path, 'a' if mode == 'cache' else 'w') as f:
973 f.write(str(os.getpid()) + ('\n' if mode == 'cache' else ''))
974 def read():
975 headers = {}
976 while True:
977 line = input_stream.readline()
978 if not line:
979 raise SystemExit(0)
980 if line == b'\r\n':
981 break
982 key, value = line.decode().split(':', 1)
983 headers[key.lower()] = value.strip()
984 body = b''
985 length = int(headers['content-length'])
986 while len(body) < length:
987 chunk = input_stream.read(length - len(body))
988 if not chunk:
989 raise SystemExit(0)
990 body += chunk
991 return json.loads(body)
992 def send(value):
993 body = json.dumps(value).encode()
994 sys.stdout.buffer.write(('Content-Length: %d\r\n\r\n' % len(body)).encode() + body)
995 sys.stdout.buffer.flush()
996 request = read()
997 assert request['method'] == 'initialize'
998 if mode == 'eof':
999 raise SystemExit(0)
1000 if mode == 'silence':
1001 time.sleep(60)
1002 if mode == 'error':
1003 send({'jsonrpc':'2.0','id':request['id'],'error':{'code':-32002,'message':'fixture rejected initialization'}})
1004 time.sleep(60)
1005 if mode == 'delayed' and select.select([input_stream], [], [], 0.1)[0]:
1006 send({'jsonrpc':'2.0','id':request['id'],'error':{'code':-32002,'message':'notification arrived before initialize reply'}})
1007 raise SystemExit(1)
1008 if mode == 'stderr':
1009 for _ in range(64):
1010 os.write(2, b'x' * 32768)
1011 send({'jsonrpc':'2.0','id':request['id'],'result':{'capabilities':{}}})
1012 assert read()['method'] == 'initialized'
1013 while True:
1014 request = read()
1015 if request.get('method') == 'fixture/exit':
1016 raise SystemExit(0)
1017 if request.get('method') == 'fixture/overflow':
1018 for _ in range(80):
1019 send({'jsonrpc':'2.0','method':'textDocument/publishDiagnostics','params':{'uri':'file:///tmp/overflow.rs','version':1,'diagnostics':[]}})
1020 time.sleep(60)
1021 elif 'id' in request:
1022 send({'jsonrpc':'2.0','id':request['id'],'result':{'ready':True}})
1023 "#;
1024
1025 #[cfg(unix)]
1026 async fn spawn_stdio_fixture(
1027 mode: &str,
1028 root: &Path,
1029 wait: Duration,
1030 ) -> Result<StdioLspTransport> {
1031 StdioLspTransport::spawn_with_timeout(
1032 "python3",
1033 &[
1034 "-u".into(),
1035 "-c".into(),
1036 STDIO_FIXTURE.into(),
1037 mode.into(),
1038 root.join("pid").to_string_lossy().into_owned(),
1039 ],
1040 "rust",
1041 root.to_path_buf(),
1042 wait,
1043 )
1044 .await
1045 }
1046
1047 #[cfg(unix)]
1048 async fn assert_fixture_exited(root: &Path) {
1049 let pid = std::fs::read_to_string(root.join("pid"))
1050 .expect("fixture started")
1051 .parse::<i32>()
1052 .unwrap();
1053 for _ in 0..100 {
1054 // Signal zero probes only this fixture PID; it never sends a signal.
1055 if unsafe { libc::kill(pid, 0) } == -1 {
1056 assert_eq!(
1057 std::io::Error::last_os_error().raw_os_error(),
1058 Some(libc::ESRCH)
1059 );
1060 return;
1061 }
1062 tokio::time::sleep(Duration::from_millis(10)).await;
1063 }
1064 panic!("fixture child remained alive after transport termination");
1065 }
1066
1067 #[cfg(unix)]
1068 #[tokio::test]
1069 async fn stdio_startup_waits_for_initialize_and_drains_stderr_pressure() {
1070 for mode in ["delayed", "stderr"] {
1071 let root = tempfile::tempdir().unwrap();
1072 let transport = spawn_stdio_fixture(mode, root.path(), Duration::from_secs(3))
1073 .await
1074 .unwrap();
1075 let result = transport
1076 .request("fixture/ready", json!({}), Duration::from_secs(1))
1077 .await
1078 .unwrap();
1079 assert_eq!(result["ready"], true, "{mode}");
1080 transport.shutdown().await;
1081 assert_fixture_exited(root.path()).await;
1082 }
1083 }
1084
1085 #[cfg(unix)]
1086 #[tokio::test]
1087 async fn stdio_startup_error_eof_and_silence_fail_and_terminate_child() {
1088 for (mode, expected) in [
1089 ("error", "fixture rejected initialization"),
1090 ("eof", "channel closed"),
1091 ("silence", "timed out"),
1092 ] {
1093 let root = tempfile::tempdir().unwrap();
1094 let error = match timeout(
1095 Duration::from_secs(3),
1096 spawn_stdio_fixture(mode, root.path(), Duration::from_millis(500)),
1097 )
1098 .await
1099 .unwrap()
1100 {
1101 Ok(_) => panic!("{mode} unexpectedly initialized"),
1102 Err(error) => error,
1103 };
1104 assert!(format!("{error:#}").contains(expected), "{mode}: {error:#}");
1105 assert_fixture_exited(root.path()).await;
1106 }
1107 }
1108
1109 #[cfg(unix)]
1110 #[tokio::test]
1111 async fn stdio_diagnostic_overflow_fails_pending_request_and_terminates_child() {
1112 let root = tempfile::tempdir().unwrap();
1113 let transport = spawn_stdio_fixture("ready", root.path(), Duration::from_secs(3))
1114 .await
1115 .unwrap();
1116 let error = transport
1117 .request("fixture/overflow", json!({}), Duration::from_secs(2))
1118 .await
1119 .unwrap_err();
1120 assert!(error.to_string().contains("channel closed"), "{error:#}");
1121 assert!(transport.pending.lock().await.is_empty());
1122 assert_fixture_exited(root.path()).await;
1123 }
1124
1125 #[cfg(unix)]
1126 #[tokio::test]
1127 async fn language_server_does_not_inherit_parent_secret_env() {
1128 use crate::test_support::{EnvVarGuard, lock_test_env};
1129 let _env_lock = lock_test_env();
1130 let _secret = EnvVarGuard::set("CODEWHALE_TEST_LSP_SECRET", "lsp-secret-value");
1131 let root = tempfile::tempdir().unwrap();
1132 let report = root.path().join("env");
1133 // A fake server that records what it can see and exits before
1134 // answering `initialize`.
1135 let script = format!(
1136 "printf 'secret=%s\\n' \"${{CODEWHALE_TEST_LSP_SECRET-unset}}\" > '{}'; \
1137 [ -n \"$PATH\" ] && printf 'path-ok\\n' >> '{}'",
1138 report.display(),
1139 report.display()
1140 );
1141 let result = StdioLspTransport::spawn_with_timeout(
1142 "/bin/sh",
1143 &["-c".into(), script],
1144 "rust",
1145 root.path().to_path_buf(),
1146 Duration::from_secs(3),
1147 )
1148 .await;
1149 assert!(result.is_err(), "fake server exits without initializing");
1150 let seen = std::fs::read_to_string(&report).expect("fake server ran");
1151 assert!(seen.contains("secret=unset"), "{seen}");
1152 assert!(!seen.contains("lsp-secret-value"), "{seen}");
1153 assert!(seen.contains("path-ok"), "{seen}");
1154 }
1155
1156 #[tokio::test]
1157 async fn semantic_request_holds_document_gate_until_reply_before_diagnostics() {
1158 timeout(Duration::from_secs(2), async {
1159 let (transport, mut outbound, diag) = diagnostic_fixture();
1160 let transport = Arc::new(transport);
1161 let semantic = {
1162 let transport = transport.clone();
1163 tokio::spawn(async move {
1164 transport
1165 .request_for_document(
1166 &fixture_path("semantic.rs"),
1167 "before",
1168 "textDocument/definition",
1169 json!({}),
1170 Duration::from_secs(1),
1171 )
1172 .await
1173 })
1174 };
1175 let open = next_document(&mut outbound).await;
1176 assert_eq!(open["method"], "textDocument/didOpen");
1177 let request = next_document(&mut outbound).await;
1178 assert_eq!(request["method"], "textDocument/definition");
1179 let diagnostics = {
1180 let transport = transport.clone();
1181 tokio::spawn(async move {
1182 transport
1183 .diagnostics_for(
1184 &fixture_path("semantic.rs"),
1185 "after",
1186 Duration::from_secs(1),
1187 )
1188 .await
1189 })
1190 };
1191 assert!(
1192 timeout(Duration::from_millis(20), outbound.recv())
1193 .await
1194 .is_err()
1195 );
1196 transport
1197 .pending
1198 .lock()
1199 .await
1200 .remove(&request["id"].as_i64().unwrap())
1201 .unwrap()
1202 .send(json!({"result":[]}))
1203 .unwrap();
1204 assert_eq!(semantic.await.unwrap().unwrap().document_version, Some(1));
1205 let change = next_document(&mut outbound).await;
1206 assert_eq!(change["params"]["textDocument"]["version"], 2);
1207 diag.send((fixture_path("semantic.rs"), Some(2), vec![]))
1208 .await
1209 .unwrap();
1210 assert_eq!(
1211 diagnostics.await.unwrap().unwrap().document_version,
1212 Some(2)
1213 );
1214 })
1215 .await
1216 .unwrap();
1217 }
1218
1219 #[tokio::test]
1220 async fn semantic_deadline_includes_gate_wait_and_cleans_pending_request() {
1221 let (transport, _outbound, _diag) = diagnostic_fixture();
1222 let gate = transport.diagnostics_gate.lock().await;
1223 let result = timeout(
1224 Duration::from_millis(250),
1225 transport.request_for_document(
1226 &fixture_path("a.rs"),
1227 "a",
1228 "textDocument/definition",
1229 json!({}),
1230 Duration::from_millis(20),
1231 ),
1232 )
1233 .await
1234 .unwrap();
1235 assert!(result.unwrap_err().to_string().contains("timed out"));
1236 drop(gate);
1237 assert!(transport.pending.lock().await.is_empty());
1238 let result = transport
1239 .request_for_document(
1240 &fixture_path("a.rs"),
1241 "a",
1242 "textDocument/definition",
1243 json!({}),
1244 Duration::from_millis(20),
1245 )
1246 .await;
1247 assert!(result.unwrap_err().to_string().contains("timed out"));
1248 assert!(transport.pending.lock().await.is_empty());
1249 }
1250
1251 #[test]
1252 fn semantic_file_uris_decode_spaces_and_reject_authority_and_query() {
1253 let path = &fixture_path("a b#🐋.rs");
1254 assert_eq!(
1255 path_from_uri(&uri_from_path(path)).as_deref(),
1256 Some(path.as_path())
1257 );
1258 for uri in [
1259 "https://example.test/a.rs",
1260 "file://remote/tmp/a.rs",
1261 "file:///tmp/a.rs?q=1",
1262 "file:///tmp/a.rs#fragment",
1263 ] {
1264 assert!(path_from_uri(uri).is_none(), "{uri}");
1265 }
1266 }
1267
1268 #[tokio::test]
1269 async fn stdio_request_deadline_includes_full_outbound_queue() {
1270 let (transport, _outbound, _diag) = diagnostic_fixture();
1271 for _ in 0..8 {
1272 transport.tx_outbound.try_send(vec![]).unwrap();
1273 }
1274 let error = timeout(
1275 Duration::from_millis(250),
1276 transport.request("fixture/blocked", json!({}), Duration::from_millis(20)),
1277 )
1278 .await
1279 .unwrap()
1280 .unwrap_err();
1281 assert!(error.to_string().contains("timed out"));
1282 assert!(transport.pending.lock().await.is_empty());
1283 }
1284
1285 #[tokio::test]
1286 async fn stdio_reader_bounds_and_malformed_frames_close_pending_replies() {
1287 let frames = [
1288 vec![b'x'; MAX_LSP_HEADER_BYTES + 1],
1289 format!("Content-Length: {}\r\n\r\n", MAX_LSP_FRAME_BYTES + 1).into_bytes(),
1290 b"Content-Length: 999999999999999999999999999999999\r\n\r\n".to_vec(),
1291 b"Content-Length: 1\r\nContent-Length: 1\r\n\r\nx".to_vec(),
1292 b"Missing-Length: 1\r\n\r\nx".to_vec(),
1293 b"Content-Length: 1\r\n\r\n{".to_vec(),
1294 ];
1295 for frame in frames {
1296 let (mut producer, reader) = tokio::io::duplex(MAX_LSP_HEADER_BYTES * 2);
1297 let (tx, rx) = mpsc::channel(8);
1298 let (diag, _diagnostics) = mpsc::channel(8);
1299 let pending = Arc::new(AsyncMutex::new(HashMap::new()));
1300 let (reply, receiver) = oneshot::channel();
1301 pending.lock().await.insert(1, reply);
1302 let read_task = tokio::spawn(reader_task(reader, tx));
1303 let dispatch = tokio::spawn(dispatcher_task(rx, diag, pending.clone()));
1304 producer.write_all(&frame).await.unwrap();
1305 assert!(
1306 timeout(Duration::from_secs(1), receiver)
1307 .await
1308 .unwrap()
1309 .is_err()
1310 );
1311 read_task.await.unwrap();
1312 dispatch.await.unwrap();
1313 assert!(pending.lock().await.is_empty());
1314 }
1315 }
1316
1317 #[tokio::test]
1318 async fn stdio_reader_preserves_fragmented_and_coalesced_valid_frames() {
1319 let (mut producer, reader) = tokio::io::duplex(128);
1320 let (tx, mut rx) = mpsc::channel(8);
1321 let task = tokio::spawn(reader_task(reader, tx));
1322 producer.write_all(b"content-length: 2\r\n").await.unwrap();
1323 producer
1324 .write_all(b"\r\n{}Content-Length: 2\r\n\r\n[]")
1325 .await
1326 .unwrap();
1327 assert_eq!(
1328 timeout(Duration::from_secs(1), rx.recv()).await.unwrap(),
1329 Some(json!({}))
1330 );
1331 assert_eq!(
1332 timeout(Duration::from_secs(1), rx.recv()).await.unwrap(),
1333 Some(json!([]))
1334 );
1335 drop(producer);
1336 task.await.unwrap();
1337 assert!(rx.recv().await.is_none());
1338 }
1339 }
1340
1340 lines RUST