Skip to content

InProcess AppServer

追踪嵌入式 App Server 的 channel transport、初始化握手、事件背压和关闭流程。

基于rust-v0.150.0
CodexRustAppServerInProcess

InProcess AppServer ​

本文承接AppServer启动与依赖装配,面向已经理解 MessageProcessor、connection state 和 outgoing queue 的读者。本文回答同进程嵌入模式如何复用完整协议、请求怎样通过 channel 路由、初始化为何在返回前完成,以及事件 lag 和 shutdown timeout 如何处理;不展开上层 app-server-client 封装。

1. 同进程边界 ​

In-process 模式没有 stdio/WebSocket framing,但仍保留 typed request、request ID、server request response、notification 和 connection capability。它用一个固定 connection ID 代替网络连接。

2. 启动参数 ​

源码位置:codex-rs/app-server/src/in_process.rs :: InProcessStartArgs

rust
pub struct InProcessStartArgs {
    pub arg0_paths: Arg0DispatchPaths,
    pub config: Arc<Config>,
    pub cli_overrides: Vec<(String, TomlValue)>,
    pub loader_overrides: LoaderOverrides,
    pub strict_config: bool,
    pub cloud_config_bundle: CloudConfigBundleLoader,
    pub thread_config_loader: Arc<dyn ThreadConfigLoader>,
    pub feedback: CodexFeedback,
    pub log_db: Option<LogDbLayer>,
    pub state_db: Option<StateDbHandle>,
    pub environment_manager: Arc<EnvironmentManager>,
    pub config_warnings: Vec<ConfigWarningNotification>,
    pub session_source: SessionSource,
    pub enable_codex_api_key_env: bool,
    pub initialize: InitializeParams,
    pub channel_capacity: usize,
}

调用方必须提前提供生产 transport 通常从进程环境装配的依赖:配置 loader、环境管理器、反馈/日志、state DB 和 session source 都由参数显式传入。channel_capacity 会被钳制为至少 1,它同时影响 client、processor、outgoing 和 event 队列。

3. 自动握手 ​

源码位置:codex-rs/app-server/src/in_process.rs :: start

rust
let initialize = args.initialize.clone();
let client = start_uninitialized(args).await?;
let initialize_response = client
    .request(ClientRequest::Initialize {
        request_id: RequestId::Integer(0),
        params: initialize,
    })
    .await?;
client.notify(ClientNotification::Initialized)?;

start() 返回前完成 initialize 和 initialized,因此调用者拿到的是 ready runtime。initialize 失败时会关闭 runtime 并返回 InvalidData,不会留下半初始化 handle。

4. 请求路由 ​

源码位置:codex-rs/app-server/src/in_process.rs :: InProcessClientSender::request

rust
let (response_tx, response_rx) = oneshot::channel();
self.try_send_client_message(InProcessClientMessage::Request {
    request: Box::new(request),
    response_tx,
})?;
response_rx.await.map_err(|err| {
    IoError::new(
        ErrorKind::BrokenPipe,
        format!("in-process request response channel closed: {err}"),
    )
})

每个 request 用 oneshot 等待 response;runtime 维护 HashMap<RequestId, Sender>。并发复用尚未完成的 ID 会返回 invalid request,避免 response 路由歧义。

5. Server事件 ​

InProcessServerEvent 区分 server request、notification 和 Lagged。Lagged 是 transport health marker,表示消费者落后且事件已经丢失,不是 Core 业务事件。

源码位置:codex-rs/app-server/src/in_process.rs :: InProcessServerEvent

rust
pub enum InProcessServerEvent {
    ServerRequest(Box<ServerRequest>),
    ServerNotification(Box<ServerNotification>),
    Lagged { skipped: usize },
}

server request 必须用当前 event stream 提供的 request ID 响应;任意 ID 不会改变 App Server 状态,还可能掩盖卡住的审批流程。

6. Backpressure ​

源码位置:codex-rs/app-server/src/in_process.rs :: try_send_client_message

rust
match self.client_tx.try_send(message) {
    Ok(()) => Ok(()),
    Err(TrySendError::Full(_)) => Err(IoError::new(
        ErrorKind::WouldBlock,
        "in-process app-server client queue is full",
    )),
    Err(TrySendError::Closed(_)) => Err(IoError::new(
        ErrorKind::BrokenPipe,
        "in-process app-server runtime is closed",
    )),
}

队列满是可重试 transport failure,队列关闭是生命周期结束。调用者不能把两者包装成业务 JSON-RPC error。

7. Runtime任务 ​

源码位置:codex-rs/app-server/src/in_process.rs :: start_uninitialized

rust
async fn start_uninitialized(args: InProcessStartArgs) -> IoResult<InProcessClientHandle> {
    args.config.auth_config().validate()?;
    let channel_capacity = args.channel_capacity.max(1);
    let installation_id = resolve_installation_id(&args.config.codex_home).await?;
    let auth_manager =
        AuthManager::shared_from_config(args.config.as_ref(), args.enable_codex_api_key_env)
            .await
            .map_err(IoError::other)?;
    let (client_tx, mut client_rx) = mpsc::channel::<InProcessClientMessage>(channel_capacity);
    let (event_tx, event_rx) = mpsc::channel::<InProcessServerEvent>(channel_capacity);

    let runtime_handle = tokio::spawn(async move {
        let (outgoing_tx, outgoing_rx) = mpsc::channel::<OutgoingEnvelope>(channel_capacity);
        let analytics_events_client =
            analytics_events_client_from_config(Arc::clone(&auth_manager), args.config.as_ref());
        let analytics_events_flush_client = analytics_events_client.clone();
        let outgoing_message_sender = Arc::new(OutgoingMessageSender::new(
            outgoing_tx,
            analytics_events_client.clone(),
        ));

这段实际启动路径先校验配置和认证,解析 installation id,再创建四类 bounded channel 的入口。analytics_events_flush_client 被保留到 runtime 退出时,用于关闭前 flush;因此它不是普通请求的临时对象。

8. 响应、通知与背压 ​

runtime 主循环对三种消息采用不同策略:response/error 必须按 request id 唤醒 oneshot;server request 队列满时把 overload error 回送给 Core;普通 notification 队列满时可以丢弃,但标记为可丢失的业务通知与终态通知采用不同发送策略。

源码位置:codex-rs/app-server/src/in_process.rs :: InProcessClientMessage::Request、OutgoingMessage::Response、OutgoingMessage::Request、OutgoingMessage::AppServerNotification

rust
match pending_request_responses.entry(request_id.clone()) {
    Entry::Vacant(entry) => {
        entry.insert(response_tx);
    }
    Entry::Occupied(_) => {
        let _ = response_tx.send(Err(invalid_request(format!(
            "duplicate request id: {request_id:?}"
        ))));
        continue;
    }
}

match processor_tx.try_send(ProcessorCommand::Request(Box::new(request))) {
    Ok(()) => {}
    Err(mpsc::error::TrySendError::Full(_)) => {
        if let Some(response_tx) = pending_request_responses.remove(&request_id) {
            let _ = response_tx.send(Err(JSONRPCErrorError {
                code: OVERLOADED_ERROR_CODE,
                message: "in-process app-server request queue is full".to_string(),
                data: None,
            }));
        }
    }
    Err(mpsc::error::TrySendError::Closed(_)) => {
        if let Some(response_tx) = pending_request_responses.remove(&request_id) {
            let _ = response_tx.send(Err(internal_error(
                "in-process app-server request processor is closed",
            )));
        }
        break;
    }
}

响应路由和业务请求路由是两张表:pending_request_responses 只保存客户端发出的 request id,server request 则通过 event stream 交给 embedder 决定。重复 request id 在进入 processor 前就失败,避免 response 错配。

9. 关闭流程 ​

源码位置:codex-rs/app-server/src/in_process.rs :: InProcessClientHandle::shutdown

rust
if self.client.client_tx.send(InProcessClientMessage::Shutdown { done_tx }).await.is_ok() {
    let _ = timeout(SHUTDOWN_ACK_TIMEOUT, done_rx).await;
}
if timeout(SHUTDOWN_TIMEOUT, &mut runtime_handle).await.is_err() {
    runtime_handle.abort();
    let _ = runtime_handle.await;
}

runtime 会取消 pending server requests、向未完成 client requests 返回 internal error、drain processor/background tasks,并 flush analytics。超时后强制 abort,不能假设所有后台任务都优雅结束。

10. 源码验证 ​

in-process tests 验证自动 initialize、thread/start、server request response、重复 ID、queue saturation 和 clean shutdown;专门测试验证 runtime 超时时 abort。它们证明 channel transport 和生命周期,不证明 stdio/WebSocket framing 行为。

源码位置:codex-rs/app-server/src/in_process.rs :: tests

text
cd codex-rs
cargo test -p codex-app-server in_process

11. 嵌入排查 ​

请求卡住时检查 request ID、pending map 和 server request 是否未回答;事件缺失时检查 Lagged 与 channel capacity;启动失败时检查 initialize response;关闭缓慢时区分 ACK timeout、processor drain 和 runtime abort。In-process 省去进程边界,但没有省去协议状态机。