From 752ed90d78ad9444dd449fbc8e71ad2121865be0 Mon Sep 17 00:00:00 2001 From: rka-oai Date: Thu, 18 Jun 2026 12:18:42 -0700 Subject: [PATCH] current time reminders impl for system clock (varlatency 2/n) (#28824) Stacked on #28822. ## Summary - add a host-injectable current-time provider with a built-in system implementation - record UTC developer reminders in history immediately before due model requests - keep cadence state per session and force a refresh after compaction This does NOT include the app server client <-> server clock logic. This PR is only for the reminder message & system clock that will be used in prod. ## Testing - `just test -p codex-core varlatency_` - `just clippy -p codex-core -p codex-app-server -p codex-mcp-server -p codex-thread-manager-sample` - `just fmt` --- 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 | 35 +++ codex-rs/core/src/context/mod.rs | 2 + codex-rs/core/src/current_time.rs | 41 +++ codex-rs/core/src/lib.rs | 3 + 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 | 61 ++++ codex-rs/core/src/session/turn.rs | 70 +++-- 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 + .../core/tests/suite/current_time_reminder.rs | 274 ++++++++++++++++++ codex-rs/core/tests/suite/mod.rs | 1 + codex-rs/mcp-server/src/message_processor.rs | 1 + codex-rs/thread-manager-sample/src/main.rs | 1 + 25 files changed, 516 insertions(+), 27 deletions(-) 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/current_time_reminder.rs diff --git a/codex-rs/app-server/src/mcp_refresh.rs b/codex-rs/app-server/src/mcp_refresh.rs index a5659174f..84a6b8e83 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_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 9e9437bd3..0f9f6d332 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_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 b7b1fe758..dab902ec3 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_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/context/current_time_reminder.rs b/codex-rs/core/src/context/current_time_reminder.rs new file mode 100644 index 000000000..0e3fe0b1e --- /dev/null +++ b/codex-rs/core/src/context/current_time_reminder.rs @@ -0,0 +1,35 @@ +use chrono::DateTime; +use chrono::Utc; + +use super::ContextualUserFragment; + +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 7fc137578..3c2f48723 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 000000000..7740552e1 --- /dev/null +++ b/codex-rs/core/src/current_time.rs @@ -0,0 +1,41 @@ +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::CurrentTimeSource; +use codex_protocol::ThreadId; + +use crate::config::CurrentTimeReminderConfig; + +pub type TimeFuture<'a> = Pin>> + Send + 'a>>; + +/// Host integration boundary for obtaining the current time. +pub trait TimeProvider: Send + Sync { + fn current_time(&self, thread_id: ThreadId) -> TimeFuture<'_>; +} + +pub(crate) struct SystemTimeProvider; + +impl TimeProvider for SystemTimeProvider { + fn current_time(&self, _thread_id: ThreadId) -> TimeFuture<'_> { + Box::pin(async { Ok(Utc::now()) }) + } +} + +pub(crate) fn resolve_time_provider( + config: Option<&CurrentTimeReminderConfig>, + 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 0e193baa1..121c8d2be 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,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::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 0da42d8a4..bd9289b77 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_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 f92bd6368..b55c4b11d 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::TimeProvider; 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_time_provider: Option>, pub(crate) inherited_multi_agent_version: Option, } @@ -522,6 +525,7 @@ impl Codex { analytics_events_client, thread_store, attestation_provider, + external_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_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 2b61468c1..af2d1abeb 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_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 time_provider = crate::current_time::resolve_time_provider( + config.current_time_reminder.as_ref(), + external_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(), + 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 b13cd029c..0b1cc045c 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_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, + time_provider: Arc::new(crate::current_time::SystemTimeProvider), 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_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_time_provider*/ None, Some(config.multi_agent_version_from_features()), ) .await?; @@ -7078,6 +7082,7 @@ where state_db, )), attestation_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 0bed03235..5de331ee1 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_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 000000000..058132a89 --- /dev/null +++ b/codex-rs/core/src/session/time_reminder.rs @@ -0,0 +1,61 @@ +use codex_protocol::error::CodexErr; +use codex_protocol::error::Result as CodexResult; + +use super::session::Session; +use super::turn_context::TurnContext; +use crate::context::ContextualUserFragment; + +#[derive(Default)] +pub(crate) struct CurrentTimeReminderState { + model_requests_since_delivery: u64, + last_window_id: Option, +} + +impl CurrentTimeReminderState { + 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); + 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()); + } + + reminder_is_due + } +} + +pub(super) async fn maybe_record_current_time_reminder( + sess: &Session, + turn_context: &TurnContext, + window_id: &str, +) -> CodexResult<()> { + let Some(config) = turn_context.config.current_time_reminder else { + return Ok(()); + }; + + let reminder_is_due = { + let mut state = sess.state.lock().await; + state + .current_time_reminder + .take_reminder_due(window_id, config.reminder_interval_model_requests) + }; + if !reminder_is_due { + return Ok(()); + } + + let current_time = sess + .services + .time_provider + .current_time(sess.thread_id) + .await + .map_err(|err| CodexErr::Fatal(format!("failed to read current time: {err:#}")))?; + + let response_item = + ContextualUserFragment::into(crate::context::CurrentTimeReminder::new(current_time)); + sess.record_conversation_items(turn_context, std::slice::from_ref(&response_item)) + .await; + + Ok(()) +} diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index 657e00ff3..25e4e1a6a 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -226,34 +226,50 @@ pub(crate) async fn run_turn( ) .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) - } - .instrument(trace_span!("run_turn.prepare_sampling_request_input")) - .await; + let sampling_request_result: CodexResult<_> = async { + super::time_reminder::maybe_record_current_time_reminder( + sess.as_ref(), + turn_context.as_ref(), + &window_id, + ) + .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)) => { + // 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, + )) + } + .await; + 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, diff --git a/codex-rs/core/src/state/service.rs b/codex-rs/core/src/state/service.rs index 914fd4c32..853b96e92 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::TimeProvider; 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) 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/state/session.rs b/codex-rs/core/src/state/session.rs index 269d3e0f6..2bd9da5ce 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 77a552002..814c8fe76 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::TimeProvider; 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_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_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_time_provider, auth_manager, session_source, installation_id, @@ -405,6 +409,7 @@ impl ThreadManager { ), thread_store, attestation_provider: None, + external_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_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 ee010b5e1..bfb7865bb 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_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_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_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_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_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_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_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_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_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_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_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 ed15a7483..eefea1eb8 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_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 b10e26bf0..42b9bf87c 100644 --- a/codex-rs/core/tests/common/test_codex.rs +++ b/codex-rs/core/tests/common/test_codex.rs @@ -17,6 +17,7 @@ use codex_config::CloudConfigBundleLoader; use codex_core::CodexThread; 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; @@ -261,6 +262,7 @@ pub struct TestCodexBuilder { extensions: Arc>, user_instructions_provider: Option>, supports_openai_form_elicitation: bool, + external_time_provider: Option>, } impl TestCodexBuilder { @@ -362,6 +364,11 @@ impl TestCodexBuilder { self } + pub fn with_external_time_provider(mut self, provider: Arc) -> Self { + self.external_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_time_provider*/ self.external_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, + external_time_provider: None, } } diff --git a/codex-rs/core/tests/suite/client.rs b/codex-rs/core/tests/suite/client.rs index 54e24d49c..e32961882 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_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 new file mode 100644 index 000000000..5b5970a0e --- /dev/null +++ b/codex-rs/core/tests/suite/current_time_reminder.rs @@ -0,0 +1,274 @@ +use std::sync::Arc; +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::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; +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; +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 TestTimeProvider(AtomicI64); + +impl Default for TestTimeProvider { + fn default() -> Self { + Self(AtomicI64::new(FIRST_TIME_UNIX_SECONDS)) + } +} + +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) + .expect("test timestamp should be valid")) + }) + } +} + +struct FailingTimeProvider; + +impl TimeProvider for FailingTimeProvider { + fn current_time(&self, _thread_id: ThreadId) -> TimeFuture<'_> { + Box::pin(async { Err(anyhow!("test clock unavailable")) }) + } +} + +fn current_time_reminders(request: &ResponsesRequest) -> Vec { + request + .message_input_texts("developer") + .into_iter() + .filter(|text| text.starts_with("It is ")) + .collect() +} + +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, + }); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn current_time_reminders_follow_request_interval_and_persist_in_history() -> Result<()> { + 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_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, CurrentTimeSource::External) + }) + .with_external_time_provider(Arc::new(TestTimeProvider::default())) + .build(&server) + .await?; + + test.submit_turn_with_permission_profile("first turn", PermissionProfile::Disabled) + .await?; + test.submit_turn("second 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] + ); + 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(())); + + 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 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_current_time_reminder(config, /*interval*/ 50, CurrentTimeSource::External); + }) + .with_external_time_provider(Arc::new(TestTimeProvider::default())) + .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(()) +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn 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, CurrentTimeSource::External) + }) + .with_external_time_provider(Arc::new(FailingTimeProvider)) + .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(()) +} diff --git a/codex-rs/core/tests/suite/mod.rs b/codex-rs/core/tests/suite/mod.rs index 0fdd958b3..1e200ab48 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; diff --git a/codex-rs/mcp-server/src/message_processor.rs b/codex-rs/mcp-server/src/message_processor.rs index ce19575a9..5f36091fe 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_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 1534ba86a..26ba71ae4 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_time_provider*/ None, ); let NewThread {