From 601a2ce39243ed98c80a1529afeeef87f8425c1a Mon Sep 17 00:00:00 2001 From: Rohit Arunachalam Date: Wed, 17 Jun 2026 17:53:10 -0700 Subject: [PATCH 1/9] Inject current-time reminders before sampling --- codex-rs/app-server/src/mcp_refresh.rs | 1 + codex-rs/app-server/src/message_processor.rs | 1 + codex-rs/core/src/codex_delegate.rs | 1 + .../core/src/context/current_time_reminder.rs | 36 ++++ codex-rs/core/src/context/mod.rs | 2 + codex-rs/core/src/current_time.rs | 53 +++++ codex-rs/core/src/lib.rs | 4 + codex-rs/core/src/prompt_debug.rs | 1 + codex-rs/core/src/session/mod.rs | 5 + codex-rs/core/src/session/session.rs | 6 + codex-rs/core/src/session/tests.rs | 5 + .../core/src/session/tests/guardian_tests.rs | 1 + codex-rs/core/src/session/time_reminder.rs | 71 +++++++ codex-rs/core/src/session/turn.rs | 14 ++ codex-rs/core/src/state/service.rs | 2 + codex-rs/core/src/state/session.rs | 3 + codex-rs/core/src/thread_manager.rs | 6 + codex-rs/core/src/thread_manager_tests.rs | 11 + .../src/tools/handlers/multi_agents_tests.rs | 1 + codex-rs/core/tests/common/test_codex.rs | 9 + codex-rs/core/tests/suite/client.rs | 1 + codex-rs/core/tests/suite/mod.rs | 1 + codex-rs/core/tests/suite/varlatency.rs | 196 ++++++++++++++++++ codex-rs/mcp-server/src/message_processor.rs | 1 + codex-rs/thread-manager-sample/src/main.rs | 1 + 25 files changed, 433 insertions(+) create mode 100644 codex-rs/core/src/context/current_time_reminder.rs create mode 100644 codex-rs/core/src/current_time.rs create mode 100644 codex-rs/core/src/session/time_reminder.rs create mode 100644 codex-rs/core/tests/suite/varlatency.rs diff --git a/codex-rs/app-server/src/mcp_refresh.rs b/codex-rs/app-server/src/mcp_refresh.rs index a5659174f465..078664b46a22 100644 --- a/codex-rs/app-server/src/mcp_refresh.rs +++ b/codex-rs/app-server/src/mcp_refresh.rs @@ -250,6 +250,7 @@ mod tests { Some(state_db.clone()), "11111111-1111-4111-8111-111111111111".to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ) }); thread_manager.start_thread(good_config).await?; diff --git a/codex-rs/app-server/src/message_processor.rs b/codex-rs/app-server/src/message_processor.rs index 9e9437bd34fd..0d1d773cdeea 100644 --- a/codex-rs/app-server/src/message_processor.rs +++ b/codex-rs/app-server/src/message_processor.rs @@ -379,6 +379,7 @@ impl MessageProcessor { outgoing.clone(), thread_state_manager.clone(), )), + /*external_current_time_provider*/ None, ) }); let models_manager = thread_manager.get_models_manager(); diff --git a/codex-rs/core/src/codex_delegate.rs b/codex-rs/core/src/codex_delegate.rs index b7b1fe758066..5becea1d313a 100644 --- a/codex-rs/core/src/codex_delegate.rs +++ b/codex-rs/core/src/codex_delegate.rs @@ -120,6 +120,7 @@ pub(crate) async fn run_codex_thread_interactive( analytics_events_client: Some(parent_session.services.analytics_events_client.clone()), thread_store: Arc::clone(&parent_session.services.thread_store), attestation_provider: parent_session.services.attestation_provider.clone(), + external_current_time_provider: parent_session.services.current_time_provider.clone(), inherited_multi_agent_version: Some(MultiAgentVersion::Disabled), })) .or_cancel(&cancel_token) diff --git a/codex-rs/core/src/context/current_time_reminder.rs b/codex-rs/core/src/context/current_time_reminder.rs new file mode 100644 index 000000000000..d81376604e62 --- /dev/null +++ b/codex-rs/core/src/context/current_time_reminder.rs @@ -0,0 +1,36 @@ +use chrono::DateTime; +use chrono::Utc; + +use super::ContextualUserFragment; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct CurrentTimeReminder { + current_time: DateTime, +} + +impl CurrentTimeReminder { + pub(crate) fn new(current_time: DateTime) -> Self { + Self { current_time } + } +} + +impl ContextualUserFragment for CurrentTimeReminder { + fn role(&self) -> &'static str { + "developer" + } + + fn markers(&self) -> (&'static str, &'static str) { + Self::type_markers() + } + + fn type_markers() -> (&'static str, &'static str) { + ("", "") + } + + fn body(&self) -> String { + format!( + "It is {}.", + self.current_time.format("%Y-%m-%d %H:%M:%S UTC") + ) + } +} diff --git a/codex-rs/core/src/context/mod.rs b/codex-rs/core/src/context/mod.rs index 7fc1375781a4..3c2f48723146 100644 --- a/codex-rs/core/src/context/mod.rs +++ b/codex-rs/core/src/context/mod.rs @@ -6,6 +6,7 @@ mod available_plugins_instructions; mod available_skills_instructions; mod collaboration_mode_instructions; mod contextual_user_message; +mod current_time_reminder; mod environment_context; mod guardian_followup_review_reminder; mod hook_additional_context; @@ -44,6 +45,7 @@ pub(crate) use codex_core_skills::SkillInstructions; 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 current_time_reminder::CurrentTimeReminder; pub(crate) use environment_context::EnvironmentContext; pub(crate) use guardian_followup_review_reminder::GuardianFollowupReviewReminder; pub(crate) use hook_additional_context::HookAdditionalContext; diff --git a/codex-rs/core/src/current_time.rs b/codex-rs/core/src/current_time.rs new file mode 100644 index 000000000000..4622a6e3a4f7 --- /dev/null +++ b/codex-rs/core/src/current_time.rs @@ -0,0 +1,53 @@ +use std::future::Future; +use std::pin::Pin; +use std::sync::Arc; + +use anyhow::Result; +use anyhow::anyhow; +use chrono::DateTime; +use chrono::Utc; +use codex_features::VarlatencyClockSource; +use codex_protocol::ThreadId; + +use crate::config::VarlatencyConfig; + +pub type CurrentTimeFuture<'a> = Pin>> + Send + 'a>>; + +/// Context supplied when Codex asks a host integration for the current time. +#[derive(Clone, Copy, Debug)] +pub struct CurrentTimeContext { + /// Thread preparing the model request. + pub thread_id: ThreadId, +} + +/// Host integration boundary for obtaining the current time. +pub trait CurrentTimeProvider: std::fmt::Debug + Send + Sync { + fn current_time(&self, context: CurrentTimeContext) -> CurrentTimeFuture<'_>; +} + +#[derive(Debug, Default)] +struct SystemCurrentTimeProvider; + +impl CurrentTimeProvider for SystemCurrentTimeProvider { + fn current_time(&self, _context: CurrentTimeContext) -> CurrentTimeFuture<'_> { + Box::pin(async { Ok(Utc::now()) }) + } +} + +pub(crate) fn resolve_current_time_provider( + config: Option<&VarlatencyConfig>, + external_provider: Option>, +) -> Result>> { + let Some(config) = config else { + return Ok(None); + }; + + match config.clock_source { + VarlatencyClockSource::System => Ok(Some(Arc::new(SystemCurrentTimeProvider))), + VarlatencyClockSource::AppServerClient => external_provider.map(Some).ok_or_else(|| { + anyhow!( + "features.varlatency.clock_source is app_server_client, but no external current-time provider is available" + ) + }), + } +} diff --git a/codex-rs/core/src/lib.rs b/codex-rs/core/src/lib.rs index 0e193baa1d37..2c4f6e311160 100644 --- a/codex-rs/core/src/lib.rs +++ b/codex-rs/core/src/lib.rs @@ -37,6 +37,7 @@ pub mod config; pub mod connectors; pub mod context; mod context_manager; +mod current_time; mod environment_selection; pub mod exec; pub mod exec_env; @@ -186,6 +187,9 @@ pub use client_common::ResponseEvent; pub use client_common::ResponseStream; pub use codex_prompts::REVIEW_PROMPT; pub use compact::content_items_to_text; +pub use current_time::CurrentTimeContext; +pub use current_time::CurrentTimeFuture; +pub use current_time::CurrentTimeProvider; pub use event_mapping::parse_turn_item; pub use exec_policy::ExecPolicyError; pub use exec_policy::check_execpolicy_for_warnings; diff --git a/codex-rs/core/src/prompt_debug.rs b/codex-rs/core/src/prompt_debug.rs index 0da42d8a4e33..31f63ccceaa6 100644 --- a/codex-rs/core/src/prompt_debug.rs +++ b/codex-rs/core/src/prompt_debug.rs @@ -60,6 +60,7 @@ pub async fn build_prompt_input( state_db.clone(), installation_id, /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let thread = thread_manager.start_thread(config).await?; diff --git a/codex-rs/core/src/session/mod.rs b/codex-rs/core/src/session/mod.rs index f92bd6368080..bb1e7bdde36a 100644 --- a/codex-rs/core/src/session/mod.rs +++ b/codex-rs/core/src/session/mod.rs @@ -31,6 +31,7 @@ use crate::context::NetworkRuleSaved; use crate::context::PermissionsInstructions; use crate::context::PersonalitySpecInstructions; use crate::context::RecommendedPluginsInstructions; +use crate::current_time::CurrentTimeProvider; use crate::default_skill_metadata_budget; use crate::environment_selection::TurnEnvironmentSnapshot; use crate::exec_policy::ExecPolicyManager; @@ -215,6 +216,7 @@ mod rollout_budget; mod rollout_reconstruction; #[allow(clippy::module_inception)] pub(crate) mod session; +pub(crate) mod time_reminder; mod token_budget; pub(crate) mod turn; pub(crate) mod turn_context; @@ -439,6 +441,7 @@ pub(crate) struct CodexSpawnArgs { pub(crate) analytics_events_client: Option, pub(crate) thread_store: Arc, pub(crate) attestation_provider: Option>, + pub(crate) external_current_time_provider: Option>, pub(crate) inherited_multi_agent_version: Option, } @@ -522,6 +525,7 @@ impl Codex { analytics_events_client, thread_store, attestation_provider, + external_current_time_provider, inherited_multi_agent_version, } = args; let (tx_sub, rx_sub) = async_channel::bounded(SUBMISSION_CHANNEL_CAPACITY); @@ -669,6 +673,7 @@ impl Codex { thread_store, parent_rollout_thread_trace, attestation_provider, + external_current_time_provider, multi_agent_version, )) .await diff --git a/codex-rs/core/src/session/session.rs b/codex-rs/core/src/session/session.rs index 2b61468c15d3..3ef4f7fdd0f5 100644 --- a/codex-rs/core/src/session/session.rs +++ b/codex-rs/core/src/session/session.rs @@ -492,6 +492,7 @@ impl Session { thread_store: Arc, parent_rollout_thread_trace: ThreadTraceContext, attestation_provider: Option>, + external_current_time_provider: Option>, multi_agent_version: Option, ) -> anyhow::Result> { debug!( @@ -516,6 +517,10 @@ impl Session { } InitialHistory::Resumed(resumed_history) => resumed_history.conversation_id, }; + let current_time_provider = crate::current_time::resolve_current_time_provider( + config.varlatency.as_ref(), + external_current_time_provider, + )?; let mcp_thread_init = thread_extension_init.clone(); let thread_extension_data = codex_extension_api::ExtensionData::new_with_init( thread_id.to_string(), @@ -1022,6 +1027,7 @@ impl Session { live_thread: live_thread_init.as_ref().cloned(), thread_store: Arc::clone(&thread_store), attestation_provider: attestation_provider.clone(), + current_time_provider, model_client: ModelClient::new( Some(Arc::clone(&auth_manager)), thread_id, diff --git a/codex-rs/core/src/session/tests.rs b/codex-rs/core/src/session/tests.rs index b13cd029ce62..56a1e368fdf5 100644 --- a/codex-rs/core/src/session/tests.rs +++ b/codex-rs/core/src/session/tests.rs @@ -4867,6 +4867,7 @@ async fn session_new_fails_when_zsh_fork_enabled_without_packaged_zsh() { )), codex_rollout_trace::ThreadTraceContext::disabled(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, Some(config.multi_agent_version_from_features()), ) .await; @@ -5032,6 +5033,7 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) { /*state_db*/ None, )), attestation_provider: None, + current_time_provider: None, model_client: ModelClient::new( Some(auth_manager.clone()), thread_id, @@ -5216,6 +5218,7 @@ async fn make_session_with_config_and_rx( )), codex_rollout_trace::ThreadTraceContext::disabled(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, Some(config.multi_agent_version_from_features()), ) .await?; @@ -5328,6 +5331,7 @@ async fn make_session_with_history_source_and_agent_control_and_rx( )), codex_rollout_trace::ThreadTraceContext::disabled(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, Some(config.multi_agent_version_from_features()), ) .await?; @@ -7078,6 +7082,7 @@ where state_db, )), attestation_provider: None, + current_time_provider: None, model_client: ModelClient::new( Some(Arc::clone(&auth_manager)), thread_id, diff --git a/codex-rs/core/src/session/tests/guardian_tests.rs b/codex-rs/core/src/session/tests/guardian_tests.rs index 0bed03235169..e3451df56863 100644 --- a/codex-rs/core/src/session/tests/guardian_tests.rs +++ b/codex-rs/core/src/session/tests/guardian_tests.rs @@ -736,6 +736,7 @@ async fn guardian_subagent_does_not_inherit_parent_exec_policy_rules() { analytics_events_client: None, thread_store, attestation_provider: None, + external_current_time_provider: None, inherited_multi_agent_version: None, }) .await diff --git a/codex-rs/core/src/session/time_reminder.rs b/codex-rs/core/src/session/time_reminder.rs new file mode 100644 index 000000000000..2a6363d2d310 --- /dev/null +++ b/codex-rs/core/src/session/time_reminder.rs @@ -0,0 +1,71 @@ +use std::sync::Arc; + +use codex_protocol::error::CodexErr; +use codex_protocol::error::Result as CodexResult; +use codex_protocol::models::ResponseItem; + +use super::session::Session; +use super::turn_context::TurnContext; +use crate::context::ContextualUserFragment; +use crate::current_time::CurrentTimeContext; + +#[derive(Default)] +pub(crate) struct CurrentTimeReminderState { + model_requests_since_delivery: u64, + last_window_id: Option, +} + +impl CurrentTimeReminderState { + fn begin_model_request(&mut self, window_id: &str, interval: u64) -> bool { + self.model_requests_since_delivery = self.model_requests_since_delivery.saturating_add(1); + self.last_window_id.as_deref() != Some(window_id) + || self.model_requests_since_delivery >= interval + } + + fn record_delivery(&mut self, window_id: &str) { + self.model_requests_since_delivery = 0; + self.last_window_id = Some(window_id.to_string()); + } +} + +pub(super) async fn maybe_record_current_time_reminder( + sess: &Session, + turn_context: &TurnContext, + window_id: &str, +) -> CodexResult<()> { + let Some(config) = turn_context.config.varlatency else { + return Ok(()); + }; + + let reminder_is_due = { + let mut state = sess.state.lock().await; + state + .current_time_reminder + .begin_model_request(window_id, config.reminder_interval_model_requests) + }; + if !reminder_is_due { + return Ok(()); + } + + let provider = sess + .services + .current_time_provider + .as_ref() + .map(Arc::clone) + .ok_or_else(|| CodexErr::Fatal("current-time provider is not configured".to_string()))?; + let current_time = provider + .current_time(CurrentTimeContext { + thread_id: sess.thread_id, + }) + .await + .map_err(|err| CodexErr::Fatal(format!("failed to read current time: {err:#}")))?; + + let response_item: ResponseItem = + ContextualUserFragment::into(crate::context::CurrentTimeReminder::new(current_time)); + sess.record_conversation_items(turn_context, std::slice::from_ref(&response_item)) + .await; + + let mut state = sess.state.lock().await; + state.current_time_reminder.record_delivery(window_id); + Ok(()) +} diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index 657e00ff3575..1e2e442a13be 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -226,6 +226,20 @@ pub(crate) async fn run_turn( ) .await; + if let Err(err) = super::time_reminder::maybe_record_current_time_reminder( + sess.as_ref(), + turn_context.as_ref(), + &window_id, + ) + .await + { + let error = err.to_codex_protocol_error(); + sess.emit_turn_error_lifecycle(turn_context.as_ref(), error.clone()) + .await; + error!("Failed to record current-time reminder"); + return None; + } + // Construct the input that we will send to the model. let sampling_request_input: Vec = async { sess.clone_history() diff --git a/codex-rs/core/src/state/service.rs b/codex-rs/core/src/state/service.rs index 914fd4c32d61..059362e6e039 100644 --- a/codex-rs/core/src/state/service.rs +++ b/codex-rs/core/src/state/service.rs @@ -8,6 +8,7 @@ use crate::attestation::AttestationProvider; use crate::client::ModelClient; use crate::config::NetworkProxyAuditMetadata; use crate::config::StartedNetworkProxy; +use crate::current_time::CurrentTimeProvider; use crate::environment_selection::ThreadEnvironments; use crate::exec_policy::ExecPolicyManager; use crate::guardian::GuardianRejection; @@ -79,6 +80,7 @@ pub(crate) struct SessionServices { pub(crate) live_thread: Option, pub(crate) thread_store: Arc, pub(crate) attestation_provider: Option>, + pub(crate) current_time_provider: Option>, /// Session-scoped model client shared across turns. pub(crate) model_client: ModelClient, pub(crate) code_mode_service: CodeModeService, diff --git a/codex-rs/core/src/state/session.rs b/codex-rs/core/src/state/session.rs index 269d3e0f607e..2bd9da5ce3de 100644 --- a/codex-rs/core/src/state/session.rs +++ b/codex-rs/core/src/state/session.rs @@ -13,6 +13,7 @@ use super::auto_compact_window::AutoCompactWindowSnapshot; use crate::context_manager::ContextManager; use crate::session::PreviousTurnSettings; use crate::session::session::SessionConfiguration; +use crate::session::time_reminder::CurrentTimeReminderState; use crate::session_startup_prewarm::SessionStartupPrewarmHandle; use codex_protocol::protocol::RateLimitSnapshot; use codex_protocol::protocol::TokenUsage; @@ -36,6 +37,7 @@ pub(crate) struct SessionState { auto_compact_window: AutoCompactWindow, /// Startup prewarmed session prepared during session initialization. pub(crate) startup_prewarm: Option, + pub(crate) current_time_reminder: CurrentTimeReminderState, pub(crate) active_connector_selection: HashSet, pub(crate) pending_session_start_sources: VecDeque, granted_permissions_by_environment_id: HashMap, @@ -56,6 +58,7 @@ impl SessionState { previous_turn_settings: None, auto_compact_window: AutoCompactWindow::new(), startup_prewarm: None, + current_time_reminder: CurrentTimeReminderState::default(), active_connector_selection: HashSet::new(), pending_session_start_sources: VecDeque::new(), granted_permissions_by_environment_id: HashMap::new(), diff --git a/codex-rs/core/src/thread_manager.rs b/codex-rs/core/src/thread_manager.rs index 77a5520026d9..3d84e5251046 100644 --- a/codex-rs/core/src/thread_manager.rs +++ b/codex-rs/core/src/thread_manager.rs @@ -4,6 +4,7 @@ use crate::attestation::AttestationProvider; use crate::codex_thread::CodexThread; use crate::config::Config; use crate::config::ThreadStoreConfig; +use crate::current_time::CurrentTimeProvider; use crate::environment_selection::TurnEnvironmentSnapshot; use crate::environment_selection::default_thread_environment_selections; use crate::mcp::McpManager; @@ -215,6 +216,7 @@ pub(crate) struct ThreadManagerState { user_instructions_provider: Arc, thread_store: Arc, attestation_provider: Option>, + external_current_time_provider: Option>, session_source: SessionSource, installation_id: String, analytics_events_client: Option, @@ -269,6 +271,7 @@ impl ThreadManager { state_db: Option, installation_id: String, attestation_provider: Option>, + external_current_time_provider: Option>, ) -> Self { let codex_home = config.codex_home.clone(); let restriction_product = session_source.restriction_product(); @@ -300,6 +303,7 @@ impl ThreadManager { user_instructions_provider, thread_store, attestation_provider, + external_current_time_provider, auth_manager, session_source, installation_id, @@ -405,6 +409,7 @@ impl ThreadManager { ), thread_store, attestation_provider: None, + external_current_time_provider: None, auth_manager, session_source: SessionSource::Exec, installation_id, @@ -1466,6 +1471,7 @@ impl ThreadManagerState { analytics_events_client: self.analytics_events_client.clone(), thread_store: Arc::clone(&self.thread_store), attestation_provider: self.attestation_provider.clone(), + external_current_time_provider: self.external_current_time_provider.clone(), inherited_multi_agent_version: multi_agent_version, })) .await?; diff --git a/codex-rs/core/src/thread_manager_tests.rs b/codex-rs/core/src/thread_manager_tests.rs index ee010b5e1d82..d843de96bbb2 100644 --- a/codex-rs/core/src/thread_manager_tests.rs +++ b/codex-rs/core/src/thread_manager_tests.rs @@ -440,6 +440,7 @@ async fn start_thread_seeds_extension_data_for_mcp_and_lifecycle_contributors() /*state_db*/ None, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let selected_root_init = |id: &str, environment_id: &str| { let mut init = codex_extension_api::ExtensionDataInit::new(); @@ -548,6 +549,7 @@ async fn resume_and_fork_do_not_restore_thread_environments_from_rollout() { /*state_db*/ None, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let selected_cwd = AbsolutePathBuf::try_from(config.cwd.as_path().join("selected")).expect("absolute path"); @@ -671,6 +673,7 @@ async fn explicit_installation_id_skips_codex_home_file() { state_db.clone(), installation_id.clone(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let thread = manager @@ -711,6 +714,7 @@ async fn resume_active_thread_from_rollout_returns_running_thread() { /*state_db*/ None, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let source = manager @@ -770,6 +774,7 @@ async fn resume_stopped_thread_from_rollout_spawns_new_thread() { /*state_db*/ None, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let source = manager @@ -836,6 +841,7 @@ async fn resume_stopped_thread_from_rollout_preserves_thread_source() { state_db.clone(), TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let source = manager @@ -929,6 +935,7 @@ async fn rollout_path_resume_and_fork_read_history_through_thread_store() { state_db, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let source = manager @@ -1033,6 +1040,7 @@ async fn new_uses_active_provider_for_model_refresh() { /*state_db*/ None, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let _ = manager.list_models(RefreshStrategy::Online).await; @@ -1254,6 +1262,7 @@ async fn interrupted_fork_snapshot_does_not_synthesize_turn_id_for_legacy_histor state_db.clone(), TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let source = manager @@ -1362,6 +1371,7 @@ async fn interrupted_fork_snapshot_preserves_explicit_turn_id() { state_db.clone(), TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let source = manager @@ -1460,6 +1470,7 @@ async fn interrupted_fork_snapshot_uses_persisted_mid_turn_history_without_live_ state_db.clone(), TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let source = manager diff --git a/codex-rs/core/src/tools/handlers/multi_agents_tests.rs b/codex-rs/core/src/tools/handlers/multi_agents_tests.rs index ed15a7483a8f..eabbcbbfa3c8 100644 --- a/codex-rs/core/src/tools/handlers/multi_agents_tests.rs +++ b/codex-rs/core/src/tools/handlers/multi_agents_tests.rs @@ -4279,6 +4279,7 @@ async fn tool_handlers_cascade_close_and_resume_and_keep_explicitly_closed_subtr state_db.clone(), "11111111-1111-4111-8111-111111111111".to_string(), /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let parent = manager diff --git a/codex-rs/core/tests/common/test_codex.rs b/codex-rs/core/tests/common/test_codex.rs index b10e26bf0c76..ba3c66c1eb2e 100644 --- a/codex-rs/core/tests/common/test_codex.rs +++ b/codex-rs/core/tests/common/test_codex.rs @@ -15,6 +15,7 @@ use anyhow::Result; use anyhow::anyhow; use codex_config::CloudConfigBundleLoader; use codex_core::CodexThread; +use codex_core::CurrentTimeProvider; use codex_core::StartThreadOptions; use codex_core::ThreadManager; use codex_core::config::Config; @@ -261,6 +262,7 @@ pub struct TestCodexBuilder { extensions: Arc>, user_instructions_provider: Option>, supports_openai_form_elicitation: bool, + current_time_provider: Option>, } impl TestCodexBuilder { @@ -362,6 +364,11 @@ impl TestCodexBuilder { self } + pub fn with_current_time_provider(mut self, provider: Arc) -> Self { + self.current_time_provider = Some(provider); + self + } + pub fn with_windows_cmd_shell(self) -> Self { if cfg!(windows) { self.with_user_shell(get_shell_by_model_provided_path(&PathBuf::from("cmd.exe"))) @@ -568,6 +575,7 @@ impl TestCodexBuilder { state_db.clone(), installation_id, /*attestation_provider*/ None, + /*external_current_time_provider*/ self.current_time_provider.clone(), ); let thread_manager = Arc::new(thread_manager); let user_shell_override = self.user_shell_override.clone(); @@ -1172,6 +1180,7 @@ pub fn test_codex() -> TestCodexBuilder { extensions: empty_extension_registry(), user_instructions_provider: None, supports_openai_form_elicitation: false, + current_time_provider: None, } } diff --git a/codex-rs/core/tests/suite/client.rs b/codex-rs/core/tests/suite/client.rs index 54e24d49cbe8..99da34e0ff9d 100644 --- a/codex-rs/core/tests/suite/client.rs +++ b/codex-rs/core/tests/suite/client.rs @@ -1259,6 +1259,7 @@ async fn prefers_apikey_when_config_prefers_apikey_even_with_chatgpt_tokens() { /*state_db*/ None, installation_id, /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let NewThread { thread: codex, .. } = thread_manager .start_thread(config.clone()) diff --git a/codex-rs/core/tests/suite/mod.rs b/codex-rs/core/tests/suite/mod.rs index 0fdd958b3983..7cba602fbf5b 100644 --- a/codex-rs/core/tests/suite/mod.rs +++ b/codex-rs/core/tests/suite/mod.rs @@ -126,6 +126,7 @@ mod unified_exec_zsh_fork_approvals; mod unstable_features_warning; mod user_notification; mod user_shell_cmd; +mod varlatency; mod view_image; mod web_search; mod websocket_fallback; diff --git a/codex-rs/core/tests/suite/varlatency.rs b/codex-rs/core/tests/suite/varlatency.rs new file mode 100644 index 000000000000..d7a92caab321 --- /dev/null +++ b/codex-rs/core/tests/suite/varlatency.rs @@ -0,0 +1,196 @@ +use std::collections::VecDeque; +use std::sync::Arc; +use std::sync::Mutex; + +use anyhow::Result; +use anyhow::anyhow; +use chrono::DateTime; +use chrono::Utc; +use codex_core::CurrentTimeContext; +use codex_core::CurrentTimeFuture; +use codex_core::CurrentTimeProvider; +use codex_core::config::VarlatencyConfig; +use codex_features::Feature; +use codex_features::VarlatencyClockSource; +use codex_model_provider_info::built_in_model_providers; +use codex_protocol::ThreadId; +use codex_protocol::protocol::EventMsg; +use codex_protocol::protocol::Op; +use core_test_support::responses::ResponsesRequest; +use core_test_support::responses::ev_assistant_message; +use core_test_support::responses::ev_completed; +use core_test_support::responses::ev_response_created; +use core_test_support::responses::mount_sse_sequence; +use core_test_support::responses::sse; +use core_test_support::responses::start_mock_server; +use core_test_support::skip_if_no_network; +use core_test_support::test_codex::test_codex; +use core_test_support::wait_for_event; +use pretty_assertions::assert_eq; + +const FIRST_REMINDER: &str = "It is 2026-06-17 17:34:15 UTC."; +const SECOND_REMINDER: &str = "It is 2026-06-17 17:35:15 UTC."; + +#[derive(Debug)] +struct TestCurrentTimeProvider { + times: Mutex>>, + requested_thread_ids: Mutex>, +} + +impl TestCurrentTimeProvider { + fn new(times: &[&str]) -> Self { + Self { + times: Mutex::new( + times + .iter() + .map(|time| { + DateTime::parse_from_rfc3339(time) + .expect("test time should be valid RFC 3339") + .with_timezone(&Utc) + }) + .collect(), + ), + requested_thread_ids: Mutex::new(Vec::new()), + } + } + + fn requested_thread_ids(&self) -> Vec { + self.requested_thread_ids + .lock() + .expect("requested thread IDs lock should not be poisoned") + .clone() + } +} + +impl CurrentTimeProvider for TestCurrentTimeProvider { + fn current_time(&self, context: CurrentTimeContext) -> CurrentTimeFuture<'_> { + self.requested_thread_ids + .lock() + .expect("requested thread IDs lock should not be poisoned") + .push(context.thread_id); + let time = self + .times + .lock() + .expect("test times lock should not be poisoned") + .pop_front() + .ok_or_else(|| anyhow!("test current-time provider is exhausted")); + Box::pin(async move { time }) + } +} + +fn current_time_reminders(request: &ResponsesRequest) -> Vec { + request + .message_input_texts("developer") + .into_iter() + .filter(|text| text.starts_with("It is ") && text.ends_with(" UTC.")) + .collect() +} + +fn enable_varlatency(config: &mut codex_core::config::Config, interval: u64) { + config + .features + .enable(Feature::Varlatency) + .expect("test config should allow varlatency"); + config.varlatency = Some(VarlatencyConfig { + reminder_interval_model_requests: interval, + clock_source: VarlatencyClockSource::AppServerClient, + }); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn varlatency_reminders_follow_request_interval_and_persist_in_history() -> Result<()> { + skip_if_no_network!(Ok(())); + + let server = start_mock_server().await; + let responses = mount_sse_sequence( + &server, + vec![ + sse(vec![ev_response_created("resp-1"), ev_completed("resp-1")]), + sse(vec![ev_response_created("resp-2"), ev_completed("resp-2")]), + sse(vec![ev_response_created("resp-3"), ev_completed("resp-3")]), + ], + ) + .await; + let provider = Arc::new(TestCurrentTimeProvider::new(&[ + "2026-06-17T17:34:15Z", + "2026-06-17T17:35:15Z", + ])); + let test = test_codex() + .with_config(|config| enable_varlatency(config, /*interval*/ 2)) + .with_current_time_provider(provider.clone()) + .build(&server) + .await?; + + test.submit_turn("first turn").await?; + test.submit_turn("second turn").await?; + test.submit_turn("third turn").await?; + + let requests = responses.requests(); + assert_eq!(requests.len(), 3); + assert_eq!(current_time_reminders(&requests[0]), vec![FIRST_REMINDER]); + assert_eq!(current_time_reminders(&requests[1]), vec![FIRST_REMINDER]); + assert_eq!( + current_time_reminders(&requests[2]), + vec![FIRST_REMINDER, SECOND_REMINDER] + ); + assert_eq!( + provider.requested_thread_ids(), + vec![test.session_configured.thread_id; 2] + ); + + Ok(()) +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn varlatency_reminder_is_refreshed_after_compaction() -> Result<()> { + skip_if_no_network!(Ok(())); + + let server = start_mock_server().await; + let responses = mount_sse_sequence( + &server, + vec![ + sse(vec![ev_response_created("resp-1"), ev_completed("resp-1")]), + sse(vec![ + ev_response_created("resp-compact"), + ev_assistant_message("msg-compact", "compact summary"), + ev_completed("resp-compact"), + ]), + sse(vec![ev_response_created("resp-2"), ev_completed("resp-2")]), + ], + ) + .await; + let provider = Arc::new(TestCurrentTimeProvider::new(&[ + "2026-06-17T17:34:15Z", + "2026-06-17T17:35:15Z", + ])); + let mut model_provider = built_in_model_providers(/*openai_base_url*/ None)["openai"].clone(); + model_provider.name = "OpenAI-compatible test provider".to_string(); + model_provider.base_url = Some(format!("{}/v1", server.uri())); + model_provider.supports_websockets = false; + let test = test_codex() + .with_config(move |config| { + config.model_provider = model_provider; + enable_varlatency(config, /*interval*/ 50); + }) + .with_current_time_provider(provider) + .build(&server) + .await?; + + test.submit_turn("before compact").await?; + test.codex.submit(Op::Compact).await?; + wait_for_event(&test.codex, |event| { + matches!(event, EventMsg::TurnComplete(_)) + }) + .await; + test.submit_turn("after compact").await?; + + let requests = responses.requests(); + assert_eq!(requests.len(), 3); + assert_eq!( + current_time_reminders(&requests[2]), + vec![SECOND_REMINDER], + "a new context window should force a fresh reminder before the next model request" + ); + + Ok(()) +} diff --git a/codex-rs/mcp-server/src/message_processor.rs b/codex-rs/mcp-server/src/message_processor.rs index ce19575a9c60..088b455cac3d 100644 --- a/codex-rs/mcp-server/src/message_processor.rs +++ b/codex-rs/mcp-server/src/message_processor.rs @@ -78,6 +78,7 @@ impl MessageProcessor { state_db.clone(), installation_id, /*attestation_provider*/ None, + /*external_current_time_provider*/ None, )); Self { outgoing, diff --git a/codex-rs/thread-manager-sample/src/main.rs b/codex-rs/thread-manager-sample/src/main.rs index 1534ba86a084..778c37b0a63a 100644 --- a/codex-rs/thread-manager-sample/src/main.rs +++ b/codex-rs/thread-manager-sample/src/main.rs @@ -136,6 +136,7 @@ async fn run_main(arg0_paths: Arg0DispatchPaths) -> anyhow::Result<()> { state_db, installation_id, /*attestation_provider*/ None, + /*external_current_time_provider*/ None, ); let NewThread { From 28a79c41c791aa8a2f8457d76077f744f701b0c6 Mon Sep 17 00:00:00 2001 From: Rohit Arunachalam Date: Wed, 17 Jun 2026 18:09:11 -0700 Subject: [PATCH 2/9] Simplify current-time reminders --- .../core/src/context/current_time_reminder.rs | 1 - codex-rs/core/src/current_time.rs | 14 +--- codex-rs/core/src/lib.rs | 1 - codex-rs/core/src/session/time_reminder.rs | 11 +-- codex-rs/core/src/session/turn.rs | 3 +- codex-rs/core/tests/suite/varlatency.rs | 75 ++++--------------- 6 files changed, 22 insertions(+), 83 deletions(-) diff --git a/codex-rs/core/src/context/current_time_reminder.rs b/codex-rs/core/src/context/current_time_reminder.rs index d81376604e62..0e3fe0b1e8a9 100644 --- a/codex-rs/core/src/context/current_time_reminder.rs +++ b/codex-rs/core/src/context/current_time_reminder.rs @@ -3,7 +3,6 @@ use chrono::Utc; use super::ContextualUserFragment; -#[derive(Debug, Clone, PartialEq, Eq)] pub(crate) struct CurrentTimeReminder { current_time: DateTime, } diff --git a/codex-rs/core/src/current_time.rs b/codex-rs/core/src/current_time.rs index 4622a6e3a4f7..36cc4b40e80f 100644 --- a/codex-rs/core/src/current_time.rs +++ b/codex-rs/core/src/current_time.rs @@ -13,23 +13,15 @@ use crate::config::VarlatencyConfig; pub type CurrentTimeFuture<'a> = Pin>> + Send + 'a>>; -/// Context supplied when Codex asks a host integration for the current time. -#[derive(Clone, Copy, Debug)] -pub struct CurrentTimeContext { - /// Thread preparing the model request. - pub thread_id: ThreadId, -} - /// Host integration boundary for obtaining the current time. -pub trait CurrentTimeProvider: std::fmt::Debug + Send + Sync { - fn current_time(&self, context: CurrentTimeContext) -> CurrentTimeFuture<'_>; +pub trait CurrentTimeProvider: Send + Sync { + fn current_time(&self, thread_id: ThreadId) -> CurrentTimeFuture<'_>; } -#[derive(Debug, Default)] struct SystemCurrentTimeProvider; impl CurrentTimeProvider for SystemCurrentTimeProvider { - fn current_time(&self, _context: CurrentTimeContext) -> CurrentTimeFuture<'_> { + fn current_time(&self, _thread_id: ThreadId) -> CurrentTimeFuture<'_> { Box::pin(async { Ok(Utc::now()) }) } } diff --git a/codex-rs/core/src/lib.rs b/codex-rs/core/src/lib.rs index 2c4f6e311160..94458e5431dc 100644 --- a/codex-rs/core/src/lib.rs +++ b/codex-rs/core/src/lib.rs @@ -187,7 +187,6 @@ pub use client_common::ResponseEvent; pub use client_common::ResponseStream; pub use codex_prompts::REVIEW_PROMPT; pub use compact::content_items_to_text; -pub use current_time::CurrentTimeContext; pub use current_time::CurrentTimeFuture; pub use current_time::CurrentTimeProvider; pub use event_mapping::parse_turn_item; diff --git a/codex-rs/core/src/session/time_reminder.rs b/codex-rs/core/src/session/time_reminder.rs index 2a6363d2d310..b17bc6a10c35 100644 --- a/codex-rs/core/src/session/time_reminder.rs +++ b/codex-rs/core/src/session/time_reminder.rs @@ -1,13 +1,9 @@ -use std::sync::Arc; - use codex_protocol::error::CodexErr; use codex_protocol::error::Result as CodexResult; -use codex_protocol::models::ResponseItem; use super::session::Session; use super::turn_context::TurnContext; use crate::context::ContextualUserFragment; -use crate::current_time::CurrentTimeContext; #[derive(Default)] pub(crate) struct CurrentTimeReminderState { @@ -51,16 +47,13 @@ pub(super) async fn maybe_record_current_time_reminder( .services .current_time_provider .as_ref() - .map(Arc::clone) .ok_or_else(|| CodexErr::Fatal("current-time provider is not configured".to_string()))?; let current_time = provider - .current_time(CurrentTimeContext { - thread_id: sess.thread_id, - }) + .current_time(sess.thread_id) .await .map_err(|err| CodexErr::Fatal(format!("failed to read current time: {err:#}")))?; - let response_item: ResponseItem = + let response_item = ContextualUserFragment::into(crate::context::CurrentTimeReminder::new(current_time)); sess.record_conversation_items(turn_context, std::slice::from_ref(&response_item)) .await; diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index 1e2e442a13be..57f3db96ed5b 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -233,8 +233,7 @@ pub(crate) async fn run_turn( ) .await { - let error = err.to_codex_protocol_error(); - sess.emit_turn_error_lifecycle(turn_context.as_ref(), error.clone()) + sess.emit_turn_error_lifecycle(turn_context.as_ref(), err.to_codex_protocol_error()) .await; error!("Failed to record current-time reminder"); return None; diff --git a/codex-rs/core/tests/suite/varlatency.rs b/codex-rs/core/tests/suite/varlatency.rs index d7a92caab321..197141327484 100644 --- a/codex-rs/core/tests/suite/varlatency.rs +++ b/codex-rs/core/tests/suite/varlatency.rs @@ -1,12 +1,10 @@ -use std::collections::VecDeque; use std::sync::Arc; -use std::sync::Mutex; +use std::sync::atomic::AtomicI64; +use std::sync::atomic::Ordering; use anyhow::Result; -use anyhow::anyhow; use chrono::DateTime; use chrono::Utc; -use codex_core::CurrentTimeContext; use codex_core::CurrentTimeFuture; use codex_core::CurrentTimeProvider; use codex_core::config::VarlatencyConfig; @@ -30,51 +28,23 @@ use pretty_assertions::assert_eq; const FIRST_REMINDER: &str = "It is 2026-06-17 17:34:15 UTC."; const SECOND_REMINDER: &str = "It is 2026-06-17 17:35:15 UTC."; +const FIRST_TIME_UNIX_SECONDS: i64 = 1_781_717_655; -#[derive(Debug)] -struct TestCurrentTimeProvider { - times: Mutex>>, - requested_thread_ids: Mutex>, -} - -impl TestCurrentTimeProvider { - fn new(times: &[&str]) -> Self { - Self { - times: Mutex::new( - times - .iter() - .map(|time| { - DateTime::parse_from_rfc3339(time) - .expect("test time should be valid RFC 3339") - .with_timezone(&Utc) - }) - .collect(), - ), - requested_thread_ids: Mutex::new(Vec::new()), - } - } +struct TestCurrentTimeProvider(AtomicI64); - fn requested_thread_ids(&self) -> Vec { - self.requested_thread_ids - .lock() - .expect("requested thread IDs lock should not be poisoned") - .clone() +impl Default for TestCurrentTimeProvider { + fn default() -> Self { + Self(AtomicI64::new(FIRST_TIME_UNIX_SECONDS)) } } impl CurrentTimeProvider for TestCurrentTimeProvider { - fn current_time(&self, context: CurrentTimeContext) -> CurrentTimeFuture<'_> { - self.requested_thread_ids - .lock() - .expect("requested thread IDs lock should not be poisoned") - .push(context.thread_id); - let time = self - .times - .lock() - .expect("test times lock should not be poisoned") - .pop_front() - .ok_or_else(|| anyhow!("test current-time provider is exhausted")); - Box::pin(async move { time }) + fn current_time(&self, _thread_id: ThreadId) -> CurrentTimeFuture<'_> { + let timestamp = self.0.fetch_add(60, Ordering::Relaxed); + Box::pin(async move { + Ok(DateTime::::from_timestamp(timestamp, 0) + .expect("test timestamp should be valid")) + }) } } @@ -82,7 +52,7 @@ fn current_time_reminders(request: &ResponsesRequest) -> Vec { request .message_input_texts("developer") .into_iter() - .filter(|text| text.starts_with("It is ") && text.ends_with(" UTC.")) + .filter(|text| text.starts_with("It is ")) .collect() } @@ -111,13 +81,9 @@ async fn varlatency_reminders_follow_request_interval_and_persist_in_history() - ], ) .await; - let provider = Arc::new(TestCurrentTimeProvider::new(&[ - "2026-06-17T17:34:15Z", - "2026-06-17T17:35:15Z", - ])); let test = test_codex() .with_config(|config| enable_varlatency(config, /*interval*/ 2)) - .with_current_time_provider(provider.clone()) + .with_current_time_provider(Arc::new(TestCurrentTimeProvider::default())) .build(&server) .await?; @@ -133,11 +99,6 @@ async fn varlatency_reminders_follow_request_interval_and_persist_in_history() - current_time_reminders(&requests[2]), vec![FIRST_REMINDER, SECOND_REMINDER] ); - assert_eq!( - provider.requested_thread_ids(), - vec![test.session_configured.thread_id; 2] - ); - Ok(()) } @@ -159,10 +120,6 @@ async fn varlatency_reminder_is_refreshed_after_compaction() -> Result<()> { ], ) .await; - let provider = Arc::new(TestCurrentTimeProvider::new(&[ - "2026-06-17T17:34:15Z", - "2026-06-17T17:35:15Z", - ])); let mut model_provider = built_in_model_providers(/*openai_base_url*/ None)["openai"].clone(); model_provider.name = "OpenAI-compatible test provider".to_string(); model_provider.base_url = Some(format!("{}/v1", server.uri())); @@ -172,7 +129,7 @@ async fn varlatency_reminder_is_refreshed_after_compaction() -> Result<()> { config.model_provider = model_provider; enable_varlatency(config, /*interval*/ 50); }) - .with_current_time_provider(provider) + .with_current_time_provider(Arc::new(TestCurrentTimeProvider::default())) .build(&server) .await?; From 994d4830f2a495495b99a5393e4fa8cef4897a9d Mon Sep 17 00:00:00 2001 From: Rohit Arunachalam Date: Wed, 17 Jun 2026 18:19:26 -0700 Subject: [PATCH 3/9] Use current-time reminder naming --- .../core/tests/suite/{varlatency.rs => current_time_reminder.rs} | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename codex-rs/core/tests/suite/{varlatency.rs => current_time_reminder.rs} (100%) diff --git a/codex-rs/core/tests/suite/varlatency.rs b/codex-rs/core/tests/suite/current_time_reminder.rs similarity index 100% rename from codex-rs/core/tests/suite/varlatency.rs rename to codex-rs/core/tests/suite/current_time_reminder.rs From 1dc4eb8f1cf3444bbe6320b407340b59f9b39d12 Mon Sep 17 00:00:00 2001 From: Rohit Arunachalam Date: Wed, 17 Jun 2026 18:19:37 -0700 Subject: [PATCH 4/9] Update current-time reminder references --- codex-rs/core/src/current_time.rs | 12 +++++----- codex-rs/core/src/session/session.rs | 2 +- codex-rs/core/src/session/time_reminder.rs | 2 +- .../core/tests/suite/current_time_reminder.rs | 22 +++++++++---------- codex-rs/core/tests/suite/mod.rs | 2 +- 5 files changed, 20 insertions(+), 20 deletions(-) diff --git a/codex-rs/core/src/current_time.rs b/codex-rs/core/src/current_time.rs index 36cc4b40e80f..091a684397fb 100644 --- a/codex-rs/core/src/current_time.rs +++ b/codex-rs/core/src/current_time.rs @@ -6,10 +6,10 @@ use anyhow::Result; use anyhow::anyhow; use chrono::DateTime; use chrono::Utc; -use codex_features::VarlatencyClockSource; +use codex_features::CurrentTimeSource; use codex_protocol::ThreadId; -use crate::config::VarlatencyConfig; +use crate::config::CurrentTimeReminderConfig; pub type CurrentTimeFuture<'a> = Pin>> + Send + 'a>>; @@ -27,7 +27,7 @@ impl CurrentTimeProvider for SystemCurrentTimeProvider { } pub(crate) fn resolve_current_time_provider( - config: Option<&VarlatencyConfig>, + config: Option<&CurrentTimeReminderConfig>, external_provider: Option>, ) -> Result>> { let Some(config) = config else { @@ -35,10 +35,10 @@ pub(crate) fn resolve_current_time_provider( }; match config.clock_source { - VarlatencyClockSource::System => Ok(Some(Arc::new(SystemCurrentTimeProvider))), - VarlatencyClockSource::AppServerClient => external_provider.map(Some).ok_or_else(|| { + CurrentTimeSource::System => Ok(Some(Arc::new(SystemCurrentTimeProvider))), + CurrentTimeSource::AppServerClient => external_provider.map(Some).ok_or_else(|| { anyhow!( - "features.varlatency.clock_source is app_server_client, but no external current-time provider is available" + "features.current_time_reminder.clock_source is app_server_client, but no external current-time provider is available" ) }), } diff --git a/codex-rs/core/src/session/session.rs b/codex-rs/core/src/session/session.rs index 3ef4f7fdd0f5..fc2ed03745ec 100644 --- a/codex-rs/core/src/session/session.rs +++ b/codex-rs/core/src/session/session.rs @@ -518,7 +518,7 @@ impl Session { InitialHistory::Resumed(resumed_history) => resumed_history.conversation_id, }; let current_time_provider = crate::current_time::resolve_current_time_provider( - config.varlatency.as_ref(), + config.current_time_reminder.as_ref(), external_current_time_provider, )?; let mcp_thread_init = thread_extension_init.clone(); diff --git a/codex-rs/core/src/session/time_reminder.rs b/codex-rs/core/src/session/time_reminder.rs index b17bc6a10c35..9e912dcb33e6 100644 --- a/codex-rs/core/src/session/time_reminder.rs +++ b/codex-rs/core/src/session/time_reminder.rs @@ -29,7 +29,7 @@ pub(super) async fn maybe_record_current_time_reminder( turn_context: &TurnContext, window_id: &str, ) -> CodexResult<()> { - let Some(config) = turn_context.config.varlatency else { + let Some(config) = turn_context.config.current_time_reminder else { return Ok(()); }; diff --git a/codex-rs/core/tests/suite/current_time_reminder.rs b/codex-rs/core/tests/suite/current_time_reminder.rs index 197141327484..3fc8a6a41232 100644 --- a/codex-rs/core/tests/suite/current_time_reminder.rs +++ b/codex-rs/core/tests/suite/current_time_reminder.rs @@ -7,9 +7,9 @@ use chrono::DateTime; use chrono::Utc; use codex_core::CurrentTimeFuture; use codex_core::CurrentTimeProvider; -use codex_core::config::VarlatencyConfig; +use codex_core::config::CurrentTimeReminderConfig; +use codex_features::CurrentTimeSource; use codex_features::Feature; -use codex_features::VarlatencyClockSource; use codex_model_provider_info::built_in_model_providers; use codex_protocol::ThreadId; use codex_protocol::protocol::EventMsg; @@ -56,19 +56,19 @@ fn current_time_reminders(request: &ResponsesRequest) -> Vec { .collect() } -fn enable_varlatency(config: &mut codex_core::config::Config, interval: u64) { +fn enable_current_time_reminder(config: &mut codex_core::config::Config, interval: u64) { config .features - .enable(Feature::Varlatency) - .expect("test config should allow varlatency"); - config.varlatency = Some(VarlatencyConfig { + .enable(Feature::CurrentTimeReminder) + .expect("test config should allow current-time reminders"); + config.current_time_reminder = Some(CurrentTimeReminderConfig { reminder_interval_model_requests: interval, - clock_source: VarlatencyClockSource::AppServerClient, + clock_source: CurrentTimeSource::AppServerClient, }); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn varlatency_reminders_follow_request_interval_and_persist_in_history() -> Result<()> { +async fn current_time_reminders_follow_request_interval_and_persist_in_history() -> Result<()> { skip_if_no_network!(Ok(())); let server = start_mock_server().await; @@ -82,7 +82,7 @@ async fn varlatency_reminders_follow_request_interval_and_persist_in_history() - ) .await; let test = test_codex() - .with_config(|config| enable_varlatency(config, /*interval*/ 2)) + .with_config(|config| enable_current_time_reminder(config, /*interval*/ 2)) .with_current_time_provider(Arc::new(TestCurrentTimeProvider::default())) .build(&server) .await?; @@ -103,7 +103,7 @@ async fn varlatency_reminders_follow_request_interval_and_persist_in_history() - } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn varlatency_reminder_is_refreshed_after_compaction() -> Result<()> { +async fn current_time_reminder_is_refreshed_after_compaction() -> Result<()> { skip_if_no_network!(Ok(())); let server = start_mock_server().await; @@ -127,7 +127,7 @@ async fn varlatency_reminder_is_refreshed_after_compaction() -> Result<()> { let test = test_codex() .with_config(move |config| { config.model_provider = model_provider; - enable_varlatency(config, /*interval*/ 50); + enable_current_time_reminder(config, /*interval*/ 50); }) .with_current_time_provider(Arc::new(TestCurrentTimeProvider::default())) .build(&server) diff --git a/codex-rs/core/tests/suite/mod.rs b/codex-rs/core/tests/suite/mod.rs index 7cba602fbf5b..1e200ab4863e 100644 --- a/codex-rs/core/tests/suite/mod.rs +++ b/codex-rs/core/tests/suite/mod.rs @@ -48,6 +48,7 @@ mod compact; mod compact_remote; mod compact_remote_parity; mod compact_resume_fork; +mod current_time_reminder; mod deprecation_notice; mod exec; mod exec_policy; @@ -126,7 +127,6 @@ mod unified_exec_zsh_fork_approvals; mod unstable_features_warning; mod user_notification; mod user_shell_cmd; -mod varlatency; mod view_image; mod web_search; mod websocket_fallback; From 7b0fa01b894f491bceb994641b7f007c4a13d4bb Mon Sep 17 00:00:00 2001 From: Rohit Arunachalam Date: Wed, 17 Jun 2026 18:53:36 -0700 Subject: [PATCH 5/9] Report current-time provider failures --- codex-rs/core/src/session/turn.rs | 5 +- .../core/tests/suite/current_time_reminder.rs | 64 +++++++++++++++++++ 2 files changed, 68 insertions(+), 1 deletion(-) diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index 57f3db96ed5b..8d4635c9c456 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -233,9 +233,12 @@ pub(crate) async fn run_turn( ) .await { + info!("Turn error: {err:#}"); sess.emit_turn_error_lifecycle(turn_context.as_ref(), err.to_codex_protocol_error()) .await; - error!("Failed to record current-time reminder"); + sess.track_turn_codex_error(turn_context.as_ref(), &err); + let event = EventMsg::Error(err.to_error_event(/*message_prefix*/ None)); + sess.send_event(&turn_context, event).await; return None; } diff --git a/codex-rs/core/tests/suite/current_time_reminder.rs b/codex-rs/core/tests/suite/current_time_reminder.rs index 3fc8a6a41232..6198e0d9256d 100644 --- a/codex-rs/core/tests/suite/current_time_reminder.rs +++ b/codex-rs/core/tests/suite/current_time_reminder.rs @@ -3,6 +3,7 @@ use std::sync::atomic::AtomicI64; use std::sync::atomic::Ordering; use anyhow::Result; +use anyhow::anyhow; use chrono::DateTime; use chrono::Utc; use codex_core::CurrentTimeFuture; @@ -12,12 +13,15 @@ use codex_features::CurrentTimeSource; use codex_features::Feature; use codex_model_provider_info::built_in_model_providers; use codex_protocol::ThreadId; +use codex_protocol::protocol::CodexErrorInfo; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::Op; +use codex_protocol::user_input::UserInput; use core_test_support::responses::ResponsesRequest; use core_test_support::responses::ev_assistant_message; use core_test_support::responses::ev_completed; use core_test_support::responses::ev_response_created; +use core_test_support::responses::mount_sse_once; use core_test_support::responses::mount_sse_sequence; use core_test_support::responses::sse; use core_test_support::responses::start_mock_server; @@ -48,6 +52,14 @@ impl CurrentTimeProvider for TestCurrentTimeProvider { } } +struct FailingCurrentTimeProvider; + +impl CurrentTimeProvider for FailingCurrentTimeProvider { + fn current_time(&self, _thread_id: ThreadId) -> CurrentTimeFuture<'_> { + Box::pin(async { Err(anyhow!("test clock unavailable")) }) + } +} + fn current_time_reminders(request: &ResponsesRequest) -> Vec { request .message_input_texts("developer") @@ -151,3 +163,55 @@ async fn current_time_reminder_is_refreshed_after_compaction() -> Result<()> { Ok(()) } + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn current_time_provider_failure_stops_before_inference() -> Result<()> { + skip_if_no_network!(Ok(())); + + let server = start_mock_server().await; + let responses = mount_sse_once( + &server, + sse(vec![ + ev_response_created("unused-response"), + ev_completed("unused-response"), + ]), + ) + .await; + let test = test_codex() + .with_config(|config| enable_current_time_reminder(config, /*interval*/ 1)) + .with_current_time_provider(Arc::new(FailingCurrentTimeProvider)) + .build(&server) + .await?; + + test.codex + .submit(Op::UserInput { + items: vec![UserInput::Text { + text: "fail before inference".into(), + text_elements: Vec::new(), + }], + final_output_json_schema: None, + responsesapi_client_metadata: None, + additional_context: Default::default(), + thread_settings: Default::default(), + }) + .await?; + + let EventMsg::Error(error) = + wait_for_event(&test.codex, |event| matches!(event, EventMsg::Error(_))).await + else { + unreachable!(); + }; + assert_eq!( + error.message, + "Fatal error: failed to read current time: test clock unavailable" + ); + assert_eq!(error.codex_error_info, Some(CodexErrorInfo::Other)); + + wait_for_event(&test.codex, |event| { + matches!(event, EventMsg::TurnComplete(_)) + }) + .await; + assert!(responses.requests().is_empty()); + + Ok(()) +} From 3b1af392926dc4e5f00383f5bef1661a0c95b003 Mon Sep 17 00:00:00 2001 From: Rohit Arunachalam Date: Thu, 18 Jun 2026 05:47:34 -0700 Subject: [PATCH 6/9] Use external current time source --- codex-rs/core/src/current_time.rs | 4 ++-- codex-rs/core/tests/suite/current_time_reminder.rs | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/codex-rs/core/src/current_time.rs b/codex-rs/core/src/current_time.rs index 091a684397fb..8b27721b79a6 100644 --- a/codex-rs/core/src/current_time.rs +++ b/codex-rs/core/src/current_time.rs @@ -36,9 +36,9 @@ pub(crate) fn resolve_current_time_provider( match config.clock_source { CurrentTimeSource::System => Ok(Some(Arc::new(SystemCurrentTimeProvider))), - CurrentTimeSource::AppServerClient => external_provider.map(Some).ok_or_else(|| { + CurrentTimeSource::External => external_provider.map(Some).ok_or_else(|| { anyhow!( - "features.current_time_reminder.clock_source is app_server_client, but no external current-time provider is available" + "features.current_time_reminder.clock_source is external, but no external current-time provider is available" ) }), } diff --git a/codex-rs/core/tests/suite/current_time_reminder.rs b/codex-rs/core/tests/suite/current_time_reminder.rs index 6198e0d9256d..027b97dc3fdc 100644 --- a/codex-rs/core/tests/suite/current_time_reminder.rs +++ b/codex-rs/core/tests/suite/current_time_reminder.rs @@ -75,7 +75,7 @@ fn enable_current_time_reminder(config: &mut codex_core::config::Config, interva .expect("test config should allow current-time reminders"); config.current_time_reminder = Some(CurrentTimeReminderConfig { reminder_interval_model_requests: interval, - clock_source: CurrentTimeSource::AppServerClient, + clock_source: CurrentTimeSource::External, }); } From 5d5fff5b9e8233b9b92728dad8610b6e31617991 Mon Sep 17 00:00:00 2001 From: Rohit Arunachalam Date: Thu, 18 Jun 2026 06:06:20 -0700 Subject: [PATCH 7/9] Simplify current time provider wiring --- codex-rs/app-server/src/mcp_refresh.rs | 2 +- codex-rs/app-server/src/message_processor.rs | 2 +- codex-rs/core/src/codex_delegate.rs | 2 +- codex-rs/core/src/current_time.rs | 28 +++-- codex-rs/core/src/lib.rs | 4 +- codex-rs/core/src/prompt_debug.rs | 2 +- codex-rs/core/src/session/mod.rs | 8 +- codex-rs/core/src/session/session.rs | 8 +- codex-rs/core/src/session/tests.rs | 10 +- .../core/src/session/tests/guardian_tests.rs | 2 +- codex-rs/core/src/session/time_reminder.rs | 7 +- codex-rs/core/src/state/service.rs | 4 +- codex-rs/core/src/thread_manager.rs | 12 +-- codex-rs/core/src/thread_manager_tests.rs | 22 ++-- .../src/tools/handlers/multi_agents_tests.rs | 2 +- codex-rs/core/tests/common/test_codex.rs | 12 +-- codex-rs/core/tests/suite/client.rs | 2 +- .../core/tests/suite/current_time_reminder.rs | 101 ++++++++++++++---- codex-rs/mcp-server/src/message_processor.rs | 2 +- codex-rs/thread-manager-sample/src/main.rs | 2 +- 20 files changed, 142 insertions(+), 92 deletions(-) diff --git a/codex-rs/app-server/src/mcp_refresh.rs b/codex-rs/app-server/src/mcp_refresh.rs index 078664b46a22..84a6b8e83b6d 100644 --- a/codex-rs/app-server/src/mcp_refresh.rs +++ b/codex-rs/app-server/src/mcp_refresh.rs @@ -250,7 +250,7 @@ mod tests { Some(state_db.clone()), "11111111-1111-4111-8111-111111111111".to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ) }); thread_manager.start_thread(good_config).await?; diff --git a/codex-rs/app-server/src/message_processor.rs b/codex-rs/app-server/src/message_processor.rs index 0d1d773cdeea..0f9f6d3321fa 100644 --- a/codex-rs/app-server/src/message_processor.rs +++ b/codex-rs/app-server/src/message_processor.rs @@ -379,7 +379,7 @@ impl MessageProcessor { outgoing.clone(), thread_state_manager.clone(), )), - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ) }); let models_manager = thread_manager.get_models_manager(); diff --git a/codex-rs/core/src/codex_delegate.rs b/codex-rs/core/src/codex_delegate.rs index 5becea1d313a..dab902ec3741 100644 --- a/codex-rs/core/src/codex_delegate.rs +++ b/codex-rs/core/src/codex_delegate.rs @@ -120,7 +120,7 @@ pub(crate) async fn run_codex_thread_interactive( analytics_events_client: Some(parent_session.services.analytics_events_client.clone()), thread_store: Arc::clone(&parent_session.services.thread_store), attestation_provider: parent_session.services.attestation_provider.clone(), - external_current_time_provider: parent_session.services.current_time_provider.clone(), + external_time_provider: Some(Arc::clone(&parent_session.services.time_provider)), inherited_multi_agent_version: Some(MultiAgentVersion::Disabled), })) .or_cancel(&cancel_token) diff --git a/codex-rs/core/src/current_time.rs b/codex-rs/core/src/current_time.rs index 8b27721b79a6..7740552e1e5c 100644 --- a/codex-rs/core/src/current_time.rs +++ b/codex-rs/core/src/current_time.rs @@ -11,32 +11,28 @@ use codex_protocol::ThreadId; use crate::config::CurrentTimeReminderConfig; -pub type CurrentTimeFuture<'a> = Pin>> + Send + 'a>>; +pub type TimeFuture<'a> = Pin>> + Send + 'a>>; /// Host integration boundary for obtaining the current time. -pub trait CurrentTimeProvider: Send + Sync { - fn current_time(&self, thread_id: ThreadId) -> CurrentTimeFuture<'_>; +pub trait TimeProvider: Send + Sync { + fn current_time(&self, thread_id: ThreadId) -> TimeFuture<'_>; } -struct SystemCurrentTimeProvider; +pub(crate) struct SystemTimeProvider; -impl CurrentTimeProvider for SystemCurrentTimeProvider { - fn current_time(&self, _thread_id: ThreadId) -> CurrentTimeFuture<'_> { +impl TimeProvider for SystemTimeProvider { + fn current_time(&self, _thread_id: ThreadId) -> TimeFuture<'_> { Box::pin(async { Ok(Utc::now()) }) } } -pub(crate) fn resolve_current_time_provider( +pub(crate) fn resolve_time_provider( config: Option<&CurrentTimeReminderConfig>, - external_provider: Option>, -) -> Result>> { - let Some(config) = config else { - return Ok(None); - }; - - match config.clock_source { - CurrentTimeSource::System => Ok(Some(Arc::new(SystemCurrentTimeProvider))), - CurrentTimeSource::External => external_provider.map(Some).ok_or_else(|| { + external_provider: Option>, +) -> Result> { + match config.map(|config| config.clock_source).unwrap_or_default() { + CurrentTimeSource::System => Ok(Arc::new(SystemTimeProvider)), + CurrentTimeSource::External => external_provider.ok_or_else(|| { anyhow!( "features.current_time_reminder.clock_source is external, but no external current-time provider is available" ) diff --git a/codex-rs/core/src/lib.rs b/codex-rs/core/src/lib.rs index 94458e5431dc..121c8d2be2ea 100644 --- a/codex-rs/core/src/lib.rs +++ b/codex-rs/core/src/lib.rs @@ -187,8 +187,8 @@ pub use client_common::ResponseEvent; pub use client_common::ResponseStream; pub use codex_prompts::REVIEW_PROMPT; pub use compact::content_items_to_text; -pub use current_time::CurrentTimeFuture; -pub use current_time::CurrentTimeProvider; +pub use current_time::TimeFuture; +pub use current_time::TimeProvider; pub use event_mapping::parse_turn_item; pub use exec_policy::ExecPolicyError; pub use exec_policy::check_execpolicy_for_warnings; diff --git a/codex-rs/core/src/prompt_debug.rs b/codex-rs/core/src/prompt_debug.rs index 31f63ccceaa6..bd9289b77030 100644 --- a/codex-rs/core/src/prompt_debug.rs +++ b/codex-rs/core/src/prompt_debug.rs @@ -60,7 +60,7 @@ pub async fn build_prompt_input( state_db.clone(), installation_id, /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let thread = thread_manager.start_thread(config).await?; diff --git a/codex-rs/core/src/session/mod.rs b/codex-rs/core/src/session/mod.rs index bb1e7bdde36a..b55c4b11d771 100644 --- a/codex-rs/core/src/session/mod.rs +++ b/codex-rs/core/src/session/mod.rs @@ -31,7 +31,7 @@ use crate::context::NetworkRuleSaved; use crate::context::PermissionsInstructions; use crate::context::PersonalitySpecInstructions; use crate::context::RecommendedPluginsInstructions; -use crate::current_time::CurrentTimeProvider; +use crate::current_time::TimeProvider; use crate::default_skill_metadata_budget; use crate::environment_selection::TurnEnvironmentSnapshot; use crate::exec_policy::ExecPolicyManager; @@ -441,7 +441,7 @@ pub(crate) struct CodexSpawnArgs { pub(crate) analytics_events_client: Option, pub(crate) thread_store: Arc, pub(crate) attestation_provider: Option>, - pub(crate) external_current_time_provider: Option>, + pub(crate) external_time_provider: Option>, pub(crate) inherited_multi_agent_version: Option, } @@ -525,7 +525,7 @@ impl Codex { analytics_events_client, thread_store, attestation_provider, - external_current_time_provider, + external_time_provider, inherited_multi_agent_version, } = args; let (tx_sub, rx_sub) = async_channel::bounded(SUBMISSION_CHANNEL_CAPACITY); @@ -673,7 +673,7 @@ impl Codex { thread_store, parent_rollout_thread_trace, attestation_provider, - external_current_time_provider, + external_time_provider, multi_agent_version, )) .await diff --git a/codex-rs/core/src/session/session.rs b/codex-rs/core/src/session/session.rs index fc2ed03745ec..af2d1abeb89f 100644 --- a/codex-rs/core/src/session/session.rs +++ b/codex-rs/core/src/session/session.rs @@ -492,7 +492,7 @@ impl Session { thread_store: Arc, parent_rollout_thread_trace: ThreadTraceContext, attestation_provider: Option>, - external_current_time_provider: Option>, + external_time_provider: Option>, multi_agent_version: Option, ) -> anyhow::Result> { debug!( @@ -517,9 +517,9 @@ impl Session { } InitialHistory::Resumed(resumed_history) => resumed_history.conversation_id, }; - let current_time_provider = crate::current_time::resolve_current_time_provider( + let time_provider = crate::current_time::resolve_time_provider( config.current_time_reminder.as_ref(), - external_current_time_provider, + external_time_provider, )?; let mcp_thread_init = thread_extension_init.clone(); let thread_extension_data = codex_extension_api::ExtensionData::new_with_init( @@ -1027,7 +1027,7 @@ impl Session { live_thread: live_thread_init.as_ref().cloned(), thread_store: Arc::clone(&thread_store), attestation_provider: attestation_provider.clone(), - current_time_provider, + time_provider, model_client: ModelClient::new( Some(Arc::clone(&auth_manager)), thread_id, diff --git a/codex-rs/core/src/session/tests.rs b/codex-rs/core/src/session/tests.rs index 56a1e368fdf5..0b1cc045c44a 100644 --- a/codex-rs/core/src/session/tests.rs +++ b/codex-rs/core/src/session/tests.rs @@ -4867,7 +4867,7 @@ async fn session_new_fails_when_zsh_fork_enabled_without_packaged_zsh() { )), codex_rollout_trace::ThreadTraceContext::disabled(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, Some(config.multi_agent_version_from_features()), ) .await; @@ -5033,7 +5033,7 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) { /*state_db*/ None, )), attestation_provider: None, - current_time_provider: None, + time_provider: Arc::new(crate::current_time::SystemTimeProvider), model_client: ModelClient::new( Some(auth_manager.clone()), thread_id, @@ -5218,7 +5218,7 @@ async fn make_session_with_config_and_rx( )), codex_rollout_trace::ThreadTraceContext::disabled(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, Some(config.multi_agent_version_from_features()), ) .await?; @@ -5331,7 +5331,7 @@ async fn make_session_with_history_source_and_agent_control_and_rx( )), codex_rollout_trace::ThreadTraceContext::disabled(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, Some(config.multi_agent_version_from_features()), ) .await?; @@ -7082,7 +7082,7 @@ where state_db, )), attestation_provider: None, - current_time_provider: None, + time_provider: Arc::new(crate::current_time::SystemTimeProvider), model_client: ModelClient::new( Some(Arc::clone(&auth_manager)), thread_id, diff --git a/codex-rs/core/src/session/tests/guardian_tests.rs b/codex-rs/core/src/session/tests/guardian_tests.rs index e3451df56863..5de331ee1fb1 100644 --- a/codex-rs/core/src/session/tests/guardian_tests.rs +++ b/codex-rs/core/src/session/tests/guardian_tests.rs @@ -736,7 +736,7 @@ async fn guardian_subagent_does_not_inherit_parent_exec_policy_rules() { analytics_events_client: None, thread_store, attestation_provider: None, - external_current_time_provider: None, + external_time_provider: None, inherited_multi_agent_version: None, }) .await diff --git a/codex-rs/core/src/session/time_reminder.rs b/codex-rs/core/src/session/time_reminder.rs index 9e912dcb33e6..d0708c5fa5e7 100644 --- a/codex-rs/core/src/session/time_reminder.rs +++ b/codex-rs/core/src/session/time_reminder.rs @@ -43,12 +43,9 @@ pub(super) async fn maybe_record_current_time_reminder( return Ok(()); } - let provider = sess + let current_time = sess .services - .current_time_provider - .as_ref() - .ok_or_else(|| CodexErr::Fatal("current-time provider is not configured".to_string()))?; - let current_time = provider + .time_provider .current_time(sess.thread_id) .await .map_err(|err| CodexErr::Fatal(format!("failed to read current time: {err:#}")))?; diff --git a/codex-rs/core/src/state/service.rs b/codex-rs/core/src/state/service.rs index 059362e6e039..853b96e921d7 100644 --- a/codex-rs/core/src/state/service.rs +++ b/codex-rs/core/src/state/service.rs @@ -8,7 +8,7 @@ use crate::attestation::AttestationProvider; use crate::client::ModelClient; use crate::config::NetworkProxyAuditMetadata; use crate::config::StartedNetworkProxy; -use crate::current_time::CurrentTimeProvider; +use crate::current_time::TimeProvider; use crate::environment_selection::ThreadEnvironments; use crate::exec_policy::ExecPolicyManager; use crate::guardian::GuardianRejection; @@ -80,7 +80,7 @@ pub(crate) struct SessionServices { pub(crate) live_thread: Option, pub(crate) thread_store: Arc, pub(crate) attestation_provider: Option>, - pub(crate) current_time_provider: Option>, + pub(crate) time_provider: Arc, /// Session-scoped model client shared across turns. pub(crate) model_client: ModelClient, pub(crate) code_mode_service: CodeModeService, diff --git a/codex-rs/core/src/thread_manager.rs b/codex-rs/core/src/thread_manager.rs index 3d84e5251046..814c8fe764d4 100644 --- a/codex-rs/core/src/thread_manager.rs +++ b/codex-rs/core/src/thread_manager.rs @@ -4,7 +4,7 @@ use crate::attestation::AttestationProvider; use crate::codex_thread::CodexThread; use crate::config::Config; use crate::config::ThreadStoreConfig; -use crate::current_time::CurrentTimeProvider; +use crate::current_time::TimeProvider; use crate::environment_selection::TurnEnvironmentSnapshot; use crate::environment_selection::default_thread_environment_selections; use crate::mcp::McpManager; @@ -216,7 +216,7 @@ pub(crate) struct ThreadManagerState { user_instructions_provider: Arc, thread_store: Arc, attestation_provider: Option>, - external_current_time_provider: Option>, + external_time_provider: Option>, session_source: SessionSource, installation_id: String, analytics_events_client: Option, @@ -271,7 +271,7 @@ impl ThreadManager { state_db: Option, installation_id: String, attestation_provider: Option>, - external_current_time_provider: Option>, + external_time_provider: Option>, ) -> Self { let codex_home = config.codex_home.clone(); let restriction_product = session_source.restriction_product(); @@ -303,7 +303,7 @@ impl ThreadManager { user_instructions_provider, thread_store, attestation_provider, - external_current_time_provider, + external_time_provider, auth_manager, session_source, installation_id, @@ -409,7 +409,7 @@ impl ThreadManager { ), thread_store, attestation_provider: None, - external_current_time_provider: None, + external_time_provider: None, auth_manager, session_source: SessionSource::Exec, installation_id, @@ -1471,7 +1471,7 @@ impl ThreadManagerState { analytics_events_client: self.analytics_events_client.clone(), thread_store: Arc::clone(&self.thread_store), attestation_provider: self.attestation_provider.clone(), - external_current_time_provider: self.external_current_time_provider.clone(), + external_time_provider: self.external_time_provider.clone(), inherited_multi_agent_version: multi_agent_version, })) .await?; diff --git a/codex-rs/core/src/thread_manager_tests.rs b/codex-rs/core/src/thread_manager_tests.rs index d843de96bbb2..bfb7865bb219 100644 --- a/codex-rs/core/src/thread_manager_tests.rs +++ b/codex-rs/core/src/thread_manager_tests.rs @@ -440,7 +440,7 @@ async fn start_thread_seeds_extension_data_for_mcp_and_lifecycle_contributors() /*state_db*/ None, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let selected_root_init = |id: &str, environment_id: &str| { let mut init = codex_extension_api::ExtensionDataInit::new(); @@ -549,7 +549,7 @@ async fn resume_and_fork_do_not_restore_thread_environments_from_rollout() { /*state_db*/ None, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let selected_cwd = AbsolutePathBuf::try_from(config.cwd.as_path().join("selected")).expect("absolute path"); @@ -673,7 +673,7 @@ async fn explicit_installation_id_skips_codex_home_file() { state_db.clone(), installation_id.clone(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let thread = manager @@ -714,7 +714,7 @@ async fn resume_active_thread_from_rollout_returns_running_thread() { /*state_db*/ None, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let source = manager @@ -774,7 +774,7 @@ async fn resume_stopped_thread_from_rollout_spawns_new_thread() { /*state_db*/ None, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let source = manager @@ -841,7 +841,7 @@ async fn resume_stopped_thread_from_rollout_preserves_thread_source() { state_db.clone(), TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let source = manager @@ -935,7 +935,7 @@ async fn rollout_path_resume_and_fork_read_history_through_thread_store() { state_db, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let source = manager @@ -1040,7 +1040,7 @@ async fn new_uses_active_provider_for_model_refresh() { /*state_db*/ None, TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let _ = manager.list_models(RefreshStrategy::Online).await; @@ -1262,7 +1262,7 @@ async fn interrupted_fork_snapshot_does_not_synthesize_turn_id_for_legacy_histor state_db.clone(), TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let source = manager @@ -1371,7 +1371,7 @@ async fn interrupted_fork_snapshot_preserves_explicit_turn_id() { state_db.clone(), TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let source = manager @@ -1470,7 +1470,7 @@ async fn interrupted_fork_snapshot_uses_persisted_mid_turn_history_without_live_ state_db.clone(), TEST_INSTALLATION_ID.to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let source = manager diff --git a/codex-rs/core/src/tools/handlers/multi_agents_tests.rs b/codex-rs/core/src/tools/handlers/multi_agents_tests.rs index eabbcbbfa3c8..eefea1eb8fed 100644 --- a/codex-rs/core/src/tools/handlers/multi_agents_tests.rs +++ b/codex-rs/core/src/tools/handlers/multi_agents_tests.rs @@ -4279,7 +4279,7 @@ async fn tool_handlers_cascade_close_and_resume_and_keep_explicitly_closed_subtr state_db.clone(), "11111111-1111-4111-8111-111111111111".to_string(), /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let parent = manager diff --git a/codex-rs/core/tests/common/test_codex.rs b/codex-rs/core/tests/common/test_codex.rs index ba3c66c1eb2e..42b9bf87ca28 100644 --- a/codex-rs/core/tests/common/test_codex.rs +++ b/codex-rs/core/tests/common/test_codex.rs @@ -15,9 +15,9 @@ use anyhow::Result; use anyhow::anyhow; use codex_config::CloudConfigBundleLoader; use codex_core::CodexThread; -use codex_core::CurrentTimeProvider; use codex_core::StartThreadOptions; use codex_core::ThreadManager; +use codex_core::TimeProvider; use codex_core::config::Config; use codex_core::resolve_installation_id; use codex_core::shell::Shell; @@ -262,7 +262,7 @@ pub struct TestCodexBuilder { extensions: Arc>, user_instructions_provider: Option>, supports_openai_form_elicitation: bool, - current_time_provider: Option>, + external_time_provider: Option>, } impl TestCodexBuilder { @@ -364,8 +364,8 @@ impl TestCodexBuilder { self } - pub fn with_current_time_provider(mut self, provider: Arc) -> Self { - self.current_time_provider = Some(provider); + pub fn with_external_time_provider(mut self, provider: Arc) -> Self { + self.external_time_provider = Some(provider); self } @@ -575,7 +575,7 @@ impl TestCodexBuilder { state_db.clone(), installation_id, /*attestation_provider*/ None, - /*external_current_time_provider*/ self.current_time_provider.clone(), + /*external_time_provider*/ self.external_time_provider.clone(), ); let thread_manager = Arc::new(thread_manager); let user_shell_override = self.user_shell_override.clone(); @@ -1180,7 +1180,7 @@ pub fn test_codex() -> TestCodexBuilder { extensions: empty_extension_registry(), user_instructions_provider: None, supports_openai_form_elicitation: false, - current_time_provider: None, + external_time_provider: None, } } diff --git a/codex-rs/core/tests/suite/client.rs b/codex-rs/core/tests/suite/client.rs index 99da34e0ff9d..e32961882858 100644 --- a/codex-rs/core/tests/suite/client.rs +++ b/codex-rs/core/tests/suite/client.rs @@ -1259,7 +1259,7 @@ async fn prefers_apikey_when_config_prefers_apikey_even_with_chatgpt_tokens() { /*state_db*/ None, installation_id, /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let NewThread { thread: codex, .. } = thread_manager .start_thread(config.clone()) diff --git a/codex-rs/core/tests/suite/current_time_reminder.rs b/codex-rs/core/tests/suite/current_time_reminder.rs index 027b97dc3fdc..5b5970a0eeb6 100644 --- a/codex-rs/core/tests/suite/current_time_reminder.rs +++ b/codex-rs/core/tests/suite/current_time_reminder.rs @@ -6,20 +6,23 @@ use anyhow::Result; use anyhow::anyhow; use chrono::DateTime; use chrono::Utc; -use codex_core::CurrentTimeFuture; -use codex_core::CurrentTimeProvider; +use codex_core::TimeFuture; +use codex_core::TimeProvider; use codex_core::config::CurrentTimeReminderConfig; use codex_features::CurrentTimeSource; use codex_features::Feature; use codex_model_provider_info::built_in_model_providers; use codex_protocol::ThreadId; +use codex_protocol::models::PermissionProfile; use codex_protocol::protocol::CodexErrorInfo; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::Op; use codex_protocol::user_input::UserInput; +use core_test_support::assert_regex_match; use core_test_support::responses::ResponsesRequest; use core_test_support::responses::ev_assistant_message; use core_test_support::responses::ev_completed; +use core_test_support::responses::ev_function_call; use core_test_support::responses::ev_response_created; use core_test_support::responses::mount_sse_once; use core_test_support::responses::mount_sse_sequence; @@ -29,21 +32,22 @@ use core_test_support::skip_if_no_network; use core_test_support::test_codex::test_codex; use core_test_support::wait_for_event; use pretty_assertions::assert_eq; +use serde_json::json; const FIRST_REMINDER: &str = "It is 2026-06-17 17:34:15 UTC."; const SECOND_REMINDER: &str = "It is 2026-06-17 17:35:15 UTC."; const FIRST_TIME_UNIX_SECONDS: i64 = 1_781_717_655; -struct TestCurrentTimeProvider(AtomicI64); +struct TestTimeProvider(AtomicI64); -impl Default for TestCurrentTimeProvider { +impl Default for TestTimeProvider { fn default() -> Self { Self(AtomicI64::new(FIRST_TIME_UNIX_SECONDS)) } } -impl CurrentTimeProvider for TestCurrentTimeProvider { - fn current_time(&self, _thread_id: ThreadId) -> CurrentTimeFuture<'_> { +impl TimeProvider for TestTimeProvider { + fn current_time(&self, _thread_id: ThreadId) -> TimeFuture<'_> { let timestamp = self.0.fetch_add(60, Ordering::Relaxed); Box::pin(async move { Ok(DateTime::::from_timestamp(timestamp, 0) @@ -52,10 +56,10 @@ impl CurrentTimeProvider for TestCurrentTimeProvider { } } -struct FailingCurrentTimeProvider; +struct FailingTimeProvider; -impl CurrentTimeProvider for FailingCurrentTimeProvider { - fn current_time(&self, _thread_id: ThreadId) -> CurrentTimeFuture<'_> { +impl TimeProvider for FailingTimeProvider { + fn current_time(&self, _thread_id: ThreadId) -> TimeFuture<'_> { Box::pin(async { Err(anyhow!("test clock unavailable")) }) } } @@ -68,14 +72,18 @@ fn current_time_reminders(request: &ResponsesRequest) -> Vec { .collect() } -fn enable_current_time_reminder(config: &mut codex_core::config::Config, interval: u64) { +fn enable_current_time_reminder( + config: &mut codex_core::config::Config, + interval: u64, + clock_source: CurrentTimeSource, +) { config .features .enable(Feature::CurrentTimeReminder) .expect("test config should allow current-time reminders"); config.current_time_reminder = Some(CurrentTimeReminderConfig { reminder_interval_model_requests: interval, - clock_source: CurrentTimeSource::External, + clock_source, }); } @@ -84,24 +92,42 @@ async fn current_time_reminders_follow_request_interval_and_persist_in_history() skip_if_no_network!(Ok(())); let server = start_mock_server().await; + let tool_args = json!({ + "command": "echo current time", + "timeout_ms": 1_000, + }); let responses = mount_sse_sequence( &server, vec![ - sse(vec![ev_response_created("resp-1"), ev_completed("resp-1")]), - sse(vec![ev_response_created("resp-2"), ev_completed("resp-2")]), + sse(vec![ + ev_response_created("resp-1"), + ev_function_call( + "current-time-tool-call", + "shell_command", + &serde_json::to_string(&tool_args)?, + ), + ev_completed("resp-1"), + ]), + sse(vec![ + ev_response_created("resp-2"), + ev_assistant_message("msg-2", "done"), + ev_completed("resp-2"), + ]), sse(vec![ev_response_created("resp-3"), ev_completed("resp-3")]), ], ) .await; let test = test_codex() - .with_config(|config| enable_current_time_reminder(config, /*interval*/ 2)) - .with_current_time_provider(Arc::new(TestCurrentTimeProvider::default())) + .with_config(|config| { + enable_current_time_reminder(config, /*interval*/ 2, CurrentTimeSource::External) + }) + .with_external_time_provider(Arc::new(TestTimeProvider::default())) .build(&server) .await?; - test.submit_turn("first turn").await?; + test.submit_turn_with_permission_profile("first turn", PermissionProfile::Disabled) + .await?; test.submit_turn("second turn").await?; - test.submit_turn("third turn").await?; let requests = responses.requests(); assert_eq!(requests.len(), 3); @@ -114,6 +140,35 @@ async fn current_time_reminders_follow_request_interval_and_persist_in_history() Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn system_time_source_adds_current_time_reminder() -> Result<()> { + skip_if_no_network!(Ok(())); + + let server = start_mock_server().await; + let responses = mount_sse_once( + &server, + sse(vec![ev_response_created("resp-1"), ev_completed("resp-1")]), + ) + .await; + let test = test_codex() + .with_config(|config| { + enable_current_time_reminder(config, /*interval*/ 1, CurrentTimeSource::System) + }) + .build(&server) + .await?; + + test.submit_turn("what time is it?").await?; + + let reminders = current_time_reminders(&responses.single_request()); + assert_eq!(reminders.len(), 1); + assert_regex_match( + r"^It is \d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2} UTC\.$", + &reminders[0], + ); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn current_time_reminder_is_refreshed_after_compaction() -> Result<()> { skip_if_no_network!(Ok(())); @@ -139,9 +194,9 @@ async fn current_time_reminder_is_refreshed_after_compaction() -> Result<()> { let test = test_codex() .with_config(move |config| { config.model_provider = model_provider; - enable_current_time_reminder(config, /*interval*/ 50); + enable_current_time_reminder(config, /*interval*/ 50, CurrentTimeSource::External); }) - .with_current_time_provider(Arc::new(TestCurrentTimeProvider::default())) + .with_external_time_provider(Arc::new(TestTimeProvider::default())) .build(&server) .await?; @@ -165,7 +220,7 @@ async fn current_time_reminder_is_refreshed_after_compaction() -> Result<()> { } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn current_time_provider_failure_stops_before_inference() -> Result<()> { +async fn time_provider_failure_stops_before_inference() -> Result<()> { skip_if_no_network!(Ok(())); let server = start_mock_server().await; @@ -178,8 +233,10 @@ async fn current_time_provider_failure_stops_before_inference() -> Result<()> { ) .await; let test = test_codex() - .with_config(|config| enable_current_time_reminder(config, /*interval*/ 1)) - .with_current_time_provider(Arc::new(FailingCurrentTimeProvider)) + .with_config(|config| { + enable_current_time_reminder(config, /*interval*/ 1, CurrentTimeSource::External) + }) + .with_external_time_provider(Arc::new(FailingTimeProvider)) .build(&server) .await?; diff --git a/codex-rs/mcp-server/src/message_processor.rs b/codex-rs/mcp-server/src/message_processor.rs index 088b455cac3d..5f36091feef2 100644 --- a/codex-rs/mcp-server/src/message_processor.rs +++ b/codex-rs/mcp-server/src/message_processor.rs @@ -78,7 +78,7 @@ impl MessageProcessor { state_db.clone(), installation_id, /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, )); Self { outgoing, diff --git a/codex-rs/thread-manager-sample/src/main.rs b/codex-rs/thread-manager-sample/src/main.rs index 778c37b0a63a..26ba71ae4669 100644 --- a/codex-rs/thread-manager-sample/src/main.rs +++ b/codex-rs/thread-manager-sample/src/main.rs @@ -136,7 +136,7 @@ async fn run_main(arg0_paths: Arg0DispatchPaths) -> anyhow::Result<()> { state_db, installation_id, /*attestation_provider*/ None, - /*external_current_time_provider*/ None, + /*external_time_provider*/ None, ); let NewThread { From 46fa2cdf58143916f3194d81c3df1c620da4fd69 Mon Sep 17 00:00:00 2001 From: Rohit Arunachalam Date: Thu, 18 Jun 2026 06:11:31 -0700 Subject: [PATCH 8/9] Simplify current time reminder state --- codex-rs/core/src/session/time_reminder.rs | 20 ++++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/codex-rs/core/src/session/time_reminder.rs b/codex-rs/core/src/session/time_reminder.rs index d0708c5fa5e7..058132a8956c 100644 --- a/codex-rs/core/src/session/time_reminder.rs +++ b/codex-rs/core/src/session/time_reminder.rs @@ -12,15 +12,17 @@ pub(crate) struct CurrentTimeReminderState { } impl CurrentTimeReminderState { - fn begin_model_request(&mut self, window_id: &str, interval: u64) -> bool { + fn take_reminder_due(&mut self, window_id: &str, interval: u64) -> bool { self.model_requests_since_delivery = self.model_requests_since_delivery.saturating_add(1); - self.last_window_id.as_deref() != Some(window_id) - || self.model_requests_since_delivery >= interval - } + let reminder_is_due = self.last_window_id.as_deref() != Some(window_id) + || self.model_requests_since_delivery >= interval; + + if reminder_is_due { + self.model_requests_since_delivery = 0; + self.last_window_id = Some(window_id.to_string()); + } - fn record_delivery(&mut self, window_id: &str) { - self.model_requests_since_delivery = 0; - self.last_window_id = Some(window_id.to_string()); + reminder_is_due } } @@ -37,7 +39,7 @@ pub(super) async fn maybe_record_current_time_reminder( let mut state = sess.state.lock().await; state .current_time_reminder - .begin_model_request(window_id, config.reminder_interval_model_requests) + .take_reminder_due(window_id, config.reminder_interval_model_requests) }; if !reminder_is_due { return Ok(()); @@ -55,7 +57,5 @@ pub(super) async fn maybe_record_current_time_reminder( sess.record_conversation_items(turn_context, std::slice::from_ref(&response_item)) .await; - let mut state = sess.state.lock().await; - state.current_time_reminder.record_delivery(window_id); Ok(()) } From ba927faa21b063a6dabd10a2d23ecd0ade301962 Mon Sep 17 00:00:00 2001 From: Rohit Arunachalam Date: Thu, 18 Jun 2026 09:08:11 -0700 Subject: [PATCH 9/9] Unify current time reminder errors --- codex-rs/core/src/session/turn.rs | 82 +++++++++++++++---------------- 1 file changed, 41 insertions(+), 41 deletions(-) diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index 8d4635c9c456..25e4e1a6afd5 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -226,50 +226,50 @@ pub(crate) async fn run_turn( ) .await; - if let Err(err) = super::time_reminder::maybe_record_current_time_reminder( - sess.as_ref(), - turn_context.as_ref(), - &window_id, - ) - .await - { - info!("Turn error: {err:#}"); - sess.emit_turn_error_lifecycle(turn_context.as_ref(), err.to_codex_protocol_error()) - .await; - sess.track_turn_codex_error(turn_context.as_ref(), &err); - let event = EventMsg::Error(err.to_error_event(/*message_prefix*/ None)); - sess.send_event(&turn_context, event).await; - return None; - } + let sampling_request_result: CodexResult<_> = async { + super::time_reminder::maybe_record_current_time_reminder( + sess.as_ref(), + turn_context.as_ref(), + &window_id, + ) + .await?; - // Construct the input that we will send to the model. - let sampling_request_input: Vec = async { - sess.clone_history() - .await - .for_prompt(&turn_context.model_info.input_modalities) + // Construct the input that we will send to the model. + let sampling_request_input: Vec = async { + sess.clone_history() + .await + .for_prompt(&turn_context.model_info.input_modalities) + } + .instrument(trace_span!("run_turn.prepare_sampling_request_input")) + .await; + + let responses_metadata = turn_context.turn_metadata_state.to_responses_metadata( + sess.installation_id.clone(), + window_id, + CodexResponsesRequestKind::Turn, + ); + let tokens_before_sampling = sess.get_total_token_usage().await; + let (sampling_request_output, sampling_request_input) = run_sampling_request( + Arc::clone(&sess), + Arc::clone(&turn_context), + Arc::clone(&turn_extension_data), + Arc::clone(&turn_diff_tracker), + &mut client_session, + &responses_metadata, + sampling_request_input, + cancellation_token.child_token(), + ) + .await?; + + Ok(( + tokens_before_sampling, + sampling_request_output, + sampling_request_input, + )) } - .instrument(trace_span!("run_turn.prepare_sampling_request_input")) .await; - - let responses_metadata = turn_context.turn_metadata_state.to_responses_metadata( - sess.installation_id.clone(), - window_id, - CodexResponsesRequestKind::Turn, - ); - let tokens_before_sampling = sess.get_total_token_usage().await; - match run_sampling_request( - Arc::clone(&sess), - Arc::clone(&turn_context), - Arc::clone(&turn_extension_data), - Arc::clone(&turn_diff_tracker), - &mut client_session, - &responses_metadata, - sampling_request_input, - cancellation_token.child_token(), - ) - .await - { - Ok((sampling_request_output, sampling_request_input)) => { + match sampling_request_result { + Ok((tokens_before_sampling, sampling_request_output, sampling_request_input)) => { let SamplingRequestResult { needs_follow_up: model_needs_follow_up, last_agent_message: sampling_request_last_agent_message,