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
184 changes: 179 additions & 5 deletions codex-rs/app-server-protocol/src/protocol/thread_history.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1245,8 +1245,16 @@ impl ThreadHistoryBuilder {
}

fn handle_turn_complete(&mut self, payload: &TurnCompleteEvent) {
let mark_completed = |turn: &mut PendingTurn| {
if matches!(turn.status, TurnStatus::Completed | TurnStatus::InProgress) {
let terminal_error = payload.error.as_ref().map(|error| V2TurnError {
message: error.message.clone(),
codex_error_info: error.codex_error_info.clone().map(Into::into),
additional_details: None,
});
Comment on lines +1248 to +1252

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.

P1 Badge Exercise failed completion replay through the public RPC

When a persisted rollout containing TurnCompleteEvent { error: Some(...) } is loaded through thread/read or thread/turns/list, the new protocol test only invokes the reducer directly, while the TUI snapshot starts with an already-failed AppServerTurn. Neither test verifies that the request-processing and serialization path returns the reconstructed failed status and error, so add an app-server integration test that seeds this rollout and asserts the JSON-RPC response.

AGENTS.md reference: AGENTS.md:L252-L257

Useful? React with 👍 / 👎.

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;
Expand All @@ -1260,7 +1268,7 @@ impl ThreadHistoryBuilder {
.as_mut()
.filter(|turn| turn.id == payload.turn_id)
{
let changed_turn = mark_completed(current_turn);
let changed_turn = apply_completion(current_turn);
self.record_changed_turn(changed_turn);
self.finish_current_turn();
return;
Expand All @@ -1271,7 +1279,10 @@ impl ThreadHistoryBuilder {
.iter_mut()
.find(|turn| turn.id == payload.turn_id)
{
if matches!(turn.status, TurnStatus::Completed | TurnStatus::InProgress) {
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;
Expand All @@ -1283,7 +1294,7 @@ impl ThreadHistoryBuilder {

// 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 = mark_completed(current_turn);
let changed_turn = apply_completion(current_turn);
self.record_changed_turn(changed_turn);
self.finish_current_turn();
}
Expand Down Expand Up @@ -3671,6 +3682,106 @@ mod tests {
assert_eq!(turns[1].items.len(), 2);
}

#[test]
fn late_turn_complete_with_embedded_error_preserves_active_turn() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
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(),
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 {
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 {
message: "Selected model is at capacity. Please try a different model."
.into(),
codex_error_info: Some(
crate::protocol::v2::CodexErrorInfo::ServerOverloaded,
),
additional_details: None,
}),
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(10_000),
},
Turn {
id: "turn-b".into(),
items_view: TurnItemsView::Full,
items: vec![ThreadItem::UserMessage {
id: "item-2".into(),
client_id: None,
content: vec![UserInput::Text {
text: "second".into(),
text_elements: Vec::new(),
}],
}],
status: TurnStatus::InProgress,
error: None,
started_at: Some(30),
completed_at: None,
duration_ms: None,
},
]
);
}

#[test]
fn late_turn_aborted_does_not_interrupt_active_turn() {
let events = vec![
Expand Down Expand Up @@ -4121,6 +4232,69 @@ mod tests {
);
}

#[test]
fn turn_complete_with_embedded_error_marks_turn_failed() {
let events = vec![
EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
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 {
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 {
message: "Selected model is at capacity. Please try a different model.".into(),
codex_error_info: Some(crate::protocol::v2::CodexErrorInfo::ServerOverloaded),
additional_details: None,
}),
started_at: Some(10),
completed_at: Some(20),
duration_ms: Some(10_000),
}]
);
}

#[test]
fn rebuilds_hook_prompt_items_from_rollout_response_items() {
let hook_prompt = build_hook_prompt_message(&[
Expand Down
48 changes: 48 additions & 0 deletions codex-rs/tui/src/chatwidget/tests/history_replay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,54 @@ async fn resumed_initial_messages_render_history() {
);
}

#[tokio::test]
async fn replayed_failed_turns_preserve_overload_warnings_between_retries() {
let (mut chat, mut rx, _ops) = make_chatwidget_manual(/*model_override*/ None).await;
let prompt = "The workspace also looks super confusing with its separator.";
let error_message = "Selected model is at capacity. Please try a different model.";
let failed_turn = |turn_id: &str, item_id: &str| AppServerTurn {
items: vec![AppServerThreadItem::UserMessage {
id: item_id.to_string(),
client_id: None,
content: vec![AppServerUserInput::Text {
text: prompt.to_string(),
text_elements: Vec::new(),
}],
}],
..app_server_turn(
turn_id,
AppServerTurnStatus::Failed,
/*duration_ms*/ None,
/*error*/
Some(AppServerTurnError {
message: error_message.to_string(),
codex_error_info: Some(CodexErrorInfo::ServerOverloaded),
additional_details: None,
}),
)
};

chat.replay_thread_turns(
vec![
failed_turn("turn-1", "user-1"),
failed_turn("turn-2", "user-2"),
],
ReplayKind::ResumeInitialMessages,
);

let rendered = drain_insert_history(&mut rx)
.into_iter()
.map(|lines| lines_to_single_string(&lines))
.collect::<String>();

assert_eq!(rendered.matches(prompt).count(), 2);
assert_eq!(rendered.matches(error_message).count(), 2);
insta::assert_snapshot!(
"replayed_failed_turns_preserve_overload_warnings_between_retries",
rendered
);
}

#[tokio::test]
async fn restored_conversation_ultra_remains_selected_after_switching_to_plan() {
let (mut chat, _rx, _ops) = make_chatwidget_manual(Some("gpt-5.4")).await;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
---
source: tui/src/chatwidget/tests/history_replay.rs
expression: rendered
---

› The workspace also looks super confusing with its separator.


⚠ Selected model is at capacity. Please try a different model.


› The workspace also looks super confusing with its separator.


⚠ Selected model is at capacity. Please try a different model.
Loading