Merge pull request #48 from mist-go/lsp-patches-and-completions

Lsp patches and completions
This commit is contained in:
2026-05-24 17:44:37 +02:00
committed by GitHub
4 changed files with 424 additions and 180 deletions
Generated
+17
View File
@@ -448,6 +448,7 @@ dependencies = [
"dashmap 6.1.0", "dashmap 6.1.0",
"mist-codegen", "mist-codegen",
"mist-parser", "mist-parser",
"ropey",
"serde", "serde",
"serde_json", "serde_json",
"tokio", "tokio",
@@ -623,6 +624,16 @@ dependencies = [
"bitflags 2.11.1", "bitflags 2.11.1",
] ]
[[package]]
name = "ropey"
version = "1.6.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "93411e420bcd1a75ddd1dc3caf18c23155eda2c090631a85af21ba19e97093b5"
dependencies = [
"smallvec",
"str_indices",
]
[[package]] [[package]]
name = "scopeguard" name = "scopeguard"
version = "1.2.0" version = "1.2.0"
@@ -742,6 +753,12 @@ version = "1.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "6ce2be8dc25455e1f91df71bfa12ad37d7af1092ae736f3a6cd0e37bc7810596" checksum = "6ce2be8dc25455e1f91df71bfa12ad37d7af1092ae736f3a6cd0e37bc7810596"
[[package]]
name = "str_indices"
version = "0.4.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d08889ec5408683408db66ad89e0e1f93dff55c73a4ccc71c427d5b277ee47e6"
[[package]] [[package]]
name = "syn" name = "syn"
version = "2.0.117" version = "2.0.117"
+1
View File
@@ -19,3 +19,4 @@ cargo_metadata = "0.23.1"
serde = { version = "1.0", features = ["derive"] } serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0" serde_json = "1.0"
ropey = "1"
+271 -105
View File
@@ -3,11 +3,11 @@ pub mod rust_analyzer;
pub mod transpiler; pub mod transpiler;
use std::collections::{HashMap, HashSet}; use std::collections::{HashMap, HashSet};
use std::fs;
use std::path::{Component, PathBuf}; use std::path::{Component, PathBuf};
use std::sync::Arc; use std::sync::Arc;
use mist_parser::rev_mapper; use mist_parser::rev_mapper;
use ropey::Rope;
use tokio::sync::Mutex; use tokio::sync::Mutex;
use tower_lsp::jsonrpc::Result; use tower_lsp::jsonrpc::Result;
use tower_lsp::lsp_types::notification::Notification; use tower_lsp::lsp_types::notification::Notification;
@@ -24,6 +24,7 @@ struct Backend {
previous_diagnostics: Arc<Mutex<HashMap<Url, Vec<Diagnostic>>>>, previous_diagnostics: Arc<Mutex<HashMap<Url, Vec<Diagnostic>>>>,
rust_analyzer: Arc<Mutex<RustAnalyzer>>, rust_analyzer: Arc<Mutex<RustAnalyzer>>,
mapping: Arc<Mutex<HashMap<PathBuf, HashSet<(rev_mapper::RustMap, rev_mapper::MistMap)>>>>, mapping: Arc<Mutex<HashMap<PathBuf, HashSet<(rev_mapper::RustMap, rev_mapper::MistMap)>>>>,
documents: Arc<Mutex<HashMap<PathBuf, Rope>>>,
} }
/// Helper function to force percent-encoding on Windows drive colons /// Helper function to force percent-encoding on Windows drive colons
@@ -93,7 +94,7 @@ impl LanguageServer for Backend {
res.capabilities.text_document_sync = Some(TextDocumentSyncCapability::Options( res.capabilities.text_document_sync = Some(TextDocumentSyncCapability::Options(
TextDocumentSyncOptions { TextDocumentSyncOptions {
open_close: Some(true), open_close: Some(true),
change: Some(TextDocumentSyncKind::INCREMENTAL), change: Some(TextDocumentSyncKind::FULL),
will_save: Some(false), will_save: Some(false),
will_save_wait_until: Some(false), will_save_wait_until: Some(false),
save: Some(SaveOptions::default().into()), save: Some(SaveOptions::default().into()),
@@ -101,9 +102,20 @@ impl LanguageServer for Backend {
)); ));
res.capabilities.completion_provider = Some(CompletionOptions { res.capabilities.completion_provider = Some(CompletionOptions {
resolve_provider: None, resolve_provider: Some(true),
trigger_characters: Some(vec![".".to_string()]), trigger_characters: Some(vec![
..Default::default() ":".to_owned(),
".".to_owned(),
"'".to_owned(),
"(".to_owned(),
]),
all_commit_characters: None,
completion_item: Some(CompletionOptionsCompletionItem {
label_details_support: None,
}),
work_done_progress_options: WorkDoneProgressOptions {
work_done_progress: None,
},
}); });
res.capabilities.definition_provider = Some(OneOf::Left(true)); res.capabilities.definition_provider = Some(OneOf::Left(true));
@@ -122,8 +134,14 @@ impl LanguageServer for Backend {
let analyzer = self.rust_analyzer.clone(); let analyzer = self.rust_analyzer.clone();
let documents = self.documents.clone();
let mapping = self.mapping.clone();
tokio::spawn(async move { tokio::spawn(async move {
if let Some(root) = &*workspace_folder.lock().await { if let Some(root) = &*workspace_folder.lock().await {
let src_root = root.join("src");
transpiler::build(root); transpiler::build(root);
analyzer analyzer
@@ -132,6 +150,32 @@ impl LanguageServer for Backend {
.initialize(root) .initialize(root)
.await .await
.expect("Failed to initialize rust analyzer"); .expect("Failed to initialize rust analyzer");
// ---- LOAD ALL .MIST FILES ----
let mut files = Vec::new();
collect_mist_files(&src_root, &mut files);
for file in files {
if let Ok(text) = std::fs::read_to_string(&file) {
if let Ok(transpiled) = transpiler::transpile_text(&text) {
let rust_path = from_mist_to_rust(file.clone());
// store document
documents
.lock()
.await
.insert(file.clone(), Rope::from_str(&text));
// store mapping
mapping
.lock()
.await
.insert(rust_path, rev_mapper::get_mapping(&transpiled));
}
}
}
eprintln!("Ready to use");
} }
}); });
@@ -139,16 +183,16 @@ impl LanguageServer for Backend {
} }
async fn initialized(&self, _: InitializedParams) { async fn initialized(&self, _: InitializedParams) {
self.client
.log_message(MessageType::INFO, "server initialized!")
.await;
self.rust_analyzer self.rust_analyzer
.lock() .lock()
.await .await
.initialized() .initialized()
.await .await
.expect("Failed to initialize rust analyzer"); .expect("Failed to initialize rust analyzer");
self.client
.log_message(MessageType::INFO, "server initialized!")
.await;
} }
async fn shutdown(&self) -> Result<()> { async fn shutdown(&self) -> Result<()> {
@@ -213,7 +257,7 @@ impl LanguageServer for Backend {
}, },
end: Position { end: Position {
line, line,
character: column + 1, character: u32::MAX,
}, },
}, },
severity: Some(severity), severity: Some(severity),
@@ -249,13 +293,22 @@ impl LanguageServer for Backend {
async fn did_open(&self, mut params: DidOpenTextDocumentParams) { async fn did_open(&self, mut params: DidOpenTextDocumentParams) {
if params.text_document.language_id == "mist" { if params.text_document.language_id == "mist" {
let original_text = params.text_document.text.clone();
self.documents.lock().await.insert(
params.text_document.uri.to_file_path().unwrap(),
Rope::from_str(&original_text),
);
params.text_document.language_id = "rust".to_string(); params.text_document.language_id = "rust".to_string();
match transpiler::transpile_text(&params.text_document.text) { match transpiler::transpile_text(&original_text) {
Ok(transpiled_text) => { Ok(transpiled_text) => {
params.text_document.text = transpiled_text; params.text_document.text = transpiled_text;
let rust_path = let rust_path =
from_mist_to_rust(params.text_document.uri.to_file_path().unwrap()); from_mist_to_rust(params.text_document.uri.to_file_path().unwrap());
params.text_document.uri = Url::from_file_path(&rust_path).unwrap(); params.text_document.uri = Url::from_file_path(&rust_path).unwrap();
self.mapping.lock().await.insert( self.mapping.lock().await.insert(
@@ -265,8 +318,12 @@ impl LanguageServer for Backend {
} }
Err(e) => { Err(e) => {
self.client self.client
.log_message(MessageType::WARNING, format!("MIST-LSP: Syntax invalid during open/change. Parsing stopped: {:?}", e)) .log_message(
MessageType::WARNING,
format!("MIST-LSP: Syntax invalid during open/change: {:?}", e),
)
.await; .await;
return; return;
} }
} }
@@ -279,6 +336,72 @@ impl LanguageServer for Backend {
} }
} }
async fn did_change(&self, mut params: DidChangeTextDocumentParams) {
let mist_path = params.text_document.uri.to_file_path().unwrap();
let rust_path = from_mist_to_rust(mist_path.clone());
let rust_uri = Url::from_file_path(&rust_path).unwrap();
if let Some(change) = params.content_changes.first_mut() {
self.documents
.lock()
.await
.insert(mist_path, Rope::from_str(&change.text));
match transpiler::transpile_text(&change.text) {
Ok(transpiled_text) => {
change.text = transpiled_text;
self.mapping
.lock()
.await
.insert(rust_path, rev_mapper::get_mapping(&change.text));
}
Err(e) => {
self.client
.log_message(
MessageType::WARNING,
format!("MIST-LSP: Syntax invalid during change: {:?}", e),
)
.await;
return;
}
}
}
params.text_document.uri = rust_uri;
if let Ok(mut ra) = self.rust_analyzer.try_lock() {
let _ = ra
.notify(notification::DidChangeTextDocument::METHOD, params)
.await;
}
}
async fn did_close(&self, mut params: DidCloseTextDocumentParams) {
self.client
.log_message(MessageType::INFO, "MIST-LSP: Processing did_close event")
.await;
let rust_path = from_mist_to_rust(params.text_document.uri.to_file_path().unwrap());
let mist_path = params.text_document.uri.to_file_path().unwrap();
self.documents.lock().await.remove(&mist_path);
self.mapping.lock().await.remove(&rust_path);
params.text_document.uri = Url::from_file_path(&rust_path).unwrap();
if let Ok(mut ra) = self.rust_analyzer.try_lock() {
let _ = ra
.notify(notification::DidCloseTextDocument::METHOD, params)
.await;
}
}
async fn goto_definition( async fn goto_definition(
&self, &self,
params: GotoDefinitionParams, params: GotoDefinitionParams,
@@ -290,9 +413,9 @@ impl LanguageServer for Backend {
.to_file_path() .to_file_path()
.unwrap(); .unwrap();
let source = match fs::read_to_string(&file_path) { let source = match self.documents.lock().await.get(&file_path) {
Ok(src) => src, Some(doc) => doc.clone(),
Err(_) => return Ok(None), None => return Ok(None),
}; };
let inject = "__mist_23"; let inject = "__mist_23";
@@ -303,10 +426,12 @@ impl LanguageServer for Backend {
&inject, &inject,
); );
let output = match transpiler::transpile_text(&injected_source) { let output = Rope::from_str(
Ok(out) => out, &match transpiler::transpile_text(&injected_source.to_string()) {
Err(_) => return Ok(None), Ok(out) => out,
}; Err(_) => return Ok(None),
},
);
let (line, character) = match find_row_col(&output, inject) { let (line, character) = match find_row_col(&output, inject) {
Some(coords) => coords, Some(coords) => coords,
@@ -323,7 +448,7 @@ impl LanguageServer for Backend {
.request::<request::GotoDefinition>(lsp_types::GotoDefinitionParams { .request::<request::GotoDefinition>(lsp_types::GotoDefinitionParams {
text_document_position_params: lsp_types::TextDocumentPositionParams { text_document_position_params: lsp_types::TextDocumentPositionParams {
position: lsp_types::Position { position: lsp_types::Position {
line: line as u32 - 1, line: line as u32,
character: character as u32, character: character as u32,
}, },
text_document: lsp_types::TextDocumentIdentifier { uri }, text_document: lsp_types::TextDocumentIdentifier { uri },
@@ -347,60 +472,81 @@ impl LanguageServer for Backend {
}) })
} }
async fn completion(&self, _params: CompletionParams) -> Result<Option<CompletionResponse>> { async fn completion(&self, mut params: CompletionParams) -> Result<Option<CompletionResponse>> {
self.client self.client
.log_message(MessageType::INFO, "getting completion") .log_message(MessageType::INFO, "COMPLETEING!")
.await; .await;
Ok(Some(CompletionResponse::Array(vec![ let file_path = params
CompletionItem::new_simple("new".to_string(), "The new keyword".to_string()), .text_document_position
]))) .text_document
.uri
.to_file_path()
.unwrap();
let source = match self.documents.lock().await.get(&file_path) {
Some(doc) => doc.clone(),
None => return Ok(None),
};
let inject = "__mist_23";
let injected_source = insert_at_position(
&source,
params.text_document_position.position.line as usize + 1,
params.text_document_position.position.character as usize,
&inject,
);
let output = Rope::from_str(
&match transpiler::transpile_text(&injected_source.to_string()) {
Ok(out) => out,
Err(_) => return Ok(None),
},
);
let (line, character) = match find_row_col(&output, inject) {
Some(coords) => coords,
None => return Ok(None),
};
let uri =
Url::from_file_path(from_mist_to_rust(file_path)).expect("failed to generate rs url");
params.text_document_position.text_document.uri = uri;
params.text_document_position.position.line = line as u32;
params.text_document_position.position.character = character as u32;
let rs_res = self
.rust_analyzer
.lock()
.await
.request::<request::Completion>(params)
.await
.expect("Failed to send to rust");
Ok(rs_res.map(|rs_res| match rs_res {
CompletionResponse::Array(items) => {
CompletionResponse::Array(items.into_iter().map(simplify_item).collect())
}
CompletionResponse::List(list) => {
CompletionResponse::Array(list.items.into_iter().map(simplify_item).collect())
}
}))
} }
// async fn completion(&self, mut params: CompletionParams) -> Result<Option<CompletionResponse>> { async fn completion_resolve(&self, params: CompletionItem) -> Result<CompletionItem> {
// self.client match self
// .log_message(MessageType::INFO, "COMPLETEING!") .rust_analyzer
// .await; .lock()
.await
// let file_path = params .request::<request::ResolveCompletionItem>(params.clone())
// .text_document_position .await
// .text_document {
// .uri Ok(o) => Ok(o),
// .to_file_path() _ => Ok(params),
// .unwrap(); }
}
// let mut source = fs::read_to_string(&file_path).expect("Failed to read source");
// let inject = "__mist_23";
// source = insert_at_position(
// &source,
// params.text_document_position.position.line as usize + 1,
// params.text_document_position.position.character as usize,
// &inject,
// );
// let output = transpiler::transpile_text(&source).expect("Failed to transpile");
// let (line, character) = find_row_col(&output, inject).unwrap();
// let uri =
// Url::from_file_path(from_mist_to_rust(file_path)).expect("failed to generate rs url");
// params.text_document_position.text_document.uri = uri;
// params.text_document_position.position.line = line as u32;
// params.text_document_position.position.character = character as u32;
// let rs_res = self
// .rust_analyzer
// .lock()
// .await
// .request::<request::Completion>(params)
// .await
// .expect("Failed to send to rust");
// Ok(rs_res)
// }
} }
#[tokio::main] #[tokio::main]
@@ -413,6 +559,7 @@ pub async fn start() {
workspace_folder: Arc::new(Mutex::new(None)), workspace_folder: Arc::new(Mutex::new(None)),
previous_diagnostics: Arc::new(Mutex::new(HashMap::new())), previous_diagnostics: Arc::new(Mutex::new(HashMap::new())),
mapping: Arc::new(Mutex::new(HashMap::new())), mapping: Arc::new(Mutex::new(HashMap::new())),
documents: Arc::new(Mutex::new(HashMap::new())),
rust_analyzer: Arc::new(Mutex::new( rust_analyzer: Arc::new(Mutex::new(
RustAnalyzer::new().expect("Failed to create rust analyzer"), RustAnalyzer::new().expect("Failed to create rust analyzer"),
)), )),
@@ -462,45 +609,64 @@ pub fn from_rust_to_mist(mut path: PathBuf) -> PathBuf {
} }
} }
fn insert_at_position(s: &str, line: usize, col: usize, insert: &str) -> String { fn insert_at_position(rope: &Rope, line: usize, col: usize, insert: &str) -> Rope {
let mut lines: Vec<String> = s.lines().map(|l| l.to_string()).collect(); let mut rope = rope.clone();
// Convert to 0-index
let line_idx = line.saturating_sub(1); let line_idx = line.saturating_sub(1);
let col_idx = col.saturating_sub(1); let col_idx = col.saturating_sub(1);
// Ensure line exists (optional behavior: extend with empty lines) let line_idx = line_idx.min(rope.len_lines().saturating_sub(1));
if line_idx >= lines.len() {
lines.resize(line_idx + 1, String::new()); let line_start = rope.line_to_char(line_idx);
let line_slice = rope.line(line_idx);
let line_len = line_slice.len_chars();
let col_idx = col_idx.min(line_len);
let idx = line_start + col_idx;
rope.insert(idx, insert);
rope
}
fn find_row_col(rope: &Rope, needle: &str) -> Option<(usize, usize)> {
let text = rope.to_string();
let byte_idx = text.find(needle)?;
let char_idx = text[..byte_idx].chars().count();
let line_idx = rope.char_to_line(char_idx);
let line_start = rope.line_to_char(line_idx);
let col_idx = char_idx - line_start;
Some((line_idx, col_idx + 1))
}
fn simplify_item(mut item: CompletionItem) -> CompletionItem {
item.text_edit = None;
item.additional_text_edits = None;
item.command = None;
item
}
fn collect_mist_files(dir: &PathBuf, out: &mut Vec<PathBuf>) {
let Ok(entries) = std::fs::read_dir(dir) else {
return;
};
for entry in entries.flatten() {
let path = entry.path();
if path.is_dir() {
collect_mist_files(&path, out);
} else if path.extension().and_then(|e| e.to_str()) == Some("mist") {
out.push(path);
}
} }
let target = &mut lines[line_idx];
// Clamp column to valid char boundary
let char_idx = target
.char_indices()
.nth(col_idx)
.map(|(i, _)| i)
.unwrap_or(target.len()); // if past end, append
target.insert_str(char_idx, insert);
lines.join("\n")
}
fn find_row_col(output: &str, inject: &str) -> Option<(usize, usize)> {
// 1. Find the flat byte index just like your original code
let byte_idx = output.find(inject)?;
// 2. Slice the string up to the match point
let prefix = &output[..byte_idx];
// 3. Row = number of newlines found before the match + 1 (1-indexed)
let row = prefix.lines().count();
// 4. Column = character count of the remaining text on the current line + 1
// (Using .chars().count() ensures it works with multi-byte UTF-8 symbols)
let col = prefix.lines().last().unwrap_or("").chars().count() + 1;
Some((row, col))
} }
+135 -75
View File
@@ -1,10 +1,11 @@
use std::{collections::HashMap, path::PathBuf, process::Stdio, sync::Arc}; use std::{collections::HashMap, path::PathBuf, process::Stdio, sync::Arc, time::Duration};
use serde::{Deserialize, Serialize, de::DeserializeOwned}; use serde::{Deserialize, Serialize, de::DeserializeOwned};
use serde_json::{Value, json}; use serde_json::{Value, json};
use tokio::{ use tokio::{
io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader}, io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader},
sync::{Mutex, oneshot}, sync::{Mutex, oneshot},
time::timeout,
}; };
use tower_lsp::lsp_types::{ use tower_lsp::lsp_types::{
self, ClientCapabilities, InitializeParams, InitializedParams, Url, WorkspaceFolder, self, ClientCapabilities, InitializeParams, InitializedParams, Url, WorkspaceFolder,
@@ -12,10 +13,13 @@ use tower_lsp::lsp_types::{
request::{self, Request}, request::{self, Request},
}; };
const MAX_CONTENT_LENGTH: usize = 50 * 1024 * 1024;
const REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
#[derive(Debug, Deserialize)] #[derive(Debug, Deserialize)]
pub struct JsonRpcResponse<T> { pub struct JsonRpcResponse<T> {
pub jsonrpc: String, pub jsonrpc: String,
pub id: Option<usize>, pub id: Option<serde_json::Value>,
pub result: Option<T>, pub result: Option<T>,
pub error: Option<JsonRpcError>, pub error: Option<JsonRpcError>,
} }
@@ -27,13 +31,16 @@ pub struct JsonRpcError {
pub data: Option<serde_json::Value>, pub data: Option<serde_json::Value>,
} }
type PendingMap = Arc<Mutex<HashMap<usize, oneshot::Sender<Value>>>>; type PendingMap = Arc<Mutex<HashMap<usize, oneshot::Sender<Result<Value, String>>>>>;
#[derive(Debug)] #[derive(Debug)]
pub struct RustAnalyzer { pub struct RustAnalyzer {
stdin: tokio::process::ChildStdin, // Wrapped in a Mutex to support safe concurrent sharing if the design expands
stdin: Arc<Mutex<tokio::process::ChildStdin>>,
pending: PendingMap, pending: PendingMap,
id: usize, id: usize,
// Keep child handler to explicitly manage child process lifecycle and prevent zombie processes
_child: tokio::process::Child,
} }
async fn send_lsp_message<W: AsyncWriteExt + Unpin>( async fn send_lsp_message<W: AsyncWriteExt + Unpin>(
@@ -49,27 +56,37 @@ async fn send_lsp_message<W: AsyncWriteExt + Unpin>(
async fn read_lsp_message<R: AsyncBufReadExt + Unpin>( async fn read_lsp_message<R: AsyncBufReadExt + Unpin>(
reader: &mut R, reader: &mut R,
) -> Result<String, Box<dyn std::error::Error>> { ) -> Result<String, Box<dyn std::error::Error + Send + Sync>> {
let mut line = String::new(); let mut line = String::new();
let mut content_length = 0; let mut content_length = 0;
// Read headers until we hit the empty separator line (\r\n) // Guard against infinite header reading attacks/bugs (Max 100 headers)
loop { for _ in 0..100 {
line.clear(); line.clear();
reader.read_line(&mut line).await?; let bytes_read = reader.read_line(&mut line).await?;
if line == "\r\n" || line.is_empty() { if bytes_read == 0 || line == "\r\n" || line.is_empty() {
break; break;
} }
if line.to_lowercase().starts_with("content-length:") { if line.to_ascii_lowercase().starts_with("content-length:") {
content_length = line["content-length:".len()..].trim().parse::<usize>()?; if let Some(val_str) = line.split(':').nth(1) {
content_length = val_str.trim().parse::<usize>()?;
}
} }
} }
// Explode early if payload size violates strict guard rails to prevent memory-exhaustion (OOM)
if content_length == 0 { if content_length == 0 {
return Err("Missing or invalid Content-Length header".into()); return Err("Missing, invalid, or zero Content-Length header".into());
}
if content_length > MAX_CONTENT_LENGTH {
return Err(format!(
"Content-Length {} exceeds maximum threshold",
content_length
)
.into());
} }
// Read the exact byte buffer payload // Explicitly secure pre-allocation limit
let mut buffer = vec![0u8; content_length]; let mut buffer = vec![0u8; content_length];
reader.read_exact(&mut buffer).await?; reader.read_exact(&mut buffer).await?;
@@ -77,19 +94,26 @@ async fn read_lsp_message<R: AsyncBufReadExt + Unpin>(
} }
impl RustAnalyzer { impl RustAnalyzer {
pub fn new() -> Result<Self, Box<dyn std::error::Error>> { pub fn new() -> Result<Self, Box<dyn std::error::Error + Send + Sync>> {
let mut child = tokio::process::Command::new("rust-analyzer") let mut child = tokio::process::Command::new("rust-analyzer")
.stdin(Stdio::piped()) .stdin(Stdio::piped())
.stdout(Stdio::piped()) .stdout(Stdio::piped())
.stderr(Stdio::null()) // Ignore logs for simplicity .stderr(Stdio::null())
.spawn()?; .spawn()?;
let stdin = child.stdin.take().unwrap(); let stdin = child
let stdout = child.stdout.take().unwrap(); .stdin
.take()
.ok_or("Failed to open child stdin pipe")?;
let stdout = child
.stdout
.take()
.ok_or("Failed to open child stdout pipe")?;
let pending: PendingMap = Arc::new(Mutex::new(HashMap::new())); let pending: PendingMap = Arc::new(Mutex::new(HashMap::new()));
let pending_clone = pending.clone(); let pending_clone = pending.clone();
// Background supervisor task loop
tokio::spawn(async move { tokio::spawn(async move {
let mut stdout = BufReader::new(stdout); let mut stdout = BufReader::new(stdout);
@@ -97,7 +121,12 @@ impl RustAnalyzer {
let raw = match read_lsp_message(&mut stdout).await { let raw = match read_lsp_message(&mut stdout).await {
Ok(v) => v, Ok(v) => v,
Err(err) => { Err(err) => {
eprintln!("LSP read error: {err}"); eprintln!("LSP fatal stream read failure: {err}");
// CRITICAL: Notify all pending channels that the bridge broke down
let mut lock = pending_clone.lock().await;
for (_, tx) in lock.drain() {
let _ = tx.send(Err(format!("LSP reader task dropped: {}", err)));
}
break; break;
} }
}; };
@@ -105,51 +134,43 @@ impl RustAnalyzer {
let value: Value = match serde_json::from_str(&raw) { let value: Value = match serde_json::from_str(&raw) {
Ok(v) => v, Ok(v) => v,
Err(err) => { Err(err) => {
eprintln!("Invalid JSON from rust-analyzer: {err}"); eprintln!("Corrupted JSON received: {err}");
continue; continue; // Keep the connection running despite malformed frame
} }
}; };
if value.get("method").is_some() { // Filter server notification frames
if value.get("method").is_some() && value.get("id").is_none() {
continue; continue;
} }
let id = value.get("id").and_then(|v| v.as_u64()).map(|v| v as usize); let id = value.get("id").and_then(|v| v.as_u64()).map(|v| v as usize);
match id { if let Some(id) = id {
Some(id) => { let tx = pending_clone.lock().await.remove(&id);
let tx = pending_clone.lock().await.remove(&id); if let Some(tx) = tx {
let _ = tx.send(Ok(value));
match tx { } else {
Some(tx) => { eprintln!("Received orphaned or delayed frame for ID: {id}");
let _ = tx.send(value);
}
None => {
eprintln!("({:?}) {}", pending_clone, value);
eprintln!("Received response for unknown request id {id}");
}
}
}
None => {
// notification or server request
eprintln!("Received server notification/request: {raw}");
} }
} else {
eprintln!("Received unhandled protocol notification framework: {raw}");
} }
} }
}); });
Ok(Self { Ok(Self {
stdin, stdin: Arc::new(Mutex::new(stdin)),
pending, pending,
id: 0, id: 0,
_child: child,
}) })
} }
pub async fn request<R: Request>( pub async fn request<R: Request>(
&mut self, &mut self,
params: R::Params, params: R::Params,
) -> Result<R::Result, Box<dyn std::error::Error>> ) -> Result<R::Result, Box<dyn std::error::Error + Send + Sync>>
where where
R::Result: DeserializeOwned, R::Result: DeserializeOwned,
{ {
@@ -160,45 +181,86 @@ impl RustAnalyzer {
let (tx, rx) = oneshot::channel(); let (tx, rx) = oneshot::channel();
self.pending.lock().await.insert(id, tx); // Scope the lock allocation tightly
{
send_lsp_message( self.pending.lock().await.insert(id, tx);
&mut self.stdin,
&json!({
"jsonrpc": "2.0",
"id": id,
"method": R::METHOD,
"params": params,
}),
)
.await?;
let value = rx.await?;
let envelope: JsonRpcResponse<R::Result> = serde_json::from_value(value)?;
if let Some(err) = envelope.error {
return Err(format!("LSP Error ({}): {}", err.code, err.message).into());
} }
envelope.result.ok_or_else(|| "missing result".into()) let payload = json!({
"jsonrpc": "2.0",
"id": id,
"method": R::METHOD,
"params": params,
});
// Acquire lock on writing stream to ensure thread safety
let mut stdin_lock = self.stdin.lock().await;
if let Err(err) = send_lsp_message(&mut *stdin_lock, &payload).await {
// Rollback the pending map operation to avoid internal memory memory-leaks if serialization/IO errors trigger
self.pending.lock().await.remove(&id);
return Err(Box::new(err));
}
// Explicit drop of write lock early so other operations can pipe messages synchronously
drop(stdin_lock);
// Enforce an absolute time constraint limit to break free from hanging processes
let response_payload = match timeout(REQUEST_TIMEOUT, rx).await {
Ok(Ok(Ok(value))) => value,
Ok(Ok(Err(task_err))) => return Err(task_err.into()),
Ok(Err(_oneshot_canceled)) => {
return Err(
"Bridge connection closed down; reader channel dropped unexpectedly".into(),
);
}
Err(_timeout_elapsed) => {
// Clear state tracking entries dynamically upon expiration failure
self.pending.lock().await.remove(&id);
return Err(
format!("Request ID {} timed out after {:?}", id, REQUEST_TIMEOUT).into(),
);
}
};
let envelope: JsonRpcResponse<R::Result> =
serde_json::from_value(response_payload.clone())?;
if let Some(err) = envelope.error {
return Err(format!("LSP Engine Error ({}): {}", err.code, err.message).into());
}
envelope.result.ok_or_else(|| {
format!(
"Missing inner structural payload result: {:?}",
response_payload
)
.into()
})
} }
pub async fn notify<T: Serialize>(&mut self, method: &str, req: T) -> std::io::Result<()> { pub async fn notify<T: Serialize>(
send_lsp_message( &mut self,
&mut self.stdin, method: &str,
&json!({ req: T,
"jsonrpc": "2.0", ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
"method": method, let payload = json!({
"params": req, "jsonrpc": "2.0",
}), "method": method,
) "params": req,
.await });
let mut stdin_lock = self.stdin.lock().await;
send_lsp_message(&mut *stdin_lock, &payload).await?;
Ok(())
} }
} }
impl RustAnalyzer { impl RustAnalyzer {
pub async fn initialize(&mut self, root: &PathBuf) -> Result<(), Box<dyn std::error::Error>> { pub async fn initialize(
&mut self,
root: &PathBuf,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let project_uri = Url::from_directory_path(root) let project_uri = Url::from_directory_path(root)
.map_err(|_| "Failed to convert path to valid file:// URL")?; .map_err(|_| "Failed to convert path to valid file:// URL")?;
@@ -224,14 +286,12 @@ impl RustAnalyzer {
}; };
self.request::<request::Initialize>(init_params).await?; self.request::<request::Initialize>(init_params).await?;
Ok(()) Ok(())
} }
pub async fn initialized(&mut self) -> Result<(), Box<dyn std::error::Error>> { pub async fn initialized(&mut self) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
self.notify(Initialized::METHOD, InitializedParams {}) self.notify(Initialized::METHOD, InitializedParams {})
.await?; .await?;
Ok(()) Ok(())
} }
} }