Simplify MCP executable monitoring

Exit the process from `main` for every command instead of returning. The
`mcp start` command parks Tokio's stdin reader on a read only the MCP host
can end, so dropping the runtime waits forever. Exiting in `main` also keeps
the CLI telemetry event, which the previous exit inside the MCP command
skipped.

That removes the reason for the `McpServerExit` enum, whose only job was to
carry an implementation detail out to the CLI so it could exit.

Watch the executable through its device and inode on Unix. That is a
complete file identity, so the length and modification time no longer add
anything. Drop the PATH scan: `current_exe` reports the symlink itself on
macOS, so it detects a Homebrew relink without it. This also drops the
`fabro-static` dependency and a clippy suppression.

Bound the shutdown wait after an upgrade is detected. The transport closes
by writing to a stdout the host may already have stopped reading, which
could hang the exit the change is supposed to trigger.

Log a warning when upgrade detection cannot start, rather than disabling it
silently.

Share one spawn helper between the two raw stdio tests, and link the test
executable instead of copying 200 MB of binary.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-08-01 09:10:28 -04:00
parent f868734f27
commit 14cf4f3e94
No known key found for this signature in database
9 changed files with 165 additions and 220 deletions

2
Cargo.lock generated
View file

@ -2838,7 +2838,6 @@ dependencies = [
"fabro-manifest",
"fabro-model",
"fabro-server",
"fabro-static",
"fabro-tool",
"fabro-types",
"fabro-util",
@ -2852,6 +2851,7 @@ dependencies = [
"tempfile",
"tokio",
"toml 0.8.23",
"tracing",
]
[[package]]

View file

@ -1,8 +1,6 @@
use std::fmt::Write as _;
use std::process;
use anyhow::{Context as _, Result};
use fabro_mcp_server::McpServerExit;
use crate::args::{McpAgent, McpCommand, McpNamespace, ServerConnectionArgs};
use crate::command_context::CommandContext;
@ -11,15 +9,7 @@ use crate::server_client;
pub(crate) async fn dispatch(ns: McpNamespace, base_ctx: &CommandContext) -> Result<()> {
match ns.command {
McpCommand::Start(args) => {
let exit =
fabro_mcp_server::start(server_settings(base_ctx, &args.connection)?).await?;
if exit == McpServerExit::ExecutableReplaced {
// Tokio's stdin worker can remain blocked after the MCP service
// closes. Exit at the CLI boundary so the host can reconnect to
// the replacement executable.
process::exit(0);
}
Ok(())
fabro_mcp_server::start(server_settings(base_ctx, &args.connection)?).await
}
McpCommand::Config(args) => {
let json = fabro_mcp_server::config_json(&config_settings(&args.connection))?;

View file

@ -120,8 +120,12 @@ async fn main() {
"{:?}",
miette::Report::new(CliDiagnostic::new(err, !json_mode))
);
std::process::exit(exit_code);
}
// Exit rather than returning. A command can leave a blocked worker thread
// behind — `mcp start` parks Tokio's stdin reader on a read only the MCP
// host can end — and dropping the runtime would wait on it forever.
std::process::exit(exit_code);
}
fn install_miette_hook() {

View file

@ -12,7 +12,7 @@ use std::collections::HashMap;
use std::fs;
use std::io::{BufRead as _, Write as _};
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::process::{Child, ChildStdin, Command, Stdio};
use std::thread;
use std::time::{Duration, Instant};
@ -519,40 +519,14 @@ async fn stdio_server_initializes_and_lists_run_tools() {
fn stdio_start_writes_only_json_rpc_to_stdout() {
let context = test_context!();
let fixture = mcp_stdio_fixture(&context, &[]);
let mut cmd = Command::new(&fixture.command[0]);
cmd.args(&fixture.command[1..])
.env_clear()
.envs(&fixture.env)
.current_dir(&fixture.current_dir)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
let mut child = cmd.spawn().unwrap();
let mut stdin = child.stdin.take().unwrap();
writeln!(
stdin,
r#"{{"jsonrpc":"2.0","id":1,"method":"initialize","params":{{"protocolVersion":"2025-06-18","capabilities":{{}},"clientInfo":{{"name":"fabro-test","version":"0.0.0"}}}}}}"#
)
.unwrap();
let (mut child, _stdin, response) =
spawn_stdio_server(&fixture, Path::new(&fixture.command[0]));
let stdout = child.stdout.take().unwrap();
let (tx, rx) = std::sync::mpsc::channel();
thread::spawn(move || {
let mut line = String::new();
let result = std::io::BufReader::new(stdout).read_line(&mut line);
let _ = tx.send(result.map(|_| line));
});
let line = rx
.recv_timeout(Duration::from_secs(5))
.expect("initialize response should arrive")
.expect("stdout should be readable");
let value: serde_json::Value = serde_json::from_str(line.trim()).unwrap();
assert_eq!(value["jsonrpc"], "2.0");
assert_eq!(value["result"]["serverInfo"]["name"], "fabro");
assert_eq!(response["jsonrpc"], "2.0");
assert_eq!(response["result"]["serverInfo"]["name"], "fabro");
assert_eq!(
value["result"]["serverInfo"]["version"],
response["result"]["serverInfo"]["version"],
env!("CARGO_PKG_VERSION")
);
@ -566,40 +540,17 @@ fn stdio_server_exits_when_executable_is_replaced() {
let context = test_context!();
let fixture = mcp_stdio_fixture(&context, &[]);
let directory = tempfile::tempdir().expect("replacement directory should exist");
// A symlink stands in for a Homebrew install: the server follows it to the
// real binary, so replacing the link changes the identity it watches without
// copying a multi-hundred-megabyte executable.
let executable = directory.path().join("fabro");
fs::copy(&fixture.command[0], &executable).expect("Fabro executable should be copied");
let mut cmd = Command::new(&executable);
cmd.args(&fixture.command[1..])
.env_clear()
.envs(&fixture.env)
.current_dir(&fixture.current_dir)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
std::os::unix::fs::symlink(&fixture.command[0], &executable)
.expect("Fabro executable should be linked");
let mut child = cmd.spawn().expect("MCP server should start");
let mut stdin = child.stdin.take().expect("MCP stdin should be available");
writeln!(
stdin,
r#"{{"jsonrpc":"2.0","id":1,"method":"initialize","params":{{"protocolVersion":"2025-06-18","capabilities":{{}},"clientInfo":{{"name":"fabro-test","version":"0.0.0"}}}}}}"#
)
.expect("initialize request should be written");
let stdout = child.stdout.take().expect("MCP stdout should be available");
let (tx, rx) = std::sync::mpsc::channel();
thread::spawn(move || {
let mut line = String::new();
let result = std::io::BufReader::new(stdout).read_line(&mut line);
let _ = tx.send(result.map(|_| line));
});
let response = rx
.recv_timeout(Duration::from_secs(5))
.expect("initialize response should arrive")
.expect("MCP stdout should be readable");
let response: serde_json::Value = serde_json::from_str(response.trim()).unwrap();
// `_stdin` holds the pipe open so replacement, not EOF, stops the server.
let (mut child, _stdin, response) = spawn_stdio_server(&fixture, &executable);
assert_eq!(response["result"]["serverInfo"]["name"], "fabro");
// Keep stdin open so replacement, rather than EOF, stops the server.
let replacement = directory.path().join("fabro-replacement");
fs::write(&replacement, b"replacement").expect("replacement file should be written");
fs::rename(replacement, &executable).expect("Fabro executable should be replaced");
@ -2537,6 +2488,45 @@ fn mcp_stdio_fixture(context: &fabro_test::TestContext, extra_args: &[&str]) ->
}
}
const MCP_INITIALIZE_REQUEST: &str = r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-06-18","capabilities":{},"clientInfo":{"name":"fabro-test","version":"0.0.0"}}}"#;
/// Starts `fabro mcp start` as a raw child process, sends `initialize`, and
/// returns the decoded response. Unlike `spawn_mcp_client`, the caller keeps
/// the `Child` and its stdin, so it can observe how and when the server exits.
fn spawn_stdio_server(
fixture: &McpStdioFixture,
program: &Path,
) -> (Child, ChildStdin, serde_json::Value) {
let mut child = Command::new(program)
.args(&fixture.command[1..])
.env_clear()
.envs(&fixture.env)
.current_dir(&fixture.current_dir)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("MCP server should start");
let mut stdin = child.stdin.take().expect("MCP stdin should be available");
writeln!(stdin, "{MCP_INITIALIZE_REQUEST}").expect("initialize request should be written");
let stdout = child.stdout.take().expect("MCP stdout should be available");
let (tx, rx) = std::sync::mpsc::channel();
thread::spawn(move || {
let mut line = String::new();
let result = std::io::BufReader::new(stdout).read_line(&mut line);
let _ = tx.send(result.map(|_| line));
});
let line = rx
.recv_timeout(Duration::from_secs(5))
.expect("initialize response should arrive")
.expect("MCP stdout should be readable");
let response = serde_json::from_str(line.trim()).expect("response should be JSON");
(child, stdin, response)
}
fn write_mcp_server_settings(
context: &mut fabro_test::TestContext,
storage_dir: &Path,

View file

@ -21,7 +21,6 @@ fabro-manifest = { path = "../../components/fabro-manifest" }
fabro-config = { path = "../../foundation/fabro-config" }
fabro-model = { path = "../../foundation/fabro-model" }
fabro-server = { path = "../fabro-server" }
fabro-static = { path = "../../foundation/fabro-static" }
fabro-tool = { path = "../../components/fabro-tool" }
fabro-types = { path = "../../foundation/fabro-types" }
fabro-util = { path = "../../foundation/fabro-util" }
@ -33,6 +32,7 @@ serde_json.workspace = true
strum.workspace = true
tokio.workspace = true
toml.workspace = true
tracing.workspace = true
[dev-dependencies]
httpmock = "0.8"

View file

@ -9,9 +9,7 @@ use anyhow::{Context as _, Result, anyhow};
use serde_json::map::Entry;
use serde_json::{Map, Value, json};
use crate::{McpAgent, McpConfigSettings, McpInitSettings};
const SERVER_NAME: &str = "fabro";
use crate::{McpAgent, McpConfigSettings, McpInitSettings, SERVER_NAME};
pub fn config_json(settings: &McpConfigSettings) -> Result<String> {
serde_json::to_string_pretty(&generic_config(settings))

View file

@ -5,143 +5,100 @@
//! process keeps its old API response decoder after the `fabro` file on disk is
//! upgraded.
use std::path::PathBuf;
use std::time::{Duration, SystemTime};
use std::{env, io};
use std::fs::Metadata;
use std::path::{self, PathBuf};
use std::time::Duration;
use std::{env, fs, io};
use fabro_static::EnvVars;
use tokio::time::Instant;
use tokio::{fs, time};
use tokio::time::{self, Instant, MissedTickBehavior};
const CHECK_INTERVAL: Duration = Duration::from_secs(1);
pub(crate) struct ExecutableMonitor {
path: PathBuf,
identity: ExecutableIdentity,
identity: Identity,
}
impl ExecutableMonitor {
pub(crate) async fn current() -> io::Result<Self> {
let path = invoked_executable_path().await?;
Self::new(path).await
pub(crate) fn current() -> io::Result<Self> {
Self::new(invoked_executable_path()?)
}
async fn new(path: PathBuf) -> io::Result<Self> {
let identity = ExecutableIdentity::from_metadata(&fs::metadata(&path).await?);
fn new(path: PathBuf) -> io::Result<Self> {
let identity = identity(&fs::metadata(&path)?);
Ok(Self { path, identity })
}
pub(crate) async fn wait_until_replaced(self) {
let mut interval = time::interval_at(Instant::now() + CHECK_INTERVAL, CHECK_INTERVAL);
interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
loop {
interval.tick().await;
if self.was_replaced().await {
if self.was_replaced() {
return;
}
}
}
async fn was_replaced(&self) -> bool {
fs::metadata(&self.path).await.map_or(true, |metadata| {
ExecutableIdentity::from_metadata(&metadata) != self.identity
})
/// Reads the identity synchronously. This is a `stat` of a page-cached
/// inode once per second, so handing it to Tokio's blocking pool would
/// cost more than the call itself and would keep a pool thread resident
/// for the life of the server.
fn was_replaced(&self) -> bool {
!fs::metadata(&self.path).is_ok_and(|metadata| identity(&metadata) == self.identity)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct ExecutableIdentity {
len: u64,
modified: Option<SystemTime>,
#[cfg(unix)]
device: u64,
#[cfg(unix)]
inode: u64,
/// Identifies the file behind an executable path. Upgrades always swap a new
/// file into place — `fabro upgrade` renames over the old one and Homebrew
/// repoints a symlink — so the identity changes even though the path does not.
#[cfg(unix)]
type Identity = (u64, u64);
#[cfg(unix)]
fn identity(metadata: &Metadata) -> Identity {
use std::os::unix::fs::MetadataExt as _;
(metadata.dev(), metadata.ino())
}
impl ExecutableIdentity {
fn from_metadata(metadata: &std::fs::Metadata) -> Self {
#[cfg(unix)]
use std::os::unix::fs::MetadataExt as _;
#[cfg(not(unix))]
type Identity = (u64, Option<std::time::SystemTime>);
Self {
len: metadata.len(),
modified: metadata.modified().ok(),
#[cfg(unix)]
device: metadata.dev(),
#[cfg(unix)]
inode: metadata.ino(),
}
}
#[cfg(not(unix))]
fn identity(metadata: &Metadata) -> Identity {
(metadata.len(), metadata.modified().ok())
}
#[expect(
clippy::disallowed_methods,
reason = "MCP startup resolves its invoked executable through the process PATH so it can detect Homebrew symlink updates"
)]
async fn invoked_executable_path() -> io::Result<PathBuf> {
let invoked = env::args_os()
.next()
.map(PathBuf::from)
.ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "process argv[0] is unavailable"))?;
if invoked.components().count() > 1 {
return absolute_path(invoked);
}
if let Some(path) = env::var_os(EnvVars::PATH) {
for directory in env::split_paths(&path) {
let candidate = absolute_path(directory.join(&invoked))?;
if fs::metadata(&candidate)
.await
.is_ok_and(|metadata| is_executable_file(&metadata))
{
return Ok(candidate);
}
}
}
env::current_exe()
}
fn absolute_path(path: PathBuf) -> io::Result<PathBuf> {
if path.is_absolute() {
Ok(path)
} else {
env::current_dir().map(|cwd| cwd.join(path))
}
}
fn is_executable_file(metadata: &std::fs::Metadata) -> bool {
if !metadata.is_file() {
return false;
}
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt as _;
metadata.permissions().mode() & 0o111 != 0
}
#[cfg(not(unix))]
{
true
/// Resolves the executable path to watch.
///
/// `argv[0]` wins when it carries a directory, because it names the path the
/// host actually launched, symlink included. MCP hosts normally launch a bare
/// `fabro` found on `PATH`, which leaves `current_exe`: it reports the symlink
/// on macOS, and the Homebrew symlink's own target on Linux.
fn invoked_executable_path() -> io::Result<PathBuf> {
match env::args_os().next().map(PathBuf::from) {
Some(invoked) if invoked.components().count() > 1 => path::absolute(invoked),
_ => env::current_exe(),
}
}
#[cfg(test)]
mod tests {
use tokio::fs as async_fs;
use super::*;
#[tokio::test]
async fn unchanged_executable_is_current() {
let directory = tempfile::tempdir().expect("temp directory should exist");
let executable = directory.path().join("fabro");
fs::write(&executable, b"current")
async_fs::write(&executable, b"current")
.await
.expect("fixture executable should be written");
let monitor = ExecutableMonitor::new(executable).await.unwrap();
let monitor = ExecutableMonitor::new(executable).unwrap();
assert!(!monitor.was_replaced().await);
assert!(!monitor.was_replaced());
}
#[tokio::test]
@ -149,42 +106,34 @@ mod tests {
let directory = tempfile::tempdir().expect("temp directory should exist");
let executable = directory.path().join("fabro");
let replacement = directory.path().join("fabro-new");
fs::write(&executable, b"old")
async_fs::write(&executable, b"old")
.await
.expect("old fixture executable should be written");
fs::write(&replacement, b"new executable")
async_fs::write(&replacement, b"new executable")
.await
.expect("new fixture executable should be written");
let monitor = ExecutableMonitor::new(executable.clone()).await.unwrap();
let monitor = ExecutableMonitor::new(executable.clone()).unwrap();
fs::rename(&replacement, &executable)
async_fs::rename(&replacement, &executable)
.await
.expect("fixture executable should be replaced");
assert!(monitor.was_replaced().await);
assert!(monitor.was_replaced());
}
#[tokio::test]
async fn removed_executable_is_detected() {
let directory = tempfile::tempdir().expect("temp directory should exist");
let executable = directory.path().join("fabro");
fs::write(&executable, b"current")
async_fs::write(&executable, b"current")
.await
.expect("fixture executable should be written");
let monitor = ExecutableMonitor::new(executable.clone()).await.unwrap();
let monitor = ExecutableMonitor::new(executable.clone()).unwrap();
fs::remove_file(executable)
async_fs::remove_file(executable)
.await
.expect("fixture executable should be removed");
assert!(monitor.was_replaced().await);
}
#[test]
fn executable_check_rejects_directories() {
let directory = tempfile::tempdir().expect("temp directory should exist");
let metadata = std::fs::metadata(directory.path()).unwrap();
assert!(!is_executable_file(&metadata));
assert!(monitor.was_replaced());
}
}

View file

@ -11,7 +11,11 @@ use std::sync::Arc;
use anyhow::Result;
pub use config::{config_json, init_agent};
use fabro_client::Client;
pub use server::{McpServerExit, start};
pub use server::start;
/// The name this MCP server reports over the wire and registers under in agent
/// config files.
pub(crate) const SERVER_NAME: &str = "fabro";
pub type FabroClientFuture = Pin<Box<dyn Future<Output = Result<Client>> + Send>>;

View file

@ -1,5 +1,6 @@
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use anyhow::Result;
use fabro_tool::fabro_client::ClientBackend;
@ -12,10 +13,12 @@ use rmcp::transport::stdio;
use rmcp::{ErrorData, ServerHandler, serve_server, tool, tool_handler, tool_router};
use serde::Serialize;
use tokio::sync::OnceCell;
use tokio::time;
use tracing::warn;
use crate::FabroMcpServerSettings;
use crate::executable_monitor::ExecutableMonitor;
use crate::manifest_builder::McpRunManifestBuilder;
use crate::{FabroMcpServerSettings, SERVER_NAME};
#[derive(Clone)]
pub(crate) struct FabroMcpServer {
@ -25,45 +28,52 @@ pub(crate) struct FabroMcpServer {
tool_router: ToolRouter<Self>,
}
/// The reason a running MCP stdio server returned to its caller.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum McpServerExit {
/// The MCP service stopped without an executable replacement.
ServiceStopped,
/// The executable on disk changed while the MCP service was running.
ExecutableReplaced,
}
/// How long to wait for the MCP service to stop after an upgrade is detected.
/// Bounded because the transport closes by writing to a stdout the host may
/// already have stopped reading.
const SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5);
pub async fn start(settings: FabroMcpServerSettings) -> Result<McpServerExit> {
let executable_monitor = ExecutableMonitor::current().await.ok();
let server = FabroMcpServer::new(Arc::new(settings));
let service = serve_server(server, stdio()).await?;
let exit = if let Some(executable_monitor) = executable_monitor {
let cancellation = service.cancellation_token();
let mut service_wait = Box::pin(service.waiting());
tokio::select! {
result = &mut service_wait => {
result?;
McpServerExit::ServiceStopped
}
() = executable_monitor.wait_until_replaced() => {
cancellation.cancel();
service_wait.await?;
McpServerExit::ExecutableReplaced
}
pub async fn start(settings: FabroMcpServerSettings) -> Result<()> {
let monitor = match ExecutableMonitor::current() {
Ok(monitor) => Some(monitor),
Err(error) => {
warn!(
%error,
"Upgrade detection is unavailable; this MCP server will keep running after an \
upgrade replaces it"
);
None
}
} else {
service.waiting().await?;
McpServerExit::ServiceStopped
};
Ok(exit)
let service = serve_server(FabroMcpServer::new(Arc::new(settings)), stdio()).await?;
let Some(monitor) = monitor else {
service.waiting().await?;
return Ok(());
};
let cancellation = service.cancellation_token();
let mut service_wait = Box::pin(service.waiting());
tokio::select! {
result = &mut service_wait => {
result?;
}
() = monitor.wait_until_replaced() => {
// An upgrade replaced the executable, so stop serving and let the
// host reconnect to the new one. The CLI exits the process rather
// than returning, because Tokio's stdin worker stays blocked on a
// read that only the host can end.
cancellation.cancel();
let _ = time::timeout(SHUTDOWN_TIMEOUT, service_wait).await;
}
}
Ok(())
}
#[tool_handler(router = self.tool_router)]
impl ServerHandler for FabroMcpServer {
fn get_info(&self) -> ServerInfo {
ServerInfo::new(ServerCapabilities::builder().enable_tools().build())
.with_server_info(Implementation::new("fabro", FABRO_VERSION).with_title("Fabro"))
.with_server_info(Implementation::new(SERVER_NAME, FABRO_VERSION).with_title("Fabro"))
.with_instructions("Use these tools to create, inspect, control, wait for, and read events from Fabro workflow runs.")
}
}