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
105 changes: 50 additions & 55 deletions codex-rs/app-server/src/bespoke_event_handling.rs
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,7 @@ use codex_core::ThreadManager;
use codex_core::review_format::format_review_findings_block;
use codex_core::review_prompts;
use codex_protocol::ThreadId;
use codex_protocol::items::TurnItem as CoreTurnItem;
use codex_protocol::items::parse_hook_prompt_message;
use codex_protocol::models::AdditionalPermissionProfile as CoreAdditionalPermissionProfile;
use codex_protocol::plan_tool::UpdatePlanArgs;
Expand Down Expand Up @@ -1027,10 +1028,37 @@ pub(crate) async fn apply_bespoke_event_handling(
.send_server_notification(ServerNotification::ItemCompleted(completed))
.await;
}
msg @ (EventMsg::ItemStarted(_)
| EventMsg::ItemCompleted(_)
| EventMsg::PatchApplyUpdated(_)
| EventMsg::TerminalInteraction(_)) => {
EventMsg::ItemStarted(event) => {
let should_emit = match &event.item {
// Approval and guardian flows can emit the command start notification before core
// emits the canonical item. Reuse the same set to suppress that duplicate.
CoreTurnItem::CommandExecution(item) => thread_state
.lock()
.await
.turn_summary
.command_execution_started

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

do we still need this?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

unfortunately we do, we synthesize a ThreadItem::CommandExecution whenever an approval is requested (user or guardian), so this helps us dedupe

.insert(item.id.clone()),
_ => true,
};
if should_emit {
let notification = item_event_to_server_notification(
EventMsg::ItemStarted(event),
&conversation_id.to_string(),
&event_turn_id,
);
outgoing.send_server_notification(notification).await;
}
}
EventMsg::ItemCompleted(event) => {
apply_canonical_item_completed_side_effects(&thread_state, &event.item).await;
let notification = item_event_to_server_notification(
EventMsg::ItemCompleted(event),
&conversation_id.to_string(),
&event_turn_id,
);
outgoing.send_server_notification(notification).await;
}
msg @ (EventMsg::PatchApplyUpdated(_) | EventMsg::TerminalInteraction(_)) => {
let notification = item_event_to_server_notification(
msg,
&conversation_id.to_string(),
Expand Down Expand Up @@ -1106,32 +1134,10 @@ pub(crate) async fn apply_bespoke_event_handling(
// Core still fans out these deprecated events for legacy clients;
// v2 clients receive the canonical FileChange item instead.
}
EventMsg::ExecCommandBegin(exec_command_begin_event) => {
if matches!(
exec_command_begin_event.source,
codex_protocol::protocol::ExecCommandSource::UnifiedExecInteraction
) {
// TerminalInteraction is the v2 surface for unified exec
// stdin/poll events. Suppress the legacy CommandExecution
// item so clients do not render the same wait twice.
return;
}
let item_id = exec_command_begin_event.call_id.clone();
let first_start = {
let mut state = thread_state.lock().await;
state
.turn_summary
.command_execution_started
.insert(item_id.clone())
};
if first_start {
let notification = item_event_to_server_notification(
EventMsg::ExecCommandBegin(exec_command_begin_event),
&conversation_id.to_string(),
&event_turn_id,
);
outgoing.send_server_notification(notification).await;
}
EventMsg::ExecCommandBegin(_) | EventMsg::ExecCommandEnd(_) => {
// Deprecated command-execution events are still fanned out for raw-event and rollout
// compatibility consumers. App-server v2 receives the canonical CommandExecution
// item lifecycle instead.
}
EventMsg::ExecCommandOutputDelta(exec_command_output_delta_event) => {
let notification = item_event_to_server_notification(
Expand All @@ -1141,31 +1147,6 @@ pub(crate) async fn apply_bespoke_event_handling(
);
outgoing.send_server_notification(notification).await;
}
EventMsg::ExecCommandEnd(exec_command_end_event) => {
let call_id = exec_command_end_event.call_id.clone();
{
let mut state = thread_state.lock().await;
state
.turn_summary
.command_execution_started
.remove(&call_id);
}
if matches!(
exec_command_end_event.source,
codex_protocol::protocol::ExecCommandSource::UnifiedExecInteraction
) {
// The paired begin event is suppressed above; keep the
// completion out of v2 as well so no orphan legacy item is
// emitted for unified exec interactions.
return;
}
let notification = item_event_to_server_notification(
EventMsg::ExecCommandEnd(exec_command_end_event),
&conversation_id.to_string(),
&event_turn_id,
);
outgoing.send_server_notification(notification).await;
}
// If this is a TurnAborted, reply to any pending interrupt requests.
EventMsg::TurnAborted(turn_aborted_event) => {
// All per-thread requests are bound to a turn, so abort them.
Expand Down Expand Up @@ -1369,6 +1350,20 @@ async fn emit_turn_completed_with_status(
.await;
}

async fn apply_canonical_item_completed_side_effects(
thread_state: &Arc<Mutex<ThreadState>>,
item: &CoreTurnItem,
) {
if let CoreTurnItem::CommandExecution(item) = item {
thread_state
.lock()
.await
.turn_summary
.command_execution_started
.remove(&item.id);
}
}

#[allow(clippy::too_many_arguments)]
async fn start_command_execution_item(
conversation_id: &ThreadId,
Expand Down
98 changes: 48 additions & 50 deletions codex-rs/core/src/tasks/user_shell.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,16 +25,15 @@ use crate::tools::runtimes::RuntimePathPrepends;
use crate::tools::runtimes::apply_package_path_prepend;
use crate::tools::runtimes::maybe_wrap_shell_lc_with_snapshot;
use crate::tools::runtimes::strip_managed_proxy_env;
use crate::turn_timing::now_unix_timestamp_ms;
use crate::user_shell_command::user_shell_command_record_item;
use codex_protocol::exec_output::ExecToolCallOutput;
use codex_protocol::exec_output::StreamOutput;
use codex_protocol::items::CommandExecutionItem;
use codex_protocol::items::CommandExecutionStatus;
use codex_protocol::items::TurnItem;
use codex_protocol::protocol::ErrorEvent;
use codex_protocol::protocol::EventMsg;
use codex_protocol::protocol::ExecCommandBeginEvent;
use codex_protocol::protocol::ExecCommandEndEvent;
use codex_protocol::protocol::ExecCommandSource;
use codex_protocol::protocol::ExecCommandStatus;
use codex_protocol::protocol::TurnStartedEvent;
use codex_sandboxing::SandboxType;
use codex_shell_command::parse_command::parse_command;
Expand Down Expand Up @@ -181,18 +180,23 @@ pub(crate) async fn execute_user_shell_command(

let parsed_cmd = parse_command(&display_command);
session
.send_event(
.emit_turn_item_started(
turn_context.as_ref(),
EventMsg::ExecCommandBegin(ExecCommandBeginEvent {
call_id: call_id.clone(),
&TurnItem::CommandExecution(CommandExecutionItem {
id: call_id.clone(),
process_id: None,
turn_id: turn_context.sub_id.clone(),
started_at_ms: now_unix_timestamp_ms(),
command: display_command.clone(),
cwd: cwd.clone().into(),
parsed_cmd: parsed_cmd.clone(),
source: ExecCommandSource::UserShell,
interaction_input: None,
status: CommandExecutionStatus::InProgress,
stdout: None,
stderr: None,
aggregated_output: None,
exit_code: None,
duration: None,
formatted_output: None,
}),
)
.await;
Expand Down Expand Up @@ -259,57 +263,53 @@ pub(crate) async fn execute_user_shell_command(
)
.await;
session
.send_event(
.emit_turn_item_completed(
turn_context.as_ref(),
EventMsg::ExecCommandEnd(ExecCommandEndEvent {
call_id,
TurnItem::CommandExecution(CommandExecutionItem {
id: call_id,
process_id: None,
turn_id: turn_context.sub_id.clone(),
completed_at_ms: now_unix_timestamp_ms(),
command: display_command.clone(),
cwd: cwd.clone().into(),
parsed_cmd: parsed_cmd.clone(),
source: ExecCommandSource::UserShell,
interaction_input: None,
stdout: String::new(),
stderr: aborted_message.clone(),
aggregated_output: aborted_message.clone(),
exit_code: -1,
duration: Duration::ZERO,
formatted_output: aborted_message,
status: ExecCommandStatus::Failed,
status: CommandExecutionStatus::Failed,
stdout: Some(String::new()),
stderr: Some(aborted_message.clone()),
aggregated_output: Some(aborted_message.clone()),
exit_code: Some(-1),
duration: Some(Duration::ZERO),
formatted_output: Some(aborted_message),
}),
)
.await;
}
Ok(Ok(output)) => {
session
.send_event(
.emit_turn_item_completed(
turn_context.as_ref(),
EventMsg::ExecCommandEnd(ExecCommandEndEvent {
call_id: call_id.clone(),
TurnItem::CommandExecution(CommandExecutionItem {
id: call_id.clone(),
process_id: None,
turn_id: turn_context.sub_id.clone(),
completed_at_ms: now_unix_timestamp_ms(),
command: display_command.clone(),
cwd: cwd.clone().into(),
parsed_cmd: parsed_cmd.clone(),
source: ExecCommandSource::UserShell,
interaction_input: None,
stdout: output.stdout.text.clone(),
stderr: output.stderr.text.clone(),
aggregated_output: output.aggregated_output.text.clone(),
exit_code: output.exit_code,
duration: output.duration,
formatted_output: format_exec_output_str(
&output,
turn_context.model_info.truncation_policy.into(),
),
status: if output.exit_code == 0 {
ExecCommandStatus::Completed
CommandExecutionStatus::Completed
} else {
ExecCommandStatus::Failed
CommandExecutionStatus::Failed
},
stdout: Some(output.stdout.text.clone()),
stderr: Some(output.stderr.text.clone()),
aggregated_output: Some(output.aggregated_output.text.clone()),
exit_code: Some(output.exit_code),
duration: Some(output.duration),
formatted_output: Some(format_exec_output_str(
&output,
turn_context.model_info.truncation_policy.into(),
)),
}),
)
.await;
Expand All @@ -329,28 +329,26 @@ pub(crate) async fn execute_user_shell_command(
timed_out: false,
};
session
.send_event(
.emit_turn_item_completed(
turn_context.as_ref(),
EventMsg::ExecCommandEnd(ExecCommandEndEvent {
call_id,
TurnItem::CommandExecution(CommandExecutionItem {
id: call_id,
process_id: None,
turn_id: turn_context.sub_id.clone(),
completed_at_ms: now_unix_timestamp_ms(),
command: display_command,
cwd: cwd.into(),
parsed_cmd,
source: ExecCommandSource::UserShell,
interaction_input: None,
stdout: exec_output.stdout.text.clone(),
stderr: exec_output.stderr.text.clone(),
aggregated_output: exec_output.aggregated_output.text.clone(),
exit_code: exec_output.exit_code,
duration: exec_output.duration,
formatted_output: format_exec_output_str(
status: CommandExecutionStatus::Failed,
stdout: Some(exec_output.stdout.text.clone()),
stderr: Some(exec_output.stderr.text.clone()),
aggregated_output: Some(exec_output.aggregated_output.text.clone()),
exit_code: Some(exec_output.exit_code),
duration: Some(exec_output.duration),
formatted_output: Some(format_exec_output_str(
&exec_output,
turn_context.model_info.truncation_policy.into(),
),
status: ExecCommandStatus::Failed,
)),
}),
)
.await;
Expand Down
Loading
Loading