diff --git a/Cargo.lock b/Cargo.lock index c66a6239a..3c48e0aa9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2699,6 +2699,18 @@ dependencies = [ "toml 0.8.23", ] +[[package]] +name = "fabro-mcp-store" +version = "0.267.0-nightly.0" +dependencies = [ + "fabro-types", + "serde", + "tempfile", + "thiserror 2.0.18", + "tokio", + "toml 0.8.23", +] + [[package]] name = "fabro-model" version = "0.275.0-nightly.0" diff --git a/lib/crates/fabro-mcp-store/Cargo.toml b/lib/crates/fabro-mcp-store/Cargo.toml new file mode 100644 index 000000000..2fe21ece3 --- /dev/null +++ b/lib/crates/fabro-mcp-store/Cargo.toml @@ -0,0 +1,24 @@ +[package] +name = "fabro-mcp-store" +edition.workspace = true +version.workspace = true +publish = false +license.workspace = true +description = "Server-managed MCP server catalog durable storage for Fabro" + +[lib] +doctest = false + +[lints] +workspace = true + +[dependencies] +fabro-types = { path = "../fabro-types" } +serde.workspace = true +thiserror.workspace = true +tokio.workspace = true +toml.workspace = true + +[dev-dependencies] +tempfile = "3" +tokio = { workspace = true, features = ["macros", "test-util"] } diff --git a/lib/crates/fabro-mcp-store/src/error.rs b/lib/crates/fabro-mcp-store/src/error.rs new file mode 100644 index 000000000..f4956b18d --- /dev/null +++ b/lib/crates/fabro-mcp-store/src/error.rs @@ -0,0 +1,86 @@ +use std::path::PathBuf; + +use fabro_types::{McpServerId, McpServerRevision, McpServerValidationError}; +use toml::de::Error as TomlDeError; +use toml::ser::Error as TomlSerError; + +#[derive(Debug, thiserror::Error)] +pub enum McpServerStoreError { + #[error("mcp server not found: {id}")] + NotFound { id: McpServerId }, + #[error("mcp server already exists: {id}")] + AlreadyExists { id: McpServerId }, + #[error("mcp server revision is stale for {id}: expected {expected}, actual {actual}")] + StaleRevision { + id: McpServerId, + expected: McpServerRevision, + actual: McpServerRevision, + }, + #[error("mcp server validation failed")] + Validation { + #[from] + source: McpServerValidationError, + }, + #[error("invalid mcp server filename at {path:?}")] + InvalidFilename { path: PathBuf, reason: String }, + #[error("failed to parse mcp server TOML at {path:?}")] + Parse { + path: PathBuf, + #[source] + source: TomlDeError, + }, + #[error("mcp server TOML at {path:?} is not UTF-8")] + InvalidUtf8 { + path: PathBuf, + #[source] + source: std::str::Utf8Error, + }, + #[error("failed to serialize mcp server TOML")] + Serialize { + #[from] + source: TomlSerError, + }, + #[error("I/O error at {path:?}")] + Io { + path: PathBuf, + #[source] + source: std::io::Error, + }, +} + +impl McpServerStoreError { + pub(crate) fn io(path: impl Into, source: std::io::Error) -> Self { + Self::Io { + path: path.into(), + source, + } + } + + pub(crate) fn parse(path: impl Into, source: TomlDeError) -> Self { + Self::Parse { + path: path.into(), + source, + } + } + + pub(crate) fn invalid_utf8(path: impl Into, source: std::str::Utf8Error) -> Self { + Self::InvalidUtf8 { + path: path.into(), + source, + } + } + + #[must_use] + pub fn kind(&self) -> &'static str { + match self { + Self::NotFound { .. } => "not_found", + Self::AlreadyExists { .. } => "already_exists", + Self::StaleRevision { .. } => "stale_revision", + Self::Validation { .. } => "validation", + Self::InvalidFilename { .. } => "invalid_filename", + Self::Parse { .. } | Self::InvalidUtf8 { .. } => "parse", + Self::Serialize { .. } => "serialize", + Self::Io { .. } => "io", + } + } +} diff --git a/lib/crates/fabro-mcp-store/src/lib.rs b/lib/crates/fabro-mcp-store/src/lib.rs new file mode 100644 index 000000000..7000a5fc5 --- /dev/null +++ b/lib/crates/fabro-mcp-store/src/lib.rs @@ -0,0 +1,14 @@ +//! Durable storage for server-managed MCP server definitions. +//! +//! Concrete [`McpServerStore`] modeled on `fabro-automation`'s +//! `AutomationStore`: per-file TOML under `{config}/mcps/{id}.toml`, in-memory +//! cache, SHA-256 revision for optimistic concurrency, and async +//! storage-agnostic methods. The domain model lives in `fabro-types`; this +//! crate owns persistence. + +mod error; +mod model; +mod store; + +pub use error::McpServerStoreError; +pub use store::McpServerStore; diff --git a/lib/crates/fabro-mcp-store/src/model.rs b/lib/crates/fabro-mcp-store/src/model.rs new file mode 100644 index 000000000..e2564284b --- /dev/null +++ b/lib/crates/fabro-mcp-store/src/model.rs @@ -0,0 +1,123 @@ +//! Construction and persistence glue for [`McpServerDefinition`]. +//! +//! The domain types (`McpServerDefinition`, `McpServerDraft`, +//! `McpServerReplace`, `McpServerId`, `McpServerRevision`) live in +//! `fabro-types` so they stay persistence-independent. This module owns the +//! store-side glue: validating, serializing to canonical TOML bytes, deriving +//! the revision, and reconstructing definitions from persisted bytes. + +use std::path::PathBuf; + +use fabro_types::settings::McpTransport; +use fabro_types::{ + McpServerDefinition, McpServerId, McpServerReplace, McpServerRevision, mcp_store, +}; +use serde::{Deserialize, Serialize}; + +use crate::error::McpServerStoreError; + +/// The on-disk body of a definition. Excludes `id`/`revision`, which are +/// derived from the filename and content hash rather than persisted. +#[derive(Debug, Clone, PartialEq, Deserialize)] +#[serde(deny_unknown_fields)] +struct PersistedMcpServer { + name: String, + #[serde(default)] + description: Option, + transport: McpTransport, + startup_timeout_secs: u64, + tool_timeout_secs: u64, +} + +#[derive(Serialize)] +struct PersistedMcpServerRef<'a> { + name: &'a str, + #[serde(default, skip_serializing_if = "Option::is_none")] + description: Option<&'a str>, + transport: &'a McpTransport, + startup_timeout_secs: u64, + tool_timeout_secs: u64, +} + +impl<'a> From<&'a McpServerReplace> for PersistedMcpServerRef<'a> { + fn from(value: &'a McpServerReplace) -> Self { + Self { + name: &value.name, + description: value.description.as_deref(), + transport: &value.transport, + startup_timeout_secs: value.startup_timeout_secs, + tool_timeout_secs: value.tool_timeout_secs, + } + } +} + +impl From for McpServerReplace { + fn from(value: PersistedMcpServer) -> Self { + Self { + name: value.name, + description: value.description, + transport: value.transport, + startup_timeout_secs: value.startup_timeout_secs, + tool_timeout_secs: value.tool_timeout_secs, + } + } +} + +/// Build a definition + its canonical persisted bytes from a replace payload. +/// +/// The revision is the SHA-256 of the freshly serialized canonical bytes, so a +/// caller can compare it to the on-disk content hash for optimistic +/// concurrency. +pub(crate) fn definition_from_replace( + id: McpServerId, + replace: McpServerReplace, +) -> Result<(McpServerDefinition, Vec), McpServerStoreError> { + mcp_store::validate_mcp_server_fields(&replace)?; + let bytes = canonical_bytes(&replace)?; + let revision = McpServerRevision::from_bytes(&bytes); + let definition = assemble(id, revision, replace); + Ok((definition, bytes)) +} + +/// Reconstruct a definition from bytes loaded off disk, deriving the revision +/// from the raw file bytes (not a re-serialization). +pub(crate) fn definition_from_persisted_path( + id: McpServerId, + bytes: &[u8], + path: impl Into, +) -> Result { + let path = path.into(); + let revision = McpServerRevision::from_bytes(bytes); + let persisted = parse_persisted(bytes, path)?; + let replace = McpServerReplace::from(persisted); + mcp_store::validate_mcp_server_fields(&replace)?; + Ok(assemble(id, revision, replace)) +} + +fn assemble( + id: McpServerId, + revision: McpServerRevision, + replace: McpServerReplace, +) -> McpServerDefinition { + McpServerDefinition { + id, + revision, + name: replace.name, + description: replace.description, + transport: replace.transport, + startup_timeout_secs: replace.startup_timeout_secs, + tool_timeout_secs: replace.tool_timeout_secs, + } +} + +pub(crate) fn canonical_bytes(replace: &McpServerReplace) -> Result, McpServerStoreError> { + let persisted = PersistedMcpServerRef::from(replace); + let toml = toml::to_string_pretty(&persisted)?; + Ok(toml.into_bytes()) +} + +fn parse_persisted(bytes: &[u8], path: PathBuf) -> Result { + let content = std::str::from_utf8(bytes) + .map_err(|err| McpServerStoreError::invalid_utf8(path.clone(), err))?; + toml::from_str(content).map_err(|err| McpServerStoreError::parse(path, err)) +} diff --git a/lib/crates/fabro-mcp-store/src/store.rs b/lib/crates/fabro-mcp-store/src/store.rs new file mode 100644 index 000000000..2440b5c5b --- /dev/null +++ b/lib/crates/fabro-mcp-store/src/store.rs @@ -0,0 +1,461 @@ +use std::collections::HashMap; +use std::io::ErrorKind; +use std::path::{Path, PathBuf}; +use std::time::{SystemTime, UNIX_EPOCH}; + +use fabro_types::{ + McpServerDefinition, McpServerDraft, McpServerId, McpServerReplace, McpServerRevision, +}; +use tokio::fs; +use tokio::io::AsyncWriteExt as _; +use tokio::sync::{Mutex, RwLock}; + +use crate::error::McpServerStoreError; +use crate::model; + +/// Durable per-file TOML store for server-managed MCP server definitions. +/// +/// Concrete by design (no trait): a future FS→SQL move is a one-time migration, +/// not a runtime backend choice. The migration seam is the async, storage- +/// agnostic method surface; only [`McpServerStore::load`] knows about the +/// filesystem. +#[derive(Debug)] +pub struct McpServerStore { + dir: PathBuf, + mutations: Mutex<()>, + defs: RwLock>, +} + +impl McpServerStore { + /// Synchronously load every persisted definition in `dir`. Returns an error + /// if any file fails to parse or validate; the caller decides startup + /// failure policy. Synchronous because it runs once at construction time + /// (typically during server startup) and is invoked from non-async code. + pub fn load(dir: impl Into) -> Result { + let dir = dir.into(); + let defs = load_definitions(&dir)?; + Ok(Self { + dir, + mutations: Mutex::new(()), + defs: RwLock::new(defs), + }) + } + + pub async fn list(&self) -> Vec { + let defs = self.defs.read().await; + let mut values = defs.values().cloned().collect::>(); + values.sort_by(|left, right| left.id.cmp(&right.id)); + values + } + + /// Sorted ids only, without cloning the (potentially sensitive) env/header + /// maps carried by full definitions. Used by missing-reference errors to + /// list available ids cheaply. + pub async fn ids(&self) -> Vec { + let defs = self.defs.read().await; + let mut ids = defs.keys().cloned().collect::>(); + ids.sort(); + ids + } + + pub async fn get(&self, id: &McpServerId) -> Option { + self.defs.read().await.get(id).cloned() + } + + pub async fn create( + &self, + draft: McpServerDraft, + ) -> Result { + let (id, replace) = draft.into(); + let _mutation = self.mutations.lock().await; + if self.defs.read().await.contains_key(&id) { + return Err(McpServerStoreError::AlreadyExists { id }); + } + let (definition, bytes) = model::definition_from_replace(id.clone(), replace)?; + + let path = definition_path(&self.dir, &id); + write_new(&self.dir, &path, &bytes) + .await + .map_err(|err| create_error_for(id.clone(), err))?; + + let mut defs = self.defs.write().await; + defs.insert(id, definition.clone()); + Ok(definition) + } + + pub async fn replace( + &self, + id: &McpServerId, + expected: &McpServerRevision, + replace: McpServerReplace, + ) -> Result { + let _mutation = self.mutations.lock().await; + { + let defs = self.defs.read().await; + check_revision(&defs, id, expected)?; + } + let (definition, bytes) = model::definition_from_replace(id.clone(), replace)?; + + write_atomic(&self.dir, &definition_path(&self.dir, id), &bytes).await?; + let mut defs = self.defs.write().await; + defs.insert(id.clone(), definition.clone()); + Ok(definition) + } + + pub async fn delete( + &self, + id: &McpServerId, + expected: &McpServerRevision, + ) -> Result<(), McpServerStoreError> { + let _mutation = self.mutations.lock().await; + { + let defs = self.defs.read().await; + check_revision(&defs, id, expected)?; + } + + let path = definition_path(&self.dir, id); + fs::remove_file(&path) + .await + .map_err(|err| McpServerStoreError::io(path, err))?; + let mut defs = self.defs.write().await; + defs.remove(id); + Ok(()) + } +} + +fn check_revision( + defs: &HashMap, + id: &McpServerId, + expected: &McpServerRevision, +) -> Result<(), McpServerStoreError> { + let current = defs + .get(id) + .ok_or_else(|| McpServerStoreError::NotFound { id: id.clone() })?; + if ¤t.revision != expected { + return Err(McpServerStoreError::StaleRevision { + id: id.clone(), + expected: expected.clone(), + actual: current.revision.clone(), + }); + } + Ok(()) +} + +#[expect( + clippy::disallowed_methods, + reason = "MCP server directory scan runs once at startup, before the runtime needs to make progress; std::fs avoids needing a Tokio runtime for the caller." +)] +fn load_definitions( + dir: &Path, +) -> Result, McpServerStoreError> { + let entries = match std::fs::read_dir(dir) { + Ok(entries) => entries, + Err(err) if err.kind() == ErrorKind::NotFound => return Ok(HashMap::new()), + Err(err) => return Err(McpServerStoreError::io(dir, err)), + }; + + let mut defs = HashMap::new(); + for entry in entries { + let entry = entry.map_err(|err| McpServerStoreError::io(dir, err))?; + let path = entry.path(); + let file_type = entry + .file_type() + .map_err(|err| McpServerStoreError::io(&path, err))?; + if !file_type.is_file() || !is_toml_file(&path) { + continue; + } + let definition = load_definition_file(&path)?; + defs.insert(definition.id.clone(), definition); + } + Ok(defs) +} + +#[expect( + clippy::disallowed_methods, + reason = "Sync sibling of `load_definitions`; only invoked from the synchronous startup load path." +)] +fn load_definition_file(path: &Path) -> Result { + let id = id_from_path(path)?; + let bytes = std::fs::read(path).map_err(|err| McpServerStoreError::io(path, err))?; + model::definition_from_persisted_path(id, &bytes, path) +} + +fn id_from_path(path: &Path) -> Result { + let stem = path + .file_stem() + .and_then(|stem| stem.to_str()) + .ok_or_else(|| McpServerStoreError::InvalidFilename { + path: path.to_path_buf(), + reason: "filename is not valid UTF-8".to_string(), + })?; + McpServerId::new(stem).map_err(|source| McpServerStoreError::InvalidFilename { + path: path.to_path_buf(), + reason: source.to_string(), + }) +} + +fn is_toml_file(path: &Path) -> bool { + path.extension() + .and_then(|extension| extension.to_str()) + .is_some_and(|extension| extension == "toml") +} + +async fn write_atomic(dir: &Path, path: &Path, bytes: &[u8]) -> Result<(), McpServerStoreError> { + let temp_path = write_temp_file(dir, path, bytes).await?; + if let Err(err) = fs::rename(&temp_path, path).await { + cleanup_temp(&temp_path).await; + return Err(McpServerStoreError::io(path, err)); + } + + Ok(()) +} + +async fn write_new(dir: &Path, path: &Path, bytes: &[u8]) -> Result<(), McpServerStoreError> { + let temp_path = write_temp_file(dir, path, bytes).await?; + if let Err(err) = fs::hard_link(&temp_path, path).await { + cleanup_temp(&temp_path).await; + return Err(McpServerStoreError::io(path, err)); + } + cleanup_temp(&temp_path).await; + Ok(()) +} + +async fn write_temp_file( + dir: &Path, + path: &Path, + bytes: &[u8], +) -> Result { + fs::create_dir_all(dir) + .await + .map_err(|err| McpServerStoreError::io(dir, err))?; + let temp_path = temp_path_for(path); + let mut file = fs::OpenOptions::new() + .write(true) + .create_new(true) + .open(&temp_path) + .await + .map_err(|err| McpServerStoreError::io(&temp_path, err))?; + + if let Err(err) = file.write_all(bytes).await { + cleanup_temp(&temp_path).await; + return Err(McpServerStoreError::io(&temp_path, err)); + } + if let Err(err) = file.sync_all().await { + cleanup_temp(&temp_path).await; + return Err(McpServerStoreError::io(&temp_path, err)); + } + drop(file); + + Ok(temp_path) +} + +async fn cleanup_temp(path: &Path) { + let _ = fs::remove_file(path).await; +} + +fn create_error_for(id: McpServerId, err: McpServerStoreError) -> McpServerStoreError { + match err { + McpServerStoreError::Io { source, .. } if source.kind() == ErrorKind::AlreadyExists => { + McpServerStoreError::AlreadyExists { id } + } + err => err, + } +} + +fn temp_path_for(path: &Path) -> PathBuf { + let parent = path.parent().unwrap_or_else(|| Path::new(".")); + let file_name = path + .file_name() + .and_then(|name| name.to_str()) + .unwrap_or("mcp-server.toml"); + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_or(0, |duration| duration.as_nanos()); + parent.join(format!(".{file_name}.{}.{}.tmp", std::process::id(), now)) +} + +fn definition_path(dir: &Path, id: &McpServerId) -> PathBuf { + dir.join(format!("{id}.toml")) +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + + use fabro_types::settings::McpTransport; + use fabro_types::settings::run::McpHttpProtocol; + use fabro_types::{McpServerDraft, McpServerId, McpServerReplace, McpServerRevision}; + use tokio::fs; + + use crate::error::McpServerStoreError; + use crate::store::McpServerStore; + + fn http_transport(url: &str) -> McpTransport { + McpTransport::Http { + protocol: McpHttpProtocol::default(), + url: url.to_string(), + headers: HashMap::new(), + } + } + + fn draft(id: &str, name: &str) -> McpServerDraft { + McpServerDraft { + id: McpServerId::new(id).unwrap(), + name: name.to_string(), + description: None, + transport: http_transport("https://example.com/mcp"), + startup_timeout_secs: 10, + tool_timeout_secs: 60, + } + } + + fn replacement(name: &str) -> McpServerReplace { + McpServerReplace { + name: name.to_string(), + description: Some("updated".to_string()), + transport: http_transport("https://example.com/mcp/v2"), + startup_timeout_secs: 15, + tool_timeout_secs: 90, + } + } + + #[tokio::test] + async fn missing_directory_loads_empty_store() { + let dir = tempfile::tempdir().unwrap(); + let store = McpServerStore::load(dir.path().join("mcps")).unwrap(); + + assert!(store.list().await.is_empty()); + assert!(store.ids().await.is_empty()); + } + + #[tokio::test] + async fn load_ignores_non_toml_files_and_keeps_valid_definitions() { + let dir = tempfile::tempdir().unwrap(); + let mcp_dir = dir.path().join("mcps"); + fs::create_dir_all(&mcp_dir).await.unwrap(); + fs::write(mcp_dir.join("notes.txt"), "ignore") + .await + .unwrap(); + fs::write( + mcp_dir.join("sentry.toml"), + r#" +name = "Sentry" +startup_timeout_secs = 10 +tool_timeout_secs = 60 + +[transport] +type = "http" +url = "https://sentry.example.com/mcp" + +[transport.headers] +"#, + ) + .await + .unwrap(); + + let store = McpServerStore::load(&mcp_dir).unwrap(); + let defs = store.list().await; + + assert_eq!(defs.len(), 1); + assert_eq!(defs[0].id.as_str(), "sentry"); + assert_eq!(defs[0].name, "Sentry"); + } + + #[tokio::test] + async fn load_fails_on_malformed_toml() { + let dir = tempfile::tempdir().unwrap(); + let mcp_dir = dir.path().join("mcps"); + fs::create_dir_all(&mcp_dir).await.unwrap(); + fs::write(mcp_dir.join("broken.toml"), "not valid toml =") + .await + .unwrap(); + + let err = McpServerStore::load(&mcp_dir).unwrap_err(); + assert!(matches!(err, McpServerStoreError::Parse { .. })); + } + + #[tokio::test] + async fn load_fails_on_invalid_filename_id() { + let dir = tempfile::tempdir().unwrap(); + let mcp_dir = dir.path().join("mcps"); + fs::create_dir_all(&mcp_dir).await.unwrap(); + fs::write(mcp_dir.join("Bad Name.toml"), "name = \"Bad\"") + .await + .unwrap(); + + let err = McpServerStore::load(&mcp_dir).unwrap_err(); + assert!(matches!(err, McpServerStoreError::InvalidFilename { .. })); + } + + #[tokio::test] + async fn create_get_list_replace_and_delete_round_trip_files_and_revisions() { + let dir = tempfile::tempdir().unwrap(); + let mcp_dir = dir.path().join("mcps"); + let store = McpServerStore::load(&mcp_dir).unwrap(); + + let created = store.create(draft("sentry", "Sentry")).await.unwrap(); + let path = mcp_dir.join("sentry.toml"); + let persisted = fs::read_to_string(&path).await.unwrap(); + assert!(persisted.contains("name = \"Sentry\"")); + assert!(!top_level_lines(&persisted).any(|line| line.starts_with("id = "))); + assert!(!top_level_lines(&persisted).any(|line| line.starts_with("revision = "))); + assert_eq!( + created.revision, + McpServerRevision::from_bytes(persisted.as_bytes()) + ); + + assert_eq!(store.get(&created.id).await.unwrap(), created); + let listed = store.list().await; + assert_eq!(listed.len(), 1); + assert_eq!(store.ids().await, vec![created.id.clone()]); + + let replaced = store + .replace(&created.id, &created.revision, replacement("Sentry v2")) + .await + .unwrap(); + assert_ne!(replaced.revision, created.revision); + assert_eq!(replaced.name, "Sentry v2"); + assert_eq!( + store.get(&created.id).await.unwrap().revision, + replaced.revision + ); + + store.delete(&created.id, &replaced.revision).await.unwrap(); + assert!(store.get(&created.id).await.is_none()); + assert!(!path.exists()); + } + + #[tokio::test] + async fn replace_with_stale_revision_is_rejected() { + let dir = tempfile::tempdir().unwrap(); + let store = McpServerStore::load(dir.path().join("mcps")).unwrap(); + let created = store.create(draft("sentry", "Sentry")).await.unwrap(); + + let stale = McpServerRevision::from_bytes(b"stale"); + let err = store + .replace(&created.id, &stale, replacement("Updated")) + .await + .unwrap_err(); + assert!(matches!(err, McpServerStoreError::StaleRevision { .. })); + + // The on-disk and in-memory definition is unchanged after a rejected replace. + assert_eq!(store.get(&created.id).await.unwrap(), created); + } + + #[tokio::test] + async fn duplicate_create_is_rejected() { + let dir = tempfile::tempdir().unwrap(); + let store = McpServerStore::load(dir.path().join("mcps")).unwrap(); + store.create(draft("sentry", "Sentry")).await.unwrap(); + + let err = store + .create(draft("sentry", "Duplicate")) + .await + .unwrap_err(); + assert!(matches!(err, McpServerStoreError::AlreadyExists { .. })); + } + + fn top_level_lines(toml: &str) -> impl Iterator { + toml.lines().take_while(|line| !line.starts_with('[')) + } +} diff --git a/lib/crates/fabro-types/src/lib.rs b/lib/crates/fabro-types/src/lib.rs index ce31e5139..7edf91066 100644 --- a/lib/crates/fabro-types/src/lib.rs +++ b/lib/crates/fabro-types/src/lib.rs @@ -16,6 +16,7 @@ mod id; pub mod interview; pub mod llm_backend; pub mod manifest_path; +pub mod mcp_store; pub mod outcome; pub mod pair; pub mod principal; @@ -75,6 +76,10 @@ pub use graph::{ pub use interview::{InterviewQuestionRecord, QuestionType}; pub use llm_backend::AgentBackend; pub use manifest_path::{ManifestPath, ManifestPathParseError}; +pub use mcp_store::{ + McpServerDefinition, McpServerDraft, McpServerId, McpServerReplace, McpServerRevision, + McpServerRevisionParseError, McpServerValidationError, validate_mcp_server_fields, +}; pub use outcome::{ FailureCategory, FailureDetail, NodeResult, Outcome, OutcomeMeta, StageOutcome, StageState, }; diff --git a/lib/crates/fabro-types/src/mcp_store.rs b/lib/crates/fabro-types/src/mcp_store.rs new file mode 100644 index 000000000..3efe0ac34 --- /dev/null +++ b/lib/crates/fabro-types/src/mcp_store.rs @@ -0,0 +1,372 @@ +//! Server-managed MCP server catalog domain model. +//! +//! These types describe MCP server definitions that are stored once on a Fabro +//! server and later referenced by name from workflow configs. They are +//! persistence-independent: the durable storage lives in the `fabro-mcp-store` +//! crate, which derives `id` (filename stem) and `revision` (content hash) and +//! never persists them inside the TOML body. +//! +//! Transport is the existing [`McpTransport`](crate::settings::McpTransport) +//! reused verbatim, so a stored definition uses the same `stdio`/`http`/ +//! `sandbox` shape as inline MCP config. + +use std::fmt; +use std::str::FromStr; + +use serde::de::Error as _; +use serde::{Deserialize, Deserializer, Serialize, Serializer}; +use sha2::{Digest, Sha256}; + +use crate::settings::McpTransport; + +/// A server-managed MCP server definition. +/// +/// `id` and `revision` are derived (filename stem + content hash of the +/// persisted TOML bytes) and are not stored in the file body. +#[derive(Debug, Clone, PartialEq)] +pub struct McpServerDefinition { + pub id: McpServerId, + pub revision: McpServerRevision, + pub name: String, + pub description: Option, + pub transport: McpTransport, + pub startup_timeout_secs: u64, + pub tool_timeout_secs: u64, +} + +/// Fields supplied when creating a new definition. Carries an `id` (the create +/// call assigns the filename) but no `revision` (the store derives it). +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct McpServerDraft { + pub id: McpServerId, + pub name: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub description: Option, + pub transport: McpTransport, + pub startup_timeout_secs: u64, + pub tool_timeout_secs: u64, +} + +/// Fields supplied when replacing an existing definition. The id is fixed by +/// the path and the revision is supplied separately for optimistic concurrency, +/// so neither appears here. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct McpServerReplace { + pub name: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub description: Option, + pub transport: McpTransport, + pub startup_timeout_secs: u64, + pub tool_timeout_secs: u64, +} + +impl From for (McpServerId, McpServerReplace) { + fn from(value: McpServerDraft) -> Self { + (value.id, McpServerReplace { + name: value.name, + description: value.description, + transport: value.transport, + startup_timeout_secs: value.startup_timeout_secs, + tool_timeout_secs: value.tool_timeout_secs, + }) + } +} + +/// Validation errors for MCP server domain fields. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum McpServerValidationError { + InvalidMcpServerId { value: String }, + EmptyName, + InvalidTransport { reason: String }, +} + +impl fmt::Display for McpServerValidationError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::InvalidMcpServerId { value } => { + write!( + f, + "mcp server id {value:?} must match [a-z0-9][a-z0-9-]{{0,62}}" + ) + } + Self::EmptyName => f.write_str("mcp server name must not be empty"), + Self::InvalidTransport { reason } => { + write!(f, "mcp server transport is invalid: {reason}") + } + } + } +} + +impl std::error::Error for McpServerValidationError {} + +/// An MCP server id: lowercase, matches `^[a-z0-9][a-z0-9-]{0,62}$`, and equals +/// the persisted file's stem. +#[derive(Clone, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct McpServerId(String); + +impl McpServerId { + pub fn new(value: impl Into) -> Result { + let value = value.into(); + if is_valid_mcp_server_id(&value) { + Ok(Self(value)) + } else { + Err(McpServerValidationError::InvalidMcpServerId { value }) + } + } + + #[must_use] + pub fn as_str(&self) -> &str { + &self.0 + } +} + +impl fmt::Display for McpServerId { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(self.as_str()) + } +} + +impl FromStr for McpServerId { + type Err = McpServerValidationError; + + fn from_str(value: &str) -> Result { + Self::new(value) + } +} + +impl Serialize for McpServerId { + fn serialize(&self, serializer: S) -> Result + where + S: Serializer, + { + serializer.serialize_str(self.as_str()) + } +} + +impl<'de> Deserialize<'de> for McpServerId { + fn deserialize(deserializer: D) -> Result + where + D: Deserializer<'de>, + { + let value = String::deserialize(deserializer)?; + value.parse().map_err(D::Error::custom) + } +} + +/// A revision: the lowercase SHA-256 hex of a definition's canonical persisted +/// TOML bytes. Used as an ETag for optimistic concurrency. +#[derive(Clone, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct McpServerRevision(String); + +impl McpServerRevision { + #[must_use] + pub fn from_bytes(bytes: &[u8]) -> Self { + Self(hex::encode(Sha256::digest(bytes))) + } + + #[must_use] + pub fn as_str(&self) -> &str { + &self.0 + } +} + +impl fmt::Display for McpServerRevision { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(self.as_str()) + } +} + +impl FromStr for McpServerRevision { + type Err = McpServerRevisionParseError; + + fn from_str(value: &str) -> Result { + if value.len() == 64 + && value + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) + { + Ok(Self(value.to_string())) + } else { + Err(McpServerRevisionParseError(value.to_string())) + } + } +} + +impl Serialize for McpServerRevision { + fn serialize(&self, serializer: S) -> Result + where + S: Serializer, + { + serializer.serialize_str(self.as_str()) + } +} + +impl<'de> Deserialize<'de> for McpServerRevision { + fn deserialize(deserializer: D) -> Result + where + D: Deserializer<'de>, + { + let value = String::deserialize(deserializer)?; + value.parse().map_err(D::Error::custom) + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct McpServerRevisionParseError(String); + +impl fmt::Display for McpServerRevisionParseError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "invalid mcp server revision: {:?}", self.0) + } +} + +impl std::error::Error for McpServerRevisionParseError {} + +fn is_valid_mcp_server_id(value: &str) -> bool { + let mut bytes = value.bytes(); + let Some(first) = bytes.next() else { + return false; + }; + if !first.is_ascii_lowercase() && !first.is_ascii_digit() { + return false; + } + if value.len() > 63 { + return false; + } + bytes.all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-') +} + +/// Validate the structural invariants of a definition's fields. +/// +/// Scope is intentionally structural for now: id format (enforced by +/// [`McpServerId`]), non-empty name, and a well-formed transport. It does not +/// reject credential-looking literal values in env vars or HTTP headers. +pub fn validate_mcp_server_fields( + replace: &McpServerReplace, +) -> Result<(), McpServerValidationError> { + if replace.name.trim().is_empty() { + return Err(McpServerValidationError::EmptyName); + } + validate_transport(&replace.transport) +} + +fn validate_transport(transport: &McpTransport) -> Result<(), McpServerValidationError> { + match transport { + McpTransport::Stdio { command, .. } | McpTransport::Sandbox { command, .. } => { + if command + .first() + .is_none_or(|program| program.trim().is_empty()) + { + return Err(McpServerValidationError::InvalidTransport { + reason: "command program must not be empty".to_string(), + }); + } + } + McpTransport::Http { url, .. } => { + if url.trim().is_empty() { + return Err(McpServerValidationError::InvalidTransport { + reason: "url must not be empty".to_string(), + }); + } + } + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + + use super::{ + McpServerId, McpServerReplace, McpServerRevision, McpTransport, validate_mcp_server_fields, + }; + use crate::settings::run::McpHttpProtocol; + + fn http_transport() -> McpTransport { + McpTransport::Http { + protocol: McpHttpProtocol::default(), + url: "https://example.com/mcp".to_string(), + headers: HashMap::new(), + } + } + + #[test] + fn mcp_server_id_validation_matches_contract() { + assert!("a".parse::().is_ok()); + assert!("a-1".parse::().is_ok()); + assert!("0".parse::().is_ok()); + assert!("sentry-dev".parse::().is_ok()); + assert!("A".parse::().is_err()); + assert!("a_1".parse::().is_err()); + assert!("-a".parse::().is_err()); + assert!("".parse::().is_err()); + assert!("a".repeat(64).parse::().is_err()); + } + + #[test] + fn revision_is_lowercase_sha256_hex() { + let revision = McpServerRevision::from_bytes(b"hello"); + assert_eq!( + revision.to_string(), + "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824" + ); + assert!(revision.to_string().parse::().is_ok()); + assert!("ABC".parse::().is_err()); + } + + #[test] + fn validation_rejects_empty_name() { + let replace = McpServerReplace { + name: " ".to_string(), + description: None, + transport: http_transport(), + startup_timeout_secs: 10, + tool_timeout_secs: 60, + }; + assert!(validate_mcp_server_fields(&replace).is_err()); + } + + #[test] + fn validation_rejects_empty_transport_command() { + let replace = McpServerReplace { + name: "Local".to_string(), + description: None, + transport: McpTransport::Stdio { + command: Vec::new(), + env: HashMap::new(), + }, + startup_timeout_secs: 10, + tool_timeout_secs: 60, + }; + assert!(validate_mcp_server_fields(&replace).is_err()); + } + + #[test] + fn validation_rejects_blank_transport_program() { + let replace = McpServerReplace { + name: "Local".to_string(), + description: None, + transport: McpTransport::Stdio { + command: vec![" ".to_string(), "--arg".to_string()], + env: HashMap::new(), + }, + startup_timeout_secs: 10, + tool_timeout_secs: 60, + }; + assert!(validate_mcp_server_fields(&replace).is_err()); + } + + #[test] + fn validation_accepts_well_formed_definition() { + let replace = McpServerReplace { + name: "Sentry".to_string(), + description: Some("Issue tracker".to_string()), + transport: http_transport(), + startup_timeout_secs: 10, + tool_timeout_secs: 60, + }; + assert!(validate_mcp_server_fields(&replace).is_ok()); + } +}