流式推理与消息Delta
本文承接 Responses流解析 和 Codex API类型。 上一篇已经说明 SSE frame 如何变成 ResponseEvent;本文继续追踪这些事件进入 run_sampling_request 之后发生的事情:哪个 output_item 拥有 delta,普通文本为什么要经过二次解析, reasoning summary 为什么有两种交付协议,工具参数增量又为什么不能直接广播给客户端。
本文不重新讲 SSE 解帧、请求构造、WebSocket transport,也不把最终 ResponseItem 的完整持久化流程当作 重点。读者需要知道 Rust 异步 Stream、ResponseEvent 和 Turn 的基本层级;读完后应能从 response.output_item.added 和某个 delta 事件定位到活动 item、对应的 EventMsg 消费者,以及 item 结束时 执行的 flush、完成或丢弃分支。
1. 两级增量
Responses 流有两个不同含义的“增量”:API 层的 ResponseEvent 是服务端 wire 事件的归一化结果,Core 层的 EventMsg 是面向 UI、工具 runtime、遥测和历史的运行时事件。二者名字相似,但 owner、字段和生效时机都不同。
| 层级 | 代表对象 | 主要所有者 | 何时产生 | 主要消费者 |
|---|---|---|---|---|
| wire 事件 | response.output_text.delta | provider | 服务端产生一段字段变化 | API parser |
| API 事件 | OutputTextDelta(String) | codex-api stream | JSON 事件字段满足最低条件 | Core Turn |
| 运行事件 | AgentMessageContentDelta | Session | active item 和 streaming policy 都允许 | TUI、App Server、SDK |
| 完整 item | ResponseItem::Message、Reasoning、tool call | Turn | response.output_item.done 到达 | 历史、工具调度、下一步 sampling |
关键不变量是:delta 不能独立决定自己的 item。Core 先由 OutputItemAdded 建立活动 item,再把后续 delta 投影到这个 item;如果活动 item 消失,普通文本和 reasoning delta 会进入保护分支,而不是凭最近一个 id 猜测。
2. 活动归属
try_run_sampling_request 在一次 sampling 内维护 active_item、是否对客户端流式发送,以及工具参数 diff consumer。OutputItemAdded 是建立归属的入口,OutputItemDone 是切换归属和完成前一个 item 的屏障。
当前实现对没有 active_item 的文本、reasoning 或 tool delta 不会凭最近 ID 猜测归属:文本和 reasoning 会 进入 error_or_panic 保护分支,tool delta 只有存在对应 diff consumer 才消费。生产构建会记录错误并继续, debug 构建则立即暴露 provider 事件顺序违反协议的问题。
源码位置:codex-rs/core/src/session/turn.rs :: try_run_sampling_request 的活动状态
let mut active_item: Option<TurnItem> = None;
let mut active_tool_argument_diff_consumer: Option<(
String,
Box<dyn ToolArgumentDiffConsumer>,
)> = None;
let mut assistant_message_stream_parsers = AssistantMessageStreamParsers::new(plan_mode);
let mut active_item_is_streaming_to_client = false;
// 每次收到完整 output item,都先结束前一个活动项,再处理新的 item。
ResponseEvent::OutputItemDone(mut item) => {
assign_missing_streamed_response_item_id(&mut item, active_item.as_ref());
if let Some((_, mut consumer)) = active_tool_argument_diff_consumer.take()
&& let Ok(Some(event)) = consumer.finish()
{
sess.send_event(&turn_context, event).await;
}
let previously_active_item = active_item.take();
let previously_streamed_item = if active_item_is_streaming_to_client {
previously_active_item
} else {
None
};
active_item_is_streaming_to_client = false;
if let Some(previous) = previously_streamed_item.as_ref()
&& matches!(previous, TurnItem::AgentMessage(_))
{
let item_id = previous.id();
flush_assistant_text_segments_for_item(
&sess,
&turn_context,
plan_mode_state.as_mut(),
&mut assistant_message_stream_parsers,
&item_id,
)
.await;
}
// 随后才把当前完整 item 交给工具/历史处理。
}同一个 ResponseItem 可能没有 server id。assign_missing_streamed_response_item_id 先尝试从前一个活动 item 取得关联,再调用 Session 的全局补 id 逻辑。这个动作不是 UI 装饰:item id 会进入后续 delta、历史记录和 工具调用关联,缺失时必须在首次消费前补齐。
OutputItemAdded 则建立新的可流式 item,并对已有文本做 seed。原因是服务端可能在 added 事件中已经带有 一段 output text,后续 delta 又从同一个 item 继续;如果直接从空 parser 开始,跨边界的隐藏 tag 会被拆坏。
源码位置:codex-rs/core/src/session/turn.rs :: OutputItemAdded 分支
ResponseEvent::OutputItemAdded(mut item) => {
assign_missing_streamed_response_item_id(&mut item, /*active_item*/ None);
if let ResponseItem::CustomToolCall {
call_id,
name,
namespace,
..
} = &item
{
let tool_name = ToolName::new(namespace.clone(), name.as_str());
active_tool_argument_diff_consumer = tool_runtime
.create_diff_consumer(&tool_name)
.map(|consumer| (call_id.clone(), consumer));
} else if matches!(&item, ResponseItem::FunctionCall { .. }) {
active_tool_argument_diff_consumer = None;
}
if let Some(turn_item) = handle_non_tool_response_item(
sess.as_ref(),
TurnItemContributorPolicy::Skip,
&item,
plan_mode,
)
.await
{
let mut turn_item = turn_item;
let stream_item_to_client = !defer_streamed_turn_items_for_contributors;
if stream_item_to_client
&& matches!(turn_item, TurnItem::AgentMessage(_))
&& let Some(raw_text) = raw_assistant_output_text_from_item(&item)
{
let item_id = turn_item.id();
let mut seeded = assistant_message_stream_parsers
.seed_item_text(&item_id, &raw_text);
if let TurnItem::AgentMessage(agent_message) = &mut turn_item {
agent_message.content = vec![
codex_protocol::items::AgentMessageContent::Text {
text: std::mem::take(&mut seeded.visible_text),
},
];
}
}
if stream_item_to_client {
sess.emit_turn_item_started(&turn_context, &turn_item).await;
}
active_item = Some(turn_item);
active_item_is_streaming_to_client = stream_item_to_client;
}
}这里有两个容易误读的边界:
handle_non_tool_response_item返回None时不会建立active_item,所以并非每个output_item.added都会 对外开始一个 Turn item。- 有 extension turn-item contributor 时,
defer_streamed_turn_items_for_contributors为 true,客户端暂时 不收到 started/delta;完整 item 会在 contributor 处理后再交付。不能把“没有 UI delta”解释成“provider 没有 发送 delta”。
3. 文本解析
普通 assistant 文本不直接从 OutputTextDelta 透传。AssistantMessageStreamParsers 为每个 item id 持有 一个 AssistantTextStreamParser,其输入依次经过 citation 隐藏标签处理,以及 plan mode 下的 <proposed_plan> 分段处理。
源码位置:codex-rs/utils/stream-parser/src/assistant_text.rs :: AssistantTextChunk、AssistantTextStreamParser
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct AssistantTextChunk {
pub visible_text: String,
pub citations: Vec<String>,
pub plan_segments: Vec<ProposedPlanSegment>,
}
#[derive(Debug, Default)]
pub struct AssistantTextStreamParser {
plan_mode: bool,
citations: CitationStreamParser,
plan: ProposedPlanParser,
}
impl AssistantTextStreamParser {
pub fn push_str(&mut self, chunk: &str) -> AssistantTextChunk {
let citation_chunk = self.citations.push_str(chunk);
let mut out = self.parse_visible_text(citation_chunk.visible_text);
out.citations = citation_chunk.extracted;
out
}
pub fn finish(&mut self) -> AssistantTextChunk {
let citation_chunk = self.citations.finish();
let mut out = self.parse_visible_text(citation_chunk.visible_text);
if self.plan_mode {
let mut tail = self.plan.finish();
if !tail.is_empty() {
out.visible_text.push_str(&tail.visible_text);
out.plan_segments.append(&mut tail.extracted);
}
}
out.citations = citation_chunk.extracted;
out
}
}finish() 是必要的,不是普通的 push_str("")。隐藏 tag 的开头可能跨两个网络 delta,plan block 的结尾也 可能只在 item done 或 response completed 时才知道。parser 直到 finish_item 才释放该 item 的缓存,随后 emit_streamed_assistant_text_delta 才把尾部可见文本发送出去。
图中的隐藏状态属于每个 item 的 parser,不属于全局 Session;因此两个并发或连续 item 的 tag 前缀不会共享。
3.1 普通模式
普通模式只把 visible_text 转成 AgentMessageContentDelta。citation payload 被提取但当前只在本地消费, 不会作为独立 protocol event 暴露;因此读者在 UI 看见的文本可能比原始 response item 短,这不是数据丢失, 而是公开文本投影主动移除了隐藏标记。
源码位置:codex-rs/core/src/session/turn.rs :: emit_streamed_assistant_text_delta
async fn emit_streamed_assistant_text_delta(
sess: &Session,
turn_context: &TurnContext,
plan_mode_state: Option<&mut PlanModeStreamState>,
item_id: &str,
parsed: ParsedAssistantTextDelta,
) {
if parsed.is_empty() {
return;
}
if !parsed.citations.is_empty() {
// citation extraction remains local; visible text omits the hidden markup.
let _citations = parsed.citations;
}
if let Some(state) = plan_mode_state {
if !parsed.plan_segments.is_empty() {
handle_plan_segments(sess, turn_context, state, item_id, parsed.plan_segments).await;
}
return;
}
if parsed.visible_text.is_empty() {
return;
}
let event = AgentMessageContentDeltaEvent {
thread_id: sess.thread_id.to_string(),
turn_id: turn_context.sub_id.clone(),
item_id: item_id.to_string(),
delta: parsed.visible_text,
};
sess.send_event(turn_context, EventMsg::AgentMessageContentDelta(event))
.await;
}3.2 计划模式
计划模式不只是“换一个 UI 标签”。AssistantTextStreamParser 将普通文本与 proposed plan 分成不同 ProposedPlanSegment;Session 再分别发送 AgentMessageContentDelta、PlanDelta 和计划 item 生命周期事件。 计划标记跨 chunk 时,parser 必须等待足够前缀才能决定它是普通文本还是计划开始。
源码位置:codex-rs/core/src/session/turn.rs :: handle_plan_segments
for segment in segments {
match segment {
ProposedPlanSegment::Normal(delta) => {
if delta.is_empty() {
continue;
}
maybe_emit_pending_agent_message_start(sess, turn_context, state, item_id).await;
let event = AgentMessageContentDeltaEvent {
thread_id: sess.thread_id.to_string(),
turn_id: turn_context.sub_id.clone(),
item_id: item_id.to_string(),
delta,
};
sess.send_event(turn_context, EventMsg::AgentMessageContentDelta(event))
.await;
}
ProposedPlanSegment::ProposedPlanStart => {
if !state.plan_item_state.completed {
state.plan_item_state.start(sess, turn_context).await;
}
}
ProposedPlanSegment::ProposedPlanDelta(delta) => {
if !state.plan_item_state.completed {
if !state.plan_item_state.started {
state.plan_item_state.start(sess, turn_context).await;
}
state.plan_item_state.push_delta(sess, turn_context, &delta).await;
}
}
ProposedPlanSegment::ProposedPlanEnd => {}
}
}普通文本和计划文本共享一个 provider delta 顺序,却进入不同的消费者。计划模式还会延迟只包含计划内容的 assistant item start;否则客户端会先看到一个空的 assistant message,随后才看到真正的 plan。
4. 文本收口
OutputItemDone 到达时,Core 先 flush 前一个 assistant item 的 parser;ResponseEvent::Completed 到达时, 再调用 flush_assistant_text_segments_all 处理尚未收到 done 的尾部 parser。两者分别对应“item 切换”和“整个 sampling 收束”,不能只在 completed 时 flush。
源码位置:codex-rs/core/src/session/turn.rs :: flush_assistant_text_segments_for_item、flush_assistant_text_segments_all
async fn flush_assistant_text_segments_for_item(
sess: &Session,
turn_context: &TurnContext,
plan_mode_state: Option<&mut PlanModeStreamState>,
parsers: &mut AssistantMessageStreamParsers,
item_id: &str,
) {
let parsed = parsers.finish_item(item_id);
emit_streamed_assistant_text_delta(sess, turn_context, plan_mode_state, item_id, parsed).await;
}
async fn flush_assistant_text_segments_all(
sess: &Session,
turn_context: &TurnContext,
mut plan_mode_state: Option<&mut PlanModeStreamState>,
parsers: &mut AssistantMessageStreamParsers,
) {
for (item_id, parsed) in parsers.drain_finished() {
emit_streamed_assistant_text_delta(
sess,
turn_context,
plan_mode_state.as_deref_mut(),
&item_id,
parsed,
)
.await;
}
}如果只收到 OutputItemAdded、若干 text delta 和 response.completed,没有 OutputItemDone,最后一次 drain_finished() 仍会把 parser 缓冲尾部交付。反过来,如果 item 已经 done,finish_item 会从 map 删除 parser,completed 阶段不会重复发送同一尾部。
5. 文本分派
Turn 收到 OutputTextDelta 后,必须同时满足三个条件才对外发送:存在 active item、该 item 允许流向客户端、 且 item 类型是 assistant message 或可承载 agent message delta 的类型。review child thread 会保留 provider 事件,但抑制客户端文本 delta,由最终 ReviewOutput 负责展示。
源码位置:codex-rs/core/src/session/turn.rs :: ResponseEvent::OutputTextDelta
ResponseEvent::OutputTextDelta(delta) => {
if let Some(active) = active_item.as_ref() {
if !active_item_is_streaming_to_client {
continue;
}
let item_id = active.id();
if matches!(active, TurnItem::AgentMessage(_)) {
let parsed = assistant_message_stream_parsers.parse_delta(&item_id, &delta);
emit_streamed_assistant_text_delta(
&sess,
&turn_context,
plan_mode_state.as_mut(),
&item_id,
parsed,
)
.await;
} else {
let event = AgentMessageContentDeltaEvent {
thread_id: sess.thread_id.to_string(),
turn_id: turn_context.sub_id.clone(),
item_id,
delta,
};
sess.send_event(&turn_context, EventMsg::AgentMessageContentDelta(event))
.await;
}
} else {
error_or_panic("OutputTextDelta without active item".to_string());
}
}“没有 active item”不是普通的空事件:它表示服务端事件顺序违反了 Core 的归属前提,因此代码选择 error_or_panic,而不是把 delta 贴到未知 item。extension contributor 的延迟路径则不同,它有明确的 active_item_is_streaming_to_client=false,因此是策略性抑制,不是顺序错误。
这张图把“顺序错误”和“策略性抑制”分开:没有活动 item 是归属不成立,streaming flag 为 false 则是已知的 consumer policy。
6. 工具参数
工具调用的参数 delta 不直接成为 AgentMessageContentDelta。OutputItemAdded 根据 custom tool 的 namespace/name 创建 ToolArgumentDiffConsumer,后续 ToolCallInputDelta 必须匹配当前 call id;匹配成功后 才生成工具专属事件。provider 可能省略 call id,Core 此时使用当前 active call id;若显式 id 不匹配,则丢弃 该 delta,避免把两个并发或连续调用的参数拼到一起。
源码位置:codex-rs/core/src/session/turn.rs :: ResponseEvent::ToolCallInputDelta
ResponseEvent::ToolCallInputDelta {
item_id: _,
call_id,
delta,
} => {
let Some((active_call_id, consumer)) = active_tool_argument_diff_consumer.as_mut()
else {
continue;
};
let call_id = match call_id {
Some(call_id) if call_id.as_str() != active_call_id.as_str() => continue,
Some(call_id) => call_id,
None => active_call_id.clone(),
};
if let Some(event) = consumer.consume_diff(turn_context.as_ref(), call_id, &delta) {
sess.send_event(&turn_context, event).await;
}
}item id 在这个分支被显式忽略,实际关联键是 tool call id。因为参数 diff consumer 由 OutputItemAdded 创建并 且一次只保存当前 consumer,Core 不允许用一个旧 item id 越过当前调用生命周期。
7. 推理摘要
推理有三类不同数据,不应合并为一个“思维文本”:
| API 事件 | Core 事件 | 关联字段 | 语义 |
|---|---|---|---|
reasoning_summary_text.delta | ReasoningContentDelta | summary_index | 可逐段显示的摘要文本 |
reasoning_summary_part.added | AgentReasoningSectionBreak | summary_index | 摘要段边界 |
reasoning_text.delta | ReasoningRawContentDelta | content_index | 原始 reasoning 内容,受配置控制 |
reasoning_summary_text.done | ReasoningContentDelta | item_id、summary_index | sequential cutoff 下的原子摘要 |
summary_index 是摘要段索引,content_index 是原始内容索引,不能互换。item_id 只在 done summary 事件中由 wire 明确提供,普通 delta 使用当前 active reasoning item。
源码位置:codex-rs/core/src/session/turn.rs :: reasoning delta 分支
ResponseEvent::ReasoningSummaryDelta {
delta,
summary_index,
} => {
if uses_sequential_cutoff_reasoning_summaries {
continue;
}
if let Some(active) = active_item.as_ref() {
if !active_item_is_streaming_to_client {
continue;
}
let event = ReasoningContentDeltaEvent {
thread_id: sess.thread_id.to_string(),
turn_id: turn_context.sub_id.clone(),
item_id: active.id(),
delta,
summary_index,
};
sess.send_event(&turn_context, EventMsg::ReasoningContentDelta(event))
.await;
} else {
error_or_panic("ReasoningSummaryDelta without active item".to_string());
}
}
ResponseEvent::ReasoningContentDelta {
delta,
content_index,
} => {
if let Some(active) = active_item.as_ref() {
if !active_item_is_streaming_to_client {
continue;
}
let event = ReasoningRawContentDeltaEvent {
thread_id: sess.thread_id.to_string(),
turn_id: turn_context.sub_id.clone(),
item_id: active.id(),
delta,
content_index,
};
sess.send_event(&turn_context, EventMsg::ReasoningRawContentDelta(event))
.await;
} else {
error_or_panic("ReasoningRawContentDelta without active item".to_string());
}
}原始 reasoning delta 是否对外出现还取决于 show_raw_agent_reasoning 等请求配置;不能从“服务端发送过 reasoning_text.delta”推断客户端一定会显示它。持久化的完整 ResponseItem::Reasoning 仍由 OutputItemDone 路径处理,实时 delta 只是展示投影。
8. 摘要截止模式
当 ConcurrentReasoningSummaries feature 开启且 provider 是 OpenAI 时,Core 使用 sequential cutoff。此时 实时的 ReasoningSummaryDelta 和 ReasoningSummaryPartAdded 被跳过,改用 ReasoningSummaryDone 的完整 段落。这样客户端不会同时看到部分摘要和最终同段摘要两份内容。
源码位置:codex-rs/core/src/session/turn.rs :: ReasoningSummaryDone
ResponseEvent::ReasoningSummaryDone {
item_id,
text,
summary_index,
} => {
if !uses_sequential_cutoff_reasoning_summaries {
continue;
}
let Some(active) = active_item.as_ref() else {
continue;
};
if !active_item_is_streaming_to_client || active.id() != item_id {
continue;
}
if summary_index > 0 {
sess.send_event(
&turn_context,
EventMsg::AgentReasoningSectionBreak(AgentReasoningSectionBreakEvent {
item_id: item_id.clone(),
summary_index,
}),
)
.await;
}
let event = ReasoningContentDeltaEvent {
thread_id: sess.thread_id.to_string(),
turn_id: turn_context.sub_id.clone(),
item_id,
delta: text,
summary_index,
};
sess.send_event(&turn_context, EventMsg::ReasoningContentDelta(event))
.await;
}这里的 active.id() != item_id 检查是过期事件保护。测试故意在 reasoning item 结束、message item 开始后再 发送旧 reasoning item 的 summary done;Core 丢弃它,不把 late step 挂到当前 assistant message。
该 feature 只在配置和 provider 条件同时满足时生效。对非 OpenAI provider,即使 feature 打开,也不能把 ReasoningSummaryDone 路径当作必然行为;请求构造测试明确验证了此时不会发送 sequential_cutoff。
9. 完整收口
delta 是可丢弃、可暂停、可重放的展示级输入;OutputItemDone 才触发完整 item 的工具、历史和 follow-up 决策。response.completed 之后还会 flush parser、等待 in-flight tool、记录 token usage,并根据 end_turn 决定是否继续下一次 sampling。
这条顺序解释了为什么“看到了文本 delta”不等于“历史已经写入该文本”,也解释了为什么工具参数 delta 只能 更新 diff consumer,而完整 function/custom call 仍需在 item done 后进入工具调度。
10. 增量事件测试
codex-utils-stream-parser 的测试 parses_citations_across_seed_and_delta_boundaries 把 <oai-mem-citation> 拆在两个输入 chunk 中,断言可见文本只保留 hello 与 world,citation payload 在 第二个 chunk 完成闭合后才提取。它证明 parser 的状态跨 delta 保持,不证明 Core 会把 citation 发送成公开事件。
core 的 assistant_message_stream_parsers_seed_plan_parser_across_added_and_delta_boundaries 先以 output_item.added seed Intro\n<proposed,再以 delta 补齐 _plan>、计划正文、结束标签和 Outro,断言输出 分别是 Normal、PlanStart、PlanDelta、PlanEnd、Normal。它证明 added 与 delta 共用同一 item parser,不证明真实 provider 一定按这个事件顺序发送。
reasoning_content_delta_has_item_metadata 通过 mock SSE 输入 reasoning item added、summary delta 和完整 reasoning item,断言对外 ReasoningContentDelta.item_id 等于 ItemStarted 的 reasoning id,证明 Core 使用 活动 item 补齐 wire delta 的归属。
sequential_cutoff_renders_done_summaries_for_active_reasoning_item 开启 feature,输入两个 done summary,随后 切换到 message item,再发送旧 reasoning item 的 late summary。测试断言只收到 step one、step two 和一个 section break,证明 partial delta 被抑制、summary index 1 产生边界、过期 item 不污染当前消息。
reasoning_raw_content_delta_respects_flag 打开 raw reasoning 配置后才断言 ReasoningRawContentDelta,说明 原始推理展示是配置条件,不是所有 reasoning_text.delta 的默认公开输出。
11. 复核路径
用下面的搜索可以从服务端事件走到三个不同消费者,不依赖文章内部编号:
rg -n "OutputTextDelta|ReasoningSummaryDelta|ReasoningContentDelta|ToolCallInputDelta" \
codex-rs/codex-api/src codex-rs/core/src/session/turn.rs
rg -n "AssistantTextStreamParser|seed_item_text|finish_item|flush_assistant_text_segments" \
codex-rs/utils/stream-parser/src codex-rs/core/src/session/turn.rs
rg -n "sequential_cutoff|late step|reasoning_content_delta_has_item_metadata" \
codex-rs/core/tests codex-rs/core/src/session/tests.rs最小测试命令:
cd codex-rs
cargo test -p codex-utils-stream-parser assistant_text::tests::parses_citations_across_seed_and_delta_boundaries -- --exact
cargo test -p codex-core session::tests::assistant_message_stream_parsers_seed_plan_parser_across_added_and_delta_boundaries -- --exact
RUST_MIN_STACK=16777216 cargo test -p codex-core --test all suite::items::reasoning_content_delta_has_item_metadata -- --exact
RUST_MIN_STACK=16777216 cargo test -p codex-core --test all suite::items::sequential_cutoff_renders_done_summaries_for_active_reasoning_item -- --exact读者可以用两个故障问题检查自己是否走通了源码:如果 UI 少了一段文本,先判断它是 citation/plan parser 缓存未 flush、active_item_is_streaming_to_client=false,还是 provider 根本没有对应 delta;如果 reasoning 摘要出现在错误的 assistant item 上,应该检查 ReasoningSummaryDone 的 item_id 比较和 OutputItemDone 的活动切换,而不是只检查网络顺序。
本文只覆盖 provider delta 到 Core 事件的投影。最终完整 item 如何写入 rollout、工具如何执行以及下一次 sampling 如何消费结果,分别由 模型上下文体系总览、工具系列和 运行时流程文章继续展开。
