diff --git a/codex-rs/exec-server/README.md b/codex-rs/exec-server/README.md index 8aeb3a73ec..af21dc568a 100644 --- a/codex-rs/exec-server/README.md +++ b/codex-rs/exec-server/README.md @@ -3,6 +3,10 @@ `codex-exec-server` is a small standalone stdio JSON-RPC server for spawning and controlling subprocesses through `codex-utils-pty`. +This PR intentionally lands only the standalone binary, client, wire protocol, +and docs. Exec and filesystem methods are stubbed server-side here and are +implemented in follow-up PRs. + It currently provides: - a standalone binary: `codex-exec-server` @@ -36,10 +40,7 @@ Each connection follows this sequence: 1. Send `initialize`. 2. Wait for the `initialize` response. 3. Send `initialized`. -4. Start and manage processes with `command/exec`, `command/exec/write`, and - `command/exec/terminate`. -5. Read streaming notifications from `command/exec/outputDelta` and - `command/exec/exited`. +4. Call exec or filesystem RPCs once the follow-up implementation PRs land. If the server receives any notification other than `initialized`, it replies with an error using request id `-1`. diff --git a/codex-rs/exec-server/src/server.rs b/codex-rs/exec-server/src/server.rs index 56d2206b0c..f802a509d0 100644 --- a/codex-rs/exec-server/src/server.rs +++ b/codex-rs/exec-server/src/server.rs @@ -1,420 +1,137 @@ -use std::collections::HashMap; -use std::collections::VecDeque; -use std::sync::Arc; -use std::sync::Mutex as StdMutex; - 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; use codex_app_server_protocol::RequestId; -use codex_utils_pty::ExecCommandSession; -use codex_utils_pty::TerminalSize; -use serde::Serialize; use tokio::io::AsyncBufReadExt; use tokio::io::AsyncWriteExt; use tokio::io::BufReader; -use tokio::io::BufWriter; -use tokio::sync::Mutex; -use crate::protocol::EXEC_EXITED_METHOD; -use crate::protocol::EXEC_METHOD; -use crate::protocol::EXEC_OUTPUT_DELTA_METHOD; -use crate::protocol::EXEC_TERMINATE_METHOD; -use crate::protocol::EXEC_WRITE_METHOD; -use crate::protocol::ExecExitedNotification; -use crate::protocol::ExecOutputDeltaNotification; -use crate::protocol::ExecOutputStream; -use crate::protocol::ExecParams; -use crate::protocol::ExecResponse; use crate::protocol::INITIALIZE_METHOD; use crate::protocol::INITIALIZED_METHOD; use crate::protocol::InitializeResponse; use crate::protocol::PROTOCOL_VERSION; -use crate::protocol::TerminateParams; -use crate::protocol::TerminateResponse; -use crate::protocol::WriteParams; -use crate::protocol::WriteResponse; - -struct RunningProcess { - session: ExecCommandSession, - tty: bool, - stdout_buffer: Arc>, - stderr_buffer: Arc>, -} - -#[derive(Debug)] -struct BoundedBytesBuffer { - max_bytes: usize, - bytes: VecDeque, -} - -impl BoundedBytesBuffer { - fn new(max_bytes: usize) -> Self { - Self { - max_bytes, - bytes: VecDeque::with_capacity(max_bytes.min(8192)), - } - } - - fn push_chunk(&mut self, chunk: &[u8]) { - if self.max_bytes == 0 { - return; - } - for byte in chunk { - self.bytes.push_back(*byte); - if self.bytes.len() > self.max_bytes { - self.bytes.pop_front(); - } - } - } - - fn snapshot(&self) -> Vec { - self.bytes.iter().copied().collect() - } -} pub async fn run_main() -> Result<(), Box> { - let writer = Arc::new(Mutex::new(BufWriter::new(tokio::io::stdout()))); - let processes = Arc::new(Mutex::new(HashMap::::new())); - let mut lines = BufReader::new(tokio::io::stdin()).lines(); + let mut stdin = BufReader::new(tokio::io::stdin()).lines(); + let mut stdout = tokio::io::stdout(); - while let Some(line) = lines.next_line().await? { + while let Some(line) = stdin.next_line().await? { if line.trim().is_empty() { continue; } let message = serde_json::from_str::(&line)?; - if let JSONRPCMessage::Request(request) = message { - handle_request(request, &writer, &processes).await; - continue; - } - - if let JSONRPCMessage::Notification(notification) = message { - if notification.method != INITIALIZED_METHOD { - send_error( - &writer, - RequestId::Integer(-1), - invalid_request(format!( - "unexpected notification method: {}", - notification.method - )), - ) - .await; + match message { + JSONRPCMessage::Request(request) => { + handle_request(request, &mut stdout).await?; + } + JSONRPCMessage::Notification(notification) => { + if notification.method != INITIALIZED_METHOD { + send_error( + &mut stdout, + RequestId::Integer(-1), + invalid_request(format!( + "unexpected notification method: {}", + notification.method + )), + ) + .await?; + } + } + JSONRPCMessage::Response(response) => { + send_error( + &mut stdout, + response.id, + invalid_request("unexpected response from client".to_string()), + ) + .await?; + } + JSONRPCMessage::Error(error) => { + send_error( + &mut stdout, + error.id, + invalid_request("unexpected error from client".to_string()), + ) + .await?; } - continue; } } - let remaining = { - let mut processes = processes.lock().await; - processes - .drain() - .map(|(_, process)| process) - .collect::>() - }; - for process in remaining { - process.session.terminate(); - } - Ok(()) } async fn handle_request( request: JSONRPCRequest, - writer: &Arc>>, - processes: &Arc>>, -) { - let response = match request.method.as_str() { - INITIALIZE_METHOD => serde_json::to_value(InitializeResponse { - protocol_version: PROTOCOL_VERSION.to_string(), - }) - .map_err(|err| internal_error(err.to_string())), - EXEC_METHOD => handle_exec_request(request.params, writer, processes).await, - EXEC_WRITE_METHOD => handle_write_request(request.params, processes).await, - EXEC_TERMINATE_METHOD => handle_terminate_request(request.params, processes).await, - other => Err(invalid_request(format!("unknown method: {other}"))), - }; + stdout: &mut tokio::io::Stdout, +) -> Result<(), std::io::Error> { + match request.method.as_str() { + INITIALIZE_METHOD => { + let result = serde_json::to_value(InitializeResponse { + protocol_version: PROTOCOL_VERSION.to_string(), + }) + .map_err(std::io::Error::other)?; - match response { - Ok(result) => { send_response( - writer, + stdout, JSONRPCResponse { id: request.id, result, }, ) - .await; - } - Err(err) => { - send_error(writer, request.id, err).await; - } - } -} - -async fn handle_exec_request( - params: Option, - writer: &Arc>>, - processes: &Arc>>, -) -> Result { - let params: ExecParams = serde_json::from_value(params.unwrap_or(serde_json::Value::Null)) - .map_err(|err| invalid_params(err.to_string()))?; - - let (program, args) = params - .argv - .split_first() - .ok_or_else(|| invalid_params("argv must not be empty".to_string()))?; - - let spawned = if params.tty { - codex_utils_pty::spawn_pty_process( - program, - args, - params.cwd.as_path(), - ¶ms.env, - ¶ms.arg0, - TerminalSize::default(), - ) - .await - } else { - codex_utils_pty::spawn_pipe_process_no_stdin( - program, - args, - params.cwd.as_path(), - ¶ms.env, - ¶ms.arg0, - ) - .await - } - .map_err(|err| internal_error(err.to_string()))?; - - let stdout_buffer = Arc::new(StdMutex::new(BoundedBytesBuffer::new( - params.output_bytes_cap, - ))); - let stderr_buffer = Arc::new(StdMutex::new(BoundedBytesBuffer::new( - params.output_bytes_cap, - ))); - - let process_id = params.process_id.clone(); - { - let mut process_map = processes.lock().await; - if process_map.contains_key(&process_id) { - spawned.session.terminate(); - return Err(invalid_request(format!( - "process {} already exists", - params.process_id - ))); - } - process_map.insert( - process_id.clone(), - RunningProcess { - session: spawned.session, - tty: params.tty, - stdout_buffer: Arc::clone(&stdout_buffer), - stderr_buffer: Arc::clone(&stderr_buffer), - }, - ); - } - - tokio::spawn(stream_output( - process_id.clone(), - ExecOutputStream::Stdout, - spawned.stdout_rx, - Arc::clone(writer), - Arc::clone(&stdout_buffer), - )); - tokio::spawn(stream_output( - process_id.clone(), - ExecOutputStream::Stderr, - spawned.stderr_rx, - Arc::clone(writer), - Arc::clone(&stderr_buffer), - )); - tokio::spawn(watch_exit( - process_id.clone(), - spawned.exit_rx, - Arc::clone(writer), - Arc::clone(processes), - )); - - serde_json::to_value(ExecResponse { - process_id, - running: true, - exit_code: None, - stdout: None, - stderr: None, - }) - .map_err(|err| internal_error(err.to_string())) -} - -async fn handle_write_request( - params: Option, - processes: &Arc>>, -) -> Result { - let params: WriteParams = serde_json::from_value(params.unwrap_or(serde_json::Value::Null)) - .map_err(|err| invalid_params(err.to_string()))?; - - let writer_tx = { - let process_map = processes.lock().await; - let process = process_map - .get(¶ms.process_id) - .ok_or_else(|| invalid_request(format!("unknown process id {}", params.process_id)))?; - if !process.tty { - return Err(invalid_request(format!( - "stdin is closed for process {}", - params.process_id - ))); - } - process.session.writer_sender() - }; - - writer_tx - .send(params.chunk.into_inner()) - .await - .map_err(|_| internal_error("failed to write to process stdin".to_string()))?; - - serde_json::to_value(WriteResponse { accepted: true }) - .map_err(|err| internal_error(err.to_string())) -} - -async fn handle_terminate_request( - params: Option, - processes: &Arc>>, -) -> Result { - let params: TerminateParams = serde_json::from_value(params.unwrap_or(serde_json::Value::Null)) - .map_err(|err| invalid_params(err.to_string()))?; - - let process = { - let mut process_map = processes.lock().await; - process_map.remove(¶ms.process_id) - }; - - if let Some(process) = process { - process.session.terminate(); - serde_json::to_value(TerminateResponse { running: true }) - .map_err(|err| internal_error(err.to_string())) - } else { - serde_json::to_value(TerminateResponse { running: false }) - .map_err(|err| internal_error(err.to_string())) - } -} - -async fn stream_output( - process_id: String, - stream: ExecOutputStream, - mut receiver: tokio::sync::mpsc::Receiver>, - writer: Arc>>, - buffer: Arc>, -) { - while let Some(chunk) = receiver.recv().await { - if let Ok(mut guard) = buffer.lock() { - guard.push_chunk(&chunk); - } - let notification = ExecOutputDeltaNotification { - process_id: process_id.clone(), - stream, - chunk: chunk.into(), - }; - if send_notification(&writer, EXEC_OUTPUT_DELTA_METHOD, ¬ification) .await - .is_err() - { - break; + } + method => { + send_error( + stdout, + request.id, + method_not_implemented(format!( + "exec-server stub does not implement `{method}` yet" + )), + ) + .await } } } -async fn watch_exit( - process_id: String, - exit_rx: tokio::sync::oneshot::Receiver, - writer: Arc>>, - processes: Arc>>, -) { - let exit_code = exit_rx.await.unwrap_or(-1); - let removed = { - let mut processes = processes.lock().await; - processes.remove(&process_id) - }; - if let Some(process) = removed { - let _ = process.stdout_buffer.lock().map(|buffer| buffer.snapshot()); - let _ = process.stderr_buffer.lock().map(|buffer| buffer.snapshot()); - } - let _ = send_notification( - &writer, - EXEC_EXITED_METHOD, - &ExecExitedNotification { - process_id, - exit_code, - }, - ) - .await; -} - async fn send_response( - writer: &Arc>>, + stdout: &mut tokio::io::Stdout, response: JSONRPCResponse, -) { - let _ = send_message(writer, JSONRPCMessage::Response(response)).await; +) -> Result<(), std::io::Error> { + send_message(stdout, &JSONRPCMessage::Response(response)).await } async fn send_error( - writer: &Arc>>, + stdout: &mut tokio::io::Stdout, id: RequestId, error: JSONRPCErrorError, -) { - let _ = send_message(writer, JSONRPCMessage::Error(JSONRPCError { error, id })).await; -} - -async fn send_notification( - writer: &Arc>>, - method: &str, - params: &T, -) -> Result<(), serde_json::Error> { - send_message( - writer, - JSONRPCMessage::Notification(JSONRPCNotification { - method: method.to_string(), - params: Some(serde_json::to_value(params)?), - }), - ) - .await - .map_err(serde_json::Error::io) +) -> Result<(), std::io::Error> { + send_message(stdout, &JSONRPCMessage::Error(JSONRPCError { id, error })).await } async fn send_message( - writer: &Arc>>, - message: JSONRPCMessage, -) -> std::io::Result<()> { - let encoded = - serde_json::to_vec(&message).map_err(|err| std::io::Error::other(err.to_string()))?; - let mut writer = writer.lock().await; - writer.write_all(&encoded).await?; - writer.write_all(b"\n").await?; - writer.flush().await + stdout: &mut tokio::io::Stdout, + message: &JSONRPCMessage, +) -> Result<(), std::io::Error> { + let encoded = serde_json::to_vec(message).map_err(std::io::Error::other)?; + stdout.write_all(&encoded).await?; + stdout.write_all(b"\n").await?; + stdout.flush().await } fn invalid_request(message: String) -> JSONRPCErrorError { JSONRPCErrorError { code: -32600, - data: None, message, + data: None, } } -fn invalid_params(message: String) -> JSONRPCErrorError { +fn method_not_implemented(message: String) -> JSONRPCErrorError { JSONRPCErrorError { - code: -32602, - data: None, - message, - } -} - -fn internal_error(message: String) -> JSONRPCErrorError { - JSONRPCErrorError { - code: -32603, - data: None, + code: -32601, message, + data: None, } } diff --git a/codex-rs/exec-server/tests/stdio_smoke.rs b/codex-rs/exec-server/tests/stdio_smoke.rs index 79f264327c..f2095129a6 100644 --- a/codex-rs/exec-server/tests/stdio_smoke.rs +++ b/codex-rs/exec-server/tests/stdio_smoke.rs @@ -8,9 +8,6 @@ use codex_app_server_protocol::JSONRPCNotification; use codex_app_server_protocol::JSONRPCRequest; use codex_app_server_protocol::JSONRPCResponse; use codex_app_server_protocol::RequestId; -use codex_exec_server::ExecParams; -use codex_exec_server::ExecServerClient; -use codex_exec_server::ExecServerLaunchCommand; use codex_exec_server::InitializeParams; use codex_exec_server::InitializeResponse; use codex_utils_cargo_bin::cargo_bin; @@ -19,7 +16,6 @@ use tokio::io::AsyncBufReadExt; use tokio::io::AsyncWriteExt; use tokio::io::BufReader; use tokio::process::Command; -use tokio::sync::broadcast; use tokio::time::timeout; #[tokio::test(flavor = "multi_thread", worker_threads = 2)] @@ -70,72 +66,62 @@ async fn exec_server_accepts_initialize_over_stdio() -> anyhow::Result<()> { } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn exec_server_client_streams_output_and_accepts_writes() -> anyhow::Result<()> { - let mut env = std::collections::HashMap::new(); - if let Some(path) = std::env::var_os("PATH") { - env.insert("PATH".to_string(), path.to_string_lossy().into_owned()); - } +async fn exec_server_stubs_command_exec_over_stdio() -> anyhow::Result<()> { + let binary = cargo_bin("codex-exec-server")?; + let mut child = Command::new(binary); + child.stdin(Stdio::piped()); + child.stdout(Stdio::piped()); + child.stderr(Stdio::inherit()); + let mut child = child.spawn()?; - let client = ExecServerClient::spawn(ExecServerLaunchCommand { - program: cargo_bin("codex-exec-server")?, - args: Vec::new(), - }) - .await?; + let mut stdin = child.stdin.take().expect("stdin"); + let stdout = child.stdout.take().expect("stdout"); + let mut stdout = BufReader::new(stdout).lines(); - let process = client - .start_process(ExecParams { - process_id: "2001".to_string(), - argv: vec![ - "bash".to_string(), - "-lc".to_string(), - "printf 'ready\\n'; while IFS= read -r line; do printf 'echo:%s\\n' \"$line\"; done" - .to_string(), - ], - cwd: std::env::current_dir()?, - env, - tty: true, - output_bytes_cap: 4096, - arg0: None, - }) + let initialize = JSONRPCMessage::Request(JSONRPCRequest { + id: RequestId::Integer(1), + method: "initialize".to_string(), + params: Some(serde_json::to_value(InitializeParams { + client_name: "exec-server-test".to_string(), + })?), + trace: None, + }); + stdin + .write_all(format!("{}\n", serde_json::to_string(&initialize)?).as_bytes()) + .await?; + let _ = timeout(Duration::from_secs(5), stdout.next_line()).await??; + + let exec = JSONRPCMessage::Request(JSONRPCRequest { + id: RequestId::Integer(2), + method: "command/exec".to_string(), + params: Some(serde_json::json!({ + "processId": "proc-1", + "argv": ["true"], + "cwd": std::env::current_dir()?, + "env": {}, + "tty": false, + "arg0": null + })), + trace: None, + }); + stdin + .write_all(format!("{}\n", serde_json::to_string(&exec)?).as_bytes()) .await?; - let mut output = process.output_receiver(); - assert!( - recv_until_contains(&mut output, "ready") - .await? - .contains("ready"), - "expected initial ready output" + let response_line = timeout(Duration::from_secs(5), stdout.next_line()).await??; + let response_line = response_line.expect("exec response line"); + let response: JSONRPCMessage = serde_json::from_str(&response_line)?; + let JSONRPCMessage::Error(codex_app_server_protocol::JSONRPCError { id, error }) = response + else { + panic!("expected command/exec stub error"); + }; + assert_eq!(id, RequestId::Integer(2)); + assert_eq!(error.code, -32601); + assert_eq!( + error.message, + "exec-server stub does not implement `command/exec` yet" ); - process - .writer_sender() - .send(b"hello\n".to_vec()) - .await - .expect("write should succeed"); - - assert!( - recv_until_contains(&mut output, "echo:hello") - .await? - .contains("echo:hello"), - "expected echoed output" - ); - - process.terminate(); + child.start_kill()?; Ok(()) } - -async fn recv_until_contains( - output: &mut broadcast::Receiver>, - needle: &str, -) -> anyhow::Result { - let deadline = tokio::time::Instant::now() + Duration::from_secs(5); - let mut collected = String::new(); - loop { - let remaining = deadline.saturating_duration_since(tokio::time::Instant::now()); - let chunk = timeout(remaining, output.recv()).await??; - collected.push_str(&String::from_utf8_lossy(&chunk)); - if collected.contains(needle) { - return Ok(collected); - } - } -}