From b1d95faa57cd84e52ca94f39e0e44e5ac361b297 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 19 Sep 2026 07:08:08 -0400 Subject: [PATCH] Hand Petri the server's environment and MCP catalogs Petri's Fabro frontend refused a bundle naming an environment it did not declare and every MCP catalog reference, so the fixtures declared `[environments.local]` and the server's catalogs never reached Petri. Pin Petri at c874b86, where the frontend reads `[environments.]` and `[run.environment]` from every settings layer (bundle over project over the host's layer, key by key), takes the environment a launch selected over the layers, and resolves `[run.agent.mcps.] id = "..."` against a catalog the host binds. The server hands Petri its environment catalog as `[environments.]` tables of the settings layer it already passes, the intent's environment as the launch's selection (`Launch::environment`, as the intent overrides the bundle in Fabro's own resolution), and its MCP catalog as `RuntimeSpec::mcp_catalog_toml`, one inline entry per definition keyed by id. Offline validation hands Petri the seeded catalog the same way, so `fabro validate` accepts `[run.environment] id = "local"`. The fixtures drop the `[environments.local]` tables they carried for this; the secrets test keeps its own, on purpose. Scenario tests cover a bundle naming a catalog environment (its image lowered, and run on Docker when the plugin and a daemon are there), a bundle's own table winning key by key, the server refusing an unknown environment before Petri, and a catalog MCP reference whose tool the agent session lists (an echo server under `test/mcp/`). Co-Authored-By: Claude Fable 5.1 --- Cargo.lock | 30 +- Cargo.toml | 14 +- .../src/commands/run/petri_worker.rs | 1 + lib/apps/fabro-cli/tests/it/cmd/dump.rs | 3 - lib/apps/fabro-cli/tests/it/cmd/run.rs | 3 - lib/apps/fabro-cli/tests/it/cmd/support.rs | 7 +- .../fabro-server/src/manifest_validation.rs | 32 +- lib/apps/fabro-server/src/petri_check.rs | 20 +- lib/apps/fabro-server/src/run_compiler.rs | 16 + lib/apps/fabro-server/src/run_manifest.rs | 15 +- .../fabro-server/src/server/handler/runs.rs | 8 +- .../fabro-server/src/server/petri_runs.rs | 208 ++++++++- .../fabro-server/tests/it/api/mcp_servers.rs | 54 ++- .../fabro-server/tests/it/scenario/petri.rs | 403 +++++++++++++++++- lib/components/fabro-petri/src/check.rs | 33 +- lib/components/fabro-petri/src/runtime.rs | 26 +- lib/components/fabro-petri/tests/check.rs | 58 ++- test/mcp/echo_server.py | 80 ++++ 18 files changed, 890 insertions(+), 121 deletions(-) create mode 100755 test/mcp/echo_server.py diff --git a/Cargo.lock b/Cargo.lock index 72ea10089..fc302fe1a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5324,7 +5324,7 @@ checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" [[package]] name = "petri-attractor-steps" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "async-trait", "globset", @@ -5355,7 +5355,7 @@ dependencies = [ [[package]] name = "petri-driver" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "async-trait", "getrandom 0.3.4", @@ -5375,7 +5375,7 @@ dependencies = [ [[package]] name = "petri-engine" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "petri-ir", "serde", @@ -5387,7 +5387,7 @@ dependencies = [ [[package]] name = "petri-execution" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "async-trait", "petri-driver", @@ -5411,7 +5411,7 @@ dependencies = [ [[package]] name = "petri-executor" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "async-trait", "libc", @@ -5426,7 +5426,7 @@ dependencies = [ [[package]] name = "petri-executor-sandbox" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "async-trait", "petri-executor", @@ -5448,7 +5448,7 @@ dependencies = [ [[package]] name = "petri-frontend" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "marked-yaml", "petri-ir", @@ -5462,7 +5462,7 @@ dependencies = [ [[package]] name = "petri-frontend-attractor" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "minijinja", "petri-frontend", @@ -5479,7 +5479,7 @@ dependencies = [ [[package]] name = "petri-frontend-fabro" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "petri-frontend", "petri-frontend-attractor", @@ -5495,7 +5495,7 @@ dependencies = [ [[package]] name = "petri-frontend-native" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "petri-frontend", "petri-ir", @@ -5506,7 +5506,7 @@ dependencies = [ [[package]] name = "petri-ir" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "regex", "serde", @@ -5519,7 +5519,7 @@ dependencies = [ [[package]] name = "petri-runtime" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "async-trait", "petri-driver", @@ -5540,7 +5540,7 @@ dependencies = [ [[package]] name = "petri-steps" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "async-trait", "petri-executor", @@ -5556,7 +5556,7 @@ dependencies = [ [[package]] name = "petri-store" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "async-trait", "getrandom 0.3.4", @@ -5571,7 +5571,7 @@ dependencies = [ [[package]] name = "petri-testkit" version = "0.1.0" -source = "git+https://github.com/lithoscomputer/petri.git?rev=639ce3e368b84ff48a3d8f695b8c5736b1c1563e#639ce3e368b84ff48a3d8f695b8c5736b1c1563e" +source = "git+https://github.com/lithoscomputer/petri.git?rev=c874b863671eec309f111e5639c2687f29342042#c874b863671eec309f111e5639c2687f29342042" dependencies = [ "async-trait", "petri-driver", diff --git a/Cargo.toml b/Cargo.toml index 4ddb5c447..f2a105ade 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -132,13 +132,13 @@ pebble-cli-core = { git = "https://github.com/lithoscomputer/pebble", rev = "a39 # lithos-llm and sandbox-driver revisions as this file, so the workspace links # one copy of each. Only `fabro-petri` may depend on these packages; the keys # carry the `petri_` prefix so the crate names say where they come from. -petri_runtime = { git = "https://github.com/lithoscomputer/petri.git", rev = "639ce3e368b84ff48a3d8f695b8c5736b1c1563e", package = "petri-runtime" } -petri_execution = { git = "https://github.com/lithoscomputer/petri.git", rev = "639ce3e368b84ff48a3d8f695b8c5736b1c1563e", package = "petri-execution" } -petri_store = { git = "https://github.com/lithoscomputer/petri.git", rev = "639ce3e368b84ff48a3d8f695b8c5736b1c1563e", package = "petri-store" } -petri_attractor_steps = { git = "https://github.com/lithoscomputer/petri.git", rev = "639ce3e368b84ff48a3d8f695b8c5736b1c1563e", package = "petri-attractor-steps" } -petri_frontend_attractor = { git = "https://github.com/lithoscomputer/petri.git", rev = "639ce3e368b84ff48a3d8f695b8c5736b1c1563e", package = "petri-frontend-attractor" } -petri_frontend_fabro = { git = "https://github.com/lithoscomputer/petri.git", rev = "639ce3e368b84ff48a3d8f695b8c5736b1c1563e", package = "petri-frontend-fabro" } -petri_testkit = { git = "https://github.com/lithoscomputer/petri.git", rev = "639ce3e368b84ff48a3d8f695b8c5736b1c1563e", package = "petri-testkit" } +petri_runtime = { git = "https://github.com/lithoscomputer/petri.git", rev = "c874b863671eec309f111e5639c2687f29342042", package = "petri-runtime" } +petri_execution = { git = "https://github.com/lithoscomputer/petri.git", rev = "c874b863671eec309f111e5639c2687f29342042", package = "petri-execution" } +petri_store = { git = "https://github.com/lithoscomputer/petri.git", rev = "c874b863671eec309f111e5639c2687f29342042", package = "petri-store" } +petri_attractor_steps = { git = "https://github.com/lithoscomputer/petri.git", rev = "c874b863671eec309f111e5639c2687f29342042", package = "petri-attractor-steps" } +petri_frontend_attractor = { git = "https://github.com/lithoscomputer/petri.git", rev = "c874b863671eec309f111e5639c2687f29342042", package = "petri-frontend-attractor" } +petri_frontend_fabro = { git = "https://github.com/lithoscomputer/petri.git", rev = "c874b863671eec309f111e5639c2687f29342042", package = "petri-frontend-fabro" } +petri_testkit = { git = "https://github.com/lithoscomputer/petri.git", rev = "c874b863671eec309f111e5639c2687f29342042", package = "petri-testkit" } sentry = { version = "0.35", default-features = false, features = ["backtrace", "contexts", "ureq", "rustls"] } fork = "0.2" exec = "0.3" diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index 9403a2d6d..43d28d4b6 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -559,6 +559,7 @@ async fn runtime_spec( }; Ok(RuntimeSpec { settings_toml: None, + mcp_catalog_toml: None, model_client, dry_run: run_state.spec.settings.run.execution.mode == RunMode::DryRun, fabro_home, diff --git a/lib/apps/fabro-cli/tests/it/cmd/dump.rs b/lib/apps/fabro-cli/tests/it/cmd/dump.rs index 7e0335c11..53a09d96b 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/dump.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/dump.rs @@ -180,9 +180,6 @@ goal = "Generate oversized command output and artifacts" [run.environment] id = "local" -[environments.local] -provider = "local" - [run.artifacts] include = ["assets/**"] "#, diff --git a/lib/apps/fabro-cli/tests/it/cmd/run.rs b/lib/apps/fabro-cli/tests/it/cmd/run.rs index 4c5a7a3c2..2931f27f9 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/run.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/run.rs @@ -784,9 +784,6 @@ goal = "Show stored artifacts" [run.environment] id = "local" -[environments.local] -provider = "local" - [run.artifacts] include = ["assets/**"] "#, diff --git a/lib/apps/fabro-cli/tests/it/cmd/support.rs b/lib/apps/fabro-cli/tests/it/cmd/support.rs index 1ff0929d8..6a7787ecb 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/support.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/support.rs @@ -550,8 +550,7 @@ fn git_backed_run( write_text_file( &workspace_dir.join("workflow.toml"), "_version = 1\n\n[workflow]\ngraph = \"story.fabro\"\n\n[run]\ngoal = \"Change the \ - story\"\n\n[run.environment]\nid = \"local\"\n\n[environments.local]\nprovider = \ - \"local\"\n", + story\"\n\n[run.environment]\nid = \"local\"\n", ); init_remote_fixture(&workspace_dir, "main"); let run = run_local_workflow(context, &workspace_dir, "workflow.toml"); @@ -594,10 +593,6 @@ goal = "Exercise sandbox commands" [run.environment] id = "local" - -[environments.local] -provider = "local" - "#, ); diff --git a/lib/apps/fabro-server/src/manifest_validation.rs b/lib/apps/fabro-server/src/manifest_validation.rs index 43c93571c..f04b8f954 100644 --- a/lib/apps/fabro-server/src/manifest_validation.rs +++ b/lib/apps/fabro-server/src/manifest_validation.rs @@ -3,7 +3,7 @@ use std::path::PathBuf; use anyhow::{Result, anyhow}; use fabro_api::types; -use fabro_config::{RunLayer, WorkflowSettingsBuilder}; +use fabro_config::{RunLayer, SettingsLayer, WorkflowSettingsBuilder}; use fabro_manifest::CollectedWorkflowClosure; use fabro_petri::runtime::RuntimeSpec; use fabro_workflow::operations::{ValidateInput, WorkflowInput, validate}; @@ -32,7 +32,7 @@ pub fn validate_manifest( &prepared, &HashMap::new(), petri_check::launch_without_catalog(&prepared.settings), - offline_runtime(), + offline_runtime(Some(manifest_run_defaults)), true, true, ) @@ -40,15 +40,25 @@ pub fn validate_manifest( Ok(run_manifest::validate_response(&prepared, &validated)) } -/// Petri's runtime for a check away from the server: no operator settings, -/// no model client, no Fabro home, no run tools. -fn offline_runtime() -> RuntimeSpec { +/// Petri's runtime for a check away from the server: the seeded environment +/// catalog and the given `[run]` layer as the settings layer, the same +/// defaults the legacy validation judges against, so a bundle that names a +/// seeded environment validates; no MCP catalog, no model client, no Fabro +/// home, no run tools. +fn offline_runtime(run: Option<&RunLayer>) -> RuntimeSpec { + let layer = SettingsLayer { + version: Some(1), + environments: fabro_environment::seeded_catalog_layer(), + run: run.cloned(), + ..SettingsLayer::default() + }; RuntimeSpec { - settings_toml: None, - model_client: None, - dry_run: false, - fabro_home: None, - run_tools: None, + settings_toml: toml::to_string(&layer).ok(), + mcp_catalog_toml: None, + model_client: None, + dry_run: false, + fabro_home: None, + run_tools: None, } } @@ -98,7 +108,7 @@ pub fn validate_collected_workflow( &settings, &HashMap::new(), petri_check::launch_without_catalog(&settings), - offline_runtime(), + offline_runtime(run_overrides), false, ) .map_err(anyhow::Error::new)?; diff --git a/lib/apps/fabro-server/src/petri_check.rs b/lib/apps/fabro-server/src/petri_check.rs index 543c43bc3..fd88e1a70 100644 --- a/lib/apps/fabro-server/src/petri_check.rs +++ b/lib/apps/fabro-server/src/petri_check.rs @@ -25,15 +25,17 @@ use lithos_llm::catalog::ProviderId; /// Fabro's rule for a model node with no provider ready to run it. pub(crate) const NO_READY_PROVIDER_RULE: &str = "fabro.model.no_ready_provider"; -/// The launch Fabro binds below the settings: the run's model and provider. -/// When the settings name neither, the default offering of the eligible -/// providers is bound as the launch model alone: a node that names no model -/// runs on it, and a node that names a model the catalog lacks stays -/// unqualified, so Petri's admission refuses it. +/// The launch Fabro binds around the settings: the run's model and provider +/// below them, and the environment the run selected above them. When the +/// settings name neither model nor provider, the default offering of the +/// eligible providers is bound as the launch model alone: a node that +/// names no model runs on it, and a node that names a model the catalog +/// lacks stays unqualified, so Petri's admission refuses it. pub(crate) fn launch( catalog: &Catalog, settings: &WorkflowSettings, eligible: &[ProviderId], + environment: Option<&str>, repository: Option, ) -> Launch { let model = settings.run.model.name.clone().or_else(|| { @@ -48,6 +50,7 @@ pub(crate) fn launch( Launch { model, provider: settings.run.model.provider.clone(), + environment: environment.map(str::to_owned), repository, } } @@ -56,9 +59,10 @@ pub(crate) fn launch( /// name, for a check away from the server. pub(crate) fn launch_without_catalog(settings: &WorkflowSettings) -> Launch { Launch { - model: settings.run.model.name.clone(), - provider: settings.run.model.provider.clone(), - repository: None, + model: settings.run.model.name.clone(), + provider: settings.run.model.provider.clone(), + environment: None, + repository: None, } } diff --git a/lib/apps/fabro-server/src/run_compiler.rs b/lib/apps/fabro-server/src/run_compiler.rs index 1abbd42f2..2bafdd33f 100644 --- a/lib/apps/fabro-server/src/run_compiler.rs +++ b/lib/apps/fabro-server/src/run_compiler.rs @@ -97,6 +97,9 @@ pub(crate) struct NormalizedRun { struct RunMetadata { run_id: Option, + /// The environment the run overrides selected, for Petri's settings + /// layer. + environment_id: Option, storage_root: PathBuf, workflow_slug: Option, workflow_version_id: Option, @@ -165,6 +168,11 @@ impl PreparedRun { self.layered.metadata.parent_id } + /// The environment the run overrides selected, when they did. + pub(crate) fn environment_id(&self) -> Option<&str> { + self.layered.metadata.environment_id.as_deref() + } + pub(crate) fn resolve_run_id(mut self) -> (Self, RunId) { let run_id = self.layered.metadata.run_id.unwrap_or_default(); self.layered.metadata.run_id = Some(run_id); @@ -296,6 +304,10 @@ pub(crate) fn normalize_source(input: RawRunCompilerInput) -> Result Result CreateRunPersistenceInput { } = pinned; let RunMetadata { run_id, + // Consumed at admission, as the launch's environment; the resolved + // settings carry the environment the run persists. + environment_id: _, storage_root, workflow_slug, workflow_version_id, diff --git a/lib/apps/fabro-server/src/run_manifest.rs b/lib/apps/fabro-server/src/run_manifest.rs index 032662fc9..02c834de2 100644 --- a/lib/apps/fabro-server/src/run_manifest.rs +++ b/lib/apps/fabro-server/src/run_manifest.rs @@ -1664,8 +1664,13 @@ mod tests { prepared: &PreparedManifest, ready_providers: &[ProviderId], ) -> Result { - let launch = - petri_check::launch(&state.catalog(), &prepared.settings, ready_providers, None); + let launch = petri_check::launch( + &state.catalog(), + &prepared.settings, + ready_providers, + None, + None, + ); let runtime = crate::server::petri_runs::runtime_spec(state, ready_providers, false); validate_prepared_manifest( prepared, @@ -2560,9 +2565,6 @@ name = "Control Plane" path: "workflow.toml".to_string(), source: r#"_version = 1 -[environments.local] -provider = "local" - [run.environment] id = "local" @@ -2608,9 +2610,6 @@ issues = "read" path: "workflow.toml".to_string(), source: r#"_version = 1 -[environments.local] -provider = "local" - [run.environment] id = "local" diff --git a/lib/apps/fabro-server/src/server/handler/runs.rs b/lib/apps/fabro-server/src/server/handler/runs.rs index 3610d3ebc..9dec4b39d 100644 --- a/lib/apps/fabro-server/src/server/handler/runs.rs +++ b/lib/apps/fabro-server/src/server/handler/runs.rs @@ -1287,7 +1287,13 @@ async fn validate_manifest_on_petri( vars: HashMap, ready_providers: &[ProviderId], ) -> Result { - let launch = petri_check::launch(&state.catalog(), &prepared.settings, ready_providers, None); + let launch = petri_check::launch( + &state.catalog(), + &prepared.settings, + ready_providers, + None, + None, + ); let dry_run = prepared.settings.run.execution.mode == RunMode::DryRun; let runtime = petri_runs::runtime_spec(state, ready_providers, dry_run); let has_ready_provider = !ready_providers.is_empty(); diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 5a14d2c13..ec4844155 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -28,6 +28,7 @@ //! workspace to the snapshot its durable state names, or reports the run //! failed when it cannot. +use std::collections::HashMap; use std::sync::Arc; use std::time::Instant; @@ -44,7 +45,8 @@ use fabro_petri::runtime::{self, RuntimeSpec}; use fabro_petri::secrets::VaultSecrets; use fabro_petri::{SqliteRunStore, admission}; use fabro_store::platform_records::{RunLifecycleKind, RunLifecycleRecord}; -use fabro_types::settings::run::{ApprovalMode, RunMode}; +use fabro_types::settings::McpTransport; +use fabro_types::settings::run::{ApprovalMode, McpServerSettings, RunMode}; use fabro_types::{PetriAdmission, RunId, RunRunnableSource, RunTarget}; use fabro_util::error as error_util; use fabro_workflow::Error as WorkflowError; @@ -63,14 +65,16 @@ use crate::petri_runs::PetriRuns; use crate::run_compiler::{PreparedRun, RunCompilerError}; /// The runtime Petri gets, at create and at execution: the server's run -/// defaults as the settings layer, the model client over the server's -/// catalog and credentials for the eligible providers, and the run mode. +/// defaults and environment catalog as the settings layer, the MCP +/// catalog, the model client over the server's catalog and credentials for +/// the eligible providers, and the run mode. pub(crate) fn runtime_spec( state: &AppState, eligible: &[ProviderId], dry_run: bool, ) -> RuntimeSpec { let settings_toml = settings_layer_toml(state); + let mcp_catalog_toml = mcp_catalog_toml(&state.mcp_server_store().catalog_settings()); let catalog = state.catalog(); let model_client = match runtime::model_client( (*catalog).clone(), @@ -86,6 +90,7 @@ pub(crate) fn runtime_spec( }; RuntimeSpec { settings_toml, + mcp_catalog_toml, model_client, dry_run, fabro_home: Some(Home::from_env().root().to_path_buf()), @@ -95,12 +100,18 @@ pub(crate) fn runtime_spec( } } -/// The server's `[run]` defaults, as the text of the operator settings -/// layer the Fabro frontend reads below `.fabro/project.toml` and -/// `workflow.toml`. +/// The operator settings layer the Fabro frontend reads below +/// `.fabro/project.toml` and `workflow.toml`, as text: the server's `[run]` +/// defaults and the server's environment catalog as `[environments.]` +/// tables, so a bundle can name any of them and Petri lowers that +/// environment's image and resources. A bundle's own table wins over the +/// catalog's key by key, as the frontend layers them; the environment the +/// run selected is bound above every layer by the launch +/// (`petri_check::launch`). fn settings_layer_toml(state: &AppState) -> Option { let layer = SettingsLayer { version: Some(1), + environments: (*state.environment_store().catalog_layer()).clone(), run: Some((*state.manifest_run_defaults()).clone()), ..SettingsLayer::default() }; @@ -113,6 +124,83 @@ fn settings_layer_toml(state: &AppState) -> Option { } } +/// The server's MCP catalog as the text the Fabro frontend resolves +/// `[run.agent.mcps.] id = "..."` references against: one table per +/// definition, keyed by its id, in the inline `[run.agent.mcps.]` +/// shape of Fabro's settings files (`type`, then `command`, `url` or +/// `port`, `env` or `headers`, `protocol`, and the timeouts as durations). +/// `None` for an empty catalog. The frontend reads the entry with the +/// rules of an inline entry, so a `{{ secrets.NAME }}` value under `env` or +/// `headers` resolves at launch as it would in a settings file. +fn mcp_catalog_toml(catalog: &HashMap) -> Option { + if catalog.is_empty() { + return None; + } + let table: toml::Table = catalog + .iter() + .map(|(id, server)| (id.clone(), toml::Value::Table(mcp_catalog_entry(server)))) + .collect(); + match toml::to_string(&table) { + Ok(text) => Some(text), + Err(err) => { + warn!(error = %err, "the MCP catalog does not serialize; Petri gets no catalog"); + None + } + } +} + +fn mcp_catalog_entry(server: &McpServerSettings) -> toml::Table { + let text = |value: &str| toml::Value::String(value.to_string()); + let strings = |values: &HashMap| { + toml::Value::Table( + values + .iter() + .map(|(key, value)| (key.clone(), text(value))) + .collect(), + ) + }; + let argv = |command: &[String]| toml::Value::Array(command.iter().map(|w| text(w)).collect()); + let mut entry = toml::Table::new(); + match &server.transport { + McpTransport::Stdio { command, env } => { + entry.insert("type".to_string(), text("stdio")); + entry.insert("command".to_string(), argv(command)); + entry.insert("env".to_string(), strings(env)); + } + McpTransport::Http { + protocol, + url, + headers, + } => { + entry.insert("type".to_string(), text("http")); + entry.insert("protocol".to_string(), text(&protocol.to_string())); + entry.insert("url".to_string(), text(url)); + entry.insert("headers".to_string(), strings(headers)); + } + McpTransport::Sandbox { + protocol, + command, + port, + env, + } => { + entry.insert("type".to_string(), text("sandbox")); + entry.insert("protocol".to_string(), text(&protocol.to_string())); + entry.insert("command".to_string(), argv(command)); + entry.insert("port".to_string(), toml::Value::Integer(i64::from(*port))); + entry.insert("env".to_string(), strings(env)); + } + } + entry.insert( + "startup_timeout".to_string(), + text(&format!("{}s", server.startup_timeout_secs)), + ); + entry.insert( + "tool_timeout".to_string(), + text(&format!("{}s", server.tool_timeout_secs)), + ); + entry +} + /// Petri compiles the run: check the bundle, map the diagnostics, and /// persist the admitted graphs. A refusal is the same validation error the /// legacy compiler raised, carrying Petri's diagnostics. @@ -126,7 +214,13 @@ pub(crate) async fn admit( Some(RunTarget::Folder { path }) => Some(path.into()), Some(RunTarget::Git(_) | RunTarget::None {}) | None => None, }; - let launch = petri_check::launch(&state.catalog(), settings, eligible, repository); + let launch = petri_check::launch( + &state.catalog(), + settings, + eligible, + prepared.environment_id(), + repository, + ); let dry_run = settings.run.execution.mode == RunMode::DryRun; let request = petri_check::check_request( prepared.workflow_bundle(), @@ -501,3 +595,103 @@ fn finish(state: &Arc, run_id: RunId, status: RunStatus, error: Option drop(runs); state.scheduler_notify.notify_one(); } + +#[cfg(test)] +mod tests { + use fabro_types::settings::run::McpHttpProtocol; + + use super::*; + + /// Every transport of the catalog serializes in the inline shape Petri's + /// Fabro frontend reads, keyed by catalog id, with the timeouts as + /// durations; an empty catalog is no text at all. + #[test] + fn the_mcp_catalog_serializes_in_the_inline_entry_shape() { + assert_eq!(mcp_catalog_toml(&HashMap::new()), None); + let catalog = HashMap::from([ + ("files".to_string(), McpServerSettings { + name: "files".to_string(), + transport: McpTransport::Stdio { + command: vec!["srv".to_string(), "--root".to_string()], + env: HashMap::from([("TOKEN".to_string(), "{{ secrets.T }}".to_string())]), + }, + startup_timeout_secs: 15, + ..McpServerSettings::default() + }), + ("remote".to_string(), McpServerSettings { + name: "remote".to_string(), + transport: McpTransport::Http { + protocol: McpHttpProtocol::Sse, + url: "https://mcp.example/sse".to_string(), + headers: HashMap::from([("X-Org".to_string(), "fabro".to_string())]), + }, + ..McpServerSettings::default() + }), + ("browser".to_string(), McpServerSettings { + name: "browser".to_string(), + transport: McpTransport::Sandbox { + protocol: McpHttpProtocol::StreamableHttp, + command: vec!["npx".to_string(), "mcp".to_string()], + port: 3100, + env: HashMap::new(), + }, + tool_timeout_secs: 90, + ..McpServerSettings::default() + }), + ]); + let text = mcp_catalog_toml(&catalog).expect("the catalog serializes"); + let table: toml::Table = text.parse().expect("the catalog text is TOML"); + assert_eq!(table["files"]["type"].as_str(), Some("stdio")); + assert_eq!( + table["files"]["command"], + toml::Value::Array(vec!["srv".into(), "--root".into()]) + ); + assert_eq!( + table["files"]["env"]["TOKEN"].as_str(), + Some("{{ secrets.T }}") + ); + assert_eq!(table["files"]["startup_timeout"].as_str(), Some("15s")); + assert_eq!(table["files"]["tool_timeout"].as_str(), Some("60s")); + assert_eq!(table["remote"]["type"].as_str(), Some("http")); + assert_eq!(table["remote"]["protocol"].as_str(), Some("sse")); + assert_eq!( + table["remote"]["url"].as_str(), + Some("https://mcp.example/sse") + ); + assert_eq!(table["remote"]["headers"]["X-Org"].as_str(), Some("fabro")); + assert_eq!(table["browser"]["type"].as_str(), Some("sandbox")); + assert_eq!( + table["browser"]["protocol"].as_str(), + Some("streamable_http") + ); + assert_eq!(table["browser"]["port"].as_integer(), Some(3100)); + assert_eq!(table["browser"]["tool_timeout"].as_str(), Some("90s")); + } + + /// The settings layer carries the server's `[run]` defaults and its + /// environment catalog as `[environments.]` tables. + #[test] + fn the_settings_layer_carries_the_environment_catalog() { + let state = crate::test_support::test_app_state(); + let text = settings_layer_toml(&state).expect("the layer serializes"); + let table: toml::Table = text.parse().expect("the layer text is TOML"); + let environments = table["environments"] + .as_table() + .expect("an environments table"); + let listed = state.environment_store().list(); + let ids: Vec<&str> = listed + .iter() + .map(|environment| environment.id.as_str()) + .map(|id| environments.contains_key(id).then_some(id)) + .map(|found| found.expect("every catalog environment is in the layer")) + .collect(); + assert!(!ids.is_empty(), "{text}"); + for id in ids { + assert!( + environments[id]["provider"].is_str(), + "`[environments.{id}]` names its provider: {text}" + ); + } + assert!(table.contains_key("run"), "{text}"); + } +} diff --git a/lib/apps/fabro-server/tests/it/api/mcp_servers.rs b/lib/apps/fabro-server/tests/it/api/mcp_servers.rs index cc6d56c59..9d1d06f73 100644 --- a/lib/apps/fabro-server/tests/it/api/mcp_servers.rs +++ b/lib/apps/fabro-server/tests/it/api/mcp_servers.rs @@ -566,42 +566,38 @@ async fn invalid_mcp_server_id_is_bad_request() { } /// A `run.agent.mcps.` entry that names a server catalog entry by -/// `id`: Petri's Fabro frontend reads `workflow.toml` itself and has no -/// server catalog to resolve the reference against, so the check refuses -/// it (`unsupported.workflow_toml.run.agent.mcps.reference`) until the -/// frontend takes the catalog. The run create path shares the gap. +/// `id` validates: the server hands Petri's Fabro frontend its catalog, so +/// the reference resolves there as it does in the server's own settings +/// resolution. An id the catalog lacks is refused by that resolution first. #[tokio::test] -async fn manifest_validation_reports_a_catalog_mcp_reference_as_unsupported() { +async fn manifest_validation_resolves_a_catalog_mcp_reference() { let (app, _temp_dir, _mcp_dir) = mcp_server_app(); create_mcp_server(&app, "sentry", "Sentry").await; - let mut manifest = minimal_manifest_json(MINIMAL_DOT); - manifest["workflows"]["workflow.fabro"]["config"] = json!({ - "path": "workflow.toml", - "source": r#" -_version = 1 + let validate = |reference: &str, expected: StatusCode| { + let mut manifest = minimal_manifest_json(MINIMAL_DOT); + manifest["workflows"]["workflow.fabro"]["config"] = json!({ + "path": "workflow.toml", + "source": format!( + "_version = 1\n\n[run.agent.mcps.sentry]\nid = \"{reference}\"\n" + ) + }); + let app = app.clone(); + async move { + let response = app + .oneshot(json_request(Method::POST, "/validate", &manifest)) + .await + .expect("manifest validation should respond"); + response_json(response, expected, "POST /api/v1/validate").await + } + }; -[run.agent.mcps.sentry] -id = "sentry" -"# - }); + let body = validate("sentry", StatusCode::OK).await; + assert_eq!(body["ok"], true, "{body}"); - let response = app - .oneshot(json_request(Method::POST, "/validate", &manifest)) - .await - .expect("manifest validation should respond"); - let body = response_json(response, StatusCode::OK, "POST /api/v1/validate").await; - - assert_eq!(body["ok"], false, "{body}"); - let rules: Vec<&str> = body["workflow"]["diagnostics"] - .as_array() - .expect("diagnostics") - .iter() - .filter_map(|diagnostic| diagnostic["rule"].as_str()) - .collect(); + let body = validate("nowhere", StatusCode::BAD_REQUEST).await; assert_eq!( - rules, - vec!["unsupported.workflow_toml.run.agent.mcps.reference"], + body["errors"][0]["detail"], "failed to resolve manifest settings", "{body}" ); } diff --git a/lib/apps/fabro-server/tests/it/scenario/petri.rs b/lib/apps/fabro-server/tests/it/scenario/petri.rs index 1c7da0bb2..2f93f82c3 100644 --- a/lib/apps/fabro-server/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-server/tests/it/scenario/petri.rs @@ -40,9 +40,9 @@ use fabro_types::{RunId, WorkflowPath, WorkflowVersion}; use tower::ServiceExt; use crate::helpers::{ - api, create_and_start_run_from_intent, read_repo_file, response_json, run_json, - settings_from_toml, test_app_state_with_options, test_app_with_scheduler, test_settings, - wait_for_run_status, + api, create_and_start_run_from_intent, minimal_manifest_json, read_repo_file, response_json, + run_json, settings_from_toml, test_app_state_with_options, test_app_with_scheduler, + test_settings, wait_for_run_status, }; const HOST_PLUGIN: &str = "sandbox-driver-host"; @@ -869,3 +869,400 @@ async fn a_runs_projection_carries_its_docker_sandbox_instance() { .status(); } } + +/// The image the server's `docker-small` environment names in the tests +/// below: a runner image with `git` for the checkpoint commit, and not the +/// plugin's default, so the container proves the catalog's image reached it. +const CATALOG_IMAGE: &str = "ghcr.io/lithoscomputer/ubuntu-22.04:slim"; + +/// A Docker environment in the server's catalog, with the image it runs. +async fn create_docker_environment(app: &axum::Router, id: &str, image: &str) { + let environment = serde_json::json!({ + "id": id, + "provider": "docker", + "image": { "docker": image, "dockerfile": null }, + "resources": { "cpu": null, "memory": null, "disk": null }, + "network": { "mode": "allow_all", "allow": [] }, + "lifecycle": { "preserve": false, "stop_on_terminal": true, "auto_stop": null }, + "labels": {}, + "env": {} + }); + let request = Request::builder() + .method("POST") + .uri(api("/environments")) + .header("content-type", "application/json") + .body(Body::from(environment.to_string())) + .expect("environment request should build"); + let response = app + .clone() + .oneshot(request) + .await + .expect("environment request routes"); + response_json( + response, + StatusCode::CREATED, + "POST /api/v1/environments".to_string(), + ) + .await; +} + +/// Create a run from `intent` without starting it: the run's id. +async fn create_run(app: &axum::Router, intent: serde_json::Value) -> String { + let request = Request::builder() + .method("POST") + .uri(api("/runs")) + .header("content-type", "application/json") + .body(Body::from(intent.to_string())) + .expect("create request should build"); + let response = app + .clone() + .oneshot(request) + .await + .expect("create request routes"); + let body = response_json(response, StatusCode::CREATED, "POST /api/v1/runs").await; + body["id"] + .as_str() + .expect("the created run's id") + .to_string() +} + +async fn start_run(app: &axum::Router, run_id: &str) { + let request = Request::builder() + .method("POST") + .uri(api(&format!("/runs/{run_id}/start"))) + .body(Body::empty()) + .expect("start request should build"); + let response = app + .clone() + .oneshot(request) + .await + .expect("start request routes"); + assert_eq!( + response.status(), + StatusCode::OK, + "POST /runs/{run_id}/start" + ); +} + +/// The root graph Petri admitted for the run, as the run's blob store holds +/// it: its `params` carry `fabro.environment` and `fabro.launch`, its nodes +/// their step configuration. +async fn admitted_root_graph(app: &axum::Router, run_id: &str) -> serde_json::Value { + let admission = run_admission(app, run_id).await; + let blob = admission["graph"]["blob"] + .as_str() + .expect("the root graph's blob hash"); + let request = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}/blobs/{blob}"))) + .body(Body::empty()) + .expect("blob request should build"); + let response = app + .clone() + .oneshot(request) + .await + .expect("blob request routes"); + assert_eq!( + response.status(), + StatusCode::OK, + "GET /runs/{run_id}/blobs/{blob}" + ); + let bytes = axum::body::to_bytes(response.into_body(), usize::MAX) + .await + .expect("the graph blob reads"); + serde_json::from_slice(&bytes).expect("the graph blob is JSON") +} + +/// A bundle that names a server environment it does not declare admits: the +/// catalog's `[environments.docker-small]` reaches Petri through the +/// settings layer, its image lands on the lowered environment, and, with +/// the Docker plugin and a daemon, the run's container runs that image. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_bundle_naming_a_catalog_environment_runs_on_docker_with_its_image() { + let settings = settings_from_toml("_version = 1\n\n[run.environment]\nid = \"local\"\n"); + let state = test_app_state_with_options(settings, 5); + let app = test_app_with_scheduler(Arc::clone(&state)); + create_docker_environment(&app, "docker-small", CATALOG_IMAGE).await; + + let version_id = register_version(&app, &[ + ("workflow.fabro", COMMAND_DOT), + ( + "workflow.toml", + "_version = 1\n\n[run.environment]\nid = \"docker-small\"\n", + ), + ]) + .await; + let intent = serde_json::json!({ + "workflow_version_id": version_id, + "target": {"kind": "none"}, + "environment_id": "docker-small", + "args": {}, + }); + let run_id = create_run(&app, intent).await; + let graph = admitted_root_graph(&app, &run_id).await; + let environment = &graph["params"]["fabro.environment"]; + assert_eq!(environment["id"], "docker-small", "{environment}"); + assert_eq!(environment["provider"], "docker", "{environment}"); + assert_eq!(environment["image"], CATALOG_IMAGE, "{environment}"); + assert_eq!( + graph["params"]["fabro.launch"]["sandbox_backend"], "docker", + "{}", + graph["params"]["fabro.launch"] + ); + + if docker_plugin().is_none() { + return; + } + start_run(&app, &run_id).await; + let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await; + let projection = settled_state(&state, &app, &run_id).await; + assert_eq!(status, "succeeded", "run: {projection}"); + let instance = &projection["sandbox"]["instance"]; + assert_eq!(instance["provider"], "docker", "{instance}"); + assert_eq!( + instance["image"], CATALOG_IMAGE, + "the container runs the catalog's image: {instance}" + ); + let output = Command::new("docker") + .args([ + "ps", + "-aq", + "--filter", + &format!("label=petri.run={run_id}"), + ]) + .output() + .expect("docker ps runs"); + for container in String::from_utf8_lossy(&output.stdout) + .lines() + .map(str::trim) + .filter(|line| !line.is_empty()) + { + let _ = Command::new("docker") + .args(["rm", "-f", container]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status(); + } +} + +/// A bundle's own `[environments.]` table wins over the server's, key +/// by key: its image replaces the catalog's on the lowered environment. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_bundles_own_environment_table_overrides_the_servers() { + let settings = settings_from_toml("_version = 1\n\n[run.environment]\nid = \"local\"\n"); + let state = test_app_state_with_options(settings, 5); + let app = test_app_with_scheduler(Arc::clone(&state)); + create_docker_environment(&app, "docker-small", CATALOG_IMAGE).await; + + let version_id = register_version(&app, &[ + ("workflow.fabro", COMMAND_DOT), + ( + "workflow.toml", + "_version = 1\n\n[run.environment]\nid = \"docker-small\"\n\n\ + [environments.docker-small]\nprovider = \"docker\"\n\n\ + [environments.docker-small.image]\ndocker = \"alpine:3.19\"\n", + ), + ]) + .await; + let intent = serde_json::json!({ + "workflow_version_id": version_id, + "target": {"kind": "none"}, + "environment_id": "docker-small", + "args": {}, + }); + let run_id = create_run(&app, intent).await; + let graph = admitted_root_graph(&app, &run_id).await; + let environment = &graph["params"]["fabro.environment"]; + assert_eq!(environment["provider"], "docker", "{environment}"); + assert_eq!( + environment["image"], "alpine:3.19", + "the bundle's image over the catalog's: {environment}" + ); +} + +/// An environment no layer declares is refused before Petri sees the +/// bundle: the server's own settings resolution refuses it at validation, +/// and an intent naming an environment the catalog lacks is refused at +/// create. Petri's own diagnostic for the same bundle is +/// `fabro-petri::check::an_unknown_environment_is_refused_and_the_launch_selects_over_the_bundle`. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn an_unknown_environment_id_is_refused_before_petri() { + let workspace = tempfile::tempdir().expect("workspace tempdir"); + let state = test_app_state_with_options(test_settings(), 5); + let app = test_app_with_scheduler(Arc::clone(&state)); + + let mut manifest = minimal_manifest_json(COMMAND_DOT); + manifest["workflows"]["workflow.fabro"]["config"] = serde_json::json!({ + "path": "workflow.toml", + "source": "_version = 1\n\n[run.environment]\nid = \"nowhere\"\n", + }); + let request = Request::builder() + .method("POST") + .uri(api("/validate")) + .header("content-type", "application/json") + .body(Body::from(manifest.to_string())) + .expect("validate request should build"); + let response = app + .clone() + .oneshot(request) + .await + .expect("validate request routes"); + let body = response_json( + response, + StatusCode::BAD_REQUEST, + "POST /api/v1/validate".to_string(), + ) + .await; + assert_eq!( + body["errors"][0]["detail"], "failed to resolve manifest settings", + "{body}" + ); + + let version_id = register_version(&app, &[ + ("workflow.fabro", COMMAND_DOT), + ("workflow.toml", PLAIN_SETTINGS), + ]) + .await; + let mut intent = intent(&version_id, workspace.path()); + intent["environment_id"] = serde_json::json!("nowhere"); + let request = Request::builder() + .method("POST") + .uri(api("/runs")) + .header("content-type", "application/json") + .body(Body::from(intent.to_string())) + .expect("create request should build"); + let response = app + .clone() + .oneshot(request) + .await + .expect("create request routes"); + let status = response.status(); + let bytes = axum::body::to_bytes(response.into_body(), usize::MAX) + .await + .expect("the response body reads"); + let text = String::from_utf8_lossy(&bytes); + assert!( + status.is_client_error() && text.contains("nowhere"), + "the server refuses an environment its catalog lacks: {status} {text}" + ); +} + +/// An agent workflow with one stage. +const AGENT_DOT: &str = r#"digraph Agent { + graph [goal="Greet with the notes server available"] + start [shape=Mdiamond] + exit [shape=Msquare] + greet [prompt="Say hello. Use no tools."] + start -> greet -> exit +}"#; + +/// A bundle referencing a catalog MCP server by id admits without declaring +/// it: the server's catalog reaches Petri, the entry lands on the agent +/// node under the reference's name, and the agent session lists the +/// server's tool to the model. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_bundle_naming_a_catalog_mcp_server_lists_its_tools_to_the_model() { + if host_plugin().is_none() { + return; + } + let workspace = tempfile::tempdir().expect("workspace tempdir"); + let twin = twin_openai().await; + let namespace = format!("{}::{}", module_path!(), line!()); + TwinScenarios::new(&namespace) + .scenario(TwinScenario::responses(OPENAI_MODEL).text("Hello.")) + .load(twin) + .await; + let settings = test_settings(); + let state = TestAppStateBuilder::new() + .runtime_settings(settings.server_settings, settings.manifest_run_defaults) + .max_concurrent_runs(5) + .in_process_execution() + .llm_overlay(llm_overlay_with_provider_base_url( + "openai", + twin.base_url.clone(), + )) + .vault_entries([(EnvVars::OPENAI_API_KEY, namespace.clone())]) + .build(); + let app = test_app_with_scheduler(Arc::clone(&state)); + + // The catalog entry: the echo server over stdio, checked in under `test/`. + let server = crate::helpers::repo_root().join("test/mcp/echo_server.py"); + let definition = serde_json::json!({ + "id": "echo-prod", + "display_name": "Echo", + "description": "The scenario tests' echo server.", + "transport": { + "type": "stdio", + "command": ["python3", server.to_string_lossy()], + "env": {} + }, + "startup_timeout_secs": 10, + "tool_timeout_secs": 60 + }); + let request = Request::builder() + .method("POST") + .uri(api("/mcp-servers")) + .header("content-type", "application/json") + .body(Body::from(definition.to_string())) + .expect("mcp server request should build"); + let response = app + .clone() + .oneshot(request) + .await + .expect("mcp server request routes"); + response_json( + response, + StatusCode::CREATED, + "POST /api/v1/mcp-servers".to_string(), + ) + .await; + + let version_id = register_version(&app, &[ + ("workflow.fabro", AGENT_DOT), + ( + "workflow.toml", + "_version = 1\n\n[run.agent.mcps.notes]\nid = \"echo-prod\"\n", + ), + ]) + .await; + let mut intent = intent(&version_id, workspace.path()); + intent["args"]["model"] = serde_json::json!(OPENAI_MODEL); + let run_id = create_and_start_run_from_intent(&app, intent).await; + + let graph = admitted_root_graph(&app, &run_id).await; + let greet = graph["nodes"] + .as_array() + .expect("the graph's nodes") + .iter() + .find(|node| node["name"] == "greet") + .unwrap_or_else(|| panic!("the greet node: {graph}")); + let mcps = &greet["step"]["config"]["mcps"]; + assert_eq!(mcps.as_array().map(Vec::len), Some(1), "{greet}"); + assert_eq!(mcps[0]["name"], "notes", "the reference's name: {mcps}"); + assert_eq!(mcps[0]["source"], "mcp-catalog:echo-prod", "{mcps}"); + assert_eq!(mcps[0]["transport"]["type"], "stdio", "{mcps}"); + + let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await; + assert_eq!( + status, + "succeeded", + "run: {}", + run_json(&app, &run_id).await + ); + // The tools the session offered the model, as Petri records them once + // per session (`attractor.tools`) and the projection lists them. + let projection = settled_state(&state, &app, &run_id).await; + let tools: Vec<&str> = projection["stages"]["greet@1"]["agent_tools"] + .as_array() + .map(|tools| { + tools + .iter() + .filter_map(|tool| tool["name"].as_str()) + .collect() + }) + .unwrap_or_default(); + assert!( + tools.contains(&"mcp__notes__echo"), + "the session lists the catalog server's tool under the reference's name: {tools:?}" + ); +} diff --git a/lib/components/fabro-petri/src/check.rs b/lib/components/fabro-petri/src/check.rs index 010952b20..6686c9aea 100644 --- a/lib/components/fabro-petri/src/check.rs +++ b/lib/components/fabro-petri/src/check.rs @@ -12,10 +12,12 @@ //! //! The launch binds the compile variables the Fabro frontend reads: //! `petri.launch_model` and `petri.launch_provider` as the model default -//! below every file layer, and `petri.repository` as the repository the root -//! `start` stage checks out. A caller with no local repository binds `null`, -//! and the run starts from an empty workspace. The server's run variables -//! (`{{ vars.NAME }}`) are bound as compile variables beside them. +//! below every file layer, `petri.launch_environment` as the environment +//! the run selected over every file layer, and `petri.repository` as the +//! repository the root `start` stage checks out. A caller with no local +//! repository binds `null`, and the run starts from an empty workspace. The +//! server's run variables (`{{ vars.NAME }}`) are bound as compile +//! variables beside them. use std::collections::BTreeMap; use std::path::PathBuf; @@ -23,7 +25,8 @@ use std::path::PathBuf; use petri_frontend_attractor::kinds::{AGENT_KIND, PROMPT_KIND}; use petri_runtime::LoadError; use petri_runtime::frontend::{ - self, CompileInputs, LAUNCH_MODEL_VAR, LAUNCH_PROVIDER_VAR, MapFiles, REPOSITORY_VAR, Severity, + self, CompileInputs, LAUNCH_ENVIRONMENT_VAR, LAUNCH_MODEL_VAR, LAUNCH_PROVIDER_VAR, MapFiles, + REPOSITORY_VAR, Severity, }; use petri_runtime::ir::Graph; use serde::{Deserialize, Serialize}; @@ -61,14 +64,20 @@ impl Bundle { } } -/// What the launch binds below the file layers. +/// What the launch binds around the file layers: the model default below +/// them, the environment selection above them, and the repository. #[derive(Clone, Debug, Default)] pub struct Launch { - pub model: Option, - pub provider: Option, + pub model: Option, + pub provider: Option, + /// The environment the run selected, by its id in the server's + /// catalog, over every layer's `[run.environment]`, as the intent's + /// selection overrides the bundle in Fabro's own resolution; `None` + /// leaves the layers to select. + pub environment: Option, /// The local repository the root `start` stage checks out into the /// workspace; `None` starts the run from an empty workspace. - pub repository: Option, + pub repository: Option, } /// One check: the bundle, the run's inputs and variables, the launch and @@ -204,6 +213,12 @@ fn compile_inputs( compile .vars .insert(LAUNCH_PROVIDER_VAR.into(), text(&launch.provider)); + if let Some(environment) = &launch.environment { + compile.vars.insert( + LAUNCH_ENVIRONMENT_VAR.into(), + Value::String(environment.clone()), + ); + } // `Runtime::check_source` uses the inputs as given, so the repository // is the host's to bind: the launch's path, or `null` for a run that // starts from an empty workspace. diff --git a/lib/components/fabro-petri/src/runtime.rs b/lib/components/fabro-petri/src/runtime.rs index ae8ba0752..b51a8a14f 100644 --- a/lib/components/fabro-petri/src/runtime.rs +++ b/lib/components/fabro-petri/src/runtime.rs @@ -32,23 +32,30 @@ use crate::host_tools; pub struct RuntimeSpec { /// The operator's settings layer, as `~/.fabro/settings.toml` text: the /// lowest of the three layers the Fabro frontend reads (`[run.model]` - /// defaults, `[[run.hooks]]`, `[run.agent.mcps]`). - pub settings_toml: Option, + /// defaults, `[[run.hooks]]`, `[run.agent.mcps]`, `[run.environment]` + /// and the `[environments.]` catalog a bundle may name). + pub settings_toml: Option, + /// The server's MCP catalog, as the TOML text the Fabro frontend + /// resolves `[run.agent.mcps.] id = "..."` references against: a + /// table keyed by catalog id, each entry in the inline + /// `[run.agent.mcps.]` shape. `None` leaves every reference + /// refused, as the standalone runner refuses it. + pub mcp_catalog_toml: Option, /// The model client the native agent and prompt steps call, and the /// catalog the admission pass resolves model selectors against. `None` /// leaves every LLM node unpinned and every model call unconfigured. - pub model_client: Option, + pub model_client: Option, /// Run the simulated step registry (Fabro's `--dry-run` handlers) /// instead of the real one. - pub dry_run: bool, + pub dry_run: bool, /// The Fabro home the skills step reads; `None` leaves it to Petri's /// own lookup (`FABRO_HOME`, else `$HOME/.fabro`). - pub fabro_home: Option, + pub fabro_home: Option, /// Fabro's run tools for every native agent session of the run, when /// the run enables them (`[run.agent] fabro_tools` and the worker /// token's `agent:run_tools` scope); `None` gives the sessions Pebble's /// tools alone. See [`crate::host_tools`]. - pub run_tools: Option, + pub run_tools: Option, } impl RuntimeSpec { @@ -57,8 +64,11 @@ impl RuntimeSpec { /// registry: only execution swaps in the stubs. #[must_use] pub fn runtime(&self, for_execution: bool) -> Runtime { - let mut runtime = Runtime::standard() - .frontend(Fabro::new().with_settings_toml(self.settings_toml.clone())); + let mut runtime = Runtime::standard().frontend( + Fabro::new() + .with_settings_toml(self.settings_toml.clone()) + .with_mcp_catalog_toml(self.mcp_catalog_toml.clone()), + ); if let Some(client) = &self.model_client { runtime = runtime.capability(PebbleClient(client.clone())); } diff --git a/lib/components/fabro-petri/tests/check.rs b/lib/components/fabro-petri/tests/check.rs index 2796f724d..9a78a5672 100644 --- a/lib/components/fabro-petri/tests/check.rs +++ b/lib/components/fabro-petri/tests/check.rs @@ -139,9 +139,10 @@ async fn a_launch_binds_the_repository_and_the_model_default() { inputs: BTreeMap::new(), vars: BTreeMap::new(), launch: Launch { - model: Some("gpt-5.4".to_string()), - provider: None, - repository: Some(repository.path().to_path_buf()), + model: Some("gpt-5.4".to_string()), + provider: None, + environment: None, + repository: Some(repository.path().to_path_buf()), }, runtime: RuntimeSpec::default(), unbound_is_warning: false, @@ -315,3 +316,54 @@ async fn a_known_model_is_pinned_at_admission() { work.step.config ); } + +/// The server's environment catalog reaches Petri as `[environments.]` +/// tables of the settings layer, and the environment the run selected as +/// the launch: a bundle naming an environment only the catalog declares +/// admits with the catalog's image; the launch's selection wins over the +/// bundle's own `[run.environment]`; an id no layer declares is refused +/// with Petri's diagnostic. +#[test] +fn an_unknown_environment_is_refused_and_the_launch_selects_over_the_bundle() { + let catalog = "[environments.local]\nprovider = \"local\"\n\ + [environments.docker-small]\nprovider = \"docker\"\n\ + [environments.docker-small.image]\ndocker = \"alpine:3.20\"\n"; + let runtime = || RuntimeSpec { + settings_toml: Some(catalog.to_string()), + ..RuntimeSpec::default() + }; + let bundle_naming = |id: &str| { + bundle(&[ + ("workflow.fabro", COMMAND_WORKFLOW), + ( + "workflow.toml", + &format!("_version = 1\n\n[run.environment]\nid = \"{id}\"\n"), + ), + ]) + }; + + let admitted = check::check(&request(bundle_naming("docker-small"), runtime())) + .expect("the catalog's environment admits"); + let environment = &admitted.graph.params["fabro.environment"]; + assert_eq!(environment["provider"], "docker"); + assert_eq!(environment["image"], "alpine:3.20"); + + let Err(CheckError::Rejected(diagnostics)) = + check::check(&request(bundle_naming("nowhere"), runtime())) + else { + panic!("an environment no layer declares should be refused"); + }; + let refusal = diagnostics + .iter() + .find(|diagnostic| diagnostic.code == "unsupported.workflow_toml.run.environment") + .unwrap_or_else(|| panic!("Petri names the unknown environment: {diagnostics:?}")); + assert!( + refusal.is_error() && refusal.message.contains("nowhere"), + "{refusal:?}" + ); + + let mut selected = request(bundle_naming("nowhere"), runtime()); + selected.launch.environment = Some("local".to_string()); + let admitted = check::check(&selected).expect("the launch's selection admits"); + assert_eq!(admitted.graph.params["fabro.environment"]["id"], "local"); +} diff --git a/test/mcp/echo_server.py b/test/mcp/echo_server.py new file mode 100755 index 000000000..1fec071c8 --- /dev/null +++ b/test/mcp/echo_server.py @@ -0,0 +1,80 @@ +#!/usr/bin/env python3 +"""A one-tool MCP server over stdio for the server's scenario tests. + +Speaks JSON-RPC 2.0, one message per line on stdin and stdout, as the MCP +specification describes for the stdio transport. It exposes ``echo(message)``, +which answers the message, so a test can see the server's tool reach an +agent session's tool list. Dependency-free. +""" + +import json +import sys + +SERVER_INFO = {"name": "fabro-test-echo", "version": "1.0.0"} +PROTOCOL_VERSION = "2025-03-26" + +TOOLS = [ + { + "name": "echo", + "description": "Echo back the message", + "inputSchema": { + "type": "object", + "properties": {"message": {"type": "string"}}, + "required": ["message"], + }, + } +] + + +def handle(request): + """Answer one request, or ``None`` for a notification.""" + method = request.get("method") + request_id = request.get("id") + params = request.get("params") or {} + if method == "initialize": + return { + "jsonrpc": "2.0", + "id": request_id, + "result": { + "protocolVersion": params.get("protocolVersion", PROTOCOL_VERSION), + "capabilities": {"tools": {}}, + "serverInfo": SERVER_INFO, + }, + } + if method == "ping": + return {"jsonrpc": "2.0", "id": request_id, "result": {}} + if method == "tools/list": + return {"jsonrpc": "2.0", "id": request_id, "result": {"tools": TOOLS}} + if method == "tools/call": + message = (params.get("arguments") or {}).get("message", "") + return { + "jsonrpc": "2.0", + "id": request_id, + "result": {"content": [{"type": "text", "text": message}]}, + } + if request_id is None: + return None + return { + "jsonrpc": "2.0", + "id": request_id, + "error": {"code": -32601, "message": f"Method not found: {method}"}, + } + + +def main(): + for line in sys.stdin: + line = line.strip() + if not line: + continue + try: + request = json.loads(line) + except json.JSONDecodeError: + continue + response = handle(request) + if response is not None: + sys.stdout.write(json.dumps(response) + "\n") + sys.stdout.flush() + + +if __name__ == "__main__": + main()