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