mirror of
https://github.com/openai/codex.git
synced 2026-09-10 20:26:47 +00:00
Fill rollout DB pages after filtering
This commit is contained in:
@@ -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");
|
||||
|
||||
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user