From cfb8ddb9fffeaa561ffdc69d203d422d94ea2fb8 Mon Sep 17 00:00:00 2001 From: Eric Traut Date: Wed, 8 Jul 2026 19:47:45 -0700 Subject: [PATCH] codex: address PR review feedback (#31176) --- codex-rs/core/src/session/input_queue.rs | 19 ++++------- codex-rs/core/src/session/turn.rs | 40 ++++++++++++++---------- 2 files changed, 30 insertions(+), 29 deletions(-) diff --git a/codex-rs/core/src/session/input_queue.rs b/codex-rs/core/src/session/input_queue.rs index e46a71892b..8a810f2949 100644 --- a/codex-rs/core/src/session/input_queue.rs +++ b/codex-rs/core/src/session/input_queue.rs @@ -54,12 +54,12 @@ impl InputQueue { Option, ) { let activity_rx = self.activity_tx.subscribe(); - let has_pending_steer = if let Some(turn_state) = turn_state { - turn_state.lock().await.pending_input.has_user_input() + let has_pending_turn_input = if let Some(turn_state) = turn_state { + !turn_state.lock().await.pending_input.items.is_empty() } else { false }; - let pending_activity = if has_pending_steer { + let pending_activity = if has_pending_turn_input { Some(InputQueueActivity::Steer) } else if self.has_pending_mailbox_items().await { Some(InputQueueActivity::Mailbox) @@ -181,6 +181,7 @@ impl InputQueue { input: Vec, ) { turn_state.lock().await.pending_input.items.extend(input); + self.activity_tx.send_replace(InputQueueActivity::Steer); } pub(crate) async fn take_pending_input_for_turn_state( @@ -252,14 +253,6 @@ impl InputQueue { } } -impl TurnInputQueue { - fn has_user_input(&self) -> bool { - self.items - .iter() - .any(|input| matches!(input, TurnInput::UserInput { .. })) - } -} - #[cfg(test)] mod tests { use super::*; @@ -313,7 +306,7 @@ mod tests { } #[tokio::test] - async fn input_queue_notifies_steer_subscribers() { + async fn input_queue_notifies_turn_input_subscribers() { let input_queue = InputQueue::new(); let turn_state = Mutex::new(TurnState::default()); let (mut activity_rx, pending_activity) = @@ -321,7 +314,7 @@ mod tests { assert_eq!(pending_activity, None); input_queue - .extend_pending_input_and_accept_mailbox_delivery_for_turn_state( + .extend_pending_input_for_turn_state( &turn_state, vec![TurnInput::UserInput { content: vec![UserInput::Text { diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index 2bf6b1cf7b..7e7b4ea6d6 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -39,7 +39,6 @@ use crate::responses_metadata::CodexResponsesMetadata; use crate::responses_metadata::CodexResponsesRequestKind; use crate::responses_retry::ResponsesStreamRequest; use crate::responses_retry::handle_retryable_response_stream_error; -use crate::session::InputQueueActivity; use crate::session::PreviousTurnSettings; use crate::session::TurnInput; use crate::session::session::Session; @@ -224,7 +223,7 @@ pub(crate) async fn run_turn( // 2. After auto-compact, when model/tool continuation needs to resume before any steer. let mut next_step_context = Some(first_step_context); - let mut retrying_sampling_without_steer = false; + let mut retrying_sampling_without_input = false; loop { // Note that pending_input would be something like a message the user // submitted through the UI while the model was running. Though the UI @@ -254,7 +253,7 @@ pub(crate) async fn run_turn( }; let sampling_request_result: CodexResult<_> = async { // A capacity-only retry should not grow model context or invalidate its cache. - if !std::mem::take(&mut retrying_sampling_without_steer) { + if !std::mem::take(&mut retrying_sampling_without_input) { super::time_reminder::maybe_record_current_time_reminder( sess.as_ref(), turn_context.as_ref(), @@ -483,26 +482,35 @@ pub(crate) async fn run_turn( .input_queue .turn_state_for_sub_id(&sess.active_turn, &turn_context.sub_id) .await; - let (mut activity_rx, pending_activity) = sess + let (mut activity_rx, _) = sess .input_queue .subscribe_activity(turn_state.as_deref()) .await; - let steered = if pending_activity == Some(InputQueueActivity::Steer) { + let retry_interrupted = if sess + .input_queue + .has_pending_input(&sess.active_turn) + .await + { true } else { - tokio::time::timeout(delay, async { - while activity_rx.changed().await.is_ok() { - // Mailbox activity does not represent user intent to retry now. - if *activity_rx.borrow_and_update() == InputQueueActivity::Steer { - return true; - } + tokio::select! { + _ = cancellation_token.cancelled() => { + return Err(CodexErr::TurnAborted); } - false - }) - .await - .unwrap_or(false) + result = tokio::time::timeout(delay, async { + while activity_rx.changed().await.is_ok() { + // Ignore mailbox activity deferred to the next turn. + if sess.input_queue.has_pending_input(&sess.active_turn).await { + return true; + } + } + false + }) => { + result.unwrap_or(false) + } + } }; - retrying_sampling_without_steer = !steered; + retrying_sampling_without_input = !retry_interrupted; can_drain_pending_input = true; continue; }