use super::lifecycle::spawn_proxy_if_needed;
use crate::{core, mcp_stdio, tools};
use anyhow::Result;
pub(super) fn run_mcp_server() -> Result<()> {
use rmcp::ServiceExt;
let started_at = std::time::Instant::now();
unsafe { std::env::set_var("LEAN_CTX_MCP_SERVER", "1") };
crate::core::startup_guard::crash_loop_backoff(crate::core::startup_guard::MCP_PROCESS_NAME);
crate::core::layout_pin::heal();
let startup_lock = crate::core::startup_guard::try_acquire_lock(
"mcp-startup",
std::time::Duration::from_secs(3),
std::time::Duration::from_secs(30),
);
let parallelism = std::thread::available_parallelism().map_or(2, std::num::NonZeroUsize::get);
let worker_threads = resolve_worker_threads(parallelism);
let max_blocking_threads = (worker_threads * 4).clamp(8, 32);
let index_threads = {
let configured = crate::core::config::Config::load().max_index_threads_effective();
if configured > 0 {
configured
} else {
let concurrent = crate::ipc::process::find_pids_by_name("lean-ctx").len() + 1;
if concurrent > 1 {
herd_aware_index_threads(parallelism, concurrent)
} else {
0
}
}
};
if index_threads > 0 {
let _ = rayon::ThreadPoolBuilder::new()
.num_threads(index_threads)
.build_global();
}
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(worker_threads)
.max_blocking_threads(max_blocking_threads)
.enable_all()
.build()?;
let server = tools::create_server();
drop(startup_lock);
let result = rt.block_on(async {
core::logging::init_mcp_logging();
core::protocol::set_mcp_context(true);
core::plugins::PluginManager::init();
core::plugins::PluginManager::notify(core::plugins::executor::HookPoint::OnSessionStart);
tracing::info!(
"lean-ctx v{} MCP server starting",
env!("CARGO_PKG_VERSION")
);
core::pathjail::warn_if_relaxed();
spawn_parent_watchdog();
let _housekeeping = tokio::task::spawn_blocking(|| {
cleanup_orphan_mcp_processes();
spawn_proxy_if_needed();
crate::cli::wrapped_publish::maybe_auto_publish_background();
});
let transport =
mcp_stdio::HybridStdioTransport::new_server(tokio::io::stdin(), tokio::io::stdout());
let server_handle = server.clone();
let service = match server.serve(transport).await {
Ok(s) => s,
Err(e) => {
let msg = e.to_string();
if msg.contains("expect initialized")
|| msg.contains("context canceled")
|| msg.contains("broken pipe")
{
tracing::debug!("Client disconnected before init: {msg}");
return Ok(());
}
return Err(e.into());
}
};
tracing::info!(
time_to_initialize_ms = started_at.elapsed().as_millis() as u64,
"MCP server initialized"
);
core::startup_guard::reset_crash_loop(core::startup_guard::MCP_PROCESS_NAME);
match service.waiting().await {
Ok(reason) => {
tracing::info!("MCP server stopped: {reason:?}");
}
Err(e) => {
let msg = e.to_string();
if msg.contains("broken pipe")
|| msg.contains("connection reset")
|| msg.contains("context canceled")
{
tracing::info!("MCP server: transport closed ({msg})");
} else {
tracing::error!("MCP server error: {msg}");
}
}
}
server_handle.shutdown().await;
if core::plugins::PluginManager::has_listener("on_session_end") {
let _ = core::plugins::PluginManager::fire_hook(
&core::plugins::executor::HookPoint::OnSessionEnd,
);
}
core::tool_lifecycle::flush_all();
core::efficacy::capture();
Ok(())
});
shutdown_runtime_bounded(rt);
result
}
fn shutdown_runtime_bounded(rt: tokio::runtime::Runtime) {
rt.shutdown_timeout(std::time::Duration::from_secs(2));
}
fn cleanup_orphan_mcp_processes() {
#[cfg(unix)]
{
let my_pid = std::process::id();
let pids = crate::ipc::process::find_pids_by_name("lean-ctx");
for pid in pids {
if pid == my_pid {
continue;
}
if !is_orphan_mcp(pid) {
continue;
}
tracing::info!("[orphan-cleanup] killing orphan MCP process {pid} (parent=1)");
let _ = crate::ipc::process::terminate_gracefully(pid);
}
}
}
#[cfg(unix)]
fn is_orphan_mcp(pid: u32) -> bool {
let Ok(output) = std::process::Command::new("ps")
.args(["-o", "ppid=,command=", "-p", &pid.to_string()])
.output()
else {
return false;
};
let text = String::from_utf8_lossy(&output.stdout);
let line = text.trim();
if line.is_empty() {
return false;
}
is_orphan_mcp_ps_line(line)
}
#[cfg(unix)]
fn is_orphan_mcp_ps_line(line: &str) -> bool {
let Some((ppid_str, command)) = split_ppid_and_command(line) else {
return false;
};
let Ok(ppid) = ppid_str.parse::<u32>() else {
return false;
};
if ppid > 1 {
return false;
}
let mut parts = command.split_whitespace();
let Some(exe) = parts.next() else {
return false;
};
if !is_lean_ctx_executable(exe) {
return false;
}
let mut first_arg = parts.next();
if first_arg == Some("(deleted)") {
first_arg = parts.next();
}
matches!(first_arg, None | Some("mcp"))
}
#[cfg(unix)]
fn split_ppid_and_command(line: &str) -> Option<(&str, &str)> {
let trimmed = line.trim();
let split = trimmed
.char_indices()
.find_map(|(idx, ch)| ch.is_whitespace().then_some(idx))?;
let (ppid, rest) = trimmed.split_at(split);
Some((ppid, rest.trim_start()))
}
#[cfg(unix)]
fn is_lean_ctx_executable(value: &str) -> bool {
value == "lean-ctx" || value.ends_with("/lean-ctx")
}
fn spawn_parent_watchdog() {
#[cfg(unix)]
{
let ppid = unsafe { libc::getppid() } as u32;
if ppid <= 1 {
return;
}
std::thread::Builder::new()
.name("parent-watchdog".into())
.spawn(move || {
loop {
std::thread::sleep(std::time::Duration::from_secs(5));
let current_ppid = unsafe { libc::getppid() } as u32;
if current_ppid != ppid || current_ppid <= 1 {
tracing::info!(
"[parent-watchdog] parent PID changed ({ppid} → {current_ppid}), \
IDE likely closed — exiting to prevent orphan"
);
core::tool_lifecycle::flush_all();
std::process::exit(0);
}
}
})
.ok();
}
}
pub(super) fn resolve_worker_threads(parallelism: usize) -> usize {
std::env::var("LEAN_CTX_WORKER_THREADS")
.ok()
.and_then(|v| v.parse::<usize>().ok())
.unwrap_or_else(|| parallelism.clamp(1, 4))
}
pub(super) fn herd_aware_index_threads(cores: usize, concurrent: usize) -> usize {
let cores = cores.max(1);
let concurrent = concurrent.max(1);
(cores / concurrent).max(1)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn lone_session_keeps_all_cores() {
assert_eq!(herd_aware_index_threads(16, 1), 16);
assert_eq!(herd_aware_index_threads(8, 1), 8);
}
#[test]
fn fleet_splits_cores_and_stays_under_core_count() {
assert_eq!(herd_aware_index_threads(16, 10), 1);
assert!(herd_aware_index_threads(16, 10) * 10 < 16);
assert_eq!(herd_aware_index_threads(16, 2), 8);
assert_eq!(herd_aware_index_threads(16, 4), 4);
}
#[test]
fn never_returns_zero_threads() {
assert_eq!(herd_aware_index_threads(4, 32), 1);
assert_eq!(herd_aware_index_threads(0, 0), 1);
}
#[test]
fn runtime_shutdown_is_bounded_despite_hung_blocking_task() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::{Duration, Instant};
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(1)
.enable_all()
.build()
.expect("runtime builds");
let stop = std::sync::Arc::new(AtomicBool::new(false));
let stop_in = stop.clone();
rt.block_on(async {
let _abandoned = tokio::task::spawn_blocking(move || {
for _ in 0..600 {
if stop_in.load(Ordering::Relaxed) {
break;
}
std::thread::sleep(Duration::from_millis(100));
}
});
tokio::time::sleep(Duration::from_millis(50)).await;
});
let started = Instant::now();
shutdown_runtime_bounded(rt);
let elapsed = started.elapsed();
stop.store(true, Ordering::Relaxed);
assert!(
elapsed < Duration::from_secs(10),
"bounded shutdown must not wait for the hung task (took {elapsed:?})"
);
}
#[cfg(unix)]
#[test]
fn orphan_mcp_cleanup_matches_only_stdio_mcp_processes() {
assert!(is_orphan_mcp_ps_line("1 /Users/me/.local/bin/lean-ctx"));
assert!(is_orphan_mcp_ps_line("1 /opt/homebrew/bin/lean-ctx mcp"));
assert!(is_orphan_mcp_ps_line(
"1 /opt/homebrew/bin/lean-ctx (deleted) mcp"
));
assert!(!is_orphan_mcp_ps_line("99 /Users/me/.local/bin/lean-ctx"));
assert!(!is_orphan_mcp_ps_line(
"1 /Users/me/.local/bin/lean-ctx proxy start --port=4444"
));
assert!(!is_orphan_mcp_ps_line(
"1 /Users/me/.local/bin/lean-ctx serve --port=8080"
));
assert!(!is_orphan_mcp_ps_line(
"1 /Users/me/.local/bin/lean-ctx daemon start"
));
assert!(!is_orphan_mcp_ps_line(
"1 /usr/bin/sandbox-exec -f /tmp/seatbelt.sb /Users/me/.local/bin/lean-ctx proxy start --port=4444"
));
}
}