feat(mcp): add server-side MCP server store (fabro-mcp-store) (#521)

## What

Adds the storage foundation for server-managed MCP servers: a durable
store plus its domain model. No server wiring, HTTP API, or UI yet —
this is standalone scaffolding that later PRs build on.

- New **`fabro-mcp-store`** crate: a concrete, filesystem-backed
`McpServerStore` — one TOML file per definition under
`{active-config-dir}/mcps/`, an in-memory cache, and a SHA-256
content-hash revision for optimistic concurrency. Modeled directly on
`AutomationStore`. Includes an id-only `ids()` accessor for cheap
listing that avoids cloning the (potentially sensitive) env/header maps
a full definition carries.
- New **`McpServerDefinition` / `McpServerDraft` / `McpServerReplace`**
domain model (plus `McpServerId` / `McpServerRevision` and structural
validation) in `fabro-types`, reusing the existing `McpTransport`. These
stay persistence-independent; the on-disk TOML DTO and the filesystem
plumbing live in `fabro-mcp-store`.

Nothing in the workspace depends on the new crate yet. Wiring
`McpServerStore` into the server, the HTTP API, and the UI are follow-up
PRs.

## Testing

- `fabro-mcp-store`: 7/7 (empty/missing dir, non-TOML ignored,
malformed/invalid-filename fail load, CRUD round-trip, stale-revision
and duplicate-create rejected).
- `fabro-types`: `mcp_store` validation and round-trip tests pass.
`cargo build --workspace`, fmt, and clippy all green.

## Notes

- The domain model derives `PartialEq` but not `Eq` because
`McpTransport` carries `HashMap`s (differs from `Automation*`, matches
the transport's capabilities).
- Validation is structural for now (id format, non-empty name,
well-formed transport); credential-literal validation is deliberately
deferred to the API layer (flagged TODO).
- The store is concrete by design (no trait): a future move off per-file
TOML is a one-time migration, not a runtime backend choice. The revision
is currently derived from the canonical TOML bytes — the one
storage-coupled detail to revisit if that move happens.
- Part of a short series adding server-managed MCP servers; independent
of the sibling PRs.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

---------

Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-06-25 15:50:02 -04:00 • committed by GitHub
parent ba56a170d8
commit 4b7c2690c8
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
8 changed files with 1097 additions and 0 deletions

12
Cargo.lock generated
View file

@ -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"

View file

@ -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"] }

View file

@ -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<PathBuf>, source: std::io::Error) -> Self {
Self::Io {
path: path.into(),
source,
}
}
pub(crate) fn parse(path: impl Into<PathBuf>, source: TomlDeError) -> Self {
Self::Parse {
path: path.into(),
source,
}
}
pub(crate) fn invalid_utf8(path: impl Into<PathBuf>, 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",
}
}
}

View file

@ -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;

View file

@ -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<String>,
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<PersistedMcpServer> 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<u8>), 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<PathBuf>,
) -> Result<McpServerDefinition, McpServerStoreError> {
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<Vec<u8>, 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<PersistedMcpServer, McpServerStoreError> {
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))
}

View file

@ -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<HashMap<McpServerId, McpServerDefinition>>,
}
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<PathBuf>) -> Result<Self, McpServerStoreError> {
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<McpServerDefinition> {
let defs = self.defs.read().await;
let mut values = defs.values().cloned().collect::<Vec<_>>();
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<McpServerId> {
let defs = self.defs.read().await;
let mut ids = defs.keys().cloned().collect::<Vec<_>>();
ids.sort();
ids
}
pub async fn get(&self, id: &McpServerId) -> Option<McpServerDefinition> {
self.defs.read().await.get(id).cloned()
}
pub async fn create(
&self,
draft: McpServerDraft,
) -> Result<McpServerDefinition, McpServerStoreError> {
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<McpServerDefinition, McpServerStoreError> {
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<McpServerId, McpServerDefinition>,
id: &McpServerId,
expected: &McpServerRevision,
) -> Result<(), McpServerStoreError> {
let current = defs
.get(id)
.ok_or_else(|| McpServerStoreError::NotFound { id: id.clone() })?;
if &current.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<HashMap<McpServerId, McpServerDefinition>, 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<McpServerDefinition, McpServerStoreError> {
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<McpServerId, McpServerStoreError> {
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<PathBuf, McpServerStoreError> {
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<Item = &str> {
toml.lines().take_while(|line| !line.starts_with('['))
}
}

View file

@ -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,
};

View file

@ -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<String>,
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<String>,
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<String>,
pub transport: McpTransport,
pub startup_timeout_secs: u64,
pub tool_timeout_secs: u64,
}
impl From<McpServerDraft> 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<String>) -> Result<Self, McpServerValidationError> {
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, Self::Err> {
Self::new(value)
}
}
impl Serialize for McpServerId {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
serializer.serialize_str(self.as_str())
}
}
impl<'de> Deserialize<'de> for McpServerId {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
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<Self, Self::Err> {
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<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
serializer.serialize_str(self.as_str())
}
}
impl<'de> Deserialize<'de> for McpServerRevision {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
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::<McpServerId>().is_ok());
assert!("a-1".parse::<McpServerId>().is_ok());
assert!("0".parse::<McpServerId>().is_ok());
assert!("sentry-dev".parse::<McpServerId>().is_ok());
assert!("A".parse::<McpServerId>().is_err());
assert!("a_1".parse::<McpServerId>().is_err());
assert!("-a".parse::<McpServerId>().is_err());
assert!("".parse::<McpServerId>().is_err());
assert!("a".repeat(64).parse::<McpServerId>().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::<McpServerRevision>().is_ok());
assert!("ABC".parse::<McpServerRevision>().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());
}
}