Unified Exec写入stdin
write_stdin 不是简单的“向子进程写字符串”。它同时承担三种工作:向仍在运行的 TTY 会话写入字节;在没有输入时轮询已有输出;在非 TTY 会话中把控制字符转换成中断。写入结束后,它还要判断进程是否退出、是否仍在 store、是否需要发送 TerminalInteractionEvent。
本文承接Unified Exec会话状态机,面向理解 Rust async、PTY/pipe 和 Arc 锁的读者。范围从模型的 write_stdin 调用到 ExecCommandToolOutput 与交互事件;不展开输出截断算法和 exec-server 完整 RPC 协议。读完后,你应能解释空 chars 为什么不是 EOF、普通字符为什么不能写入非 TTY,以及一次写入失败如何被重新分类。
1. 输入入口
1.1 模型字段
公开工具使用 session_id,内部请求使用 process_id;handler 只负责解析和转发,不在这里判断 TTY 或进程状态。chars 默认空字符串,使同一个工具既能写入又能轮询。
源码位置:codex-rs/core/src/tools/handlers/unified_exec/write_stdin.rs :: WriteStdinArgs、WriteStdinHandler::handle_call
#[derive(Debug, Deserialize)]
struct WriteStdinArgs {
// The model is trained on `session_id`.
session_id: i32,
#[serde(default)]
chars: String,
#[serde(default = "super::default_write_stdin_yield_time_ms")]
yield_time_ms: u64,
#[serde(default)]
max_output_tokens: Option<usize>,
}
let args: WriteStdinArgs = parse_arguments(&arguments)?;
let response = session
.services
.unified_exec_manager
.write_stdin(WriteStdinRequest {
process_id: args.session_id,
input: &args.chars,
yield_time_ms: args.yield_time_ms,
max_output_tokens: args.max_output_tokens,
truncation_policy: turn.model_info.truncation_policy.into(),
interaction_event: Some(WriteStdinInteractionEvent {
session: &session,
turn: &turn,
}),
})
.await
.map_err(|err| {
FunctionCallError::RespondToModel(format!("write_stdin failed: {err}"))
})?;输入字符串借用到 manager 调用结束;handler 不保存它,也不把 session_id 转成 OS pid。解析错误直接返回模型错误,没有进程 I/O。
1.2 Hook边界
write_stdin 不发第二次 PreToolUse hook,因为原始 exec_command 已经以 Bash 工具名执行过;但完成轮询可以触发匹配的 PostToolUse。这个设计避免把一次交互误记成两次独立命令。
源码位置:codex-rs/core/src/tools/handlers/unified_exec/write_stdin.rs :: CoreToolRuntime
fn pre_tool_use_payload(&self, _invocation: &ToolInvocation) -> Option<PreToolUsePayload> {
// `write_stdin` is transport for an existing exec session.
None
}
fn post_tool_use_payload(
&self,
invocation: &ToolInvocation,
result: &dyn ToolOutput,
) -> Option<PostToolUsePayload> {
post_unified_exec_tool_use_payload(invocation, result)
}2. 同会话交互
2.1 先锁后借
manager 先从 store 克隆 Arc<UnifiedExecProcess>,再取得该进程的 interaction_lock,最后重新准备句柄。不同会话可以并发,同一会话的输出 drain、写入和终止必须串行。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: write_stdin
let locked_process = {
let store = self.process_store.lock().await;
let entry = store
.processes
.get(&process_id)
.ok_or(UnifiedExecError::UnknownProcessId { process_id })?;
Arc::clone(&entry.process)
};
let _interaction_guard = locked_process.interaction_lock().lock_owned().await;
let PreparedProcessHandles {
process,
output,
pause_state,
session,
network_approval,
call_id,
hook_command,
process_id,
tty,
..
} = self
.prepare_process_handles(process_id, &locked_process)
.await?;两次访问 store 之间条目可能被移除或替换,prepare_process_handles 用 Arc::ptr_eq 防止旧任务操作新进程。锁只保护交互阶段,不会把 store mutex 带进等待输出。
2.2 输入分支
空输入完全跳过写入,直接进入输出收集;TTY 写入原始字节;非 TTY 只接受内部 INTERRUPT 常量,否则返回 StdinClosed。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: write_stdin
if !request.input.is_empty() {
if !tty {
if request.input == INTERRUPT {
process.interrupt().await?;
} else {
return Err(UnifiedExecError::StdinClosed);
}
} else {
match process.write(request.input.as_bytes()).await {
Ok(()) => {
tokio::time::sleep(Duration::from_millis(100)).await;
}
Err(err) => {
let status = self.refresh_process_state(process_id).await;
if matches!(status, ProcessStatus::Exited { .. }) {
status_after_write = Some(status);
} else if matches!(err, UnifiedExecError::ProcessFailed { .. }) {
process.terminate();
self.release_process_id(process_id).await;
return Err(err);
} else {
return Err(err);
}
}
}
}
}100ms 不是命令超时,而是给远程进程一个短暂反应窗口,使随后的 poll 更可能收集到回显;真正的等待仍由 yield_time_ms 决定。
2.3 跨会话并行
WriteStdinHandler::supports_parallel_tool_calls 返回 true,但并行单位不是“同一个终端中的多个写操作”。工具路由可以同时分派不同 session,manager 再用每个 UnifiedExecProcess 自己的 interaction_lock 串行化同一 session。
源码位置:
codex-rs/core/src/tools/handlers/unified_exec/write_stdin.rs :: WriteStdinHandler::supports_parallel_tool_callscodex-rs/core/src/unified_exec/process.rs :: UnifiedExecProcess::interaction_lock
fn supports_parallel_tool_calls(&self) -> bool {
true
}
pub(super) fn interaction_lock(&self) -> Arc<Mutex<()>> {
Arc::clone(&self.interaction_lock)
}这形成两级并发控制:router 允许跨 session 并行,process lock 保证单 session 的 write、poll、terminate 与终态发布不重叠。
3. 等待与输出
3.1 两种等待上限
空 poll 使用至少 MIN_EMPTY_YIELD_TIME_MS,并受 manager 配置的后台上限约束;非空写入只受 MAX_YIELD_TIME_MS 限制,避免交互输入长时间阻塞。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: write_stdin
let yield_time_ms = {
let time_ms = request.yield_time_ms.max(MIN_YIELD_TIME_MS);
if request.input.is_empty() {
time_ms.clamp(MIN_EMPTY_YIELD_TIME_MS, self.max_write_stdin_yield_time_ms)
} else {
time_ms.min(MAX_YIELD_TIME_MS)
}
};
let start = Instant::now();
let deadline = start + Duration::from_millis(yield_time_ms);
let collected_output =
Self::collect_output_until_deadline(&output, pause_state, deadline).await;这一区分使后台命令可以用较长空 poll 等待,而交互输入保持响应。暂停状态还可以延长 deadline,避免用户界面暂停期间误判超时。
3.2 恢复写入标识
exec-server client 为每次逻辑写入分配递增 write_id。如果 RPC transport 在响应前关闭,Session::write 等待恢复并重试,但复用相同的 chunk 和 ID;只有下一次独立写入才分配新 ID。
源码位置:codex-rs/exec-server/src/client.rs :: SessionState::next_write_id、Session::write
fn next_write_id(&self) -> String {
self.next_write_id
.fetch_add(1, Ordering::Relaxed)
.to_string()
}
pub(crate) async fn write(&self, chunk: Vec<u8>) -> Result<WriteResponse, ExecServerError> {
let write_id = self.state.next_write_id();
loop {
match self
.client
.write(&self.process_id, chunk.clone(), write_id.clone())
.await
{
Ok(response) => return Ok(response),
Err(error)
if is_transport_closed_error(&error) && !self.client.inner.is_failed() =>
{
continue;
}
Err(error) => return Err(error),
}
}
}没有 write_id 时,连接断开发生在“服务端已写入、客户端未收到响应”这一窗口,重试就可能把同一批字节写入两次。ID 把一次逻辑写入从某个具体 RPC request 中分离出来。
3.3 服务端去重
协议要求 WriteParams 同时携带 process ID、字节块和非空 write ID。每个运行中进程维护一个有界 AcceptedStdinWriteIds:已经接受的 ID 直接返回 Accepted;新 ID 在取得 writer permit 后同步发送,并在任何后续 await 之前写入缓存。
源码位置:
codex-rs/exec-server-protocol/src/protocol.rs :: WriteParams、WriteStatuscodex-rs/exec-server/src/local_process.rs :: LocalProcess::exec_write
if accepted_stdin_write_ids
.lock()
.await
.contains(¶ms.write_id)
{
return Ok(WriteResponse {
status: WriteStatus::Accepted,
});
}
let permit = writer_tx
.reserve()
.await
.map_err(|_| internal_error("failed to write to process stdin".to_string()))?;
let mut accepted_stdin_write_ids = accepted_stdin_write_ids.lock().await;
if accepted_stdin_write_ids.contains(¶ms.write_id) {
return Ok(WriteResponse {
status: WriteStatus::Accepted,
});
}
permit.send(params.chunk.into_inner());
accepted_stdin_write_ids.remember(params.write_id);发送后立即记忆 ID 是取消安全边界:如果 handler 在下一次 await 前被取消,恢复后的同 ID 请求仍只会得到确认,不会再次写入。缓存有界,因此它保证近期恢复窗口内的幂等,而不是为进程保存无限历史。
3.4 状态回到Core
远程 ExecProcess::write 返回 WriteStatus。Accepted 继续 poll;UnknownProcess 和 StdinClosed 把统一快照标记为 exited;Starting 也不能当作成功写入。
源码位置:codex-rs/core/src/unified_exec/process.rs :: UnifiedExecProcess::write
match process_handle.write(data.to_vec()).await {
Ok(response) => match response.status {
WriteStatus::Accepted => Ok(()),
WriteStatus::UnknownProcess | WriteStatus::StdinClosed => {
let state = self.state_rx.borrow().clone();
let _ = self.state_tx.send_replace(state.exited(state.exit_code));
self.output.cancellation_token.cancel();
Err(UnifiedExecError::WriteToStdin)
}
WriteStatus::Starting => Err(UnifiedExecError::WriteToStdin),
},
Err(err) => Err(UnifiedExecError::process_failed(err.to_string())),
}这解释了为什么 manager 在写入错误后还要调用 refresh_process_state:transport 错误可能与远程进程刚刚退出同时发生。write_id 解决的是“是否重复写入”,WriteStatus 解决的是“进程当前能否接受写入”,两者不能互相替代。
4. 终态回灌
4.1 Alive与Exited
轮询完成后,refresh_process_state 决定输出是否保留 process_id。Alive 返回 ID 以便下一次 write_stdin;Exited 移除条目、结束网络审批,并把 process_id 置为 None。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: write_stdin
let (process_id, exit_code, event_call_id) = match status {
ProcessStatus::Alive {
exit_code,
call_id,
process_id,
} => (Some(process_id), exit_code, call_id),
ProcessStatus::Exited { exit_code, entry } => {
let call_id = entry.call_id.clone();
finish_network_approval_after_process_exit_for_entry(&entry).await?;
(None, exit_code, call_id)
}
ProcessStatus::Unknown => {
if process.has_exited() {
(None, process.exit_code(), call_id)
} else {
return Err(UnifiedExecError::UnknownProcessId {
process_id: request.process_id,
});
}
}
};Unknown + 当前 Arc 已退出 是并发 fallback,不是第三种持久状态。它允许一个已经取得句柄的 poll 在另一个任务完成 map 清理后仍返回正确 exit metadata。
4.2 交互事件
只要本次有输入,或进程仍然存活,manager 就发送 TerminalInteractionEvent。纯空 poll 且进程已退出时不会伪造一次交互。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: write_stdin
let should_emit_interaction = !request.input.is_empty() || response.process_id.is_some();
if should_emit_interaction
&& let Some(WriteStdinInteractionEvent { session, turn }) = request.interaction_event
{
let interaction = TerminalInteractionEvent {
call_id: response.event_call_id.clone(),
process_id: response
.process_id
.unwrap_or(request.process_id)
.to_string(),
stdin: request.input.to_string(),
};
session
.send_event(turn.as_ref(), EventMsg::TerminalInteraction(interaction))
.await;
}process_id 在进程已退出时使用请求 ID 作为事件关联键,响应本身仍以 None 表示不可继续交互。事件的消费者是 session event stream,不是 OS stdin。
5. 失败与清理
网络审批取消或进程 failure message 出现时,manager 先收集已有输出,再结束审批、释放 ID,并把失败消息转成工具错误。失败清理和普通 exit 清理不是同一条分支,因为 failure 需要保留根因。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: write_stdin
if network_approval
.as_ref()
.is_some_and(DeferredNetworkApproval::is_cancelled)
{
let message =
network_denial_message_for_session(session.as_ref(), network_approval.clone()).await;
self.release_process_id(process_id).await;
return Err(fail_process_with_message(process.as_ref(), message));
}
if let Some(message) = process.failure_message() {
let finish_result = finish_deferred_network_approval_for_session(
session.as_ref(),
network_approval.clone(),
)
.await;
self.release_process_id(process_id).await;
if let Err(message) = finish_result {
return Err(fail_process_with_message(process.as_ref(), message));
}
return Err(UnifiedExecError::process_failed(message));
}清理顺序保证 store 不再接受后续写入,同时网络审批不会悬挂;底层进程若仍存在,则 fail_process_with_message 负责终止。
6. 四组验证
6.1 交互与退出
unified_exec_emits_terminal_interaction_for_write_stdin 启动 TTY shell,写入一条 echo 命令,断言事件中的 call ID、process ID 和 stdin 与请求一致。write_stdin_returns_exit_metadata_and_clears_session 再用 cat 验证运行中写入保留同一 process ID,发送 EOF 后响应移除 ID、返回 exit code 0,且 TerminalInteraction 先于 ExecCommandEnd。
源码位置:
codex-rs/core/tests/suite/unified_exec.rs :: unified_exec_emits_terminal_interaction_for_write_stdincodex-rs/core/tests/suite/unified_exec.rs :: write_stdin_returns_exit_metadata_and_clears_session
cd codex-rs
RUST_MIN_STACK=8388608 cargo test -p codex-core --test all unified_exec_emits_terminal_interaction_for_write_stdin -- --test-threads=1
RUST_MIN_STACK=8388608 cargo test -p codex-core --test all write_stdin_returns_exit_metadata_and_clears_session -- --test-threads=1这些测试证明 handler、manager、事件和模型工具输出的关联顺序,不证明终端程序对任意输入的业务解释。
6.2 Ctrl-C与非TTY
两个集成测试分别让非 TTY shell 捕获 SIGINT 并返回 42,以及使用默认中断语义返回 130;两者都断言响应移除 process ID 并保留真实 exit code。remote_write_closed_stdin_marks_process_exited 则构造 exec-server StdinClosed 响应,断言 Core 将统一进程标记为 exited。
源码位置:
codex-rs/core/tests/suite/unified_exec.rs :: write_stdin_ctrl_c_interrupts_non_tty_sessioncodex-rs/core/tests/suite/unified_exec.rs :: write_stdin_ctrl_c_default_interrupt_reports_130_for_non_tty_sessioncodex-rs/core/src/unified_exec/process_tests.rs :: remote_write_closed_stdin_marks_process_exited
cd codex-rs
RUST_MIN_STACK=8388608 cargo test -p codex-core --test all write_stdin_ctrl_c_interrupts_non_tty_session -- --test-threads=1
RUST_MIN_STACK=8388608 cargo test -p codex-core --test all write_stdin_ctrl_c_default_interrupt_reports_130_for_non_tty_session -- --test-threads=1
cargo test -p codex-core --lib remote_write_closed_stdin_marks_process_exited -- --test-threads=1Unix 测试证明当前平台的 SIGINT 与 shell trap 语义;Windows 使用单独的终止测试,不能把退出码 42 或 130 外推到所有平台。
6.3 跨会话并行
write_stdin_calls_run_in_parallel_across_sessions 同时创建两个后台 session,并在同一模型响应中发出两个 write_stdin。断言两个调用能并行完成,证明 handler 的 parallel 标记与 per-process lock 没有退化成全局串行。
源码位置:codex-rs/core/tests/suite/unified_exec.rs :: write_stdin_calls_run_in_parallel_across_sessions
cd codex-rs
RUST_MIN_STACK=8388608 cargo test -p codex-core --test all 'suite::unified_exec::write_stdin_calls_run_in_parallel_across_sessions' -- --exact --test-threads=16.4 重试幂等
client 测试在第一次 process/write 发出后关闭 WebSocket,恢复 session 后断言重试请求保留完全相同的 write ID。server 测试依次发送 write-1、重复的 write-1 和 write-2,子进程最终只读到 first 与 second 两行,证明重复 ID 只确认一次。
源码位置:
codex-rs/exec-server/src/client.rs :: session_write_retries_same_write_id_after_recoverycodex-rs/exec-server/tests/process.rs :: exec_server_dedupes_retried_process_write_ids
cd codex-rs
cargo test -p codex-exec-server --lib session_write_retries_same_write_id_after_recovery -- --test-threads=1
cargo test -p codex-exec-server --test process exec_server_dedupes_retried_process_write_ids -- --test-threads=17. 读源码方法
遇到写入无效时,先看 WriteStdinArgs.chars 是否为空,再看 ProcessEntry.tty,然后区分 process.write 的 transport 状态和 manager 的 refresh_process_state 结果。遇到事件重复时,检查 pre_tool_use_payload 是否仍为 None、should_emit_interaction 是否被满足,以及原始 exec_command 的 PostToolUse 是否被复用。
rg -n "struct WriteStdinArgs|write_stdin\(|INTERRUPT|StdinClosed" codex-rs/core/src
rg -n "TerminalInteractionEvent|should_emit_interaction|WriteStatus" codex-rs/core/src
rg -n "emits_terminal_interaction|ctrl_c_interrupts_non_tty|remote_write_closed" codex-rs/core本篇的核心结论是:空输入是 poll,TTY 输入是 bytes,非 TTY 只允许中断;写入后必须重新观察状态,响应中的 process_id 决定模型能否继续交互,而事件中的 process ID 只承担关联作用。
