Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion codex-rs/app-server-protocol/src/protocol/item_builders.rs
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,10 @@ pub fn build_command_execution_end_item(payload: &ExecCommandEndEvent) -> Thread
}
}

fn command_actions_for_path_uri(parsed_cmd: &[ParsedCommand], cwd: &PathUri) -> Vec<CommandAction> {
pub(crate) fn command_actions_for_path_uri(
parsed_cmd: &[ParsedCommand],
cwd: &PathUri,
) -> Vec<CommandAction> {
// TODO(anp): Carry PathUri into CommandAction so foreign Read actions retain resolved paths.
// Until then, omit those actions rather than project a foreign cwd onto the host.
let native_cwd = if cwd.infer_path_convention() == Some(PathConvention::native()) {
Expand Down
138 changes: 94 additions & 44 deletions codex-rs/app-server-protocol/src/protocol/thread_history.rs
Original file line number Diff line number Diff line change
Expand Up @@ -574,52 +574,25 @@ impl ThreadHistoryBuilder {
}

fn handle_item_started(&mut self, payload: &ItemStartedEvent) {
match &payload.item {
codex_protocol::items::TurnItem::Plan(plan) => {
if plan.text.is_empty() {
return;
}
self.upsert_item_in_turn_id(
&payload.turn_id,
ThreadItem::from(payload.item.clone()),
);
}
codex_protocol::items::TurnItem::Sleep(_) => {
self.upsert_item_in_turn_id(
&payload.turn_id,
ThreadItem::from(payload.item.clone()),
);
}
codex_protocol::items::TurnItem::UserMessage(_)
| codex_protocol::items::TurnItem::HookPrompt(_)
| 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(_) => {}
}
self.handle_materialized_item_lifecycle(&payload.turn_id, &payload.item);
}

fn handle_item_completed(&mut self, payload: &ItemCompletedEvent) {
match &payload.item {
codex_protocol::items::TurnItem::Plan(plan) => {
if plan.text.is_empty() {
return;
}
self.upsert_item_in_turn_id(
&payload.turn_id,
ThreadItem::from(payload.item.clone()),
);
}
codex_protocol::items::TurnItem::Sleep(_) => {
self.upsert_item_in_turn_id(
&payload.turn_id,
ThreadItem::from(payload.item.clone()),
);
}
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 should_upsert = match item {
codex_protocol::items::TurnItem::Plan(plan) => !plan.text.is_empty(),
codex_protocol::items::TurnItem::Sleep(_)
| codex_protocol::items::TurnItem::CommandExecution(_)
| codex_protocol::items::TurnItem::DynamicToolCall(_)
| codex_protocol::items::TurnItem::CollabAgentToolCall(_)
| codex_protocol::items::TurnItem::SubAgentActivity(_) => true,
codex_protocol::items::TurnItem::UserMessage(_)
| codex_protocol::items::TurnItem::HookPrompt(_)
| codex_protocol::items::TurnItem::AgentMessage(_)
Expand All @@ -629,7 +602,11 @@ impl ThreadHistoryBuilder {
| codex_protocol::items::TurnItem::ImageGeneration(_)
| codex_protocol::items::TurnItem::FileChange(_)
| codex_protocol::items::TurnItem::McpToolCall(_)
| codex_protocol::items::TurnItem::ContextCompaction(_) => {}
| codex_protocol::items::TurnItem::ContextCompaction(_) => false,
};

if should_upsert {
self.upsert_item_in_turn_id(turn_id, ThreadItem::from(item.clone()));
}
}

Expand Down Expand Up @@ -1573,6 +1550,8 @@ mod tests {
use crate::protocol::v2::CommandExecutionSource;
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::HookPromptFragment as CoreHookPromptFragment;
use codex_protocol::items::SleepItem as CoreSleepItem;
use codex_protocol::items::TurnItem as CoreTurnItem;
Expand Down Expand Up @@ -1867,6 +1846,77 @@ mod tests {
);
}

#[test]
fn rebuilds_command_execution_item_from_persisted_completion() {
let turn_id = "turn-1";
let thread_id = ThreadId::new();
let command_item = CoreTurnItem::CommandExecution(CoreCommandExecutionItem {
id: "exec-1".to_string(),
process_id: Some("pid-1".to_string()),
command: vec!["echo".to_string(), "hello world".to_string()],
cwd: test_path_buf("/tmp").abs().into(),
parsed_cmd: vec![ParsedCommand::Unknown {
cmd: "echo hello world".to_string(),
}],
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(),
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: command_item,
completed_at_ms: 1_000,
}),
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: turn_id.to_string(),
last_agent_message: 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::CommandExecution {
id: "exec-1".to_string(),
command: "echo 'hello world'".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: "echo hello world".to_string(),
}],
aggregated_output: Some("hello world\n".to_string()),
exit_code: Some(0),
duration_ms: Some(12),
}]
);
}

#[test]
fn preserves_user_message_client_id_from_legacy_event() {
let turn_id = "turn-1";
Expand Down
122 changes: 122 additions & 0 deletions codex-rs/app-server-protocol/src/protocol/v2/item.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,12 +8,17 @@ use super::NetworkPolicyAmendment;
use super::RequestPermissionProfile;
use super::UserInput;
use super::shared::v2_enum_from_core;
use crate::protocol::item_builders::command_actions_for_path_uri;
use crate::protocol::item_builders::convert_patch_changes;
use codex_experimental_api_macros::ExperimentalApi;
use codex_protocol::approvals::GuardianAssessmentAction as CoreGuardianAssessmentAction;
use codex_protocol::approvals::GuardianAssessmentDecisionSource as CoreGuardianAssessmentDecisionSource;
use codex_protocol::approvals::GuardianCommandSource as CoreGuardianCommandSource;
use codex_protocol::items::AgentMessageContent as CoreAgentMessageContent;
use codex_protocol::items::CollabAgentTool as CoreCollabAgentTool;
use codex_protocol::items::CollabAgentToolCallStatus as CoreCollabAgentToolCallStatus;
use codex_protocol::items::CommandExecutionStatus as CoreCommandExecutionStatus;
use codex_protocol::items::DynamicToolCallStatus as CoreDynamicToolCallStatus;
use codex_protocol::items::McpToolCallStatus as CoreMcpToolCallStatus;
use codex_protocol::items::TurnItem as CoreTurnItem;
use codex_protocol::memory_citation::MemoryCitation as CoreMemoryCitation;
Expand All @@ -30,6 +35,7 @@ use codex_protocol::protocol::GuardianUserAuthorization as CoreGuardianUserAutho
use codex_protocol::protocol::PatchApplyStatus as CorePatchApplyStatus;
use codex_protocol::protocol::ReviewDecision as CoreReviewDecision;
use codex_protocol::protocol::SubAgentActivityKind as CoreSubAgentActivityKind;
use codex_shell_command::parse_command::shlex_join;
use codex_utils_absolute_path::AbsolutePathBuf;
use codex_utils_path_uri::LegacyAppPathString;
use schemars::JsonSchema;
Expand Down Expand Up @@ -854,6 +860,64 @@ impl From<CoreTurnItem> for ThreadItem {
summary: reasoning.summary_text,
content: reasoning.raw_content,
},
CoreTurnItem::CommandExecution(command) => ThreadItem::CommandExecution {
id: command.id,
command: shlex_join(&command.command),
cwd: command.cwd.clone().into(),
process_id: command.process_id,
source: command.source.into(),
status: command.status.into(),
command_actions: command_actions_for_path_uri(&command.parsed_cmd, &command.cwd),
aggregated_output: command
.aggregated_output
.filter(|output| !output.is_empty()),
exit_code: command.exit_code,
duration_ms: command
.duration
.and_then(|duration| i64::try_from(duration.as_millis()).ok()),
},
CoreTurnItem::DynamicToolCall(call) => ThreadItem::DynamicToolCall {
id: call.id,
namespace: call.namespace,
tool: call.tool,
arguments: call.arguments,
status: call.status.into(),
content_items: call.content_items.map(|items| {
items
.into_iter()
.map(DynamicToolCallOutputContentItem::from)
.collect()
}),
success: call.success,
duration_ms: call
.duration
.and_then(|duration| i64::try_from(duration.as_millis()).ok()),
},
CoreTurnItem::CollabAgentToolCall(call) => ThreadItem::CollabAgentToolCall {
id: call.id,
tool: call.tool.into(),
status: call.status.into(),
sender_thread_id: call.sender_thread_id.to_string(),
receiver_thread_ids: call
.receiver_thread_ids
.into_iter()
.map(String::from)
.collect(),
prompt: call.prompt,
model: call.model,
reasoning_effort: call.reasoning_effort,
agents_states: call
.agents_states
.into_iter()
.map(|(thread_id, status)| (thread_id.to_string(), status.into()))
.collect(),
},
CoreTurnItem::SubAgentActivity(activity) => ThreadItem::SubAgentActivity {
id: activity.id,
kind: activity.kind.into(),
agent_thread_id: activity.agent_thread_id.to_string(),
agent_path: String::from(activity.agent_path),
},
CoreTurnItem::WebSearch(search) => ThreadItem::WebSearch {
id: search.id,
query: search.query,
Expand Down Expand Up @@ -941,6 +1005,17 @@ impl From<CoreExecCommandStatus> for CommandExecutionStatus {
}
}

impl From<CoreCommandExecutionStatus> for CommandExecutionStatus {
fn from(value: CoreCommandExecutionStatus) -> Self {
match value {
CoreCommandExecutionStatus::InProgress => Self::InProgress,
CoreCommandExecutionStatus::Completed => Self::Completed,
CoreCommandExecutionStatus::Failed => Self::Failed,
CoreCommandExecutionStatus::Declined => Self::Declined,
}
}
}

impl From<&CoreExecCommandStatus> for CommandExecutionStatus {
fn from(value: &CoreExecCommandStatus) -> Self {
match value {
Expand Down Expand Up @@ -1028,6 +1103,16 @@ impl From<CoreMcpToolCallStatus> for McpToolCallStatus {
}
}

impl From<CoreDynamicToolCallStatus> for DynamicToolCallStatus {
fn from(value: CoreDynamicToolCallStatus) -> Self {
match value {
CoreDynamicToolCallStatus::InProgress => Self::InProgress,
CoreDynamicToolCallStatus::Completed => Self::Completed,
CoreDynamicToolCallStatus::Failed => Self::Failed,
}
}
}

#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export_to = "v2/")]
Expand Down Expand Up @@ -1055,6 +1140,28 @@ pub enum CollabAgentToolCallStatus {
Failed,
}

impl From<CoreCollabAgentTool> for CollabAgentTool {
fn from(value: CoreCollabAgentTool) -> Self {
match value {
CoreCollabAgentTool::SpawnAgent => Self::SpawnAgent,
CoreCollabAgentTool::SendInput => Self::SendInput,
CoreCollabAgentTool::ResumeAgent => Self::ResumeAgent,
CoreCollabAgentTool::Wait => Self::Wait,
CoreCollabAgentTool::CloseAgent => Self::CloseAgent,
}
}
}

impl From<CoreCollabAgentToolCallStatus> for CollabAgentToolCallStatus {
fn from(value: CoreCollabAgentToolCallStatus) -> Self {
match value {
CoreCollabAgentToolCallStatus::InProgress => Self::InProgress,
CoreCollabAgentToolCallStatus::Completed => Self::Completed,
CoreCollabAgentToolCallStatus::Failed => Self::Failed,
}
}
}

#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq, JsonSchema, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export_to = "v2/")]
Expand Down Expand Up @@ -1462,6 +1569,21 @@ pub enum DynamicToolCallOutputContentItem {
InputImage { image_url: String },
}

impl From<codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem>
for DynamicToolCallOutputContentItem
{
fn from(item: codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem) -> Self {
match item {
codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem::InputText { text } => {
Self::InputText { text }
}
codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem::InputImage {
image_url,
} => Self::InputImage { image_url },
}
}
}

impl From<DynamicToolCallOutputContentItem>
for codex_protocol::dynamic_tools::DynamicToolCallOutputContentItem
{
Expand Down
Loading
Loading