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
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
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
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
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
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
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
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
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
cd codex-rs
cargo test -p codex-app-server in_process11. 嵌入排查
请求卡住时检查 request ID、pending map 和 server request 是否未回答;事件缺失时检查 Lagged 与 channel capacity;启动失败时检查 initialize response;关闭缓慢时区分 ACK timeout、processor drain 和 runtime abort。In-process 省去进程边界,但没有省去协议状态机。
