App Server 不是给 Core 套一层 JSON 序列化。它同时承担协议版本转换、输入校验、Thread 状态投影、订阅 路由、反向请求关联和异步完成语义。若把它实现成“收到 method 就 submit Op”,客户端很快会在 Turn id、 中断确认和审批生命周期上出错。
同一个入口服务两类客户端
网络客户端发送 JSONRPCRequest,MessageProcessor::process_request 先反序列化为 typed
ClientRequest。进程内 TUI 通过 process_client_request 跳过 JSON 文本解析,但两条路径都会进入
同一个 handle_client_request。
这保证 Embedded 与 WebSocket 客户端共享业务语义,而不是维护两套“几乎相同”的 handler。
每个请求还会被扩展为 ConnectionRequestId {connection_id, request_id}。JSON-RPC 只要求 id 在一条连接
内可关联;两个连接都可能发送 id=1。若服务端只用裸 request id,响应、pending interrupt 或审批回调就会
串到另一客户端。
async def process_request(connection_id, json_request):
key = ConnectionRequestId(connection_id, json_request.id)
context = RequestContext(
key=key,
trace_parent=parse_w3c_trace(json_request.trace),
)
try:
typed = deserialize_client_request(json_request)
await handle_client_request(context, typed)
except JsonRpcError as error:
await outgoing.send_error(key, error)
RequestContext 的 span 会存活到 response,Turn start 后还记录 request→turn 的映射,使 telemetry 能把一个 JSON-RPC 调用与后续异步 Turn 联系起来。
路由到领域 Processor
handle_client_request 是大的 typed dispatcher,但实际逻辑继续分到 TurnProcessor、ThreadProcessor、
ThreadLifecycleProcessor 等领域组件。它们共享进程级 ThreadManager、ThreadStore、ThreadStateManager 和
OutgoingMessageSender。
并非所有 method 都对应 Core Op:
| App Server method | Core 边界 |
|---|---|
turn/start | 构造并提交 Op::UserInput |
turn/steer | 直接调用 CodexThread::steer_input |
turn/interrupt | 提交 Op::Interrupt,延迟响应 |
thread/settings/update | 提交设置 Op,先确认排队 |
thread/read | 从 live state/store 读取,不创建 Op |
thread/resume | ThreadManager 恢复并注册 listener |
| 审批 response | 找到 pending server request,再唤醒 Core waiter |
App Server 的作用是选择正确的 Core 原语,而不是强迫所有操作走同一种队列。
turn/start 的校验顺序
turn_start_inner 先加载 live Thread,再检查当前 Thread 是否允许客户端直接输入。MultiAgentV2 的
ThreadSpawn 子 Agent 不接受 App Client 直接注入消息;它的输入所有者是父 Agent。
随后执行的关键校验与映射包括:
- 限制所有 Text input 的总字符数;
- 拒绝未经允许的远程图片 URL;
- 设置 app-server client info 和 form elicitation capability;
- 解析 runtime workspace roots、environment selections 与 cwd;
- 将 V2 UserInput 和 additional context 映射为 Core 类型;
- 构造并预览
ThreadSettingsOverrides; - 提交
Op::UserInput。
Permissions 与 legacy sandboxPolicy 不能同时指定。显式设置预览会应用 managed constraints,因此客户端 不能借一次 Turn override 绕开管理员策略。
async def turn_start(params, request):
thread = await load_live_thread(params.thread_id)
require_direct_input_owner(thread)
validate_total_text_chars(params.input)
validate_image_urls(params.input)
environments = resolve_environments(
params.cwd,
params.runtime_workspace_roots,
params.environments,
)
settings = await build_managed_settings_override(
permissions=params.permissions,
legacy_sandbox=params.sandbox_policy,
model=params.model,
environment=environments,
)
op = UserInputOp(
items=map_v2_input(params.input),
additional_context=map_context(params.additional_context),
thread_settings=settings,
output_schema=params.output_schema,
)
turn_id = await thread.submit_user_input_with_client_user_message_id(
op, request.trace, params.client_user_message_id
)
request.bind_turn(turn_id)
return Turn(id=turn_id, status="inProgress", items=[], items_view="notLoaded")
submission id 为什么就是 turn id
Core submit 为每个 Submission 生成 id;Op::UserInput 进入普通 Turn 路径时,这个 id 贯穿 Event.id
与 Turn lifecycle。App Server 直接把它作为公开 Turn id 返回,客户端无需等 TurnStarted 才知道标识。
但响应中的 status=InProgress 只表示请求已接受并拥有 id。此时 items 为空且
items_view=NotLoaded,不表示模型已经收到 prompt,更不表示任何工具已经执行。
这是第一种时序:排队成功后立即响应。
设置更新的确认也只是排队
thread/settings/update 的 response 同样不能被理解为设置已生效。Core 顺序处理 Submission,真正应用
后会发 thread/settings/updated notification。
如果下一操作依赖多项设置,客户端应把它们合在一个 TurnStart override 中,或等待 updated notification; 连续发送多个 update 再立即 start,会让协议正确性依赖队列时序和客户端猜测。
turn/steer 为什么不排 Op
Steering 要把输入追加到当前可 steering 的 ActiveTurn mailbox。它必须原子检查:
- 当前确有 ActiveTurn;
- expectedTurnId 与实际一致;
- Task kind 不是 Review/Compact;
- 输入非空。
因此 App Server 直接调用 CodexThread::steer_input。若先排成普通 Op,等队列处理时活动 Turn 可能早已
变化,expected id 的并发保护就失效。
错误被映射为明确类别:NoActiveTurn、ExpectedTurnMismatch、NonSteerableReview、
NonSteerableCompact 和 EmptyInput。不可 steering 错误还在 JSON-RPC error.data 中携带结构化
TurnError,TUI 不必解析英文 message 才知道能否排队。
turn/interrupt 为什么延迟响应
正常中断会先校验公开 ThreadState:
- active snapshot 存在时,id 必须匹配;
- 若 last terminal turn 就是目标,或 Agent 已不 Running,返回 no active turn;
- 校验通过,把完整 ConnectionRequestId 存入
pending_interrupts; - 提交
Op::Interrupt,暂不发送 response。
当 listener 收到 Core TurnAborted,bespoke handling 才取出 pending ids 并逐一回复
TurnInterruptResponse。因此客户端收到成功,确认的是“目标 Turn 已进入 abort 终态”,不是“服务端
看见了取消按钮”。
async def turn_interrupt(request, thread_id, turn_id):
if turn_id == "":
await thread.submit(InterruptOp())
return immediate_ok() # startup cancellation 没有 TurnAborted
state = await states.lock(thread_id)
state.require_active_turn(turn_id)
state.pending_interrupts.append(request.connection_request_id)
try:
await thread.submit(InterruptOp())
except Exception:
state.pending_interrupts.remove(request.connection_request_id)
raise
return DEFER_RESPONSE
async def on_turn_aborted(thread_id):
pending = states[thread_id].take_pending_interrupts()
for request_id in pending:
await outgoing.respond(request_id, TurnInterruptResponse())
空 turn id 是 startup interrupt,用来取消尚未形成 Turn 的启动过程,因此提交成功后立即响应。这是正常 Turn interrupt 的特例。
Event listener 先维护状态,再对外投影
每个 live Thread 都有 listener 消费 CodexThread::next_event。ThreadState 跟踪 active turn、items、
last terminal id、TurnSummary、pending interrupts 和 pending server requests。
Core EventMsg 不是 App Server notification 的 wire schema。Bespoke handling 会:
- TurnStarted:清理旧反向请求、更新 watcher,发送 versioned TurnStartedNotification;
- item start/delta/complete:映射为 ThreadItem 与增量 notification;
- token count:映射成 ThreadTokenUsageUpdated;
- TurnComplete/Aborted:计算状态、error、last agent message 和完成时间;
- raw response item:只有客户端声明 capability 时才额外发送。
这种投影让 Core 可以演化内部 Event,而 App Server 保持版本化外部协议,也避免把 Session 内部对象和 不安全路径直接暴露给远端。
Thread-scoped 与 global 消息
Thread 事件经 ThreadScopedOutgoingMessageSender 只发送给订阅该 Thread 的连接;账户状态、模型列表等
全局通知走不同边界。多客户端同时连接时,“知道 Thread id”不等于自动订阅其所有输出。
Thread resume 还要与 listener 串行化:恢复 response 若在历史与 live event 交界处乱序,客户端可能先看见 delta,再收到一个不含该 delta 的旧快照。Listener command 通道用于把这类操作插入事件序列。
审批为什么是 JSON-RPC request
当 Core 发出 command execution approval、file change approval、tool user input、MCP elicitation 或 permission request 时,App Server 必须等待回答。Notification 没有 response 语义,因此服务端向客户端 发起 JSON-RPC request。
Outgoing sender 为 request 分配 id并保留 oneshot callback;response 到达后,异步任务将结果映射为
Op::ExecApproval、Op::PatchApproval 等,Core 再按 call id 唤醒 waiter。
TurnStarted 与 TurnComplete 都会 abort pending server requests。双端清理看似重复,却防止上一个 Turn 遗留的请求在边界竞态中存活。
async def project_approval(core_event):
callback = await outgoing.send_server_request(
method="commandExecution/requestApproval",
params=map_approval(core_event),
scope=(core_event.thread_id, core_event.turn_id),
)
try:
decision = await callback
except TurnBoundaryCancelled:
return
await thread.submit(ExecApprovalOp(core_event.call_id, decision))
Dynamic tool call 等能力必须由客户端显式声明;不支持时 App Server 不能假装已处理,而要返回对应拒绝或 错误,避免 Core waiter 永远悬挂。
四种出站时序
| 类型 | 例子 | 完成点 |
|---|---|---|
| 立即 JSON-RPC response | turn/start、thread/read | 请求校验并接受/读取完成 |
| 延迟 JSON-RPC response | 正常 turn/interrupt | Core 发 TurnAborted |
| Server notification | item delta、TurnComplete | 无单独客户端 response |
| Server request | approval、request_user_input | 客户端必须回答或被 Turn 取消 |
这四类消息共用 JSON-RPC envelope,却不能用同一种超时、重试和 UI 状态处理。
失败与一致性矩阵
| 场景 | 拒绝位置 | 状态结果 |
|---|---|---|
| 文本总量超限 | TurnProcessor 前置校验 | 不提交 Op |
| 远程图片 URL 不允许 | 输入映射前 | 不污染历史 |
| permissions 与 sandboxPolicy 并存 | settings builder | 不应用部分设置 |
| ThreadSpawn 子 Agent 被直接输入 | ownership check | 父子协议不被绕开 |
| steer expected id 错 | Core 原子检查 | 返回 actual id |
| interrupt 提交失败 | pending queue 回滚 | 不留下永远等待的 response |
| Turn 结束时仍有审批 | listener abort | 客户端收到取消原因 |
| 未订阅 Thread | outgoing scope | 不发送 Thread notification |
应锁定的测试
async def test_request_ids_are_connection_scoped(server):
a = await server.connect()
b = await server.connect()
await parallel(a.request(id=1), b.request(id=1))
assert a.response(1).connection == a
assert b.response(1).connection == b
async def test_interrupt_response_waits_for_aborted(server, thread):
pending = spawn(server.turn_interrupt(thread.id, thread.turn_id))
assert not pending.done()
await thread.emit_turn_aborted()
assert (await pending).is_ok()
async def test_turn_boundary_cancels_approval(server, thread):
approval = await thread.emit_exec_approval()
await thread.emit_turn_complete()
assert approval.callback.cancelled
assert not thread.received_approval_op
App Server 的真正职责是维护“外部协议可观察状态”和“Core 内部运行状态”之间的一致性。JSON-RPC 只是 载体;连接级 id、领域校验、状态投影和异步完成点才是这一层的设计主体。
评论
登录后即可评论