ExecServer网络RPC
exec-server 的网络能力由环境所有:调用方只持有 HttpClient trait,真正发请求的进程由 local/remote 环境选择。远程路径中,ExecServerClient 通过 http/request 发送请求,executor 端的 RouteAwareHttpClient 校验 URL、构造 headers、选择重定向策略并调用共享 route-aware transport; 流式响应先返回状态和 headers,再用 http/request/bodyDelta notification 发送有序 body chunk。
本文承接ExecServer消息模型、ExecServer连接与握手 和ExecServer进程RPC。前两篇解释 RPC envelope、initialized 门槛和 notification,进程篇解释 executor 资源所有权;这里聚焦 HTTP 请求字段、local/remote transport、 body stream 生命周期和 network policy callback,不展开 HTTP 协议基础。读完后,你应能定位一次 http/request 的真实执行进程,解释为何 stream response 要求独立 request ID,并处理 body delta 乱序、背压、断连和策略拒绝。
1. 能力边界
1.1 trait facade
HttpClient 只暴露缓冲和流式两种能力。local 环境使用 RouteAwareHttpClient,remote 环境使用 ExecServerClient;上层不依赖连接类型,因此环境切换不会改变调用方的数据结构。
源码位置:codex-rs/exec-server/src/client_api.rs :: HttpClient
pub trait HttpClient: Send + Sync {
fn http_request(
&self,
params: HttpRequestParams,
) -> BoxFuture<'_, Result<HttpRequestResponse, ExecServerError>>;
fn http_request_stream(
&self,
params: HttpRequestParams,
) -> BoxFuture<'_, Result<(HttpRequestResponse, HttpResponseBodyStream), ExecServerError>>;
}HttpResponseBodyStream 对 local body 和 remote body 使用同一 recv API。local 直接读取 HTTP client 的 byte stream;remote 则从 request ID 对应的 mpsc 队列重建通知序列。
源码位置:codex-rs/exec-server/src/client/http_response_body_stream.rs :: HttpResponseBodyStream
enum HttpResponseBodyStreamInner {
Local { body: Pin<Box<dyn Stream<Item = Result<Bytes, HttpError>> + Send>> },
Remote {
inner: Arc<Inner>,
request_id: String,
next_seq: u64,
rx: mpsc::Receiver<QueuedHttpBodyDelta>,
pending_eof: bool,
closed: bool,
},
}2. 请求信封
2.1 参数字段
HttpRequestParams 是 transport-shaped 的协议对象:method 和绝对 URL 必填,headers 保持顺序, body 用 ByteChunk,timeout 可省略,redirect policy 选择两套 client pool,request_id 只在流式 响应中承担 body 路由键,stream_response 决定初始 response 是否携带 body。
源码位置:codex-rs/exec-server-protocol/src/protocol.rs :: HttpHeader、HttpRedirectPolicy、HttpRequestParams、HttpRequestResponse
pub struct HttpHeader {
pub name: String,
pub value: String,
pub value_env_var: Option<String>,
}
pub enum HttpRedirectPolicy {
Follow,
Stop,
}
pub struct HttpRequestParams {
pub method: String,
pub url: String,
pub headers: Vec<HttpHeader>,
pub body: Option<ByteChunk>,
pub timeout_ms: Option<u64>,
pub redirect_policy: HttpRedirectPolicy,
pub request_id: String,
pub stream_response: bool,
}
pub struct HttpRequestResponse {
pub status: u16,
pub headers: Vec<HttpHeader>,
pub body: ByteChunk,
}2.2 handler门槛
server registry 用 request_with_id 注册 http/request,因此 handler 同时拿到 JSON-RPC envelope 的 request ID 和 HTTP 参数。只有 stream_response 为真时才登记 params.request_id;登记失败表示同一 body stream 已经活跃。HTTP 请求完成后,handler 先用原 request ID 返回 status/headers,再异步启动 body stream task。
源码位置:
codex-rs/exec-server/src/server/registry.rs::build_router中的HTTP_REQUEST_METHODroutecodex-rs/exec-server/src/server/handler.rs::ExecServerHandler::http_request、reserve_http_body_stream、start_http_body_stream
pub(crate) async fn http_request(
self: &Arc<Self>,
request_id: RequestId,
params: HttpRequestParams,
) -> Result<(), JSONRPCErrorError> {
self.require_initialized_for("http")?;
let stream_response = params.stream_response;
let http_request_id = params.request_id.clone();
if stream_response {
self.reserve_http_body_stream(&http_request_id).await?;
}
let response = self
.http_client
.runner(params.redirect_policy)
.run(params)
.await;
if response.is_err() && stream_response {
self.release_http_body_stream(&http_request_id).await;
}
let (response, mut pending_stream) = response?;
let result = serde_json::to_value(response)
.map_err(|err| internal_error(err.to_string()))?;
self.notifications.response(request_id, result).await?;
if let Some(pending_stream) = pending_stream.take() {
self.start_http_body_stream(pending_stream).await;
}
Ok(())
}3. Executor请求
3.1 URL和headers
RouteAwareHttpRequestRunner::run 在创建请求前解析 method 和 URL,只接受 http/https scheme。 headers 逐项转换为 HeaderMap,可选的 value_env_var 在 executor 进程解析;认证 token、API key、 云身份等变量被 denylist 拒绝,空变量也会产生 invalid-params。trace headers 在请求发送前注入。
源码位置:codex-rs/exec-server/src/client/route_aware_http_client.rs :: RouteAwareHttpRequestRunner::run、build_headers
let method = Method::from_bytes(params.method.as_bytes())
.map_err(|error| invalid_params(format!("http/request method is invalid: {error}")))?;
let url = Url::parse(¶ms.url)
.map_err(|error| invalid_params(format!("http/request url is invalid: {error}")))?;
match url.scheme() {
"http" | "https" => {}
scheme => return Err(invalid_params(format!(
"http/request only supports http and https URLs, got {scheme}"
))),
}
let mut headers = Self::build_headers(params.headers)?;
codex_otel::inject_span_w3c_trace_headers(&request_span, &mut headers);3.2 缓冲响应
RouteAwareHttpClient 为 Follow 和 Stop 各维护一个 RouteAwareClientPool。请求成功后,runner 记录 status 和可转换的 response headers;非流式路径读取完整 body 并返回 HttpRequestResponse,读取 body 失败则返回 internal error。
源码位置:codex-rs/exec-server/src/client/route_aware_http_client.rs :: RouteAwareHttpClient::runner、RouteAwareHttpRequestRunner::run
if params.stream_response {
return Ok((
HttpRequestResponse {
status,
headers,
body: Vec::new().into(),
},
Some(PendingRouteAwareHttpBodyStream {
request_id: params.request_id,
response,
}),
));
}
let body = response.bytes().await.map_err(|error| {
internal_error(format!("failed to read http/request response body: {error}"))
})?;
Ok((HttpRequestResponse { status, headers, body: body.to_vec().into() }, None))4. 流式响应
4.1 executor发送
流式 runner 在拿到 headers 后返回 PendingRouteAwareHttpBodyStream。后台 task 读取 response 的 bytes_stream,将每个底层 chunk 再按 MAX_HTTP_BODY_DELTA_BYTES = 1 MiB 切分,序号从 1 递增。 正常结束发送 done: true 的空 delta;底层读取失败发送带 error 的 terminal delta。
源码位置:codex-rs/exec-server/src/client/route_aware_http_client.rs :: RouteAwareHttpRequestRunner::stream_body
let mut seq = 1;
let mut body = response.bytes_stream();
while let Some(chunk) = body.next().await {
match chunk {
Ok(bytes) => {
for chunk in bytes.chunks(MAX_HTTP_BODY_DELTA_BYTES) {
if !send_body_delta(¬ifications, HttpRequestBodyDeltaNotification {
request_id: request_id.clone(),
seq,
delta: chunk.to_vec().into(),
done: false,
error: None,
}).await {
return;
}
seq += 1;
}
}
Err(error) => {
let _ = send_body_delta(¬ifications, HttpRequestBodyDeltaNotification {
request_id,
seq,
delta: Vec::new().into(),
done: true,
error: Some(error.to_string()),
}).await;
return;
}
}
}4.2 client重建
远程 client 在发送 request 前生成连接级 request ID、创建 256 容量的 mpsc 队列并登记 stream route。 收到 bodyDelta 后,Inner::handle_http_body_delta_notification 校验 base64 解码后的大小、申请 最多 16 MiB 的排队字节预算并尝试入队;未知 ID 被有意忽略,队列满或预算耗尽则记录 stream failure 并移除 route。
源码位置:
codex-rs/exec-server/src/client/rpc_http_client.rs::ExecServerClient::http_request_streamcodex-rs/exec-server/src/client/http_response_body_stream.rs::Inner::handle_http_body_delta_notification
params.stream_response = true;
let request_id = self.inner.next_http_body_stream_request_id();
params.request_id = request_id.clone();
let (tx, rx) = mpsc::channel(HTTP_BODY_DELTA_CHANNEL_CAPACITY);
self.inner.insert_http_body_stream(request_id.clone(), tx).await?;
let response = self.call_rpc(&rpc_client, HTTP_REQUEST_METHOD, ¶ms).await?;
Ok((response, HttpResponseBodyStream::remote(
Arc::clone(&self.inner), request_id, rx,
)))HttpResponseBodyStream::recv 要求 delta.seq == next_seq。遇到 gap、error 或 done,会先移除 route; done 携带非空最后 chunk 时,当前调用先返回该 chunk,下一次 recv 才返回 EOF。
源码位置:codex-rs/exec-server/src/client/http_response_body_stream.rs :: HttpResponseBodyStream::recv
if delta.seq != *next_seq {
finish_remote_stream(inner, request_id, closed).await;
return Err(ExecServerError::Protocol(format!(
"http response stream `{request_id}` received seq {}, expected {}",
delta.seq, *next_seq
)));
}
*next_seq += 1;
if delta.done {
finish_remote_stream(inner, request_id, closed).await;
if chunk.is_empty() {
return Ok(None);
}
*pending_eof = true;
}
Ok(Some(chunk))4.3 drop和断连
如果消费者在 EOF 前 drop stream,Drop for HttpResponseBodyStream 会异步移除 route。连接断开时, fail_all_http_body_streams 向每个队列注入失败通知,避免 caller 永久等待;stream route 的释放与 HTTP request 的 JSON-RPC response 是两个独立时刻。
源码位置:codex-rs/exec-server/src/client/http_response_body_stream.rs :: Drop for HttpResponseBodyStream、Inner::fail_all_http_body_streams
5. 网络策略
5.1 executor回调
process/start 可配置 managed network proxy 和 policy decision timeout。LocalProcess::start_process 为每个 process 创建 cancellation token 和 network_policy_decider,由 executor-side 网络请求触发 反向 RPC;callback timeout 为零或 process ID 超出限制会在启动前被拒绝。HTTP capability 本身不 偷偷绕过该策略,route-aware transport 只负责发送。
源码位置:
codex-rs/exec-server/src/local_process.rs::LocalProcess::start_processcodex-rs/exec-server/src/network_policy_decisions.rs::network_policy_decidercodex-rs/exec-server-protocol/src/network_policy.rs::NetworkPolicyRequestParams、NetworkPolicyRequestResponse
let network_policy_shutdown = policy_decision_timeout.map(|_| CancellationToken::new());
let network_policy_decider = network_policy_shutdown
.as_ref()
.zip(policy_decision_timeout)
.map(|(process_shutdown, controller_timeout)| {
network_policy_decider(
process_id.clone(),
Arc::clone(&self.inner.requests),
controller_timeout,
process_shutdown.clone(),
)
});5.2 生命周期清理
process terminate 或 output close 时会取消 policy decision;managed proxy handle 由 process owner 在 maybe_emit_closed 中 shutdown。HTTP body stream 则有自己的 route、队列和字节 permit,终端 delta、 consumer drop、队列背压和 transport failure 都会释放它们。两套清理不能互相替代。
6. 验证
HTTP request 测试使用本地 listener 检查缓冲响应的 status、headers、body、fragment 移除和 stream body delta;duplicate request ID 测试确认活跃 stream 不会被覆盖。HTTP client 单测覆盖 header 环境变量 denylist、URL scheme、redirect policy 和流式序号;network policy 测试覆盖 allow/deny/ask、超时和 审计通知的 process attribution。
源码位置:
codex-rs/exec-server/tests/http_request.rs:: HTTP request 集成测试codex-rs/exec-server/src/client.rs:: HTTP body delta 与 network policy client testscodex-rs/exec-server/src/client/tests/network_policy_tests.rs:: network policy 测试
cd codex-rs
cargo test -p codex-exec-server --test http_request -- --test-threads=1
cargo test -p codex-exec-server --test http_client -- --test-threads=1
cargo test -p codex-exec-server --lib client::tests::network_policy_tests -- --test-threads=1这些测试不证明所有外部代理、TLS backend 和跨平台网络实现都在同一环境执行;它们证明的是协议字段、 stream 序号/背压、route 清理和策略回调的当前 contract。
沿源码排查“只有 headers 没有 body”时,先看 executor 的 stream_body 是否发出了从 1 开始的 seq, 再看 orchestrator 是否已登记同一 request_id、队列是否因 16 MiB 或 256 frame 上限终止,最后检查 recv 是否因 seq gap、terminal error 或 consumer drop 提前移除 route。排查“请求被拒绝”时,按 method/URL/header 解码、route/policy、transport send 和 response body 四层区分错误来源。
下一篇本地与远程执行Backend会比较 process、filesystem 和 HTTP 三种 能力在 local/remote backend 中的所有权与恢复策略。
