diff --git a/Cargo.lock b/Cargo.lock index a485c5cc4..5ae6ec28d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2838,6 +2838,7 @@ dependencies = [ "fabro-manifest", "fabro-model", "fabro-server", + "fabro-static", "fabro-tool", "fabro-types", "fabro-util", diff --git a/lib/apps/fabro-cli/src/commands/mcp/mod.rs b/lib/apps/fabro-cli/src/commands/mcp/mod.rs index f6a64a0ca..ef7233a66 100644 --- a/lib/apps/fabro-cli/src/commands/mcp/mod.rs +++ b/lib/apps/fabro-cli/src/commands/mcp/mod.rs @@ -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))?; diff --git a/lib/apps/fabro-cli/tests/it/cmd/mcp.rs b/lib/apps/fabro-cli/tests/it/cmd/mcp.rs index dc1052ac0..1d2ed68d2 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/mcp.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/mcp.rs @@ -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!(); diff --git a/lib/apps/fabro-mcp-server/Cargo.toml b/lib/apps/fabro-mcp-server/Cargo.toml index cdbf6b65f..da0b9978b 100644 --- a/lib/apps/fabro-mcp-server/Cargo.toml +++ b/lib/apps/fabro-mcp-server/Cargo.toml @@ -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" \ No newline at end of file +tempfile = "3" diff --git a/lib/apps/fabro-mcp-server/src/executable_monitor.rs b/lib/apps/fabro-mcp-server/src/executable_monitor.rs new file mode 100644 index 000000000..9ba92149b --- /dev/null +++ b/lib/apps/fabro-mcp-server/src/executable_monitor.rs @@ -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 { + let path = invoked_executable_path().await?; + Self::new(path).await + } + + async fn new(path: PathBuf) -> io::Result { + 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, + #[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 { + 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 { + 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)); + } +} diff --git a/lib/apps/fabro-mcp-server/src/lib.rs b/lib/apps/fabro-mcp-server/src/lib.rs index ad6264dcb..7aa598e80 100644 --- a/lib/apps/fabro-mcp-server/src/lib.rs +++ b/lib/apps/fabro-mcp-server/src/lib.rs @@ -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> + Send>>; diff --git a/lib/apps/fabro-mcp-server/src/server.rs b/lib/apps/fabro-mcp-server/src/server.rs index c1b932308..7fdc64d2e 100644 --- a/lib/apps/fabro-mcp-server/src/server.rs +++ b/lib/apps/fabro-mcp-server/src/server.rs @@ -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, } -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 { + 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 {