Skip to content

AppServer架构总览

从连接、消息处理器、Core bridge、Thread state 到 outgoing queue,建立 App Server 源码地图。

基于rust-v0.150.0
CodexRustAppServerArchitecture

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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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_children
  • codex-rs/app-server/src/in_process.rs :: process_client_request
  • codex-rs/app-server/src/transport_tests.rs :: experimental_notifications_are_dropped_without_capability
  • codex-rs/app-server/src/transport_tests.rs :: broadcast_does_not_block_on_slow_connection
text
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=1

9. 阅读路径 ​

继续阅读时,先理解初始化握手与 ClientInfo,再进入 MessageProcessor 读循环、Outgoing 写队列和断连清理。排查具体 API 时,沿着“transport → connection state → MessageProcessor → processor → Core → event mapper → outgoing queue”回溯,不要从单个 handler 直接推断整个 App Server 行为。