From dc5f035527c1f98d5065291ea93cfc323828ef2e Mon Sep 17 00:00:00 2001 From: starr-openai Date: Tue, 17 Mar 2026 01:52:06 +0000 Subject: [PATCH] test(exec-server): add unit coverage for transport and handshake Co-authored-by: Codex --- codex-rs/exec-server/src/client.rs | 246 +++++++++++++++++++ codex-rs/exec-server/src/connection.rs | 155 ++++++++++++ codex-rs/exec-server/src/server/processor.rs | 198 +++++++++++++++ 3 files changed, 599 insertions(+) diff --git a/codex-rs/exec-server/src/client.rs b/codex-rs/exec-server/src/client.rs index d68eaf3277..8aa780dd34 100644 --- a/codex-rs/exec-server/src/client.rs +++ b/codex-rs/exec-server/src/client.rs @@ -527,3 +527,249 @@ async fn handle_transport_shutdown(inner: &Arc) { process.status.mark_exited(None); } } + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + use std::time::Duration; + + use pretty_assertions::assert_eq; + use tokio::io::AsyncBufReadExt; + use tokio::io::AsyncWriteExt; + use tokio::io::BufReader; + use tokio::time::timeout; + + use super::ExecServerClient; + use super::ExecServerClientConnectOptions; + use super::ExecServerError; + use crate::protocol::EXEC_METHOD; + use crate::protocol::ExecParams; + use crate::protocol::INITIALIZE_METHOD; + use crate::protocol::INITIALIZED_METHOD; + use crate::protocol::PROTOCOL_VERSION; + use codex_app_server_protocol::JSONRPCError; + use codex_app_server_protocol::JSONRPCErrorError; + use codex_app_server_protocol::JSONRPCMessage; + use codex_app_server_protocol::JSONRPCNotification; + use codex_app_server_protocol::JSONRPCRequest; + use codex_app_server_protocol::JSONRPCResponse; + + async fn read_jsonrpc_line(lines: &mut tokio::io::Lines>) -> JSONRPCMessage + where + R: tokio::io::AsyncRead + Unpin, + { + let next_line = timeout(Duration::from_secs(1), lines.next_line()).await; + let line_result = match next_line { + Ok(line_result) => line_result, + Err(err) => panic!("timed out waiting for JSON-RPC line: {err}"), + }; + let maybe_line = match line_result { + Ok(maybe_line) => maybe_line, + Err(err) => panic!("failed to read JSON-RPC line: {err}"), + }; + let line = match maybe_line { + Some(line) => line, + None => panic!("server connection closed before JSON-RPC line arrived"), + }; + match serde_json::from_str::(&line) { + Ok(message) => message, + Err(err) => panic!("failed to parse JSON-RPC line: {err}"), + } + } + + async fn write_jsonrpc_line(writer: &mut W, message: JSONRPCMessage) + where + W: tokio::io::AsyncWrite + Unpin, + { + let encoded = match serde_json::to_string(&message) { + Ok(encoded) => encoded, + Err(err) => panic!("failed to encode JSON-RPC message: {err}"), + }; + if let Err(err) = writer.write_all(format!("{encoded}\n").as_bytes()).await { + panic!("failed to write JSON-RPC line: {err}"); + } + } + + #[tokio::test] + async fn connect_stdio_performs_initialize_handshake() { + let (client_stdin, server_reader) = tokio::io::duplex(4096); + let (mut server_writer, client_stdout) = tokio::io::duplex(4096); + + let server = tokio::spawn(async move { + let mut lines = BufReader::new(server_reader).lines(); + + let initialize = read_jsonrpc_line(&mut lines).await; + let JSONRPCMessage::Request(request) = initialize else { + panic!("expected initialize request"); + }; + assert_eq!(request.method, INITIALIZE_METHOD); + assert_eq!( + request.params, + Some(serde_json::json!({ "clientName": "test-client" })) + ); + write_jsonrpc_line( + &mut server_writer, + JSONRPCMessage::Response(JSONRPCResponse { + id: request.id, + result: serde_json::json!({ "protocolVersion": PROTOCOL_VERSION }), + }), + ) + .await; + + let initialized = read_jsonrpc_line(&mut lines).await; + let JSONRPCMessage::Notification(JSONRPCNotification { method, params }) = initialized + else { + panic!("expected initialized notification"); + }; + assert_eq!(method, INITIALIZED_METHOD); + assert_eq!(params, Some(serde_json::json!({}))); + }); + + let client = ExecServerClient::connect_stdio( + client_stdin, + client_stdout, + ExecServerClientConnectOptions { + client_name: "test-client".to_string(), + }, + ) + .await; + if let Err(err) = client { + panic!("failed to connect test client: {err}"); + } + + if let Err(err) = server.await { + panic!("server task failed: {err}"); + } + } + + #[tokio::test] + async fn connect_stdio_returns_initialize_errors() { + let (client_stdin, server_reader) = tokio::io::duplex(4096); + let (mut server_writer, client_stdout) = tokio::io::duplex(4096); + + tokio::spawn(async move { + let mut lines = BufReader::new(server_reader).lines(); + + let initialize = read_jsonrpc_line(&mut lines).await; + let JSONRPCMessage::Request(request) = initialize else { + panic!("expected initialize request"); + }; + write_jsonrpc_line( + &mut server_writer, + JSONRPCMessage::Error(JSONRPCError { + id: request.id, + error: JSONRPCErrorError { + code: -32600, + message: "rejected".to_string(), + data: None, + }, + }), + ) + .await; + }); + + let result = ExecServerClient::connect_stdio( + client_stdin, + client_stdout, + ExecServerClientConnectOptions { + client_name: "test-client".to_string(), + }, + ) + .await; + + match result { + Err(ExecServerError::Server { code, message }) => { + assert_eq!(code, -32600); + assert_eq!(message, "rejected"); + } + Err(err) => panic!("unexpected initialize failure: {err}"), + Ok(_) => panic!("expected initialize failure"), + } + } + + #[tokio::test] + async fn start_process_cleans_up_registered_process_after_request_error() { + let (client_stdin, server_reader) = tokio::io::duplex(4096); + let (mut server_writer, client_stdout) = tokio::io::duplex(4096); + + tokio::spawn(async move { + let mut lines = BufReader::new(server_reader).lines(); + + let initialize = read_jsonrpc_line(&mut lines).await; + let JSONRPCMessage::Request(initialize_request) = initialize else { + panic!("expected initialize request"); + }; + write_jsonrpc_line( + &mut server_writer, + JSONRPCMessage::Response(JSONRPCResponse { + id: initialize_request.id, + result: serde_json::json!({ "protocolVersion": PROTOCOL_VERSION }), + }), + ) + .await; + + let initialized = read_jsonrpc_line(&mut lines).await; + let JSONRPCMessage::Notification(notification) = initialized else { + panic!("expected initialized notification"); + }; + assert_eq!(notification.method, INITIALIZED_METHOD); + + let exec_request = read_jsonrpc_line(&mut lines).await; + let JSONRPCMessage::Request(JSONRPCRequest { id, method, .. }) = exec_request else { + panic!("expected exec request"); + }; + assert_eq!(method, EXEC_METHOD); + write_jsonrpc_line( + &mut server_writer, + JSONRPCMessage::Error(JSONRPCError { + id, + error: JSONRPCErrorError { + code: -32600, + message: "duplicate process".to_string(), + data: None, + }, + }), + ) + .await; + }); + + let client = match ExecServerClient::connect_stdio( + client_stdin, + client_stdout, + ExecServerClientConnectOptions { + client_name: "test-client".to_string(), + }, + ) + .await + { + Ok(client) => client, + Err(err) => panic!("failed to connect test client: {err}"), + }; + + let result = client + .start_process(ExecParams { + process_id: "proc-1".to_string(), + argv: vec!["bash".to_string(), "-lc".to_string(), "true".to_string()], + cwd: std::env::current_dir().unwrap_or_else(|err| panic!("missing cwd: {err}")), + env: HashMap::new(), + tty: true, + output_bytes_cap: 4096, + arg0: None, + }) + .await; + + match result { + Err(ExecServerError::Server { code, message }) => { + assert_eq!(code, -32600); + assert_eq!(message, "duplicate process"); + } + Err(err) => panic!("unexpected start_process failure: {err}"), + Ok(_) => panic!("expected start_process failure"), + } + + assert!( + client.inner.processes.lock().await.is_empty(), + "failed requests should not leave registered process state behind" + ); + } +} diff --git a/codex-rs/exec-server/src/connection.rs b/codex-rs/exec-server/src/connection.rs index f9a1ec669f..691a3dce72 100644 --- a/codex-rs/exec-server/src/connection.rs +++ b/codex-rs/exec-server/src/connection.rs @@ -260,3 +260,158 @@ where fn serialize_jsonrpc_message(message: &JSONRPCMessage) -> Result { serde_json::to_string(message) } + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use codex_app_server_protocol::JSONRPCMessage; + use codex_app_server_protocol::JSONRPCRequest; + use codex_app_server_protocol::JSONRPCResponse; + use codex_app_server_protocol::RequestId; + use pretty_assertions::assert_eq; + use tokio::io::AsyncBufReadExt; + use tokio::io::AsyncWriteExt; + use tokio::io::BufReader; + use tokio::sync::mpsc; + use tokio::time::timeout; + + use super::JsonRpcConnection; + use super::JsonRpcConnectionEvent; + use super::serialize_jsonrpc_message; + + async fn recv_event( + incoming_rx: &mut mpsc::Receiver, + ) -> JsonRpcConnectionEvent { + let recv_result = timeout(Duration::from_secs(1), incoming_rx.recv()).await; + let maybe_event = match recv_result { + Ok(maybe_event) => maybe_event, + Err(err) => panic!("timed out waiting for connection event: {err}"), + }; + match maybe_event { + Some(event) => event, + None => panic!("connection event stream ended unexpectedly"), + } + } + + async fn read_jsonrpc_line(lines: &mut tokio::io::Lines>) -> JSONRPCMessage + where + R: tokio::io::AsyncRead + Unpin, + { + let next_line = timeout(Duration::from_secs(1), lines.next_line()).await; + let line_result = match next_line { + Ok(line_result) => line_result, + Err(err) => panic!("timed out waiting for JSON-RPC line: {err}"), + }; + let maybe_line = match line_result { + Ok(maybe_line) => maybe_line, + Err(err) => panic!("failed to read JSON-RPC line: {err}"), + }; + let line = match maybe_line { + Some(line) => line, + None => panic!("connection closed before JSON-RPC line arrived"), + }; + match serde_json::from_str::(&line) { + Ok(message) => message, + Err(err) => panic!("failed to parse JSON-RPC line: {err}"), + } + } + + #[tokio::test] + async fn stdio_connection_reads_and_writes_jsonrpc_messages() { + let (mut writer_to_connection, connection_reader) = tokio::io::duplex(1024); + let (connection_writer, reader_from_connection) = tokio::io::duplex(1024); + let connection = + JsonRpcConnection::from_stdio(connection_reader, connection_writer, "test".to_string()); + let (outgoing_tx, mut incoming_rx) = connection.into_parts(); + + let incoming_message = JSONRPCMessage::Request(JSONRPCRequest { + id: RequestId::Integer(7), + method: "initialize".to_string(), + params: Some(serde_json::json!({ "clientName": "test-client" })), + trace: None, + }); + let encoded = match serialize_jsonrpc_message(&incoming_message) { + Ok(encoded) => encoded, + Err(err) => panic!("failed to serialize incoming message: {err}"), + }; + if let Err(err) = writer_to_connection + .write_all(format!("{encoded}\n").as_bytes()) + .await + { + panic!("failed to write to connection: {err}"); + } + + let event = recv_event(&mut incoming_rx).await; + match event { + JsonRpcConnectionEvent::Message(message) => { + assert_eq!(message, incoming_message); + } + JsonRpcConnectionEvent::Disconnected { reason } => { + panic!("unexpected disconnect event: {reason:?}"); + } + } + + let outgoing_message = JSONRPCMessage::Response(JSONRPCResponse { + id: RequestId::Integer(7), + result: serde_json::json!({ "protocolVersion": "exec-server.v0" }), + }); + if let Err(err) = outgoing_tx.send(outgoing_message.clone()).await { + panic!("failed to queue outgoing message: {err}"); + } + + let mut lines = BufReader::new(reader_from_connection).lines(); + let message = read_jsonrpc_line(&mut lines).await; + assert_eq!(message, outgoing_message); + } + + #[tokio::test] + async fn stdio_connection_reports_parse_errors() { + let (mut writer_to_connection, connection_reader) = tokio::io::duplex(1024); + let (connection_writer, _reader_from_connection) = tokio::io::duplex(1024); + let connection = + JsonRpcConnection::from_stdio(connection_reader, connection_writer, "test".to_string()); + let (_outgoing_tx, mut incoming_rx) = connection.into_parts(); + + if let Err(err) = writer_to_connection.write_all(b"not-json\n").await { + panic!("failed to write invalid JSON: {err}"); + } + + let event = recv_event(&mut incoming_rx).await; + match event { + JsonRpcConnectionEvent::Disconnected { reason } => { + let reason = match reason { + Some(reason) => reason, + None => panic!("expected a parse error reason"), + }; + assert!( + reason.contains("failed to parse JSON-RPC message from test"), + "unexpected disconnect reason: {reason}" + ); + } + JsonRpcConnectionEvent::Message(message) => { + panic!("unexpected JSON-RPC message: {message:?}"); + } + } + } + + #[tokio::test] + async fn stdio_connection_reports_clean_disconnect() { + let (writer_to_connection, connection_reader) = tokio::io::duplex(1024); + let (connection_writer, _reader_from_connection) = tokio::io::duplex(1024); + let connection = + JsonRpcConnection::from_stdio(connection_reader, connection_writer, "test".to_string()); + let (_outgoing_tx, mut incoming_rx) = connection.into_parts(); + drop(writer_to_connection); + + let event = recv_event(&mut incoming_rx).await; + match event { + JsonRpcConnectionEvent::Disconnected { reason } => { + assert_eq!(reason, None); + } + JsonRpcConnectionEvent::Message(message) => { + panic!("unexpected JSON-RPC message: {message:?}"); + } + } + } +} diff --git a/codex-rs/exec-server/src/server/processor.rs b/codex-rs/exec-server/src/server/processor.rs index d0f9252e3a..a7d1dcbcb3 100644 --- a/codex-rs/exec-server/src/server/processor.rs +++ b/codex-rs/exec-server/src/server/processor.rs @@ -466,3 +466,201 @@ fn internal_error(message: String) -> JSONRPCErrorError { message, } } + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use pretty_assertions::assert_eq; + use serde_json::json; + use tokio::time::timeout; + + use super::ExecServerConnectionProcessor; + use crate::protocol::EXEC_METHOD; + use crate::protocol::INITIALIZE_METHOD; + use crate::protocol::INITIALIZED_METHOD; + use crate::protocol::PROTOCOL_VERSION; + use codex_app_server_protocol::JSONRPCMessage; + use codex_app_server_protocol::JSONRPCNotification; + use codex_app_server_protocol::JSONRPCRequest; + use codex_app_server_protocol::RequestId; + + fn request(id: i64, method: &str, params: serde_json::Value) -> JSONRPCMessage { + JSONRPCMessage::Request(JSONRPCRequest { + id: RequestId::Integer(id), + method: method.to_string(), + params: Some(params), + trace: None, + }) + } + + async fn recv_outgoing_json( + outgoing_rx: &mut tokio::sync::mpsc::Receiver, + ) -> serde_json::Value { + let recv_result = timeout(Duration::from_secs(1), outgoing_rx.recv()).await; + let maybe_message = match recv_result { + Ok(maybe_message) => maybe_message, + Err(err) => panic!("timed out waiting for processor output: {err}"), + }; + let message = match maybe_message { + Some(message) => message, + None => panic!("processor output channel closed unexpectedly"), + }; + serde_json::to_value(message) + .unwrap_or_else(|err| panic!("failed to serialize processor output: {err}")) + } + + #[tokio::test] + async fn initialize_response_reports_protocol_version() { + let (outgoing_tx, mut outgoing_rx) = tokio::sync::mpsc::channel(1); + let mut processor = ExecServerConnectionProcessor::new(outgoing_tx); + + if let Err(err) = processor + .handle_message(request( + 1, + INITIALIZE_METHOD, + json!({ "clientName": "test" }), + )) + .await + { + panic!("initialize should succeed: {err}"); + } + + let outgoing = recv_outgoing_json(&mut outgoing_rx).await; + assert_eq!( + outgoing, + json!({ + "id": 1, + "result": { + "protocolVersion": PROTOCOL_VERSION + } + }) + ); + } + + #[tokio::test] + async fn exec_methods_require_initialize() { + let (outgoing_tx, mut outgoing_rx) = tokio::sync::mpsc::channel(1); + let mut processor = ExecServerConnectionProcessor::new(outgoing_tx); + + if let Err(err) = processor + .handle_message(request(7, EXEC_METHOD, json!({ "processId": "proc-1" }))) + .await + { + panic!("request handling should not fail the connection: {err}"); + } + + let outgoing = recv_outgoing_json(&mut outgoing_rx).await; + assert_eq!( + outgoing, + json!({ + "id": 7, + "error": { + "code": -32600, + "message": "client must call initialize before using exec methods" + } + }) + ); + } + + #[tokio::test] + async fn exec_methods_require_initialized_notification_after_initialize() { + let (outgoing_tx, mut outgoing_rx) = tokio::sync::mpsc::channel(2); + let mut processor = ExecServerConnectionProcessor::new(outgoing_tx); + + if let Err(err) = processor + .handle_message(request( + 1, + INITIALIZE_METHOD, + json!({ "clientName": "test" }), + )) + .await + { + panic!("initialize should succeed: {err}"); + } + let _ = recv_outgoing_json(&mut outgoing_rx).await; + + if let Err(err) = processor + .handle_message(request(2, EXEC_METHOD, json!({ "processId": "proc-1" }))) + .await + { + panic!("request handling should not fail the connection: {err}"); + } + + let outgoing = recv_outgoing_json(&mut outgoing_rx).await; + assert_eq!( + outgoing, + json!({ + "id": 2, + "error": { + "code": -32600, + "message": "client must send initialized before using exec methods" + } + }) + ); + } + + #[tokio::test] + async fn initialized_before_initialize_is_a_protocol_error() { + let (outgoing_tx, _outgoing_rx) = tokio::sync::mpsc::channel(1); + let mut processor = ExecServerConnectionProcessor::new(outgoing_tx); + + let result = processor + .handle_message(JSONRPCMessage::Notification(JSONRPCNotification { + method: INITIALIZED_METHOD.to_string(), + params: Some(json!({})), + })) + .await; + + match result { + Err(err) => { + assert_eq!( + err, + "received `initialized` notification before `initialize`" + ); + } + Ok(()) => panic!("expected protocol error for early initialized notification"), + } + } + + #[tokio::test] + async fn initialize_may_only_be_sent_once_per_connection() { + let (outgoing_tx, mut outgoing_rx) = tokio::sync::mpsc::channel(2); + let mut processor = ExecServerConnectionProcessor::new(outgoing_tx); + + if let Err(err) = processor + .handle_message(request( + 1, + INITIALIZE_METHOD, + json!({ "clientName": "test" }), + )) + .await + { + panic!("initialize should succeed: {err}"); + } + let _ = recv_outgoing_json(&mut outgoing_rx).await; + + if let Err(err) = processor + .handle_message(request( + 2, + INITIALIZE_METHOD, + json!({ "clientName": "test" }), + )) + .await + { + panic!("duplicate initialize should not fail the connection: {err}"); + } + + let outgoing = recv_outgoing_json(&mut outgoing_rx).await; + assert_eq!( + outgoing, + json!({ + "id": 2, + "error": { + "code": -32600, + "message": "initialize may only be sent once per connection" + } + }) + ); + } +}