From 117ce787837a34fc4e5208da24916e0d1a8289ec Mon Sep 17 00:00:00 2001 From: Eric Traut Date: Mon, 6 Apr 2026 18:20:48 -0700 Subject: [PATCH] Fill rollout DB pages after filtering --- codex-rs/rollout/src/recorder_tests.rs | 80 ++++++++++++++++ codex-rs/rollout/src/state_db.rs | 121 +++++++++++++++---------- 2 files changed, 152 insertions(+), 49 deletions(-) diff --git a/codex-rs/rollout/src/recorder_tests.rs b/codex-rs/rollout/src/recorder_tests.rs index 9efd255393..9a597257d2 100644 --- a/codex-rs/rollout/src/recorder_tests.rs +++ b/codex-rs/rollout/src/recorder_tests.rs @@ -515,6 +515,86 @@ async fn list_threads_includes_alarm_only_sessions_without_user_messages() -> st Ok(()) } +#[tokio::test] +async fn list_threads_db_enabled_fills_page_after_filtering_empty_threads() -> std::io::Result<()> { + let home = TempDir::new().expect("temp dir"); + let config = test_config(home.path()); + + let empty_uuid = Uuid::from_u128(9015); + let valid_uuid = Uuid::from_u128(9016); + let empty_thread_id = + ThreadId::from_string(&empty_uuid.to_string()).expect("valid empty thread id"); + let valid_thread_id = ThreadId::from_string(&valid_uuid.to_string()).expect("valid thread id"); + let empty_path = + write_session_file_without_user_message(home.path(), "2025-01-03T15-00-00", empty_uuid)?; + let valid_path = write_session_file(home.path(), "2025-01-03T14-00-00", valid_uuid)?; + + let runtime = codex_state::StateRuntime::init( + home.path().to_path_buf(), + config.model_provider_id.clone(), + ) + .await + .expect("state db should initialize"); + runtime + .mark_backfill_complete(/*last_watermark*/ None) + .await + .expect("backfill should be complete"); + + let empty_created_at = chrono::Utc + .with_ymd_and_hms(2025, 1, 3, 15, 0, 0) + .single() + .expect("valid empty thread datetime"); + let mut empty_builder = codex_state::ThreadMetadataBuilder::new( + empty_thread_id, + empty_path, + empty_created_at, + SessionSource::Cli, + ); + empty_builder.model_provider = Some(config.model_provider_id.clone()); + empty_builder.cwd = home.path().to_path_buf(); + let empty_metadata = empty_builder.build(config.model_provider_id.as_str()); + runtime + .upsert_thread(&empty_metadata) + .await + .expect("state db empty upsert should succeed"); + + let valid_created_at = chrono::Utc + .with_ymd_and_hms(2025, 1, 3, 14, 0, 0) + .single() + .expect("valid thread datetime"); + let mut valid_builder = codex_state::ThreadMetadataBuilder::new( + valid_thread_id, + valid_path.clone(), + valid_created_at, + SessionSource::Cli, + ); + valid_builder.model_provider = Some(config.model_provider_id.clone()); + valid_builder.cwd = home.path().to_path_buf(); + let mut valid_metadata = valid_builder.build(config.model_provider_id.as_str()); + valid_metadata.first_user_message = Some("Hello from user".to_string()); + runtime + .upsert_thread(&valid_metadata) + .await + .expect("state db valid upsert should succeed"); + + let default_provider = config.model_provider_id.clone(); + let page = RolloutRecorder::list_threads( + &config, + /*page_size*/ 1, + /*cursor*/ None, + ThreadSortKey::CreatedAt, + &[], + /*model_providers*/ None, + default_provider.as_str(), + /*search_term*/ None, + ) + .await?; + assert_eq!(page.items.len(), 1); + assert_eq!(page.items[0].path, valid_path); + assert!(page.next_cursor.is_none()); + Ok(()) +} + #[tokio::test] async fn resume_candidate_matches_cwd_reads_latest_turn_context() -> std::io::Result<()> { let home = TempDir::new().expect("temp dir"); diff --git a/codex-rs/rollout/src/state_db.rs b/codex-rs/rollout/src/state_db.rs index 8f4b169088..f95525c928 100644 --- a/codex-rs/rollout/src/state_db.rs +++ b/codex-rs/rollout/src/state_db.rs @@ -227,58 +227,81 @@ pub async fn list_threads_db( }) .collect(); let model_providers = model_providers.map(<[String]>::to_vec); - match ctx - .list_threads( - page_size, - anchor.as_ref(), - match sort_key { - ThreadSortKey::CreatedAt => codex_state::SortKey::CreatedAt, - ThreadSortKey::UpdatedAt => codex_state::SortKey::UpdatedAt, - }, - allowed_sources.as_slice(), - model_providers.as_deref(), - archived, - search_term, - ) - .await - { - Ok(mut page) => { - let mut valid_items = Vec::with_capacity(page.items.len()); - for mut item in page.items { - if tokio::fs::try_exists(&item.rollout_path) - .await - .unwrap_or(false) - { - let missing_preview = - item.first_user_message.as_deref().is_none_or(str::is_empty); - if missing_preview { - if let Some(alarm_preview) = - thread_preview_from_alarm_sidecar(&item.rollout_path).await - { - item.first_user_message = Some(alarm_preview); - } else { - continue; - } - } - valid_items.push(item); - } else { - warn!( - "state db list_threads returned stale rollout path for thread {}: {}", - item.id, - item.rollout_path.display() - ); - warn!("state db discrepancy during list_threads_db: stale_db_path_dropped"); - let _ = ctx.delete_thread(item.id).await; - } - } - page.items = valid_items; - Some(page) + let db_sort_key = match sort_key { + ThreadSortKey::CreatedAt => codex_state::SortKey::CreatedAt, + ThreadSortKey::UpdatedAt => codex_state::SortKey::UpdatedAt, + }; + let mut db_anchor = anchor; + let mut valid_items = Vec::with_capacity(page_size); + let mut next_anchor = None; + let mut num_scanned_rows = 0; + loop { + let remaining = page_size.saturating_sub(valid_items.len()); + if remaining == 0 { + break; } - Err(err) => { - warn!("state db list_threads failed: {err}"); - None + let mut page = match ctx + .list_threads( + remaining, + db_anchor.as_ref(), + db_sort_key, + allowed_sources.as_slice(), + model_providers.as_deref(), + archived, + search_term, + ) + .await + { + Ok(page) => page, + Err(err) => { + warn!("state db list_threads failed: {err}"); + return None; + } + }; + num_scanned_rows += page.num_scanned_rows; + let page_next_anchor = page.next_anchor.take(); + for mut item in page.items { + if tokio::fs::try_exists(&item.rollout_path) + .await + .unwrap_or(false) + { + let missing_preview = item.first_user_message.as_deref().is_none_or(str::is_empty); + if missing_preview { + if let Some(alarm_preview) = + thread_preview_from_alarm_sidecar(&item.rollout_path).await + { + item.first_user_message = Some(alarm_preview); + } else { + continue; + } + } + valid_items.push(item); + } else { + warn!( + "state db list_threads returned stale rollout path for thread {}: {}", + item.id, + item.rollout_path.display() + ); + warn!("state db discrepancy during list_threads_db: stale_db_path_dropped"); + let _ = ctx.delete_thread(item.id).await; + } + } + + if valid_items.len() == page_size { + next_anchor = page_next_anchor; + break; + } + + match page_next_anchor { + Some(anchor) => db_anchor = Some(anchor), + None => break, } } + Some(codex_state::ThreadsPage { + items: valid_items, + next_anchor, + num_scanned_rows, + }) } /// Look up the rollout path for a thread id using SQLite.