From 6fe44104f3c450361abb721a3586cc177ae95bba Mon Sep 17 00:00:00 2001 From: Eric Traut Date: Fri, 10 Apr 2026 20:59:16 -0700 Subject: [PATCH] codex: restore timers after resume history replay (#17380) --- codex-rs/core/src/codex.rs | 3 +- codex-rs/core/tests/suite/timers.rs | 107 ++++++++++++++++++++++++++++ 2 files changed, 108 insertions(+), 2 deletions(-) diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index 682d96fef9..dbc76e2548 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -2137,8 +2137,6 @@ impl Session { for event in events { sess.send_event_raw(event).await; } - sess.restore_timers_from_db().await; - // Start the watcher after SessionConfigured so it cannot emit earlier events. sess.start_skills_watcher_listener(); // Construct sandbox_state before MCP startup so it can be sent to each @@ -2233,6 +2231,7 @@ impl Session { let mut state = sess.state.lock().await; state.set_pending_session_start_source(Some(session_start_source)); } + sess.restore_timers_from_db().await; memories::start_memories_startup_task( &sess, diff --git a/codex-rs/core/tests/suite/timers.rs b/codex-rs/core/tests/suite/timers.rs index 2ec532e108..e899fe7ae3 100644 --- a/codex-rs/core/tests/suite/timers.rs +++ b/codex-rs/core/tests/suite/timers.rs @@ -8,6 +8,7 @@ use codex_core::timers::ThreadTimerTrigger; use codex_core::timers::TimerDelivery; use codex_features::Feature; use codex_protocol::protocol::EventMsg; +use codex_protocol::protocol::Op; use core_test_support::responses::ev_assistant_message; use core_test_support::responses::ev_completed; use core_test_support::responses::ev_response_created; @@ -184,6 +185,112 @@ async fn create_timer_rejects_ephemeral_thread() -> Result<()> { Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn resume_due_timer_runs_after_history_reconstruction() -> Result<()> { + let server = start_mock_server().await; + let mock = mount_sse_sequence( + &server, + vec![ + sse(vec![ + ev_response_created("resp-1"), + ev_assistant_message("msg-1", "recorded before resume"), + ev_completed("resp-1"), + ]), + sse(vec![ + ev_response_created("resp-2"), + ev_assistant_message("msg-2", "timer after resume"), + ev_completed("resp-2"), + ]), + ], + ) + .await; + + let mut builder = test_codex().with_config(|config| { + config + .features + .enable(Feature::Timers) + .unwrap_or_else(|err| panic!("test config should allow feature update: {err}")); + config + .features + .enable(Feature::Sqlite) + .unwrap_or_else(|err| panic!("test config should allow feature update: {err}")); + }); + let initial = builder.build(&server).await?; + initial.submit_turn("context before resume").await?; + + let db = initial.codex.state_db().expect("state db enabled"); + let home = initial.home.clone(); + let rollout_path = initial + .codex + .rollout_path() + .expect("rollout path should be materialized"); + let thread_id = initial.session_configured.session_id.to_string(); + + initial.codex.submit(Op::Shutdown).await?; + wait_for_event(&initial.codex, |event| { + matches!(event, EventMsg::ShutdownComplete) + }) + .await; + + let now = Utc::now().timestamp(); + db.create_thread_timer(&codex_state::ThreadTimerCreateParams { + id: "resume-due-timer".to_string(), + thread_id, + source: "external".to_string(), + client_id: "codex-test".to_string(), + trigger_json: r#"{"kind":"delay","seconds":0,"repeat":false}"#.to_string(), + content: "resume timer".to_string(), + instructions: None, + meta_json: "{}".to_string(), + delivery: TimerDelivery::AfterTurn.as_str().to_string(), + created_at: now - 10, + next_run_at: Some(now - 1), + last_run_at: None, + pending_run: false, + }) + .await?; + + let mut resume_builder = test_codex().with_config(|config| { + config + .features + .enable(Feature::Timers) + .unwrap_or_else(|err| panic!("test config should allow feature update: {err}")); + config + .features + .enable(Feature::Sqlite) + .unwrap_or_else(|err| panic!("test config should allow feature update: {err}")); + }); + let resumed = resume_builder.resume(&server, home, rollout_path).await?; + + wait_for_event_with_timeout( + &resumed.codex, + |event| match event { + EventMsg::InjectedMessage(event) => { + event.source == "timer resume-due-timer" + && event.content == "Timer fired: resume timer" + } + _ => false, + }, + Duration::from_secs(20), + ) + .await; + wait_for_event_with_timeout( + &resumed.codex, + |event| matches!(event, EventMsg::TurnComplete(_)), + Duration::from_secs(20), + ) + .await; + + let requests = mock.requests(); + assert_eq!(requests.len(), 2); + let resumed_request = &requests[1]; + assert!(resumed_request.body_contains_text("context before resume")); + assert!(resumed_request.body_contains_text("recorded before resume")); + assert!(resumed_request.body_contains_text("Timer fired: resume timer")); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn list_timers_discovers_externally_inserted_timer() -> Result<()> { let server = start_mock_server().await;