From 1224f44c15cb15b644b59c71b45a65d6c1de6118 Mon Sep 17 00:00:00 2001 From: Eric Traut Date: Wed, 27 May 2026 01:04:54 -0700 Subject: [PATCH 1/2] Add extension hooks for goal extraction --- .../src/context/contextual_user_message.rs | 4 + .../context/contextual_user_message_tests.rs | 10 ++ .../core/src/context/extension_context.rs | 31 +++++ codex-rs/core/src/context/mod.rs | 2 + codex-rs/core/src/session/turn.rs | 32 +++++- codex-rs/core/src/tasks/idle_extension.rs | 19 ++++ codex-rs/core/src/tasks/lifecycle.rs | 19 ++++ codex-rs/core/src/tasks/mod.rs | 107 ++++++++++++++++++ codex-rs/core/src/tools/router.rs | 11 +- codex-rs/core/src/tools/router_tests.rs | 2 +- .../ext/extension-api/src/contributors.rs | 50 ++++++++ .../src/contributors/turn_lifecycle.rs | 15 +++ codex-rs/ext/extension-api/src/lib.rs | 5 + codex-rs/ext/extension-api/src/registry.rs | 15 +++ 14 files changed, 313 insertions(+), 9 deletions(-) create mode 100644 codex-rs/core/src/context/extension_context.rs create mode 100644 codex-rs/core/src/tasks/idle_extension.rs diff --git a/codex-rs/core/src/context/contextual_user_message.rs b/codex-rs/core/src/context/contextual_user_message.rs index 940a938fce..8523b083eb 100644 --- a/codex-rs/core/src/context/contextual_user_message.rs +++ b/codex-rs/core/src/context/contextual_user_message.rs @@ -4,6 +4,7 @@ use codex_protocol::models::ContentItem; use super::AdditionalContextUserFragment; use super::EnvironmentContext; +use super::ExtensionContext; use super::FragmentRegistration; use super::FragmentRegistrationProxy; use super::GoalContext; @@ -22,6 +23,8 @@ static ENVIRONMENT_CONTEXT_REGISTRATION: FragmentRegistrationProxy = FragmentRegistrationProxy::new(); +static EXTENSION_CONTEXT_REGISTRATION: FragmentRegistrationProxy = + FragmentRegistrationProxy::new(); static SKILL_INSTRUCTIONS_REGISTRATION: FragmentRegistrationProxy = FragmentRegistrationProxy::new(); static USER_SHELL_COMMAND_REGISTRATION: FragmentRegistrationProxy = @@ -46,6 +49,7 @@ static CONTEXTUAL_USER_FRAGMENTS: &[&dyn FragmentRegistration] = &[ &USER_INSTRUCTIONS_REGISTRATION, &ENVIRONMENT_CONTEXT_REGISTRATION, &ADDITIONAL_CONTEXT_REGISTRATION, + &EXTENSION_CONTEXT_REGISTRATION, &SKILL_INSTRUCTIONS_REGISTRATION, &USER_SHELL_COMMAND_REGISTRATION, &TURN_ABORTED_REGISTRATION, diff --git a/codex-rs/core/src/context/contextual_user_message_tests.rs b/codex-rs/core/src/context/contextual_user_message_tests.rs index 195b189929..58714cb6da 100644 --- a/codex-rs/core/src/context/contextual_user_message_tests.rs +++ b/codex-rs/core/src/context/contextual_user_message_tests.rs @@ -1,5 +1,6 @@ use super::*; use crate::context::ContextualUserFragment; +use crate::context::ExtensionContext; use crate::context::GoalContext; use crate::context::SubagentNotification; use codex_protocol::items::HookPromptFragment; @@ -38,6 +39,15 @@ fn detects_goal_context_fragment() { })); } +#[test] +fn detects_extension_context_fragment() { + let text = ExtensionContext::new("\nContinue working.\n").render(); + + assert!(is_contextual_user_fragment(&ContentItem::InputText { + text + })); +} + #[test] fn contextual_user_fragment_is_dyn_compatible() { let fragment: Box = Box::new(GoalContext::new( diff --git a/codex-rs/core/src/context/extension_context.rs b/codex-rs/core/src/context/extension_context.rs new file mode 100644 index 0000000000..5f40992237 --- /dev/null +++ b/codex-rs/core/src/context/extension_context.rs @@ -0,0 +1,31 @@ +use super::ContextualUserFragment; + +/// Hidden user-context fragment for extension-owned steering prompts. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ExtensionContext { + body: String, +} + +impl ExtensionContext { + pub fn new(body: impl Into) -> Self { + Self { body: body.into() } + } +} + +impl ContextualUserFragment for ExtensionContext { + fn role() -> &'static str { + "user" + } + + fn markers(&self) -> (&'static str, &'static str) { + Self::type_markers() + } + + fn type_markers() -> (&'static str, &'static str) { + ("", "") + } + + fn body(&self) -> String { + format!("\n{}\n", self.body) + } +} diff --git a/codex-rs/core/src/context/mod.rs b/codex-rs/core/src/context/mod.rs index 6d88df87e2..3a5553ba86 100644 --- a/codex-rs/core/src/context/mod.rs +++ b/codex-rs/core/src/context/mod.rs @@ -7,6 +7,7 @@ mod available_skills_instructions; mod collaboration_mode_instructions; mod contextual_user_message; mod environment_context; +mod extension_context; mod fragment; mod fragments; mod goal_context; @@ -38,6 +39,7 @@ pub(crate) use collaboration_mode_instructions::CollaborationModeInstructions; pub(crate) use contextual_user_message::is_contextual_user_fragment; pub(crate) use contextual_user_message::parse_visible_hook_prompt_message; pub(crate) use environment_context::EnvironmentContext; +pub use extension_context::ExtensionContext; pub use fragment::ContextualUserFragment; pub(crate) use fragment::FragmentRegistration; pub(crate) use fragment::FragmentRegistrationProxy; diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index 2b9fa02529..d3d2fe8d5a 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -145,7 +145,15 @@ pub(crate) async fn run_turn( // diffs/full reinjection + user input) and trigger compaction preemptively // when they would push the thread over the compaction threshold. if let Err(err) = run_pre_sampling_compact(&sess, &turn_context, &mut client_session).await { - if err.to_codex_protocol_error() == CodexErrorInfo::UsageLimitExceeded + let error_info = err.to_codex_protocol_error(); + if error_info == CodexErrorInfo::UsageLimitExceeded { + sess.emit_turn_error_lifecycle( + turn_context.as_ref(), + CodexErrorInfo::UsageLimitExceeded, + ) + .await; + } + if error_info == CodexErrorInfo::UsageLimitExceeded && let Err(err) = sess .goal_runtime_apply(GoalRuntimeEvent::UsageLimitReached { turn_context: turn_context.as_ref(), @@ -292,7 +300,15 @@ pub(crate) async fn run_turn( ) .await { - if err.to_codex_protocol_error() == CodexErrorInfo::UsageLimitExceeded + let error_info = err.to_codex_protocol_error(); + if error_info == CodexErrorInfo::UsageLimitExceeded { + sess.emit_turn_error_lifecycle( + turn_context.as_ref(), + CodexErrorInfo::UsageLimitExceeded, + ) + .await; + } + if error_info == CodexErrorInfo::UsageLimitExceeded && let Err(err) = sess .goal_runtime_apply(GoalRuntimeEvent::UsageLimitReached { turn_context: turn_context.as_ref(), @@ -381,7 +397,15 @@ pub(crate) async fn run_turn( } Err(e) => { info!("Turn error: {e:#}"); - if e.to_codex_protocol_error() == CodexErrorInfo::UsageLimitExceeded + let error_info = e.to_codex_protocol_error(); + if error_info == CodexErrorInfo::UsageLimitExceeded { + sess.emit_turn_error_lifecycle( + turn_context.as_ref(), + CodexErrorInfo::UsageLimitExceeded, + ) + .await; + } + if error_info == CodexErrorInfo::UsageLimitExceeded && let Err(err) = sess .goal_runtime_apply(GoalRuntimeEvent::UsageLimitReached { turn_context: turn_context.as_ref(), @@ -1102,7 +1126,7 @@ pub(crate) async fn built_tools( mcp_tools, deferred_mcp_tools, discoverable_tools, - extension_tool_executors: extension_tool_executors(sess), + extension_tool_executors: extension_tool_executors(sess, turn_context), dynamic_tools: turn_context.dynamic_tools.as_slice(), }, ))) diff --git a/codex-rs/core/src/tasks/idle_extension.rs b/codex-rs/core/src/tasks/idle_extension.rs new file mode 100644 index 0000000000..03a3953b34 --- /dev/null +++ b/codex-rs/core/src/tasks/idle_extension.rs @@ -0,0 +1,19 @@ +use std::sync::Arc; + +use crate::session::session::Session; + +pub(super) fn schedule_turn(session: &Arc) { + if session + .services + .extensions + .idle_turn_contributors() + .is_empty() + { + return; + } + + let session = Arc::clone(session); + let _handle = tokio::spawn(async move { + session.maybe_start_idle_extension_turn().await; + }); +} diff --git a/codex-rs/core/src/tasks/lifecycle.rs b/codex-rs/core/src/tasks/lifecycle.rs index 8b934175cf..882e787608 100644 --- a/codex-rs/core/src/tasks/lifecycle.rs +++ b/codex-rs/core/src/tasks/lifecycle.rs @@ -1,4 +1,5 @@ use codex_extension_api::ExtensionData; +use codex_protocol::protocol::CodexErrorInfo; use codex_protocol::protocol::TokenUsage; use codex_protocol::protocol::TurnAbortReason; @@ -53,4 +54,22 @@ impl Session { .await; } } + + pub(crate) async fn emit_turn_error_lifecycle( + &self, + turn_context: &TurnContext, + error: CodexErrorInfo, + ) { + for contributor in self.services.extensions.turn_lifecycle_contributors() { + contributor + .on_turn_error(codex_extension_api::TurnErrorInput { + turn_id: turn_context.sub_id.as_str(), + error: error.clone(), + session_store: &self.services.session_extension_data, + thread_store: &self.services.thread_extension_data, + turn_store: turn_context.extension_data.as_ref(), + }) + .await; + } + } } diff --git a/codex-rs/core/src/tasks/mod.rs b/codex-rs/core/src/tasks/mod.rs index aa161c1274..3a5484fe75 100644 --- a/codex-rs/core/src/tasks/mod.rs +++ b/codex-rs/core/src/tasks/mod.rs @@ -1,4 +1,5 @@ mod compact; +mod idle_extension; mod lifecycle; mod regular; mod review; @@ -33,6 +34,7 @@ use crate::session::turn_context::TurnContext; use crate::state::ActiveTurn; use crate::state::RunningTask; use crate::state::TaskKind; +use crate::state::TurnState; use codex_analytics::TurnTokenUsageFact; use codex_login::AuthManager; use codex_models_manager::manager::SharedModelsManager; @@ -495,6 +497,110 @@ impl Session { .await; } + pub(crate) async fn maybe_start_idle_extension_turn(self: &Arc) { + if self.services.extensions.idle_turn_contributors().is_empty() { + return; + } + self.maybe_start_turn_for_pending_work().await; + if self.active_turn.lock().await.is_some() { + return; + } + if self + .input_queue + .has_queued_response_items_for_next_turn() + .await + || self.input_queue.has_trigger_turn_mailbox_items().await + { + return; + } + + let items = { + let collaboration_mode = self.collaboration_mode().await; + let mut requested_items = None; + for contributor in self.services.extensions.idle_turn_contributors() { + if let Some(items) = contributor + .next_idle_turn(codex_extension_api::IdleTurnInput { + collaboration_mode: &collaboration_mode, + session_store: &self.services.session_extension_data, + thread_store: &self.services.thread_extension_data, + }) + .await + { + requested_items = Some(items); + break; + } + } + requested_items + }; + let Some(items) = items.filter(|items| !items.is_empty()) else { + return; + }; + + let turn_state = { + let mut active_turn = self.active_turn.lock().await; + if active_turn.is_some() { + return; + } + let active_turn = active_turn.get_or_insert_with(ActiveTurn::default); + Arc::clone(&active_turn.turn_state) + }; + if self + .input_queue + .has_queued_response_items_for_next_turn() + .await + || self.input_queue.has_trigger_turn_mailbox_items().await + { + self.clear_reserved_idle_extension_turn(&turn_state).await; + self.maybe_start_turn_for_pending_work().await; + return; + } + + self.input_queue + .extend_pending_input_for_turn_state( + turn_state.as_ref(), + items + .into_iter() + .map(TurnInput::ResponseInputItem) + .collect(), + ) + .await; + + let turn_context = self + .new_default_turn_with_sub_id(uuid::Uuid::new_v4().to_string()) + .await; + self.maybe_emit_unknown_model_warning_for_turn(turn_context.as_ref()) + .await; + let still_reserved = { + let active_turn = self.active_turn.lock().await; + active_turn.as_ref().is_some_and(|active_turn| { + active_turn.task.is_none() && Arc::ptr_eq(&active_turn.turn_state, &turn_state) + }) + }; + if !still_reserved { + self.input_queue + .take_pending_input_for_turn_state(turn_state.as_ref()) + .await; + self.clear_reserved_idle_extension_turn(&turn_state).await; + return; + } + drop(turn_state); + self.start_task(turn_context, Vec::new(), RegularTask::new()) + .await; + } + + async fn clear_reserved_idle_extension_turn( + &self, + turn_state: &Arc>, + ) { + let mut active_turn_guard = self.active_turn.lock().await; + if let Some(active_turn) = active_turn_guard.as_ref() + && active_turn.task.is_none() + && Arc::ptr_eq(&active_turn.turn_state, turn_state) + { + *active_turn_guard = None; + } + } + pub async fn abort_all_tasks(self: &Arc, reason: TurnAbortReason) { let mut aborted_turn = false; let mut active_turn_to_clear = None; @@ -809,6 +915,7 @@ impl Session { { warn!("failed to apply goal runtime maybe-continue event: {err}"); } + idle_extension::schedule_turn(self); } async fn take_active_turn(&self) -> Option { diff --git a/codex-rs/core/src/tools/router.rs b/codex-rs/core/src/tools/router.rs index cfd4776890..dd55c622d5 100644 --- a/codex-rs/core/src/tools/router.rs +++ b/codex-rs/core/src/tools/router.rs @@ -219,6 +219,7 @@ impl ToolRouter { pub(crate) fn extension_tool_executors( session: &Session, + turn_context: &TurnContext, ) -> Vec>> { session .services @@ -226,10 +227,12 @@ pub(crate) fn extension_tool_executors( .tool_contributors() .iter() .flat_map(|contributor| { - contributor.tools( - &session.services.session_extension_data, - &session.services.thread_extension_data, - ) + contributor.tools_for_turn(codex_extension_api::ToolContributionInput { + session_store: &session.services.session_extension_data, + thread_store: &session.services.thread_extension_data, + session_source: &turn_context.session_source, + persistent_thread: !turn_context.config.ephemeral && session.state_db().is_some(), + }) }) .collect() } diff --git a/codex-rs/core/src/tools/router_tests.rs b/codex-rs/core/src/tools/router_tests.rs index 9d685e53b6..f94b37be8b 100644 --- a/codex-rs/core/src/tools/router_tests.rs +++ b/codex-rs/core/src/tools/router_tests.rs @@ -347,7 +347,7 @@ async fn extension_tool_executors_are_model_visible_and_dispatchable() -> anyhow deferred_mcp_tools: None, mcp_tools: None, discoverable_tools: None, - extension_tool_executors: extension_tool_executors(&session), + extension_tool_executors: extension_tool_executors(&session, &turn), dynamic_tools: turn.dynamic_tools.as_slice(), }, ); diff --git a/codex-rs/ext/extension-api/src/contributors.rs b/codex-rs/ext/extension-api/src/contributors.rs index 2113fdef8e..ee15b82bc3 100644 --- a/codex-rs/ext/extension-api/src/contributors.rs +++ b/codex-rs/ext/extension-api/src/contributors.rs @@ -1,8 +1,12 @@ use std::future::Future; +use std::pin::Pin; use std::sync::Arc; +use codex_protocol::config_types::CollaborationMode; use codex_protocol::items::TurnItem; +use codex_protocol::models::ResponseInputItem; use codex_protocol::protocol::ReviewDecision; +use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::TokenUsageInfo; use codex_tools::ToolCall; use codex_tools::ToolExecutor; @@ -25,6 +29,7 @@ pub use tool_lifecycle::ToolFinishInput; pub use tool_lifecycle::ToolLifecycleFuture; pub use tool_lifecycle::ToolStartInput; pub use turn_lifecycle::TurnAbortInput; +pub use turn_lifecycle::TurnErrorInput; pub use turn_lifecycle::TurnStartInput; pub use turn_lifecycle::TurnStopInput; @@ -71,6 +76,10 @@ pub trait TurnLifecycleContributor: Send + Sync { /// Called after the host aborts a running turn. async fn on_turn_abort(&self, _input: TurnAbortInput<'_>) {} + + /// Called when a running turn encounters a host-classified error before it + /// can complete normally. + async fn on_turn_error(&self, _input: TurnErrorInput<'_>) {} } /// Contributor for host-owned configuration changes. @@ -115,6 +124,47 @@ pub trait ToolContributor: Send + Sync { session_store: &ExtensionData, thread_store: &ExtensionData, ) -> Vec>>; + + /// Returns the native tools visible for the current turn, when the host has + /// turn context available. + fn tools_for_turn( + &self, + input: ToolContributionInput<'_>, + ) -> Vec>> { + self.tools(input.session_store, input.thread_store) + } +} + +/// Context supplied while collecting extension-owned tools for a turn. +pub struct ToolContributionInput<'a> { + /// Store scoped to the host session runtime. + pub session_store: &'a ExtensionData, + /// Store scoped to this thread runtime. + pub thread_store: &'a ExtensionData, + /// Source for the current turn's session. + pub session_source: &'a SessionSource, + /// Whether the current thread has persistent state available. + pub persistent_thread: bool, +} + +/// Context supplied when the host is idle and extensions may request work. +pub struct IdleTurnInput<'a> { + /// Effective collaboration mode for the next default turn. + pub collaboration_mode: &'a CollaborationMode, + /// Store scoped to the host session runtime. + pub session_store: &'a ExtensionData, + /// Store scoped to this thread runtime. + pub thread_store: &'a ExtensionData, +} + +pub type IdleTurnFuture<'a> = + Pin>> + Send + 'a>>; + +/// Extension contribution that can request host-owned work while a thread is idle. +pub trait IdleTurnContributor: Send + Sync { + fn next_idle_turn<'a>(&'a self, _input: IdleTurnInput<'a>) -> IdleTurnFuture<'a> { + Box::pin(std::future::ready(None)) + } } /// Contributor for host-owned tool lifecycle gates. diff --git a/codex-rs/ext/extension-api/src/contributors/turn_lifecycle.rs b/codex-rs/ext/extension-api/src/contributors/turn_lifecycle.rs index bbd3ae8f39..de1df81cdb 100644 --- a/codex-rs/ext/extension-api/src/contributors/turn_lifecycle.rs +++ b/codex-rs/ext/extension-api/src/contributors/turn_lifecycle.rs @@ -1,4 +1,5 @@ use codex_protocol::config_types::CollaborationMode; +use codex_protocol::protocol::CodexErrorInfo; use codex_protocol::protocol::TokenUsage; use codex_protocol::protocol::TurnAbortReason; @@ -41,3 +42,17 @@ pub struct TurnAbortInput<'a> { /// Store scoped to this turn runtime. pub turn_store: &'a ExtensionData, } + +/// Input supplied when a running turn encounters a host-classified error. +pub struct TurnErrorInput<'a> { + /// Stable host-owned turn identifier. + pub turn_id: &'a str, + /// Stable error classification reported by the host. + pub error: CodexErrorInfo, + /// Store scoped to the host session runtime. + pub session_store: &'a ExtensionData, + /// Store scoped to this thread runtime. + pub thread_store: &'a ExtensionData, + /// Store scoped to this turn runtime. + pub turn_store: &'a ExtensionData, +} diff --git a/codex-rs/ext/extension-api/src/lib.rs b/codex-rs/ext/extension-api/src/lib.rs index f3c3dc526e..f92aca46b6 100644 --- a/codex-rs/ext/extension-api/src/lib.rs +++ b/codex-rs/ext/extension-api/src/lib.rs @@ -25,6 +25,9 @@ pub use codex_tools::parse_tool_input_schema_without_compaction; pub use contributors::ApprovalReviewContributor; pub use contributors::ConfigContributor; pub use contributors::ContextContributor; +pub use contributors::IdleTurnContributor; +pub use contributors::IdleTurnFuture; +pub use contributors::IdleTurnInput; pub use contributors::PromptFragment; pub use contributors::PromptSlot; pub use contributors::ThreadLifecycleContributor; @@ -34,12 +37,14 @@ pub use contributors::ThreadStopInput; pub use contributors::TokenUsageContributor; pub use contributors::ToolCallOutcome; pub use contributors::ToolCallSource; +pub use contributors::ToolContributionInput; pub use contributors::ToolContributor; pub use contributors::ToolFinishInput; pub use contributors::ToolLifecycleContributor; pub use contributors::ToolLifecycleFuture; pub use contributors::ToolStartInput; pub use contributors::TurnAbortInput; +pub use contributors::TurnErrorInput; pub use contributors::TurnItemContributor; pub use contributors::TurnLifecycleContributor; pub use contributors::TurnStartInput; diff --git a/codex-rs/ext/extension-api/src/registry.rs b/codex-rs/ext/extension-api/src/registry.rs index a359726257..1dd2f9e620 100644 --- a/codex-rs/ext/extension-api/src/registry.rs +++ b/codex-rs/ext/extension-api/src/registry.rs @@ -7,6 +7,7 @@ use crate::ConfigContributor; use crate::ContextContributor; use crate::ExtensionData; use crate::ExtensionEventSink; +use crate::IdleTurnContributor; use crate::NoopExtensionEventSink; use crate::ThreadLifecycleContributor; use crate::TokenUsageContributor; @@ -26,6 +27,7 @@ pub struct ExtensionRegistryBuilder { tool_contributors: Vec>, tool_lifecycle_contributors: Vec>, turn_item_contributors: Vec>, + idle_turn_contributors: Vec>, approval_review_contributors: Vec>, } @@ -42,6 +44,7 @@ impl Default for ExtensionRegistryBuilder { tool_contributors: Vec::new(), tool_lifecycle_contributors: Vec::new(), turn_item_contributors: Vec::new(), + idle_turn_contributors: Vec::new(), } } } @@ -113,6 +116,11 @@ impl ExtensionRegistryBuilder { self.turn_item_contributors.push(contributor); } + /// Registers one idle-turn contributor. + pub fn idle_turn_contributor(&mut self, contributor: Arc) { + self.idle_turn_contributors.push(contributor); + } + /// Finishes construction and returns the immutable registry. pub fn build(self) -> ExtensionRegistry { ExtensionRegistry { @@ -126,6 +134,7 @@ impl ExtensionRegistryBuilder { tool_contributors: self.tool_contributors, tool_lifecycle_contributors: self.tool_lifecycle_contributors, turn_item_contributors: self.turn_item_contributors, + idle_turn_contributors: self.idle_turn_contributors, } } } @@ -141,6 +150,7 @@ pub struct ExtensionRegistry { tool_contributors: Vec>, tool_lifecycle_contributors: Vec>, turn_item_contributors: Vec>, + idle_turn_contributors: Vec>, approval_review_contributors: Vec>, } @@ -209,6 +219,11 @@ impl ExtensionRegistry { pub fn turn_item_contributors(&self) -> &[Arc] { &self.turn_item_contributors } + + /// Returns the registered idle-turn contributors. + pub fn idle_turn_contributors(&self) -> &[Arc] { + &self.idle_turn_contributors + } } /// Creates an empty shared registry for hosts that do not register contributions. From db1111720acb2909b0d94fde53075df6b803d411 Mon Sep 17 00:00:00 2001 From: Eric Traut Date: Wed, 27 May 2026 06:58:57 -0700 Subject: [PATCH 2/2] Simplify idle extension turn scheduling --- codex-rs/core/src/tasks/idle_extension.rs | 85 +++++++++++++- codex-rs/core/src/tasks/mod.rs | 105 ------------------ .../ext/extension-api/src/contributors.rs | 2 + 3 files changed, 86 insertions(+), 106 deletions(-) diff --git a/codex-rs/core/src/tasks/idle_extension.rs b/codex-rs/core/src/tasks/idle_extension.rs index 03a3953b34..71d15632b3 100644 --- a/codex-rs/core/src/tasks/idle_extension.rs +++ b/codex-rs/core/src/tasks/idle_extension.rs @@ -1,6 +1,12 @@ +//! Scheduling glue for extension-owned turns that should run when a session is idle. + use std::sync::Arc; +use crate::session::TurnInput; use crate::session::session::Session; +use crate::state::ActiveTurn; + +use super::RegularTask; pub(super) fn schedule_turn(session: &Arc) { if session @@ -14,6 +20,83 @@ pub(super) fn schedule_turn(session: &Arc) { let session = Arc::clone(session); let _handle = tokio::spawn(async move { - session.maybe_start_idle_extension_turn().await; + maybe_start_turn(session).await; }); } + +async fn maybe_start_turn(session: Arc) { + // Give queued user-visible work the first chance to wake the session. + session.maybe_start_turn_for_pending_work().await; + if has_active_or_pending_work(&session).await { + return; + } + + // Ask extensions for one idle turn only while the session still looks quiet. + let Some(input) = next_idle_turn_input(&session).await else { + return; + }; + + // Extension callbacks can race with user input or mailbox delivery. Re-check before + // claiming the idle slot, and let the normal pending-work path handle anything that appeared. + if has_active_or_pending_work(&session).await { + session.maybe_start_turn_for_pending_work().await; + return; + } + + // Reserve the active turn after the final quiet check so another task cannot start while + // the turn context is being built. + { + let mut active_turn = session.active_turn.lock().await; + if active_turn.is_some() { + return; + } + *active_turn = Some(ActiveTurn::default()); + } + + // Treat extension-provided items as the new turn's initial input rather than stashing them + // in turn-state pending input; this keeps rollback logic out of the scheduler. + let turn_context = session + .new_default_turn_with_sub_id(uuid::Uuid::new_v4().to_string()) + .await; + session + .maybe_emit_unknown_model_warning_for_turn(turn_context.as_ref()) + .await; + session + .start_task(turn_context, input, RegularTask::new()) + .await; +} + +async fn has_active_or_pending_work(session: &Session) -> bool { + session.active_turn.lock().await.is_some() + || session + .input_queue + .has_queued_response_items_for_next_turn() + .await + || session.input_queue.has_trigger_turn_mailbox_items().await +} + +async fn next_idle_turn_input(session: &Session) -> Option> { + let collaboration_mode = session.collaboration_mode().await; + for contributor in session.services.extensions.idle_turn_contributors() { + let Some(items) = contributor + .next_idle_turn(codex_extension_api::IdleTurnInput { + collaboration_mode: &collaboration_mode, + session_store: &session.services.session_extension_data, + thread_store: &session.services.thread_extension_data, + }) + .await + else { + continue; + }; + if !items.is_empty() { + return Some( + items + .into_iter() + .map(TurnInput::ResponseInputItem) + .collect(), + ); + } + } + + None +} diff --git a/codex-rs/core/src/tasks/mod.rs b/codex-rs/core/src/tasks/mod.rs index 3a5484fe75..bda60e0121 100644 --- a/codex-rs/core/src/tasks/mod.rs +++ b/codex-rs/core/src/tasks/mod.rs @@ -34,7 +34,6 @@ use crate::session::turn_context::TurnContext; use crate::state::ActiveTurn; use crate::state::RunningTask; use crate::state::TaskKind; -use crate::state::TurnState; use codex_analytics::TurnTokenUsageFact; use codex_login::AuthManager; use codex_models_manager::manager::SharedModelsManager; @@ -497,110 +496,6 @@ impl Session { .await; } - pub(crate) async fn maybe_start_idle_extension_turn(self: &Arc) { - if self.services.extensions.idle_turn_contributors().is_empty() { - return; - } - self.maybe_start_turn_for_pending_work().await; - if self.active_turn.lock().await.is_some() { - return; - } - if self - .input_queue - .has_queued_response_items_for_next_turn() - .await - || self.input_queue.has_trigger_turn_mailbox_items().await - { - return; - } - - let items = { - let collaboration_mode = self.collaboration_mode().await; - let mut requested_items = None; - for contributor in self.services.extensions.idle_turn_contributors() { - if let Some(items) = contributor - .next_idle_turn(codex_extension_api::IdleTurnInput { - collaboration_mode: &collaboration_mode, - session_store: &self.services.session_extension_data, - thread_store: &self.services.thread_extension_data, - }) - .await - { - requested_items = Some(items); - break; - } - } - requested_items - }; - let Some(items) = items.filter(|items| !items.is_empty()) else { - return; - }; - - let turn_state = { - let mut active_turn = self.active_turn.lock().await; - if active_turn.is_some() { - return; - } - let active_turn = active_turn.get_or_insert_with(ActiveTurn::default); - Arc::clone(&active_turn.turn_state) - }; - if self - .input_queue - .has_queued_response_items_for_next_turn() - .await - || self.input_queue.has_trigger_turn_mailbox_items().await - { - self.clear_reserved_idle_extension_turn(&turn_state).await; - self.maybe_start_turn_for_pending_work().await; - return; - } - - self.input_queue - .extend_pending_input_for_turn_state( - turn_state.as_ref(), - items - .into_iter() - .map(TurnInput::ResponseInputItem) - .collect(), - ) - .await; - - let turn_context = self - .new_default_turn_with_sub_id(uuid::Uuid::new_v4().to_string()) - .await; - self.maybe_emit_unknown_model_warning_for_turn(turn_context.as_ref()) - .await; - let still_reserved = { - let active_turn = self.active_turn.lock().await; - active_turn.as_ref().is_some_and(|active_turn| { - active_turn.task.is_none() && Arc::ptr_eq(&active_turn.turn_state, &turn_state) - }) - }; - if !still_reserved { - self.input_queue - .take_pending_input_for_turn_state(turn_state.as_ref()) - .await; - self.clear_reserved_idle_extension_turn(&turn_state).await; - return; - } - drop(turn_state); - self.start_task(turn_context, Vec::new(), RegularTask::new()) - .await; - } - - async fn clear_reserved_idle_extension_turn( - &self, - turn_state: &Arc>, - ) { - let mut active_turn_guard = self.active_turn.lock().await; - if let Some(active_turn) = active_turn_guard.as_ref() - && active_turn.task.is_none() - && Arc::ptr_eq(&active_turn.turn_state, turn_state) - { - *active_turn_guard = None; - } - } - pub async fn abort_all_tasks(self: &Arc, reason: TurnAbortReason) { let mut aborted_turn = false; let mut active_turn_to_clear = None; diff --git a/codex-rs/ext/extension-api/src/contributors.rs b/codex-rs/ext/extension-api/src/contributors.rs index ee15b82bc3..e01552c9c1 100644 --- a/codex-rs/ext/extension-api/src/contributors.rs +++ b/codex-rs/ext/extension-api/src/contributors.rs @@ -1,3 +1,5 @@ +//! Contributor traits and inputs that let extensions hook into host-owned workflows. + use std::future::Future; use std::pin::Pin; use std::sync::Arc;