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 { 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, pub completed_at_ms: Option, } /// 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, pub status: TurnStatus, pub error: Option, pub started_at: Option, pub completed_at: Option, pub duration_ms: Option, } /// Incremental changes produced by opt-in `ThreadHistoryBuilder` handlers. #[derive(Debug, Default, Clone, PartialEq)] pub struct ThreadHistoryChangeSet { pub changed_items: Vec, pub changed_turns: Vec, pub removed_turn_ids: Vec, } 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>, changed_item_indexes: HashMap<(String, String), usize>, changed_turns: Vec>, changed_turn_indexes: HashMap, removed_turn_ids: Vec, removed_turn_indexes: HashMap, } 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, current_turn: Option, next_item_index: i64, current_rollout_index: usize, next_rollout_index: usize, active_change_set: Option, } 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 { self.finish_current_turn(); self.turns.into_iter().map(Turn::from).collect() } pub fn active_turn_snapshot(&self) -> Option { 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 { 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 { 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 { self.current_turn .as_ref() .filter(|turn| turn.opened_explicitly) .map(|turn| turn.id.clone()) } pub fn active_turn_start_index(&self) -> Option { 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 = 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) -> 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) { 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 { 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 { 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>, } impl TurnItemIndex { fn push(&mut self, items: &mut Vec, 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, 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, items: Vec, item_index: TurnItemIndex, error: Option, status: TurnStatus, started_at: Option, completed_at: Option, duration_ms: Option, /// 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) -> Self { self.started_at = started_at; self } } impl From 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); let turns = build_turns_from_rollout_items(&items); assert_eq!(turns, Vec::::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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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::>(); 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()], } ); } }