雨天小六

读懂 Codex(2.14):Exec Server 与远程执行环境的协议边界

· 更新于 2026-08-02 · 专栏:读懂 Codex

#Codex#Agent Runtime#软件架构#Exec Server#Remote Environment

这里首先要消除一个命名误会:codex exec 是一次非交互 Agent 运行入口,codex exec-server 则是为 Codex 提供进程、文件系统和部分 HTTP 能力的执行环境服务。前者消费用户任务,后者不调用模型,只管理 “代码到底在哪台机器上执行”。

Exec Server 的设计目标不是把 std::process::Command 包成 RPC,而是让 Core 的 UnifiedExec 在本地和 远端共享进程生命周期、流式输出、沙箱、网络策略与恢复语义。

Environment 是 Core 看见的边界

EnvironmentManager 保存多个具名 Environment 和一个 default id。Environment 向 Core 暴露:

  • ExecBackend:启动和控制进程;
  • ExecutorFileSystem:远端路径的文件操作;
  • HttpClient:按环境路由的 HTTP;
  • capability root discovery;
  • readiness 与 connection state。

Local Environment 直接组合 LocalProcess、LocalFileSystem 与本地 HTTP;Remote Environment 组合 LazyRemoteExecServerClient、RemoteProcess 与 RemoteFileSystem。Core 工具层依赖 trait,不需要到处判断 WebSocket URL。

Core UnifiedExec 经过 EnvironmentManager 和 ExecBackend 选择本地后端或远端 Exec Server
图 2.14-1:远端执行是 Environment 后端选择,不是 shell tool 内临时拼一个 SSH 命令。
class Environment:
    exec_backend: ExecBackend
    filesystem: ExecutorFileSystem
    http_client: HttpClient
    remote_client: Optional[LazyExecServerClient]


async def open_process(turn_environment, request):
    env = turn_environment.selected()
    if env.is_remote:
        params = map_unified_request_to_exec_server(request)
        started = await env.exec_backend.start(params)
        return await UnifiedExecProcess.from_exec_server(started)
    return await spawn_local_with_native_paths(request)

配置、默认值与禁用语义

Environment 可来自 CODEX_HOME/environments.toml、旧的 CODEX_EXEC_SERVER_URL,或 Noise 环境变量。 Provider snapshot 决定:

  • environments:id → transport;
  • default:某个 environment id 或 Disabled;
  • include_local:是否额外注册保留 id local

CODEX_EXEC_SERVER_URL=none 会让 default environment 为空并省略 local。上层以 default_environment().is_some() 判断是否向模型暴露 shell/filesystem 工具,这不是一次普通连接失败, 而是能力明确关闭。

Manager 验证空 id、重复 id、保留名和不存在的 default。构建成功后会后台触发 remote connect,但 status() 只观察现状,不启动、不等待也不恢复连接,避免“查看健康状态”本身产生远程副作用。

初始连接可以 eager,也可以 deferred

普通远程 Environment 在注册后开始连接。Deferred Noise Environment 则先获得 DeferredEnvironmentRegistration,等环境所有者发布 ready info 并完成 one-shot registration 后才能 建立传输。

Ready info 可携带选中的 capability roots;实现限制最多 256 项,并验证 id 非空、同环境归属且不重复。 这样技能或插件根目录不会被错误地挂到另一执行环境。

def register_deferred_environment(manager, id, provider):
    identity = create_noise_identity()
    readiness_tx, readiness_rx = oneshot()
    env = Environment.remote(
        Deferred(readiness_rx, NoiseRendezvous(provider, identity))
    )
    manager.insert(id, env)
    return DeferredRegistration(readiness_tx, id, env.ready_info)

三步握手是协议门

ExecServerClient 建立 transport 后执行:

  1. initialize {clientName,resumeSessionId?}
  2. 等待 InitializeResponse {sessionId}
  3. 发送 initialized notification。

Server 的 ConnectionProcessor 顺序消费 inbound event,保证 initialize/initialized 不会被并发 handler 重排。Handler 用两个原子状态区分“initialize 已请求”和“initialized 已完成”;process/fs method 都调用 require_initialized_for

async def connect_transport(transport, resume_session_id=None):
    rpc = await RpcClient.connect(transport)
    response = await timeout(
        rpc.call("initialize", {
            "clientName": CLIENT_NAME,
            "resumeSessionId": resume_session_id,
        }),
        INITIALIZE_TIMEOUT,
    )
    verify_session_id_is_stable(response.sessionId)
    await rpc.notify("initialized", {})
    return rpc

同一连接第二次 initialize 会被拒绝;initialized 早于 initialize 也会导致协议错误。未知 request 返回 method-not-found,意外 notification、无法关联的 response/error 会关闭连接,而不是悄悄忽略状态机破坏。

Server 的每连接组件

ConnectionProcessor 拆出传输收发 task,再建立有界 outbound channel、RpcNotificationSender 和 ExecServerHandler。Inbound 主循环按序处理消息;单个 request route 可以异步执行,但与 transport disconnect 做 select

Handler 组合 SessionHandle、ProcessHandler、FileSystemHandler、route-aware HTTP 和后台任务 tracker。 连接结束时关闭 pending server requests、停止 handler background tasks、detach session,并中止 transport task。

这里的“detach”不等于立即 kill 所有进程。当前 SessionRegistry 保留 SessionEntry 供短时恢复。

Session resume 改变了断线语义

首次 initialize 未给 resume id 时,SessionRegistry 创建 UUID session id,并把通知 sender 附到当前 connection id。断开时:

  • 清除当前 attachment;
  • 从 ProcessHandler 移除 notification sender;
  • 记录 detached connection id;
  • 设置 30 秒 TTL。

新连接在 TTL 内用 resumeSessionId attach,可以复用同一 ProcessHandler 和仍在运行的进程;同一 session 若仍有活动连接,则拒绝二次 attach。TTL 到期后 registry 才 shutdown process;整个 exec-server shutdown 也会清空所有 sessions。

async def attach(registry, resume_id, connection, notifications):
    if resume_id is None:
        session = registry.create(uuid4(), ProcessHandler(notifications))
    else:
        session = registry.lookup(resume_id)
        require(not session.expired)
        require(not session.has_active_connection)
        session.process.replace_notification_sender(notifications)
    session.attach(connection.id)
    return session


async def detach(session):
    session.process.replace_notification_sender(None)
    session.mark_detached(deadline=now() + 30_seconds)
    spawn(expire_and_shutdown_if_still_detached(session))

仓库 README 仍写着 WebSocket 关闭会立即终止连接所属进程,这与当前 server/session_registry.rs 的 resume/TTL 实现不一致。研究运行机制时应以源码和对应测试为准,并把 文档差异记录下来,而不是为叙述整齐忽略它。

本地 WebSocket 与 Noise rendezvous

本地 WebSocket 每个 frame 承载一个 JSON-RPC message。远程 Noise 模式则在 rendezvous WebSocket 上发送 二进制 protobuf RelayMessageFrame,payload 是端到端加密记录。

Relay frame 的关键字段包括:

  • stream_id:一条虚拟 harness/environment JSON-RPC session;
  • body:handshake、data、ack_frame、resume、reset 或 heartbeat;
  • seqsegment_indexsegment_count:消息分段;
  • ackack_bits:最高连续确认与其后的位图;
  • next_seq:resume 协商;
  • payload/reason。

Harness 生成 UUIDv4 stream id;Environment 端按 stream id demux,每个 stream 启动独立 ConnectionProcessor。Rendezvous 只路由 frame,不解密 JSON-RPC,也不替端点保证可靠性。

def receive_relay_frame(frame):
    stream = sessions.get_or_create(frame.stream_id)
    stream.acknowledge(frame.ack, frame.ack_bits)

    if stream.seen(frame.seq):
        return send_redundant_ack()

    stream.store_segment(frame.seq, frame.segment_index, frame.payload)
    if stream.message_complete(frame.seq, frame.segment_count):
        ciphertext = stream.reassemble_message(frame.seq)
        jsonrpc_bytes = noise.decrypt(ciphertext)
        stream.connection_processor.accept(jsonrpc_bytes)

ACK 本身不再被 ACK,但 ack 与 ack_bits 冗余附在每个 outbound frame;重试、去重、重组都由端点负责。 这种设计让不可信或功能简单的 rendezvous 无需理解应用消息。

process/start 的原子登记

ExecParams 包含 caller-chosen process id、argv、PathUri cwd、env policy/overlay、tty、pipeStdin、arg0、 sandbox、managed network 与 network proxy 描述。

LocalProcess 启动前先校验:

  • network policy decision timeout 不能为 0;
  • 启用 callback 时 process id 必须非空且不超过协议上限;
  • argv 不能为空;
  • process id 在 Session 内不能已存在。

实现先在 map 中插入 Starting(ProcessStart token),然后异步 spawn。失败时只删除仍指向同一 token 的 占位;成功后也再次验证 token,才替换为 Running。这样并发 terminate 能取消 Starting,迟到的 spawn 不会把已取消进程“复活”。

async def process_start(params):
    validate(params)
    token = ProcessStart()
    async with process_map.lock():
        require(params.id not in process_map)
        process_map[params.id] = Starting(token)

    try:
        child = await sandbox_spawn(params)
    except Exception:
        async with process_map.lock():
            remove_only_if_same_start_token(params.id, token)
        raise

    async with process_map.lock():
        if process_map.get(params.id) != Starting(token):
            child.terminate()
            raise StartWasCancelled()
        process_map[params.id] = Running(child, next_seq=1)
    spawn(stdout_reader, stderr_reader, exit_watcher)

Sandbox retry 会复用 UnifiedExec 的整数 id,但 Core 在存在 exec-server sandbox 时追加 UUID,生成新的 executor process id,避免前一次被拒进程与重试冲突。

输出为什么同时支持 Push 与 Read

stdout、stderr 或 PTY bytes 被包装为 ProcessOutputChunk {seq,stream,chunk};Exited 和 Closed 也消耗 同一单调 seq。语义上:

  • Output:一段字节到达;
  • Exited:子进程退出码已知;
  • Closed:所有输出流关闭,不会再有输出;
  • Failed:客户端合成的 session/transport 失败,无 seq。

ExecProcessEventLog 保留有界 replay history,并向 live broadcast fan-out。订阅者先读 replay,再接 live;若 broadcast 返回 Lagged,调用者使用最后 seq 调 process/read 补洞。

Exec Server 进程输出通过序列事件推送,落后时用游标读取补齐,传输断开时用 session id 恢复
图 2.14-2:push 提供低延迟,cursor read 提供可恢复性;两者不是重复 API。

process/read(afterSeq,maxBytes,waitMs) 返回保留的较新 chunks、nextSeq、exited、exitCode、closed 和 sandboxDenied。waitMs 允许 long poll。maxBytes 至少会返回一个匹配 chunk,即使它本身超过预算,避免 一个大 chunk 永远无法前进。

async def follow_process(process):
    last_seq = 0
    events = process.subscribe_events()
    while True:
        try:
            event = await events.recv()
            consume_in_order(event)
            last_seq = event.seq or last_seq
        except BroadcastLagged:
            page = await process.read(after_seq=last_seq, wait_ms=0)
            for event in page.as_events():
                consume_in_order(event)
                last_seq = event.seq

Client 端还要处理不同接收 task 让 output/exited/closed 通知乱序到达的情况:先按 seq 缓冲,只有更低序号 均已交付才向 UnifiedExec 发布。

stdin 写入为什么需要 writeId

网络重试可能让同一 process/write 到达两次。若 payload 是 shell 命令,重复写入可能执行两遍。 所以每次写带非空 writeId;Server 保存最近接受的 id,相同 id 再来直接返回 Accepted。

实现先 reserve writer channel permit,再次检查 id,随后同步 send,并在任何下一次 await 前记录 id。 这个顺序保证 RPC task 在取消点重试时,不会出现“字节已写,但 id 尚未登记”的窗口。

非 PTY 进程只有 pipeStdin=true 才可写;不存在、Starting、stdin closed 等返回结构化 WriteStatus, 不全部变成 transport error。Terminate 返回 running: bool,使清理操作天然可重复。

UnifiedExec 怎样吸收远端差异

Core 将远端 StartedExecProcess 包装成 ProcessHandle::ExecServer,启动 output task 把 ExecProcessEvent 写入与本地进程相同的 OutputBuffer、broadcast channel 和 ProcessState watch。上层 shell tool 因此仍用 统一的 poll/write/terminate 逻辑。

但边界差异必须提前拒绝。Remote exec-server 不支持 inherited file descriptors;ProcessManager 在 调用 backend 之前检查,返回 create-process error。远端 cwd 保持 PathUri;本地分支才转 native absolute path。

Environment connection 还区分 Connected、Recovering 和 Failed。普通调用可按 RecoveryPolicy 等待恢复, 健康检查使用 fail-fast 现有连接,不能因为 status probe 偷偷启动 recovery。

失败矩阵

失败恢复或终态
initialize超时/错误 session id连接失败,不开放 RPC
protocol notificationinitialized 之外的意外方法关闭连接
transport短时断开Recovering + resumeSessionId
session30 秒内未重连shutdown 所有 managed process
relaysegment 丢失/重复ACK bits、retry、dedupe、reassembly
process startspawn 失败删除同 token Starting
process subscriberbroadcast Laggedprocess/read cursor 补齐
stdin writeresponse 丢失后重试writeId 去重
remote launchinherited FD 非空Core 边界拒绝
outputExited 已到但流未关等 Closed 再判定无后续字节

应锁定的测试

async def test_cancelled_start_cannot_resurrect(server):
    start = spawn(server.process_start(slow_params("p1")))
    await server.terminate("p1")
    await assert_raises(StartWasCancelled, start)
    assert "p1" not in server.running_processes


async def test_resume_preserves_process_within_ttl(server):
    first = await server.connect()
    session_id = await first.initialize()
    await first.start(long_running("p1"))
    await first.disconnect()
    second = await server.connect(resume_session_id=session_id)
    assert await second.read("p1").failure is None


async def test_duplicate_write_id_is_idempotent(process):
    await process.write(b"echo once\n", write_id="w1")
    await process.write(b"echo once\n", write_id="w1")
    assert (await process.output()).count(b"once") == 1


async def test_lagged_subscriber_recovers_by_cursor(process):
    last = await consume_until_lagged(process.events)
    page = await process.read(after_seq=last)
    assert contiguous(page.chunks, starts_at=last + 1)

Exec Server 把“远端机器”抽象成可恢复的执行环境,而不是隐藏成一次 RPC。Environment 选择、握手状态机、 Session resume、进程原子登记、序列化输出和幂等写入共同决定了远程 shell 是否能像本地 shell 一样可靠。

评论


← 返回文章列表