Skip to content

Connection核心状态

解释 App Server connection/session state、RPC gate、outgoing state 与断连回收边界。

基于rust-v0.150.0
CodexRustAppServerConnection

Connection核心状态 ​

本文承接AppServer架构总览,面向已经理解 transport 和 MessageProcessor 的读者。本文回答连接状态由哪些对象拥有、initialize 后哪些字段生效、RPC gate 如何阻止迟到请求,以及 outgoing state 和 cleanup 如何配合;不展开 initialize 参数的完整字段。

1. 状态分层 ​

App Server 把连接状态拆成 inbound session、outbound writer 和 cleanup task 三层。inbound 保存协议能力,outbound 保存写队列和过滤状态,cleanup 负责异步资源回收。

2. Session字段 ​

源码位置:codex-rs/app-server/src/message_processor.rs :: ConnectionSessionState、InitializedConnectionSessionState

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,
}

OnceLock 保证 initialize 状态只设置一次。调用方通过 accessor 读取未初始化时的默认值;experimental、MCP extension、request attestation 和 opt-out 都是连接 scope,不应写进全局 Config。mcp_event_streams 在 initialize 之前就存在,因为订阅任务的清理必须绑定连接本身,而不是绑定初始化 payload。

3. Inbound连接状态 ​

源码位置:codex-rs/app-server/src/transport.rs :: ConnectionState、ConnectionState::new

rust
pub(crate) struct ConnectionState {
    pub(crate) origin: ConnectionOrigin,
    pub(crate) outbound_initialized: Arc<AtomicBool>,
    pub(crate) outbound_experimental_api_enabled: Arc<AtomicBool>,
    pub(crate) outbound_opted_out_notification_methods: Arc<RwLock<HashSet<String>>>,
    pub(crate) session: Arc<ConnectionSessionState>,
}

impl ConnectionState {
    pub(crate) fn new(
        origin: ConnectionOrigin,
        outbound_initialized: Arc<AtomicBool>,
        outbound_experimental_api_enabled: Arc<AtomicBool>,
        outbound_opted_out_notification_methods: Arc<RwLock<HashSet<String>>>,
    ) -> Self {
        Self {
            origin,
            outbound_initialized,
            outbound_experimental_api_enabled,
            outbound_opted_out_notification_methods,
            session: Arc::new(ConnectionSessionState::new()),
        }
    }
}

inbound ConnectionState 保存 transport origin 和 session,同时持有 outbound 状态的共享镜像。初始化处理器写入 session 后,会把 capability 与 opt-out 同步到这些原子量/锁中,使 outbound router 不需要读取 MessageProcessor 内部对象。

4. Outbound字段 ​

源码位置:codex-rs/app-server/src/transport.rs :: OutboundConnectionState

rust
pub(crate) struct OutboundConnectionState {
    pub(crate) initialized: Arc<AtomicBool>,
    pub(crate) experimental_api_enabled: Arc<AtomicBool>,
    pub(crate) opted_out_notification_methods: Arc<RwLock<HashSet<String>>>,
    pub(crate) writer: mpsc::Sender<QueuedOutgoingMessage>,
    disconnect_sender: Option<CancellationToken>,
}

outbound state 使用原子变量和共享集合,让写路由无需持有 MessageProcessor 的 session lock;disconnect_sender=None 表示某些嵌入式连接不可由 transport 主动断开。

5. RPC Gate ​

源码位置:codex-rs/app-server/src/connection_rpc_gate.rs :: ConnectionRpcGate

rust
pub(crate) async fn run<F>(&self, future: F)
where
    F: Future<Output = ()>,
{
    let token = {
        let accepting = self.accepting.lock().await;
        if !*accepting { return; }
        self.tasks.token()
    };
    future.await;
    drop(token);
}

pub(crate) async fn close(&self) {
    let mut accepting = self.accepting.lock().await;
    *accepting = false;
    self.tasks.close();
}

close 只阻止尚未获取 token 的 handler;已经开始的 handler 继续执行。shutdown 再等待 TaskTracker,形成“拒绝迟到请求、等待在途请求”的边界。

6. Notification过滤 ​

源码位置:codex-rs/app-server/src/transport.rs :: should_skip_notification_for_connection

rust
if envelope.notification.experimental_reason().is_some()
    && !connection_state.experimental_api_enabled.load(Ordering::Acquire)
{
    return true;
}
let method = envelope.notification.to_string();
opted_out_notification_methods.contains(method.as_str())

过滤发生在写入 queue 前;Core 已产生 notification 不代表该 connection 一定观察到。实验 capability 和用户 opt-out 是两种独立过滤原因。

7. Cleanup任务 ​

源码位置:codex-rs/app-server/src/connection_cleanup.rs :: ConnectionCleanupTasks

rust
pub(crate) fn spawn(&mut self, future: impl Future<Output = ()> + Send + 'static) {
    self.tasks.spawn(future);
}

pub(crate) async fn reap_next(&mut self) {
    if self.tasks.is_empty() {
        pending::<()>().await;
    }
    if let Some(result) = self.tasks.join_next().await {
        log_cleanup_result(result);
    }
}

pub(crate) async fn drain(&mut self) {
    while let Some(result) = self.tasks.join_next().await {
        log_cleanup_result(result);
    }
}

pub(crate) fn abort(&mut self) {
    self.tasks.abort_all();
}

cleanup task 失败只记录 warning;reap_next 让主事件循环在有任务时逐个回收结果,drain 等待全部正常完成,abort 用于强制关闭。资源回收与 RPC gate 的在途 handler 等待是两个层次。

源码位置:codex-rs/app-server/src/lib.rs :: TransportEvent::ConnectionOpened、TransportEvent::ConnectionClosed

rust
TransportEvent::ConnectionOpened {
    connection_id,
    origin,
    writer,
    disconnect_sender,
} => {
    let outbound_initialized = Arc::new(AtomicBool::new(false));
    let outbound_experimental_api_enabled = Arc::new(AtomicBool::new(false));
    let outbound_opted_out_notification_methods =
        Arc::new(RwLock::new(HashSet::new()));
    outbound_control_tx
        .send(OutboundControlEvent::Opened {
            connection_id,
            writer,
            disconnect_sender,
            initialized: Arc::clone(&outbound_initialized),
            experimental_api_enabled: Arc::clone(&outbound_experimental_api_enabled),
            opted_out_notification_methods: Arc::clone(
                &outbound_opted_out_notification_methods,
            ),
        })
        .await?;
    connections.insert(
        connection_id,
        ConnectionState::new(
            origin,
            outbound_initialized,
            outbound_experimental_api_enabled,
            outbound_opted_out_notification_methods,
        ),
    );
}

连接打开时 inbound map 与 outbound router map 使用同一组 Arc 状态,但分别拥有 session 与 writer。关闭时主循环先移除 inbound state、关闭 RPC gate,再通知 outbound router 删除 writer,最后异步调用 MessageProcessor::connection_closed 清理 processor 资源。

8. 失败与竞态 ​

initialize 重复设置会失败;close 与新 RPC 竞态由 gate token 决定;writer queue 满会让消息无法排队;cleanup task panic 不会自动恢复业务状态。连接断开后,Core thread 是否继续由 Thread owner 决定。

9. 源码验证 ​

ConnectionRpcGate 测试验证 close 不轮询迟到 future、shutdown 等待已启动 handler、late run 被丢弃;transport 测试验证 experimental/opt-out notification 过滤;cleanup 测试验证 drain、abort 和 JoinError 日志。它们证明连接级状态边界,不证明业务 handler 自身可取消。

源码位置:

  • codex-rs/app-server/src/connection_rpc_gate.rs :: shutdown_drops_late_runs_while_waiting_for_inflight_work
  • codex-rs/app-server/src/transport.rs :: should_skip_notification_for_connection
  • codex-rs/app-server/src/connection_cleanup.rs :: ConnectionCleanupTasks
text
cd codex-rs
cargo test -p codex-app-server connection_rpc_gate
cargo test -p codex-app-server notification_filter
cargo test -p codex-app-server connection_cleanup

10. 状态排查 ​

遇到请求被静默丢弃,检查 gate 是否 closed;遇到 notification 不见,检查 experimental/opt-out filter 和 writer queue;遇到断连资源未回收,检查 outbound map、cleanup JoinSet 和 Core thread owner。连接状态排查必须区分 inbound、outbound 和 cleanup 三层。