From 7a5c29804c59115b35711a4146a65aa9d86fbf66 Mon Sep 17 00:00:00 2001 From: Michael Bolin Date: Wed, 13 Aug 2025 12:49:39 -0700 Subject: [PATCH] feat: support traditional JSON-RPC request/response in MCP server --- .../mcp-server/src/codex_message_processor.rs | 128 ++++++++++++++++++ codex-rs/mcp-server/src/codex_tool_config.rs | 4 +- codex-rs/mcp-server/src/lib.rs | 2 + codex-rs/mcp-server/src/message_processor.rs | 26 +++- codex-rs/mcp-server/src/wire_format.rs | 105 ++++++++++++++ 5 files changed, 261 insertions(+), 4 deletions(-) create mode 100644 codex-rs/mcp-server/src/codex_message_processor.rs create mode 100644 codex-rs/mcp-server/src/wire_format.rs diff --git a/codex-rs/mcp-server/src/codex_message_processor.rs b/codex-rs/mcp-server/src/codex_message_processor.rs new file mode 100644 index 0000000000..73f0c6b403 --- /dev/null +++ b/codex-rs/mcp-server/src/codex_message_processor.rs @@ -0,0 +1,128 @@ +use std::path::PathBuf; +use std::sync::Arc; + +use codex_core::ConversationManager; +use codex_core::NewConversation; +use codex_core::config::Config; +use codex_core::config::ConfigOverrides; +use mcp_types::JSONRPCErrorError; +use mcp_types::RequestId; + +use crate::error_code::INTERNAL_ERROR_CODE; +use crate::error_code::INVALID_REQUEST_ERROR_CODE; +use crate::json_to_toml::json_to_toml; +use crate::outgoing_message::OutgoingMessageSender; +use crate::wire_format::CodexRequest; +use crate::wire_format::ConversationId; +use crate::wire_format::NewConversationParams; +use crate::wire_format::NewConversationResponse; + +/// Handles JSON-RPC messages for Codex conversations. +pub(crate) struct CodexMessageProcessor { + conversation_manager: Arc, + outgoing: Arc, + codex_linux_sandbox_exe: Option, +} + +impl CodexMessageProcessor { + pub fn new( + conversation_manager: Arc, + outgoing: Arc, + codex_linux_sandbox_exe: Option, + ) -> Self { + Self { + conversation_manager, + outgoing, + codex_linux_sandbox_exe, + } + } + + pub async fn process_request(&self, request: CodexRequest) { + match request { + CodexRequest::NewConversation { + id: request_id, + params, + } => { + // Do not tokio::spawn() to process new_conversation() + // asynchronously because we need to ensure the conversation is + // created before processing any subsequent messages. + self.process_new_conversation(request_id, params).await; + } + } + } + + async fn process_new_conversation(&self, request_id: RequestId, params: NewConversationParams) { + let config = match derive_config(params, self.codex_linux_sandbox_exe.clone()) { + Ok(config) => config, + Err(err) => { + let error = JSONRPCErrorError { + code: INVALID_REQUEST_ERROR_CODE, + message: format!("Error deriving config: {err}"), + data: None, + }; + self.outgoing.send_error(request_id, error).await; + return; + } + }; + + match self.conversation_manager.new_conversation(config).await { + Ok(conversation_id) => { + let NewConversation { + conversation_id, + session_configured, + .. + } = conversation_id; + let response = NewConversationResponse { + conversation_id: ConversationId(conversation_id), + model: session_configured.model, + }; + self.outgoing.send_response(request_id, response).await; + } + Err(err) => { + let error = JSONRPCErrorError { + code: INTERNAL_ERROR_CODE, + message: format!("Error creating conversation: {err}"), + data: None, + }; + self.outgoing.send_error(request_id, error).await; + } + } + } +} + +fn derive_config( + params: NewConversationParams, + codex_linux_sandbox_exe: Option, +) -> std::io::Result { + let NewConversationParams { + model, + profile, + cwd, + approval_policy, + sandbox, + config: cli_overrides, + base_instructions, + include_plan_tool, + } = params; + let overrides = ConfigOverrides { + model, + config_profile: profile, + cwd: cwd.map(PathBuf::from), + approval_policy: approval_policy.map(Into::into), + sandbox_mode: sandbox.map(Into::into), + model_provider: None, + codex_linux_sandbox_exe, + base_instructions, + include_plan_tool, + disable_response_storage: None, + show_raw_agent_reasoning: None, + }; + + let cli_overrides = cli_overrides + .unwrap_or_default() + .into_iter() + .map(|(k, v)| (k, json_to_toml(v))) + .collect(); + + Config::load_with_cli_overrides(cli_overrides, overrides) +} diff --git a/codex-rs/mcp-server/src/codex_tool_config.rs b/codex-rs/mcp-server/src/codex_tool_config.rs index 4af3e29c48..1d77eb4b82 100644 --- a/codex-rs/mcp-server/src/codex_tool_config.rs +++ b/codex-rs/mcp-server/src/codex_tool_config.rs @@ -58,7 +58,7 @@ pub struct CodexToolCallParam { /// Custom enum mirroring [`AskForApproval`], but has an extra dependency on /// [`JsonSchema`]. -#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)] +#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)] #[serde(rename_all = "kebab-case")] pub enum CodexToolCallApprovalPolicy { Untrusted, @@ -80,7 +80,7 @@ impl From for AskForApproval { /// Custom enum mirroring [`SandboxMode`] from config_types.rs, but with /// `JsonSchema` support. -#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)] +#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)] #[serde(rename_all = "kebab-case")] pub enum CodexToolCallSandboxMode { ReadOnly, diff --git a/codex-rs/mcp-server/src/lib.rs b/codex-rs/mcp-server/src/lib.rs index abbe3b94a0..aa8583b77d 100644 --- a/codex-rs/mcp-server/src/lib.rs +++ b/codex-rs/mcp-server/src/lib.rs @@ -15,6 +15,7 @@ use tracing::error; use tracing::info; use tracing_subscriber::EnvFilter; +mod codex_message_processor; mod codex_tool_config; mod codex_tool_runner; mod conversation_loop; @@ -26,6 +27,7 @@ pub(crate) mod message_processor; mod outgoing_message; mod patch_approval; pub(crate) mod tool_handlers; +mod wire_format; use crate::message_processor::MessageProcessor; use crate::outgoing_message::OutgoingMessage; diff --git a/codex-rs/mcp-server/src/message_processor.rs b/codex-rs/mcp-server/src/message_processor.rs index 98beb15460..d043493c7b 100644 --- a/codex-rs/mcp-server/src/message_processor.rs +++ b/codex-rs/mcp-server/src/message_processor.rs @@ -3,6 +3,7 @@ use std::collections::HashSet; use std::path::PathBuf; use std::sync::Arc; +use crate::codex_message_processor::CodexMessageProcessor; use crate::codex_tool_config::CodexToolCallParam; use crate::codex_tool_config::CodexToolCallReplyParam; use crate::codex_tool_config::create_tool_for_codex_tool_call_param; @@ -14,6 +15,7 @@ use crate::mcp_protocol::ToolCallResponseResult; use crate::outgoing_message::OutgoingMessageSender; use crate::tool_handlers::create_conversation::handle_create_conversation; use crate::tool_handlers::send_message::handle_send_message; +use crate::wire_format::CodexRequest; use codex_core::ConversationManager; use codex_core::config::Config as CodexConfig; @@ -40,6 +42,7 @@ use tokio::task; use uuid::Uuid; pub(crate) struct MessageProcessor { + codex_message_processor: CodexMessageProcessor, outgoing: Arc, initialized: bool, codex_linux_sandbox_exe: Option, @@ -55,11 +58,19 @@ impl MessageProcessor { outgoing: OutgoingMessageSender, codex_linux_sandbox_exe: Option, ) -> Self { + let outgoing = Arc::new(outgoing); + let conversation_manager = Arc::new(ConversationManager::default()); + let codex_message_processor = CodexMessageProcessor::new( + conversation_manager.clone(), + outgoing.clone(), + codex_linux_sandbox_exe.clone(), + ); Self { - outgoing: Arc::new(outgoing), + codex_message_processor, + outgoing, initialized: false, codex_linux_sandbox_exe, - conversation_manager: Arc::new(ConversationManager::default()), + conversation_manager, running_requests_id_to_codex_uuid: Arc::new(Mutex::new(HashMap::new())), running_session_ids: Arc::new(Mutex::new(HashSet::new())), } @@ -78,6 +89,17 @@ impl MessageProcessor { } pub(crate) async fn process_request(&mut self, request: JSONRPCRequest) { + if let Ok(request_json) = serde_json::to_value(request.clone()) + && let Ok(codex_request) = serde_json::from_value::(request_json) + { + // If the request is a Codex request, handle it with the Codex + // message processor. + self.codex_message_processor + .process_request(codex_request) + .await; + return; + } + // Hold on to the ID so we can respond. let request_id = request.id.clone(); diff --git a/codex-rs/mcp-server/src/wire_format.rs b/codex-rs/mcp-server/src/wire_format.rs new file mode 100644 index 0000000000..c30f8a1c12 --- /dev/null +++ b/codex-rs/mcp-server/src/wire_format.rs @@ -0,0 +1,105 @@ +use std::collections::HashMap; + +use mcp_types::RequestId; +use serde::Deserialize; +use serde::Serialize; +use strum_macros::Display; + +use crate::codex_tool_config::CodexToolCallApprovalPolicy; +use crate::codex_tool_config::CodexToolCallSandboxMode; +use uuid::Uuid; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(transparent)] +pub(crate) struct ConversationId(pub Uuid); + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Display)] +#[serde(tag = "method", rename_all = "camelCase")] +pub(crate) enum CodexRequest { + NewConversation { + id: RequestId, + params: NewConversationParams, + }, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq)] +#[serde(rename_all = "camelCase")] +pub(crate) struct NewConversationParams { + /// Optional override for the model name (e.g. "o3", "o4-mini"). + #[serde(skip_serializing_if = "Option::is_none")] + pub(crate) model: Option, + + /// Configuration profile from config.toml to specify default options. + #[serde(skip_serializing_if = "Option::is_none")] + pub(crate) profile: Option, + + /// Working directory for the session. If relative, it is resolved against + /// the server process's current working directory. + #[serde(skip_serializing_if = "Option::is_none")] + pub(crate) cwd: Option, + + /// Approval policy for shell commands generated by the model: + /// `untrusted`, `on-failure`, `on-request`, `never`. + #[serde(skip_serializing_if = "Option::is_none")] + pub(crate) approval_policy: Option, + + /// Sandbox mode: `read-only`, `workspace-write`, or `danger-full-access`. + #[serde(skip_serializing_if = "Option::is_none")] + pub(crate) sandbox: Option, + + /// Individual config settings that will override what is in + /// CODEX_HOME/config.toml. + #[serde(skip_serializing_if = "Option::is_none")] + pub(crate) config: Option>, + + /// The set of instructions to use instead of the default ones. + #[serde(skip_serializing_if = "Option::is_none")] + pub(crate) base_instructions: Option, + + /// Whether to include the plan tool in the conversation. + #[serde(skip_serializing_if = "Option::is_none")] + pub(crate) include_plan_tool: Option, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq)] +#[serde(rename_all = "camelCase")] +pub(crate) struct NewConversationResponse { + pub(crate) conversation_id: ConversationId, + pub(crate) model: String, +} + +#[allow(clippy::unwrap_used)] +#[cfg(test)] +mod tests { + use super::*; + use pretty_assertions::assert_eq; + use serde_json::json; + + #[test] + fn serialize_new_conversation() { + let request = CodexRequest::NewConversation { + id: RequestId::Integer(42), + params: NewConversationParams { + model: None, + profile: None, + cwd: None, + approval_policy: Some(CodexToolCallApprovalPolicy::OnRequest), + sandbox: None, + config: None, + base_instructions: None, + include_plan_tool: None, + }, + }; + assert_eq!( + json!({ + "method": "newConversation", + "id": 42, + "params": { + "prompt": "Hello, Codex!", + "approvalPolicy": "on-request" + } + }), + serde_json::to_value(&request).unwrap(), + ); + } +}