AppServer架构总览
本文是 App Server 系列的入口,面向已经读过JSON-RPC封装与错误模型和协议映射与事件转换的读者。默认你了解 request/response、notification、Thread 和 Turn。本文回答一条 App Server 请求如何穿过 transport、connection state、message processor、业务 processor 和 Core,再沿 outgoing queue 返回;不展开某个具体 Thread 或 transport 的全部实现。
1. 组件地图
App Server 的关键边界不是 crate 名称,而是状态所有权:transport 管连接和写队列,message processor 管协议分派,request processor 管业务请求,ThreadManager/Core 管执行状态,ThreadState 管向客户端投影的增量。
2. Connection状态
源码位置:codex-rs/app-server/src/message_processor.rs :: ConnectionSessionState、codex-rs/app-server/src/transport.rs :: ConnectionState
pub(crate) struct ConnectionSessionState {
pub(crate) rpc_gate: Arc<ConnectionRpcGate>,
pub(crate) mcp_event_streams: McpEventStreams,
initialized: OnceLock<InitializedConnectionSessionState>,
}
pub(crate) struct InitializedConnectionSessionState {
pub(crate) experimental_api_enabled: bool,
pub(crate) opted_out_notification_methods: HashSet<String>,
pub(crate) app_server_client_name: String,
pub(crate) client_version: String,
pub(crate) request_attestation: bool,
pub(crate) client_mcp_extensions: ClientMcpExtensions,
}初始化前后能力不同;experimental capability、notification opt-out 和 client info 都是 connection/session scope,而不是 Core 全局设置。
源码位置:codex-rs/app-server/src/message_processor.rs :: ConnectionSessionState::new、initialized
impl ConnectionSessionState {
pub(crate) fn new() -> Self {
Self {
rpc_gate: Arc::new(ConnectionRpcGate::new()),
mcp_event_streams: McpEventStreams::default(),
initialized: OnceLock::new(),
}
}
pub(crate) fn initialized(&self) -> bool {
self.initialized.get().is_some()
}
pub(crate) fn experimental_api_enabled(&self) -> bool {
self.initialized
.get()
.is_some_and(|session| session.experimental_api_enabled)
}
}OnceLock 让 initialize 只能成功一次;request gate 和 MCP event streams 则从连接创建时就存在,分别负责请求排队/关闭和连接级 MCP 订阅生命周期。
3. MessageProcessor
源码位置:codex-rs/app-server/src/message_processor.rs :: MessageProcessor
pub(crate) struct MessageProcessor {
outgoing: Arc<OutgoingMessageSender>,
account_processor: AccountRequestProcessor,
config_processor: ConfigRequestProcessor,
fs_processor: FsRequestProcessor,
mcp_processor: McpRequestProcessor,
thread_processor: ThreadRequestProcessor,
turn_processor: TurnRequestProcessor,
process_exec_processor: ProcessExecRequestProcessor,
request_serialization_queues: RequestSerializationQueues,
}MessageProcessor 是装配和分派中心,不拥有所有业务状态。它把 typed request 交给专门 processor,并持有 request serialization queues 处理需要串行化的操作。
4. 请求路径
源码位置:codex-rs/app-server/src/message_processor.rs :: deserialize_client_request、request dispatch
fn deserialize_client_request(
request: JSONRPCRequest,
) -> Result<ClientRequest, JSONRPCErrorError> {
ClientRequest::try_from(request)
.map_err(|err| invalid_request(format!("Invalid request: {err}")))
}请求先从 JSON-RPC envelope 转成 typed request,再经过 experimental gate、初始化状态和 serialization queue,最后进入 processor。解析失败不会触碰 Core。
源码位置:codex-rs/app-server/src/message_processor.rs :: process_request
pub(crate) async fn process_request(
self: &Arc<Self>,
connection_id: ConnectionId,
request: JSONRPCRequest,
transport: &AppServerTransport,
session: Arc<ConnectionSessionState>,
) {
let request_id = ConnectionRequestId {
connection_id,
request_id: request.id.clone(),
};
let request_span =
crate::app_server_tracing::request_span(&request, transport, connection_id, &session);
let request_trace = request.trace.as_ref().map(|trace| W3cTraceContext {
traceparent: trace.traceparent.clone(),
tracestate: trace.tracestate.clone(),
});
let request_context = RequestContext::new(request_id.clone(), request_span, request_trace);
Self::run_request_with_context(
Arc::clone(&self.outgoing),
request_context.clone(),
async {
let result = match deserialize_client_request(request) {
Ok(codex_request) => {
self.handle_client_request(
request_id.clone(),
codex_request,
Arc::clone(&session),
None,
request_context,
)
.await
}
Err(error) => Err(error),
};
if let Err(error) = result {
self.outgoing.send_error(request_id, error).await;
}
},
)
.await;
}process_request 的职责是建立 request span、构造 RequestContext、解码 envelope 和统一发送错误;它不直接决定某个 API 的业务行为。in-process 调用走 process_client_request,跳过 JSON 解码但复用同一个 handle_client_request 语义。
5. Core与线程状态
Core 返回 EventMsg,App Server 的 bespoke event handling 将其映射为 V2 notification 或 server request;ThreadState 保存 listener、Turn summary、当前 agent message 等客户端投影所需状态。Rollout store 是持久化来源,但不是所有实时通知都从 store 读取。
源码位置:codex-rs/app-server/src/message_processor.rs :: dispatch_initialized_client_request
if !session.initialized() {
return Err(invalid_request("Not initialized"));
}
if let Some(reason) = codex_request.experimental_reason()
&& !session.experimental_api_enabled()
{
return Err(invalid_request(experimental_required_message(reason)));
}
let serialization_scope = codex_request.serialization_scope();
let request = QueuedInitializedRequest::new(
Arc::clone(&session.rpc_gate),
async move {
processor
.handle_initialized_client_request(
connection_request_id,
codex_request,
request_context,
session,
event_stream_ready,
)
.await
},
);
if let Some(scope) = serialization_scope {
let (key, access) = RequestSerializationQueueKey::from_scope(connection_id, scope);
processor.request_serialization_queues.enqueue(key, access, request).await;
} else {
tokio::spawn(async move { request.run().await; });
}初始化后的请求先经过 capability gate,再根据 serialization scope 选择串行队列或直接 spawn。scope 约束同一 Thread、进程句柄、FS watch 或 MCP OAuth 资源的并发顺序,但不会出现在 JSON wire payload 中。
6. Outgoing队列
源码位置:codex-rs/app-server/src/transport.rs :: send_message_to_connection
let message = filter_outgoing_message_for_connection(connection_state, message);
if should_skip_notification_for_connection(connection_state, &message) {
return false;
}
let queued_message = QueuedOutgoingMessage { message, write_complete_tx };
match writer.try_send(queued_message) {
Ok(()) => false,
Err(mpsc::error::TrySendError::Full(_)) => {
warn!("outgoing queue is full");
true
}
Err(mpsc::error::TrySendError::Closed(_)) => true,
}发送前按 connection capability 和 opt-out 过滤;写队列满或关闭时消息可能丢失,不能把 Core 已产生事件等同于客户端已观察到。
7. 断连边界
断连会移除 connection-scoped outbound state,并触发 disconnect cleanup;Thread/Core 资源是否继续运行由 Thread owner 和 request 类型决定。连接断开不是自动回滚业务副作用的事务边界。
源码位置:codex-rs/app-server/src/message_processor.rs :: connection_closed
pub(crate) async fn connection_closed(
&self,
connection_id: ConnectionId,
session_state: &ConnectionSessionState,
) {
session_state.rpc_gate.close().await;
session_state.mcp_event_streams.clear().await;
if timeout(
CONNECTION_RPC_DRAIN_TIMEOUT,
session_state.rpc_gate.shutdown(),
)
.await
.is_err()
{
tracing::warn!(?connection_id, "timed out waiting for connection RPCs to drain");
}
self.outgoing.connection_closed(connection_id).await;
self.fs_processor.connection_closed(connection_id).await;
self.command_exec_processor.connection_closed(connection_id).await;
self.process_exec_processor.connection_closed(connection_id).await;
self.thread_processor.connection_closed(connection_id).await;
}清理顺序先阻止新 RPC、清空 MCP subscriptions,再等待已排队 RPC 结束,最后按资源 owner 清理 outgoing、文件 watch、command/process session 和 Thread listener。连接关闭因此是资源回收触发点,但不是 Core Turn 的回滚点。
8. 源码验证
message processor tracing 测试验证 request method 和连接上下文;in-process 测试验证 Thread start response 与 notification;transport 测试验证 experimental/opt-out notification 过滤和 writer queue 行为。它们证明组件边界和消息路径,不证明所有 processor 的业务错误都已覆盖。
源码位置:
codex-rs/app-server/src/message_processor_tracing_tests.rs :: thread_start_jsonrpc_span_exports_server_span_and_parents_childrencodex-rs/app-server/src/in_process.rs :: process_client_requestcodex-rs/app-server/src/transport_tests.rs :: experimental_notifications_are_dropped_without_capabilitycodex-rs/app-server/src/transport_tests.rs :: broadcast_does_not_block_on_slow_connection
cd codex-rs
cargo test -p codex-app-server thread_start_jsonrpc_span_exports_server_span_and_parents_children -- --nocapture --test-threads=1
cargo test -p codex-app-server experimental_notifications_are_dropped_without_capability --lib -- --nocapture --test-threads=1
cargo test -p codex-app-server broadcast_does_not_block_on_slow_connection --lib -- --nocapture --test-threads=19. 阅读路径
继续阅读时,先理解初始化握手与 ClientInfo,再进入 MessageProcessor 读循环、Outgoing 写队列和断连清理。排查具体 API 时,沿着“transport → connection state → MessageProcessor → processor → Core → event mapper → outgoing queue”回溯,不要从单个 handler 直接推断整个 App Server 行为。
