From 02eaeef78f4069ee965873304042df03e75f2d52 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 11 May 2026 12:57:04 -0400 Subject: [PATCH] refactor(mcp): simplify run tooling cleanup Reuse shared run status classification and manifest override construction across CLI, MCP, and server code. Normalize MCP run inputs during validation, avoid redundant run refetches after mutations, and bound/direct lookup paths where the API supports it. Verification: - cargo +nightly-2026-04-14 fmt --all - cargo check -p fabro-mcp-server -p fabro-client -p fabro-cli -p fabro-server -p fabro-manifest -p fabro-types - cargo nextest run -p fabro-mcp-server -p fabro-manifest -p fabro-types - cargo nextest run -p fabro-cli mcp - cargo +nightly-2026-04-14 clippy -p fabro-mcp-server -p fabro-client -p fabro-cli -p fabro-server -p fabro-manifest -p fabro-types --all-targets -- -D warnings --- Cargo.lock | 2 + lib/crates/fabro-cli/src/commands/mcp/mod.rs | 9 +- .../fabro-cli/src/commands/run/overrides.rs | 2 + .../fabro-cli/src/commands/run/start.rs | 2 +- .../fabro-cli/src/commands/runs/archive.rs | 2 +- lib/crates/fabro-cli/src/shared/utilities.rs | 13 +- lib/crates/fabro-cli/tests/it/cmd/mcp.rs | 233 ++++++++++-------- lib/crates/fabro-client/src/client.rs | 106 +++++--- lib/crates/fabro-manifest/src/lib.rs | 24 +- lib/crates/fabro-mcp-server/Cargo.toml | 1 + lib/crates/fabro-mcp-server/src/config.rs | 42 ++-- lib/crates/fabro-mcp-server/src/lib.rs | 5 +- .../fabro-mcp-server/src/run_tools/common.rs | 13 +- .../fabro-mcp-server/src/run_tools/create.rs | 84 +++++-- .../fabro-mcp-server/src/run_tools/events.rs | 85 ++++--- .../fabro-mcp-server/src/run_tools/gather.rs | 25 +- .../src/run_tools/interact.rs | 216 ++++++++++------ .../src/run_tools/manifest.rs | 37 +-- .../fabro-mcp-server/src/run_tools/search.rs | 66 ++++- lib/crates/fabro-mcp-server/src/server.rs | 12 +- lib/crates/fabro-server/Cargo.toml | 1 + lib/crates/fabro-server/src/run_manifest.rs | 57 +---- lib/crates/fabro-types/src/lib.rs | 3 +- lib/crates/fabro-types/src/status.rs | 50 ++++ 24 files changed, 678 insertions(+), 412 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 75b55d988..aefa4a469 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2054,6 +2054,7 @@ dependencies = [ "fabro-server", "fabro-types", "fabro-util", + "futures", "rmcp", "schemars 1.2.1", "serde", @@ -2198,6 +2199,7 @@ dependencies = [ "fabro-interview", "fabro-llm", "fabro-macros", + "fabro-manifest", "fabro-model", "fabro-proc", "fabro-redact", diff --git a/lib/crates/fabro-cli/src/commands/mcp/mod.rs b/lib/crates/fabro-cli/src/commands/mcp/mod.rs index a18067ada..6f047a73f 100644 --- a/lib/crates/fabro-cli/src/commands/mcp/mod.rs +++ b/lib/crates/fabro-cli/src/commands/mcp/mod.rs @@ -31,18 +31,11 @@ fn server_settings( let server_target = user_config::resolve_nondefault_server_target( &connection.target, connection_ctx.user_settings(), - )? - .map(|target| { - target - .as_unix_socket_path() - .map_or_else(|| target.to_string(), |path| path.display().to_string()) - }); + )?; Ok(fabro_mcp_server::FabroMcpServerSettings { - config: config_settings(connection), server_target, storage_dir: connection_ctx.storage_dir().to_path_buf(), config_path: connection_ctx.base_config_path().to_path_buf(), - home_dir: home_dir()?, cwd: base_ctx.cwd().to_path_buf(), }) } diff --git a/lib/crates/fabro-cli/src/commands/run/overrides.rs b/lib/crates/fabro-cli/src/commands/run/overrides.rs index b61371ace..d4de0183f 100644 --- a/lib/crates/fabro-cli/src/commands/run/overrides.rs +++ b/lib/crates/fabro-cli/src/commands/run/overrides.rs @@ -84,6 +84,7 @@ pub(crate) fn run_args_overrides(args: &RunArgs) -> Result Result Result<()> { - client.start_run(run_id, resume).await + client.start_run(run_id, resume).await.map(|_| ()) } diff --git a/lib/crates/fabro-cli/src/commands/runs/archive.rs b/lib/crates/fabro-cli/src/commands/runs/archive.rs index 9694bd829..38a5308fe 100644 --- a/lib/crates/fabro-cli/src/commands/runs/archive.rs +++ b/lib/crates/fabro-cli/src/commands/runs/archive.rs @@ -71,7 +71,7 @@ async fn run_bulk(action: Action, identifiers: &[String], ctx: &CommandContext) Action::Unarchive => client.unarchive_run(&run_id).await, }; match result { - Ok(()) => { + Ok(_) => { let run_id_string = run_id.to_string(); changed.push(run_id_string.clone()); if !json { diff --git a/lib/crates/fabro-cli/src/shared/utilities.rs b/lib/crates/fabro-cli/src/shared/utilities.rs index 433886932..7f3c9f20f 100644 --- a/lib/crates/fabro-cli/src/shared/utilities.rs +++ b/lib/crates/fabro-cli/src/shared/utilities.rs @@ -127,18 +127,7 @@ pub(crate) fn color_if(use_color: bool, color: Color) -> Option { } pub(crate) fn run_status_kind(status: RunStatus) -> &'static str { - match status { - RunStatus::Submitted => "submitted", - RunStatus::Queued => "queued", - RunStatus::Starting => "starting", - RunStatus::Running => "running", - RunStatus::Blocked { .. } => "blocked", - RunStatus::Paused { .. } => "paused", - RunStatus::Removing => "removing", - RunStatus::Succeeded { .. } => "succeeded", - RunStatus::Failed { .. } => "failed", - RunStatus::Dead => "dead", - } + status.kind().into() } pub(crate) fn split_run_path(s: &str) -> Option<(&str, &str)> { diff --git a/lib/crates/fabro-cli/tests/it/cmd/mcp.rs b/lib/crates/fabro-cli/tests/it/cmd/mcp.rs index 79b37aaf8..309141fb2 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/mcp.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/mcp.rs @@ -229,7 +229,7 @@ fn init_cursor_writes_idempotent_config() { } #[test] -fn init_claude_writes_platform_config() { +fn init_claude_writes_desktop_and_code_configs() { let context = test_context!(); context .command() @@ -237,10 +237,11 @@ fn init_claude_writes_platform_config() { .assert() .success(); - let config_path = expected_claude_config_path(&context.home_dir); - let config: serde_json::Value = - serde_json::from_str(&std::fs::read_to_string(config_path).unwrap()).unwrap(); - fabro_json_snapshot!(context, config, @r#" + let desktop_config: serde_json::Value = serde_json::from_str( + &std::fs::read_to_string(expected_claude_desktop_config_path(&context.home_dir)).unwrap(), + ) + .unwrap(); + fabro_json_snapshot!(context, desktop_config, @r#" { "mcpServers": { "fabro": { @@ -253,6 +254,62 @@ fn init_claude_writes_platform_config() { } } "#); + + let code_config: serde_json::Value = serde_json::from_str( + &std::fs::read_to_string(context.home_dir.join(".claude.json")).unwrap(), + ) + .unwrap(); + fabro_json_snapshot!(context, code_config, @r#" + { + "mcpServers": { + "fabro": { + "command": "fabro", + "args": [ + "mcp", + "start" + ] + } + } + } + "#); +} + +#[test] +fn init_claude_preserves_existing_claude_code_config() { + let context = test_context!(); + let claude_code_path = context.home_dir.join(".claude.json"); + std::fs::write( + &claude_code_path, + r#"{"numStartups":42,"mcpServers":{"other":{"type":"http","url":"https://example.test/mcp"}}}"#, + ) + .unwrap(); + + context + .command() + .args(["mcp", "init", "claude"]) + .assert() + .success(); + + let config: serde_json::Value = + serde_json::from_str(&std::fs::read_to_string(&claude_code_path).unwrap()).unwrap(); + fabro_json_snapshot!(context, config, @r#" + { + "numStartups": 42, + "mcpServers": { + "other": { + "type": "http", + "url": "https://example.test/mcp" + }, + "fabro": { + "command": "fabro", + "args": [ + "mcp", + "start" + ] + } + } + } + "#); } #[test] @@ -625,19 +682,8 @@ async fn mcp_search_includes_archived_runs_by_default() { ); archived["lifecycle"]["archived"] = serde_json::json!(true); archived["lifecycle"]["archived_at"] = serde_json::json!("2026-04-05T12:02:00Z"); - let list_runs = server.mock(|when, then| { - when.method(GET) - .path("/api/v1/runs") - .query_param("include_archived", "true") - .query_param("page[limit]", "100") - .query_param("page[offset]", "0"); - then.status(200) - .header("Content-Type", "application/json") - .json_body(serde_json::json!({ - "data": [active, archived], - "meta": { "has_more": false } - })); - }); + let active_resolve = mock_resolved_run_json(&server, &active_id, active, None); + let archived_resolve = mock_resolved_run_json(&server, &archived_id, archived, None); let client = spawn_mcp_client(&context, &["--server", &target_url]).await; let result = call_tool_json( @@ -655,7 +701,8 @@ async fn mcp_search_includes_archived_runs_by_default() { .iter() .any(|run| run["archived"] == true) ); - list_runs.assert(); + active_resolve.assert(); + archived_resolve.assert(); client .shutdown() .await @@ -677,7 +724,8 @@ async fn mcp_search_refreshes_expired_oauth_token() { let run_id = unique_run_id(); let expired_access = server.mock(|when, then| { when.method(GET) - .path("/api/v1/runs") + .path("/api/v1/runs/resolve") + .query_param("selector", run_id.clone()) .header("authorization", "Bearer expired-access"); then.status(401) .header("Content-Type", "application/json") @@ -710,24 +758,19 @@ async fn mcp_search_refreshes_expired_oauth_token() { }); let fresh_access = server.mock(|when, then| { when.method(GET) - .path("/api/v1/runs") + .path("/api/v1/runs/resolve") .header("authorization", "Bearer fresh-access") - .query_param("include_archived", "true") - .query_param("page[limit]", "100") - .query_param("page[offset]", "0"); + .query_param("selector", run_id.clone()); then.status(200) .header("Content-Type", "application/json") - .json_body(serde_json::json!({ - "data": [remote_run_summary_json( - &run_id, - "Simple", - "simple", - "OAuth refreshed", - &serde_json::json!({ "kind": "submitted" }), - "2026-04-05T12:00:00Z", - )], - "meta": { "has_more": false } - })); + .json_body(remote_run_summary_json( + &run_id, + "Simple", + "simple", + "OAuth refreshed", + &serde_json::json!({ "kind": "submitted" }), + "2026-04-05T12:00:00Z", + )); }); let client = spawn_mcp_client(&context, &["--server", &target_url]).await; @@ -773,27 +816,20 @@ async fn mcp_search_uses_fabro_auth_file_override() { ) .expect("custom auth store should be seeded"); let run_id = unique_run_id(); - let list_runs = server.mock(|when, then| { - when.method(GET) - .path("/api/v1/runs") - .header("authorization", format!("Bearer {TEST_DEV_TOKEN}")) - .query_param("include_archived", "true") - .query_param("page[limit]", "100") - .query_param("page[offset]", "0"); - then.status(200) - .header("Content-Type", "application/json") - .json_body(serde_json::json!({ - "data": [remote_run_summary_json( - &run_id, - "Simple", - "simple", - "Custom auth file", - &serde_json::json!({ "kind": "submitted" }), - "2026-04-05T12:00:00Z", - )], - "meta": { "has_more": false } - })); - }); + let authorization = format!("Bearer {TEST_DEV_TOKEN}"); + let resolve = mock_resolved_run_json( + &server, + &run_id, + remote_run_summary_json( + &run_id, + "Simple", + "simple", + "Custom auth file", + &serde_json::json!({ "kind": "submitted" }), + "2026-04-05T12:00:00Z", + ), + Some(&authorization), + ); let mut fixture = mcp_stdio_fixture(&context, &["--server", &target_url]); fixture.env.insert( "FABRO_AUTH_FILE".to_string(), @@ -809,7 +845,7 @@ async fn mcp_search_uses_fabro_auth_file_override() { .await; assert_eq!(result["runs"][0]["run_id"], run_id); - list_runs.assert(); + resolve.assert(); client .shutdown() .await @@ -842,19 +878,8 @@ async fn mcp_search_orders_by_started_timestamp_before_created_timestamp() { "2026-04-05T12:00:00Z", ); running["timestamps"]["started_at"] = serde_json::json!("2026-04-05T12:20:00Z"); - let list_runs = server.mock(|when, then| { - when.method(GET) - .path("/api/v1/runs") - .query_param("include_archived", "true") - .query_param("page[limit]", "100") - .query_param("page[offset]", "0"); - then.status(200) - .header("Content-Type", "application/json") - .json_body(serde_json::json!({ - "data": [submitted, running], - "meta": { "has_more": false } - })); - }); + let submitted_resolve = mock_resolved_run_json(&server, &submitted_id, submitted, None); + let running_resolve = mock_resolved_run_json(&server, &running_id, running, None); let client = spawn_mcp_client(&context, &["--server", &target_url]).await; let result = call_tool_json( @@ -866,7 +891,8 @@ async fn mcp_search_orders_by_started_timestamp_before_created_timestamp() { assert_eq!(result["runs"][0]["run_id"], running_id); assert_eq!(result["runs"][1]["run_id"], submitted_id); - list_runs.assert(); + submitted_resolve.assert(); + running_resolve.assert(); client .shutdown() .await @@ -900,19 +926,8 @@ async fn mcp_search_orders_submitted_runs_by_created_timestamp_not_run_id_timest "2026-04-05T12:10:00Z", ); older_created["timestamps"]["started_at"] = serde_json::Value::Null; - let list_runs = server.mock(|when, then| { - when.method(GET) - .path("/api/v1/runs") - .query_param("include_archived", "true") - .query_param("page[limit]", "100") - .query_param("page[offset]", "0"); - then.status(200) - .header("Content-Type", "application/json") - .json_body(serde_json::json!({ - "data": [newer_created, older_created], - "meta": { "has_more": false } - })); - }); + let newer_resolve = mock_resolved_run_json(&server, &newer_created_id, newer_created, None); + let older_resolve = mock_resolved_run_json(&server, &older_created_id, older_created, None); let client = spawn_mcp_client(&context, &["--server", &target_url]).await; let result = call_tool_json( @@ -924,7 +939,8 @@ async fn mcp_search_orders_submitted_runs_by_created_timestamp_not_run_id_timest assert_eq!(result["runs"][0]["run_id"], newer_created_id); assert_eq!(result["runs"][1]["run_id"], older_created_id); - list_runs.assert(); + newer_resolve.assert(); + older_resolve.assert(); client .shutdown() .await @@ -1277,7 +1293,7 @@ async fn mcp_interact_actions_resolve_selector_and_call_expected_endpoints() { assert_eq!(message_result["result"]["interrupt"], true); assert_eq!(cancel_result["result"]["summary"]["run_id"], run_id); resolve.assert_calls(4); - retrieve.assert_calls(3); + retrieve.assert_calls(1); projection.assert(); start.assert(); message.assert(); @@ -1821,30 +1837,18 @@ async fn mcp_events_offset_beyond_fetch_cap_reaches_later_pages() { }) }) .collect::>(); - let first_200_events = events.iter().take(200).cloned().collect::>(); - let capped_events = server.mock(|when, then| { + let first_251_events = events.iter().take(251).cloned().collect::>(); + let bounded_events = server.mock(|when, then| { when.method(GET) .path(format!("/api/v1/runs/{run_id}/events")) - .query_param("limit", "200"); + .query_param("limit", "251"); then.status(200) .header("Content-Type", "application/json") .json_body(serde_json::json!({ - "data": first_200_events, + "data": first_251_events, "meta": { "has_more": true } })); }); - let full_events = server.mock(|when, then| { - when.method(GET) - .path(format!("/api/v1/runs/{run_id}/events")) - .query_param_missing("limit") - .query_param_missing("since_seq"); - then.status(200) - .header("Content-Type", "application/json") - .json_body(serde_json::json!({ - "data": events, - "meta": { "has_more": false } - })); - }); let client = spawn_mcp_client(&context, &["--server", &target_url]).await; let paged = call_tool_json( @@ -1862,8 +1866,7 @@ async fn mcp_events_offset_beyond_fetch_cap_reaches_later_pages() { assert_eq!(paged["events"][0]["event_id"], "evt-251"); assert_eq!(paged["next_cursor"], 252); resolve.assert(); - capped_events.assert_calls(0); - full_events.assert(); + bounded_events.assert(); client .shutdown() .await @@ -1898,7 +1901,7 @@ async fn mcp_tool_auth_error_mentions_login() { harness.shutdown().await; } -fn expected_claude_config_path(home_dir: &Path) -> PathBuf { +fn expected_claude_desktop_config_path(home_dir: &Path) -> PathBuf { #[cfg(target_os = "macos")] { home_dir @@ -2120,3 +2123,23 @@ fn run_id_with_timestamp(timestamp: &str, sequence: u128) -> String { .with_timezone(&Utc); RunId::with_timestamp(timestamp, sequence).to_string() } + +fn mock_resolved_run_json<'a>( + server: &'a MockServer, + selector: &str, + body: serde_json::Value, + authorization: Option<&str>, +) -> httpmock::Mock<'a> { + server.mock(|when, then| { + let when = when + .method(GET) + .path("/api/v1/runs/resolve") + .query_param("selector", selector); + if let Some(authorization) = authorization { + when.header("authorization", authorization); + } + then.status(200) + .header("Content-Type", "application/json") + .json_body(body); + }) +} diff --git a/lib/crates/fabro-client/src/client.rs b/lib/crates/fabro-client/src/client.rs index 97627466b..3654fd740 100644 --- a/lib/crates/fabro-client/src/client.rs +++ b/lib/crates/fabro-client/src/client.rs @@ -794,25 +794,27 @@ impl Client { Ok(bytes) } - pub async fn start_run(&self, run_id: &RunId, resume: bool) -> Result<()> { - self.send_api(|client| async move { - client - .start_run() - .id(run_id.to_string()) - .body(types::StartRunRequest { resume }) - .send() - .await - }) - .await?; - Ok(()) + pub async fn start_run(&self, run_id: &RunId, resume: bool) -> Result { + let response = self + .send_api(|client| async move { + client + .start_run() + .id(run_id.to_string()) + .body(types::StartRunRequest { resume }) + .send() + .await + }) + .await?; + convert_type(response.into_inner()) } - pub async fn cancel_run(&self, run_id: &RunId) -> Result<()> { - self.send_api( - |client| async move { client.cancel_run().id(run_id.to_string()).send().await }, - ) - .await?; - Ok(()) + pub async fn cancel_run(&self, run_id: &RunId) -> Result { + let response = self + .send_api( + |client| async move { client.cancel_run().id(run_id.to_string()).send().await }, + ) + .await?; + convert_type(response.into_inner()) } pub async fn interrupt_run(&self, run_id: &RunId) -> Result<()> { @@ -844,20 +846,22 @@ impl Client { Ok(()) } - pub async fn archive_run(&self, run_id: &RunId) -> Result<()> { - self.send_api( - |client| async move { client.archive_run().id(run_id.to_string()).send().await }, - ) - .await?; - Ok(()) + pub async fn archive_run(&self, run_id: &RunId) -> Result { + let response = self + .send_api( + |client| async move { client.archive_run().id(run_id.to_string()).send().await }, + ) + .await?; + convert_type(response.into_inner()) } - pub async fn unarchive_run(&self, run_id: &RunId) -> Result<()> { - self.send_api(|client| async move { - client.unarchive_run().id(run_id.to_string()).send().await - }) - .await?; - Ok(()) + pub async fn unarchive_run(&self, run_id: &RunId) -> Result { + let response = self + .send_api( + |client| async move { client.unarchive_run().id(run_id.to_string()).send().await }, + ) + .await?; + convert_type(response.into_inner()) } pub async fn rewind_run( @@ -1123,6 +1127,50 @@ impl Client { Ok(all_events) } + pub async fn list_run_events_until( + &self, + run_id: &RunId, + since_seq: Option, + max_events: usize, + ) -> Result> { + if max_events == 0 { + return Ok(Vec::new()); + } + + let mut next_since_seq = since_seq; + let mut all_events = Vec::new(); + while all_events.len() < max_events { + let remaining = max_events - all_events.len(); + let response = self + .send_api(|client| async move { + let mut request = client + .list_run_events() + .id(run_id.to_string()) + .limit(remaining.min(1000) as u64); + if let Some(seq) = next_since_seq.and_then(non_zero_u64_from_u32) { + request = request.since_seq(seq); + } + request.send().await + }) + .await?; + let parsed = response.into_inner(); + let page_events = parsed + .data + .into_iter() + .map(convert_type::<_, EventEnvelope>) + .collect::>>()?; + let next_page_since_seq = page_events.last().map(|event| event.seq.saturating_add(1)); + all_events.extend(page_events); + + if !parsed.meta.has_more || next_page_since_seq.is_none() { + break; + } + next_since_seq = next_page_since_seq; + } + + Ok(all_events) + } + pub async fn attach_run_events( &self, run_id: &RunId, diff --git a/lib/crates/fabro-manifest/src/lib.rs b/lib/crates/fabro-manifest/src/lib.rs index 7cd5ebb17..6d6944ffc 100644 --- a/lib/crates/fabro-manifest/src/lib.rs +++ b/lib/crates/fabro-manifest/src/lib.rs @@ -11,8 +11,8 @@ use fabro_api::types; use fabro_config::project::{self, discover_project_config, resolve_workflow_path}; use fabro_config::run::{resolve_run_goal_from_layer, resolve_run_goal_from_namespace}; use fabro_config::{ - CliLayer, DaytonaDockerfileLayer, ReplaceMap, RunExecutionLayer, RunGoalLayer, RunLayer, - RunModelLayer, RunSandboxLayer, WorkflowSettingsBuilder, + CliLayer, DaytonaDockerfileLayer, DockerSandboxLayer, ReplaceMap, RunExecutionLayer, + RunGoalLayer, RunLayer, RunModelLayer, RunSandboxLayer, WorkflowSettingsBuilder, }; use fabro_graphviz::graph::AttrValue; use fabro_graphviz::parser; @@ -51,6 +51,7 @@ pub struct RunOverrideInput<'a> { pub model: Option<&'a str>, pub provider: Option<&'a str>, pub sandbox: Option<&'a str>, + pub docker_image: Option<&'a str>, pub preserve_sandbox: Option, pub dry_run: Option, pub auto_approve: Option, @@ -67,12 +68,18 @@ pub fn build_run_overrides(input: RunOverrideInput<'_>) -> RunLayer { name: input.model.map(InterpString::parse), fallbacks: Vec::new(), }); - let sandbox = - (input.sandbox.is_some() || input.preserve_sandbox.is_some()).then(|| RunSandboxLayer { - provider: input.sandbox.map(ToOwned::to_owned), - preserve: input.preserve_sandbox, - ..RunSandboxLayer::default() - }); + let sandbox = (input.sandbox.is_some() + || input.docker_image.is_some() + || input.preserve_sandbox.is_some()) + .then(|| RunSandboxLayer { + provider: input.sandbox.map(ToOwned::to_owned), + docker: input.docker_image.map(|image| DockerSandboxLayer { + image: Some(image.to_string()), + ..DockerSandboxLayer::default() + }), + preserve: input.preserve_sandbox, + ..RunSandboxLayer::default() + }); let execution = (input.dry_run.is_some() || input.auto_approve.is_some()).then(|| RunExecutionLayer { mode: input.dry_run.map(|dry_run| { @@ -707,6 +714,7 @@ mod tests { model: Some("gpt-5.4-mini"), provider: Some("openai"), sandbox: Some("local"), + docker_image: None, preserve_sandbox: Some(true), dry_run: Some(true), auto_approve: Some(false), diff --git a/lib/crates/fabro-mcp-server/Cargo.toml b/lib/crates/fabro-mcp-server/Cargo.toml index fc40b6047..c10f1e0e7 100644 --- a/lib/crates/fabro-mcp-server/Cargo.toml +++ b/lib/crates/fabro-mcp-server/Cargo.toml @@ -24,6 +24,7 @@ fabro-config = { path = "../fabro-config" } fabro-server = { path = "../fabro-server" } fabro-types = { path = "../fabro-types" } fabro-util = { path = "../fabro-util" } +futures.workspace = true rmcp = { workspace = true, features = ["server", "macros", "schemars", "transport-io"] } schemars = "1.2.1" serde.workspace = true diff --git a/lib/crates/fabro-mcp-server/src/config.rs b/lib/crates/fabro-mcp-server/src/config.rs index f0318ae7f..9eb8a67c2 100644 --- a/lib/crates/fabro-mcp-server/src/config.rs +++ b/lib/crates/fabro-mcp-server/src/config.rs @@ -20,9 +20,10 @@ pub fn config_json(settings: &McpConfigSettings) -> Result { } pub fn init_agent(settings: &McpInitSettings) -> Result<()> { - let path = agent_config_path(settings.agent, &settings.home_dir); let entry = server_entry(&settings.config); - merge_server_entry(&path, entry)?; + for path in agent_config_paths(settings.agent, &settings.home_dir) { + merge_server_entry(&path, entry.clone())?; + } Ok(()) } @@ -60,13 +61,11 @@ fn merge_server_entry(path: &Path, entry: Value) -> Result<()> { .with_context(|| format!("failed to create {}", parent.display()))?; } - let mut root = if path.exists() { - let contents = std::fs::read_to_string(path) - .with_context(|| format!("failed to read {}", path.display()))?; - serde_json::from_str::(&contents) - .with_context(|| format!("failed to parse MCP config {}", path.display()))? - } else { - Value::Object(Map::new()) + let mut root = match std::fs::read_to_string(path) { + Ok(contents) => serde_json::from_str::(&contents) + .with_context(|| format!("failed to parse MCP config {}", path.display()))?, + Err(err) if err.kind() == std::io::ErrorKind::NotFound => Value::Object(Map::new()), + Err(err) => return Err(err).with_context(|| format!("failed to read {}", path.display())), }; let root_object = root @@ -91,18 +90,27 @@ fn merge_server_entry(path: &Path, entry: Value) -> Result<()> { std::fs::write(path, rendered).with_context(|| format!("failed to write {}", path.display())) } -fn agent_config_path(agent: McpAgent, home_dir: &Path) -> PathBuf { +fn agent_config_paths(agent: McpAgent, home_dir: &Path) -> Vec { match agent { - McpAgent::Claude => claude_config_path(home_dir), - McpAgent::Cursor => home_dir.join(".cursor").join("mcp.json"), - McpAgent::Windsurf => home_dir - .join(".codeium") - .join("windsurf") - .join("mcp_config.json"), + McpAgent::Claude => vec![ + claude_desktop_config_path(home_dir), + claude_code_config_path(home_dir), + ], + McpAgent::Cursor => vec![home_dir.join(".cursor").join("mcp.json")], + McpAgent::Windsurf => vec![ + home_dir + .join(".codeium") + .join("windsurf") + .join("mcp_config.json"), + ], } } -fn claude_config_path(home_dir: &Path) -> PathBuf { +fn claude_code_config_path(home_dir: &Path) -> PathBuf { + home_dir.join(".claude.json") +} + +fn claude_desktop_config_path(home_dir: &Path) -> PathBuf { #[cfg(target_os = "macos")] { home_dir diff --git a/lib/crates/fabro-mcp-server/src/lib.rs b/lib/crates/fabro-mcp-server/src/lib.rs index db94592c0..4c2ed94d3 100644 --- a/lib/crates/fabro-mcp-server/src/lib.rs +++ b/lib/crates/fabro-mcp-server/src/lib.rs @@ -5,15 +5,14 @@ mod server; use std::path::PathBuf; pub use config::{config_json, init_agent}; +use fabro_client::ServerTarget; pub use server::start; #[derive(Debug, Clone)] pub struct FabroMcpServerSettings { - pub config: McpConfigSettings, - pub server_target: Option, + pub server_target: Option, pub storage_dir: PathBuf, pub config_path: PathBuf, - pub home_dir: PathBuf, pub cwd: PathBuf, } diff --git a/lib/crates/fabro-mcp-server/src/run_tools/common.rs b/lib/crates/fabro-mcp-server/src/run_tools/common.rs index 8fb1d27bf..d708ed3ad 100644 --- a/lib/crates/fabro-mcp-server/src/run_tools/common.rs +++ b/lib/crates/fabro-mcp-server/src/run_tools/common.rs @@ -127,18 +127,7 @@ pub(super) fn parse_datetime_filter(name: &str, raw: &str) -> ToolResult &'static str { - match status { - RunStatus::Submitted => "submitted", - RunStatus::Queued => "queued", - RunStatus::Starting => "starting", - RunStatus::Running => "running", - RunStatus::Blocked { .. } => "blocked", - RunStatus::Paused { .. } => "paused", - RunStatus::Removing => "removing", - RunStatus::Succeeded { .. } => "succeeded", - RunStatus::Failed { .. } => "failed", - RunStatus::Dead => "dead", - } + status.kind().into() } fn format_tool_error(err: &anyhow::Error) -> String { diff --git a/lib/crates/fabro-mcp-server/src/run_tools/create.rs b/lib/crates/fabro-mcp-server/src/run_tools/create.rs index e27cda9cc..4d5b48097 100644 --- a/lib/crates/fabro-mcp-server/src/run_tools/create.rs +++ b/lib/crates/fabro-mcp-server/src/run_tools/create.rs @@ -3,6 +3,7 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use fabro_client::Client; +use fabro_types::RunId; use schemars::JsonSchema; use serde::{Deserialize, Serialize}; use serde_json::Value; @@ -36,7 +37,24 @@ pub(crate) struct CreateRunSpec { #[derive(Debug)] pub(crate) struct ValidatedCreateRuns { - pub(crate) runs: Vec, + pub(crate) runs: Vec, +} + +#[derive(Debug)] +pub(crate) struct ValidatedCreateRunSpec { + pub(crate) workflow: String, + pub(crate) cwd: Option, + pub(crate) run_id: Option, + pub(crate) goal: Option, + pub(crate) inputs: HashMap, + pub(crate) labels: HashMap, + pub(crate) dry_run: Option, + pub(crate) auto_approve: Option, + pub(crate) model: Option, + pub(crate) provider: Option, + pub(crate) sandbox: Option, + pub(crate) preserve_sandbox: Option, + pub(crate) start: Option, } impl TryFrom for ValidatedCreateRuns { @@ -44,12 +62,49 @@ impl TryFrom for ValidatedCreateRuns { fn try_from(params: FabroRunCreateParams) -> Result { common::validate_len("runs", params.runs.len(), 1, 50)?; - for spec in ¶ms.runs { - for (key, value) in &spec.inputs { - manifest::json_to_toml_value(key, value)?; - } - } - Ok(Self { runs: params.runs }) + let runs = params + .runs + .into_iter() + .map(ValidatedCreateRunSpec::try_from) + .collect::, _>>()?; + Ok(Self { runs }) + } +} + +impl TryFrom for ValidatedCreateRunSpec { + type Error = ToolError; + + fn try_from(spec: CreateRunSpec) -> Result { + let run_id = spec + .run_id + .as_deref() + .map(str::parse::) + .transpose() + .map_err(|err| { + ToolError::message(format!("run_id must be a valid Fabro run id: {err}")) + })?; + let inputs = spec + .inputs + .iter() + .map(|(key, value)| { + manifest::json_to_toml_value(key, value).map(|value| (key.clone(), value)) + }) + .collect::>>()?; + Ok(Self { + workflow: spec.workflow, + cwd: spec.cwd, + run_id, + goal: spec.goal, + inputs, + labels: spec.labels, + dry_run: spec.dry_run, + auto_approve: spec.auto_approve, + model: spec.model, + provider: spec.provider, + sandbox: spec.sandbox, + preserve_sandbox: spec.preserve_sandbox, + start: spec.start, + }) } } @@ -81,16 +136,17 @@ pub(crate) async fn create_runs( .await .map_err(|err| ToolError::from_anyhow(&err))?; let started = spec.start.unwrap_or(true); - if started { + let summary = if started { client .start_run(&run_id, false) .await - .map_err(|err| ToolError::from_anyhow(&err))?; - } - let summary = client - .retrieve_run(&run_id) - .await - .map_err(|err| ToolError::from_anyhow(&err))?; + .map_err(|err| ToolError::from_anyhow(&err))? + } else { + client + .retrieve_run(&run_id) + .await + .map_err(|err| ToolError::from_anyhow(&err))? + }; created.push(CreatedRunResult { run_id: summary.id.to_string(), workflow: spec.workflow, diff --git a/lib/crates/fabro-mcp-server/src/run_tools/events.rs b/lib/crates/fabro-mcp-server/src/run_tools/events.rs index a49d2c98f..42b7c44cc 100644 --- a/lib/crates/fabro-mcp-server/src/run_tools/events.rs +++ b/lib/crates/fabro-mcp-server/src/run_tools/events.rs @@ -1,5 +1,6 @@ use std::sync::Arc; +use chrono::{DateTime, Utc}; use fabro_client::Client; use fabro_types::EventEnvelope; use schemars::JsonSchema; @@ -37,7 +38,11 @@ pub(crate) struct FabroRunEventsParams { #[derive(Debug)] pub(crate) struct ValidatedRunEvents { - pub(crate) raw: FabroRunEventsParams, + pub(crate) raw: FabroRunEventsParams, + pub(crate) descending: bool, + pub(crate) first: usize, + pub(crate) created_after: Option>, + pub(crate) created_before: Option>, } impl TryFrom for ValidatedRunEvents { @@ -51,17 +56,21 @@ impl TryFrom for ValidatedRunEvents { if first > 200 { return Err(ToolError::message("first must be <= 200")); } - if let Some(direction) = params.direction.as_deref() { - if !matches!(direction, "asc" | "desc") { - return Err(ToolError::message("direction must be `asc` or `desc`")); - } - } - if let Some(created_after) = params.created_after.as_deref() { - common::parse_datetime_filter("created_after", created_after)?; - } - if let Some(created_before) = params.created_before.as_deref() { - common::parse_datetime_filter("created_before", created_before)?; - } + let descending = match params.direction.as_deref() { + None | Some("asc") => false, + Some("desc") => true, + Some(_) => return Err(ToolError::message("direction must be `asc` or `desc`")), + }; + let created_after = params + .created_after + .as_deref() + .map(|created_after| common::parse_datetime_filter("created_after", created_after)) + .transpose()?; + let created_before = params + .created_before + .as_deref() + .map(|created_before| common::parse_datetime_filter("created_before", created_before)) + .transpose()?; if matches!(params.action, RunEventsAction::Details) && params.event_ids.as_ref().is_none_or(Vec::is_empty) { @@ -77,7 +86,13 @@ impl TryFrom for ValidatedRunEvents { { return Err(ToolError::message("query is required for search action")); } - Ok(Self { raw: params }) + Ok(Self { + raw: params, + descending, + first, + created_after, + created_before, + }) } } @@ -101,29 +116,35 @@ pub(crate) async fn run_events( client: Arc, params: ValidatedRunEvents, ) -> ToolResult { + let descending = params.descending; + let first = params.first; + let created_after = params.created_after; + let created_before = params.created_before; let raw = params.raw; let run_id = client .resolve_run(&raw.run_id) .await .map_err(|err| ToolError::from_anyhow(&err))? .id; - let descending = raw.direction.as_deref() == Some("desc"); let fetch_after = if descending { None } else { raw.after }; - let mut events = client - .list_run_events(&run_id, fetch_after, event_fetch_limit(&raw)) - .await - .map_err(|err| ToolError::from_anyhow(&err))?; + let mut events = if let Some(limit) = event_fetch_limit(&raw, first) { + client + .list_run_events_until(&run_id, fetch_after, limit) + .await + } else { + client.list_run_events(&run_id, fetch_after, None).await + } + .map_err(|err| ToolError::from_anyhow(&err))?; if descending { if let Some(after) = raw.after { events.retain(|event| event.seq < after); } } - filter_events(&mut events, &raw)?; + filter_events(&mut events, &raw, created_after, created_before); if descending { events.reverse(); } let offset = raw.offset.unwrap_or(0); - let first = raw.first.or(raw.limit).unwrap_or(50).min(200); let page = events .into_iter() .skip(offset) @@ -154,7 +175,7 @@ pub(crate) fn run_events_text(result: &RunEventsResult) -> String { format!("returned {} Fabro event(s)", result.events.len()) } -fn event_fetch_limit(params: &FabroRunEventsParams) -> Option { +fn event_fetch_limit(params: &FabroRunEventsParams, first: usize) -> Option { let needs_full_scan = params.event_ids.is_some() || params.event_types.is_some() || params.categories.is_some() @@ -169,15 +190,16 @@ fn event_fetch_limit(params: &FabroRunEventsParams) -> Option { return None; } - let requested = params - .first - .or(params.limit) - .unwrap_or(50) - .saturating_add(params.offset.unwrap_or(0)); - (requested <= 200).then_some(requested.max(1)) + let requested = first.saturating_add(params.offset.unwrap_or(0)); + Some(requested.max(1)) } -fn filter_events(events: &mut Vec, params: &FabroRunEventsParams) -> ToolResult<()> { +fn filter_events( + events: &mut Vec, + params: &FabroRunEventsParams, + created_after: Option>, + created_before: Option>, +) { if let Some(event_ids) = params.event_ids.as_ref() { events.retain(|event| event_ids.contains(&event.event.id)); } @@ -199,12 +221,10 @@ fn filter_events(events: &mut Vec, params: &FabroRunEventsParams) categories.iter().any(|candidate| candidate == category) }); } - if let Some(created_after) = params.created_after.as_deref() { - let cutoff = common::parse_datetime_filter("created_after", created_after)?; + if let Some(cutoff) = created_after { events.retain(|event| event.event.ts >= cutoff); } - if let Some(created_before) = params.created_before.as_deref() { - let cutoff = common::parse_datetime_filter("created_before", created_before)?; + if let Some(cutoff) = created_before { events.retain(|event| event.event.ts <= cutoff); } if matches!(params.action, RunEventsAction::Search) { @@ -214,7 +234,6 @@ fn filter_events(events: &mut Vec, params: &FabroRunEventsParams) }); } } - Ok(()) } fn run_event_result( diff --git a/lib/crates/fabro-mcp-server/src/run_tools/gather.rs b/lib/crates/fabro-mcp-server/src/run_tools/gather.rs index 2774475b7..48e59b33f 100644 --- a/lib/crates/fabro-mcp-server/src/run_tools/gather.rs +++ b/lib/crates/fabro-mcp-server/src/run_tools/gather.rs @@ -2,6 +2,7 @@ use std::sync::Arc; use std::time::{Duration, Instant}; use fabro_client::Client; +use futures::future::try_join_all; use schemars::JsonSchema; use serde::{Deserialize, Serialize}; use tokio::time; @@ -58,22 +59,24 @@ pub(crate) async fn gather_runs( ) -> ToolResult { let start = Instant::now(); let deadline = start + Duration::from_secs(params.timeout_seconds); - let mut run_ids = Vec::with_capacity(params.run_ids.len()); - for selector in params.run_ids { - run_ids.push( + let run_ids = try_join_all(params.run_ids.into_iter().map(|selector| { + let client = Arc::clone(&client); + async move { client .resolve_run(&selector) .await - .map_err(|err| ToolError::from_anyhow(&err))? - .id, - ); - } + .map(|run| run.id) + .map_err(|err| ToolError::from_anyhow(&err)) + } + })) + .await?; loop { - let mut summaries = Vec::with_capacity(run_ids.len()); - for run_id in &run_ids { - summaries.push(common::retrieve_run(&client, run_id).await?); - } + let summaries = try_join_all(run_ids.iter().map(|run_id| { + let client = Arc::clone(&client); + async move { common::retrieve_run(&client, run_id).await } + })) + .await?; if summaries .iter() .all(|run| run.lifecycle.status.is_terminal()) diff --git a/lib/crates/fabro-mcp-server/src/run_tools/interact.rs b/lib/crates/fabro-mcp-server/src/run_tools/interact.rs index 6ab6e9bc8..2ce94f8ef 100644 --- a/lib/crates/fabro-mcp-server/src/run_tools/interact.rs +++ b/lib/crates/fabro-mcp-server/src/run_tools/interact.rs @@ -35,7 +35,41 @@ pub(crate) struct FabroRunInteractParams { #[derive(Debug)] pub(crate) struct ValidatedInteractRun { - pub(crate) raw: FabroRunInteractParams, + pub(crate) run_id: String, + pub(crate) action: ValidatedInteractAction, +} + +#[derive(Debug)] +pub(crate) enum ValidatedInteractAction { + Get, + Start, + Message { + message: String, + interrupt: bool, + }, + Cancel, + Archive, + Unarchive, + GetQuestions, + Answer { + question_id: String, + body: types::SubmitAnswerRequest, + }, +} + +impl ValidatedInteractAction { + fn action(&self) -> RunInteractAction { + match self { + Self::Get => RunInteractAction::Get, + Self::Start => RunInteractAction::Start, + Self::Message { .. } => RunInteractAction::Message, + Self::Cancel => RunInteractAction::Cancel, + Self::Archive => RunInteractAction::Archive, + Self::Unarchive => RunInteractAction::Unarchive, + Self::GetQuestions => RunInteractAction::GetQuestions, + Self::Answer { .. } => RunInteractAction::Answer, + } + } } impl TryFrom for ValidatedInteractRun { @@ -45,26 +79,51 @@ impl TryFrom for ValidatedInteractRun { if params.run_id.trim().is_empty() { return Err(ToolError::message("run_id is required")); } - if matches!(params.action, RunInteractAction::Message) - && params - .message - .as_deref() - .is_none_or(|message| message.trim().is_empty()) - { - return Err(ToolError::message("message is required for action message")); - } - if matches!(params.action, RunInteractAction::Answer) { - if params.question_id.as_deref().is_none_or(str::is_empty) { - return Err(ToolError::message( - "question_id is required for action answer", - )); + let action = match params.action { + RunInteractAction::Get => ValidatedInteractAction::Get, + RunInteractAction::Start => ValidatedInteractAction::Start, + RunInteractAction::Message => { + let Some(message) = params + .message + .as_deref() + .map(str::trim) + .filter(|message| !message.is_empty()) + else { + return Err(ToolError::message("message is required for action message")); + }; + ValidatedInteractAction::Message { + message: message.to_string(), + interrupt: params.interrupt.unwrap_or(false), + } } - let Some(answer) = params.answer.as_ref() else { - return Err(ToolError::message("answer is required for action answer")); - }; - answer_to_submit_request(answer.clone())?; - } - Ok(Self { raw: params }) + RunInteractAction::Cancel => ValidatedInteractAction::Cancel, + RunInteractAction::Archive => ValidatedInteractAction::Archive, + RunInteractAction::Unarchive => ValidatedInteractAction::Unarchive, + RunInteractAction::GetQuestions => ValidatedInteractAction::GetQuestions, + RunInteractAction::Answer => { + let Some(question_id) = params + .question_id + .as_deref() + .map(str::trim) + .filter(|question_id| !question_id.is_empty()) + else { + return Err(ToolError::message( + "question_id is required for action answer", + )); + }; + let Some(answer) = params.answer else { + return Err(ToolError::message("answer is required for action answer")); + }; + ValidatedInteractAction::Answer { + question_id: question_id.to_string(), + body: answer_to_submit_request(answer)?, + } + } + }; + Ok(Self { + run_id: params.run_id.trim().to_string(), + action, + }) } } @@ -79,68 +138,57 @@ pub(crate) async fn interact_run( client: Arc, params: ValidatedInteractRun, ) -> ToolResult { - let raw = params.raw; let run_id = client - .resolve_run(&raw.run_id) + .resolve_run(¶ms.run_id) .await .map_err(|err| ToolError::from_anyhow(&err))? .id; - let result = match raw.action { - RunInteractAction::Get => interact_get(&client, &run_id).await?, - RunInteractAction::Start => { - client + let action = params.action.action(); + let result = match params.action { + ValidatedInteractAction::Get => interact_get(&client, &run_id).await?, + ValidatedInteractAction::Start => { + let summary = client .start_run(&run_id, false) .await .map_err(|err| ToolError::from_anyhow(&err))?; - json!({ "summary": common::run_summary_result(&common::retrieve_run(&client, &run_id).await?) }) + json!({ "summary": common::run_summary_result(&summary) }) } - RunInteractAction::Message => { - let message = raw - .message - .expect("validated message action has a message") - .trim() - .to_string(); + ValidatedInteractAction::Message { message, interrupt } => { client - .steer_run(&run_id, message.clone(), raw.interrupt.unwrap_or(false)) + .steer_run(&run_id, message.clone(), interrupt) .await .map_err(|err| ToolError::from_anyhow(&err))?; - json!({ "message": message, "interrupt": raw.interrupt.unwrap_or(false) }) + json!({ "message": message, "interrupt": interrupt }) } - RunInteractAction::Cancel => { - client + ValidatedInteractAction::Cancel => { + let summary = client .cancel_run(&run_id) .await .map_err(|err| ToolError::from_anyhow(&err))?; - json!({ "summary": common::run_summary_result(&common::retrieve_run(&client, &run_id).await?) }) + json!({ "summary": common::run_summary_result(&summary) }) } - RunInteractAction::Archive => { - client + ValidatedInteractAction::Archive => { + let summary = client .archive_run(&run_id) .await .map_err(|err| ToolError::from_anyhow(&err))?; - json!({ "summary": common::run_summary_result(&common::retrieve_run(&client, &run_id).await?) }) + json!({ "summary": common::run_summary_result(&summary) }) } - RunInteractAction::Unarchive => { - client + ValidatedInteractAction::Unarchive => { + let summary = client .unarchive_run(&run_id) .await .map_err(|err| ToolError::from_anyhow(&err))?; - json!({ "summary": common::run_summary_result(&common::retrieve_run(&client, &run_id).await?) }) + json!({ "summary": common::run_summary_result(&summary) }) } - RunInteractAction::GetQuestions => { + ValidatedInteractAction::GetQuestions => { let questions = client .list_run_questions(&run_id) .await .map_err(|err| ToolError::from_anyhow(&err))?; json!({ "questions": questions }) } - RunInteractAction::Answer => { - let question_id = raw - .question_id - .expect("validated answer action has a question_id"); - let body = answer_to_submit_request( - raw.answer.expect("validated answer action has an answer"), - )?; + ValidatedInteractAction::Answer { question_id, body } => { client .submit_run_answer(&run_id, &question_id, body) .await @@ -151,7 +199,7 @@ pub(crate) async fn interact_run( Ok(InteractRunResult { run_id: run_id.to_string(), - action: raw.action, + action, result, }) } @@ -176,31 +224,59 @@ async fn interact_get(client: &Client, run_id: &RunId) -> ToolResult { } fn answer_to_submit_request(answer: Value) -> ToolResult { - let payload = match answer { - Value::Bool(true) => json!({ "kind": "yes" }), - Value::Bool(false) => json!({ "kind": "no" }), - Value::String(text) => json!({ "kind": "text", "text": text }), + match answer { + Value::Bool(true) => Ok(types::SubmitAnswerYesRequest { + kind: types::SubmitAnswerYesRequestKind::Yes, + } + .into()), + Value::Bool(false) => Ok(types::SubmitAnswerNoRequest { + kind: types::SubmitAnswerNoRequestKind::No, + } + .into()), + Value::String(text) => Ok(text_answer_request(text)), Value::Object(mut object) => { if let Some(option) = object.remove("option") { - json!({ "kind": "selected", "option_key": option }) + let option_key = serde_json::from_value::(option).map_err(|err| { + ToolError::message(format!("answer option must be a string: {err}")) + })?; + Ok(types::SubmitAnswerSelectedRequest { + kind: types::SubmitAnswerSelectedRequestKind::Selected, + option_key, + } + .into()) } else if let Some(options) = object.remove("options") { - json!({ "kind": "multi_selected", "option_keys": options }) + let option_keys = + serde_json::from_value::>(options).map_err(|err| { + ToolError::message(format!("answer options must be strings: {err}")) + })?; + Ok(types::SubmitAnswerMultiSelectedRequest { + kind: types::SubmitAnswerMultiSelectedRequestKind::MultiSelected, + option_keys, + } + .into()) } else if let Some(text) = object.remove("text") { - json!({ "kind": "text", "text": text }) + let text = serde_json::from_value::(text).map_err(|err| { + ToolError::message(format!("answer text must be a string: {err}")) + })?; + Ok(text_answer_request(text)) } else { - return Err(ToolError::message( + Err(ToolError::message( "answer object must contain one of: option, options, text", - )); + )) } } - other => { - return Err(ToolError::message(format!( - "unsupported answer value: {other}; expected boolean, string, or object", - ))); - } - }; - serde_json::from_value(payload) - .map_err(|err| ToolError::message(format!("failed to build submit-answer request: {err}"))) + other => Err(ToolError::message(format!( + "unsupported answer value: {other}; expected boolean, string, or object", + ))), + } +} + +fn text_answer_request(text: String) -> types::SubmitAnswerRequest { + types::SubmitAnswerTextRequest { + kind: types::SubmitAnswerTextRequestKind::Text, + text, + } + .into() } #[cfg(test)] diff --git a/lib/crates/fabro-mcp-server/src/run_tools/manifest.rs b/lib/crates/fabro-mcp-server/src/run_tools/manifest.rs index aa7607c65..af4a1c295 100644 --- a/lib/crates/fabro-mcp-server/src/run_tools/manifest.rs +++ b/lib/crates/fabro-mcp-server/src/run_tools/manifest.rs @@ -1,46 +1,27 @@ -use std::collections::HashMap; use std::path::{Path, PathBuf}; use fabro_api::types; use fabro_config::{CliLayer, RunLayer}; use fabro_manifest::{self, ManifestBuildInput, RunOverrideInput}; use fabro_server::manifest_validation; -use fabro_types::RunId; use serde_json::Value; use super::common::{ToolError, ToolResult}; -use super::create::CreateRunSpec; +use super::create::ValidatedCreateRunSpec; pub(super) fn build_mcp_run_manifest( - spec: &CreateRunSpec, + spec: &ValidatedCreateRunSpec, cwd: &Path, user_settings_path: &Path, ) -> ToolResult { - if let Some(run_id) = spec.run_id.as_deref() { - run_id.parse::().map_err(|err| { - ToolError::message(format!("run_id must be a valid Fabro run id: {err}")) - })?; - } - let built = fabro_manifest::build_run_manifest(ManifestBuildInput { workflow: PathBuf::from(&spec.workflow), cwd: cwd.to_path_buf(), run_overrides: mcp_run_overrides(spec), cli_overrides: Some(CliLayer::default()), - input_overrides: spec - .inputs - .iter() - .map(|(key, value)| json_to_toml_value(key, value).map(|value| (key.clone(), value))) - .collect::>>()?, + input_overrides: spec.inputs.clone(), args: mcp_manifest_args(spec), - run_id: spec - .run_id - .as_deref() - .map(str::parse::) - .transpose() - .map_err(|err| { - ToolError::message(format!("run_id must be a valid Fabro run id: {err}")) - })?, + run_id: spec.run_id, user_settings_path: Some(user_settings_path.to_path_buf()), }) .map_err(|err| ToolError::from_anyhow(&err))?; @@ -85,7 +66,7 @@ pub(super) fn json_to_toml_value(key: &str, value: &Value) -> ToolResult Option { +fn mcp_manifest_args(spec: &ValidatedCreateRunSpec) -> Option { let mut input = spec .inputs .iter() @@ -113,12 +94,13 @@ fn mcp_manifest_args(spec: &CreateRunSpec) -> Option { (!fabro_manifest::manifest_args_is_empty(&payload)).then_some(payload) } -fn mcp_run_overrides(spec: &CreateRunSpec) -> Option { +fn mcp_run_overrides(spec: &ValidatedCreateRunSpec) -> Option { fabro_manifest::build_sparse_run_overrides(RunOverrideInput { goal: spec.goal.as_deref(), model: spec.model.as_deref(), provider: spec.provider.as_deref(), sandbox: spec.sandbox.as_deref(), + docker_image: None, preserve_sandbox: spec.preserve_sandbox, dry_run: spec.dry_run, auto_approve: spec.auto_approve, @@ -173,7 +155,7 @@ mod tests { #[test] fn mcp_manifest_args_preserve_input_provenance() { - let args = mcp_manifest_args(&CreateRunSpec { + let spec = ValidatedCreateRunSpec::try_from(CreateRunSpec { workflow: "simple".to_string(), run_id: None, cwd: None, @@ -191,7 +173,8 @@ mod tests { preserve_sandbox: None, start: None, }) - .expect("input args should be present"); + .expect("create spec should validate"); + let args = mcp_manifest_args(&spec).expect("input args should be present"); assert_eq!(args.input, vec![r"count=3", r#"decision="approve""#]); } diff --git a/lib/crates/fabro-mcp-server/src/run_tools/search.rs b/lib/crates/fabro-mcp-server/src/run_tools/search.rs index 14336f4a2..fe4e83a05 100644 --- a/lib/crates/fabro-mcp-server/src/run_tools/search.rs +++ b/lib/crates/fabro-mcp-server/src/run_tools/search.rs @@ -2,6 +2,8 @@ use std::collections::HashMap; use std::sync::Arc; use fabro_client::Client; +use fabro_types::{Run, RunStatusKind}; +use futures::future::try_join_all; use schemars::JsonSchema; use serde::{Deserialize, Serialize}; @@ -23,7 +25,8 @@ pub(crate) struct FabroRunSearchParams { #[derive(Debug)] pub(crate) struct ValidatedSearchRuns { - pub(crate) raw: FabroRunSearchParams, + pub(crate) raw: FabroRunSearchParams, + pub(crate) status: Option>, } impl TryFrom for ValidatedSearchRuns { @@ -33,13 +36,33 @@ impl TryFrom for ValidatedSearchRuns { if params.first.is_some_and(|first| first > 100) { return Err(ToolError::message("first must be <= 100")); } + if let Some(run_ids) = params.run_ids.as_ref() { + common::validate_len("run_ids", run_ids.len(), 1, 100)?; + } + let status = params + .status + .as_ref() + .map(|statuses| { + statuses + .iter() + .map(|status| { + status.parse::().map_err(|_| { + ToolError::message(format!("unknown run status `{status}`")) + }) + }) + .collect::>>() + }) + .transpose()?; if let Some(created_after) = params.created_after.as_deref() { common::parse_datetime_filter("created_after", created_after)?; } if let Some(created_before) = params.created_before.as_deref() { common::parse_datetime_filter("created_before", created_before)?; } - Ok(Self { raw: params }) + Ok(Self { + raw: params, + status, + }) } } @@ -53,11 +76,16 @@ pub(crate) async fn search_runs( client: Arc, params: ValidatedSearchRuns, ) -> ToolResult { + let status = params.status; let raw = params.raw; - let mut runs = client - .list_store_runs() - .await - .map_err(|err| ToolError::from_anyhow(&err))?; + let mut runs = if let Some(run_ids) = raw.run_ids.as_ref() { + resolve_requested_runs(&client, run_ids).await? + } else { + client + .list_store_runs() + .await + .map_err(|err| ToolError::from_anyhow(&err))? + }; runs.sort_by(|a, b| { let a_sort_time = a.timestamps.started_at.unwrap_or(a.timestamps.created_at); let b_sort_time = b.timestamps.started_at.unwrap_or(b.timestamps.created_at); @@ -70,9 +98,6 @@ pub(crate) async fn search_runs( } } - if let Some(run_ids) = raw.run_ids.as_ref() { - runs.retain(|run| run_ids.iter().any(|id| id == &run.id.to_string())); - } if let Some(workflow) = raw.workflow.as_deref() { runs.retain(|run| { run.workflow.name == workflow || run.workflow.slug.as_deref() == Some(workflow) @@ -85,11 +110,11 @@ pub(crate) async fn search_runs( .all(|(key, value)| run.labels.get(key) == Some(value)) }); } - if let Some(status) = raw.status.as_ref() { + if let Some(status) = status.as_ref() { runs.retain(|run| { status .iter() - .any(|status| status == common::run_status_kind(run.lifecycle.status)) + .any(|status| *status == run.lifecycle.status.kind()) }); } if let Some(archived) = raw.archived { @@ -119,3 +144,22 @@ pub(crate) async fn search_runs( pub(crate) fn search_runs_text(result: &SearchRunsResult) -> String { format!("found {} Fabro run(s)", result.runs.len()) } + +async fn resolve_requested_runs(client: &Arc, run_ids: &[String]) -> ToolResult> { + let runs = try_join_all(run_ids.iter().map(|run_id| { + let client = Arc::clone(client); + async move { + client + .resolve_run(run_id) + .await + .map_err(|err| ToolError::from_anyhow(&err)) + } + })) + .await?; + + let mut unique = HashMap::new(); + for run in runs { + unique.entry(run.id).or_insert(run); + } + Ok(unique.into_values().collect()) +} diff --git a/lib/crates/fabro-mcp-server/src/server.rs b/lib/crates/fabro-mcp-server/src/server.rs index c381b9348..74fe1e7e5 100644 --- a/lib/crates/fabro-mcp-server/src/server.rs +++ b/lib/crates/fabro-mcp-server/src/server.rs @@ -194,20 +194,22 @@ async fn client_from_settings(settings: &FabroMcpServerSettings) -> Result Result { - let target: ServerTarget = server.parse()?; +async fn connect_target( + target: &ServerTarget, + settings: &FabroMcpServerSettings, +) -> Result { let auth_store = AuthStore::default(); - let mut credential = resolve_target_credential_with_store(&target, &auth_store)?; + let mut credential = resolve_target_credential_with_store(target, &auth_store)?; if credential.is_none() && target.is_unix_socket() { let runtime_token_path = Storage::new(&settings.storage_dir) .runtime_directory() .dev_token_path(); credential = dev_token::read_dev_token_file(&runtime_token_path).map(Credential::DevToken); } - let oauth_session = refreshable_oauth(&target, &auth_store, credential.as_ref()); + let oauth_session = refreshable_oauth(target, &auth_store, credential.as_ref()); let mut builder = Client::builder() .target(target.clone()) - .transport_connector(target_transport_connector(target)) + .transport_connector(target_transport_connector(target.clone())) .request_timeout(CLIENT_REQUEST_TIMEOUT); if let Some(credential) = credential { builder = builder.credential(credential); diff --git a/lib/crates/fabro-server/Cargo.toml b/lib/crates/fabro-server/Cargo.toml index e78933660..304ae45ec 100644 --- a/lib/crates/fabro-server/Cargo.toml +++ b/lib/crates/fabro-server/Cargo.toml @@ -35,6 +35,7 @@ fabro-sandbox = { path = "../fabro-sandbox", features = ["daytona", "docker"] } fabro-github = { path = "../fabro-github" } fabro-agent = { path = "../fabro-agent" } fabro-llm = { path = "../fabro-llm" } +fabro-manifest = { path = "../fabro-manifest" } fabro-model = { path = "../fabro-model" } fabro-proc = { path = "../fabro-proc" } fabro-types = { path = "../fabro-types" } diff --git a/lib/crates/fabro-server/src/run_manifest.rs b/lib/crates/fabro-server/src/run_manifest.rs index 0bced72f6..2178196ee 100644 --- a/lib/crates/fabro-server/src/run_manifest.rs +++ b/lib/crates/fabro-server/src/run_manifest.rs @@ -9,8 +9,7 @@ use fabro_api::types; use fabro_auth::auth_issue_message; use fabro_config::run::parse_run_layer_from_settings_toml; use fabro_config::{ - CliLayer, CliOutputLayer, DaytonaDockerfileLayer, DockerSandboxLayer, ReplaceMap, - RunExecutionLayer, RunLayer, RunModelLayer, RunSandboxLayer, WorkflowSettingsBuilder, + CliLayer, CliOutputLayer, DaytonaDockerfileLayer, RunLayer, WorkflowSettingsBuilder, parse_input_overrides, }; use fabro_graphviz::graph::{Graph, is_llm_handler_type}; @@ -28,8 +27,8 @@ use fabro_static::EnvVars; use fabro_types::settings::cli::OutputVerbosity; use fabro_types::settings::interp::InterpString; use fabro_types::settings::run::{ - ApprovalMode, DaytonaNetworkLayer, DaytonaSettings, DockerSettings, DockerfileSource, RunGoal, - RunMode, RunNamespace, + DaytonaNetworkLayer, DaytonaSettings, DockerSettings, DockerfileSource, RunGoal, RunMode, + RunNamespace, }; use fabro_types::{RunId, WorkflowSettings}; use fabro_util::check_report::{CheckDetail, CheckReport, CheckResult, CheckSection, CheckStatus}; @@ -316,46 +315,16 @@ fn manifest_args_overrides( return Ok(ManifestSettingsOverrides::default()); }; - let model = (args.model.is_some() || args.provider.is_some()).then(|| RunModelLayer { - provider: args.provider.as_deref().map(InterpString::parse), - name: args.model.as_deref().map(InterpString::parse), - fallbacks: Vec::new(), - }); - let sandbox = - (args.sandbox.is_some() || args.preserve_sandbox.is_some() || args.docker_image.is_some()) - .then(|| RunSandboxLayer { - provider: args.sandbox.clone(), - preserve: args.preserve_sandbox, - docker: args.docker_image.as_ref().map(|image| DockerSandboxLayer { - image: Some(image.clone()), - ..DockerSandboxLayer::default() - }), - ..RunSandboxLayer::default() - }); - - let execution_has_any = args.dry_run.is_some() || args.auto_approve.is_some(); - let execution = execution_has_any.then(|| RunExecutionLayer { - mode: args - .dry_run - .map(|d| if d { RunMode::DryRun } else { RunMode::Normal }), - approval: args.auto_approve.map(|a| { - if a { - ApprovalMode::Auto - } else { - ApprovalMode::Prompt - } - }), - }); - - let run_has_any = - model.is_some() || sandbox.is_some() || execution.is_some() || !args.label.is_empty(); - - let run = run_has_any.then(|| RunLayer { - model, - sandbox, - execution, - metadata: ReplaceMap::from(parse_labels(&args.label)), - ..RunLayer::default() + let run = fabro_manifest::build_sparse_run_overrides(fabro_manifest::RunOverrideInput { + goal: None, + model: args.model.as_deref(), + provider: args.provider.as_deref(), + sandbox: args.sandbox.as_deref(), + docker_image: args.docker_image.as_deref(), + preserve_sandbox: args.preserve_sandbox, + dry_run: args.dry_run, + auto_approve: args.auto_approve, + labels: parse_labels(&args.label), }); // Verbose is a CLI output concern in v2; route it through cli.output.verbosity. diff --git a/lib/crates/fabro-types/src/lib.rs b/lib/crates/fabro-types/src/lib.rs index dcac895f2..c1a9afeb7 100644 --- a/lib/crates/fabro-types/src/lib.rs +++ b/lib/crates/fabro-types/src/lib.rs @@ -102,5 +102,6 @@ pub use stage_id::{InvalidStageVisit, ParallelBranchId, StageId}; pub use start::StartRecord; pub use status::{ BlockedReason, FailureReason, InvalidTransition, ParseFailureReasonError, - ParseSuccessReasonError, RunControlAction, RunStatus, SuccessReason, TerminalStatus, + ParseSuccessReasonError, RunControlAction, RunStatus, RunStatusKind, SuccessReason, + TerminalStatus, }; diff --git a/lib/crates/fabro-types/src/status.rs b/lib/crates/fabro-types/src/status.rs index d59f89b65..cf2d7d1bb 100644 --- a/lib/crates/fabro-types/src/status.rs +++ b/lib/crates/fabro-types/src/status.rs @@ -2,6 +2,35 @@ use std::fmt; use std::str::FromStr; use serde::{Deserialize, Serialize}; +use strum::{Display, EnumString, IntoStaticStr}; + +#[derive( + Debug, + Clone, + Copy, + PartialEq, + Eq, + Hash, + Serialize, + Deserialize, + Display, + EnumString, + IntoStaticStr, +)] +#[serde(rename_all = "snake_case")] +#[strum(serialize_all = "snake_case")] +pub enum RunStatusKind { + Submitted, + Queued, + Starting, + Running, + Blocked, + Paused, + Removing, + Succeeded, + Failed, + Dead, +} #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case")] @@ -19,6 +48,10 @@ pub enum RunStatus { } impl RunStatus { + pub fn kind(self) -> RunStatusKind { + self.into() + } + /// Whether the run has reached a terminal outcome and stops poll loops, /// finalization, and similar "done" handling. pub fn is_terminal(self) -> bool { @@ -138,6 +171,23 @@ impl RunStatus { } } +impl From for RunStatusKind { + fn from(status: RunStatus) -> Self { + match status { + RunStatus::Submitted => Self::Submitted, + RunStatus::Queued => Self::Queued, + RunStatus::Starting => Self::Starting, + RunStatus::Running => Self::Running, + RunStatus::Blocked { .. } => Self::Blocked, + RunStatus::Paused { .. } => Self::Paused, + RunStatus::Removing => Self::Removing, + RunStatus::Succeeded { .. } => Self::Succeeded, + RunStatus::Failed { .. } => Self::Failed, + RunStatus::Dead => Self::Dead, + } + } +} + impl fmt::Display for RunStatus { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self {