use std::future::Future;
use std::ops::ControlFlow;
use std::path::{Path, PathBuf};
use std::process::Stdio;
use std::sync::mpsc::Sender;
use std::sync::Arc;
use async_lsp::lsp_types::notification::{LogMessage, PublishDiagnostics, ShowMessage};
use async_lsp::lsp_types::{InitializeParams, InitializedParams};
use async_lsp::router::Router;
use tokio_util::compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt};
use super::queue::{self, WireEnv};
use super::sync::{self, FlushError};
use super::trace_io::{Direction, Observed};
use super::Client;
use crate::caps::ServerCaps;
use crate::convert::diag_from_lsp;
use crate::protocol::*;
use crate::registry;
use crate::target::Workspace;
pub(crate) struct ClientState {
tx: Sender<LspEvent>,
id: ServerId,
caps: ServerCaps,
sync: Arc<parking_lot::Mutex<sync::SyncState>>,
config: Option<serde_json::Value>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SpawnError {
RootUri,
Startup(String),
}
impl std::fmt::Display for SpawnError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::RootUri => write!(f, "the workspace root is not an absolute path"),
Self::Startup(message) => write!(f, "{message}"),
}
}
}
impl std::error::Error for SpawnError {}
enum Launch {
Local {
cmd: String,
args: Vec<String>,
cwd: PathBuf,
},
Container(strop_containers::AdmittedExec),
Remote(RemoteLaunch),
}
const REMOTE_EXIT_GRACE: std::time::Duration = std::time::Duration::from_secs(3);
const STDERR_TAIL_CAP: usize = 8192;
struct RemoteLaunch {
command: std::process::Command,
supervision: strop_remote::SupervisionKey,
}
fn remote_launch(
endpoint: &strop_workspace::RemoteEndpoint,
spec: ®istry::ServerSpec<'_>,
root: &Path,
) -> Result<RemoteLaunch, SpawnError> {
let args: Vec<std::ffi::OsString> = spec
.args
.iter()
.map(|arg| std::ffi::OsString::from(arg.as_str()))
.collect();
let command = strop_remote::RemoteCommand::new(spec.command, args, root).map_err(|error| {
SpawnError::Startup(format!("remote command rejected for {endpoint}: {error}"))
})?;
let (mut ssh, supervision) =
strop_remote::command_supervised(endpoint, &command, strop_remote::StdinMode::Relayed)
.map_err(|error| SpawnError::Startup(format!("ssh for {endpoint}: {error}")))?;
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
ssh.process_group(0);
}
Ok(RemoteLaunch {
command: ssh,
supervision,
})
}
fn container_launch(
id: &strop_workspace::ContainerId,
spec: ®istry::ServerSpec<'_>,
root: &Path,
token: &strop_core::worker::CancelToken,
) -> Result<strop_containers::AdmittedExec, SpawnError> {
let engine = strop_containers::engine(token)
.map_err(|error| SpawnError::Startup(format!("container engine: {error}")))?;
strop_containers::ExecSpec::resolve(&engine, id, spec.command, spec.args, root, token)
.map_err(|error| SpawnError::Startup(format!("container exec admission: {error}")))
}
pub(crate) fn client_router(
tx: Sender<LspEvent>,
id: ServerId,
caps: ServerCaps,
sync: Arc<parking_lot::Mutex<sync::SyncState>>,
workspace: Workspace,
name: String,
config: Option<serde_json::Value>,
) -> Router<ClientState> {
let mut router = Router::new(ClientState {
tx,
id,
caps,
sync,
config,
});
let diag_workspace = workspace;
router.notification::<PublishDiagnostics>(move |st, params| {
let Some(path) = diag_workspace.decode(¶ms.uri) else {
return ControlFlow::Continue(());
};
let context = st.sync.lock().diagnostic_context(
&path,
params.version.map(WireVersion::new),
st.id,
st.caps.encoding(),
);
if let Some(context) = context {
let diags = params.diagnostics.iter().map(diag_from_lsp).collect();
let _ = st.tx.send(LspEvent::Diagnostics {
context,
doc: strop_workspace::ResourceLocation {
filesystem: diag_workspace.target(),
path,
},
diags,
});
}
ControlFlow::Continue(())
});
router.notification::<ShowMessage>(move |st, params| {
strop_trace::record_with(strop_trace::EventKind::LspMessage, || {
serde_json::json!({"service":"lsp","server":st.id,"method":"window/showMessage",
"message":strop_trace::preview(¶ms.message)})
});
let _ = st.tx.send(LspEvent::ServerMessage {
server: st.id,
name: name.clone(),
text: params.message,
});
ControlFlow::Continue(())
});
router.notification::<LogMessage>(move |st, params| {
strop_trace::record_with(strop_trace::EventKind::LspMessage, || {
serde_json::json!({"service":"lsp","server":st.id,"method":"window/logMessage",
"message":strop_trace::preview(¶ms.message)})
});
ControlFlow::Continue(())
});
router.request::<async_lsp::lsp_types::request::WorkspaceConfiguration, _>(|st, params| {
let config = st.config.clone();
async move {
let answer: Vec<serde_json::Value> = params
.items
.iter()
.map(|item| match (&config, &item.section) {
(Some(config), Some(section)) => config
.get(section)
.cloned()
.unwrap_or(serde_json::Value::Null),
(Some(config), None) => config.clone(),
(None, _) => serde_json::Value::Null,
})
.collect();
Ok(answer)
}
});
router.unhandled_notification(|st, notif| {
strop_trace::record_with(strop_trace::EventKind::LspMessage, || {
serde_json::json!({"service":"lsp","server":st.id,"method":notif.method,"ignored":true})
});
ControlFlow::Continue(())
});
router
}
impl Client {
pub fn spawn(
spec: ®istry::ServerSpec<'_>,
workspace: Workspace,
tx: Sender<LspEvent>,
token: &strop_core::worker::CancelToken,
) -> Result<Self, SpawnError> {
let root_uri = workspace.uri(workspace.root()).ok_or(SpawnError::RootUri)?;
let launch = match &workspace {
Workspace::Local { root } => Launch::Local {
cmd: spec.command.to_string(),
args: spec.args.to_vec(),
cwd: root.clone(),
},
Workspace::Remote { endpoint, root } => {
Launch::Remote(remote_launch(endpoint, spec, root)?)
}
Workspace::Container { container, root } => {
Launch::Container(container_launch(container, spec, root, token)?)
}
};
let remote = workspace.endpoint().is_some();
let in_container = matches!(workspace, Workspace::Container { .. });
let label_workspace = workspace.label();
let endpoint_display = workspace.endpoint().map(|e| e.to_string());
let supervision = match &launch {
Launch::Remote(remote) => Some(remote.supervision.clone()),
Launch::Local { .. } | Launch::Container(_) => None,
};
let id = ServerId::allocate();
let self_caps = ServerCaps::default();
let sync = Arc::new(parking_lot::Mutex::new(sync::SyncState::default()));
let diag_workspace = workspace.clone();
let name = spec.name.to_string();
let (mainloop, socket) = async_lsp::MainLoop::new_client({
let tx = tx.clone();
let caps = self_caps.clone();
let sync = sync.clone();
let workspace = diag_workspace.clone();
let name = name.clone();
let config = spec.init_options.cloned();
move |_server| client_router(tx, id, caps, sync, workspace, name, config)
});
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|error| {
SpawnError::Startup(format!("cannot build the LSP runtime: {error}"))
})?;
let handle = rt.handle().clone();
let cmd = spec.command.to_string();
let tx_fail = tx.clone();
let hint = spec
.install_hint
.map(ToString::to_string)
.unwrap_or_else(|| format!("install `{cmd}` or fix the command in languages.toml"));
let name_loop = name.clone();
let hint_loop = hint.clone();
let quitting = Arc::new(std::sync::atomic::AtomicBool::new(false));
let quitting_mainloop = quitting.clone();
let closed = Arc::new(std::sync::atomic::AtomicBool::new(false));
let closed_mainloop = closed.clone();
let stderr_tail = Arc::new(parking_lot::Mutex::new(Vec::new()));
let env = WireEnv {
id,
name: name.clone(),
hint: hint.clone(),
socket: socket.clone(),
handle: handle.clone(),
tx: tx.clone(),
caps: self_caps.clone(),
workspace: workspace.clone(),
sync: sync.clone(),
quitting: quitting.clone(),
closed,
};
let queue = queue::start(env)
.ok_or_else(|| SpawnError::Startup("cannot start the LSP wire worker".into()))?;
let (stop_signal, stopping) = tokio::sync::oneshot::channel();
let thread = std::thread::Builder::new()
.name("strop-lsp-client".into())
.spawn(move || {
let name = name_loop;
let hint = hint_loop;
rt.block_on(async move {
let mut command = match launch {
Launch::Local { cmd, args, cwd } => {
let mut command = tokio::process::Command::new(&cmd);
command.args(&args).current_dir(&cwd);
command
}
Launch::Remote(remote) => {
let mut command = tokio::process::Command::from(remote.command);
command.kill_on_drop(true);
command
}
Launch::Container(exec) => {
let mut command =
tokio::process::Command::from(exec.command());
command.kill_on_drop(true);
command
}
};
command
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
command.kill_on_drop(true);
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
command.as_std_mut().process_group(0);
}
match command.spawn() {
Ok(child) => {
let mut c = super::process::ServerProcess::new(child);
let Ok((stdout, stdin, mut stderr)) = c.take_io() else {
let _ = tx_fail.send(LspEvent::Failed {
server: id, name: name.clone(), hint: hint.clone(),
});
return;
};
let stderr_drain = {
let name_stderr = name.clone();
let tail = stderr_tail.clone();
tokio::spawn(async move {
use tokio::io::AsyncReadExt;
let mut chunk = [0; 4096];
loop {
match stderr.read(&mut chunk).await {
Ok(0) => break,
Ok(bytes) => {
let mut guard = tail.lock();
guard.extend_from_slice(&chunk[..bytes]);
let overflow = guard.len().saturating_sub(STDERR_TAIL_CAP);
if overflow > 0 {
guard.drain(..overflow);
}
strop_trace::record_with(strop_trace::EventKind::Error, || serde_json::json!({
"source":"lsp_stderr", "server":name_stderr, "bytes":bytes,
"message":strop_trace::preview(&String::from_utf8_lossy(&chunk[..bytes])),
}));
}
Err(error) => {
strop_trace::record_with(strop_trace::EventKind::Error, || serde_json::json!({
"source":"lsp_stderr_read", "server":name_stderr,"message":error.to_string(),
}));
break;
}
}
}
})
};
let label = format!("{name}@{label_workspace}");
strop_trace::record_with(strop_trace::EventKind::JobStarted, || {
let mut event = serde_json::json!({"service":"lsp","server":label,"pid":c.id()});
if let Some(endpoint) = &endpoint_display {
event["remote"] = serde_json::json!(endpoint);
}
event
});
let result = {
let mut run = std::pin::pin!(mainloop.run_buffered(
Observed::new(stdout, &label, Direction::Rx).compat(),
Observed::new(stdin, &label, Direction::Tx).compat_write(),
));
let mut stopping = std::pin::pin!(stopping);
std::future::poll_fn(|context| {
if stopping.as_mut().poll(context).is_ready() {
return std::task::Poll::Ready(Ok(()));
}
run.as_mut().poll(context)
}).await
};
closed_mainloop.store(true, std::sync::atomic::Ordering::Relaxed);
strop_trace::record_with(strop_trace::EventKind::JobFinished, || serde_json::json!({
"service":"lsp","server":label,"error":result.as_ref().err().map(ToString::to_string),
}));
let grace = if remote { REMOTE_EXIT_GRACE } else { std::time::Duration::ZERO };
if let Err(error) = c.finish(Some(stderr_drain), grace).await {
strop_trace::record_with(strop_trace::EventKind::Error, || serde_json::json!({
"source":"lsp_process_cleanup", "server":label, "message":error.to_string(),
}));
}
if !quitting_mainloop.load(std::sync::atomic::Ordering::Relaxed) {
let mut detail = match &result {
Err(error) => format!(": {error}"),
Ok(()) => String::new(),
};
{
let guard = stderr_tail.lock();
let text = String::from_utf8_lossy(&guard).trim().to_string();
if !text.is_empty() {
detail = format!("{detail}; stderr: {text}");
}
if let Some(key) = &supervision {
if let Some(outcome) =
key.records(&guard).last().map(|o| format!("{o:?}"))
{
detail = format!(" ({outcome}){detail}");
}
}
}
let hint = format!("{name} exited unexpectedly{detail} — {hint}");
let _ = tx_fail.send(LspEvent::Failed {
server: id,
name: name.clone(),
hint,
});
}
}
Err(error) => {
let where_ = if remote || in_container {
format!(" on {label_workspace}")
} else {
String::new()
};
let reason = format!("cannot run `{cmd}`{where_}: {error}");
strop_trace::record_with(strop_trace::EventKind::Error, || serde_json::json!({
"source":"lsp_spawn","server":name,"command":&cmd,"message":error.to_string(),
}));
let _ = tx_fail.send(LspEvent::Failed {
server: id, name: name.clone(), hint: format!("{reason} — {hint}"),
});
}
}
});
})
.map_err(|error| {
SpawnError::Startup(format!("cannot start the LSP client thread: {error}"))
})?;
let root_name = workspace
.root()
.file_name()
.map(|name| name.to_string_lossy().into_owned())
.unwrap_or_else(|| "root".into());
let client = Self {
id,
next_request: Arc::new(std::sync::atomic::AtomicU64::new(0)),
sync,
socket,
handle,
tx,
workspace,
thread: Arc::new(std::sync::Mutex::new(Some(thread))),
caps: self_caps,
quitting,
queue,
stop: Arc::new(super::ServiceStop(parking_lot::Mutex::new(Some(
stop_signal,
)))),
};
let params = InitializeParams {
#[allow(deprecated)] root_uri: Some(root_uri.clone()),
workspace_folders: Some(vec![async_lsp::lsp_types::WorkspaceFolder {
name: root_name,
uri: root_uri,
}]),
initialization_options: spec.init_options.cloned(),
capabilities: async_lsp::lsp_types::ClientCapabilities {
text_document: Some(async_lsp::lsp_types::TextDocumentClientCapabilities {
synchronization: Some(Default::default()),
publish_diagnostics: Some(Default::default()),
hover: Some(Default::default()),
definition: Some(Default::default()),
..Default::default()
}),
workspace: Some(async_lsp::lsp_types::WorkspaceClientCapabilities {
workspace_folders: Some(true),
configuration: Some(true),
..Default::default()
}),
general: Some(async_lsp::lsp_types::GeneralClientCapabilities {
position_encodings: Some(vec![
async_lsp::lsp_types::PositionEncodingKind::UTF8,
async_lsp::lsp_types::PositionEncodingKind::UTF16,
]),
..Default::default()
}),
..Default::default()
},
..Default::default()
};
let initializing = client.clone();
client.handle.spawn(async move {
let init = tokio::time::timeout(
std::time::Duration::from_secs(10),
initializing
.socket
.request::<async_lsp::lsp_types::request::Initialize>(params),
)
.await;
match init {
Ok(Ok(response)) => {
initializing.caps.set(response.capabilities);
let initialized = initializing
.socket
.notify::<async_lsp::lsp_types::notification::Initialized>(
InitializedParams {},
);
let flushed = initializing.finish_initialize();
if initialized.is_err() || flushed.is_err() {
let reason = match flushed {
Err(FlushError::VersionExhausted) => {
"document versions exhausted".to_string()
}
_ => String::new(),
};
let hint = if reason.is_empty() {
hint
} else {
format!("{hint} ({reason})")
};
let _ = initializing.tx.send(LspEvent::Failed {
server: id,
name: name.clone(),
hint,
});
return;
}
let _ = initializing.tx.send(LspEvent::Ready {
server: id,
name: name.clone(),
});
}
outcome => {
initializing.stop.halt();
let reason = match outcome {
Ok(Err(error)) => format!("initialize refused: {error}"),
_ => "initialize timed out".to_string(),
};
let _ = initializing.tx.send(LspEvent::Failed {
server: id,
name: name.clone(),
hint: format!("{reason} — {hint}"),
});
}
}
});
Ok(client)
}
}