Fix stale MCP servers after upgrades

This commit is contained in:
Bryan Helmkamp 2026-07-31 13:18:12 -04:00
parent 14fc5d446c
commit f868734f27
No known key found for this signature in database
7 changed files with 425 additions and 18 deletions

1
Cargo.lock generated
View file

@ -2838,6 +2838,7 @@ dependencies = [
"fabro-manifest",
"fabro-model",
"fabro-server",
"fabro-static",
"fabro-tool",
"fabro-types",
"fabro-util",

View file

@ -1,6 +1,8 @@
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;
@ -9,7 +11,15 @@ use crate::server_client;
pub(crate) async fn dispatch(ns: McpNamespace, base_ctx: &CommandContext) -> Result<()> {
match ns.command {
McpCommand::Start(args) => {
fabro_mcp_server::start(server_settings(base_ctx, &args.connection)?).await
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(())
}
McpCommand::Config(args) => {
let json = fabro_mcp_server::config_json(&config_settings(&args.connection))?;

View file

@ -8,16 +8,20 @@
)]
use std::collections::HashMap;
#[cfg(unix)]
use std::fs;
use std::io::{BufRead as _, Write as _};
use std::path::{Path, PathBuf};
use std::process::Stdio;
use std::process::{Command, Stdio};
use std::thread;
use std::time::{Duration, Instant};
use chrono::{DateTime, Duration as ChronoDuration, Utc};
use fabro_client::{AuthEntry, AuthStore, DevTokenEntry, OAuthEntry, StoredSubject};
use fabro_mcp::client::McpClient;
use fabro_mcp::config::{McpServerSettings, McpTransport};
use fabro_test::{fabro_json_snapshot, fabro_snapshot, test_context};
use fabro_types::RunId;
use fabro_types::{Graph, RunId, WorkflowSettings, test_support};
use httpmock::Method::{GET, POST};
use httpmock::MockServer;
@ -515,7 +519,7 @@ 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 = std::process::Command::new(&fixture.command[0]);
let mut cmd = Command::new(&fixture.command[0]);
cmd.args(&fixture.command[1..])
.env_clear()
.envs(&fixture.env)
@ -534,31 +538,96 @@ fn stdio_start_writes_only_json_rpc_to_stdout() {
let stdout = child.stdout.take().unwrap();
let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
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(std::time::Duration::from_secs(5))
.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!(
value["result"]["serverInfo"]["version"],
env!("CARGO_PKG_VERSION")
);
let _ = child.kill();
let _ = child.wait();
}
#[cfg(unix)]
#[test]
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");
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());
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();
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");
let deadline = Instant::now() + Duration::from_secs(5);
let status = loop {
if let Some(status) = child.try_wait().expect("MCP server should be polled") {
break status;
}
if Instant::now() >= deadline {
let _ = child.kill();
let _ = child.wait();
panic!("MCP server did not exit after its executable was replaced");
}
thread::sleep(Duration::from_millis(50));
};
assert!(status.success(), "MCP server should exit successfully");
}
#[tokio::test(flavor = "multi_thread")]
async fn stdio_startup_and_list_tools_is_fast() {
let context = test_context!();
let start = std::time::Instant::now();
let start = Instant::now();
let client = spawn_mcp_client(&context, &[]).await;
let tools = client.list_tools().await.unwrap();
assert_eq!(tools.len(), MCP_RUN_TOOL_NAMES.len());
assert!(start.elapsed() < std::time::Duration::from_secs(2));
assert!(start.elapsed() < Duration::from_secs(2));
client
.shutdown()
.await
@ -1602,12 +1671,21 @@ async fn mcp_get_resolves_selector_and_returns_summary_projection_and_questions(
let projection = server.mock(|when, then| {
when.method(GET)
.path(format!("/api/v1/runs/{run_id}/state"));
let mut body = run_projection_json(&run_id, &serde_json::json!({ "kind": "running" }));
body["spec"]["settings"]["run"]["model"] = serde_json::json!({
"provider": "openai",
"name": "gpt-5.6-sol",
"fallbacks": {
"gpt-5.6-sol": ["gpt-5.6-terra"]
},
"controls": {
"reasoning_effort": null,
"speed": null
}
});
then.status(200)
.header("Content-Type", "application/json")
.json_body(run_projection_json(
&run_id,
&serde_json::json!({ "kind": "running" }),
));
.json_body(body);
});
let questions = server.mock(|when, then| {
when.method(GET)
@ -1644,6 +1722,10 @@ async fn mcp_get_resolves_selector_and_returns_summary_projection_and_questions(
assert_eq!(get["summary"]["workflow_name"], "Simple");
assert_eq!(get["summary"]["workflow_slug"], "simple");
assert_eq!(get["projection"]["status"]["kind"], "running");
assert_eq!(
get["projection"]["spec"]["settings"]["run"]["model"]["fallbacks"]["gpt-5.6-sol"][0],
"gpt-5.6-terra"
);
assert_eq!(get["questions"][0]["id"], "q-1");
resolve.assert();
retrieve.assert();
@ -2001,6 +2083,82 @@ async fn mcp_events_filters_find_matches_beyond_first_page() {
.expect("MCP client should shut down");
}
#[tokio::test(flavor = "multi_thread")]
async fn mcp_events_decodes_run_created_with_model_keyed_fallbacks() {
let context = test_context!();
let server = MockServer::start();
let target_url = format!("{}/api/v1", server.base_url());
let target: fabro_client::ServerTarget = target_url.parse().unwrap();
seed_dev_token_auth(&context.home_dir, &target, TEST_DEV_TOKEN);
let run_id = unique_run_id();
let resolve = mock_resolved_run(&server, "nightly", &run_id);
let mut settings = serde_json::to_value(WorkflowSettings::default())
.expect("workflow settings should serialize");
settings["run"]["model"] = serde_json::json!({
"provider": "openai",
"name": "gpt-5.6-sol",
"fallbacks": {
"gpt-5.6-sol": ["gpt-5.6-terra"]
},
"controls": {
"reasoning_effort": null,
"speed": null
}
});
let event = serde_json::json!({
"seq": 1,
"id": "evt-created",
"ts": "2026-04-05T12:00:00Z",
"run_id": run_id,
"event": "run.created",
"properties": {
"settings": settings,
"graph": Graph::new("Remote Workflow"),
"labels": {},
"run_dir": "/tmp/run",
"source_directory": "/srv/repo",
"provenance": test_support::test_run_provenance()
},
"actor": null
});
let events = server.mock(|when, then| {
when.method(GET)
.path(format!("/api/v1/runs/{run_id}/events"))
.query_param_missing("limit");
then.status(200)
.header("Content-Type", "application/json")
.json_body(serde_json::json!({
"data": [event],
"meta": { "has_more": false }
}));
});
let client = spawn_mcp_client(&context, &["--server", &target_url]).await;
let result = call_tool_json(
&client,
"fabro_run_events",
serde_json::json!({
"run_id": "nightly",
"action": "search",
"query": "gpt-5.6-terra",
"first": 1
}),
)
.await;
assert_eq!(
result["events"][0]["event"]["properties"]["settings"]["run"]["model"]["fallbacks"]["gpt-5.6-sol"]
[0],
"gpt-5.6-terra"
);
resolve.assert();
events.assert();
client
.shutdown()
.await
.expect("MCP client should shut down");
}
#[tokio::test(flavor = "multi_thread")]
async fn mcp_events_requires_action_specific_inputs_before_auth() {
let context = test_context!();

View file

@ -21,6 +21,7 @@ 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" }
@ -35,4 +36,4 @@ toml.workspace = true
[dev-dependencies]
httpmock = "0.8"
tempfile = "3"
tempfile = "3"

View file

@ -0,0 +1,190 @@
//! Detects when an upgrade replaces the executable that launched this MCP
//! server.
//!
//! MCP hosts can keep stdio servers alive for days. Without this check, an old
//! 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 fabro_static::EnvVars;
use tokio::time::Instant;
use tokio::{fs, time};
const CHECK_INTERVAL: Duration = Duration::from_secs(1);
pub(crate) struct ExecutableMonitor {
path: PathBuf,
identity: ExecutableIdentity,
}
impl ExecutableMonitor {
pub(crate) async fn current() -> io::Result<Self> {
let path = invoked_executable_path().await?;
Self::new(path).await
}
async fn new(path: PathBuf) -> io::Result<Self> {
let identity = ExecutableIdentity::from_metadata(&fs::metadata(&path).await?);
Ok(Self { path, identity })
}
pub(crate) async fn wait_until_replaced(self) {
let mut interval = time::interval_at(Instant::now() + CHECK_INTERVAL, CHECK_INTERVAL);
loop {
interval.tick().await;
if self.was_replaced().await {
return;
}
}
}
async fn was_replaced(&self) -> bool {
fs::metadata(&self.path).await.map_or(true, |metadata| {
ExecutableIdentity::from_metadata(&metadata) != self.identity
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct ExecutableIdentity {
len: u64,
modified: Option<SystemTime>,
#[cfg(unix)]
device: u64,
#[cfg(unix)]
inode: u64,
}
impl ExecutableIdentity {
fn from_metadata(metadata: &std::fs::Metadata) -> Self {
#[cfg(unix)]
use std::os::unix::fs::MetadataExt as _;
Self {
len: metadata.len(),
modified: metadata.modified().ok(),
#[cfg(unix)]
device: metadata.dev(),
#[cfg(unix)]
inode: metadata.ino(),
}
}
}
#[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
}
}
#[cfg(test)]
mod tests {
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")
.await
.expect("fixture executable should be written");
let monitor = ExecutableMonitor::new(executable).await.unwrap();
assert!(!monitor.was_replaced().await);
}
#[tokio::test]
async fn atomic_executable_replacement_is_detected() {
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")
.await
.expect("old fixture executable should be written");
fs::write(&replacement, b"new executable")
.await
.expect("new fixture executable should be written");
let monitor = ExecutableMonitor::new(executable.clone()).await.unwrap();
fs::rename(&replacement, &executable)
.await
.expect("fixture executable should be replaced");
assert!(monitor.was_replaced().await);
}
#[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")
.await
.expect("fixture executable should be written");
let monitor = ExecutableMonitor::new(executable.clone()).await.unwrap();
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));
}
}

View file

@ -1,4 +1,5 @@
mod config;
mod executable_monitor;
mod manifest_builder;
mod server;
@ -10,7 +11,7 @@ use std::sync::Arc;
use anyhow::Result;
pub use config::{config_json, init_agent};
use fabro_client::Client;
pub use server::start;
pub use server::{McpServerExit, start};
pub type FabroClientFuture = Pin<Box<dyn Future<Output = Result<Client>> + Send>>;

View file

@ -4,15 +4,17 @@ use std::sync::Arc;
use anyhow::Result;
use fabro_tool::fabro_client::ClientBackend;
use fabro_tool::{self as run_tools, FabroToolBackend};
use fabro_util::version::FABRO_VERSION;
use rmcp::handler::server::router::tool::ToolRouter;
use rmcp::handler::server::wrapper::Parameters;
use rmcp::model::{CallToolResult, Content, ServerCapabilities, ServerInfo};
use rmcp::model::{CallToolResult, Content, Implementation, ServerCapabilities, ServerInfo};
use rmcp::transport::stdio;
use rmcp::{ErrorData, ServerHandler, serve_server, tool, tool_handler, tool_router};
use serde::Serialize;
use tokio::sync::OnceCell;
use crate::FabroMcpServerSettings;
use crate::executable_monitor::ExecutableMonitor;
use crate::manifest_builder::McpRunManifestBuilder;
#[derive(Clone)]
@ -23,17 +25,45 @@ pub(crate) struct FabroMcpServer {
tool_router: ToolRouter<Self>,
}
pub async fn start(settings: FabroMcpServerSettings) -> Result<()> {
/// 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,
}
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?;
service.waiting().await?;
Ok(())
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
}
}
} else {
service.waiting().await?;
McpServerExit::ServiceStopped
};
Ok(exit)
}
#[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_instructions("Use these tools to create, inspect, control, wait for, and read events from Fabro workflow runs.")
}
}
@ -251,6 +281,22 @@ mod tests {
use super::*;
use crate::FabroMcpServerSettings;
#[test]
fn server_info_reports_fabro_version() {
let settings = FabroMcpServerSettings {
cwd: PathBuf::from("."),
config_path: PathBuf::from("fabro.toml"),
client_factory: Arc::new(|| {
Box::pin(async { panic!("client should not be constructed while reading info") })
}),
};
let info = FabroMcpServer::new(Arc::new(settings)).get_info();
assert_eq!(info.server_info.name, "fabro");
assert_eq!(info.server_info.title.as_deref(), Some("Fabro"));
assert_eq!(info.server_info.version, FABRO_VERSION);
}
#[test]
fn fabro_run_pair_tool_is_registered_with_stage_based_schema() {
let settings = FabroMcpServerSettings {