Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions crates/mcp/src/core/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -215,9 +215,10 @@ pub struct McpServerConfig {

/// Whether tools from this server should be hidden from standard client-visible output.
///
/// Limitation: `internal: true` does not currently redact or block live streaming
/// Server-Sent Events (SSE) tool events. Internal tool activity may still appear in
/// real-time streams until streaming redaction is implemented.
/// Behavior:
/// - Internal non-builtin server metadata/tool-trace items are redacted from assembled
/// responses and streaming SSE output.
/// - Builtin-routed tool result items (for example `web_search_call`) remain visible.
#[serde(default)]
pub internal: bool,
}
Expand Down
19 changes: 18 additions & 1 deletion crates/mcp/src/core/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -78,7 +78,7 @@ pub struct McpToolSession<'a> {
tenant_ctx: TenantContext,
/// All MCP servers in this session (including builtin).
all_mcp_servers: Vec<McpServerBinding>,
/// Non-builtin MCP servers only — used for `mcp_list_tools` output.
/// Non-builtin MCP servers only — default source for `mcp_list_tools` output.
mcp_servers: Vec<McpServerBinding>,
mcp_tools: Vec<ToolEntry>,
exposed_name_map: HashMap<String, ExposedToolBinding>,
Expand Down Expand Up @@ -213,6 +213,17 @@ impl<'a> McpToolSession<'a> {
self.exposed_name_map.contains_key(tool_name)
}

/// Returns true when a function call should be intercepted and executed as MCP.
///
/// User-defined function tools take precedence on name collisions.
pub fn should_intercept_function_call(
&self,
tool_name: &str,
user_function_names: &HashSet<String>,
) -> bool {
self.has_exposed_tool(tool_name) && !user_function_names.contains(tool_name)
}

/// Returns the session's qualified-name -> exposed-name mapping.
///
/// Router adapters should use this with response bridge builders.
Expand Down Expand Up @@ -462,6 +473,12 @@ impl<'a> McpToolSession<'a> {
build_mcp_list_tools_item(server_label, &tools)
}

/// Returns true when a per-server streaming `mcp_list_tools` item should be emitted.
pub fn should_emit_streaming_mcp_list_tools(&self, server_label: &str) -> bool {
!self.is_builtin_server_label(server_label)
&& !self.is_internal_non_builtin_server_label(server_label)
}

/// Inject MCP metadata into a response output array.
///
/// Standardized ordering:
Expand Down
13 changes: 6 additions & 7 deletions docs/reference/mcp-internal-servers.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,14 +17,13 @@ servers:
```

In the current implementation, `internal: true` applies only to self-provided
MCP servers declared under `servers:`. It affects final assembled,
non-streaming MCP responses by allowing higher layers to strip internal server
tool lists and tool-call trace items before the response is returned to the
client.
MCP servers declared under `servers:`. It affects both final assembled
responses and streaming output by allowing higher layers to strip
client-visible internal server tool lists and internal non-builtin tool-call
trace items.

This flag does not currently hide streaming output, and it does not apply to
builtin-routed MCP results such as `web_search_call`, `code_interpreter_call`,
or `file_search_call`.
Builtin-routed MCP results such as `web_search_call`, `code_interpreter_call`,
or `file_search_call` remain client-visible.

This flag is generic. It does not imply any vendor-specific behavior and does
not change transport setup or tool execution on its own.
6 changes: 5 additions & 1 deletion model_gateway/src/routers/grpc/common/responses/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,10 @@ pub(crate) mod utils;
// Re-export commonly used items
pub(crate) use context::ResponsesContext;
pub(crate) use streaming::build_sse_response;
pub(crate) use utils::{ensure_mcp_connection, persist_response_if_needed};
pub(crate) use utils::{
emit_visible_mcp_list_tools_sequence, ensure_mcp_connection, persist_response_if_needed,
retain_client_visible_output_items, retain_client_visible_request_tools,
retain_client_visible_response_tools,
};

pub(crate) use crate::routers::common::mcp_utils::collect_user_function_names;
12 changes: 11 additions & 1 deletion model_gateway/src/routers/grpc/common/responses/streaming.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
//! Streaming infrastructure for /v1/responses endpoint

use std::collections::HashSet;

use axum::{body::Body, http::StatusCode, response::Response};
use bytes::Bytes;
use http::header::{HeaderValue, CONTENT_TYPE};
Expand All @@ -15,7 +17,7 @@ use openai_protocol::{
},
};
use serde_json::json;
use smg_mcp::{self as mcp, ResponseFormat};
use smg_mcp::{self as mcp, McpToolSession, ResponseFormat};
use tokio::sync::mpsc;
use tokio_stream::wrappers::UnboundedReceiverStream;
use tracing::warn;
Expand Down Expand Up @@ -1010,3 +1012,11 @@ pub(crate) fn attach_mcp_server_label(
item["server_label"] = json!(label);
}
}

pub(crate) fn is_client_visible_streaming_output_item(
session: Option<&McpToolSession<'_>>,
item: &serde_json::Value,
user_function_names: &HashSet<String>,
) -> bool {
session.is_none_or(|s| !s.should_hide_output_item_json(item, user_function_names))
}
138 changes: 135 additions & 3 deletions model_gateway/src/routers/grpc/common/responses/utils.rs
Original file line number Diff line number Diff line change
@@ -1,20 +1,26 @@
//! Utility functions for /v1/responses endpoint

use std::sync::Arc;
use std::{collections::HashSet, sync::Arc};

use axum::response::Response;
use bytes::Bytes;
use openai_protocol::{
common::Tool,
responses::{ResponseTool, ResponsesRequest, ResponsesResponse},
responses::{
ResponseOutputItem, ResponseTool, ResponsesRequest, ResponsesResponse, ResponsesToolChoice,
ToolChoiceOptions,
},
};
use serde_json::to_value;
use smg_data_connector::{
ConversationItemStorage, ConversationStorage, RequestContext as StorageRequestContext,
ResponseStorage,
};
use smg_mcp::{McpOrchestrator, McpServerBinding};
use smg_mcp::{McpOrchestrator, McpServerBinding, McpToolSession};
use tokio::sync::mpsc;
use tracing::{debug, error, warn};

use super::streaming::ResponseStreamEventEmitter;
use crate::{
routers::{
common::{
Expand Down Expand Up @@ -166,3 +172,129 @@ pub(crate) async fn persist_response_if_needed(
}
}
}

/// Retain only client-visible output items based on session hide policy.
///
/// This keeps non-streaming redaction behavior consistent across regular and
/// harmony response paths.
pub(crate) fn retain_client_visible_output_items(
output: &mut Vec<ResponseOutputItem>,
session: &McpToolSession<'_>,
user_function_names: &HashSet<String>,
) {
output.retain(|item| {
let Ok(json) = to_value(item) else {
return true;
};
!session.should_hide_output_item_json(&json, user_function_names)
});
Comment on lines +185 to +190

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The redaction logic here is 'fail-open'. If to_value(item) fails, the item is retained and potentially leaked to the client. For security and privacy features like redaction, it is safer to 'fail-closed' by returning false on error. Additionally, as per repository guidelines, failures in serialization should be logged as warnings to aid in debugging rather than failing silently.

    output.retain(|item| {
        match to_value(item) {
            Ok(json) => !session.should_hide_output_item_json(&json, user_function_names),
            Err(e) => {
                tracing::warn!(error = ?e, "Failed to serialize item for redaction check; hiding item");
                false
            }
        }
    });
References
  1. Instead of silently ignoring potential failures (e.g., from serialization), log them as warnings to aid in debugging. In Rust, prefer using unwrap_or_else to log an error over unwrap_or_default which would fail silently.

}

fn retain_client_visible_tools_in_place(
tools: &mut Vec<ResponseTool>,
session: &McpToolSession<'_>,
user_function_names: &HashSet<String>,
) {
tools.retain(|tool| {
let Ok(json) = to_value(tool) else {
return true;
};
!session.should_hide_tool_json(&json, user_function_names)
});
}

fn has_function_tool(tools: &[ResponseTool], tool_name: &str) -> bool {
tools.iter().any(|tool| {
matches!(
tool,
ResponseTool::Function(function_tool) if function_tool.function.name == tool_name
)
})
}

fn normalize_request_tool_choice(
tool_choice: &mut Option<ResponsesToolChoice>,
tools: &[ResponseTool],
) {
if tools.is_empty() {
*tool_choice = None;
return;
}

if let Some(selected_name) = tool_choice
.as_ref()
.and_then(ResponsesToolChoice::function_name)
{
if !has_function_tool(tools, selected_name) {
*tool_choice = Some(ResponsesToolChoice::Options(ToolChoiceOptions::Auto));
}
}
}

fn normalize_response_tool_choice(tool_choice: &mut String, tools: &[ResponseTool]) {
if tools.is_empty() {
*tool_choice = "auto".to_string();
return;
}

let selected_name = serde_json::from_str::<ResponsesToolChoice>(tool_choice)
.ok()
.and_then(|choice| choice.function_name().map(str::to_string));
if let Some(selected_name) = selected_name {
if !has_function_tool(tools, &selected_name) {
*tool_choice = "auto".to_string();
}
}
}

/// Retain only client-visible tools and normalize request `tool_choice`.
///
/// Used on streaming paths before copying original request fields into
/// `response.completed` payloads.
pub(crate) fn retain_client_visible_request_tools(
request: &mut ResponsesRequest,
session: &McpToolSession<'_>,
user_function_names: &HashSet<String>,
) {
if let Some(tools) = request.tools.as_mut() {
retain_client_visible_tools_in_place(tools, session, user_function_names);
normalize_request_tool_choice(&mut request.tool_choice, tools);
if tools.is_empty() {
request.tools = None;
}
} else {
request.tool_choice = None;
}
}

/// Retain only client-visible tools and normalize response `tool_choice`.
///
/// Keeps gRPC non-streaming responses aligned with session hide policy.
pub(crate) fn retain_client_visible_response_tools(
response: &mut ResponsesResponse,
session: &McpToolSession<'_>,
user_function_names: &HashSet<String>,
) {
retain_client_visible_tools_in_place(&mut response.tools, session, user_function_names);
normalize_response_tool_choice(&mut response.tool_choice, &response.tools);
}

/// Emit the visible `mcp_list_tools` streaming sequence for all server bindings
/// in the session.
///
/// Visibility is gated by `session.should_emit_streaming_mcp_list_tools` so
/// hidden bindings do not appear in client-facing streams.
pub(crate) fn emit_visible_mcp_list_tools_sequence(
session: &McpToolSession<'_>,
emitter: &mut ResponseStreamEventEmitter,
tx: &mpsc::UnboundedSender<Result<Bytes, std::io::Error>>,
) -> Result<(), String> {
for binding in session.mcp_servers() {
if !session.should_emit_streaming_mcp_list_tools(&binding.label) {
continue;
}
let tools_for_server = session.list_tools_for_server(&binding.server_key);
emitter.emit_mcp_list_tools_sequence(&binding.label, &tools_for_server, tx)?;
}
Ok(())
}
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ use crate::{
grpc::{
common::responses::{
collect_user_function_names, ensure_mcp_connection, persist_response_if_needed,
retain_client_visible_output_items, retain_client_visible_response_tools,
ResponsesContext,
},
harmony::processor::ResponsesIterationResult,
Expand Down Expand Up @@ -180,12 +181,14 @@ async fn execute_with_mcp_loop(
"Tool calls found - separating MCP and function tools"
);

// Separate MCP and function tool calls based on session exposure.
// MCP tools are exposed to the model as function tools, so the only reliable
// discriminator is whether the name belongs to the MCP session.
let (mcp_tool_calls, function_tool_calls): (Vec<_>, Vec<_>) = tool_calls
.into_iter()
.partition(|tc| session.has_exposed_tool(tc.function.name.as_str()));
// Separate MCP and function tool calls based on session policy.
let (mcp_tool_calls, function_tool_calls): (Vec<_>, Vec<_>) =
tool_calls.into_iter().partition(|tc| {
session.should_intercept_function_call(
tc.function.name.as_str(),
&user_function_names,
)
});

debug!(
mcp_calls = mcp_tool_calls.len(),
Expand Down Expand Up @@ -245,6 +248,16 @@ async fn execute_with_mcp_loop(
&user_function_names,
);
}
retain_client_visible_output_items(
&mut response.output,
&session,
&user_function_names,
);
retain_client_visible_response_tools(
&mut response,
&session,
&user_function_names,
);

return Ok(response);
}
Expand Down Expand Up @@ -295,6 +308,11 @@ async fn execute_with_mcp_loop(
&user_function_names,
);
}
retain_client_visible_response_tools(
&mut response,
&session,
&user_function_names,
);

return Ok(response);
}
Expand Down Expand Up @@ -329,6 +347,7 @@ async fn execute_with_mcp_loop(

// Restore original tools (hide internal MCP tools from response)
response.tools = original_tools.take().unwrap_or_default();
retain_client_visible_response_tools(&mut response, &session, &user_function_names);

debug!(
mcp_calls = mcp_tracking.total_calls(),
Expand Down
Loading
Loading