Skip to content
Open
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
16 changes: 0 additions & 16 deletions crates/agentic-server-core/src/executor/accumulator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -605,22 +605,6 @@ mod tests {
}
}

#[test]
fn test_process_event_mcp_tool_call_done_accumulates_output() {
let lines = vec![
r#"data: {"type":"response.output_item.added","output_index":0,"item":{"type":"mcp_tool_call","id":"mcp_1","server":"repo","tool":"read_mcp_resource","arguments":{"server":"repo"},"status":"in_progress"}}"#.to_string(),
r#"data: {"type":"response.mcp_tool_call.in_progress","item_id":"mcp_1","output_index":0}"#.to_string(),
r#"data: {"type":"response.mcp_tool_call.completed","item_id":"mcp_1","output_index":0,"item":{"type":"mcp_tool_call","id":"mcp_1","server":"repo","tool":"read_mcp_resource","arguments":{"server":"repo"},"status":"completed","result":{"contents":[]}}}"#.to_string(),
r#"data: {"type":"response.output_item.done","output_index":0,"item":{"type":"mcp_tool_call","id":"mcp_1","server":"repo","tool":"read_mcp_resource","arguments":{"server":"repo"},"status":"completed","result":{"contents":[]}}}"#.to_string(),
r#"data: {"type":"response.done","response":{"id":"resp_1","status":"completed","usage":{"input_tokens":5,"output_tokens":2,"total_tokens":7}}}"#.to_string(),
];

let acc = ResponseAccumulator::from_sse_lines(lines, None);
assert_eq!(acc.status, ResponseStatus::Completed);
assert_eq!(acc.output.len(), 1);
assert!(matches!(acc.output[0], OutputItem::McpToolCall(_)));
}

#[test]
fn test_process_event_web_search_done_accumulates_output() {
let lines = vec![
Expand Down
5 changes: 3 additions & 2 deletions crates/agentic-server-core/src/executor/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -99,8 +99,9 @@ async fn run_until_gateway_tools_complete(
stream_upstream: bool,
mut stream: Option<(&mut GatewayStreamAccumulator, &mpsc::UnboundedSender<StreamEvent>)>,
) -> ExecutorResult<(ResponsePayload, RequestContext)> {
let registry: ToolRegistry = match ctx.enriched_request.tools.as_ref() {
Some(tools) => ToolRegistry::build_with_handlers(tools, &exec_ctx.gateway_executors).await?,
let mut executors = exec_ctx.gateway_executors.request_scoped();
let registry: ToolRegistry = match ctx.enriched_request.tools.as_mut() {
Some(tools) => ToolRegistry::build_with_handlers(tools, &mut executors).await?,
None => ToolRegistry::default(),
};
let mut combined_output: Vec<crate::OutputItem> = Vec::new();
Expand Down
18 changes: 13 additions & 5 deletions crates/agentic-server-core/src/executor/gateway.rs
Original file line number Diff line number Diff line change
Expand Up @@ -172,7 +172,7 @@ async fn execute_gateway_call_with_timeout(
(execution_error_output(&call, &message)?, GatewayCallStatus::Failed)
}
};
let public_output = gateway_public_output(dispatch.tool_type, &call, &output, status);
let public_output = gateway_public_output(dispatch.tool_type, &call, &output, status, registry);
Ok(GatewayCallResult {
call,
input_item: InputItem::FunctionCallOutput(output.into()),
Expand All @@ -185,10 +185,13 @@ fn gateway_public_output(
call: &FunctionToolCall,
output: &ToolOutput,
status: GatewayCallStatus,
registry: &ToolRegistry,
) -> Option<OutputItem> {
match tool_type {
ToolType::WebSearch => Some(crate::tool::web_search::output_item(call, output, status)),
ToolType::Mcp => Some(crate::tool::mcp::handler::output_item(call, output, status)),
ToolType::Mcp => registry
.mcp_tool_ref(&call.name)
.map(|tool_ref| crate::tool::mcp::handler::output_item(call, output, status, tool_ref)),
ToolType::Function | ToolType::CodexNamespace | ToolType::FileSearch | ToolType::CodeInterpreter => None,
}
}
Expand Down Expand Up @@ -252,7 +255,9 @@ pub(super) fn gateway_event_plans(
output_index: u32::try_from(output_index).unwrap_or(u32::MAX),
started_output: match entry.tool_type {
ToolType::WebSearch => Some(crate::tool::web_search::started_output_item(call)),
ToolType::Mcp => Some(crate::tool::mcp::handler::started_output_item(call)),
ToolType::Mcp => registry
.mcp_tool_ref(&call.name)
.map(|tool_ref| crate::tool::mcp::handler::started_output_item(call, tool_ref)),
ToolType::Function
| ToolType::CodexNamespace
| ToolType::FileSearch
Expand Down Expand Up @@ -573,7 +578,8 @@ mod tests {
serde_json::from_value(serde_json::json!({"type": "web_search_preview"})).expect("web_search tool param");
let mut executors = GatewayExecutors::default();
executors.insert(Arc::new(SlowExecutor));
let registry = ToolRegistry::build_with_handlers(&[web_search], &executors)
let mut tools = [web_search];
let registry = ToolRegistry::build_with_handlers(&mut tools, &mut executors)
.await
.expect("registry builds");

Expand Down Expand Up @@ -609,7 +615,9 @@ mod tests {
// not fail the whole request.
let web_search: ResponsesTool =
serde_json::from_value(serde_json::json!({"type": "web_search_preview"})).expect("web_search tool param");
let registry = ToolRegistry::build_with_handlers(&[web_search], &GatewayExecutors::default())
let mut tools = [web_search];
let mut executors = GatewayExecutors::default();
let registry = ToolRegistry::build_with_handlers(&mut tools, &mut executors)
.await
.expect("registry builds");

Expand Down
11 changes: 5 additions & 6 deletions crates/agentic-server-core/src/executor/messages_stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -665,11 +665,10 @@ mod tests {
/// the ONLY way `execute_gateway_calls` can produce a non-"no handler" result
/// for a malformed input is by rejecting the args before dispatch (the fix).
async fn no_op_registry() -> ToolRegistry {
ToolRegistry::build_with_handlers(
&[],
&crate::tool::GatewayExecutors::from_env(std::sync::Arc::new(reqwest::Client::new())),
)
.await
.unwrap()
let mut tools = [];
let mut executors = crate::tool::GatewayExecutors::from_env(std::sync::Arc::new(reqwest::Client::new()));
ToolRegistry::build_with_handlers(&mut tools, &mut executors)
.await
.unwrap()
}
}
11 changes: 10 additions & 1 deletion crates/agentic-server-core/src/executor/persist.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ pub(crate) async fn persist_if_needed(
/// Returns [`ExecutorError`] if the storage operation fails.
pub async fn persist_response(
payload: ResponsePayload,
ctx: RequestContext,
mut ctx: RequestContext,
conv_handler: ConversationHandler,
resp_handler: ResponseHandler,
) -> ExecutorResult<()> {
Expand All @@ -52,6 +52,15 @@ pub async fn persist_response(
return Ok(());
}

// MCP headers and authorization are request-scoped runtime credentials.
// Tool execution has completed at this point, so remove them before either
// persistence mode builds and serializes effective tool metadata.
if let Some(tools) = ctx.enriched_request.tools.as_mut() {
for tool in tools {
tool.redact_runtime_credentials();
}
}

// Move output items from payload; handlers build ResponseMetadata from ctx internally.
let output_items = payload.output;

Expand Down
4 changes: 2 additions & 2 deletions crates/agentic-server-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,8 @@ pub use storage::{
models::{Conversation as DbConversation, Item as DbItem, Response as DbResponse},
};
pub use tool::{
CodexNamespaceHandler, FunctionHandler, GatewayExecutor, McpServerEntry, ToolEntry, ToolError, ToolHandler,
ToolOutput, ToolRegistry, ToolType, WebSearchHandler,
CodexNamespaceHandler, FunctionHandler, GatewayExecutor, GatewayExecutorRegistration, McpServerEntry, ToolEntry,
ToolError, ToolHandler, ToolOutput, ToolRegistry, ToolType, WebSearchHandler,
};
pub use types::{
CodeInterpreterToolParam, CodexNamespaceMember, CodexNamespaceToolParam, CustomToolCall,
Expand Down
65 changes: 36 additions & 29 deletions crates/agentic-server-core/src/tool/codex.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ use serde_json::{Map, Value};
use crate::events::WireEvent;
use crate::types::io::{FunctionTool, FunctionToolCall, OutputItem, ToolChoice};
use crate::types::tools::{CodexNamespaceMember, CodexNamespaceToolParam, NonEmptyToolName, ResponsesTool};
use crate::utils::common::serialize_to_value_or_custom_default;

use super::handler::{ToolError, ToolHandler};
use super::registry::{ToolEntry, ToolType};
Expand Down Expand Up @@ -45,27 +46,33 @@ pub fn model_visible_namespace_member_name(namespace: &str, member: &str) -> Str
/// namespace members to those flat names first (see
/// [`CodexNamespaceHandler::resolve_namespace_members`]).
pub(crate) fn insert_namespace_entries(entries: &mut HashMap<String, ToolEntry>, p: &CodexNamespaceToolParam) {
let config = serde_json::to_value(p).expect("serialization of known struct is infallible");
for member in &p.tools {
let CodexNamespaceMember::Function(function) = member else {
continue;
};
let name = function.name.as_str().to_owned();
if entries
.insert(
name.clone(),
ToolEntry {
tool_type: ToolType::CodexNamespace,
config: config.clone(),
server_label: Some(p.name.clone()),
handler: None,
},
)
.is_some()
{
tracing::warn!(name = %name, namespace = %p.name, "duplicate tool name - previous definition overwritten");
}
}
serialize_to_value_or_custom_default(
p,
"namespace tool config serialization failed",
|config| {
for member in &p.tools {
let CodexNamespaceMember::Function(function) = member else {
continue;
};
let name = function.name.as_str().to_owned();
if entries
.insert(
name.clone(),
ToolEntry {
tool_type: ToolType::CodexNamespace,
config: config.clone(),
server_label: Some(p.name.clone()),
handler: None,
},
)
.is_some()
{
tracing::warn!(name = %name, namespace = %p.name, "duplicate tool name - previous definition overwritten");
}
}
},
(),
);
}

#[derive(Clone, Debug, Eq, Hash, PartialEq)]
Expand Down Expand Up @@ -412,11 +419,13 @@ fn typed_top_level_registry_keys(tools: &[ResponsesTool]) -> HashMap<String, Too
.filter_map(|tool| {
let registry_key = match tool {
ResponsesTool::Function(function) => function.name.as_str().to_owned(),
ResponsesTool::Mcp(mcp) => mcp.name.as_str().to_owned(),
ResponsesTool::WebSearch(_) => "web_search".to_owned(),
ResponsesTool::FileSearch(_) => "file_search".to_owned(),
ResponsesTool::CodeInterpreter(_) => "code_interpreter".to_owned(),
ResponsesTool::Namespace(_) | ResponsesTool::Custom(_) | ResponsesTool::Unknown => return None,
ResponsesTool::Mcp(_)
| ResponsesTool::Namespace(_)
| ResponsesTool::Custom(_)
| ResponsesTool::Unknown => return None,
};
tool.tool_type().map(|tool_type| (registry_key, tool_type))
})
Expand Down Expand Up @@ -774,10 +783,9 @@ mod tests {
}

#[test]
fn resolve_namespace_members_rejects_shortened_name_collision_with_later_mcp_tool() {
fn resolve_namespace_members_accepts_native_mcp_without_static_registry_key() {
let namespace = "mcp__codex_apps__github";
let member = "_remove_reaction_from_pr_review_comment";
let shortened_name = model_visible_namespace_member_name(namespace, member);
let tools: Vec<ResponsesTool> = serde_json::from_value(serde_json::json!([
{
"type": "namespace",
Expand All @@ -786,16 +794,15 @@ mod tests {
},
{
"type": "mcp",
"name": shortened_name,
"server_label": "fixture",
"server_url": "http://127.0.0.1:1/mcp"
}
]))
.unwrap();

let err = CodexNamespaceHandler.resolve_namespace_members(&tools).unwrap_err();

assert!(err.to_string().contains("collides with a declared MCP tool"));
CodexNamespaceHandler
.resolve_namespace_members(&tools)
.expect("native MCP registry keys are derived after discovery");
}

#[test]
Expand Down
Loading