mirror of
https://github.com/openai/codex.git
synced 2026-08-24 13:20:07 +00:00
## Why Sequential dispatch lets a long-running request block unrelated health checks and cleanup on the same connection. ## What changed - Add `--concurrent-requests <COUNT>` for local and remote exec-server connections, while retaining sequential dispatch when the option is omitted or set to `1`. - Preserve handshake ordering before enabling concurrent dispatch. - Reserve separate capacity for status, signal, terminate, and close requests so they remain responsive when ordinary request capacity is saturated. - Drain queued client responses during disconnect and cancel outstanding request tasks during connection shutdown. ## Testing - Cover CLI parsing and concurrency-limit validation. - Verify default sequential behavior, pipelined handshake ordering, concurrent request progress, control-request responsiveness, and disconnect handling. GitOrigin-RevId: 48e4b092e318204ed635543f01f9ee0e7df095fc
99 lines
3.2 KiB
Rust
99 lines
3.2 KiB
Rust
mod common;
|
|
|
|
use codex_exec_server::InitializeParams;
|
|
use codex_exec_server::InitializeResponse;
|
|
use codex_exec_server_protocol::JSONRPCError;
|
|
use codex_exec_server_protocol::JSONRPCErrorError;
|
|
use codex_exec_server_protocol::JSONRPCMessage;
|
|
use codex_exec_server_protocol::JSONRPCResponse;
|
|
use common::exec_server::exec_server;
|
|
use common::exec_server::exec_server_with_env;
|
|
use pretty_assertions::assert_eq;
|
|
use uuid::Uuid;
|
|
|
|
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
|
async fn exec_server_accepts_initialize() -> anyhow::Result<()> {
|
|
let mut server = exec_server().await?;
|
|
let initialize_id = server
|
|
.send_request(
|
|
"initialize",
|
|
serde_json::to_value(InitializeParams {
|
|
client_name: "exec-server-test".to_string(),
|
|
resume_session_id: None,
|
|
})?,
|
|
)
|
|
.await?;
|
|
|
|
let response = server.next_event().await?;
|
|
let JSONRPCMessage::Response(JSONRPCResponse { id, result }) = response else {
|
|
panic!("expected initialize response");
|
|
};
|
|
assert_eq!(id, initialize_id);
|
|
let initialize_response: InitializeResponse = serde_json::from_value(result)?;
|
|
Uuid::parse_str(&initialize_response.session_id)?;
|
|
|
|
server.shutdown().await?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Requests retain their wire-order initialization errors even when later handshake messages are pipelined.
|
|
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
|
async fn exec_server_rejects_pipelined_requests_before_initialized() -> anyhow::Result<()> {
|
|
let mut server = exec_server_with_env(
|
|
std::iter::empty::<(&str, &str)>(),
|
|
&["--concurrent-requests", "32"],
|
|
)
|
|
.await?;
|
|
let before_initialize_id = server
|
|
.send_request("environment/info", serde_json::json!({}))
|
|
.await?;
|
|
let initialize_id = server
|
|
.send_request(
|
|
"initialize",
|
|
serde_json::to_value(InitializeParams {
|
|
client_name: "exec-server-test".to_string(),
|
|
resume_session_id: None,
|
|
})?,
|
|
)
|
|
.await?;
|
|
|
|
assert_eq!(
|
|
server.next_event().await?,
|
|
JSONRPCMessage::Error(JSONRPCError {
|
|
id: before_initialize_id,
|
|
error: JSONRPCErrorError {
|
|
code: -32600,
|
|
data: None,
|
|
message: "client must call initialize before using environment info methods"
|
|
.to_string(),
|
|
},
|
|
})
|
|
);
|
|
let JSONRPCMessage::Response(JSONRPCResponse { id, .. }) = server.next_event().await? else {
|
|
panic!("expected initialize response");
|
|
};
|
|
assert_eq!(id, initialize_id);
|
|
|
|
let before_initialized_id = server
|
|
.send_request("environment/info", serde_json::json!({}))
|
|
.await?;
|
|
server
|
|
.send_notification("initialized", serde_json::json!({}))
|
|
.await?;
|
|
assert_eq!(
|
|
server.next_event().await?,
|
|
JSONRPCMessage::Error(JSONRPCError {
|
|
id: before_initialized_id,
|
|
error: JSONRPCErrorError {
|
|
code: -32600,
|
|
data: None,
|
|
message: "client must send initialized before using environment info methods"
|
|
.to_string(),
|
|
},
|
|
})
|
|
);
|
|
|
|
server.shutdown().await?;
|
|
Ok(())
|
|
}
|