codex / codex-rs /app-server-protocol /src /protocol /thread_history.rs
SaylorTwift's picture
SaylorTwift HF Staff
Add files using upload-large-folder tool
afa0cbf verified
Raw History Blame Contribute Delete
203 kB
use crate::protocol::item_builders::build_command_execution_begin_item;
use crate::protocol::item_builders::build_command_execution_end_item;
use crate::protocol::item_builders::build_file_change_approval_request_item;
use crate::protocol::item_builders::build_file_change_begin_item;
use crate::protocol::item_builders::build_file_change_end_item;
use crate::protocol::item_builders::build_item_from_guardian_event;
use crate::protocol::item_builders::review_output_text;
use crate::protocol::v2::CollabAgentState;
use crate::protocol::v2::CollabAgentTool;
use crate::protocol::v2::CollabAgentToolCallStatus;
use crate::protocol::v2::CommandExecutionStatus;
use crate::protocol::v2::DynamicToolCallOutputContentItem;
use crate::protocol::v2::DynamicToolCallStatus;
use crate::protocol::v2::ImageReference;
use crate::protocol::v2::McpToolCallAppContext;
use crate::protocol::v2::McpToolCallError;
use crate::protocol::v2::McpToolCallResult;
use crate::protocol::v2::McpToolCallStatus;
use crate::protocol::v2::ThreadItem;
use crate::protocol::v2::Turn;
use crate::protocol::v2::TurnError as V2TurnError;
use crate::protocol::v2::TurnError;
use crate::protocol::v2::TurnItemsView;
use crate::protocol::v2::TurnStatus;
use crate::protocol::v2::UserInput;
#[cfg(test)]
use crate::protocol::v2::WebSearchAction;
use crate::protocol::v2::WebSearchItem;
use crate::protocol::v2::web_search_action_from_core;
use codex_extension_items::image_generation::ImageGenerationItem;
use codex_protocol::items::parse_hook_prompt_message;
use codex_protocol::protocol::AgentMessageEvent;
use codex_protocol::protocol::AgentReasoningEvent;
use codex_protocol::protocol::AgentReasoningRawContentEvent;
use codex_protocol::protocol::AgentStatus;
use codex_protocol::protocol::ApplyPatchApprovalRequestEvent;
use codex_protocol::protocol::ContextCompactedEvent;
use codex_protocol::protocol::DynamicToolCallResponseEvent;
use codex_protocol::protocol::ErrorEvent;
use codex_protocol::protocol::EventMsg;
use codex_protocol::protocol::ExecCommandBeginEvent;
use codex_protocol::protocol::ExecCommandEndEvent;
use codex_protocol::protocol::GuardianAssessmentEvent;
use codex_protocol::protocol::GuardianAssessmentStatus;
use codex_protocol::protocol::ImageGenerationBeginEvent;
use codex_protocol::protocol::ImageGenerationEndEvent;
use codex_protocol::protocol::ItemCompletedEvent;
use codex_protocol::protocol::ItemStartedEvent;
use codex_protocol::protocol::McpToolCallBeginEvent;
use codex_protocol::protocol::McpToolCallEndEvent;
use codex_protocol::protocol::PatchApplyBeginEvent;
use codex_protocol::protocol::PatchApplyEndEvent;
use codex_protocol::protocol::ThreadRolledBackEvent;
use codex_protocol::protocol::TurnAbortedEvent;
use codex_protocol::protocol::TurnCompleteEvent;
use codex_protocol::protocol::TurnStartedEvent;
use codex_protocol::protocol::UserMessageEvent;
use codex_protocol::protocol::UserMessageImageKind;
use codex_protocol::protocol::ViewImageToolCallEvent;
use codex_protocol::protocol::WebSearchBeginEvent;
use codex_protocol::protocol::WebSearchEndEvent;
#[cfg(test)]
use codex_protocol::review_format::REVIEW_FALLBACK_MESSAGE;
use codex_rollout::CompactedItem;
use codex_rollout::RolloutItem;
use std::collections::HashMap;
use tracing::warn;
use uuid::Uuid;
#[cfg(test)]
use crate::protocol::v2::CommandAction;
#[cfg(test)]
use crate::protocol::v2::FileUpdateChange;
#[cfg(test)]
use crate::protocol::v2::PatchApplyStatus;
#[cfg(test)]
use crate::protocol::v2::PatchChangeKind;
#[cfg(test)]
use codex_protocol::protocol::ExecCommandStatus as CoreExecCommandStatus;
#[cfg(test)]
use codex_protocol::protocol::PatchApplyStatus as CorePatchApplyStatus;
/// Convert persisted [`RolloutItem`] entries into a sequence of [`Turn`] values.
///
/// When available, this uses `TurnContext.turn_id` as the canonical turn id so
/// resumed/rebuilt thread history preserves the original turn identifiers.
pub fn build_turns_from_rollout_items(items: &[RolloutItem]) -> Vec<Turn> {
let mut builder = ThreadHistoryBuilder::new();
for item in items {
builder.handle_rollout_item(item);
}
builder.finish()
}
/// A materialized `ThreadItem` snapshot that changed while handling one input.
#[derive(Debug, Clone, PartialEq)]
pub struct ThreadHistoryItemChange {
pub turn_id: String,
pub item: ThreadItem,
pub started_at_ms: Option<i64>,
pub completed_at_ms: Option<i64>,
}
/// Lightweight turn metadata snapshot for projectors that track turn status without
/// re-reading the full item list.
#[derive(Debug, Clone, PartialEq)]
pub struct ThreadHistoryTurnChange {
pub turn_id: String,
pub root_turn_id: Option<String>,
pub status: TurnStatus,
pub error: Option<TurnError>,
pub started_at: Option<i64>,
pub completed_at: Option<i64>,
pub duration_ms: Option<i64>,
}
/// Incremental changes produced by opt-in `ThreadHistoryBuilder` handlers.
#[derive(Debug, Default, Clone, PartialEq)]
pub struct ThreadHistoryChangeSet {
pub changed_items: Vec<ThreadHistoryItemChange>,
pub changed_turns: Vec<ThreadHistoryTurnChange>,
pub removed_turn_ids: Vec<String>,
}
impl ThreadHistoryChangeSet {
pub fn is_empty(&self) -> bool {
self.changed_items.is_empty()
&& self.changed_turns.is_empty()
&& self.removed_turn_ids.is_empty()
}
}
impl ThreadHistoryTurnChange {
fn from_pending_turn(turn: &PendingTurn) -> Self {
Self {
turn_id: turn.id.clone(),
root_turn_id: turn.root_turn_id.clone(),
status: turn.status.clone(),
error: turn.error.clone(),
started_at: turn.started_at,
completed_at: turn.completed_at,
duration_ms: turn.duration_ms,
}
}
}
/// Coalesces per-rollout-item changes into an end-of-batch view. It preserves
/// first-change order while replacing repeated item/turn snapshots with their
/// latest value, and drops accumulated changes for turns removed by rollback.
#[derive(Default)]
struct ThreadHistoryChangeAccumulator {
changed_items: Vec<Option<ThreadHistoryItemChange>>,
changed_item_indexes: HashMap<(String, String), usize>,
changed_turns: Vec<Option<ThreadHistoryTurnChange>>,
changed_turn_indexes: HashMap<String, usize>,
removed_turn_ids: Vec<String>,
removed_turn_indexes: HashMap<String, usize>,
}
impl ThreadHistoryChangeAccumulator {
fn push(&mut self, changes: ThreadHistoryChangeSet) {
for turn_id in changes.removed_turn_ids {
self.push_removed_turn_id(turn_id);
}
for item_change in changes.changed_items {
self.push_item_change(item_change);
}
for turn_change in changes.changed_turns {
self.push_turn_change(turn_change);
}
}
fn finish(self) -> ThreadHistoryChangeSet {
ThreadHistoryChangeSet {
changed_items: self.changed_items.into_iter().flatten().collect(),
changed_turns: self.changed_turns.into_iter().flatten().collect(),
removed_turn_ids: self.removed_turn_ids,
}
}
fn push_item_change(&mut self, change: ThreadHistoryItemChange) {
let key = (change.turn_id.clone(), change.item.id().to_string());
if let Some(index) = self.changed_item_indexes.get(&key).copied() {
self.changed_items[index] = Some(change);
return;
}
self.changed_item_indexes
.insert(key, self.changed_items.len());
self.changed_items.push(Some(change));
}
fn push_turn_change(&mut self, change: ThreadHistoryTurnChange) {
if let Some(index) = self.changed_turn_indexes.get(&change.turn_id).copied() {
self.changed_turns[index] = Some(change);
return;
}
self.changed_turn_indexes
.insert(change.turn_id.clone(), self.changed_turns.len());
self.changed_turns.push(Some(change));
}
fn push_removed_turn_id(&mut self, turn_id: String) {
if !self.removed_turn_indexes.contains_key(&turn_id) {
self.removed_turn_indexes
.insert(turn_id.clone(), self.removed_turn_ids.len());
self.removed_turn_ids.push(turn_id.clone());
}
if let Some(index) = self.changed_turn_indexes.remove(&turn_id) {
self.changed_turns[index] = None;
}
let removed_item_keys: Vec<(String, String)> = self
.changed_item_indexes
.keys()
.filter(|(item_turn_id, _)| item_turn_id == &turn_id)
.cloned()
.collect();
for key in removed_item_keys {
if let Some(index) = self.changed_item_indexes.remove(&key) {
self.changed_items[index] = None;
}
}
}
}
pub struct ThreadHistoryBuilder {
// Retain the builder representation so late completions can reuse each
// finished turn's item index without adding it to the public Turn type.
turns: Vec<PendingTurn>,
current_turn: Option<PendingTurn>,
next_item_index: i64,
current_rollout_index: usize,
next_rollout_index: usize,
active_change_set: Option<ThreadHistoryChangeSet>,
}
impl Default for ThreadHistoryBuilder {
fn default() -> Self {
Self::new()
}
}
impl ThreadHistoryBuilder {
pub fn new() -> Self {
Self {
turns: Vec::new(),
current_turn: None,
next_item_index: 1,
current_rollout_index: 0,
next_rollout_index: 0,
active_change_set: None,
}
}
pub fn reset(&mut self) {
*self = Self::new();
}
pub fn finish(mut self) -> Vec<Turn> {
self.finish_current_turn();
self.turns.into_iter().map(Turn::from).collect()
}
pub fn active_turn_snapshot(&self) -> Option<Turn> {
self.current_turn
.as_ref()
.map(Turn::from)
.or_else(|| self.turns.last().map(Turn::from))
}
/// Returns the id of the active turn without materializing its items.
pub fn active_turn_id(&self) -> Option<&str> {
self.current_turn
.as_ref()
.map(|turn| turn.id.as_str())
.or_else(|| self.turns.last().map(|turn| turn.id.as_str()))
}
pub fn turn_snapshot(&self, turn_id: &str) -> Option<Turn> {
self.current_turn
.as_ref()
.filter(|turn| turn.id == turn_id)
.map(Turn::from)
.or_else(|| {
self.turns
.iter()
.find(|turn| turn.id == turn_id)
.map(Turn::from)
})
}
/// Returns the index of the active turn snapshot within the finished turn list.
///
/// When a turn is still open, this is the index it will occupy after
/// `finish`. When no turn is open, it is the index of the last finished turn.
pub fn active_turn_position(&self) -> Option<usize> {
if self.current_turn.is_some() {
Some(self.turns.len())
} else if self.turns.is_empty() {
None
} else {
Some(self.turns.len() - 1)
}
}
pub fn has_active_turn(&self) -> bool {
self.current_turn.is_some()
}
pub fn active_turn_id_if_explicit(&self) -> Option<String> {
self.current_turn
.as_ref()
.filter(|turn| turn.opened_explicitly)
.map(|turn| turn.id.clone())
}
pub fn active_turn_start_index(&self) -> Option<usize> {
self.current_turn
.as_ref()
.map(|turn| turn.rollout_start_index)
}
/// Shared reducer for persisted rollout replay and in-memory current-turn
/// tracking used by running thread resume/rejoin.
///
/// This function should handle all EventMsg variants that can be persisted in a rollout file.
/// See `should_persist_event_msg` in `codex-rs/core/rollout/policy.rs`.
pub fn handle_event(&mut self, event: &EventMsg) {
match event {
EventMsg::UserMessage(payload) => self.handle_user_message(payload),
EventMsg::AgentMessage(payload) => self.handle_agent_message(payload),
EventMsg::AgentReasoning(payload) => self.handle_agent_reasoning(payload),
EventMsg::AgentReasoningRawContent(payload) => {
self.handle_agent_reasoning_raw_content(payload)
}
EventMsg::WebSearchBegin(payload) => self.handle_web_search_begin(payload),
EventMsg::WebSearchEnd(payload) => self.handle_web_search_end(payload),
EventMsg::ExecCommandBegin(payload) => self.handle_exec_command_begin(payload),
EventMsg::ExecCommandEnd(payload) => self.handle_exec_command_end(payload),
EventMsg::GuardianAssessment(payload) => self.handle_guardian_assessment(payload),
EventMsg::ApplyPatchApprovalRequest(payload) => {
self.handle_apply_patch_approval_request(payload)
}
EventMsg::PatchApplyBegin(payload) => self.handle_patch_apply_begin(payload),
EventMsg::PatchApplyEnd(payload) => self.handle_patch_apply_end(payload),
EventMsg::DynamicToolCallRequest(payload) => {
self.handle_dynamic_tool_call_request(payload)
}
EventMsg::DynamicToolCallResponse(payload) => {
self.handle_dynamic_tool_call_response(payload)
}
EventMsg::McpToolCallBegin(payload) => self.handle_mcp_tool_call_begin(payload),
EventMsg::McpToolCallEnd(payload) => self.handle_mcp_tool_call_end(payload),
EventMsg::ViewImageToolCall(payload) => self.handle_view_image_tool_call(payload),
EventMsg::ImageGenerationBegin(payload) => self.handle_image_generation_begin(payload),
EventMsg::ImageGenerationEnd(payload) => self.handle_image_generation_end(payload),
EventMsg::CollabAgentSpawnBegin(payload) => {
self.handle_collab_agent_spawn_begin(payload)
}
EventMsg::CollabAgentSpawnEnd(payload) => self.handle_collab_agent_spawn_end(payload),
EventMsg::CollabAgentInteractionBegin(payload) => {
self.handle_collab_agent_interaction_begin(payload)
}
EventMsg::CollabAgentInteractionEnd(payload) => {
self.handle_collab_agent_interaction_end(payload)
}
EventMsg::SubAgentActivity(payload) => self.handle_sub_agent_activity(payload),
EventMsg::CollabWaitingBegin(payload) => self.handle_collab_waiting_begin(payload),
EventMsg::CollabWaitingEnd(payload) => self.handle_collab_waiting_end(payload),
EventMsg::CollabCloseBegin(payload) => self.handle_collab_close_begin(payload),
EventMsg::CollabCloseEnd(payload) => self.handle_collab_close_end(payload),
EventMsg::CollabResumeBegin(payload) => self.handle_collab_resume_begin(payload),
EventMsg::CollabResumeEnd(payload) => self.handle_collab_resume_end(payload),
EventMsg::ContextCompacted(payload) => self.handle_context_compacted(payload),
EventMsg::EnteredReviewMode(payload) => self.handle_entered_review_mode(payload),
EventMsg::ExitedReviewMode(payload) => self.handle_exited_review_mode(payload),
EventMsg::ItemStarted(payload) => self.handle_item_started(payload),
EventMsg::ItemCompleted(payload) => self.handle_item_completed(payload),
EventMsg::HookStarted(_) | EventMsg::HookCompleted(_) => {}
EventMsg::Error(payload) => self.handle_error(payload),
EventMsg::TokenCount(_) => {}
EventMsg::ThreadRolledBack(payload) => self.handle_thread_rollback(payload),
EventMsg::TurnAborted(payload) => self.handle_turn_aborted(payload),
EventMsg::TurnStarted(payload) => self.handle_turn_started(payload),
EventMsg::TurnComplete(payload) => self.handle_turn_complete(payload),
_ => {}
}
}
pub fn handle_rollout_item(&mut self, item: &RolloutItem) {
self.current_rollout_index = self.next_rollout_index;
self.next_rollout_index += 1;
match item {
RolloutItem::EventMsg(event) => self.handle_event(event),
RolloutItem::Compacted(payload) => self.handle_compacted(payload),
RolloutItem::ResponseItem(item) => self.handle_response_item(&item.item),
RolloutItem::InterAgentCommunication(_)
| RolloutItem::InterAgentCommunicationMetadata { .. }
| RolloutItem::TurnContext(_)
| RolloutItem::TokenUsageRecord(_)
| RolloutItem::WorldState(_)
| RolloutItem::RealtimeItem(_)
| RolloutItem::RetainedContext(_)
| RolloutItem::SecurityRiskScore(_)
| RolloutItem::SessionMeta(_) => {}
}
}
/// Handles one event and returns the materialized items or turn metadata
/// changed by that event.
pub fn handle_event_with_changes(&mut self, event: &EventMsg) -> ThreadHistoryChangeSet {
self.collect_changes(|builder| builder.handle_event(event))
}
/// Handles a rollout item and returns the materialized items or turn metadata
/// changed by that one append.
pub fn handle_rollout_item_with_changes(
&mut self,
item: &RolloutItem,
) -> ThreadHistoryChangeSet {
self.collect_changes(|builder| builder.handle_rollout_item(item))
}
/// Handles rollout items in order and returns a coalesced end-of-batch
/// change set. Multiple changes to the same item or turn are deduplicated
/// so only the latest snapshot is emitted.
pub fn handle_rollout_items_with_changes(
&mut self,
items: &[RolloutItem],
) -> ThreadHistoryChangeSet {
let mut accumulator = ThreadHistoryChangeAccumulator::default();
for item in items {
accumulator.push(self.handle_rollout_item_with_changes(item));
}
accumulator.finish()
}
fn collect_changes(&mut self, handle: impl FnOnce(&mut Self)) -> ThreadHistoryChangeSet {
debug_assert!(self.active_change_set.is_none());
self.active_change_set = Some(ThreadHistoryChangeSet::default());
handle(self);
self.active_change_set.take().unwrap_or_default()
}
fn handle_response_item(&mut self, item: &codex_protocol::models::ResponseItem) {
let codex_protocol::models::ResponseItem::Message {
role, content, id, ..
} = item
else {
return;
};
if role != "user" {
return;
}
let Some(hook_prompt) = parse_hook_prompt_message(id.as_deref(), content) else {
return;
};
self.push_item_in_current_turn(ThreadItem::HookPrompt {
id: hook_prompt.id,
fragments: hook_prompt
.fragments
.into_iter()
.map(crate::protocol::v2::HookPromptFragment::from)
.collect(),
});
}
fn handle_user_message(&mut self, payload: &UserMessageEvent) {
// User messages should stay in explicitly opened turns. For backward
// compatibility with older streams that did not open turns explicitly,
// close any implicit/inactive turn and start a fresh one for this input.
if let Some(turn) = self.current_turn.as_ref()
&& !turn.opened_explicitly
&& !(turn.saw_compaction && turn.items.is_empty())
{
self.finish_current_turn();
}
let id = self.next_item_id();
let content = self.build_user_inputs(payload);
self.push_item_in_current_turn(ThreadItem::UserMessage {
id,
client_id: payload.client_id.clone(),
content,
});
}
fn handle_agent_message(&mut self, payload: &AgentMessageEvent) {
if payload.message.is_empty() {
return;
}
let id = self.next_item_id();
self.push_item_in_current_turn(ThreadItem::AgentMessage {
id,
text: payload.message.clone(),
phase: payload.phase.clone(),
memory_citation: payload.memory_citation.clone().map(Into::into),
delivery: payload.delivery,
questions: payload.questions.clone(),
});
}
fn handle_agent_reasoning(&mut self, payload: &AgentReasoningEvent) {
if payload.text.is_empty() {
return;
}
// If the last item is a reasoning item, add the new text to the summary.
let existing_item_change = {
let tracking_changes = self.is_tracking_changes();
let turn = self.ensure_turn();
if let Some(ThreadItem::Reasoning { summary, .. }) = turn.items.last_mut() {
summary.push(payload.text.clone());
let changed_item = if tracking_changes {
turn.items
.last()
.cloned()
.map(|item| (turn.id.clone(), item))
} else {
None
};
Some(changed_item)
} else {
None
}
};
if let Some(changed_item) = existing_item_change {
if let Some((turn_id, item)) = changed_item {
self.record_changed_item(turn_id, item);
}
return;
}
// Otherwise, create a new reasoning item.
let id = self.next_item_id();
self.push_item_in_current_turn(ThreadItem::Reasoning {
id,
summary: vec![payload.text.clone()],
content: Vec::new(),
});
}
fn handle_agent_reasoning_raw_content(&mut self, payload: &AgentReasoningRawContentEvent) {
if payload.text.is_empty() {
return;
}
// If the last item is a reasoning item, add the new text to the content.
let existing_item_change = {
let tracking_changes = self.is_tracking_changes();
let turn = self.ensure_turn();
if let Some(ThreadItem::Reasoning { content, .. }) = turn.items.last_mut() {
content.push(payload.text.clone());
let changed_item = if tracking_changes {
turn.items
.last()
.cloned()
.map(|item| (turn.id.clone(), item))
} else {
None
};
Some(changed_item)
} else {
None
}
};
if let Some(changed_item) = existing_item_change {
if let Some((turn_id, item)) = changed_item {
self.record_changed_item(turn_id, item);
}
return;
}
// Otherwise, create a new reasoning item.
let id = self.next_item_id();
self.push_item_in_current_turn(ThreadItem::Reasoning {
id,
summary: Vec::new(),
content: vec![payload.text.clone()],
});
}
fn handle_item_started(&mut self, payload: &ItemStartedEvent) {
self.handle_materialized_item_lifecycle(&payload.turn_id, &payload.item);
}
fn handle_item_completed(&mut self, payload: &ItemCompletedEvent) {
self.handle_materialized_item_lifecycle(&payload.turn_id, &payload.item);
}
fn handle_materialized_item_lifecycle(
&mut self,
turn_id: &str,
item: &codex_protocol::items::TurnItem,
) {
let is_review_mode_item = matches!(
item,
codex_protocol::items::TurnItem::EnteredReviewMode(_)
| codex_protocol::items::TurnItem::ExitedReviewMode(_)
);
let should_upsert = match item {
codex_protocol::items::TurnItem::Plan(plan) => !plan.text.is_empty(),
codex_protocol::items::TurnItem::HookPrompt(_)
| codex_protocol::items::TurnItem::FunctionCallOutput(_)
| codex_protocol::items::TurnItem::CommandExecution(_)
| codex_protocol::items::TurnItem::DynamicToolCall(_)
| codex_protocol::items::TurnItem::CollabAgentToolCall(_)
| codex_protocol::items::TurnItem::SubAgentActivity(_)
| codex_protocol::items::TurnItem::Extension(_)
| codex_protocol::items::TurnItem::EnteredReviewMode(_)
| codex_protocol::items::TurnItem::ExitedReviewMode(_) => true,
codex_protocol::items::TurnItem::UserMessage(_)
| codex_protocol::items::TurnItem::AgentMessage(_)
| codex_protocol::items::TurnItem::Reasoning(_)
| codex_protocol::items::TurnItem::WebSearch(_)
| codex_protocol::items::TurnItem::ImageView(_)
| codex_protocol::items::TurnItem::ImageGeneration(_)
| codex_protocol::items::TurnItem::FileChange(_)
| codex_protocol::items::TurnItem::McpToolCall(_)
| codex_protocol::items::TurnItem::ContextCompaction(_) => false,
};
if should_upsert {
let item = ThreadItem::from(item.clone());
if is_review_mode_item {
self.upsert_review_mode_item(Some(turn_id), item);
} else {
self.upsert_item_in_turn_id(turn_id, item);
}
}
}
fn handle_web_search_begin(&mut self, payload: &WebSearchBeginEvent) {
let item = ThreadItem::WebSearch(WebSearchItem {
id: payload.call_id.clone(),
query: String::new(),
action: None,
results: None,
});
self.upsert_item_in_current_turn(item);
}
fn handle_web_search_end(&mut self, payload: &WebSearchEndEvent) {
let item = ThreadItem::WebSearch(WebSearchItem {
id: payload.call_id.clone(),
query: payload.query.clone(),
action: Some(web_search_action_from_core(payload.action.clone())),
results: payload.results.clone(),
});
self.upsert_item_in_current_turn(item);
}
fn handle_exec_command_begin(&mut self, payload: &ExecCommandBeginEvent) {
let item = build_command_execution_begin_item(payload);
self.upsert_item_in_turn_id(&payload.turn_id, item);
}
fn handle_exec_command_end(&mut self, payload: &ExecCommandEndEvent) {
let item = build_command_execution_end_item(payload);
// Command completions can arrive out of order. Unified exec may return
// while a PTY is still running, then emit ExecCommandEnd later from a
// background exit watcher when that process finally exits. By then, a
// newer user turn may already have started. Route by event turn_id so
// replay preserves the original turn association.
self.upsert_item_in_turn_id(&payload.turn_id, item);
}
fn handle_guardian_assessment(&mut self, payload: &GuardianAssessmentEvent) {
let status = match payload.status {
GuardianAssessmentStatus::InProgress => CommandExecutionStatus::InProgress,
GuardianAssessmentStatus::Denied | GuardianAssessmentStatus::Aborted => {
CommandExecutionStatus::Declined
}
GuardianAssessmentStatus::TimedOut => CommandExecutionStatus::Failed,
GuardianAssessmentStatus::Approved => return,
};
let Some(item) = build_item_from_guardian_event(payload, status) else {
return;
};
if payload.turn_id.is_empty() {
self.upsert_item_in_current_turn(item);
} else {
self.upsert_item_in_turn_id(&payload.turn_id, item);
}
}
fn handle_apply_patch_approval_request(&mut self, payload: &ApplyPatchApprovalRequestEvent) {
let item = build_file_change_approval_request_item(payload);
if payload.turn_id.is_empty() {
self.upsert_item_in_current_turn(item);
} else {
self.upsert_item_in_turn_id(&payload.turn_id, item);
}
}
fn handle_patch_apply_begin(&mut self, payload: &PatchApplyBeginEvent) {
let item = build_file_change_begin_item(payload);
if payload.turn_id.is_empty() {
self.upsert_item_in_current_turn(item);
} else {
self.upsert_item_in_turn_id(&payload.turn_id, item);
}
}
fn handle_patch_apply_end(&mut self, payload: &PatchApplyEndEvent) {
let item = build_file_change_end_item(payload);
if payload.turn_id.is_empty() {
self.upsert_item_in_current_turn(item);
} else {
self.upsert_item_in_turn_id(&payload.turn_id, item);
}
}
fn handle_dynamic_tool_call_request(
&mut self,
payload: &codex_protocol::dynamic_tools::DynamicToolCallRequest,
) {
let item = ThreadItem::DynamicToolCall {
id: payload.call_id.clone(),
namespace: payload.namespace.clone(),
tool: payload.tool.clone(),
arguments: payload.arguments.clone(),
status: DynamicToolCallStatus::InProgress,
content_items: None,
success: None,
duration_ms: None,
};
if payload.turn_id.is_empty() {
self.upsert_item_in_current_turn(item);
} else {
self.upsert_item_in_turn_id(&payload.turn_id, item);
}
}
fn handle_dynamic_tool_call_response(&mut self, payload: &DynamicToolCallResponseEvent) {
let status = if payload.success {
DynamicToolCallStatus::Completed
} else {
DynamicToolCallStatus::Failed
};
let duration_ms = i64::try_from(payload.duration.as_millis()).ok();
let item = ThreadItem::DynamicToolCall {
id: payload.call_id.clone(),
namespace: payload.namespace.clone(),
tool: payload.tool.clone(),
arguments: payload.arguments.clone(),
status,
content_items: Some(convert_dynamic_tool_content_items(&payload.content_items)),
success: Some(payload.success),
duration_ms,
};
if payload.turn_id.is_empty() {
self.upsert_item_in_current_turn(item);
} else {
self.upsert_item_in_turn_id(&payload.turn_id, item);
}
}
fn handle_mcp_tool_call_begin(&mut self, payload: &McpToolCallBeginEvent) {
let item = ThreadItem::McpToolCall {
id: payload.call_id.clone(),
server: payload.invocation.server.clone(),
tool: payload.invocation.tool.clone(),
status: McpToolCallStatus::InProgress,
arguments: payload
.invocation
.arguments
.clone()
.unwrap_or(serde_json::Value::Null),
app_context: payload
.connector_id
.clone()
.map(|connector_id| McpToolCallAppContext {
connector_id,
link_id: payload.link_id.clone(),
resource_uri: payload.mcp_app_resource_uri.clone(),
app_name: payload.app_name.clone(),
action_name: payload.action_name.clone(),
}),
mcp_app_resource_uri: payload.mcp_app_resource_uri.clone(),
mcp_app_ui: payload.mcp_app_ui.clone(),
plugin_id: payload.plugin_id.clone(),
read_only_hint: payload.read_only_hint,
result: None,
error: None,
duration_ms: None,
};
if payload.turn_id.is_empty() {
self.upsert_item_in_current_turn(item);
} else {
self.upsert_item_in_turn_id(&payload.turn_id, item);
}
}
fn handle_mcp_tool_call_end(&mut self, payload: &McpToolCallEndEvent) {
let status = if payload.is_success() {
McpToolCallStatus::Completed
} else {
McpToolCallStatus::Failed
};
let duration_ms = i64::try_from(payload.duration.as_millis()).ok();
let (result, error) = match &payload.result {
Ok(value) => (
Some(Box::new(McpToolCallResult {
content: value.content.clone(),
structured_content: value.structured_content.clone(),
meta: value.meta.clone(),
})),
None,
),
Err(message) => (
None,
Some(McpToolCallError {
message: message.clone(),
}),
),
};
let item = ThreadItem::McpToolCall {
id: payload.call_id.clone(),
server: payload.invocation.server.clone(),
tool: payload.invocation.tool.clone(),
status,
arguments: payload
.invocation
.arguments
.clone()
.unwrap_or(serde_json::Value::Null),
app_context: payload
.connector_id
.clone()
.map(|connector_id| McpToolCallAppContext {
connector_id,
link_id: payload.link_id.clone(),
resource_uri: payload.mcp_app_resource_uri.clone(),
app_name: payload.app_name.clone(),
action_name: payload.action_name.clone(),
}),
mcp_app_resource_uri: payload.mcp_app_resource_uri.clone(),
mcp_app_ui: payload.mcp_app_ui.clone(),
plugin_id: payload.plugin_id.clone(),
read_only_hint: payload.read_only_hint,
result,
error,
duration_ms,
};
if payload.turn_id.is_empty() {
self.upsert_item_in_current_turn(item);
} else {
self.upsert_item_in_turn_id(&payload.turn_id, item);
}
}
fn handle_view_image_tool_call(&mut self, payload: &ViewImageToolCallEvent) {
let item = ThreadItem::ImageView {
id: payload.call_id.clone(),
path: payload.path.clone().into(),
};
self.upsert_item_in_current_turn(item);
}
fn handle_image_generation_begin(&mut self, payload: &ImageGenerationBeginEvent) {
let item = ThreadItem::ImageGeneration(ImageGenerationItem {
id: payload.call_id.clone(),
status: String::new(),
revised_prompt: None,
result: String::new(),
transparent_background: None,
failure: None,
saved_path: None,
imagegen_request_id: None,
generation_id: None,
});
self.upsert_item_in_current_turn(item);
}
fn handle_image_generation_end(&mut self, payload: &ImageGenerationEndEvent) {
let item = ThreadItem::ImageGeneration(ImageGenerationItem {
id: payload.call_id.clone(),
status: payload.status.clone(),
revised_prompt: payload.revised_prompt.clone(),
result: payload.result.clone(),
transparent_background: payload.transparent_background,
failure: payload.failure.clone(),
saved_path: payload.saved_path.clone(),
imagegen_request_id: None,
generation_id: None,
});
self.upsert_item_in_current_turn(item);
}
fn handle_collab_agent_spawn_begin(
&mut self,
payload: &codex_protocol::protocol::CollabAgentSpawnBeginEvent,
) {
let item = ThreadItem::CollabAgentToolCall {
id: payload.call_id.clone(),
tool: CollabAgentTool::SpawnAgent,
status: CollabAgentToolCallStatus::InProgress,
sender_thread_id: payload.sender_thread_id.to_string(),
receiver_thread_ids: Vec::new(),
prompt: Some(payload.prompt.clone()),
model: Some(payload.model.clone()),
reasoning_effort: Some(payload.reasoning_effort.clone()),
agents_states: HashMap::new(),
};
self.upsert_item_in_current_turn(item);
}
fn handle_collab_agent_spawn_end(
&mut self,
payload: &codex_protocol::protocol::CollabAgentSpawnEndEvent,
) {
let has_receiver = payload.new_thread_id.is_some();
let status = match &payload.status {
AgentStatus::Errored(_) | AgentStatus::NotFound => CollabAgentToolCallStatus::Failed,
_ if has_receiver => CollabAgentToolCallStatus::Completed,
_ => CollabAgentToolCallStatus::Failed,
};
let (receiver_thread_ids, agents_states) = match &payload.new_thread_id {
Some(id) => {
let receiver_id = id.to_string();
let received_status = CollabAgentState::from(payload.status.clone());
(
vec![receiver_id.clone()],
[(receiver_id, received_status)].into_iter().collect(),
)
}
None => (Vec::new(), HashMap::new()),
};
self.upsert_item_in_current_turn(ThreadItem::CollabAgentToolCall {
id: payload.call_id.clone(),
tool: CollabAgentTool::SpawnAgent,
status,
sender_thread_id: payload.sender_thread_id.to_string(),
receiver_thread_ids,
prompt: Some(payload.prompt.clone()),
model: Some(payload.model.clone()),
reasoning_effort: Some(payload.reasoning_effort.clone()),
agents_states,
});
}
fn handle_collab_agent_interaction_begin(
&mut self,
payload: &codex_protocol::protocol::CollabAgentInteractionBeginEvent,
) {
let item = ThreadItem::CollabAgentToolCall {
id: payload.call_id.clone(),
tool: CollabAgentTool::SendInput,
status: CollabAgentToolCallStatus::InProgress,
sender_thread_id: payload.sender_thread_id.to_string(),
receiver_thread_ids: vec![payload.receiver_thread_id.to_string()],
prompt: Some(payload.prompt.clone()),
model: None,
reasoning_effort: None,
agents_states: HashMap::new(),
};
self.upsert_item_in_current_turn(item);
}
fn handle_collab_agent_interaction_end(
&mut self,
payload: &codex_protocol::protocol::CollabAgentInteractionEndEvent,
) {
let status = match &payload.status {
AgentStatus::Errored(_) | AgentStatus::NotFound => CollabAgentToolCallStatus::Failed,
_ => CollabAgentToolCallStatus::Completed,
};
let receiver_id = payload.receiver_thread_id.to_string();
let received_status = CollabAgentState::from(payload.status.clone());
self.upsert_item_in_current_turn(ThreadItem::CollabAgentToolCall {
id: payload.call_id.clone(),
tool: CollabAgentTool::SendInput,
status,
sender_thread_id: payload.sender_thread_id.to_string(),
receiver_thread_ids: vec![receiver_id.clone()],
prompt: Some(payload.prompt.clone()),
model: None,
reasoning_effort: None,
agents_states: [(receiver_id, received_status)].into_iter().collect(),
});
}
fn handle_sub_agent_activity(
&mut self,
payload: &codex_protocol::protocol::SubAgentActivityEvent,
) {
self.upsert_item_in_current_turn(ThreadItem::SubAgentActivity {
id: payload.event_id.clone(),
kind: payload.kind.into(),
agent_thread_id: payload.agent_thread_id.to_string(),
agent_path: String::from(payload.agent_path.clone()),
});
}
fn handle_collab_waiting_begin(
&mut self,
payload: &codex_protocol::protocol::CollabWaitingBeginEvent,
) {
let item = ThreadItem::CollabAgentToolCall {
id: payload.call_id.clone(),
tool: CollabAgentTool::Wait,
status: CollabAgentToolCallStatus::InProgress,
sender_thread_id: payload.sender_thread_id.to_string(),
receiver_thread_ids: payload
.receiver_thread_ids
.iter()
.map(ToString::to_string)
.collect(),
prompt: None,
model: None,
reasoning_effort: None,
agents_states: HashMap::new(),
};
self.upsert_item_in_current_turn(item);
}
fn handle_collab_waiting_end(
&mut self,
payload: &codex_protocol::protocol::CollabWaitingEndEvent,
) {
let status = if payload
.statuses
.values()
.any(|status| matches!(status, AgentStatus::Errored(_) | AgentStatus::NotFound))
{
CollabAgentToolCallStatus::Failed
} else {
CollabAgentToolCallStatus::Completed
};
let mut receiver_thread_ids: Vec<String> =
payload.statuses.keys().map(ToString::to_string).collect();
receiver_thread_ids.sort();
let agents_states = payload
.statuses
.iter()
.map(|(id, status)| (id.to_string(), CollabAgentState::from(status.clone())))
.collect();
self.upsert_item_in_current_turn(ThreadItem::CollabAgentToolCall {
id: payload.call_id.clone(),
tool: CollabAgentTool::Wait,
status,
sender_thread_id: payload.sender_thread_id.to_string(),
receiver_thread_ids,
prompt: None,
model: None,
reasoning_effort: None,
agents_states,
});
}
fn handle_collab_close_begin(
&mut self,
payload: &codex_protocol::protocol::CollabCloseBeginEvent,
) {
let item = ThreadItem::CollabAgentToolCall {
id: payload.call_id.clone(),
tool: CollabAgentTool::CloseAgent,
status: CollabAgentToolCallStatus::InProgress,
sender_thread_id: payload.sender_thread_id.to_string(),
receiver_thread_ids: vec![payload.receiver_thread_id.to_string()],
prompt: None,
model: None,
reasoning_effort: None,
agents_states: HashMap::new(),
};
self.upsert_item_in_current_turn(item);
}
fn handle_collab_close_end(&mut self, payload: &codex_protocol::protocol::CollabCloseEndEvent) {
let status = match &payload.status {
AgentStatus::Errored(_) | AgentStatus::NotFound => CollabAgentToolCallStatus::Failed,
_ => CollabAgentToolCallStatus::Completed,
};
let receiver_id = payload.receiver_thread_id.to_string();
let agents_states = [(
receiver_id.clone(),
CollabAgentState::from(payload.status.clone()),
)]
.into_iter()
.collect();
self.upsert_item_in_current_turn(ThreadItem::CollabAgentToolCall {
id: payload.call_id.clone(),
tool: CollabAgentTool::CloseAgent,
status,
sender_thread_id: payload.sender_thread_id.to_string(),
receiver_thread_ids: vec![receiver_id],
prompt: None,
model: None,
reasoning_effort: None,
agents_states,
});
}
fn handle_collab_resume_begin(
&mut self,
payload: &codex_protocol::protocol::CollabResumeBeginEvent,
) {
let item = ThreadItem::CollabAgentToolCall {
id: payload.call_id.clone(),
tool: CollabAgentTool::ResumeAgent,
status: CollabAgentToolCallStatus::InProgress,
sender_thread_id: payload.sender_thread_id.to_string(),
receiver_thread_ids: vec![payload.receiver_thread_id.to_string()],
prompt: None,
model: None,
reasoning_effort: None,
agents_states: HashMap::new(),
};
self.upsert_item_in_current_turn(item);
}
fn handle_collab_resume_end(
&mut self,
payload: &codex_protocol::protocol::CollabResumeEndEvent,
) {
let status = match &payload.status {
AgentStatus::Errored(_) | AgentStatus::NotFound => CollabAgentToolCallStatus::Failed,
_ => CollabAgentToolCallStatus::Completed,
};
let receiver_id = payload.receiver_thread_id.to_string();
let agents_states = [(
receiver_id.clone(),
CollabAgentState::from(payload.status.clone()),
)]
.into_iter()
.collect();
self.upsert_item_in_current_turn(ThreadItem::CollabAgentToolCall {
id: payload.call_id.clone(),
tool: CollabAgentTool::ResumeAgent,
status,
sender_thread_id: payload.sender_thread_id.to_string(),
receiver_thread_ids: vec![receiver_id],
prompt: None,
model: None,
reasoning_effort: None,
agents_states,
});
}
fn handle_context_compacted(&mut self, _payload: &ContextCompactedEvent) {
let id = self.next_item_id();
self.push_item_in_current_turn(ThreadItem::ContextCompaction { id });
}
fn handle_entered_review_mode(
&mut self,
payload: &codex_protocol::protocol::EnteredReviewModeEvent,
) {
let review = payload
.user_facing_hint
.clone()
.unwrap_or_else(|| "Review requested.".to_string());
let id = payload
.item_id
.clone()
.unwrap_or_else(|| self.next_item_id());
self.upsert_review_mode_item(
payload.turn_id.as_deref(),
ThreadItem::EnteredReviewMode { id, review },
);
}
fn handle_exited_review_mode(
&mut self,
payload: &codex_protocol::protocol::ExitedReviewModeEvent,
) {
let review = review_output_text(payload.review_output.as_ref());
let id = payload
.item_id
.clone()
.unwrap_or_else(|| self.next_item_id());
self.upsert_review_mode_item(
payload.turn_id.as_deref(),
ThreadItem::ExitedReviewMode { id, review },
);
}
fn upsert_review_mode_item(&mut self, turn_id: Option<&str>, item: ThreadItem) {
let Some(turn_id) = turn_id else {
self.upsert_item_in_current_turn(item);
return;
};
let current_turn_matches = self
.current_turn
.as_ref()
.is_some_and(|turn| turn.id == turn_id);
if !current_turn_matches && !self.turns.iter().any(|turn| turn.id == turn_id) {
self.finish_current_turn();
let turn = self.new_turn(Some(turn_id.to_string()));
self.record_changed_pending_turn(&turn);
self.current_turn = Some(turn);
}
self.upsert_item_in_turn_id(turn_id, item);
}
fn handle_error(&mut self, payload: &ErrorEvent) {
if !payload.affects_turn_status() {
return;
}
let tracking_changes = self.is_tracking_changes();
let changed_turn = if let Some(turn) = self.current_turn.as_mut() {
turn.status = TurnStatus::Failed;
turn.error = Some(V2TurnError {
misalignment: payload.misalignment.clone().map(Into::into),
message: payload.message.clone(),
codex_error_info: payload.codex_error_info.clone().map(Into::into),
additional_details: None,
});
tracking_changes.then(|| ThreadHistoryTurnChange::from_pending_turn(turn))
} else {
None
};
if let Some(changed_turn) = changed_turn {
self.record_changed_turn(changed_turn);
}
}
fn handle_turn_aborted(&mut self, payload: &TurnAbortedEvent) {
let apply_abort = |turn: &mut PendingTurn| {
turn.status = TurnStatus::Interrupted;
turn.completed_at = payload.completed_at;
turn.duration_ms = payload.duration_ms;
ThreadHistoryTurnChange::from_pending_turn(turn)
};
if let Some(turn_id) = payload.turn_id.as_deref() {
// Prefer an exact ID match so we interrupt the turn explicitly targeted by the event.
if let Some(turn) = self.current_turn.as_mut().filter(|turn| turn.id == turn_id) {
let changed_turn = apply_abort(turn);
self.record_changed_turn(changed_turn);
return;
}
if let Some(turn) = self.turns.iter_mut().find(|turn| turn.id == turn_id) {
turn.status = TurnStatus::Interrupted;
turn.completed_at = payload.completed_at;
turn.duration_ms = payload.duration_ms;
let changed_turn = ThreadHistoryTurnChange::from_pending_turn(turn);
self.record_changed_turn(changed_turn);
return;
}
}
// If the event has no ID (or refers to an unknown turn), fall back to the active turn.
if let Some(turn) = self.current_turn.as_mut() {
let changed_turn = apply_abort(turn);
self.record_changed_turn(changed_turn);
}
}
fn handle_turn_started(&mut self, payload: &TurnStartedEvent) {
self.finish_current_turn();
let mut turn = self
.new_turn(Some(payload.turn_id.clone()))
.with_status(TurnStatus::InProgress)
.with_started_at(payload.started_at)
.opened_explicitly();
turn.root_turn_id = payload.root_turn_id.clone();
self.record_changed_pending_turn(&turn);
self.current_turn = Some(turn);
}
fn handle_turn_complete(&mut self, payload: &TurnCompleteEvent) {
let terminal_error = payload.error.as_ref().map(|error| V2TurnError {
misalignment: error.misalignment.clone().map(Into::into),
message: error.message.clone(),
codex_error_info: error.codex_error_info.clone().map(Into::into),
additional_details: None,
});
let apply_completion = |turn: &mut PendingTurn| {
if let Some(error) = terminal_error.as_ref() {
turn.status = TurnStatus::Failed;
turn.error = Some(error.clone());
} else if matches!(turn.status, TurnStatus::Completed | TurnStatus::InProgress) {
turn.status = TurnStatus::Completed;
}
turn.completed_at = payload.completed_at;
turn.duration_ms = payload.duration_ms;
ThreadHistoryTurnChange::from_pending_turn(turn)
};
// Prefer an exact ID match from the active turn and then close it.
if let Some(current_turn) = self
.current_turn
.as_mut()
.filter(|turn| turn.id == payload.turn_id)
{
let changed_turn = apply_completion(current_turn);
self.record_changed_turn(changed_turn);
self.finish_current_turn();
return;
}
if let Some(turn) = self
.turns
.iter_mut()
.find(|turn| turn.id == payload.turn_id)
{
if let Some(error) = terminal_error.as_ref() {
turn.status = TurnStatus::Failed;
turn.error = Some(error.clone());
} else if matches!(turn.status, TurnStatus::Completed | TurnStatus::InProgress) {
turn.status = TurnStatus::Completed;
}
turn.completed_at = payload.completed_at;
turn.duration_ms = payload.duration_ms;
let changed_turn = ThreadHistoryTurnChange::from_pending_turn(turn);
self.record_changed_turn(changed_turn);
return;
}
// If the completion event cannot be matched, apply it to the active turn.
if let Some(current_turn) = self.current_turn.as_mut() {
let changed_turn = apply_completion(current_turn);
self.record_changed_turn(changed_turn);
self.finish_current_turn();
}
}
/// Marks the current turn as containing a persisted compaction marker.
///
/// This keeps compaction-only legacy turns from being dropped by
/// `finish_current_turn` when they have no renderable items and were not
/// explicitly opened.
fn handle_compacted(&mut self, _payload: &CompactedItem) {
self.ensure_turn().saw_compaction = true;
}
fn handle_thread_rollback(&mut self, payload: &ThreadRolledBackEvent) {
self.finish_current_turn();
let n = usize::try_from(payload.num_turns).unwrap_or(usize::MAX);
let removed_turn_ids = if n >= self.turns.len() {
self.turns.iter().map(|turn| turn.id.clone()).collect()
} else if n == 0 {
Vec::new()
} else {
self.turns[self.turns.len() - n..]
.iter()
.map(|turn| turn.id.clone())
.collect()
};
self.record_removed_turn_ids(removed_turn_ids);
if n >= self.turns.len() {
self.turns.clear();
} else {
self.turns.truncate(self.turns.len().saturating_sub(n));
}
let item_count: usize = self.turns.iter().map(|t| t.items.len()).sum();
self.next_item_index = i64::try_from(item_count.saturating_add(1)).unwrap_or(i64::MAX);
}
fn finish_current_turn(&mut self) {
if let Some(turn) = self.current_turn.take() {
if turn.items.is_empty() && !turn.opened_explicitly && !turn.saw_compaction {
return;
}
self.turns.push(turn);
}
}
fn new_turn(&mut self, id: Option<String>) -> PendingTurn {
let id = id.unwrap_or_else(|| {
if self.next_rollout_index == 0 {
Uuid::now_v7().to_string()
} else {
format!("rollout-{}", self.current_rollout_index)
}
});
PendingTurn {
id,
root_turn_id: None,
items: Vec::new(),
item_index: TurnItemIndex::default(),
error: None,
status: TurnStatus::Completed,
started_at: None,
completed_at: None,
duration_ms: None,
opened_explicitly: false,
saw_compaction: false,
rollout_start_index: self.current_rollout_index,
}
}
fn ensure_turn(&mut self) -> &mut PendingTurn {
if self.current_turn.is_none() {
let turn = self.new_turn(/*id*/ None);
self.record_changed_pending_turn(&turn);
self.current_turn = Some(turn);
}
if let Some(turn) = self.current_turn.as_mut() {
return turn;
}
unreachable!("current turn must exist after initialization");
}
fn push_item_in_current_turn(&mut self, item: ThreadItem) {
let tracking_changes = self.is_tracking_changes();
let changed_item = {
let turn = self.ensure_turn();
let changed_item = tracking_changes.then(|| (turn.id.clone(), item.clone()));
turn.item_index.push(&mut turn.items, item);
changed_item
};
if let Some((turn_id, item)) = changed_item {
self.record_changed_item(turn_id, item);
}
}
fn upsert_item_in_turn_id(&mut self, turn_id: &str, item: ThreadItem) {
let tracking_changes = self.is_tracking_changes();
if let Some(turn) = self.current_turn.as_mut()
&& turn.id == turn_id
{
let changed_item = {
let item = turn.item_index.upsert(&mut turn.items, item);
tracking_changes.then(|| (turn.id.clone(), item.clone()))
};
if let Some((turn_id, item)) = changed_item {
self.record_changed_item(turn_id, item);
}
return;
}
if let Some(turn) = self.turns.iter_mut().find(|turn| turn.id == turn_id) {
let changed_item = {
let item = turn.item_index.upsert(&mut turn.items, item);
tracking_changes.then(|| (turn.id.clone(), item.clone()))
};
if let Some((turn_id, item)) = changed_item {
self.record_changed_item(turn_id, item);
}
return;
}
warn!(
item_id = item.id(),
"dropping turn-scoped item for unknown turn id `{turn_id}`"
);
}
fn upsert_item_in_current_turn(&mut self, item: ThreadItem) {
let tracking_changes = self.is_tracking_changes();
let changed_item = {
let turn = self.ensure_turn();
let item = turn.item_index.upsert(&mut turn.items, item);
tracking_changes.then(|| (turn.id.clone(), item.clone()))
};
if let Some((turn_id, item)) = changed_item {
self.record_changed_item(turn_id, item);
}
}
fn is_tracking_changes(&self) -> bool {
self.active_change_set.is_some()
}
fn record_changed_item(&mut self, turn_id: String, item: ThreadItem) {
if let Some(change_set) = self.active_change_set.as_mut() {
change_set.changed_items.push(ThreadHistoryItemChange {
turn_id,
item,
// Legacy events used by ThreadHistoryBuilder don't have timestamps
started_at_ms: None,
completed_at_ms: None,
});
}
}
fn record_changed_pending_turn(&mut self, turn: &PendingTurn) {
if self.is_tracking_changes() {
self.record_changed_turn(ThreadHistoryTurnChange::from_pending_turn(turn));
}
}
fn record_changed_turn(&mut self, turn: ThreadHistoryTurnChange) {
if let Some(change_set) = self.active_change_set.as_mut() {
change_set.changed_turns.push(turn);
}
}
fn record_removed_turn_ids(&mut self, removed_turn_ids: Vec<String>) {
if let Some(change_set) = self.active_change_set.as_mut() {
change_set.removed_turn_ids.extend(removed_turn_ids);
}
}
fn next_item_id(&mut self) -> String {
let id = format!("item-{}", self.next_item_index);
self.next_item_index += 1;
id
}
fn build_user_inputs(&self, payload: &UserMessageEvent) -> Vec<UserInput> {
let mut content = Vec::new();
if !payload.message.trim().is_empty() {
content.push(UserInput::Text {
text: payload.message.clone(),
text_elements: payload
.text_elements
.iter()
.cloned()
.map(Into::into)
.collect(),
});
}
let has_complete_image_order = payload.has_complete_image_order();
if has_complete_image_order {
let mut inline_index = 0;
let mut file_index = 0;
for image_kind in &payload.image_order {
let (image, detail) = match image_kind {
UserMessageImageKind::Inline => {
let Some(image) = payload
.images
.as_deref()
.and_then(|images| images.get(inline_index))
else {
continue;
};
let detail = payload.image_details.get(inline_index).copied().flatten();
inline_index += 1;
(ImageReference::Inline { url: image.clone() }, detail)
}
UserMessageImageKind::File => {
let Some(file_id) = payload
.file_ids
.as_deref()
.and_then(|file_ids| file_ids.get(file_index))
else {
continue;
};
let detail = payload.file_id_details.get(file_index).copied().flatten();
file_index += 1;
(
ImageReference::File {
file_id: file_id.clone(),
},
detail,
)
}
};
content.push(UserInput::Image { image, detail });
}
} else if let Some(images) = &payload.images {
for (idx, image) in images.iter().enumerate() {
content.push(UserInput::Image {
image: ImageReference::Inline { url: image.clone() },
detail: payload.image_details.get(idx).copied().flatten(),
});
}
}
if !has_complete_image_order && let Some(file_ids) = &payload.file_ids {
for (idx, file_id) in file_ids.iter().enumerate() {
content.push(UserInput::Image {
image: ImageReference::File {
file_id: file_id.clone(),
},
detail: payload.file_id_details.get(idx).copied().flatten(),
});
}
}
for (idx, path) in payload.local_images.iter().enumerate() {
content.push(UserInput::LocalImage {
path: path.clone(),
detail: payload.local_image_details.get(idx).copied().flatten(),
});
}
if let Some(audio) = &payload.audio {
content.extend(audio.iter().cloned().map(|url| UserInput::Audio { url }));
}
content.extend(
payload
.local_audio
.iter()
.cloned()
.map(|path| UserInput::LocalAudio { path }),
);
content
}
}
fn convert_dynamic_tool_content_items(
items: &[codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem],
) -> Vec<DynamicToolCallOutputContentItem> {
items
.iter()
.cloned()
.map(|item| match item {
codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem::InputText { text } => {
DynamicToolCallOutputContentItem::InputText { text }
}
codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem::InputImage {
image_url,
} => DynamicToolCallOutputContentItem::InputImage { image_url },
codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem::InputAudio {
audio_url,
} => DynamicToolCallOutputContentItem::InputAudio { audio_url },
})
.collect()
}
const TURN_ITEM_INDEX_THRESHOLD: usize = 32;
/// Lazily indexes a turn's append-only item list. Replacements keep the same ID
/// and position; duplicate IDs continue to resolve to their first occurrence.
/// Keep this index with its turn, including after completion and across rollback.
#[derive(Default)]
struct TurnItemIndex {
positions: Option<HashMap<String, usize>>,
}
impl TurnItemIndex {
fn push(&mut self, items: &mut Vec<ThreadItem>, item: ThreadItem) {
if let Some(positions) = &mut self.positions {
positions
.entry(item.id().to_string())
.or_insert(items.len());
}
items.push(item);
}
fn upsert<'a>(&mut self, items: &'a mut Vec<ThreadItem>, item: ThreadItem) -> &'a ThreadItem {
if self.positions.is_none() && items.len() >= TURN_ITEM_INDEX_THRESHOLD {
let mut positions = HashMap::with_capacity(items.len());
for (index, existing) in items.iter().enumerate() {
positions.entry(existing.id().to_string()).or_insert(index);
}
self.positions = Some(positions);
}
let existing_index = match &self.positions {
Some(positions) => positions.get(item.id()).copied(),
None => items.iter().position(|existing| existing.id() == item.id()),
};
if let Some(index) = existing_index {
items[index] = item;
&items[index]
} else {
let index = items.len();
self.push(items, item);
&items[index]
}
}
}
struct PendingTurn {
id: String,
root_turn_id: Option<String>,
items: Vec<ThreadItem>,
item_index: TurnItemIndex,
error: Option<TurnError>,
status: TurnStatus,
started_at: Option<i64>,
completed_at: Option<i64>,
duration_ms: Option<i64>,
/// True when this turn originated from an explicit `turn_started`/`turn_complete`
/// boundary, so we preserve it even if it has no renderable items.
opened_explicitly: bool,
/// True when this turn includes a persisted `RolloutItem::Compacted`, which
/// should keep the turn from being dropped even without normal items.
saw_compaction: bool,
/// Index of the rollout item that opened this turn during replay.
rollout_start_index: usize,
}
impl PendingTurn {
fn opened_explicitly(mut self) -> Self {
self.opened_explicitly = true;
self
}
fn with_status(mut self, status: TurnStatus) -> Self {
self.status = status;
self
}
fn with_started_at(mut self, started_at: Option<i64>) -> Self {
self.started_at = started_at;
self
}
}
impl From<PendingTurn> for Turn {
fn from(value: PendingTurn) -> Self {
Self {
id: value.id,
items: value.items,
items_view: TurnItemsView::Full,
error: value.error,
status: value.status,
started_at: value.started_at,
completed_at: value.completed_at,
duration_ms: value.duration_ms,
}
}
}
impl From<&PendingTurn> for Turn {
fn from(value: &PendingTurn) -> Self {
Self {
id: value.id.clone(),
items: value.items.clone(),
items_view: TurnItemsView::Full,
error: value.error.clone(),
status: value.status.clone(),
started_at: value.started_at,
completed_at: value.completed_at,
duration_ms: value.duration_ms,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::protocol::v2::AgentMessageDelivery;
use crate::protocol::v2::AsyncUserInputQuestion;
use crate::protocol::v2::CommandExecutionSource;
use codex_extension_items::ExtensionItem as CoreExtensionItem;
use codex_extension_items::sleep::SleepItem as CoreSleepItem;
use codex_protocol::ThreadId;
use codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem as CoreDynamicToolCallOutputContentItem;
use codex_protocol::items::CommandExecutionItem as CoreCommandExecutionItem;
use codex_protocol::items::CommandExecutionStatus as CoreCommandExecutionStatus;
use codex_protocol::items::EnteredReviewModeItem as CoreEnteredReviewModeItem;
use codex_protocol::items::ExitedReviewModeItem as CoreExitedReviewModeItem;
use codex_protocol::items::HookPromptFragment as CoreHookPromptFragment;
use codex_protocol::items::SubAgentActivityItem as CoreSubAgentActivityItem;
use codex_protocol::items::TurnItem as CoreTurnItem;
use codex_protocol::items::UserMessageItem as CoreUserMessageItem;
use codex_protocol::items::build_hook_prompt_message;
use codex_protocol::mcp::CallToolResult;
use codex_protocol::models::ImageDetail;
use codex_protocol::models::MessagePhase as CoreMessagePhase;
use codex_protocol::models::WebSearchAction as CoreWebSearchAction;
use codex_protocol::parse_command::ParsedCommand;
use codex_protocol::protocol::AgentReasoningEvent;
use codex_protocol::protocol::AgentReasoningRawContentEvent;
use codex_protocol::protocol::ApplyPatchApprovalRequestEvent;
use codex_protocol::protocol::CodexErrorInfo;
use codex_protocol::protocol::DynamicToolCallResponseEvent;
use codex_protocol::protocol::EnteredReviewModeEvent;
use codex_protocol::protocol::ExecCommandBeginEvent;
use codex_protocol::protocol::ExecCommandEndEvent;
use codex_protocol::protocol::ExecCommandSource;
use codex_protocol::protocol::ExitedReviewModeEvent;
use codex_protocol::protocol::ItemStartedEvent;
use codex_protocol::protocol::McpInvocation;
use codex_protocol::protocol::McpToolCallEndEvent;
use codex_protocol::protocol::PatchApplyBeginEvent;
use codex_protocol::protocol::ReviewTarget;
use codex_protocol::protocol::SubAgentActivityKind as CoreSubAgentActivityKind;
use codex_protocol::protocol::ThreadRolledBackEvent;
use codex_protocol::protocol::TurnAbortReason;
use codex_protocol::protocol::TurnAbortedEvent;
use codex_protocol::protocol::TurnCompleteEvent;
use codex_protocol::protocol::TurnStartedEvent;
use codex_protocol::protocol::UserMessageEvent;
use codex_protocol::protocol::WebSearchBeginEvent;
use codex_protocol::protocol::WebSearchEndEvent;
use codex_rollout::CompactedItem;
use codex_utils_absolute_path::test_support::PathBufExt;
use codex_utils_absolute_path::test_support::test_path_buf;
use pretty_assertions::assert_eq;
use std::path::PathBuf;
use std::time::Duration;
use uuid::Uuid;
fn indexed_sleep_item(id: &str, duration_ms: u64) -> ThreadItem {
ThreadItem::Sleep(CoreSleepItem {
id: id.to_string(),
duration_ms,
})
}
#[test]
fn pushed_duplicates_resolve_to_first_occurrence_before_and_after_indexing() {
for count in [2, TURN_ITEM_INDEX_THRESHOLD, TURN_ITEM_INDEX_THRESHOLD * 2] {
let mut index = TurnItemIndex::default();
let mut items = Vec::new();
let original = indexed_sleep_item("duplicate", /*duration_ms*/ 1);
for _ in 0..count {
index.push(&mut items, original.clone());
}
let updated = indexed_sleep_item("duplicate", /*duration_ms*/ 2);
index.upsert(&mut items, updated.clone());
// Appends after index activation must also keep the first position.
index.push(&mut items, original.clone());
index.upsert(&mut items, updated.clone());
let appended = indexed_sleep_item("appended", /*duration_ms*/ 3);
index.push(
&mut items,
indexed_sleep_item("appended", /*duration_ms*/ 1),
);
index.upsert(&mut items, appended.clone());
let mut expected = vec![original; count + 1];
expected[0] = updated;
expected.push(appended);
assert_eq!(items, expected);
}
}
#[test]
fn builds_multiple_turns_with_reasoning_items() {
let events = vec![
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "First turn".into(),
images: Some(vec!["https://example.com/one.png".into()]),
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "Hi there".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
EventMsg::AgentReasoning(AgentReasoningEvent {
text: "thinking".into(),
}),
EventMsg::AgentReasoningRawContent(AgentReasoningRawContentEvent {
text: "full reasoning".into(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "Second turn".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "Reply two".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
];
let mut builder = ThreadHistoryBuilder::new();
for event in &events {
builder.handle_event(event);
}
let turns = builder.finish();
assert_eq!(turns.len(), 2);
let first = &turns[0];
assert!(Uuid::parse_str(&first.id).is_ok());
assert_eq!(first.status, TurnStatus::Completed);
assert_eq!(first.items.len(), 3);
assert_eq!(
first.items[0],
ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![
UserInput::Text {
text: "First turn".into(),
text_elements: Vec::new(),
},
UserInput::Image {
image: ImageReference::Inline {
url: "https://example.com/one.png".into(),
},
detail: None,
}
],
}
);
assert_eq!(
first.items[1],
ThreadItem::AgentMessage {
id: "item-2".into(),
text: "Hi there".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}
);
assert_eq!(
first.items[2],
ThreadItem::Reasoning {
id: "item-3".into(),
summary: vec!["thinking".into()],
content: vec!["full reasoning".into()],
}
);
let second = &turns[1];
assert!(Uuid::parse_str(&second.id).is_ok());
assert_ne!(first.id, second.id);
assert_eq!(second.items.len(), 2);
assert_eq!(
second.items[0],
ThreadItem::UserMessage {
id: "item-4".into(),
client_id: None,
content: vec![UserInput::Text {
text: "Second turn".into(),
text_elements: Vec::new(),
}],
}
);
assert_eq!(
second.items[1],
ThreadItem::AgentMessage {
id: "item-5".into(),
text: "Reply two".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}
);
}
#[test]
fn review_mode_events_replay_persisted_ids() {
let events = vec![
EventMsg::EnteredReviewMode(EnteredReviewModeEvent {
target: ReviewTarget::Custom {
instructions: "review this".into(),
},
user_facing_hint: Some("Review requested.".into()),
turn_id: Some("turn-1".into()),
item_id: Some("entered-review".into()),
}),
EventMsg::ExitedReviewMode(ExitedReviewModeEvent {
turn_id: Some("turn-1".into()),
item_id: Some("exited-review".into()),
review_output: None,
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-1".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let mut builder = ThreadHistoryBuilder::new();
for event in &events {
builder.handle_event(event);
}
let turns = builder.finish();
assert_eq!(turns[0].id, "turn-1");
assert_eq!(
turns[0].items,
vec![
ThreadItem::EnteredReviewMode {
id: "entered-review".into(),
review: "Review requested.".into(),
},
ThreadItem::ExitedReviewMode {
id: "exited-review".into(),
review: REVIEW_FALLBACK_MESSAGE.into(),
},
]
);
}
#[test]
fn review_mode_items_replay_without_turn_started() {
let thread_id = ThreadId::new();
let entered = CoreTurnItem::EnteredReviewMode(CoreEnteredReviewModeItem {
id: "entered-review".into(),
target: ReviewTarget::Custom {
instructions: "review this".into(),
},
user_facing_hint: "Review requested.".into(),
});
let exited = CoreTurnItem::ExitedReviewMode(CoreExitedReviewModeItem {
id: "exited-review".into(),
review_output: None,
});
let events = vec![
EventMsg::ItemCompleted(ItemCompletedEvent {
thread_id,
turn_id: "turn-1".into(),
item: entered,
started_at_ms: Some(0),
completed_at_ms: 0,
}),
EventMsg::ItemCompleted(ItemCompletedEvent {
thread_id,
turn_id: "turn-1".into(),
item: exited,
started_at_ms: Some(0),
completed_at_ms: 0,
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-1".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let mut builder = ThreadHistoryBuilder::new();
for event in &events {
builder.handle_event(event);
}
let turns = builder.finish();
assert_eq!(turns[0].id, "turn-1");
assert_eq!(
turns[0].items,
vec![
ThreadItem::EnteredReviewMode {
id: "entered-review".into(),
review: "Review requested.".into(),
},
ThreadItem::ExitedReviewMode {
id: "exited-review".into(),
review: REVIEW_FALLBACK_MESSAGE.into(),
},
]
);
}
#[test]
fn rebuilds_user_message_attachments_from_legacy_events() {
let local_image_path = PathBuf::from("/tmp/local.png");
let local_audio_path = PathBuf::from("/tmp/local.wav");
let events = vec![RolloutItem::EventMsg(EventMsg::UserMessage(
UserMessageEvent {
client_id: None,
message: "inspect these".into(),
images: Some(vec!["https://example.com/image.png".into()]),
image_details: vec![Some(ImageDetail::Original)],
file_ids: Some(vec!["file_123".into()]),
file_id_details: vec![Some(ImageDetail::High)],
image_order: vec![UserMessageImageKind::File, UserMessageImageKind::Inline],
local_images: vec![local_image_path.clone()],
local_image_details: vec![Some(ImageDetail::Original)],
audio: Some(vec!["https://example.com/audio.mp3".into()]),
local_audio: vec![local_audio_path.clone()],
text_elements: Vec::new(),
},
))];
let turns = build_turns_from_rollout_items(&events);
assert_eq!(turns.len(), 1);
assert_eq!(
turns[0].items[0],
ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![
UserInput::Text {
text: "inspect these".into(),
text_elements: Vec::new(),
},
UserInput::Image {
image: ImageReference::File {
file_id: "file_123".into(),
},
detail: Some(ImageDetail::High),
},
UserInput::Image {
image: ImageReference::Inline {
url: "https://example.com/image.png".into(),
},
detail: Some(ImageDetail::Original),
},
UserInput::LocalImage {
path: local_image_path,
detail: Some(ImageDetail::Original),
},
UserInput::Audio {
url: "https://example.com/audio.mp3".into(),
},
UserInput::LocalAudio {
path: local_audio_path,
},
],
}
);
}
/// Incomplete ordering metadata must not drop references from the legacy split arrays.
#[test]
fn incomplete_image_order_falls_back_to_legacy_order() {
let events = vec![RolloutItem::EventMsg(EventMsg::UserMessage(
UserMessageEvent {
images: Some(vec!["https://example.com/image.png".into()]),
file_ids: Some(vec!["file_123".into()]),
image_order: vec![UserMessageImageKind::File],
..Default::default()
},
))];
let turns = build_turns_from_rollout_items(&events);
assert_eq!(
turns[0].items[0],
ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![
UserInput::Image {
image: ImageReference::Inline {
url: "https://example.com/image.png".into(),
},
detail: None,
},
UserInput::Image {
image: ImageReference::File {
file_id: "file_123".into(),
},
detail: None,
},
],
}
);
}
#[test]
fn ignores_user_message_item_lifecycle_events() {
let turn_id = "turn-1";
let thread_id = ThreadId::new();
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: turn_id.to_string(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "hello".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::ItemStarted(ItemStartedEvent {
thread_id,
turn_id: turn_id.to_string(),
item: CoreTurnItem::UserMessage(CoreUserMessageItem {
id: "user-item-id".to_string(),
client_id: None,
content: Vec::new(),
}),
started_at_ms: 0,
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: turn_id.to_string(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].items.len(), 1);
assert_eq!(
turns[0].items[0],
ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![UserInput::Text {
text: "hello".into(),
text_elements: Vec::new(),
}],
}
);
}
#[test]
fn rebuilds_sleep_item_from_persisted_completion() {
let turn_id = "turn-1";
let thread_id = ThreadId::new();
let sleep_item = CoreTurnItem::Extension(CoreExtensionItem::Sleep(CoreSleepItem {
id: "sleep-1".to_string(),
duration_ms: 1_000,
}));
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: turn_id.to_string(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::ItemCompleted(ItemCompletedEvent {
thread_id,
turn_id: turn_id.to_string(),
item: sleep_item,
started_at_ms: Some(0),
completed_at_ms: 1_000,
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: turn_id.to_string(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(
turns[0].items,
vec![ThreadItem::Sleep(CoreSleepItem {
id: "sleep-1".to_string(),
duration_ms: 1_000,
})]
);
}
#[test]
fn indexed_items_stay_with_their_turn_after_rollback() {
let mut builder = ThreadHistoryBuilder::new();
let mut expected: Vec<_> = (0..64)
.map(|i| indexed_sleep_item(&i.to_string(), /*duration_ms*/ 1))
.collect();
for turn_id in ["retained", "discarded"] {
builder.ensure_turn().id = turn_id.into();
for item in &expected {
builder.upsert_item_in_current_turn(item.clone());
}
builder.finish_current_turn();
// Reuse IDs at different positions in the other turn.
expected.reverse();
}
builder.handle_thread_rollback(&ThreadRolledBackEvent { num_turns: 1 });
expected[0] = indexed_sleep_item("0", /*duration_ms*/ 2);
builder.upsert_item_in_turn_id("retained", expected[0].clone());
assert_eq!(builder.finish()[0].items, expected);
}
#[test]
fn rebuilds_extension_image_generation_item_from_persisted_completion() {
let turn_id = "turn-1";
let thread_id = ThreadId::new();
let saved_path = test_path_buf("/tmp/image-1.png").abs();
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: turn_id.to_string(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::ItemCompleted(ItemCompletedEvent {
thread_id,
turn_id: turn_id.to_string(),
item: CoreTurnItem::Extension(CoreExtensionItem::ImageGeneration(
ImageGenerationItem {
id: "image-1".to_string(),
status: "completed".to_string(),
revised_prompt: Some("A blue square".to_string()),
result: "cG5n".to_string(),
transparent_background: Some(true),
failure: None,
saved_path: Some(saved_path.clone()),
imagegen_request_id: None,
generation_id: None,
},
)),
started_at_ms: Some(0),
completed_at_ms: 1_000,
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: turn_id.to_string(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(
turns[0].items,
vec![ThreadItem::ImageGeneration(ImageGenerationItem {
id: "image-1".to_string(),
status: "completed".to_string(),
revised_prompt: Some("A blue square".to_string()),
result: "cG5n".to_string(),
transparent_background: Some(true),
failure: None,
saved_path: Some(saved_path),
imagegen_request_id: None,
generation_id: None,
})]
);
}
#[test]
fn preserves_command_plugin_id_and_redacts_secrets_across_legacy_upsert() {
let turn_id = "turn-1";
let thread_id = ThreadId::new();
let command = vec![
"git".to_string(),
"-c".to_string(),
"http.extraHeader=Authorization: Bearer example_synthetic_bearer_token_123456"
.to_string(),
"push".to_string(),
];
let parsed_cmd = vec![ParsedCommand::Unknown {
cmd: "git -c 'http.extraHeader=Authorization: Bearer example_synthetic_bearer_token_123456' push"
.to_string(),
}];
let command_item = CoreTurnItem::CommandExecution(CoreCommandExecutionItem {
model_context: None,
id: "exec-1".to_string(),
plugin_id: Some("sample@openai-curated".to_string()),
script_path: Some("scripts/run.py".to_string()),
process_id: Some("pid-1".to_string()),
command: command.clone(),
cwd: test_path_buf("/tmp").abs().into(),
parsed_cmd: parsed_cmd.clone(),
source: ExecCommandSource::Agent,
interaction_input: None,
status: CoreCommandExecutionStatus::Completed,
stdout: Some("hello world\n".to_string()),
stderr: Some(String::new()),
aggregated_output: Some("hello world\n".to_string()),
exit_code: Some(0),
duration: Some(Duration::from_millis(12)),
formatted_output: Some("hello world\n".to_string()),
});
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: turn_id.to_string(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::ExecCommandBegin(ExecCommandBeginEvent {
call_id: "exec-1".to_string(),
plugin_id: Some("sample@openai-curated".to_string()),
script_path: Some("scripts/run.py".to_string()),
process_id: Some("pid-1".to_string()),
turn_id: turn_id.to_string(),
started_at_ms: 0,
command: command.clone(),
cwd: test_path_buf("/tmp").abs().into(),
parsed_cmd: parsed_cmd.clone(),
source: ExecCommandSource::Agent,
interaction_input: None,
}),
EventMsg::ItemCompleted(ItemCompletedEvent {
thread_id,
turn_id: turn_id.to_string(),
item: command_item,
started_at_ms: Some(0),
completed_at_ms: 1_000,
}),
EventMsg::ExecCommandEnd(ExecCommandEndEvent {
call_id: "exec-1".to_string(),
plugin_id: Some("sample@openai-curated".to_string()),
script_path: Some("scripts/run.py".to_string()),
process_id: Some("pid-1".to_string()),
turn_id: turn_id.to_string(),
completed_at_ms: 1_000,
command,
cwd: test_path_buf("/tmp").abs().into(),
parsed_cmd,
source: ExecCommandSource::Agent,
interaction_input: None,
stdout: "hello world\n".to_string(),
stderr: String::new(),
aggregated_output: "hello world\n".to_string(),
exit_code: 0,
duration: Duration::from_millis(12),
formatted_output: "hello world\n".to_string(),
status: CoreExecCommandStatus::Completed,
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: turn_id.to_string(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
assert_eq!(
build_turns_from_rollout_items(&items[..2])[0].items,
vec![ThreadItem::CommandExecution {
model_context: None,
id: "exec-1".to_string(),
plugin_id: Some("sample@openai-curated".to_string()),
script_path: Some("scripts/run.py".to_string()),
command: "git -c 'http.extraHeader=Authorization: Bearer [REDACTED_SECRET]' push"
.to_string(),
cwd: test_path_buf("/tmp").abs().into(),
process_id: Some("pid-1".to_string()),
source: CommandExecutionSource::Agent,
status: CommandExecutionStatus::InProgress,
command_actions: vec![CommandAction::Unknown {
command:
"git -c 'http.extraHeader=Authorization: Bearer [REDACTED_SECRET]' push"
.to_string(),
}],
aggregated_output: None,
exit_code: None,
duration_ms: None,
}]
);
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(
turns[0].items,
vec![ThreadItem::CommandExecution {
model_context: None,
id: "exec-1".to_string(),
plugin_id: Some("sample@openai-curated".to_string()),
script_path: Some("scripts/run.py".to_string()),
command: "git -c 'http.extraHeader=Authorization: Bearer [REDACTED_SECRET]' push"
.to_string(),
cwd: test_path_buf("/tmp").abs().into(),
process_id: Some("pid-1".to_string()),
source: CommandExecutionSource::Agent,
status: CommandExecutionStatus::Completed,
command_actions: vec![CommandAction::Unknown {
command:
"git -c 'http.extraHeader=Authorization: Bearer [REDACTED_SECRET]' push"
.to_string(),
}],
aggregated_output: Some("hello world\n".to_string()),
exit_code: Some(0),
duration_ms: Some(12),
}]
);
}
#[test]
fn preserves_user_message_client_id_from_legacy_event() {
let turn_id = "turn-1";
let thread_id = ThreadId::new();
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: turn_id.to_string(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::ItemStarted(ItemStartedEvent {
thread_id,
turn_id: turn_id.to_string(),
item: CoreTurnItem::UserMessage(CoreUserMessageItem {
id: "user-item-id".to_string(),
client_id: Some("client-message-1".to_string()),
content: vec![codex_protocol::user_input::UserInput::Text {
text: "hello".into(),
text_elements: Vec::new(),
}],
}),
started_at_ms: 0,
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: Some("client-message-1".to_string()),
message: "hello".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: turn_id.to_string(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(
turns[0].items,
vec![ThreadItem::UserMessage {
id: "item-1".into(),
client_id: Some("client-message-1".to_string()),
content: vec![UserInput::Text {
text: "hello".into(),
text_elements: Vec::new(),
}],
}]
);
}
#[test]
fn preserves_agent_message_phase_delivery_and_questions_in_history() {
let questions = vec![
AsyncUserInputQuestion {
title: "Which environment?".into(),
options: Some(vec!["Staging".into(), "Production".into()]),
},
AsyncUserInputQuestion {
title: "Anything else?".into(),
options: None,
},
];
let events = vec![EventMsg::AgentMessage(AgentMessageEvent {
message: "Final reply".into(),
phase: Some(CoreMessagePhase::FinalAnswer),
memory_citation: None,
delivery: Some(AgentMessageDelivery::Async),
questions: Some(questions.clone()),
})];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(
turns[0].items[0],
ThreadItem::AgentMessage {
id: "item-1".into(),
text: "Final reply".into(),
phase: Some(CoreMessagePhase::FinalAnswer),
memory_citation: None,
delivery: Some(AgentMessageDelivery::Async),
questions: Some(questions),
}
);
}
#[test]
fn replays_image_generation_end_events_into_turn_history() {
let items = vec![
RolloutItem::EventMsg(EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-image".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
})),
RolloutItem::EventMsg(EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "generate an image".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
})),
RolloutItem::EventMsg(EventMsg::ImageGenerationEnd(ImageGenerationEndEvent {
call_id: "ig_123".into(),
status: "completed".into(),
revised_prompt: Some("final prompt".into()),
result: "Zm9v".into(),
transparent_background: Some(true),
failure: Some(
codex_extension_items::image_generation::ImageGenerationFailure::UsageLimitExceeded {
limit_id: "image_gen".into(),
resets_at: Some(1_786_150_800),
},
),
saved_path: Some(test_path_buf("/tmp/ig_123.png").abs()),
})),
RolloutItem::EventMsg(EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-image".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
})),
];
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(
turns[0],
Turn {
id: "turn-image".into(),
status: TurnStatus::Completed,
error: None,
started_at: None,
completed_at: None,
duration_ms: None,
items_view: TurnItemsView::Full,
items: vec![
ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![UserInput::Text {
text: "generate an image".into(),
text_elements: Vec::new(),
}],
},
ThreadItem::ImageGeneration(ImageGenerationItem {
id: "ig_123".into(),
status: "completed".into(),
revised_prompt: Some("final prompt".into()),
result: "Zm9v".into(),
transparent_background: Some(true),
failure: Some(
codex_extension_items::image_generation::ImageGenerationFailure::UsageLimitExceeded {
limit_id: "image_gen".into(),
resets_at: Some(1_786_150_800),
},
),
saved_path: Some(test_path_buf("/tmp/ig_123.png").abs()),
imagegen_request_id: None,
generation_id: None,
}),
],
}
);
}
#[test]
fn splits_reasoning_when_interleaved() {
let events = vec![
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "Turn start".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentReasoning(AgentReasoningEvent {
text: "first summary".into(),
}),
EventMsg::AgentReasoningRawContent(AgentReasoningRawContentEvent {
text: "first content".into(),
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "interlude".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
EventMsg::AgentReasoning(AgentReasoningEvent {
text: "second summary".into(),
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
let turn = &turns[0];
assert_eq!(turn.items.len(), 4);
assert_eq!(
turn.items[1],
ThreadItem::Reasoning {
id: "item-2".into(),
summary: vec!["first summary".into()],
content: vec!["first content".into()],
}
);
assert_eq!(
turn.items[3],
ThreadItem::Reasoning {
id: "item-4".into(),
summary: vec!["second summary".into()],
content: Vec::new(),
}
);
}
#[test]
fn marks_turn_as_interrupted_when_aborted() {
let events = vec![
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "Please do the thing".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "Working...".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
EventMsg::TurnAborted(TurnAbortedEvent {
turn_id: Some("turn-1".into()),
started_at: None,
reason: TurnAbortReason::Replaced,
completed_at: None,
duration_ms: None,
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "Let's try again".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "Second attempt complete.".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 2);
let first_turn = &turns[0];
assert_eq!(first_turn.status, TurnStatus::Interrupted);
assert_eq!(first_turn.items.len(), 2);
assert_eq!(
first_turn.items[0],
ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![UserInput::Text {
text: "Please do the thing".into(),
text_elements: Vec::new(),
}],
}
);
assert_eq!(
first_turn.items[1],
ThreadItem::AgentMessage {
id: "item-2".into(),
text: "Working...".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}
);
let second_turn = &turns[1];
assert_eq!(second_turn.status, TurnStatus::Completed);
assert_eq!(second_turn.items.len(), 2);
assert_eq!(
second_turn.items[0],
ThreadItem::UserMessage {
id: "item-3".into(),
client_id: None,
content: vec![UserInput::Text {
text: "Let's try again".into(),
text_elements: Vec::new(),
}],
}
);
assert_eq!(
second_turn.items[1],
ThreadItem::AgentMessage {
id: "item-4".into(),
text: "Second attempt complete.".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}
);
}
#[test]
fn drops_last_turns_on_thread_rollback() {
let events = vec![
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "First".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "A1".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "Second".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "A2".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
EventMsg::ThreadRolledBack(ThreadRolledBackEvent { num_turns: 1 }),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "Third".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "A3".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 2);
assert_eq!(turns[0].id, "rollout-0");
assert_eq!(turns[1].id, "rollout-5");
assert_ne!(turns[0].id, turns[1].id);
assert_eq!(turns[0].status, TurnStatus::Completed);
assert_eq!(turns[1].status, TurnStatus::Completed);
assert_eq!(
turns[0].items,
vec![
ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![UserInput::Text {
text: "First".into(),
text_elements: Vec::new(),
}],
},
ThreadItem::AgentMessage {
id: "item-2".into(),
text: "A1".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
},
]
);
assert_eq!(
turns[1].items,
vec![
ThreadItem::UserMessage {
id: "item-3".into(),
client_id: None,
content: vec![UserInput::Text {
text: "Third".into(),
text_elements: Vec::new(),
}],
},
ThreadItem::AgentMessage {
id: "item-4".into(),
text: "A3".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
},
]
);
}
#[test]
fn thread_rollback_clears_all_turns_when_num_turns_exceeds_history() {
let events = vec![
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "One".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "A1".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "Two".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "A2".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
EventMsg::ThreadRolledBack(ThreadRolledBackEvent { num_turns: 99 }),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns, Vec::<Turn>::new());
}
#[test]
fn uses_explicit_turn_boundaries_for_mid_turn_steering() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "Start".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "Steer".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].id, "turn-a");
assert_eq!(
turns[0].items,
vec![
ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![UserInput::Text {
text: "Start".into(),
text_elements: Vec::new(),
}],
},
ThreadItem::UserMessage {
id: "item-2".into(),
client_id: None,
content: vec![UserInput::Text {
text: "Steer".into(),
text_elements: Vec::new(),
}],
},
]
);
}
#[test]
fn reconstructs_tool_items_from_persisted_completion_events() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-1".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "run tools".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::WebSearchEnd(WebSearchEndEvent {
call_id: "search-1".into(),
query: "codex".into(),
action: CoreWebSearchAction::Search {
query: Some("codex".into()),
queries: None,
},
results: Some(vec![serde_json::json!({
"type": "text_result",
"ref_id": "turn0search0",
"url": "https://example.com/codex",
})]),
}),
EventMsg::ExecCommandEnd(ExecCommandEndEvent {
call_id: "exec-1".into(),
plugin_id: None,
script_path: None,
process_id: Some("pid-1".into()),
turn_id: "turn-1".into(),
completed_at_ms: 0,
command: vec!["echo".into(), "hello world".into()],
cwd: test_path_buf("/tmp").abs().into(),
parsed_cmd: vec![ParsedCommand::Unknown {
cmd: "echo hello world".into(),
}],
source: ExecCommandSource::Agent,
interaction_input: None,
stdout: String::new(),
stderr: String::new(),
aggregated_output: "hello world\n".into(),
exit_code: 0,
duration: Duration::from_millis(12),
formatted_output: String::new(),
status: CoreExecCommandStatus::Completed,
}),
EventMsg::McpToolCallEnd(McpToolCallEndEvent {
turn_id: String::new(),
call_id: "mcp-1".into(),
invocation: McpInvocation {
server: "docs".into(),
tool: "lookup".into(),
arguments: Some(serde_json::json!({"id":"123"})),
},
connector_id: None,
mcp_app_resource_uri: None,
mcp_app_ui: None,
link_id: None,
app_name: None,
action_name: None,
plugin_id: None,
read_only_hint: None,
duration: Duration::from_millis(8),
result: Err("boom".into()),
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].items.len(), 4);
assert_eq!(
turns[0].items[1],
ThreadItem::WebSearch(WebSearchItem {
id: "search-1".into(),
query: "codex".into(),
action: Some(WebSearchAction::Search {
query: Some("codex".into()),
queries: None,
}),
results: Some(vec![serde_json::json!({
"type": "text_result",
"ref_id": "turn0search0",
"url": "https://example.com/codex",
})]),
})
);
assert_eq!(
turns[0].items[2],
ThreadItem::CommandExecution {
model_context: None,
id: "exec-1".into(),
plugin_id: None,
script_path: None,
command: "echo 'hello world'".into(),
cwd: test_path_buf("/tmp").abs().into(),
process_id: Some("pid-1".into()),
source: CommandExecutionSource::Agent,
status: CommandExecutionStatus::Completed,
command_actions: vec![CommandAction::Unknown {
command: "echo hello world".into(),
}],
aggregated_output: Some("hello world\n".into()),
exit_code: Some(0),
duration_ms: Some(12),
}
);
assert_eq!(
turns[0].items[3],
ThreadItem::McpToolCall {
id: "mcp-1".into(),
server: "docs".into(),
tool: "lookup".into(),
status: McpToolCallStatus::Failed,
arguments: serde_json::json!({"id":"123"}),
app_context: None,
mcp_app_resource_uri: None,
mcp_app_ui: None,
plugin_id: None,
read_only_hint: None,
result: None,
error: Some(McpToolCallError {
message: "boom".into(),
}),
duration_ms: Some(8),
}
);
}
#[test]
fn reconstructs_mcp_tool_result_meta_from_persisted_completion_events() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-1".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::McpToolCallEnd(McpToolCallEndEvent {
turn_id: String::new(),
call_id: "mcp-1".into(),
invocation: McpInvocation {
server: "docs".into(),
tool: "lookup".into(),
arguments: Some(serde_json::json!({"id":"123"})),
},
connector_id: Some("calendar".into()),
mcp_app_resource_uri: Some("ui://widget/lookup.html".into()),
mcp_app_ui: None,
link_id: Some("link_calendar".into()),
app_name: Some("Calendar".into()),
action_name: Some("lookup".into()),
plugin_id: Some("sample@test".into()),
read_only_hint: Some(false),
duration: Duration::from_millis(8),
result: Ok(CallToolResult {
content: vec![serde_json::json!({
"type": "text",
"text": "result"
})],
structured_content: Some(serde_json::json!({"id":"123"})),
is_error: Some(false),
meta: Some(serde_json::json!({
"ui/resourceUri": "ui://widget/lookup.html"
})),
}),
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(
turns[0].items[0],
ThreadItem::McpToolCall {
id: "mcp-1".into(),
server: "docs".into(),
tool: "lookup".into(),
status: McpToolCallStatus::Completed,
arguments: serde_json::json!({"id":"123"}),
app_context: Some(McpToolCallAppContext {
connector_id: "calendar".into(),
link_id: Some("link_calendar".into()),
resource_uri: Some("ui://widget/lookup.html".into()),
app_name: Some("Calendar".into()),
action_name: Some("lookup".into()),
}),
mcp_app_resource_uri: Some("ui://widget/lookup.html".into()),
mcp_app_ui: None,
plugin_id: Some("sample@test".into()),
read_only_hint: Some(false),
result: Some(Box::new(McpToolCallResult {
content: vec![serde_json::json!({
"type": "text",
"text": "result"
})],
structured_content: Some(serde_json::json!({"id":"123"})),
meta: Some(serde_json::json!({
"ui/resourceUri": "ui://widget/lookup.html"
})),
})),
error: None,
duration_ms: Some(8),
}
);
}
#[test]
fn reconstructs_dynamic_tool_items_from_request_and_response_events() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-1".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "run dynamic tool".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::DynamicToolCallRequest(
codex_protocol::dynamic_tools::DynamicToolCallRequest {
call_id: "dyn-1".into(),
turn_id: "turn-1".into(),
started_at_ms: 0,
namespace: Some("codex_app".into()),
tool: "lookup_ticket".into(),
arguments: serde_json::json!({"id":"ABC-123"}),
},
),
EventMsg::DynamicToolCallResponse(DynamicToolCallResponseEvent {
call_id: "dyn-1".into(),
turn_id: "turn-1".into(),
completed_at_ms: 0,
namespace: Some("codex_app".into()),
tool: "lookup_ticket".into(),
arguments: serde_json::json!({"id":"ABC-123"}),
content_items: vec![
CoreDynamicToolCallOutputContentItem::InputText {
text: "Ticket is open".into(),
},
CoreDynamicToolCallOutputContentItem::InputImage {
image_url: "data:image/png;base64,AAA".into(),
},
CoreDynamicToolCallOutputContentItem::InputAudio {
audio_url: "data:audio/wav;base64,YXVkaW8=".into(),
},
],
success: true,
error: None,
duration: Duration::from_millis(42),
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].items.len(), 2);
assert_eq!(
turns[0].items[1],
ThreadItem::DynamicToolCall {
id: "dyn-1".into(),
namespace: Some("codex_app".into()),
tool: "lookup_ticket".into(),
arguments: serde_json::json!({"id":"ABC-123"}),
status: DynamicToolCallStatus::Completed,
content_items: Some(vec![
DynamicToolCallOutputContentItem::InputText {
text: "Ticket is open".into(),
},
DynamicToolCallOutputContentItem::InputImage {
image_url: "data:image/png;base64,AAA".into(),
},
DynamicToolCallOutputContentItem::InputAudio {
audio_url: "data:audio/wav;base64,YXVkaW8=".into(),
},
]),
success: Some(true),
duration_ms: Some(42),
}
);
}
#[test]
fn reconstructs_declined_exec_and_patch_items() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-1".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "run tools".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::ExecCommandEnd(ExecCommandEndEvent {
call_id: "exec-declined".into(),
plugin_id: None,
script_path: None,
process_id: Some("pid-2".into()),
turn_id: "turn-1".into(),
completed_at_ms: 0,
command: vec!["ls".into()],
cwd: test_path_buf("/tmp").abs().into(),
parsed_cmd: vec![ParsedCommand::Unknown { cmd: "ls".into() }],
source: ExecCommandSource::Agent,
interaction_input: None,
stdout: String::new(),
stderr: "exec command rejected by user".into(),
aggregated_output: "exec command rejected by user".into(),
exit_code: -1,
duration: Duration::ZERO,
formatted_output: String::new(),
status: CoreExecCommandStatus::Declined,
}),
EventMsg::PatchApplyEnd(PatchApplyEndEvent {
call_id: "patch-declined".into(),
turn_id: "turn-1".into(),
stdout: String::new(),
stderr: "patch rejected by user".into(),
success: false,
changes: [(
PathBuf::from("README.md"),
codex_protocol::protocol::FileChange::Add {
content: "hello\n".into(),
},
)]
.into_iter()
.collect(),
status: CorePatchApplyStatus::Declined,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].items.len(), 3);
assert_eq!(
turns[0].items[1],
ThreadItem::CommandExecution {
model_context: None,
id: "exec-declined".into(),
plugin_id: None,
script_path: None,
command: "ls".into(),
cwd: test_path_buf("/tmp").abs().into(),
process_id: Some("pid-2".into()),
source: CommandExecutionSource::Agent,
status: CommandExecutionStatus::Declined,
command_actions: vec![CommandAction::Unknown {
command: "ls".into(),
}],
aggregated_output: Some("exec command rejected by user".into()),
exit_code: Some(-1),
duration_ms: Some(0),
}
);
assert_eq!(
turns[0].items[2],
ThreadItem::FileChange {
id: "patch-declined".into(),
changes: vec![FileUpdateChange {
path: "README.md".into(),
kind: PatchChangeKind::Add,
diff: "hello\n".into(),
}],
status: PatchApplyStatus::Declined,
}
);
}
#[test]
fn reconstructs_declined_guardian_command_item() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-1".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "review this command".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::GuardianAssessment(GuardianAssessmentEvent {
review_reason: None,
model_context: None,
id: "review-guardian-exec".into(),
target_item_id: Some("guardian-exec".into()),
plugin_id: Some("sample@openai-curated".into()),
script_path: Some("scripts/run.py".into()),
turn_id: "turn-1".into(),
started_at_ms: 1_000,
completed_at_ms: None,
status: GuardianAssessmentStatus::InProgress,
risk_level: None,
user_authorization: None,
rationale: None,
decision_source: None,
action: serde_json::from_value(serde_json::json!({
"type": "command",
"source": "shell",
"command": "rm -rf /tmp/guardian",
"cwd": test_path_buf("/tmp"),
}))
.expect("guardian action"),
}),
EventMsg::GuardianAssessment(GuardianAssessmentEvent {
review_reason: None,
model_context: None,
id: "review-guardian-exec".into(),
target_item_id: Some("guardian-exec".into()),
plugin_id: Some("sample@openai-curated".into()),
script_path: Some("scripts/run.py".into()),
turn_id: "turn-1".into(),
started_at_ms: 1_000,
completed_at_ms: Some(1_042),
status: GuardianAssessmentStatus::Denied,
risk_level: Some(codex_protocol::protocol::GuardianRiskLevel::High),
user_authorization: Some(codex_protocol::protocol::GuardianUserAuthorization::Low),
rationale: Some("Would delete user data.".into()),
decision_source: Some(
codex_protocol::protocol::GuardianAssessmentDecisionSource::Agent,
),
action: serde_json::from_value(serde_json::json!({
"type": "command",
"source": "shell",
"command": "rm -rf /tmp/guardian",
"cwd": test_path_buf("/tmp"),
}))
.expect("guardian action"),
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].items.len(), 2);
assert_eq!(
turns[0].items[1],
ThreadItem::CommandExecution {
model_context: None,
id: "guardian-exec".into(),
plugin_id: Some("sample@openai-curated".into()),
script_path: Some("scripts/run.py".into()),
command: "rm -rf /tmp/guardian".into(),
cwd: test_path_buf("/tmp").abs().into(),
process_id: None,
source: CommandExecutionSource::Agent,
status: CommandExecutionStatus::Declined,
command_actions: vec![CommandAction::Unknown {
command: "rm -rf /tmp/guardian".into(),
}],
aggregated_output: None,
exit_code: None,
duration_ms: None,
}
);
}
#[test]
fn reconstructs_in_progress_guardian_execve_item() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-1".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "run a subcommand".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::GuardianAssessment(GuardianAssessmentEvent {
review_reason: None,
model_context: None,
id: "review-guardian-execve".into(),
target_item_id: Some("guardian-execve".into()),
plugin_id: Some("sample@openai-curated".into()),
script_path: Some("scripts/run.py".into()),
turn_id: "turn-1".into(),
started_at_ms: 2_000,
completed_at_ms: None,
status: GuardianAssessmentStatus::InProgress,
risk_level: None,
user_authorization: None,
rationale: None,
decision_source: None,
action: serde_json::from_value(serde_json::json!({
"type": "execve",
"source": "shell",
"program": "/bin/rm",
"argv": ["/usr/bin/rm", "-f", "/tmp/file.sqlite"],
"cwd": test_path_buf("/tmp"),
}))
.expect("guardian action"),
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].items.len(), 2);
assert_eq!(
turns[0].items[1],
ThreadItem::CommandExecution {
model_context: None,
id: "guardian-execve".into(),
plugin_id: Some("sample@openai-curated".into()),
script_path: Some("scripts/run.py".into()),
command: "/bin/rm -f /tmp/file.sqlite".into(),
cwd: test_path_buf("/tmp").abs().into(),
process_id: None,
source: CommandExecutionSource::Agent,
status: CommandExecutionStatus::InProgress,
command_actions: vec![CommandAction::Unknown {
command: "/bin/rm -f /tmp/file.sqlite".into(),
}],
aggregated_output: None,
exit_code: None,
duration_ms: None,
}
);
}
#[test]
fn assigns_late_legacy_mcp_completion_to_original_turn() {
use codex_protocol::items::McpToolCallItem;
use codex_protocol::items::McpToolCallStatus as CoreMcpToolCallStatus;
use codex_protocol::protocol::HasLegacyEvent;
let completion = ItemCompletedEvent {
thread_id: ThreadId::new(),
turn_id: "turn-a".into(),
started_at_ms: None,
completed_at_ms: 0,
item: CoreTurnItem::McpToolCall(McpToolCallItem {
id: "mcp-1".into(),
server: "slack".into(),
tool: "search".into(),
arguments: serde_json::json!({"query": "incident"}),
connector_id: None,
mcp_app_resource_uri: None,
mcp_app_ui: None,
link_id: None,
app_name: None,
action_name: None,
plugin_id: None,
read_only_hint: Some(true),
status: CoreMcpToolCallStatus::Completed,
result: Some(CallToolResult {
content: vec![serde_json::json!({"type": "text", "text": "Found incident"})],
structured_content: None,
is_error: None,
meta: None,
}),
error: None,
duration: Some(Duration::from_millis(8)),
}),
};
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
message: "Find related incident discussions".into(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-b".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
message: "Explain retry logic while that runs".into(),
..Default::default()
}),
];
let mut items: Vec<_> = events.into_iter().map(RolloutItem::EventMsg).collect();
let mut expected = build_turns_from_rollout_items(&items);
expected[0]
.items
.push(ThreadItem::from(completion.item.clone()));
let legacy = completion.as_legacy_events(/*show_raw_agent_reasoning*/ false);
let serialized = serde_json::to_string(&legacy[0]).unwrap();
items.push(RolloutItem::EventMsg(
serde_json::from_str(&serialized).unwrap(),
));
assert_eq!(build_turns_from_rollout_items(&items), expected);
// Old records have no owner, so retain their current-turn fallback.
let mut old_record: serde_json::Value = serde_json::from_str(&serialized).unwrap();
old_record.as_object_mut().unwrap().remove("turn_id");
*items.last_mut().unwrap() =
RolloutItem::EventMsg(serde_json::from_value(old_record).unwrap());
let mut expected_old = expected;
let mcp_item = expected_old[0].items.pop().unwrap();
expected_old[1].items.push(mcp_item);
assert_eq!(build_turns_from_rollout_items(&items), expected_old);
}
#[test]
fn assigns_late_exec_completion_to_original_turn() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "first".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-b".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "second".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::ExecCommandEnd(ExecCommandEndEvent {
call_id: "exec-late".into(),
plugin_id: None,
script_path: None,
process_id: Some("pid-42".into()),
turn_id: "turn-a".into(),
completed_at_ms: 0,
command: vec!["echo".into(), "done".into()],
cwd: test_path_buf("/tmp").abs().into(),
parsed_cmd: vec![ParsedCommand::Unknown {
cmd: "echo done".into(),
}],
source: ExecCommandSource::Agent,
interaction_input: None,
stdout: "done\n".into(),
stderr: String::new(),
aggregated_output: "done\n".into(),
exit_code: 0,
duration: Duration::from_millis(5),
formatted_output: "done\n".into(),
status: CoreExecCommandStatus::Completed,
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-b".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 2);
assert_eq!(turns[0].id, "turn-a");
assert_eq!(turns[1].id, "turn-b");
assert_eq!(turns[0].items.len(), 2);
assert_eq!(turns[1].items.len(), 1);
assert_eq!(
turns[0].items[1],
ThreadItem::CommandExecution {
model_context: None,
id: "exec-late".into(),
plugin_id: None,
script_path: None,
command: "echo done".into(),
cwd: test_path_buf("/tmp").abs().into(),
process_id: Some("pid-42".into()),
source: CommandExecutionSource::Agent,
status: CommandExecutionStatus::Completed,
command_actions: vec![CommandAction::Unknown {
command: "echo done".into(),
}],
aggregated_output: Some("done\n".into()),
exit_code: Some(0),
duration_ms: Some(5),
}
);
}
#[test]
fn drops_late_turn_scoped_item_for_unknown_turn_id() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "first".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-b".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "second".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::ExecCommandEnd(ExecCommandEndEvent {
call_id: "exec-unknown-turn".into(),
plugin_id: None,
script_path: None,
process_id: Some("pid-42".into()),
turn_id: "turn-missing".into(),
completed_at_ms: 0,
command: vec!["echo".into(), "done".into()],
cwd: test_path_buf("/tmp").abs().into(),
parsed_cmd: vec![ParsedCommand::Unknown {
cmd: "echo done".into(),
}],
source: ExecCommandSource::Agent,
interaction_input: None,
stdout: "done\n".into(),
stderr: String::new(),
aggregated_output: "done\n".into(),
exit_code: 0,
duration: Duration::from_millis(5),
formatted_output: "done\n".into(),
status: CoreExecCommandStatus::Completed,
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-b".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let mut builder = ThreadHistoryBuilder::new();
for event in &events {
builder.handle_event(event);
}
let turns = builder.finish();
assert_eq!(turns.len(), 2);
assert_eq!(turns[0].id, "turn-a");
assert_eq!(turns[1].id, "turn-b");
assert_eq!(turns[0].items.len(), 1);
assert_eq!(turns[1].items.len(), 1);
assert_eq!(
turns[1].items[0],
ThreadItem::UserMessage {
id: "item-2".into(),
client_id: None,
content: vec![UserInput::Text {
text: "second".into(),
text_elements: Vec::new(),
}],
}
);
}
#[test]
fn patch_apply_begin_updates_active_turn_snapshot_with_file_change() {
let turn_id = "turn-1";
let mut builder = ThreadHistoryBuilder::new();
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: turn_id.to_string(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "apply patch".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::PatchApplyBegin(PatchApplyBeginEvent {
call_id: "patch-call".into(),
turn_id: turn_id.to_string(),
auto_approved: false,
changes: [(
PathBuf::from("README.md"),
codex_protocol::protocol::FileChange::Add {
content: "hello\n".into(),
},
)]
.into_iter()
.collect(),
}),
];
for event in &events {
builder.handle_event(event);
}
let snapshot = builder
.active_turn_snapshot()
.expect("active turn snapshot");
assert_eq!(snapshot.id, turn_id);
assert_eq!(snapshot.status, TurnStatus::InProgress);
assert_eq!(
snapshot.items,
vec![
ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![UserInput::Text {
text: "apply patch".into(),
text_elements: Vec::new(),
}],
},
ThreadItem::FileChange {
id: "patch-call".into(),
changes: vec![FileUpdateChange {
path: "README.md".into(),
kind: PatchChangeKind::Add,
diff: "hello\n".into(),
}],
status: PatchApplyStatus::InProgress,
},
]
);
}
#[test]
fn apply_patch_approval_request_updates_active_turn_snapshot_with_file_change() {
let turn_id = "turn-1";
let mut builder = ThreadHistoryBuilder::new();
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: turn_id.to_string(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "apply patch".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::ApplyPatchApprovalRequest(ApplyPatchApprovalRequestEvent {
call_id: "patch-call".into(),
turn_id: turn_id.to_string(),
started_at_ms: 0,
changes: [(
PathBuf::from("README.md"),
codex_protocol::protocol::FileChange::Add {
content: "hello\n".into(),
},
)]
.into_iter()
.collect(),
reason: None,
grant_root: None,
}),
];
for event in &events {
builder.handle_event(event);
}
let snapshot = builder
.active_turn_snapshot()
.expect("active turn snapshot");
assert_eq!(snapshot.id, turn_id);
assert_eq!(snapshot.status, TurnStatus::InProgress);
assert_eq!(
snapshot.items,
vec![
ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![UserInput::Text {
text: "apply patch".into(),
text_elements: Vec::new(),
}],
},
ThreadItem::FileChange {
id: "patch-call".into(),
changes: vec![FileUpdateChange {
path: "README.md".into(),
kind: PatchChangeKind::Add,
diff: "hello\n".into(),
}],
status: PatchApplyStatus::InProgress,
},
]
);
}
#[test]
fn late_turn_complete_does_not_close_active_turn() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "first".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-b".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "second".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "still in b".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-b".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 2);
assert_eq!(turns[0].id, "turn-a");
assert_eq!(turns[1].id, "turn-b");
assert_eq!(turns[1].items.len(), 2);
}
#[test]
fn late_turn_complete_with_embedded_error_preserves_active_turn() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: Some(10),
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "first".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-b".into(),
root_turn_id: None,
trace_id: None,
started_at: Some(30),
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "second".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: Some(10),
last_agent_message: None,
error: Some(ErrorEvent {
misalignment: None,
message: "Selected model is at capacity. Please try a different model.".into(),
codex_error_info: Some(CodexErrorInfo::ServerOverloaded),
}),
completed_at: Some(20),
duration_ms: Some(10_000),
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
assert_eq!(
build_turns_from_rollout_items(&items),
vec![
Turn {
id: "turn-a".into(),
items_view: TurnItemsView::Full,
items: vec![ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![UserInput::Text {
text: "first".into(),
text_elements: Vec::new(),
}],
}],
status: TurnStatus::Failed,
error: Some(TurnError {
misalignment: None,
message: "Selected model is at capacity. Please try a different model."
.into(),
codex_error_info: Some(
crate::protocol::v2::CodexErrorInfo::ServerOverloaded,
),
additional_details: None,
}),
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(10_000),
},
Turn {
id: "turn-b".into(),
items_view: TurnItemsView::Full,
items: vec![ThreadItem::UserMessage {
id: "item-2".into(),
client_id: None,
content: vec![UserInput::Text {
text: "second".into(),
text_elements: Vec::new(),
}],
}],
status: TurnStatus::InProgress,
error: None,
started_at: Some(30),
completed_at: None,
duration_ms: None,
},
]
);
}
#[test]
fn late_turn_aborted_does_not_interrupt_active_turn() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "first".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-b".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "second".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnAborted(TurnAbortedEvent {
turn_id: Some("turn-a".into()),
started_at: None,
reason: TurnAbortReason::Replaced,
completed_at: None,
duration_ms: None,
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "still in b".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 2);
assert_eq!(turns[0].id, "turn-a");
assert_eq!(turns[1].id, "turn-b");
assert_eq!(turns[1].status, TurnStatus::InProgress);
assert_eq!(turns[1].items.len(), 2);
}
#[test]
fn preserves_compaction_only_turn() {
let items = vec![
RolloutItem::EventMsg(EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-compact".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
})),
RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: None,
retained_context: None,
guardian_history: None,
mcp_resource_origins: None,
window_number: None,
first_window_id: None,
previous_window_id: None,
window_id: None,
compaction_response_id: None,
latest_token_usage_record: None,
}),
RolloutItem::EventMsg(EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-compact".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
})),
];
let turns = build_turns_from_rollout_items(&items);
assert_eq!(
turns,
vec![Turn {
id: "turn-compact".into(),
status: TurnStatus::Completed,
error: None,
started_at: None,
completed_at: None,
duration_ms: None,
items_view: TurnItemsView::Full,
items: Vec::new(),
}]
);
}
#[test]
fn reconstructs_collab_resume_end_item() {
let events = vec![
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "resume agent".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::CollabResumeEnd(codex_protocol::protocol::CollabResumeEndEvent {
call_id: "resume-1".into(),
completed_at_ms: 0,
sender_thread_id: ThreadId::try_from("00000000-0000-0000-0000-000000000001")
.expect("valid sender thread id"),
receiver_thread_id: ThreadId::try_from("00000000-0000-0000-0000-000000000002")
.expect("valid receiver thread id"),
receiver_agent_nickname: None,
receiver_agent_role: None,
status: AgentStatus::Completed(None),
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].items.len(), 2);
assert_eq!(
turns[0].items[1],
ThreadItem::CollabAgentToolCall {
id: "resume-1".into(),
tool: CollabAgentTool::ResumeAgent,
status: CollabAgentToolCallStatus::Completed,
sender_thread_id: "00000000-0000-0000-0000-000000000001".into(),
receiver_thread_ids: vec!["00000000-0000-0000-0000-000000000002".into()],
prompt: None,
model: None,
reasoning_effort: None,
agents_states: [(
"00000000-0000-0000-0000-000000000002".into(),
CollabAgentState {
status: crate::protocol::v2::CollabAgentStatus::Completed,
message: None,
},
)]
.into_iter()
.collect(),
}
);
}
#[test]
fn reconstructs_collab_spawn_end_item_with_model_metadata() {
let sender_thread_id = ThreadId::try_from("00000000-0000-0000-0000-000000000001")
.expect("valid sender thread id");
let spawned_thread_id = ThreadId::try_from("00000000-0000-0000-0000-000000000002")
.expect("valid receiver thread id");
let events = vec![
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "spawn agent".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::CollabAgentSpawnEnd(codex_protocol::protocol::CollabAgentSpawnEndEvent {
call_id: "spawn-1".into(),
completed_at_ms: 0,
sender_thread_id,
new_thread_id: Some(spawned_thread_id),
new_agent_nickname: Some("Scout".into()),
new_agent_role: Some("explorer".into()),
prompt: "inspect the repo".into(),
model: "gpt-5.4-mini".into(),
reasoning_effort: codex_protocol::openai_models::ReasoningEffort::Medium,
status: AgentStatus::Running,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].items.len(), 2);
assert_eq!(
turns[0].items[1],
ThreadItem::CollabAgentToolCall {
id: "spawn-1".into(),
tool: CollabAgentTool::SpawnAgent,
status: CollabAgentToolCallStatus::Completed,
sender_thread_id: "00000000-0000-0000-0000-000000000001".into(),
receiver_thread_ids: vec!["00000000-0000-0000-0000-000000000002".into()],
prompt: Some("inspect the repo".into()),
model: Some("gpt-5.4-mini".into()),
reasoning_effort: Some(codex_protocol::openai_models::ReasoningEffort::Medium),
agents_states: [(
"00000000-0000-0000-0000-000000000002".into(),
CollabAgentState {
status: crate::protocol::v2::CollabAgentStatus::Running,
message: None,
},
)]
.into_iter()
.collect(),
}
);
}
#[test]
fn reconstructs_interrupted_send_input_as_completed_collab_call() {
// `send_input(interrupt=true)` first stops the child's active turn, then redirects it with
// new input. The transient interrupted status should remain visible in agent state, but the
// collab tool call itself is still a successful redirect rather than a failed operation.
let sender = ThreadId::try_from("00000000-0000-0000-0000-000000000001")
.expect("valid sender thread id");
let receiver = ThreadId::try_from("00000000-0000-0000-0000-000000000002")
.expect("valid receiver thread id");
let events = vec![
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "redirect".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::CollabAgentInteractionBegin(
codex_protocol::protocol::CollabAgentInteractionBeginEvent {
call_id: "send-1".into(),
started_at_ms: 0,
sender_thread_id: sender,
receiver_thread_id: receiver,
prompt: "new task".into(),
},
),
EventMsg::CollabAgentInteractionEnd(
codex_protocol::protocol::CollabAgentInteractionEndEvent {
call_id: "send-1".into(),
completed_at_ms: 0,
sender_thread_id: sender,
receiver_thread_id: receiver,
receiver_agent_nickname: None,
receiver_agent_role: None,
prompt: "new task".into(),
status: AgentStatus::Interrupted,
},
),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].items.len(), 2);
assert_eq!(
turns[0].items[1],
ThreadItem::CollabAgentToolCall {
id: "send-1".into(),
tool: CollabAgentTool::SendInput,
status: CollabAgentToolCallStatus::Completed,
sender_thread_id: sender.to_string(),
receiver_thread_ids: vec![receiver.to_string()],
prompt: Some("new task".into()),
model: None,
reasoning_effort: None,
agents_states: [(
receiver.to_string(),
CollabAgentState {
status: crate::protocol::v2::CollabAgentStatus::Interrupted,
message: None,
},
)]
.into_iter()
.collect(),
}
);
}
#[test]
fn rollback_failed_error_does_not_mark_turn_failed() {
let events = vec![
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "hello".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::AgentMessage(AgentMessageEvent {
message: "done".into(),
phase: None,
memory_citation: None,
delivery: None,
questions: None,
}),
EventMsg::Error(ErrorEvent {
misalignment: None,
message: "rollback failed".into(),
codex_error_info: Some(CodexErrorInfo::ThreadRollbackFailed),
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].status, TurnStatus::Completed);
assert_eq!(turns[0].error, None);
}
#[test]
fn out_of_turn_error_does_not_create_or_fail_a_turn() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "hello".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
EventMsg::Error(ErrorEvent {
misalignment: None,
message: "request-level failure".into(),
codex_error_info: Some(CodexErrorInfo::BadRequest),
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(
turns[0],
Turn {
id: "turn-a".into(),
status: TurnStatus::Completed,
error: None,
started_at: None,
completed_at: None,
duration_ms: None,
items_view: TurnItemsView::Full,
items: vec![ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![UserInput::Text {
text: "hello".into(),
text_elements: Vec::new(),
}],
}],
}
);
}
#[test]
fn error_then_turn_complete_preserves_failed_status() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "hello".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::Error(ErrorEvent {
misalignment: None,
message: "stream failure".into(),
codex_error_info: Some(CodexErrorInfo::ResponseStreamDisconnected {
http_status_code: Some(502),
}),
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].id, "turn-a");
assert_eq!(turns[0].status, TurnStatus::Failed);
assert_eq!(
turns[0].error,
Some(TurnError {
misalignment: None,
message: "stream failure".into(),
codex_error_info: Some(
crate::protocol::v2::CodexErrorInfo::ResponseStreamDisconnected {
http_status_code: Some(502),
}
),
additional_details: None,
})
);
}
#[test]
fn turn_complete_with_embedded_error_marks_turn_failed() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: Some(10),
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "retry me".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: Some(10),
last_agent_message: None,
error: Some(ErrorEvent {
misalignment: None,
message: "Selected model is at capacity. Please try a different model.".into(),
codex_error_info: Some(CodexErrorInfo::ServerOverloaded),
}),
completed_at: Some(20),
duration_ms: Some(10_000),
time_to_first_token_ms: None,
}),
];
let items = events
.into_iter()
.map(RolloutItem::EventMsg)
.collect::<Vec<_>>();
assert_eq!(
build_turns_from_rollout_items(&items),
vec![Turn {
id: "turn-a".into(),
items_view: TurnItemsView::Full,
items: vec![ThreadItem::UserMessage {
id: "item-1".into(),
client_id: None,
content: vec![UserInput::Text {
text: "retry me".into(),
text_elements: Vec::new(),
}],
}],
status: TurnStatus::Failed,
error: Some(TurnError {
misalignment: None,
message: "Selected model is at capacity. Please try a different model.".into(),
codex_error_info: Some(crate::protocol::v2::CodexErrorInfo::ServerOverloaded),
additional_details: None,
}),
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(10_000),
}]
);
}
#[test]
fn rebuilds_hook_prompt_items_from_rollout_response_items() {
let hook_prompt = build_hook_prompt_message(&[
CoreHookPromptFragment::from_single_hook("Retry with tests.", "hook-run-1"),
CoreHookPromptFragment::from_single_hook("Then summarize cleanly.", "hook-run-2"),
])
.expect("hook prompt message");
let items = vec![
RolloutItem::EventMsg(EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
})),
RolloutItem::EventMsg(EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "hello".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
})),
RolloutItem::ResponseItem(hook_prompt.into()),
RolloutItem::EventMsg(EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
})),
];
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].items.len(), 2);
assert_eq!(
turns[0].items[1],
ThreadItem::HookPrompt {
id: turns[0].items[1].id().to_string(),
fragments: vec![
crate::protocol::v2::HookPromptFragment {
text: "Retry with tests.".into(),
hook_run_id: "hook-run-1".into(),
},
crate::protocol::v2::HookPromptFragment {
text: "Then summarize cleanly.".into(),
hook_run_id: "hook-run-2".into(),
},
],
}
);
}
#[test]
fn canonical_hook_prompt_completion_updates_turn_history() {
let hook_prompt = CoreTurnItem::HookPrompt(codex_protocol::items::HookPromptItem {
id: "hook-prompt-1".into(),
fragments: vec![CoreHookPromptFragment::from_single_hook(
"Retry with tests.",
"hook-run-1",
)],
});
let expected_item = ThreadItem::from(hook_prompt.clone());
let mut builder = ThreadHistoryBuilder::new();
builder.handle_event(&EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}));
builder.handle_event(&EventMsg::ItemCompleted(ItemCompletedEvent {
thread_id: ThreadId::new(),
turn_id: "turn-a".into(),
item: hook_prompt,
started_at_ms: Some(0),
completed_at_ms: 0,
}));
assert_eq!(
builder.active_turn_snapshot().expect("active turn").items,
vec![expected_item]
);
}
#[test]
fn completed_sub_agent_activity_updates_completed_parent_turn() {
let child_thread_id = ThreadId::new();
let child_path = codex_protocol::AgentPath::root()
.join("worker")
.expect("worker path");
let items = vec![
RolloutItem::EventMsg(EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
})),
RolloutItem::EventMsg(EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
})),
RolloutItem::EventMsg(EventMsg::ItemCompleted(ItemCompletedEvent {
thread_id: ThreadId::new(),
turn_id: "turn-a".into(),
item: CoreTurnItem::SubAgentActivity(CoreSubAgentActivityItem {
id: "child-turn-completed".into(),
kind: CoreSubAgentActivityKind::Completed,
agent_thread_id: child_thread_id,
agent_path: child_path,
}),
started_at_ms: None,
completed_at_ms: 0,
})),
];
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert_eq!(
turns[0].items,
vec![ThreadItem::SubAgentActivity {
id: "child-turn-completed".into(),
kind: crate::protocol::v2::SubAgentActivityKind::Completed,
agent_thread_id: child_thread_id.to_string(),
agent_path: "/root/worker".into(),
}]
);
}
#[test]
fn ignores_plain_user_response_items_in_rollout_replay() {
let items = vec![
RolloutItem::EventMsg(EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
})),
RolloutItem::ResponseItem(
codex_protocol::models::ResponseItem::Message {
id: Some(codex_protocol::ResponseItemId::with_suffix("msg", "1")),
role: "user".into(),
content: vec![codex_protocol::models::ContentItem::InputText {
text: "plain text".into(),
}],
phase: None,
internal_chat_message_metadata_passthrough: None,
}
.into(),
),
RolloutItem::EventMsg(EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: None,
duration_ms: None,
time_to_first_token_ms: None,
})),
];
let turns = build_turns_from_rollout_items(&items);
assert_eq!(turns.len(), 1);
assert!(turns[0].items.is_empty());
}
#[test]
fn changed_rollout_item_reports_new_item_snapshot() {
let mut builder = ThreadHistoryBuilder::new();
let changes = builder.handle_rollout_item_with_changes(&RolloutItem::EventMsg(
EventMsg::UserMessage(UserMessageEvent {
client_id: Some("client-message-1".into()),
message: "hello".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
}),
));
assert_eq!(
changes,
ThreadHistoryChangeSet {
changed_items: vec![ThreadHistoryItemChange {
turn_id: "rollout-0".into(),
item: ThreadItem::UserMessage {
id: "item-1".into(),
client_id: Some("client-message-1".into()),
content: vec![UserInput::Text {
text: "hello".into(),
text_elements: Vec::new(),
}],
},
started_at_ms: None,
completed_at_ms: None,
}],
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: "rollout-0".into(),
root_turn_id: None,
status: TurnStatus::Completed,
error: None,
started_at: None,
completed_at: None,
duration_ms: None,
}],
removed_turn_ids: Vec::new(),
}
);
}
#[test]
fn changed_rollout_item_reports_updated_existing_item_snapshot() {
let mut builder = ThreadHistoryBuilder::new();
builder.handle_rollout_item_with_changes(&RolloutItem::EventMsg(EventMsg::WebSearchBegin(
WebSearchBeginEvent {
call_id: "search-1".into(),
},
)));
let changes = builder.handle_rollout_item_with_changes(&RolloutItem::EventMsg(
EventMsg::WebSearchEnd(WebSearchEndEvent {
call_id: "search-1".into(),
query: "codex".into(),
action: CoreWebSearchAction::Search {
query: Some("codex".into()),
queries: None,
},
results: None,
}),
));
assert_eq!(
changes,
ThreadHistoryChangeSet {
changed_items: vec![ThreadHistoryItemChange {
turn_id: "rollout-0".into(),
item: ThreadItem::WebSearch(WebSearchItem {
id: "search-1".into(),
query: "codex".into(),
action: Some(WebSearchAction::Search {
query: Some("codex".into()),
queries: None,
}),
results: None,
}),
started_at_ms: None,
completed_at_ms: None,
}],
changed_turns: Vec::new(),
removed_turn_ids: Vec::new(),
}
);
}
#[test]
fn changed_rollout_item_reports_streaming_item_mutation() {
let mut builder = ThreadHistoryBuilder::new();
builder.handle_rollout_item_with_changes(&RolloutItem::EventMsg(EventMsg::AgentReasoning(
AgentReasoningEvent {
text: "summary".into(),
},
)));
let changes = builder.handle_rollout_item_with_changes(&RolloutItem::EventMsg(
EventMsg::AgentReasoningRawContent(AgentReasoningRawContentEvent {
text: "raw content".into(),
}),
));
assert_eq!(
changes,
ThreadHistoryChangeSet {
changed_items: vec![ThreadHistoryItemChange {
turn_id: "rollout-0".into(),
item: ThreadItem::Reasoning {
id: "item-1".into(),
summary: vec!["summary".into()],
content: vec!["raw content".into()],
},
started_at_ms: None,
completed_at_ms: None,
}],
changed_turns: Vec::new(),
removed_turn_ids: Vec::new(),
}
);
}
#[test]
fn changed_rollout_item_reports_turn_completion_metadata() {
let mut builder = ThreadHistoryBuilder::new();
let start_changes = builder.handle_rollout_item_with_changes(&RolloutItem::EventMsg(
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: Some("root-turn".into()),
trace_id: None,
started_at: Some(10),
model_context_window: None,
collaboration_mode_kind: Default::default(),
}),
));
assert_eq!(
start_changes,
ThreadHistoryChangeSet {
changed_items: Vec::new(),
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: "turn-a".into(),
root_turn_id: Some("root-turn".into()),
status: TurnStatus::InProgress,
error: None,
started_at: Some(10),
completed_at: None,
duration_ms: None,
}],
removed_turn_ids: Vec::new(),
}
);
builder.handle_rollout_item_with_changes(&RolloutItem::EventMsg(EventMsg::UserMessage(
UserMessageEvent {
client_id: None,
message: "hello".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
},
)));
let complete_changes = builder.handle_rollout_item_with_changes(&RolloutItem::EventMsg(
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: Some(20),
duration_ms: Some(123),
time_to_first_token_ms: None,
}),
));
assert_eq!(
complete_changes,
ThreadHistoryChangeSet {
changed_items: Vec::new(),
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: "turn-a".into(),
root_turn_id: Some("root-turn".into()),
status: TurnStatus::Completed,
error: None,
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(123),
}],
removed_turn_ids: Vec::new(),
}
);
}
#[test]
fn changed_rollout_items_dedupe_updated_item_snapshots() {
let mut builder = ThreadHistoryBuilder::new();
let changes = builder.handle_rollout_items_with_changes(&[
RolloutItem::EventMsg(EventMsg::WebSearchBegin(WebSearchBeginEvent {
call_id: "search-1".into(),
})),
RolloutItem::EventMsg(EventMsg::WebSearchEnd(WebSearchEndEvent {
call_id: "search-1".into(),
query: "codex".into(),
action: CoreWebSearchAction::Search {
query: Some("codex".into()),
queries: None,
},
results: None,
})),
]);
assert_eq!(
changes,
ThreadHistoryChangeSet {
changed_items: vec![ThreadHistoryItemChange {
turn_id: "rollout-0".into(),
item: ThreadItem::WebSearch(WebSearchItem {
id: "search-1".into(),
query: "codex".into(),
action: Some(WebSearchAction::Search {
query: Some("codex".into()),
queries: None,
}),
results: None,
}),
started_at_ms: None,
completed_at_ms: None,
}],
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: "rollout-0".into(),
root_turn_id: None,
status: TurnStatus::Completed,
error: None,
started_at: None,
completed_at: None,
duration_ms: None,
}],
removed_turn_ids: Vec::new(),
}
);
}
#[test]
fn changed_rollout_items_dedupe_turn_metadata_snapshots() {
let mut builder = ThreadHistoryBuilder::new();
let changes = builder.handle_rollout_items_with_changes(&[
RolloutItem::EventMsg(EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: Some("root-turn".into()),
trace_id: None,
started_at: Some(10),
model_context_window: None,
collaboration_mode_kind: Default::default(),
})),
RolloutItem::EventMsg(EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
started_at: None,
last_agent_message: None,
error: None,
completed_at: Some(20),
duration_ms: Some(123),
time_to_first_token_ms: None,
})),
]);
assert_eq!(
changes,
ThreadHistoryChangeSet {
changed_items: Vec::new(),
changed_turns: vec![ThreadHistoryTurnChange {
turn_id: "turn-a".into(),
root_turn_id: Some("root-turn".into()),
status: TurnStatus::Completed,
error: None,
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(123),
}],
removed_turn_ids: Vec::new(),
}
);
}
#[test]
fn changed_rollout_items_drop_prior_changes_for_removed_turns() {
let mut builder = ThreadHistoryBuilder::new();
let changes = builder.handle_rollout_items_with_changes(&[
RolloutItem::EventMsg(EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
root_turn_id: None,
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
})),
RolloutItem::EventMsg(EventMsg::UserMessage(UserMessageEvent {
client_id: None,
message: "hello".into(),
images: None,
text_elements: Vec::new(),
local_images: Vec::new(),
..Default::default()
})),
RolloutItem::EventMsg(EventMsg::ThreadRolledBack(ThreadRolledBackEvent {
num_turns: 1,
})),
]);
assert_eq!(
changes,
ThreadHistoryChangeSet {
changed_items: Vec::new(),
changed_turns: Vec::new(),
removed_turn_ids: vec!["turn-a".into()],
}
);
}
}