diff --git a/codex-rs/core/src/goals.rs b/codex-rs/core/src/goals.rs index 4b3e237d5..73e82fdcf 100644 --- a/codex-rs/core/src/goals.rs +++ b/codex-rs/core/src/goals.rs @@ -7,6 +7,7 @@ use crate::StateDbHandle; use crate::context::ContextualUserFragment; use crate::context::GoalContext; +use crate::session::TurnInput; use crate::session::session::Session; use crate::session::turn_context::TurnContext; use crate::state::ActiveTurn; @@ -1347,7 +1348,14 @@ impl Session { return; } self.input_queue - .extend_pending_input_for_turn_state(turn_state.as_ref(), candidate.items) + .extend_pending_input_for_turn_state( + turn_state.as_ref(), + candidate + .items + .into_iter() + .map(TurnInput::ResponseInputItem) + .collect(), + ) .await; let turn_context = self diff --git a/codex-rs/core/src/hook_runtime.rs b/codex-rs/core/src/hook_runtime.rs index da883b1e6..8cba0f9fa 100644 --- a/codex-rs/core/src/hook_runtime.rs +++ b/codex-rs/core/src/hook_runtime.rs @@ -20,7 +20,7 @@ use codex_hooks::UserPromptSubmitRequest; use codex_otel::HOOK_RUN_DURATION_METRIC; use codex_otel::HOOK_RUN_METRIC; use codex_protocol::items::TurnItem; -use codex_protocol::models::ResponseInputItem; +use codex_protocol::items::UserMessageItem; use codex_protocol::models::ResponseItem; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::EventMsg; @@ -30,12 +30,12 @@ use codex_protocol::protocol::HookRunStatus; use codex_protocol::protocol::HookRunSummary; use codex_protocol::protocol::HookSource; use codex_protocol::protocol::HookStartedEvent; -use codex_protocol::user_input::UserInput; use serde_json::Value; use crate::context::ContextualUserFragment; use crate::context::HookAdditionalContext; use crate::event_mapping::parse_turn_item; +use crate::session::TurnInput; use crate::session::session::Session; use crate::session::turn_context::TurnContext; use crate::tools::hook_names::HookToolName; @@ -51,22 +51,6 @@ pub(crate) enum PreToolUseHookResult { Blocked(String), } -pub(crate) enum PendingInputHookDisposition { - Accepted(Box), - Blocked { additional_contexts: Vec }, -} - -pub(crate) enum PendingInputRecord { - UserMessage { - content: Vec, - response_item: ResponseItem, - additional_contexts: Vec, - }, - ConversationItem { - response_item: ResponseItem, - }, -} - struct ContextInjectingHookOutcome { hook_events: Vec, outcome: HookRuntimeOutcome, @@ -338,32 +322,6 @@ pub(crate) async fn run_post_compact_hooks( } } -pub(crate) async fn run_user_prompt_submit_hooks( - sess: &Arc, - turn_context: &Arc, - prompt: String, -) -> HookRuntimeOutcome { - let request = UserPromptSubmitRequest { - session_id: sess.session_id().into(), - turn_id: turn_context.sub_id.clone(), - #[allow(deprecated)] - cwd: turn_context.cwd.clone(), - transcript_path: sess.hook_transcript_path().await, - model: turn_context.model_info.slug.clone(), - permission_mode: hook_permission_mode(turn_context), - prompt, - }; - let hooks = sess.hooks(); - let preview_runs = hooks.preview_user_prompt_submit(&request); - run_context_injecting_hook( - sess, - turn_context, - preview_runs, - hooks.run_user_prompt_submit(request), - ) - .await -} - pub(crate) async fn run_stop_hooks( sess: &Arc, turn_context: &Arc, @@ -458,54 +416,55 @@ pub(crate) async fn run_legacy_after_agent_hook( pub(crate) async fn inspect_pending_input( sess: &Arc, turn_context: &Arc, - pending_input_item: ResponseInputItem, -) -> PendingInputHookDisposition { - let response_item = ResponseItem::from(pending_input_item); - if let Some(TurnItem::UserMessage(user_message)) = parse_turn_item(&response_item) { - let user_prompt_submit_outcome = - run_user_prompt_submit_hooks(sess, turn_context, user_message.message()).await; - if user_prompt_submit_outcome.should_stop { - PendingInputHookDisposition::Blocked { - additional_contexts: user_prompt_submit_outcome.additional_contexts, - } - } else { - PendingInputHookDisposition::Accepted(Box::new(PendingInputRecord::UserMessage { - content: user_message.content, - response_item, - additional_contexts: user_prompt_submit_outcome.additional_contexts, - })) + pending_input_item: &TurnInput, +) -> HookRuntimeOutcome { + match pending_input_item { + TurnInput::UserInput(content) => { + let request = UserPromptSubmitRequest { + session_id: sess.session_id().into(), + turn_id: turn_context.sub_id.clone(), + #[allow(deprecated)] + cwd: turn_context.cwd.clone(), + transcript_path: sess.hook_transcript_path().await, + model: turn_context.model_info.slug.clone(), + permission_mode: hook_permission_mode(turn_context), + prompt: UserMessageItem::new(content).message(), + }; + let hooks = sess.hooks(); + let preview_runs = hooks.preview_user_prompt_submit(&request); + run_context_injecting_hook( + sess, + turn_context, + preview_runs, + hooks.run_user_prompt_submit(request), + ) + .await } - } else { - PendingInputHookDisposition::Accepted(Box::new(PendingInputRecord::ConversationItem { - response_item, - })) + TurnInput::ResponseInputItem(_) => HookRuntimeOutcome { + should_stop: false, + additional_contexts: Vec::new(), + }, } } pub(crate) async fn record_pending_input( sess: &Arc, turn_context: &Arc, - pending_input: PendingInputRecord, + pending_input: TurnInput, + additional_contexts: Vec, ) { match pending_input { - PendingInputRecord::UserMessage { - content, - response_item, - additional_contexts, - } => { - sess.record_user_prompt_and_emit_turn_item( - turn_context.as_ref(), - content.as_slice(), - response_item, - ) - .await; - record_additional_contexts(sess, turn_context, additional_contexts).await; + TurnInput::UserInput(content) => { + sess.record_user_prompt_and_emit_turn_item(turn_context.as_ref(), content.as_slice()) + .await; } - PendingInputRecord::ConversationItem { response_item } => { + TurnInput::ResponseInputItem(input) => { + let response_item = ResponseItem::from(input); sess.record_conversation_items(turn_context, std::slice::from_ref(&response_item)) .await; } } + record_additional_contexts(sess, turn_context, additional_contexts).await; } async fn run_context_injecting_hook( diff --git a/codex-rs/core/src/session/input_queue.rs b/codex-rs/core/src/session/input_queue.rs index 2e8a6ded9..5f92322c8 100644 --- a/codex-rs/core/src/session/input_queue.rs +++ b/codex-rs/core/src/session/input_queue.rs @@ -3,15 +3,22 @@ use crate::state::MailboxDeliveryPhase; use crate::state::TurnState; use codex_protocol::models::ResponseInputItem; use codex_protocol::protocol::InterAgentCommunication; +use codex_protocol::user_input::UserInput; use std::collections::VecDeque; use std::sync::Arc; use tokio::sync::Mutex; use tokio::sync::watch; +#[derive(Clone, Debug, PartialEq)] +pub(crate) enum TurnInput { + UserInput(Vec), + ResponseInputItem(ResponseInputItem), +} + /// Turn-local pending input storage owned by the input queue flow. #[derive(Default)] pub(crate) struct TurnInputQueue { - items: Vec, + items: Vec, } /// Session-scoped pending input storage and active-turn mailbox delivery coordination. @@ -151,7 +158,7 @@ impl InputQueue { pub(super) async fn push_pending_input_and_accept_mailbox_delivery_for_turn_state( &self, turn_state: &Mutex, - input: ResponseInputItem, + input: TurnInput, ) { let mut turn_state = turn_state.lock().await; turn_state.pending_input.items.push(input); @@ -161,7 +168,7 @@ impl InputQueue { pub(crate) async fn extend_pending_input_for_turn_state( &self, turn_state: &Mutex, - input: Vec, + input: Vec, ) { turn_state.lock().await.pending_input.items.extend(input); } @@ -169,7 +176,7 @@ impl InputQueue { pub(crate) async fn take_pending_input_for_turn_state( &self, turn_state: &Mutex, - ) -> Vec { + ) -> Vec { turn_state.lock().await.pending_input.items.split_off(0) } @@ -185,43 +192,20 @@ impl InputQueue { let mut active = active_turn.lock().await; match active.as_mut() { Some(active_turn) => { - active_turn - .turn_state - .lock() - .await - .pending_input - .items - .extend(input); + self.extend_pending_input_for_turn_state( + active_turn.turn_state.as_ref(), + input + .into_iter() + .map(TurnInput::ResponseInputItem) + .collect(), + ) + .await; Ok(()) } None => Err(input), } } - #[expect( - clippy::await_holding_invalid_type, - reason = "active turn checks and turn state updates must remain atomic" - )] - pub(crate) async fn prepend_pending_input( - &self, - active_turn: &Mutex>, - mut input: Vec, - ) -> Result<(), ()> { - let mut active = active_turn.lock().await; - match active.as_mut() { - Some(active_turn) => { - let mut turn_state = active_turn.turn_state.lock().await; - if !input.is_empty() { - let pending_input = &mut turn_state.pending_input; - input.append(&mut pending_input.items); - pending_input.items = input; - } - Ok(()) - } - None => Err(()), - } - } - #[expect( clippy::await_holding_invalid_type, reason = "active turn checks and turn state updates must remain atomic" @@ -229,7 +213,7 @@ impl InputQueue { pub(crate) async fn get_pending_input( &self, active_turn: &Mutex>, - ) -> Vec { + ) -> Vec { let (pending_input, accepts_mailbox_delivery) = { let mut active = active_turn.lock().await; match active.as_mut() { @@ -246,11 +230,13 @@ impl InputQueue { if !accepts_mailbox_delivery { return pending_input; } - let mailbox_items = self.drain_mailbox_input_items().await; + let mailbox_items = self + .drain_mailbox_input_items() + .await + .into_iter() + .map(TurnInput::ResponseInputItem); if pending_input.is_empty() { - mailbox_items - } else if mailbox_items.is_empty() { - pending_input + mailbox_items.collect() } else { let mut pending_input = pending_input; pending_input.extend(mailbox_items); diff --git a/codex-rs/core/src/session/mod.rs b/codex-rs/core/src/session/mod.rs index 9609a9d1b..8f0d2f1c0 100644 --- a/codex-rs/core/src/session/mod.rs +++ b/codex-rs/core/src/session/mod.rs @@ -207,6 +207,7 @@ use self::config_lock::validate_config_lock_if_configured; #[cfg(test)] use self::handlers::submission_dispatch_span; use self::handlers::submission_loop; +pub(crate) use self::input_queue::TurnInput; pub(crate) use self::input_queue::TurnInputQueue; use self::review::spawn_review_thread; use self::session::AppServerClientMetadata; @@ -3099,11 +3100,11 @@ impl Session { &self, turn_context: &TurnContext, input: &[UserInput], - response_item: ResponseItem, ) { // Persist the user message to history, but emit the turn item from `UserInput` so // UI-only `text_elements` are preserved. `ResponseItem::Message` does not carry // those spans, and `record_response_item_and_emit_turn_item` would drop them. + let response_item = ResponseItem::from(ResponseInputItem::from(input.to_vec())); self.record_conversation_items(turn_context, std::slice::from_ref(&response_item)) .await; let turn_item = TurnItem::UserMessage(UserMessageItem::new(input)); @@ -3192,7 +3193,7 @@ impl Session { self.input_queue .push_pending_input_and_accept_mailbox_delivery_for_turn_state( active_turn.turn_state.as_ref(), - input.into(), + TurnInput::UserInput(input), ) .await; Ok(active_turn_id.clone()) diff --git a/codex-rs/core/src/session/tests.rs b/codex-rs/core/src/session/tests.rs index f0cdc70fd..c41d5e348 100644 --- a/codex-rs/core/src/session/tests.rs +++ b/codex-rs/core/src/session/tests.rs @@ -7724,15 +7724,21 @@ async fn task_finish_emits_turn_item_lifecycle_for_leftover_pending_user_input() while rx.try_recv().is_ok() {} - sess.inject_response_items(vec![ResponseInputItem::Message { - role: "user".to_string(), - content: vec![ContentItem::InputText { - text: "late pending input".to_string(), - }], - phase: None, - }]) + let text_element = codex_protocol::user_input::TextElement::new( + codex_protocol::user_input::ByteRange { start: 5, end: 12 }, + Some("pending marker".to_string()), + ); + let pending_user_input = vec![UserInput::Text { + text: "late pending input".to_string(), + text_elements: vec![text_element.clone()], + }]; + sess.steer_input( + pending_user_input.clone(), + Some(&tc.sub_id), + /*responsesapi_client_metadata*/ None, + ) .await - .expect("inject pending input into active turn"); + .expect("steer pending input into active turn"); sess.on_task_finished(Arc::clone(&tc), /*last_agent_message*/ None) .await; @@ -7766,10 +7772,7 @@ async fn task_finish_emits_turn_item_lifecycle_for_leftover_pending_user_input() EventMsg::ItemStarted(ItemStartedEvent { item: TurnItem::UserMessage(UserMessageItem { content, .. }), .. - }) if content == vec![UserInput::Text { - text: "late pending input".to_string(), - text_elements: Vec::new(), - }] + }) if content == pending_user_input )); let third = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv()) @@ -7781,10 +7784,7 @@ async fn task_finish_emits_turn_item_lifecycle_for_leftover_pending_user_input() EventMsg::ItemCompleted(ItemCompletedEvent { item: TurnItem::UserMessage(UserMessageItem { content, .. }), .. - }) if content == vec![UserInput::Text { - text: "late pending input".to_string(), - text_elements: Vec::new(), - }] + }) if content == pending_user_input )); let fourth = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv()) @@ -7801,7 +7801,7 @@ async fn task_finish_emits_turn_item_lifecycle_for_leftover_pending_user_input() .. }) if message == "late pending input" && images == Some(Vec::new()) - && text_elements.is_empty() + && text_elements == vec![text_element] && local_images.is_empty() )); @@ -7954,69 +7954,6 @@ async fn steer_input_returns_active_turn_id() { assert!(sess.input_queue.has_pending_input(&sess.active_turn).await); } -#[tokio::test] -async fn prepend_pending_input_keeps_older_tail_ahead_of_newer_input() { - let (sess, tc, _rx) = make_session_and_context_with_rx().await; - let input = vec![UserInput::Text { - text: "hello".to_string(), - text_elements: Vec::new(), - }]; - sess.spawn_task( - Arc::clone(&tc), - input, - NeverEndingTask { - kind: TaskKind::Regular, - listen_to_cancellation_token: false, - }, - ) - .await; - - let blocked = ResponseInputItem::Message { - role: "user".to_string(), - content: vec![ContentItem::InputText { - text: "blocked queued prompt".to_string(), - }], - phase: None, - }; - let later = ResponseInputItem::Message { - role: "user".to_string(), - content: vec![ContentItem::InputText { - text: "later queued prompt".to_string(), - }], - phase: None, - }; - let newer = ResponseInputItem::Message { - role: "user".to_string(), - content: vec![ContentItem::InputText { - text: "newer queued prompt".to_string(), - }], - phase: None, - }; - - sess.inject_response_items(vec![blocked.clone(), later.clone()]) - .await - .expect("inject initial pending input into active turn"); - - let drained = sess.input_queue.get_pending_input(&sess.active_turn).await; - assert_eq!(drained, vec![blocked, later.clone()]); - - sess.inject_response_items(vec![newer.clone()]) - .await - .expect("inject newer pending input into active turn"); - - let mut drained_iter = drained.into_iter(); - let _blocked = drained_iter.next().expect("blocked prompt should exist"); - sess.input_queue - .prepend_pending_input(&sess.active_turn, drained_iter.collect()) - .await - .expect("requeue later pending input at the front of the queue"); - - assert_eq!( - sess.input_queue.get_pending_input(&sess.active_turn).await, - vec![later, newer] - ); -} - #[tokio::test] async fn queued_response_items_for_next_turn_move_into_next_active_turn() { let (sess, tc, _rx) = make_session_and_context_with_rx().await; @@ -8044,7 +7981,7 @@ async fn queued_response_items_for_next_turn_move_into_next_active_turn() { assert_eq!( sess.input_queue.get_pending_input(&sess.active_turn).await, - vec![queued_item] + vec![TurnInput::ResponseInputItem(queued_item)] ); } @@ -8089,7 +8026,10 @@ async fn abort_empty_active_turn_preserves_pending_input() { Arc::clone(&active_turn.turn_state) }; sess.input_queue - .extend_pending_input_for_turn_state(turn_state.as_ref(), vec![pending_item.clone()]) + .extend_pending_input_for_turn_state( + turn_state.as_ref(), + vec![TurnInput::ResponseInputItem(pending_item.clone())], + ) .await; sess.abort_all_tasks(TurnAbortReason::Replaced).await; @@ -8099,7 +8039,7 @@ async fn abort_empty_active_turn_preserves_pending_input() { sess.input_queue .take_pending_input_for_turn_state(turn_state.as_ref()) .await, - vec![pending_item] + vec![TurnInput::ResponseInputItem(pending_item)] ); } @@ -8523,7 +8463,9 @@ async fn budget_limited_accounting_steers_active_turn_without_aborting() -> anyh .await?; let pending_input = sess.input_queue.get_pending_input(&sess.active_turn).await; - let [ResponseInputItem::Message { role, content, .. }] = pending_input.as_slice() else { + let [TurnInput::ResponseInputItem(ResponseInputItem::Message { role, content, .. })] = + pending_input.as_slice() + else { panic!("expected one budget-limit steering message, got {pending_input:#?}"); }; assert_eq!("user", role); @@ -8748,7 +8690,7 @@ async fn external_objective_change_steers_active_turn() -> anyhow::Result<()> { pending_input.iter().any(|item| { matches!( item, - ResponseInputItem::Message { role, content, .. } + TurnInput::ResponseInputItem(ResponseInputItem::Message { role, content, .. }) if role == "user" && content.iter().any(|content| matches!( content, @@ -8971,7 +8913,9 @@ async fn queue_only_mailbox_mail_waits_for_next_turn_after_answer_boundary() { assert_eq!( sess.input_queue.get_pending_input(&sess.active_turn).await, - vec![communication.to_response_input_item()], + vec![TurnInput::ResponseInputItem( + communication.to_response_input_item() + )], ); } @@ -9051,11 +8995,11 @@ async fn steered_input_reopens_mailbox_delivery_for_current_turn() { assert_eq!( sess.input_queue.get_pending_input(&sess.active_turn).await, vec![ - ResponseInputItem::from(vec![UserInput::Text { + TurnInput::UserInput(vec![UserInput::Text { text: "follow up".to_string(), text_elements: Vec::new(), }]), - communication.to_response_input_item(), + TurnInput::ResponseInputItem(communication.to_response_input_item()), ], ); } @@ -9104,11 +9048,11 @@ async fn stale_defer_mailbox_delivery_does_not_override_steered_input() { assert_eq!( sess.input_queue.get_pending_input(&sess.active_turn).await, vec![ - ResponseInputItem::from(vec![UserInput::Text { + TurnInput::UserInput(vec![UserInput::Text { text: "follow up".to_string(), text_elements: Vec::new(), }]), - communication.to_response_input_item(), + TurnInput::ResponseInputItem(communication.to_response_input_item()), ], ); } @@ -9163,7 +9107,9 @@ async fn tool_calls_reopen_mailbox_delivery_for_current_turn() { assert!(output.tool_future.is_some()); assert_eq!( sess.input_queue.get_pending_input(&sess.active_turn).await, - vec![communication.to_response_input_item()], + vec![TurnInput::ResponseInputItem( + communication.to_response_input_item() + )], ); } diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index fd646631f..4ac511f4f 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -18,14 +18,12 @@ use crate::connectors; use crate::context::ContextualUserFragment; use crate::feedback_tags; use crate::goals::GoalRuntimeEvent; -use crate::hook_runtime::PendingInputHookDisposition; use crate::hook_runtime::inspect_pending_input; use crate::hook_runtime::record_additional_contexts; use crate::hook_runtime::record_pending_input; use crate::hook_runtime::run_legacy_after_agent_hook; use crate::hook_runtime::run_pending_session_start_hooks; use crate::hook_runtime::run_stop_hooks; -use crate::hook_runtime::run_user_prompt_submit_hooks; use crate::injection::ToolMentionKind; use crate::injection::app_id_from_path; use crate::injection::tool_kind_for_path; @@ -38,6 +36,7 @@ use crate::mentions::collect_explicit_plugin_mentions; use crate::mentions::collect_tool_mentions_from_messages; use crate::plugins::build_plugin_injections; use crate::session::PreviousTurnSettings; +use crate::session::TurnInput; use crate::session::session::Session; use crate::session::turn_context::TurnContext; use crate::stream_events_utils::HandleOutputCtx; @@ -78,7 +77,6 @@ use codex_protocol::error::CodexErr; use codex_protocol::error::Result as CodexResult; use codex_protocol::items::PlanItem; use codex_protocol::items::TurnItem; -use codex_protocol::items::UserMessageItem; use codex_protocol::items::build_hook_prompt_message; use codex_protocol::models::BaseInstructions; use codex_protocol::models::ContentItem; @@ -174,17 +172,10 @@ pub(crate) async fn run_turn( if run_pending_session_start_hooks(&sess, &turn_context).await { return None; } - let additional_contexts = if input.is_empty() { - Vec::new() - } else { - let initial_input_for_turn: ResponseInputItem = ResponseInputItem::from(input.clone()); - let response_item: ResponseItem = initial_input_for_turn.clone().into(); - let user_prompt_submit_outcome = run_user_prompt_submit_hooks( - &sess, - &turn_context, - UserMessageItem::new(&input).message(), - ) - .await; + if !input.is_empty() { + let initial_turn_input = TurnInput::UserInput(input.clone()); + let user_prompt_submit_outcome = + inspect_pending_input(&sess, &turn_context, &initial_turn_input).await; if user_prompt_submit_outcome.should_stop { record_additional_contexts( &sess, @@ -194,13 +185,17 @@ pub(crate) async fn run_turn( .await; return None; } - sess.record_user_prompt_and_emit_turn_item(turn_context.as_ref(), &input, response_item) - .await; - user_prompt_submit_outcome.additional_contexts - }; + record_pending_input( + &sess, + &turn_context, + initial_turn_input, + user_prompt_submit_outcome.additional_contexts, + ) + .await; + } + sess.merge_connector_selection(explicitly_enabled_connectors.clone()) .await; - record_additional_contexts(&sess, &turn_context, additional_contexts).await; sess.set_previous_turn_settings(Some(PreviousTurnSettings { model: turn_context.model_info.slug.clone(), realtime_active: Some(turn_context.realtime_active), @@ -252,45 +247,26 @@ pub(crate) async fn run_turn( }; let mut blocked_pending_input = false; - let mut blocked_pending_input_contexts = Vec::new(); - let mut requeued_pending_input = false; - let mut accepted_pending_input = Vec::new(); - if !pending_input.is_empty() { - let mut pending_input_iter = pending_input.into_iter(); - while let Some(pending_input_item) = pending_input_iter.next() { - match inspect_pending_input(&sess, &turn_context, pending_input_item).await { - PendingInputHookDisposition::Accepted(pending_input) => { - accepted_pending_input.push(*pending_input); - } - PendingInputHookDisposition::Blocked { - additional_contexts, - } => { - let remaining_pending_input = pending_input_iter.collect::>(); - if !remaining_pending_input.is_empty() { - let _ = sess - .input_queue - .prepend_pending_input(&sess.active_turn, remaining_pending_input) - .await; - requeued_pending_input = true; - } - blocked_pending_input_contexts = additional_contexts; - blocked_pending_input = true; - break; - } - } + let mut accepted_pending_input = false; + for pending_input_item in pending_input { + let hook_outcome = + inspect_pending_input(&sess, &turn_context, &pending_input_item).await; + if hook_outcome.should_stop { + blocked_pending_input = true; + record_additional_contexts(&sess, &turn_context, hook_outcome.additional_contexts) + .await; + } else { + accepted_pending_input = true; + record_pending_input( + &sess, + &turn_context, + pending_input_item, + hook_outcome.additional_contexts, + ) + .await; } } - - let has_accepted_pending_input = !accepted_pending_input.is_empty(); - for pending_input in accepted_pending_input { - record_pending_input(&sess, &turn_context, pending_input).await; - } - record_additional_contexts(&sess, &turn_context, blocked_pending_input_contexts).await; - - if blocked_pending_input && !has_accepted_pending_input { - if requeued_pending_input { - continue; - } + if blocked_pending_input && !accepted_pending_input { break; } diff --git a/codex-rs/core/src/tasks/mod.rs b/codex-rs/core/src/tasks/mod.rs index ba7f70792..164e65f07 100644 --- a/codex-rs/core/src/tasks/mod.rs +++ b/codex-rs/core/src/tasks/mod.rs @@ -24,10 +24,10 @@ use tracing::warn; use crate::config::Config; use crate::context::ContextualUserFragment; use crate::goals::GoalRuntimeEvent; -use crate::hook_runtime::PendingInputHookDisposition; use crate::hook_runtime::inspect_pending_input; use crate::hook_runtime::record_additional_contexts; use crate::hook_runtime::record_pending_input; +use crate::session::TurnInput; use crate::session::session::Session; use crate::session::turn_context::TurnContext; use crate::state::ActiveTurn; @@ -42,7 +42,6 @@ use codex_otel::TURN_MEMORY_METRIC; use codex_otel::TURN_NETWORK_PROXY_METRIC; use codex_otel::TURN_TOKEN_USAGE_METRIC; use codex_otel::TURN_TOOL_CALL_METRIC; -use codex_protocol::models::ResponseInputItem; use codex_protocol::models::ResponseItem; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::RolloutItem; @@ -359,7 +358,10 @@ impl Session { Arc::clone(&turn.turn_state) }; turn_state.lock().await.token_usage_at_turn_start = token_usage_at_turn_start; - let mut pending_items = queued_response_items; + let mut pending_items = queued_response_items + .into_iter() + .map(TurnInput::ResponseInputItem) + .collect::>(); pending_items.extend(mailbox_items); self.input_queue .extend_pending_input_for_turn_state(turn_state.as_ref(), pending_items) @@ -587,7 +589,7 @@ impl Session { .turn_metadata_state .cancel_git_enrichment_task(); - let mut pending_input = Vec::::new(); + let mut pending_input = Vec::::new(); let mut should_clear_active_turn = false; let mut token_usage_at_turn_start = None; let mut turn_had_memory_citation = false; @@ -622,15 +624,23 @@ impl Session { } if !pending_input.is_empty() { for pending_input_item in pending_input { - match inspect_pending_input(self, &turn_context, pending_input_item).await { - PendingInputHookDisposition::Accepted(pending_input) => { - record_pending_input(self, &turn_context, *pending_input).await; - } - PendingInputHookDisposition::Blocked { - additional_contexts, - } => { - record_additional_contexts(self, &turn_context, additional_contexts).await; - } + let hook_outcome = + inspect_pending_input(self, &turn_context, &pending_input_item).await; + if hook_outcome.should_stop { + record_additional_contexts( + self, + &turn_context, + hook_outcome.additional_contexts, + ) + .await; + } else { + record_pending_input( + self, + &turn_context, + pending_input_item, + hook_outcome.additional_contexts, + ) + .await; } } }