feat: add bounded MCP server management endpoints to runtime API
- Add `mcp_server_management: true` capability to `RuntimeCapabilities`
in the protocol crate so clients can discover support via
`GET /v1/runtime/info`.
- Expose 7 new routes on the runtime API:
- `POST /v1/apps/mcp/servers` create
- `GET /v1/apps/mcp/servers/{name}` read (redacted)
- `PATCH /v1/apps/mcp/servers/{name}` update (partial)
- `DELETE /v1/apps/mcp/servers/{name}` delete
- `POST /v1/apps/mcp/servers/{name}/enable` enable
- `POST /v1/apps/mcp/servers/{name}/disable` disable
- `POST /v1/apps/mcp/servers/{name}/reconnect` drop pool / re-init
- Credential redaction: `McpServerDetail` response type never returns
header values, env variable values, bearer-token env var names, or
OAuth client secrets; callers see only key names and boolean flags.
- Make `validate_mcp_transport` public so it can be called from the
new route handlers.
- Add 4 tests covering: full CRUD lifecycle, input validation (400 on
missing command/url, 400 on missing name, 409 on duplicate), credential
redaction, and capability advertisement.
Closes #5071
This commit is contained in:
committed by
Hmbown
parent
34f28e9785
commit
d16b83284e
@@ -59,6 +59,11 @@ pub struct RuntimeCapabilities {
|
||||
pub fleet_event_stream: bool,
|
||||
#[serde(default)]
|
||||
pub fleet_local_target: bool,
|
||||
/// Whether the runtime supports create/update/enable/disable/reconnect/delete
|
||||
/// operations on MCP server configuration via the `POST|GET|PATCH|DELETE
|
||||
/// /v1/apps/mcp/servers` family of endpoints.
|
||||
#[serde(default)]
|
||||
pub mcp_server_management: bool,
|
||||
}
|
||||
|
||||
/// Experimental opt-in flags advertised by `GET /v1/runtime/info`.
|
||||
|
||||
@@ -1632,7 +1632,7 @@ fn is_legacy_sse_transport(config: &McpServerConfig) -> bool {
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
fn validate_mcp_transport(transport: Option<&str>) -> Result<()> {
|
||||
pub fn validate_mcp_transport(transport: Option<&str>) -> Result<()> {
|
||||
let Some(transport) = transport else {
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
@@ -401,6 +401,7 @@ fn default_runtime_capabilities() -> RuntimeCapabilities {
|
||||
fleet_event_replay: true,
|
||||
fleet_event_stream: true,
|
||||
fleet_local_target: true,
|
||||
mcp_server_management: true,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -454,6 +455,122 @@ struct McpToolsResponse {
|
||||
tools: Vec<McpToolEntry>,
|
||||
}
|
||||
|
||||
/// Request body for `POST /v1/apps/mcp/servers` (create) and
|
||||
/// `PATCH /v1/apps/mcp/servers/{name}` (update).
|
||||
///
|
||||
/// Either `command` **or** `url` must be set on create. On update, only
|
||||
/// supplied fields are applied; absent fields leave the existing value in
|
||||
/// place.
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct McpServerWriteRequest {
|
||||
/// stdio command binary (e.g. `"npx"`).
|
||||
command: Option<String>,
|
||||
/// Arguments for the stdio command.
|
||||
args: Option<Vec<String>>,
|
||||
/// Environment variables injected into the stdio child process.
|
||||
/// Values are stored as-is; use `${VAR}` syntax to reference environment
|
||||
/// variables at runtime instead of embedding secrets here.
|
||||
env: Option<std::collections::HashMap<String, String>>,
|
||||
/// HTTP(S) endpoint for streamable-HTTP or SSE MCP servers.
|
||||
url: Option<String>,
|
||||
/// Explicit transport override (`"sse"` or `"streamable_http"`).
|
||||
transport: Option<String>,
|
||||
/// Override the server-level connect timeout in seconds.
|
||||
connect_timeout: Option<u64>,
|
||||
/// Override the server-level execute timeout in seconds.
|
||||
execute_timeout: Option<u64>,
|
||||
/// Override the server-level read timeout in seconds.
|
||||
read_timeout: Option<u64>,
|
||||
/// Whether the server is enabled. Defaults to `true` on create.
|
||||
enabled: Option<bool>,
|
||||
/// Whether a connection failure for this server is fatal.
|
||||
required: Option<bool>,
|
||||
/// Allowlist of tool names to expose (empty = expose all).
|
||||
enabled_tools: Option<Vec<String>>,
|
||||
/// Denylist of tool names to hide.
|
||||
disabled_tools: Option<Vec<String>>,
|
||||
/// Variable names whose runtime values are injected as HTTP headers.
|
||||
/// The key in this map is the HTTP header name; the value is the
|
||||
/// environment variable whose value supplies the header value at
|
||||
/// request time. Credentials remain in the environment, not on disk.
|
||||
env_headers: Option<std::collections::HashMap<String, String>>,
|
||||
/// Environment variable that contains a bearer token for URL-based servers.
|
||||
bearer_token_env_var: Option<String>,
|
||||
/// OAuth scopes requested during `codewhale mcp login`.
|
||||
scopes: Option<Vec<String>>,
|
||||
/// RFC 8707 resource parameter for the OAuth authorization URL.
|
||||
oauth_resource: Option<String>,
|
||||
}
|
||||
|
||||
/// Response returned by MCP server management endpoints.
|
||||
///
|
||||
/// Sensitive fields (`headers`, `env_headers`, `bearer_token_env_var`,
|
||||
/// `env`, OAuth client secrets) are intentionally omitted or redacted so
|
||||
/// the API never echoes credentials back to callers.
|
||||
#[derive(Debug, Serialize)]
|
||||
struct McpServerDetail {
|
||||
name: String,
|
||||
enabled: bool,
|
||||
required: bool,
|
||||
command: Option<String>,
|
||||
args: Vec<String>,
|
||||
/// Environment variable names injected into the process.
|
||||
/// Values are **not** returned — callers see only the keys.
|
||||
env_keys: Vec<String>,
|
||||
url: Option<String>,
|
||||
transport: Option<String>,
|
||||
connect_timeout: Option<u64>,
|
||||
execute_timeout: Option<u64>,
|
||||
read_timeout: Option<u64>,
|
||||
enabled_tools: Vec<String>,
|
||||
disabled_tools: Vec<String>,
|
||||
/// HTTP header names that are read from environment variables.
|
||||
/// The corresponding environment variable values are **not** returned.
|
||||
env_header_keys: Vec<String>,
|
||||
/// Whether a `bearer_token_env_var` is configured (value not returned).
|
||||
has_bearer_token_env_var: bool,
|
||||
scopes: Vec<String>,
|
||||
oauth_resource: Option<String>,
|
||||
/// Live connection state from the in-memory pool (if the pool is active).
|
||||
connected: bool,
|
||||
}
|
||||
|
||||
impl McpServerDetail {
|
||||
fn from_config(name: &str, cfg: &crate::mcp::McpServerConfig, connected: bool) -> Self {
|
||||
let mut env_keys: Vec<String> = cfg.env.keys().cloned().collect();
|
||||
env_keys.sort();
|
||||
let mut env_header_keys: Vec<String> = cfg.env_headers.keys().cloned().collect();
|
||||
env_header_keys.sort();
|
||||
Self {
|
||||
name: name.to_string(),
|
||||
enabled: cfg.is_enabled(),
|
||||
required: cfg.required,
|
||||
command: cfg.command.clone(),
|
||||
args: cfg.args.clone(),
|
||||
env_keys,
|
||||
url: cfg.url.clone(),
|
||||
transport: cfg.transport.clone(),
|
||||
connect_timeout: cfg.connect_timeout,
|
||||
execute_timeout: cfg.execute_timeout,
|
||||
read_timeout: cfg.read_timeout,
|
||||
enabled_tools: cfg.enabled_tools.clone(),
|
||||
disabled_tools: cfg.disabled_tools.clone(),
|
||||
env_header_keys,
|
||||
has_bearer_token_env_var: cfg.bearer_token_env_var.is_some(),
|
||||
scopes: cfg.scopes.clone(),
|
||||
oauth_resource: cfg.oauth_resource.clone(),
|
||||
connected,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
struct McpServerActionReceipt {
|
||||
name: String,
|
||||
action: &'static str,
|
||||
ok: bool,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct AutomationRunsQuery {
|
||||
limit: Option<usize>,
|
||||
@@ -761,7 +878,22 @@ pub fn build_router(state: RuntimeApiState) -> Router {
|
||||
.route("/v1/tasks/{id}/cancel", post(cancel_task))
|
||||
.route("/v1/skills", get(list_skills))
|
||||
.route("/v1/skills/{name}", post(set_skill_enabled))
|
||||
.route("/v1/apps/mcp/servers", get(list_mcp_servers))
|
||||
.route("/v1/apps/mcp/servers", get(list_mcp_servers).post(create_mcp_server))
|
||||
.route(
|
||||
"/v1/apps/mcp/servers/{name}",
|
||||
get(get_mcp_server)
|
||||
.patch(update_mcp_server)
|
||||
.delete(delete_mcp_server),
|
||||
)
|
||||
.route("/v1/apps/mcp/servers/{name}/enable", post(enable_mcp_server))
|
||||
.route(
|
||||
"/v1/apps/mcp/servers/{name}/disable",
|
||||
post(disable_mcp_server),
|
||||
)
|
||||
.route(
|
||||
"/v1/apps/mcp/servers/{name}/reconnect",
|
||||
post(reconnect_mcp_server),
|
||||
)
|
||||
.route("/v1/apps/mcp/tools", get(list_mcp_tools))
|
||||
.route(
|
||||
"/v1/automations",
|
||||
@@ -2320,6 +2452,351 @@ async fn list_mcp_tools(
|
||||
Ok(Json(McpToolsResponse { tools }))
|
||||
}
|
||||
|
||||
/// `GET /v1/apps/mcp/servers/{name}` — fetch a single server's redacted config.
|
||||
async fn get_mcp_server(
|
||||
State(state): State<RuntimeApiState>,
|
||||
Path(name): Path<String>,
|
||||
) -> Result<Json<McpServerDetail>, ApiError> {
|
||||
let mcp_config_path = state.config.read().mcp_config_path();
|
||||
let plugin_registry = state
|
||||
.plugin_discovery
|
||||
.registry_for_workspace(&state.workspace);
|
||||
let config = crate::mcp::load_config_with_workspace_and_plugins(
|
||||
&mcp_config_path,
|
||||
&state.workspace,
|
||||
plugin_registry.as_ref(),
|
||||
)
|
||||
.map_err(|e| ApiError::internal(format!("Failed to load MCP config: {e}")))?;
|
||||
|
||||
let server_cfg = config
|
||||
.servers
|
||||
.get(&name)
|
||||
.ok_or_else(|| ApiError::not_found(format!("MCP server '{name}' not found")))?;
|
||||
|
||||
let connected = {
|
||||
let pool_slot = state.mcp_pool.lock().await;
|
||||
pool_slot.as_ref().map_or(false, |pool_handle| {
|
||||
let pool = pool_handle.try_lock();
|
||||
pool.map_or(false, |p| p.connected_servers().contains(&name.as_str()))
|
||||
})
|
||||
};
|
||||
|
||||
Ok(Json(McpServerDetail::from_config(&name, server_cfg, connected)))
|
||||
}
|
||||
|
||||
/// `POST /v1/apps/mcp/servers` — add a new server to the persistent config.
|
||||
///
|
||||
/// Body: JSON object with all `McpServerWriteRequest` fields **plus** a
|
||||
/// required top-level `"name"` string that will be the server key.
|
||||
async fn create_mcp_server(
|
||||
State(state): State<RuntimeApiState>,
|
||||
Json(body): Json<serde_json::Value>,
|
||||
) -> Result<(StatusCode, Json<McpServerDetail>), ApiError> {
|
||||
let name = body
|
||||
.get("name")
|
||||
.and_then(|v| v.as_str())
|
||||
.ok_or_else(|| ApiError::bad_request("'name' is required"))?
|
||||
.to_string();
|
||||
|
||||
if name.trim().is_empty() {
|
||||
return Err(ApiError::bad_request("'name' must not be empty"));
|
||||
}
|
||||
|
||||
let req: McpServerWriteRequest = serde_json::from_value(body)
|
||||
.map_err(|e| ApiError::bad_request(format!("Invalid request body: {e}")))?;
|
||||
|
||||
if req.command.is_none() && req.url.is_none() {
|
||||
return Err(ApiError::bad_request(
|
||||
"Either 'command' or 'url' is required to create an MCP server",
|
||||
));
|
||||
}
|
||||
|
||||
if let Some(ref transport) = req.transport {
|
||||
crate::mcp::validate_mcp_transport(Some(transport.as_str()))
|
||||
.map_err(|e| ApiError::bad_request(e.to_string()))?;
|
||||
}
|
||||
|
||||
let mcp_config_path = state.config.read().mcp_config_path();
|
||||
|
||||
// Build the config entry from the request.
|
||||
let new_cfg = mcp_server_config_from_write_request(req, None);
|
||||
|
||||
// Persist to the global MCP config.
|
||||
{
|
||||
let mut cfg = crate::mcp::load_config(&mcp_config_path)
|
||||
.map_err(|e| ApiError::internal(format!("Failed to load MCP config: {e}")))?;
|
||||
if cfg.servers.contains_key(&name) {
|
||||
return Err(ApiError {
|
||||
status: StatusCode::CONFLICT,
|
||||
message: format!("MCP server '{name}' already exists"),
|
||||
});
|
||||
}
|
||||
cfg.servers.insert(name.clone(), new_cfg.clone());
|
||||
crate::mcp::save_config(&mcp_config_path, &cfg)
|
||||
.map_err(|e| ApiError::internal(format!("Failed to save MCP config: {e}")))?;
|
||||
}
|
||||
|
||||
// Invalidate the in-memory pool so the next tool call reloads from disk.
|
||||
{
|
||||
let mut pool_slot = state.mcp_pool.lock().await;
|
||||
*pool_slot = None;
|
||||
}
|
||||
|
||||
Ok((
|
||||
StatusCode::CREATED,
|
||||
Json(McpServerDetail::from_config(&name, &new_cfg, false)),
|
||||
))
|
||||
}
|
||||
|
||||
/// `PATCH /v1/apps/mcp/servers/{name}` — update an existing server's config.
|
||||
async fn update_mcp_server(
|
||||
State(state): State<RuntimeApiState>,
|
||||
Path(name): Path<String>,
|
||||
Json(req): Json<McpServerWriteRequest>,
|
||||
) -> Result<Json<McpServerDetail>, ApiError> {
|
||||
if let Some(ref transport) = req.transport {
|
||||
crate::mcp::validate_mcp_transport(Some(transport.as_str()))
|
||||
.map_err(|e| ApiError::bad_request(e.to_string()))?;
|
||||
}
|
||||
|
||||
let mcp_config_path = state.config.read().mcp_config_path();
|
||||
|
||||
let updated_cfg = {
|
||||
let mut cfg = crate::mcp::load_config(&mcp_config_path)
|
||||
.map_err(|e| ApiError::internal(format!("Failed to load MCP config: {e}")))?;
|
||||
let existing = cfg
|
||||
.servers
|
||||
.get_mut(&name)
|
||||
.ok_or_else(|| ApiError::not_found(format!("MCP server '{name}' not found")))?;
|
||||
apply_write_request_to_config(req, existing);
|
||||
let updated = existing.clone();
|
||||
crate::mcp::save_config(&mcp_config_path, &cfg)
|
||||
.map_err(|e| ApiError::internal(format!("Failed to save MCP config: {e}")))?;
|
||||
updated
|
||||
};
|
||||
|
||||
// Invalidate the in-memory pool.
|
||||
{
|
||||
let mut pool_slot = state.mcp_pool.lock().await;
|
||||
*pool_slot = None;
|
||||
}
|
||||
|
||||
Ok(Json(McpServerDetail::from_config(&name, &updated_cfg, false)))
|
||||
}
|
||||
|
||||
/// `DELETE /v1/apps/mcp/servers/{name}` — remove a server from the persistent config.
|
||||
async fn delete_mcp_server(
|
||||
State(state): State<RuntimeApiState>,
|
||||
Path(name): Path<String>,
|
||||
) -> Result<Json<McpServerActionReceipt>, ApiError> {
|
||||
let mcp_config_path = state.config.read().mcp_config_path();
|
||||
|
||||
crate::mcp::remove_server_config(&mcp_config_path, &name)
|
||||
.map_err(|e| {
|
||||
let msg = e.to_string();
|
||||
if msg.contains("not found") {
|
||||
ApiError::not_found(msg)
|
||||
} else {
|
||||
ApiError::internal(msg)
|
||||
}
|
||||
})?;
|
||||
|
||||
// Invalidate the in-memory pool.
|
||||
{
|
||||
let mut pool_slot = state.mcp_pool.lock().await;
|
||||
*pool_slot = None;
|
||||
}
|
||||
|
||||
Ok(Json(McpServerActionReceipt {
|
||||
name,
|
||||
action: "deleted",
|
||||
ok: true,
|
||||
}))
|
||||
}
|
||||
|
||||
/// `POST /v1/apps/mcp/servers/{name}/enable` — enable a configured server.
|
||||
async fn enable_mcp_server(
|
||||
State(state): State<RuntimeApiState>,
|
||||
Path(name): Path<String>,
|
||||
) -> Result<Json<McpServerActionReceipt>, ApiError> {
|
||||
let mcp_config_path = state.config.read().mcp_config_path();
|
||||
|
||||
crate::mcp::set_server_enabled(&mcp_config_path, &name, true).map_err(|e| {
|
||||
let msg = e.to_string();
|
||||
if msg.contains("not found") {
|
||||
ApiError::not_found(msg)
|
||||
} else {
|
||||
ApiError::internal(msg)
|
||||
}
|
||||
})?;
|
||||
|
||||
// Invalidate the in-memory pool so the enabled server participates next time.
|
||||
{
|
||||
let mut pool_slot = state.mcp_pool.lock().await;
|
||||
*pool_slot = None;
|
||||
}
|
||||
|
||||
Ok(Json(McpServerActionReceipt {
|
||||
name,
|
||||
action: "enabled",
|
||||
ok: true,
|
||||
}))
|
||||
}
|
||||
|
||||
/// `POST /v1/apps/mcp/servers/{name}/disable` — disable a configured server.
|
||||
async fn disable_mcp_server(
|
||||
State(state): State<RuntimeApiState>,
|
||||
Path(name): Path<String>,
|
||||
) -> Result<Json<McpServerActionReceipt>, ApiError> {
|
||||
let mcp_config_path = state.config.read().mcp_config_path();
|
||||
|
||||
crate::mcp::set_server_enabled(&mcp_config_path, &name, false).map_err(|e| {
|
||||
let msg = e.to_string();
|
||||
if msg.contains("not found") {
|
||||
ApiError::not_found(msg)
|
||||
} else {
|
||||
ApiError::internal(msg)
|
||||
}
|
||||
})?;
|
||||
|
||||
// Invalidate the in-memory pool so the disabled server is excluded next time.
|
||||
{
|
||||
let mut pool_slot = state.mcp_pool.lock().await;
|
||||
*pool_slot = None;
|
||||
}
|
||||
|
||||
Ok(Json(McpServerActionReceipt {
|
||||
name,
|
||||
action: "disabled",
|
||||
ok: true,
|
||||
}))
|
||||
}
|
||||
|
||||
/// `POST /v1/apps/mcp/servers/{name}/reconnect` — drop the cached pool entry
|
||||
/// for this server so it re-initializes on the next call that needs tools.
|
||||
async fn reconnect_mcp_server(
|
||||
State(state): State<RuntimeApiState>,
|
||||
Path(name): Path<String>,
|
||||
) -> Result<Json<McpServerActionReceipt>, ApiError> {
|
||||
// Verify the server exists in the config.
|
||||
let mcp_config_path = state.config.read().mcp_config_path();
|
||||
let plugin_registry = state
|
||||
.plugin_discovery
|
||||
.registry_for_workspace(&state.workspace);
|
||||
let config = crate::mcp::load_config_with_workspace_and_plugins(
|
||||
&mcp_config_path,
|
||||
&state.workspace,
|
||||
plugin_registry.as_ref(),
|
||||
)
|
||||
.map_err(|e| ApiError::internal(format!("Failed to load MCP config: {e}")))?;
|
||||
|
||||
if !config.servers.contains_key(&name) {
|
||||
return Err(ApiError::not_found(format!(
|
||||
"MCP server '{name}' not found"
|
||||
)));
|
||||
}
|
||||
|
||||
// Drop the whole pool so the next connect_all call recreates all
|
||||
// connections from the current on-disk config.
|
||||
{
|
||||
let mut pool_slot = state.mcp_pool.lock().await;
|
||||
*pool_slot = None;
|
||||
}
|
||||
|
||||
Ok(Json(McpServerActionReceipt {
|
||||
name,
|
||||
action: "reconnect_scheduled",
|
||||
ok: true,
|
||||
}))
|
||||
}
|
||||
|
||||
/// Build a fresh [`McpServerConfig`] from a create request.
|
||||
fn mcp_server_config_from_write_request(
|
||||
req: McpServerWriteRequest,
|
||||
_existing: Option<&crate::mcp::McpServerConfig>,
|
||||
) -> crate::mcp::McpServerConfig {
|
||||
let enabled = req.enabled.unwrap_or(true);
|
||||
crate::mcp::McpServerConfig {
|
||||
command: req.command,
|
||||
args: req.args.unwrap_or_default(),
|
||||
env: req.env.unwrap_or_default(),
|
||||
cwd: None,
|
||||
url: req.url,
|
||||
transport: req.transport,
|
||||
connect_timeout: req.connect_timeout,
|
||||
execute_timeout: req.execute_timeout,
|
||||
read_timeout: req.read_timeout,
|
||||
disabled: !enabled,
|
||||
enabled,
|
||||
required: req.required.unwrap_or(false),
|
||||
enabled_tools: req.enabled_tools.unwrap_or_default(),
|
||||
disabled_tools: req.disabled_tools.unwrap_or_default(),
|
||||
headers: std::collections::HashMap::new(),
|
||||
env_headers: req.env_headers.unwrap_or_default(),
|
||||
bearer_token_env_var: req.bearer_token_env_var,
|
||||
scopes: req.scopes.unwrap_or_default(),
|
||||
oauth: None,
|
||||
oauth_resource: req.oauth_resource,
|
||||
reviewed_plugin: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Apply a partial update from a PATCH request onto an existing config entry.
|
||||
fn apply_write_request_to_config(
|
||||
req: McpServerWriteRequest,
|
||||
cfg: &mut crate::mcp::McpServerConfig,
|
||||
) {
|
||||
if let Some(v) = req.command {
|
||||
cfg.command = Some(v);
|
||||
}
|
||||
if let Some(v) = req.args {
|
||||
cfg.args = v;
|
||||
}
|
||||
if let Some(v) = req.env {
|
||||
cfg.env = v;
|
||||
}
|
||||
if let Some(v) = req.url {
|
||||
cfg.url = Some(v);
|
||||
}
|
||||
if req.transport.is_some() {
|
||||
cfg.transport = req.transport;
|
||||
}
|
||||
if req.connect_timeout.is_some() {
|
||||
cfg.connect_timeout = req.connect_timeout;
|
||||
}
|
||||
if req.execute_timeout.is_some() {
|
||||
cfg.execute_timeout = req.execute_timeout;
|
||||
}
|
||||
if req.read_timeout.is_some() {
|
||||
cfg.read_timeout = req.read_timeout;
|
||||
}
|
||||
if let Some(v) = req.enabled {
|
||||
cfg.enabled = v;
|
||||
cfg.disabled = !v;
|
||||
}
|
||||
if let Some(v) = req.required {
|
||||
cfg.required = v;
|
||||
}
|
||||
if let Some(v) = req.enabled_tools {
|
||||
cfg.enabled_tools = v;
|
||||
}
|
||||
if let Some(v) = req.disabled_tools {
|
||||
cfg.disabled_tools = v;
|
||||
}
|
||||
if let Some(v) = req.env_headers {
|
||||
cfg.env_headers = v;
|
||||
}
|
||||
if req.bearer_token_env_var.is_some() {
|
||||
cfg.bearer_token_env_var = req.bearer_token_env_var;
|
||||
}
|
||||
if let Some(v) = req.scopes {
|
||||
cfg.scopes = v;
|
||||
}
|
||||
if req.oauth_resource.is_some() {
|
||||
cfg.oauth_resource = req.oauth_resource;
|
||||
}
|
||||
}
|
||||
|
||||
async fn list_automations(
|
||||
State(state): State<RuntimeApiState>,
|
||||
) -> Result<Json<Vec<AutomationRecord>>, ApiError> {
|
||||
|
||||
@@ -1206,6 +1206,260 @@ async fn mcp_tools_endpoint_is_passive_until_connect_requested() -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn mcp_server_management_crud() -> Result<()> {
|
||||
let root =
|
||||
std::env::temp_dir().join(format!("codewhale-mcp-mgmt-{}", Uuid::new_v4()));
|
||||
let sessions_dir = root.join("sessions");
|
||||
fs::create_dir_all(&root)?;
|
||||
|
||||
let Some((addr, _runtime_threads, handle)) =
|
||||
spawn_test_server_with_root(root.clone(), sessions_dir).await?
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
let client = crate::tls::reqwest_client();
|
||||
let base = format!("http://{addr}/v1/apps/mcp/servers");
|
||||
|
||||
// 1. Create a new server.
|
||||
let created: serde_json::Value = client
|
||||
.post(&base)
|
||||
.json(&serde_json::json!({
|
||||
"name": "test-stdio",
|
||||
"command": "echo",
|
||||
"args": ["hello"],
|
||||
}))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?
|
||||
.json()
|
||||
.await?;
|
||||
assert_eq!(created["name"], "test-stdio");
|
||||
assert_eq!(created["enabled"], true);
|
||||
assert_eq!(created["command"], "echo");
|
||||
|
||||
// 2. GET the server back — config should be on disk.
|
||||
let fetched: serde_json::Value = client
|
||||
.get(format!("{base}/test-stdio"))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?
|
||||
.json()
|
||||
.await?;
|
||||
assert_eq!(fetched["name"], "test-stdio");
|
||||
assert_eq!(fetched["command"], "echo");
|
||||
|
||||
// 3. Duplicate create returns 409 Conflict.
|
||||
let conflict = client
|
||||
.post(&base)
|
||||
.json(&serde_json::json!({
|
||||
"name": "test-stdio",
|
||||
"command": "echo",
|
||||
}))
|
||||
.send()
|
||||
.await?;
|
||||
assert_eq!(conflict.status(), 409);
|
||||
|
||||
// 4. PATCH (update) the server.
|
||||
let updated: serde_json::Value = client
|
||||
.patch(format!("{base}/test-stdio"))
|
||||
.json(&serde_json::json!({
|
||||
"args": ["updated"],
|
||||
"required": true,
|
||||
}))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?
|
||||
.json()
|
||||
.await?;
|
||||
assert_eq!(updated["required"], true);
|
||||
assert_eq!(updated["args"][0], "updated");
|
||||
|
||||
// 5. Disable the server.
|
||||
let disabled: serde_json::Value = client
|
||||
.post(format!("{base}/test-stdio/disable"))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?
|
||||
.json()
|
||||
.await?;
|
||||
assert_eq!(disabled["action"], "disabled");
|
||||
assert_eq!(disabled["ok"], true);
|
||||
|
||||
// Verify disabled state via GET.
|
||||
let after_disable: serde_json::Value = client
|
||||
.get(format!("{base}/test-stdio"))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?
|
||||
.json()
|
||||
.await?;
|
||||
assert_eq!(after_disable["enabled"], false);
|
||||
|
||||
// 6. Re-enable the server.
|
||||
let enabled: serde_json::Value = client
|
||||
.post(format!("{base}/test-stdio/enable"))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?
|
||||
.json()
|
||||
.await?;
|
||||
assert_eq!(enabled["action"], "enabled");
|
||||
assert_eq!(enabled["ok"], true);
|
||||
|
||||
// 7. Reconnect (schedules a reconnect — no live pool present so always succeeds).
|
||||
let reconnected: serde_json::Value = client
|
||||
.post(format!("{base}/test-stdio/reconnect"))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?
|
||||
.json()
|
||||
.await?;
|
||||
assert_eq!(reconnected["action"], "reconnect_scheduled");
|
||||
assert_eq!(reconnected["ok"], true);
|
||||
|
||||
// 8. Delete the server.
|
||||
let deleted: serde_json::Value = client
|
||||
.delete(format!("{base}/test-stdio"))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?
|
||||
.json()
|
||||
.await?;
|
||||
assert_eq!(deleted["action"], "deleted");
|
||||
assert_eq!(deleted["ok"], true);
|
||||
|
||||
// 9. GET after delete returns 404.
|
||||
let not_found = client
|
||||
.get(format!("{base}/test-stdio"))
|
||||
.send()
|
||||
.await?;
|
||||
assert_eq!(not_found.status(), 404);
|
||||
|
||||
handle.abort();
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn mcp_server_management_create_requires_command_or_url() -> Result<()> {
|
||||
let root = std::env::temp_dir().join(format!(
|
||||
"codewhale-mcp-mgmt-validation-{}",
|
||||
Uuid::new_v4()
|
||||
));
|
||||
let sessions_dir = root.join("sessions");
|
||||
fs::create_dir_all(&root)?;
|
||||
|
||||
let Some((addr, _runtime_threads, handle)) =
|
||||
spawn_test_server_with_root(root.clone(), sessions_dir).await?
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
let client = crate::tls::reqwest_client();
|
||||
|
||||
// Missing both command and url → 400.
|
||||
let resp = client
|
||||
.post(format!("http://{addr}/v1/apps/mcp/servers"))
|
||||
.json(&serde_json::json!({ "name": "bad-server" }))
|
||||
.send()
|
||||
.await?;
|
||||
assert_eq!(resp.status(), 400);
|
||||
|
||||
// Missing name → 400.
|
||||
let resp = client
|
||||
.post(format!("http://{addr}/v1/apps/mcp/servers"))
|
||||
.json(&serde_json::json!({ "command": "echo" }))
|
||||
.send()
|
||||
.await?;
|
||||
assert_eq!(resp.status(), 400);
|
||||
|
||||
handle.abort();
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn mcp_server_management_redacts_credentials() -> Result<()> {
|
||||
let root = std::env::temp_dir().join(format!(
|
||||
"codewhale-mcp-mgmt-redact-{}",
|
||||
Uuid::new_v4()
|
||||
));
|
||||
let sessions_dir = root.join("sessions");
|
||||
fs::create_dir_all(&root)?;
|
||||
|
||||
let Some((addr, _runtime_threads, handle)) =
|
||||
spawn_test_server_with_root(root.clone(), sessions_dir).await?
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
let client = crate::tls::reqwest_client();
|
||||
let base = format!("http://{addr}/v1/apps/mcp/servers");
|
||||
|
||||
// Create a server with sensitive fields.
|
||||
let created: serde_json::Value = client
|
||||
.post(&base)
|
||||
.json(&serde_json::json!({
|
||||
"name": "secret-server",
|
||||
"url": "https://example.com/mcp",
|
||||
"env_headers": { "Authorization": "MY_SECRET_ENV_VAR" },
|
||||
"bearer_token_env_var": "MY_BEARER_ENV",
|
||||
}))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?
|
||||
.json()
|
||||
.await?;
|
||||
|
||||
// The response must NOT contain the actual header value or env variable value.
|
||||
assert!(
|
||||
!created.to_string().contains("MY_SECRET_ENV_VAR_VALUE"),
|
||||
"credential values must be redacted from API responses"
|
||||
);
|
||||
// env_header_keys should list the header name (not the env-var value).
|
||||
assert_eq!(
|
||||
created["env_header_keys"]
|
||||
.as_array()
|
||||
.map(|a| a.iter().any(|v| v.as_str() == Some("Authorization"))),
|
||||
Some(true),
|
||||
"env_header_keys should list header names"
|
||||
);
|
||||
// has_bearer_token_env_var should be true, not the env var name.
|
||||
assert_eq!(created["has_bearer_token_env_var"], true);
|
||||
|
||||
handle.abort();
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn runtime_info_advertises_mcp_server_management() -> Result<()> {
|
||||
let root = std::env::temp_dir().join(format!(
|
||||
"codewhale-mcp-capability-{}",
|
||||
Uuid::new_v4()
|
||||
));
|
||||
let sessions_dir = root.join("sessions");
|
||||
let Some((addr, _runtime_threads, handle)) =
|
||||
spawn_test_server_with_root(root, sessions_dir).await?
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
let client = crate::tls::reqwest_client();
|
||||
|
||||
let info: serde_json::Value = client
|
||||
.get(format!("http://{addr}/v1/runtime/info"))
|
||||
.send()
|
||||
.await?
|
||||
.error_for_status()?
|
||||
.json()
|
||||
.await?;
|
||||
|
||||
assert_eq!(
|
||||
info["capabilities"]["mcp_server_management"],
|
||||
true,
|
||||
"runtime/info must advertise mcp_server_management capability"
|
||||
);
|
||||
|
||||
handle.abort();
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn runtime_token_guard_protects_v1_routes() -> Result<()> {
|
||||
let root = std::env::temp_dir().join(format!("deepseek-runtime-api-{}", Uuid::new_v4()));
|
||||
|
||||
Reference in New Issue
Block a user