//! Persist Codex session rollouts (.jsonl) so sessions can be replayed or inspected later. use std::fs; use std::fs::File; use std::io::Error as IoError; use std::path::Path; use std::path::PathBuf; use chrono::SecondsFormat; use chrono::Utc; use codex_protocol::ThreadId; use codex_protocol::dynamic_tools::DynamicToolSpec; use codex_protocol::models::BaseInstructions; use codex_utils_string::truncate_middle_chars; use serde_json::Value; use time::OffsetDateTime; use time::format_description::FormatItem; use time::format_description::well_known::Rfc3339; use time::macros::format_description; use tokio::io::AsyncWriteExt; use tokio::sync::mpsc; use tokio::sync::mpsc::Sender; use tokio::sync::oneshot; use tracing::info; use tracing::trace; use tracing::warn; use super::ARCHIVED_SESSIONS_SUBDIR; use super::SESSIONS_SUBDIR; use super::list::Cursor; use super::list::ThreadItem; use super::list::ThreadListConfig; use super::list::ThreadListLayout; use super::list::ThreadSortKey; use super::list::ThreadsPage; use super::list::get_threads; use super::list::get_threads_in_root; use super::list::parse_cursor; use super::list::parse_timestamp_uuid_from_filename; use super::metadata; use super::policy::EventPersistenceMode; use super::policy::is_persisted_response_item; use crate::config::RolloutConfigView; use crate::default_client::originator; use crate::state_db; use crate::state_db::StateDbHandle; use codex_git_utils::collect_git_info; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::GitInfo as ProtocolGitInfo; use codex_protocol::protocol::InitialHistory; use codex_protocol::protocol::ResumedHistory; use codex_protocol::protocol::RolloutItem; use codex_protocol::protocol::RolloutLine; use codex_protocol::protocol::SessionMeta; use codex_protocol::protocol::SessionMetaLine; use codex_protocol::protocol::SessionSource; use codex_state::StateRuntime; use codex_state::ThreadMetadataBuilder; use codex_utils_path as path_utils; /// Records all [`ResponseItem`]s for a session and flushes them to disk after /// every update. /// /// Rollouts are recorded as JSONL and can be inspected with tools such as: /// /// ```ignore /// $ jq -C . ~/.codex/sessions/rollout-2025-05-07T17-24-21-5973b6c0-94b8-487b-a530-2aeb6098ae0e.jsonl /// $ fx ~/.codex/sessions/rollout-2025-05-07T17-24-21-5973b6c0-94b8-487b-a530-2aeb6098ae0e.jsonl /// ``` #[derive(Clone)] pub struct RolloutRecorder { tx: Sender, pub(crate) rollout_path: PathBuf, state_db: Option, event_persistence_mode: EventPersistenceMode, } #[derive(Clone)] pub enum RolloutRecorderParams { Create { conversation_id: ThreadId, forked_from_id: Option, source: SessionSource, base_instructions: BaseInstructions, dynamic_tools: Vec, event_persistence_mode: EventPersistenceMode, }, Resume { path: PathBuf, event_persistence_mode: EventPersistenceMode, }, } enum RolloutCmd { AddItems(Vec), Persist { ack: oneshot::Sender<()>, }, /// Ensure all prior writes are processed; respond when flushed. Flush { ack: oneshot::Sender<()>, }, Shutdown { ack: oneshot::Sender<()>, }, } impl RolloutRecorderParams { pub fn new( conversation_id: ThreadId, forked_from_id: Option, source: SessionSource, base_instructions: BaseInstructions, dynamic_tools: Vec, event_persistence_mode: EventPersistenceMode, ) -> Self { Self::Create { conversation_id, forked_from_id, source, base_instructions, dynamic_tools, event_persistence_mode, } } pub fn resume(path: PathBuf, event_persistence_mode: EventPersistenceMode) -> Self { Self::Resume { path, event_persistence_mode, } } } const PERSISTED_EXEC_AGGREGATED_OUTPUT_MAX_BYTES: usize = 10_000; fn sanitize_rollout_item_for_persistence( item: RolloutItem, mode: EventPersistenceMode, ) -> RolloutItem { if mode != EventPersistenceMode::Extended { return item; } match item { RolloutItem::EventMsg(EventMsg::ExecCommandEnd(mut event)) => { // Persist only a bounded aggregated summary of command output. event.aggregated_output = truncate_middle_chars( &event.aggregated_output, PERSISTED_EXEC_AGGREGATED_OUTPUT_MAX_BYTES, ); // Drop unnecessary fields from rollout storage since aggregated_output is all we need. event.stdout.clear(); event.stderr.clear(); event.formatted_output.clear(); RolloutItem::EventMsg(EventMsg::ExecCommandEnd(event)) } _ => item, } } impl RolloutRecorder { /// List threads (rollout files) under the provided Codex home directory. #[allow(clippy::too_many_arguments)] pub async fn list_threads( config: &impl RolloutConfigView, page_size: usize, cursor: Option<&Cursor>, sort_key: ThreadSortKey, allowed_sources: &[SessionSource], model_providers: Option<&[String]>, default_provider: &str, search_term: Option<&str>, ) -> std::io::Result { Self::list_threads_with_db_fallback( config, page_size, cursor, sort_key, allowed_sources, model_providers, default_provider, /*archived*/ false, search_term, ) .await } /// List archived threads (rollout files) under the archived sessions directory. #[allow(clippy::too_many_arguments)] pub async fn list_archived_threads( config: &impl RolloutConfigView, page_size: usize, cursor: Option<&Cursor>, sort_key: ThreadSortKey, allowed_sources: &[SessionSource], model_providers: Option<&[String]>, default_provider: &str, search_term: Option<&str>, ) -> std::io::Result { Self::list_threads_with_db_fallback( config, page_size, cursor, sort_key, allowed_sources, model_providers, default_provider, /*archived*/ true, search_term, ) .await } #[allow(clippy::too_many_arguments)] async fn list_threads_with_db_fallback( config: &impl RolloutConfigView, page_size: usize, cursor: Option<&Cursor>, sort_key: ThreadSortKey, allowed_sources: &[SessionSource], model_providers: Option<&[String]>, default_provider: &str, archived: bool, search_term: Option<&str>, ) -> std::io::Result { let codex_home = config.codex_home(); // Filesystem-first listing intentionally overfetches so we can repair stale/missing // SQLite rollout paths before the final DB-backed page is returned. let fs_page_size = page_size.saturating_mul(2).max(page_size); let fs_page = if archived { let root = codex_home.join(ARCHIVED_SESSIONS_SUBDIR); get_threads_in_root( root, fs_page_size, cursor, sort_key, ThreadListConfig { allowed_sources, model_providers, default_provider, layout: ThreadListLayout::Flat, }, ) .await? } else { get_threads( codex_home, fs_page_size, cursor, sort_key, allowed_sources, model_providers, default_provider, ) .await? }; let state_db_ctx = state_db::get_state_db(config).await; if state_db_ctx.is_none() { // Keep legacy behavior when SQLite is unavailable: return filesystem results // at the requested page size. return Ok(truncate_fs_page(fs_page, page_size, sort_key)); } // Warm the DB by repairing every filesystem hit before querying SQLite. for item in &fs_page.items { state_db::read_repair_rollout_path( state_db_ctx.as_deref(), item.thread_id, Some(archived), item.path.as_path(), ) .await; } if let Some(db_page) = state_db::list_threads_db( state_db_ctx.as_deref(), codex_home, page_size, cursor, sort_key, allowed_sources, model_providers, archived, search_term, ) .await { return Ok(db_page.into()); } // If SQLite listing still fails, return the filesystem page rather than failing the list. tracing::error!("Falling back on rollout system"); tracing::warn!("state db discrepancy during list_threads_with_db_fallback: falling_back"); Ok(truncate_fs_page(fs_page, page_size, sort_key)) } /// Find the newest recorded thread path, optionally filtering to a matching cwd. #[allow(clippy::too_many_arguments)] pub async fn find_latest_thread_path( config: &impl RolloutConfigView, page_size: usize, cursor: Option<&Cursor>, sort_key: ThreadSortKey, allowed_sources: &[SessionSource], model_providers: Option<&[String]>, default_provider: &str, filter_cwd: Option<&Path>, ) -> std::io::Result> { let codex_home = config.codex_home(); let state_db_ctx = state_db::get_state_db(config).await; if state_db_ctx.is_some() { let mut db_cursor = cursor.cloned(); loop { let Some(db_page) = state_db::list_threads_db( state_db_ctx.as_deref(), codex_home, page_size, db_cursor.as_ref(), sort_key, allowed_sources, model_providers, /*archived*/ false, /*search_term*/ None, ) .await else { break; }; if let Some(path) = select_resume_path_from_db_page(&db_page, filter_cwd, default_provider).await { return Ok(Some(path)); } db_cursor = db_page.next_anchor.map(Into::into); if db_cursor.is_none() { break; } } } let mut cursor = cursor.cloned(); loop { let page = get_threads( codex_home, page_size, cursor.as_ref(), sort_key, allowed_sources, model_providers, default_provider, ) .await?; if let Some(path) = select_resume_path(&page, filter_cwd, default_provider).await { return Ok(Some(path)); } cursor = page.next_cursor; if cursor.is_none() { return Ok(None); } } } /// Attempt to create a new [`RolloutRecorder`]. /// /// For newly created sessions, this precomputes path/metadata and defers /// file creation/open until an explicit `persist()` call. /// /// For resumed sessions, this immediately opens the existing rollout file. pub async fn new( config: &impl RolloutConfigView, params: RolloutRecorderParams, state_db_ctx: Option, state_builder: Option, ) -> std::io::Result { let (file, deferred_log_file_info, rollout_path, meta, event_persistence_mode) = match params { RolloutRecorderParams::Create { conversation_id, forked_from_id, source, base_instructions, dynamic_tools, event_persistence_mode, } => { let log_file_info = precompute_log_file_info(config, conversation_id)?; let path = log_file_info.path.clone(); let session_id = log_file_info.conversation_id; let started_at = log_file_info.timestamp; let timestamp_format: &[FormatItem] = format_description!( "[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]Z" ); let timestamp = started_at .to_offset(time::UtcOffset::UTC) .format(timestamp_format) .map_err(|e| IoError::other(format!("failed to format timestamp: {e}")))?; let session_meta = SessionMeta { id: session_id, forked_from_id, timestamp, cwd: config.cwd().to_path_buf(), originator: originator().value, cli_version: env!("CARGO_PKG_VERSION").to_string(), agent_nickname: source.get_nickname(), agent_role: source.get_agent_role(), agent_path: source.get_agent_path().map(Into::into), source, model_provider: Some(config.model_provider_id().to_string()), base_instructions: Some(base_instructions), dynamic_tools: if dynamic_tools.is_empty() { None } else { Some(dynamic_tools) }, memory_mode: (!config.generate_memories()) .then_some("disabled".to_string()), }; ( None, Some(log_file_info), path, Some(session_meta), event_persistence_mode, ) } RolloutRecorderParams::Resume { path, event_persistence_mode, } => ( Some( tokio::fs::OpenOptions::new() .append(true) .open(&path) .await?, ), None, path, None, event_persistence_mode, ), }; // Clone the cwd for the spawned task to collect git info asynchronously let cwd = config.cwd().to_path_buf(); // A reasonably-sized bounded channel. If the buffer fills up the send // future will yield, which is fine – we only need to ensure we do not // perform *blocking* I/O on the caller's thread. let (tx, rx) = mpsc::channel::(256); // Spawn a Tokio task that owns the file handle and performs async // writes. Using `tokio::fs::File` keeps everything on the async I/O // driver instead of blocking the runtime. tokio::task::spawn(rollout_writer( file, deferred_log_file_info, rx, meta, cwd, rollout_path.clone(), state_db_ctx.clone(), state_builder, config.model_provider_id().to_string(), config.generate_memories(), )); Ok(Self { tx, rollout_path, state_db: state_db_ctx, event_persistence_mode, }) } pub fn rollout_path(&self) -> &Path { self.rollout_path.as_path() } pub fn state_db(&self) -> Option { self.state_db.clone() } pub async fn record_items(&self, items: &[RolloutItem]) -> std::io::Result<()> { let mut filtered = Vec::new(); for item in items { // Note that function calls may look a bit strange if they are // "fully qualified MCP tool calls," so we could consider // reformatting them in that case. if is_persisted_response_item(item, self.event_persistence_mode) { filtered.push(sanitize_rollout_item_for_persistence( item.clone(), self.event_persistence_mode, )); } } if filtered.is_empty() { return Ok(()); } self.tx .send(RolloutCmd::AddItems(filtered)) .await .map_err(|e| IoError::other(format!("failed to queue rollout items: {e}"))) } /// Materialize the rollout file and persist all buffered items. /// /// This is idempotent; after first materialization, repeated calls are no-ops. pub async fn persist(&self) -> std::io::Result<()> { let (tx, rx) = oneshot::channel(); self.tx .send(RolloutCmd::Persist { ack: tx }) .await .map_err(|e| IoError::other(format!("failed to queue rollout persist: {e}")))?; rx.await .map_err(|e| IoError::other(format!("failed waiting for rollout persist: {e}"))) } /// Flush all queued writes and wait until they are committed by the writer task. pub async fn flush(&self) -> std::io::Result<()> { let (tx, rx) = oneshot::channel(); self.tx .send(RolloutCmd::Flush { ack: tx }) .await .map_err(|e| IoError::other(format!("failed to queue rollout flush: {e}")))?; rx.await .map_err(|e| IoError::other(format!("failed waiting for rollout flush: {e}"))) } pub async fn load_rollout_items( path: &Path, ) -> std::io::Result<(Vec, Option, usize)> { trace!("Resuming rollout from {path:?}"); let text = tokio::fs::read_to_string(path).await?; if text.trim().is_empty() { return Err(IoError::other("empty session file")); } let mut items: Vec = Vec::new(); let mut thread_id: Option = None; let mut parse_errors = 0usize; for line in text.lines() { if line.trim().is_empty() { continue; } let v: Value = match serde_json::from_str(line) { Ok(v) => v, Err(e) => { warn!("failed to parse line as JSON: {line:?}, error: {e}"); parse_errors = parse_errors.saturating_add(1); continue; } }; // Parse the rollout line structure match serde_json::from_value::(v.clone()) { Ok(rollout_line) => match rollout_line.item { RolloutItem::SessionMeta(session_meta_line) => { // Use the FIRST SessionMeta encountered in the file as the canonical // thread id and main session information. Keep all items intact. if thread_id.is_none() { thread_id = Some(session_meta_line.meta.id); } items.push(RolloutItem::SessionMeta(session_meta_line)); } RolloutItem::ResponseItem(item) => { items.push(RolloutItem::ResponseItem(item)); } RolloutItem::Compacted(item) => { items.push(RolloutItem::Compacted(item)); } RolloutItem::TurnContext(item) => { items.push(RolloutItem::TurnContext(item)); } RolloutItem::EventMsg(_ev) => { items.push(RolloutItem::EventMsg(_ev)); } }, Err(e) => { trace!("failed to parse rollout line: {e}"); parse_errors = parse_errors.saturating_add(1); } } } tracing::debug!( "Resumed rollout with {} items, thread ID: {:?}, parse errors: {}", items.len(), thread_id, parse_errors, ); Ok((items, thread_id, parse_errors)) } pub async fn get_rollout_history(path: &Path) -> std::io::Result { let (items, thread_id, _parse_errors) = Self::load_rollout_items(path).await?; let conversation_id = thread_id .ok_or_else(|| IoError::other("failed to parse thread ID from rollout file"))?; if items.is_empty() { return Ok(InitialHistory::New); } info!("Resumed rollout successfully from {path:?}"); Ok(InitialHistory::Resumed(ResumedHistory { conversation_id, history: items, rollout_path: path.to_path_buf(), })) } pub async fn shutdown(&self) -> std::io::Result<()> { let (tx_done, rx_done) = oneshot::channel(); match self.tx.send(RolloutCmd::Shutdown { ack: tx_done }).await { Ok(_) => rx_done .await .map_err(|e| IoError::other(format!("failed waiting for rollout shutdown: {e}")))?, Err(e) => { warn!("failed to send rollout shutdown command: {e}"); return Err(IoError::other(format!( "failed to send rollout shutdown command: {e}" ))); } }; Ok(()) } } fn truncate_fs_page( mut page: ThreadsPage, page_size: usize, sort_key: ThreadSortKey, ) -> ThreadsPage { if page.items.len() <= page_size { return page; } page.items.truncate(page_size); page.next_cursor = page.items.last().and_then(|item| { let file_name = item.path.file_name()?.to_str()?; let (created_at, id) = parse_timestamp_uuid_from_filename(file_name)?; let cursor_token = match sort_key { ThreadSortKey::CreatedAt => format!("{}|{id}", created_at.format(&Rfc3339).ok()?), ThreadSortKey::UpdatedAt => format!("{}|{id}", item.updated_at.as_deref()?), }; parse_cursor(cursor_token.as_str()) }); page } struct LogFileInfo { /// Full path to the rollout file. path: PathBuf, /// Session ID (also embedded in filename). conversation_id: ThreadId, /// Timestamp for the start of the session. timestamp: OffsetDateTime, } fn precompute_log_file_info( config: &impl RolloutConfigView, conversation_id: ThreadId, ) -> std::io::Result { // Resolve ~/.codex/sessions/YYYY/MM/DD path. let timestamp = OffsetDateTime::now_local() .map_err(|e| IoError::other(format!("failed to get local time: {e}")))?; let mut dir = config.codex_home().to_path_buf(); dir.push(SESSIONS_SUBDIR); dir.push(timestamp.year().to_string()); dir.push(format!("{:02}", u8::from(timestamp.month()))); dir.push(format!("{:02}", timestamp.day())); // Custom format for YYYY-MM-DDThh-mm-ss. Use `-` instead of `:` for // compatibility with filesystems that do not allow colons in filenames. let format: &[FormatItem] = format_description!("[year]-[month]-[day]T[hour]-[minute]-[second]"); let date_str = timestamp .format(format) .map_err(|e| IoError::other(format!("failed to format timestamp: {e}")))?; let filename = format!("rollout-{date_str}-{conversation_id}.jsonl"); let path = dir.join(filename); Ok(LogFileInfo { path, conversation_id, timestamp, }) } fn open_log_file(path: &Path) -> std::io::Result { let Some(parent) = path.parent() else { return Err(IoError::other(format!( "rollout path has no parent: {}", path.display() ))); }; fs::create_dir_all(parent)?; std::fs::OpenOptions::new() .append(true) .create(true) .open(path) } #[allow(clippy::too_many_arguments)] async fn rollout_writer( file: Option, mut deferred_log_file_info: Option, mut rx: mpsc::Receiver, mut meta: Option, cwd: std::path::PathBuf, rollout_path: PathBuf, state_db_ctx: Option, mut state_builder: Option, default_provider: String, generate_memories: bool, ) -> std::io::Result<()> { let mut writer = file.map(|file| JsonlWriter { file }); let mut buffered_items = Vec::::new(); if let Some(builder) = state_builder.as_mut() { builder.rollout_path = rollout_path.clone(); } // Resumed sessions already have a file handle open, so session metadata can // be written immediately if present. if writer.is_some() && let Some(session_meta) = meta.take() { write_session_meta( writer.as_mut(), session_meta, &cwd, &rollout_path, state_db_ctx.as_deref(), &mut state_builder, default_provider.as_str(), generate_memories, ) .await?; } // Process rollout commands while let Some(cmd) = rx.recv().await { match cmd { RolloutCmd::AddItems(items) => { if items.is_empty() { continue; } if writer.is_none() { buffered_items.extend(items); continue; } write_and_reconcile_items( writer.as_mut(), items.as_slice(), &rollout_path, state_db_ctx.as_deref(), state_builder.as_ref(), default_provider.as_str(), ) .await?; } RolloutCmd::Persist { ack } => { if writer.is_none() { let result = async { let Some(log_file_info) = deferred_log_file_info.take() else { return Err(IoError::other( "deferred rollout recorder missing log file metadata", )); }; let file = open_log_file(log_file_info.path.as_path())?; writer = Some(JsonlWriter { file: tokio::fs::File::from_std(file), }); if let Some(session_meta) = meta.take() { write_session_meta( writer.as_mut(), session_meta, &cwd, &rollout_path, state_db_ctx.as_deref(), &mut state_builder, default_provider.as_str(), generate_memories, ) .await?; } if !buffered_items.is_empty() { write_and_reconcile_items( writer.as_mut(), buffered_items.as_slice(), &rollout_path, state_db_ctx.as_deref(), state_builder.as_ref(), default_provider.as_str(), ) .await?; buffered_items.clear(); } Ok(()) } .await; if let Err(err) = result { let _ = ack.send(()); return Err(err); } } let _ = ack.send(()); } RolloutCmd::Flush { ack } => { // Deferred fresh threads may not have an initialized file yet. if let Some(writer) = writer.as_mut() && let Err(e) = writer.file.flush().await { let _ = ack.send(()); return Err(e); } let _ = ack.send(()); } RolloutCmd::Shutdown { ack } => { let _ = ack.send(()); } } } Ok(()) } #[allow(clippy::too_many_arguments)] async fn write_session_meta( mut writer: Option<&mut JsonlWriter>, session_meta: SessionMeta, cwd: &Path, rollout_path: &Path, state_db_ctx: Option<&StateRuntime>, state_builder: &mut Option, default_provider: &str, generate_memories: bool, ) -> std::io::Result<()> { let git_info = collect_git_info(cwd).await.map(|info| ProtocolGitInfo { commit_hash: info.commit_hash, branch: info.branch, repository_url: info.repository_url, }); let session_meta_line = SessionMetaLine { meta: session_meta, git: git_info, }; if state_db_ctx.is_some() { *state_builder = metadata::builder_from_session_meta(&session_meta_line, rollout_path); } let rollout_item = RolloutItem::SessionMeta(session_meta_line); if let Some(writer) = writer.as_mut() { writer.write_rollout_item(&rollout_item).await?; } sync_thread_state_after_write( state_db_ctx, rollout_path, state_builder.as_ref(), std::slice::from_ref(&rollout_item), default_provider, (!generate_memories).then_some("disabled"), ) .await; Ok(()) } async fn write_and_reconcile_items( mut writer: Option<&mut JsonlWriter>, items: &[RolloutItem], rollout_path: &Path, state_db_ctx: Option<&StateRuntime>, state_builder: Option<&ThreadMetadataBuilder>, default_provider: &str, ) -> std::io::Result<()> { if let Some(writer) = writer.as_mut() { for item in items { writer.write_rollout_item(item).await?; } } sync_thread_state_after_write( state_db_ctx, rollout_path, state_builder, items, default_provider, /*new_thread_memory_mode*/ None, ) .await; Ok(()) } async fn sync_thread_state_after_write( state_db_ctx: Option<&StateRuntime>, rollout_path: &Path, state_builder: Option<&ThreadMetadataBuilder>, items: &[RolloutItem], default_provider: &str, new_thread_memory_mode: Option<&str>, ) { let updated_at = Utc::now(); if new_thread_memory_mode.is_some() || items .iter() .any(codex_state::rollout_item_affects_thread_metadata) { state_db::apply_rollout_items( state_db_ctx, rollout_path, default_provider, state_builder, items, "rollout_writer", new_thread_memory_mode, Some(updated_at), ) .await; return; } let thread_id = state_builder .map(|builder| builder.id) .or_else(|| metadata::builder_from_items(items, rollout_path).map(|builder| builder.id)); if state_db::touch_thread_updated_at(state_db_ctx, thread_id, updated_at, "rollout_writer") .await { return; } state_db::apply_rollout_items( state_db_ctx, rollout_path, default_provider, state_builder, items, "rollout_writer", new_thread_memory_mode, Some(updated_at), ) .await; } struct JsonlWriter { file: tokio::fs::File, } #[derive(serde::Serialize)] struct RolloutLineRef<'a> { timestamp: String, #[serde(flatten)] item: &'a RolloutItem, } impl JsonlWriter { async fn write_rollout_item(&mut self, rollout_item: &RolloutItem) -> std::io::Result<()> { let timestamp_format: &[FormatItem] = format_description!( "[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]Z" ); let timestamp = OffsetDateTime::now_utc() .format(timestamp_format) .map_err(|e| IoError::other(format!("failed to format timestamp: {e}")))?; let line = RolloutLineRef { timestamp, item: rollout_item, }; self.write_line(&line).await } async fn write_line(&mut self, item: &impl serde::Serialize) -> std::io::Result<()> { let mut json = serde_json::to_string(item)?; json.push('\n'); self.file.write_all(json.as_bytes()).await?; self.file.flush().await?; Ok(()) } } impl From for ThreadsPage { fn from(db_page: codex_state::ThreadsPage) -> Self { let items = db_page .items .into_iter() .map(|item| ThreadItem { path: item.rollout_path, thread_id: Some(item.id), first_user_message: item.first_user_message, cwd: Some(item.cwd), git_branch: item.git_branch, git_sha: item.git_sha, git_origin_url: item.git_origin_url, source: Some( serde_json::from_str(item.source.as_str()) .or_else(|_| serde_json::from_value(Value::String(item.source))) .unwrap_or(SessionSource::Unknown), ), agent_nickname: item.agent_nickname, agent_role: item.agent_role, model_provider: Some(item.model_provider), cli_version: Some(item.cli_version), created_at: Some(item.created_at.to_rfc3339_opts(SecondsFormat::Secs, true)), updated_at: Some(item.updated_at.to_rfc3339_opts(SecondsFormat::Secs, true)), }) .collect(); Self { items, next_cursor: db_page.next_anchor.map(Into::into), num_scanned_files: db_page.num_scanned_rows, reached_scan_cap: false, } } } async fn select_resume_path( page: &ThreadsPage, filter_cwd: Option<&Path>, default_provider: &str, ) -> Option { match filter_cwd { Some(cwd) => { for item in &page.items { if resume_candidate_matches_cwd( item.path.as_path(), item.cwd.as_deref(), cwd, default_provider, ) .await { return Some(item.path.clone()); } } None } None => page.items.first().map(|item| item.path.clone()), } } async fn resume_candidate_matches_cwd( rollout_path: &Path, cached_cwd: Option<&Path>, cwd: &Path, default_provider: &str, ) -> bool { if cached_cwd.is_some_and(|session_cwd| cwd_matches(session_cwd, cwd)) { return true; } if let Ok((items, _, _)) = RolloutRecorder::load_rollout_items(rollout_path).await && let Some(latest_turn_context_cwd) = items.iter().rev().find_map(|item| match item { RolloutItem::TurnContext(turn_context) => Some(turn_context.cwd.as_path()), RolloutItem::SessionMeta(_) | RolloutItem::ResponseItem(_) | RolloutItem::Compacted(_) | RolloutItem::EventMsg(_) => None, }) { return cwd_matches(latest_turn_context_cwd, cwd); } metadata::extract_metadata_from_rollout(rollout_path, default_provider) .await .is_ok_and(|outcome| cwd_matches(outcome.metadata.cwd.as_path(), cwd)) } async fn select_resume_path_from_db_page( page: &codex_state::ThreadsPage, filter_cwd: Option<&Path>, default_provider: &str, ) -> Option { match filter_cwd { Some(cwd) => { for item in &page.items { if resume_candidate_matches_cwd( item.rollout_path.as_path(), Some(item.cwd.as_path()), cwd, default_provider, ) .await { return Some(item.rollout_path.clone()); } } None } None => page.items.first().map(|item| item.rollout_path.clone()), } } fn cwd_matches(session_cwd: &Path, cwd: &Path) -> bool { if let (Ok(ca), Ok(cb)) = ( path_utils::normalize_for_path_comparison(session_cwd), path_utils::normalize_for_path_comparison(cwd), ) { return ca == cb; } session_cwd == cwd } #[cfg(test)] #[path = "recorder_tests.rs"] mod tests;