From 6d8e12ac42508c10bbe1d1769aee7f49d150ba5c Mon Sep 17 00:00:00 2001 From: stefanstokic-oai Date: Mon, 8 Jun 2026 14:16:32 -0400 Subject: [PATCH] [codex] Speed up external agent session imports (#26637) ## Why Importing large external-agent session histories currently starts a full live Codex thread for every imported session. This initializes unrelated runtime systems and repeats expensive transcript, metadata, hashing, and ledger work. On a 50-session, 238 MiB fixture, the existing path took roughly 70 seconds to complete the import and 77 seconds end to end. ## What changed - Persist imported sessions directly through `ThreadStore` instead of starting full live threads. - Process imports through a bounded five-session pipeline. - Parse, extract, and hash each source file in one pass. - Move blocking source preparation onto the blocking thread pool. - Reuse prepared content hashes and update the import ledger once per batch. - Avoid metadata readback for newly written rollouts. - Preserve imported conversation history and visible thread metadata. - Keep the implementation out of `codex-core` and avoid changes to the public `ThreadStore` trait. ## Performance For the same 50-session, 238 MiB fixture: | Path | Import completion | End to end | | --- | ---: | ---: | | Existing import | 69.61s | 76.62s | | This change | 5.95s | 6.58s | All 50 sessions imported successfully with no warnings or contention signals. ## Validation - `just test -p codex-external-agent-sessions` - `just test -p codex-app-server external_agent_config_import` - Verified imports do not initialize unrelated required MCP servers. - Verified previously imported source versions are skipped and changed sources can be imported again. - Verified imported rollouts remain readable through thread listing and history APIs. --- codex-rs/app-server/src/message_processor.rs | 1 + codex-rs/app-server/src/request_processors.rs | 1 + .../external_agent_config_processor.rs | 143 +--------- .../external_agent_session_import.rs | 260 ++++++++++++++++++ .../tests/suite/v2/external_agent_config.rs | 92 +++++++ .../external-agent-sessions/src/export.rs | 191 +++++++------ .../external-agent-sessions/src/ledger.rs | 67 +++-- .../src/ledger_tests.rs | 85 ++++++ codex-rs/external-agent-sessions/src/lib.rs | 192 +++++-------- .../external-agent-sessions/src/records.rs | 155 +++++++---- 10 files changed, 793 insertions(+), 394 deletions(-) create mode 100644 codex-rs/app-server/src/request_processors/external_agent_session_import.rs create mode 100644 codex-rs/external-agent-sessions/src/ledger_tests.rs diff --git a/codex-rs/app-server/src/message_processor.rs b/codex-rs/app-server/src/message_processor.rs index 65c57bf26..84d3d88dd 100644 --- a/codex-rs/app-server/src/message_processor.rs +++ b/codex-rs/app-server/src/message_processor.rs @@ -475,6 +475,7 @@ impl MessageProcessor { let external_agent_config_processor = ExternalAgentConfigRequestProcessor::new( outgoing.clone(), Arc::clone(&thread_manager), + Arc::clone(&thread_store), config_manager.clone(), config_processor.clone(), arg0_paths, diff --git a/codex-rs/app-server/src/request_processors.rs b/codex-rs/app-server/src/request_processors.rs index 57205798b..79da49fba 100644 --- a/codex-rs/app-server/src/request_processors.rs +++ b/codex-rs/app-server/src/request_processors.rs @@ -461,6 +461,7 @@ mod command_exec_processor; mod config_processor; mod environment_processor; mod external_agent_config_processor; +mod external_agent_session_import; mod feedback_doctor_report; mod feedback_processor; mod fs_processor; diff --git a/codex-rs/app-server/src/request_processors/external_agent_config_processor.rs b/codex-rs/app-server/src/request_processors/external_agent_config_processor.rs index 87956182d..5d4c0594e 100644 --- a/codex-rs/app-server/src/request_processors/external_agent_config_processor.rs +++ b/codex-rs/app-server/src/request_processors/external_agent_config_processor.rs @@ -26,53 +26,47 @@ use codex_app_server_protocol::MigrationDetails; use codex_app_server_protocol::PluginsMigration; use codex_app_server_protocol::ServerNotification; use codex_arg0::Arg0DispatchPaths; -use codex_core::StartThreadOptions; use codex_core::ThreadManager; -use codex_core::config::ConfigOverrides; use codex_external_agent_sessions::ExternalAgentSessionMigration as CoreSessionMigration; -use codex_external_agent_sessions::ImportedExternalAgentSession; -use codex_external_agent_sessions::PendingSessionImport; -use codex_external_agent_sessions::prepare_validated_session_imports; -use codex_external_agent_sessions::record_imported_session; -use codex_protocol::ThreadId; -use codex_protocol::protocol::InitialHistory; -use codex_thread_store::ThreadMetadataPatch; +use codex_thread_store::ThreadStore; use std::collections::HashSet; use std::path::PathBuf; -use tokio::sync::Semaphore; use super::ConfigRequestProcessor; +use super::external_agent_session_import::ExternalAgentSessionImporter; #[derive(Clone)] pub(crate) struct ExternalAgentConfigRequestProcessor { outgoing: Arc, - codex_home: PathBuf, migration_service: ExternalAgentConfigService, - session_import_permits: Arc, + session_importer: ExternalAgentSessionImporter, thread_manager: Arc, - config_manager: ConfigManager, config_processor: ConfigRequestProcessor, - arg0_paths: Arg0DispatchPaths, } impl ExternalAgentConfigRequestProcessor { pub(crate) fn new( outgoing: Arc, thread_manager: Arc, + thread_store: Arc, config_manager: ConfigManager, config_processor: ConfigRequestProcessor, arg0_paths: Arg0DispatchPaths, codex_home: PathBuf, ) -> Self { + let session_importer = ExternalAgentSessionImporter::new( + codex_home.clone(), + Arc::clone(&thread_manager), + thread_store, + config_manager, + arg0_paths, + ); Self { outgoing, - migration_service: ExternalAgentConfigService::new(codex_home.clone()), - codex_home, - session_import_permits: Arc::new(Semaphore::new(1)), + migration_service: ExternalAgentConfigService::new(codex_home), + session_importer, thread_manager, - config_manager, config_processor, - arg0_paths, } } @@ -207,42 +201,12 @@ impl ExternalAgentConfigRequestProcessor { return Ok(()); } - let session_import_permits = Arc::clone(&self.session_import_permits); - let session_processor = self.clone(); + let session_importer = self.session_importer.clone(); let plugin_processor = self.clone(); let outgoing = Arc::clone(&self.outgoing); let thread_manager = Arc::clone(&self.thread_manager); tokio::spawn(async move { - let session_imports = async move { - if !pending_session_imports.is_empty() { - let Ok(_session_import_permit) = session_import_permits.acquire_owned().await - else { - return; - }; - let pending_session_imports = session_processor - .prepare_validated_session_imports(pending_session_imports); - for pending_session_import in pending_session_imports { - match session_processor - .import_external_agent_session(pending_session_import.session) - .await - { - Ok(imported_thread_id) => { - session_processor.record_imported_session( - &pending_session_import.source_path, - imported_thread_id, - ); - } - Err(error) => { - tracing::warn!( - error = %error.message, - path = %pending_session_import.source_path.display(), - "external agent session import failed" - ); - } - } - } - } - }; + let session_imports = session_importer.import_sessions(pending_session_imports); let plugin_imports = async move { for pending_plugin_import in pending_plugin_imports { match plugin_processor @@ -274,65 +238,6 @@ impl ExternalAgentConfigRequestProcessor { Ok(()) } - async fn import_external_agent_session( - &self, - session: ImportedExternalAgentSession, - ) -> Result { - let ImportedExternalAgentSession { - cwd, - title, - rollout_items, - } = session; - let config = self - .config_manager - .load_with_overrides( - /*request_overrides*/ None, - ConfigOverrides { - cwd: Some(PathBuf::from(cwd.to_string_lossy().into_owned())), - codex_linux_sandbox_exe: self.arg0_paths.codex_linux_sandbox_exe.clone(), - main_execve_wrapper_exe: self.arg0_paths.main_execve_wrapper_exe.clone(), - ..Default::default() - }, - ) - .await - .map_err(|err| { - internal_error(format!("failed to load imported session config: {err}")) - })?; - let environments = self - .thread_manager - .default_environment_selections(&config.cwd); - let imported_thread = self - .thread_manager - .start_thread_with_options(StartThreadOptions { - config, - initial_history: InitialHistory::Forked(rollout_items), - session_source: None, - thread_source: None, - dynamic_tools: Vec::new(), - metrics_service_name: None, - parent_trace: None, - environments, - }) - .await - .map_err(|err| internal_error(format!("failed to import session: {err}")))?; - if let Some(title) = title - && let Some(name) = codex_core::util::normalize_thread_name(&title) - { - imported_thread - .thread - .update_thread_metadata( - ThreadMetadataPatch { - name: Some(Some(name)), - ..Default::default() - }, - /*include_archived*/ false, - ) - .await - .map_err(|err| internal_error(format!("failed to name imported session: {err}")))?; - } - Ok(imported_thread.thread_id) - } - fn validate_pending_session_imports( &self, params: &ExternalAgentConfigImportParams, @@ -371,24 +276,6 @@ impl ExternalAgentConfigRequestProcessor { Ok(selected_sessions) } - fn prepare_validated_session_imports( - &self, - sessions: Vec, - ) -> Vec { - prepare_validated_session_imports(&self.codex_home, sessions) - } - - fn record_imported_session(&self, source_path: &std::path::Path, imported_thread_id: ThreadId) { - if let Err(err) = record_imported_session(&self.codex_home, source_path, imported_thread_id) - { - tracing::warn!( - error = %err, - path = %source_path.display(), - "external agent session import ledger update failed" - ); - } - } - async fn import_external_agent_config( &self, params: ExternalAgentConfigImportParams, diff --git a/codex-rs/app-server/src/request_processors/external_agent_session_import.rs b/codex-rs/app-server/src/request_processors/external_agent_session_import.rs new file mode 100644 index 000000000..163210f16 --- /dev/null +++ b/codex-rs/app-server/src/request_processors/external_agent_session_import.rs @@ -0,0 +1,260 @@ +use std::path::PathBuf; +use std::sync::Arc; + +use chrono::Utc; +use codex_arg0::Arg0DispatchPaths; +use codex_core::ThreadManager; +use codex_core::config::ConfigOverrides; +use codex_external_agent_sessions::CompletedExternalAgentSessionImport; +use codex_external_agent_sessions::ExternalAgentSessionMigration; +use codex_external_agent_sessions::ImportedExternalAgentSession; +use codex_external_agent_sessions::PendingSessionImport; +use codex_external_agent_sessions::prepare_validated_session_import; +use codex_external_agent_sessions::record_completed_session_imports; +use codex_models_manager::manager::RefreshStrategy; +use codex_protocol::ThreadId; +use codex_protocol::models::BaseInstructions; +use codex_protocol::protocol::MultiAgentVersion; +use codex_protocol::protocol::ThreadMemoryMode; +use codex_rollout::is_persisted_rollout_item; +use codex_thread_store::AppendThreadItemsParams; +use codex_thread_store::CreateThreadParams; +use codex_thread_store::ThreadMetadataPatch; +use codex_thread_store::ThreadPersistenceMetadata; +use codex_thread_store::ThreadStore; +use codex_thread_store::UpdateThreadMetadataParams; +use futures::StreamExt; +use tokio::sync::Semaphore; + +use crate::config_manager::ConfigManager; + +const SESSION_IMPORT_CONCURRENCY: usize = 5; + +#[derive(Clone)] +pub(super) struct ExternalAgentSessionImporter { + codex_home: PathBuf, + permits: Arc, + thread_manager: Arc, + thread_store: Arc, + config_manager: ConfigManager, + arg0_paths: Arg0DispatchPaths, +} + +impl ExternalAgentSessionImporter { + pub(super) fn new( + codex_home: PathBuf, + thread_manager: Arc, + thread_store: Arc, + config_manager: ConfigManager, + arg0_paths: Arg0DispatchPaths, + ) -> Self { + Self { + codex_home, + permits: Arc::new(Semaphore::new(1)), + thread_manager, + thread_store, + config_manager, + arg0_paths, + } + } + + pub(super) async fn import_sessions(&self, sessions: Vec) { + if sessions.is_empty() { + return; + } + let Ok(_permit) = self.permits.acquire().await else { + return; + }; + let import_results = futures::stream::iter(sessions) + .map(|session| { + let importer = self.clone(); + async move { importer.import_requested_session(session).await } + }) + .buffer_unordered(SESSION_IMPORT_CONCURRENCY); + futures::pin_mut!(import_results); + + let mut completed_imports = Vec::new(); + while let Some(result) = import_results.next().await { + match result { + Ok(Some(completed_import)) => completed_imports.push(completed_import), + Ok(None) => {} + Err(failure) => { + tracing::warn!( + error = %failure.message, + path = %failure.source_path.display(), + "external agent session import failed" + ); + } + } + } + if let Err(err) = record_completed_session_imports(&self.codex_home, completed_imports) { + tracing::warn!( + error = %err, + "external agent session import ledger update failed" + ); + } + } + + async fn import_requested_session( + &self, + session: ExternalAgentSessionMigration, + ) -> Result, SessionImportFailure> { + let source_path = session.path.clone(); + let Some(pending_import) = + self.prepare_session_import(session) + .await + .map_err(|message| SessionImportFailure { + source_path: source_path.clone(), + message, + })? + else { + return Ok(None); + }; + let imported_thread_id = + self.persist_session(pending_import.session) + .await + .map_err(|message| SessionImportFailure { + source_path: pending_import.source_path.clone(), + message, + })?; + Ok(Some(CompletedExternalAgentSessionImport { + source_path: pending_import.source_path, + source_content_sha256: pending_import.source_content_sha256, + imported_thread_id, + })) + } + + async fn prepare_session_import( + &self, + session: ExternalAgentSessionMigration, + ) -> Result, String> { + let codex_home = self.codex_home.clone(); + tokio::task::spawn_blocking(move || prepare_validated_session_import(&codex_home, session)) + .await + .map_err(|err| format!("external agent session preparation task failed: {err}"))? + .map_err(|err| format!("failed to prepare external agent session: {err}")) + } + + async fn persist_session( + &self, + session: ImportedExternalAgentSession, + ) -> Result { + let ImportedExternalAgentSession { + cwd, + title, + first_user_message, + mut rollout_items, + } = session; + let config = self + .config_manager + .load_with_overrides( + /*request_overrides*/ None, + ConfigOverrides { + cwd: Some(cwd), + codex_linux_sandbox_exe: self.arg0_paths.codex_linux_sandbox_exe.clone(), + main_execve_wrapper_exe: self.arg0_paths.main_execve_wrapper_exe.clone(), + ..Default::default() + }, + ) + .await + .map_err(|err| format!("failed to load imported session config: {err}"))?; + let models_manager = self.thread_manager.get_models_manager(); + let model = models_manager + .get_default_model(&config.model, RefreshStrategy::Offline) + .await; + let model_info = models_manager + .get_model_info(model.as_str(), &config.to_models_manager_config()) + .await; + let thread_id = ThreadId::new(); + let source = self.thread_manager.session_source(); + let cwd = config.cwd.to_path_buf(); + let model_provider = config.model_provider_id.clone(); + let memory_mode = if config.memories.generate_memories { + ThreadMemoryMode::Enabled + } else { + ThreadMemoryMode::Disabled + }; + let now = Utc::now(); + let create_params = CreateThreadParams { + thread_id, + forked_from_id: None, + parent_thread_id: None, + source: source.clone(), + thread_source: None, + base_instructions: BaseInstructions { + text: config + .base_instructions + .clone() + .unwrap_or_else(|| model_info.get_model_instructions(config.personality)), + }, + dynamic_tools: Vec::new(), + multi_agent_version: Some(MultiAgentVersion::V1), + metadata: ThreadPersistenceMetadata { + cwd: Some(cwd.clone()), + model_provider: model_provider.clone(), + memory_mode, + }, + }; + rollout_items.retain(is_persisted_rollout_item); + let title = title + .as_deref() + .and_then(codex_core::util::normalize_thread_name); + let metadata = ThreadMetadataPatch { + title, + preview: first_user_message.clone(), + model_provider: Some(model_provider), + created_at: Some(now), + updated_at: Some(now), + source: Some(source.clone()), + thread_source: Some(None), + agent_nickname: Some(source.get_nickname()), + agent_role: Some(source.get_agent_role()), + agent_path: Some(source.get_agent_path().map(Into::into)), + cwd: Some(cwd), + cli_version: Some(env!("CARGO_PKG_VERSION").to_string()), + first_user_message, + memory_mode: Some(memory_mode), + ..Default::default() + }; + + self.thread_store + .create_thread(create_params) + .await + .map_err(|err| format!("failed to import session: {err}"))?; + if !rollout_items.is_empty() + && let Err(err) = self + .thread_store + .append_items(AppendThreadItemsParams { + thread_id, + items: rollout_items, + }) + .await + { + let _ = self.thread_store.discard_thread(thread_id).await; + return Err(format!("failed to import session: {err}")); + } + + self.thread_store + .update_thread_metadata(UpdateThreadMetadataParams { + thread_id, + patch: metadata, + include_archived: false, + }) + .await + .map_err(|err| format!("failed to update imported session: {err}"))?; + self.thread_store + .persist_thread(thread_id) + .await + .map_err(|err| format!("failed to persist imported session: {err}"))?; + self.thread_store + .shutdown_thread(thread_id) + .await + .map_err(|err| format!("failed to shutdown imported session: {err}"))?; + Ok(thread_id) + } +} + +struct SessionImportFailure { + source_path: PathBuf, + message: String, +} diff --git a/codex-rs/app-server/tests/suite/v2/external_agent_config.rs b/codex-rs/app-server/tests/suite/v2/external_agent_config.rs index fdb996c2f..991e12dc7 100644 --- a/codex-rs/app-server/tests/suite/v2/external_agent_config.rs +++ b/codex-rs/app-server/tests/suite/v2/external_agent_config.rs @@ -435,6 +435,98 @@ async fn external_agent_config_import_creates_session_rollouts() -> Result<()> { Ok(()) } +#[tokio::test] +async fn external_agent_config_import_does_not_initialize_required_mcp() -> Result<()> { + let server = create_mock_responses_server_repeating_assistant("unused").await; + let codex_home = TempDir::new()?; + create_config_toml(codex_home.path(), &server.uri())?; + let mut config = std::fs::read_to_string(codex_home.path().join("config.toml"))?; + config.push_str( + r#" +[mcp_servers.required_broken] +command = "this-command-does-not-exist" +required = true +"#, + ); + std::fs::write(codex_home.path().join("config.toml"), config)?; + let project_root = codex_home.path().join("repo"); + let recent_timestamp = chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true); + let session_dir = codex_home.path().join(".claude/projects/repo"); + let session_path = session_dir.join("session.jsonl"); + std::fs::create_dir_all(&project_root)?; + std::fs::create_dir_all(&session_dir)?; + std::fs::write( + &session_path, + serde_json::json!({ + "type": "user", + "cwd": &project_root, + "timestamp": &recent_timestamp, + "message": { "content": "first request" }, + }) + .to_string(), + )?; + + let home_dir = codex_home.path().display().to_string(); + let mut mcp = + TestAppServer::new_with_env(codex_home.path(), &[("HOME", Some(home_dir.as_str()))]) + .await?; + timeout(DEFAULT_TIMEOUT, mcp.initialize()).await??; + + let request_id = mcp + .send_raw_request( + "externalAgentConfig/import", + Some(serde_json::json!({ + "migrationItems": [{ + "itemType": "SESSIONS", + "description": "Migrate recent sessions", + "cwd": null, + "details": { + "sessions": [{ + "path": session_path, + "cwd": project_root, + "title": "first request" + }] + } + }] + })), + ) + .await?; + timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(request_id)), + ) + .await??; + timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_notification_message("externalAgentConfig/import/completed"), + ) + .await??; + + let request_id = mcp + .send_thread_list_request(ThreadListParams { + cursor: None, + limit: None, + sort_key: None, + sort_direction: None, + model_providers: None, + source_kinds: None, + archived: None, + cwd: None, + use_state_db_only: false, + search_term: None, + }) + .await?; + let response: JSONRPCResponse = timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(request_id)), + ) + .await??; + let response: ThreadListResponse = to_response(response)?; + assert_eq!(response.data.len(), 1); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn external_agent_config_import_accepts_detected_session_payload_after_restart() -> Result<()> { diff --git a/codex-rs/external-agent-sessions/src/export.rs b/codex-rs/external-agent-sessions/src/export.rs index e09805d4d..05a3004ee 100644 --- a/codex-rs/external-agent-sessions/src/export.rs +++ b/codex-rs/external-agent-sessions/src/export.rs @@ -1,10 +1,7 @@ use crate::ConversationMessage; use crate::ImportedExternalAgentSession; use crate::MessageRole; -use crate::records::conversation_messages; -use crate::records::project_root_from_records; -use crate::records::read_records; -use crate::records::source_title_from_records; +use crate::records::read_session_import; use crate::summarize_for_label; use codex_protocol::models::ContentItem; use codex_protocol::models::ResponseItem; @@ -23,44 +20,55 @@ use std::path::Path; const EXTERNAL_SESSION_IMPORTED_MARKER: &str = ""; -pub fn load_session_for_import(path: &Path) -> io::Result> { - let records = read_records(path)?; - let Some(cwd) = project_root_from_records(&records) else { +#[cfg(test)] +fn load_session_for_import(path: &Path) -> io::Result> { + Ok( + load_session_for_import_with_content_sha256(path)? + .map(|(session, _content_sha256)| session), + ) +} + +pub(crate) fn load_session_for_import_with_content_sha256( + path: &Path, +) -> io::Result> { + let parsed = read_session_import(path)?; + let Some(cwd) = parsed.cwd else { return Ok(None); }; - let messages = conversation_messages(&records); - let rollout_items = rollout_items_from_messages(&messages); + let messages = parsed.messages; + let first_user_message = messages + .iter() + .find(|message| message.role == MessageRole::User) + .map(|message| summarize_for_label(&message.text)); + let title = parsed.source_title.or_else(|| first_user_message.clone()); + let rollout_items = rollout_items_from_messages(messages); if rollout_items.is_empty() { return Ok(None); } - let title = source_title_from_records(&records).or_else(|| { - messages - .iter() - .find(|message| message.role == MessageRole::User) - .map(|message| summarize_for_label(&message.text)) - }); - Ok(Some(ImportedExternalAgentSession { - cwd, - title, - rollout_items, - })) + Ok(Some(( + ImportedExternalAgentSession { + cwd, + title, + first_user_message, + rollout_items, + }, + parsed.content_sha256, + ))) } -fn rollout_items_from_messages(messages: &[ConversationMessage]) -> Vec { +fn rollout_items_from_messages(messages: Vec) -> Vec { let mut items = Vec::new(); - let mut response_items = Vec::new(); - let mut current_turn: Option<(String, Option)> = None; + let mut current_turn = None; + let mut response_item_bytes = 0i64; + let mut last_model_visible_tokens = 0i64; let mut user_turn_count = 0usize; + let completed_at = messages.last().and_then(|message| message.timestamp); for message in messages { match message.role { MessageRole::User => { - if let Some((turn_id, last_agent_message)) = current_turn.take() { - items.push(turn_complete_item( - turn_id, - last_agent_message, - /*completed_at*/ None, - )); + if let Some(turn_id) = current_turn.take() { + items.push(turn_complete_item(turn_id, /*completed_at*/ None)); } user_turn_count += 1; let turn_id = format!("external-import-turn-{user_turn_count}"); @@ -73,28 +81,24 @@ fn rollout_items_from_messages(messages: &[ConversationMessage]) -> Vec { - let Some((_, last_agent_message)) = current_turn.as_mut() else { + if current_turn.is_none() { continue; - }; - let response_item = response_item(message); - response_items.push(response_item.clone()); - items.push(RolloutItem::ResponseItem(response_item)); + } + response_item_bytes = + response_item_bytes.saturating_add(message_byte_count(&message)); + last_model_visible_tokens = approx_tokens_from_byte_count_i64(response_item_bytes); items.push(RolloutItem::EventMsg(EventMsg::AgentMessage( AgentMessageEvent { message: message.text.clone(), @@ -102,20 +106,15 @@ fn rollout_items_from_messages(messages: &[ConversationMessage]) -> Vec RolloutItem { })) } -fn response_item(message: &ConversationMessage) -> ResponseItem { +fn response_item(message: ConversationMessage) -> ResponseItem { let content = match message.role { - MessageRole::Assistant => ContentItem::OutputText { - text: message.text.clone(), - }, - MessageRole::User => ContentItem::InputText { - text: message.text.clone(), - }, + MessageRole::Assistant => ContentItem::OutputText { text: message.text }, + MessageRole::User => ContentItem::InputText { text: message.text }, }; ResponseItem::Message { id: None, @@ -149,13 +144,11 @@ fn response_item(message: &ConversationMessage) -> ResponseItem { } } -fn token_count_item(response_items: &[ResponseItem]) -> RolloutItem { - let last_model_generated = response_items.iter().rposition( - |item| matches!(item, ResponseItem::Message { role, .. } if role == "assistant"), - ); - let last_model_visible_tokens = last_model_generated - .map(|index| estimate_response_items_token_count(&response_items[..=index])) - .unwrap_or_default(); +fn message_byte_count(message: &ConversationMessage) -> i64 { + i64::try_from(message.text.len()).unwrap_or(i64::MAX) +} + +fn token_count_item(last_model_visible_tokens: i64) -> RolloutItem { let usage = TokenUsage { total_tokens: last_model_visible_tokens, ..TokenUsage::default() @@ -170,26 +163,10 @@ fn token_count_item(response_items: &[ResponseItem]) -> RolloutItem { })) } -fn estimate_response_items_token_count(response_items: &[ResponseItem]) -> i64 { - response_items - .iter() - .map(|item| { - serde_json::to_string(item) - .map(|serialized| i64::try_from(serialized.len()).unwrap_or(i64::MAX)) - .map(approx_tokens_from_byte_count_i64) - .unwrap_or_default() - }) - .fold(0i64, i64::saturating_add) -} - -fn turn_complete_item( - turn_id: String, - last_agent_message: Option, - completed_at: Option, -) -> RolloutItem { +fn turn_complete_item(turn_id: String, completed_at: Option) -> RolloutItem { RolloutItem::EventMsg(EventMsg::TurnComplete(TurnCompleteEvent { turn_id, - last_agent_message, + last_agent_message: None, completed_at, duration_ms: None, time_to_first_token_ms: None, @@ -241,7 +218,7 @@ mod tests { } #[test] - fn adds_import_marker_without_replacing_last_agent_message() { + fn adds_import_marker_without_copying_last_agent_message() { let root = TempDir::new().expect("tempdir"); let project_root = root.path().join("repo"); std::fs::create_dir_all(&project_root).expect("project root"); @@ -280,10 +257,54 @@ mod tests { }); assert_eq!( last_turn_complete.and_then(|event| event.last_agent_message.as_deref()), - Some("first answer") + None ); } + #[test] + fn stores_imported_messages_as_response_items_and_visible_events() { + let root = TempDir::new().expect("tempdir"); + let project_root = root.path().join("repo"); + std::fs::create_dir_all(&project_root).expect("project root"); + let path = root.path().join("session.jsonl"); + let request = "r".repeat(1_000); + let answer = "a".repeat(1_000); + std::fs::write( + &path, + jsonl(&[ + record("user", &request, &project_root), + record("assistant", &answer, &project_root), + ]), + ) + .expect("session"); + + let imported = load_session_for_import(&path) + .expect("load") + .expect("session"); + let response_message_count = imported + .rollout_items + .iter() + .filter(|item| { + matches!( + item, + RolloutItem::ResponseItem(ResponseItem::Message { .. }) + ) + }) + .count(); + let visible_message_event_count = imported + .rollout_items + .iter() + .filter(|item| match item { + RolloutItem::EventMsg(EventMsg::UserMessage(event)) => event.message == request, + RolloutItem::EventMsg(EventMsg::AgentMessage(event)) => event.message == answer, + _ => false, + }) + .count(); + + assert_eq!(response_message_count, 2); + assert_eq!(visible_message_event_count, 2); + } + #[test] fn loads_custom_title_for_imported_session() { let root = TempDir::new().expect("tempdir"); diff --git a/codex-rs/external-agent-sessions/src/ledger.rs b/codex-rs/external-agent-sessions/src/ledger.rs index 45ff97bcf..9a3b2042b 100644 --- a/codex-rs/external-agent-sessions/src/ledger.rs +++ b/codex-rs/external-agent-sessions/src/ledger.rs @@ -30,6 +30,13 @@ struct ImportedExternalAgentSessionRecord { source_modified_at: Option, } +#[derive(Debug, PartialEq, Eq)] +pub struct CompletedExternalAgentSessionImport { + pub source_path: PathBuf, + pub source_content_sha256: String, + pub imported_thread_id: ThreadId, +} + #[derive(Debug, Clone, Copy)] pub(super) struct ImportedSourceState { pub source_modified_at: Option, @@ -43,29 +50,50 @@ pub fn has_current_session_been_imported( load_import_ledger(codex_home)?.contains_current_source(source_path) } -pub fn record_imported_session( +#[cfg(test)] +pub(crate) fn record_imported_session( codex_home: &Path, source_path: &Path, imported_thread_id: ThreadId, ) -> io::Result<()> { - let mut ledger = load_import_ledger(codex_home)?; let source_path = canonical_source_path(source_path)?; - let content_sha256 = session_content_sha256(&source_path)?; - let source_modified_at = session_modified_at(&source_path)?; - if let Some(index) = ledger.records.iter().rposition(|record| { - record.source_path == source_path && record.content_sha256 == content_sha256 - }) { - let mut record = ledger.records.remove(index); - record.imported_thread_id = imported_thread_id; - record.imported_at = now_unix_seconds(); - record.source_modified_at = source_modified_at; - ledger.records.push(record); - } else { - ledger.records.push(ImportedExternalAgentSessionRecord { + record_completed_session_imports( + codex_home, + vec![CompletedExternalAgentSessionImport { + source_content_sha256: session_content_sha256(&source_path)?, source_path, - content_sha256, imported_thread_id, - imported_at: now_unix_seconds(), + }], + ) +} + +pub fn record_completed_session_imports( + codex_home: &Path, + imports: Vec, +) -> io::Result<()> { + if imports.is_empty() { + return Ok(()); + } + let mut ledger = load_import_ledger(codex_home)?; + let imported_at = now_unix_seconds(); + for import in imports { + let source_modified_at = session_modified_at(&import.source_path).ok().flatten(); + if let Some(index) = ledger.records.iter().rposition(|record| { + record.source_path == import.source_path + && record.content_sha256 == import.source_content_sha256 + }) { + let mut record = ledger.records.remove(index); + record.imported_thread_id = import.imported_thread_id; + record.imported_at = imported_at; + record.source_modified_at = source_modified_at.or(record.source_modified_at); + ledger.records.push(record); + continue; + } + ledger.records.push(ImportedExternalAgentSessionRecord { + source_path: import.source_path, + content_sha256: import.source_content_sha256, + imported_thread_id: import.imported_thread_id, + imported_at, source_modified_at, }); } @@ -88,6 +116,9 @@ impl ImportedExternalAgentSessionLedger { } pub(super) fn contains_current_source(&self, source_path: &Path) -> io::Result { + if self.records.is_empty() { + return Ok(false); + } let source_path = canonical_source_path(source_path)?; if !self .records @@ -188,3 +219,7 @@ fn session_modified_at(path: &Path) -> io::Result> { .ok() .and_then(|duration| i64::try_from(duration.as_nanos()).ok())) } + +#[cfg(test)] +#[path = "ledger_tests.rs"] +mod tests; diff --git a/codex-rs/external-agent-sessions/src/ledger_tests.rs b/codex-rs/external-agent-sessions/src/ledger_tests.rs new file mode 100644 index 000000000..f4b9dd6f7 --- /dev/null +++ b/codex-rs/external-agent-sessions/src/ledger_tests.rs @@ -0,0 +1,85 @@ +use super::CompletedExternalAgentSessionImport; +use super::ImportedExternalAgentSessionLedger; +use super::record_completed_session_imports; +use codex_protocol::ThreadId; +use sha2::Digest; +use sha2::Sha256; +use tempfile::TempDir; + +#[test] +fn empty_ledger_does_not_read_source() { + let root = TempDir::new().expect("tempdir"); + let missing_source = root.path().join("missing-session.jsonl"); + + assert!( + !ImportedExternalAgentSessionLedger::default() + .contains_current_source(&missing_source) + .expect("empty ledger cannot contain sources") + ); +} + +#[test] +fn completed_imports_do_not_read_source_files() { + let root = TempDir::new().expect("tempdir"); + let codex_home = root.path().join("codex-home"); + let source_path = root.path().join("session.jsonl"); + let contents = b"session contents"; + std::fs::write(&source_path, contents).expect("source"); + let source_path = std::fs::canonicalize(&source_path).expect("canonical source"); + std::fs::remove_file(&source_path).expect("remove source"); + let imported_thread_id = ThreadId::new(); + + record_completed_session_imports( + &codex_home, + vec![CompletedExternalAgentSessionImport { + source_path: source_path.clone(), + source_content_sha256: format!("{:x}", Sha256::digest(contents)), + imported_thread_id, + }], + ) + .expect("record completed imports"); + + let ledger = super::load_import_ledger(&codex_home).expect("ledger"); + assert_eq!(ledger.records.len(), 1); + assert_eq!(ledger.records[0].source_path, source_path); + assert_eq!(ledger.records[0].imported_thread_id, imported_thread_id); + assert_eq!(ledger.records[0].source_modified_at, None); +} + +#[test] +fn completed_import_refreshes_existing_record_metadata() { + let root = TempDir::new().expect("tempdir"); + let codex_home = root.path().join("codex-home"); + let source_path = root.path().join("session.jsonl"); + let contents = b"session contents"; + std::fs::write(&source_path, contents).expect("source"); + let source_path = std::fs::canonicalize(source_path).expect("canonical source"); + let content_sha256 = format!("{:x}", Sha256::digest(contents)); + let first_thread_id = ThreadId::new(); + let second_thread_id = ThreadId::new(); + + record_completed_session_imports( + &codex_home, + vec![CompletedExternalAgentSessionImport { + source_path: source_path.clone(), + source_content_sha256: content_sha256.clone(), + imported_thread_id: first_thread_id, + }], + ) + .expect("record first import"); + record_completed_session_imports( + &codex_home, + vec![CompletedExternalAgentSessionImport { + source_path: source_path.clone(), + source_content_sha256: content_sha256, + imported_thread_id: second_thread_id, + }], + ) + .expect("record replacement import"); + + let ledger = super::load_import_ledger(&codex_home).expect("ledger"); + assert_eq!(ledger.records.len(), 1); + assert_eq!(ledger.records[0].source_path, source_path); + assert_eq!(ledger.records[0].imported_thread_id, second_thread_id); + assert!(ledger.records[0].source_modified_at.is_some()); +} diff --git a/codex-rs/external-agent-sessions/src/lib.rs b/codex-rs/external-agent-sessions/src/lib.rs index fe9699f0c..0b7a4eb2b 100644 --- a/codex-rs/external-agent-sessions/src/lib.rs +++ b/codex-rs/external-agent-sessions/src/lib.rs @@ -6,15 +6,15 @@ mod ledger; mod records; use codex_protocol::protocol::RolloutItem; -use std::collections::HashSet; use std::io; use std::path::Path; use std::path::PathBuf; pub use detect::detect_recent_sessions; -pub use export::load_session_for_import; +use export::load_session_for_import_with_content_sha256; +pub use ledger::CompletedExternalAgentSessionImport; pub use ledger::has_current_session_been_imported; -pub use ledger::record_imported_session; +pub use ledger::record_completed_session_imports; pub use records::SessionSummary; pub use records::summarize_session; @@ -31,105 +31,51 @@ pub struct ExternalAgentSessionMigration { pub struct ImportedExternalAgentSession { pub cwd: PathBuf, pub title: Option, + pub first_user_message: Option, pub rollout_items: Vec, } #[derive(Debug, Clone)] pub struct PendingSessionImport { pub source_path: PathBuf, + pub source_content_sha256: String, pub session: ImportedExternalAgentSession, } -#[derive(Debug)] -pub enum PrepareSessionImportsError { - SessionNotDetected(PathBuf), -} - -impl std::fmt::Display for PrepareSessionImportsError { - fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - match self { - PrepareSessionImportsError::SessionNotDetected(path) => { - write!( - formatter, - "external agent session was not detected for import: {}", - path.display() - ) - } - } - } -} - -impl std::error::Error for PrepareSessionImportsError {} - -pub fn prepare_pending_session_imports( - codex_home: &Path, - requested_sessions: Vec, - detected_sessions: Vec, -) -> Result, PrepareSessionImportsError> { - let detected_session_paths = detected_sessions - .into_iter() - .map(|session| session.path) - .collect::>(); - let mut pending_session_imports = Vec::new(); - for session in requested_sessions { - let has_been_imported = match has_current_session_been_imported(codex_home, &session.path) { - Ok(has_been_imported) => has_been_imported, - Err(_) => continue, - }; - if !detected_session_paths.contains(&session.path) && !has_been_imported { - return Err(PrepareSessionImportsError::SessionNotDetected(session.path)); - } - if has_been_imported { - continue; - } - let imported_session = match load_importable_session(&session.path) { - Ok(Some(imported_session)) => imported_session, - Ok(None) | Err(_) => continue, - }; - pending_session_imports.push(PendingSessionImport { - source_path: session.path, - session: imported_session, - }); - } - Ok(pending_session_imports) -} - -pub fn prepare_validated_session_imports( - codex_home: &Path, - requested_sessions: Vec, -) -> Vec { - requested_sessions - .into_iter() - .filter_map(|session| pending_session_import(codex_home, session)) - .collect() -} - -fn pending_session_import( +pub fn prepare_validated_session_import( codex_home: &Path, session: ExternalAgentSessionMigration, -) -> Option { - let has_been_imported = match has_current_session_been_imported(codex_home, &session.path) { - Ok(has_been_imported) => has_been_imported, - Err(_) => return None, - }; +) -> io::Result> { + let has_been_imported = has_current_session_been_imported(codex_home, &session.path)?; if has_been_imported { - return None; + return Ok(None); } - let imported_session = match load_importable_session(&session.path) { - Ok(Some(imported_session)) => imported_session, - Ok(None) | Err(_) => return None, - }; - Some(PendingSessionImport { - source_path: session.path, - session: imported_session, - }) -} - -fn load_importable_session(path: &Path) -> io::Result> { - let Some(imported_session) = load_session_for_import(path)? else { + let Some((source_path, imported_session, source_content_sha256)) = + load_importable_session(&session.path)? + else { return Ok(None); }; - Ok(imported_session.cwd.is_dir().then_some(imported_session)) + Ok(Some(PendingSessionImport { + source_path, + source_content_sha256, + session: imported_session, + })) +} + +fn load_importable_session( + path: &Path, +) -> io::Result> { + let source_path = std::fs::canonicalize(path)?; + let Some((imported_session, source_content_sha256)) = + load_session_for_import_with_content_sha256(&source_path)? + else { + return Ok(None); + }; + Ok(imported_session.cwd.is_dir().then_some(( + source_path, + imported_session, + source_content_sha256, + ))) } #[derive(Debug, Clone)] @@ -172,45 +118,59 @@ fn now_unix_seconds() -> i64 { mod tests { use super::*; use codex_protocol::ThreadId; + use sha2::Digest; + use sha2::Sha256; use tempfile::TempDir; - #[test] - fn rejects_session_that_was_not_detected() { - let root = TempDir::new().expect("tempdir"); - let codex_home = root.path().join("codex-home"); - let source_path = root.path().join("session.jsonl"); - std::fs::write(&source_path, "{}\n").expect("session"); - - let err = prepare_pending_session_imports( - &codex_home, - vec![session_migration(&source_path)], - Vec::new(), - ) - .expect_err("undetected session should be rejected"); - - match err { - PrepareSessionImportsError::SessionNotDetected(path) => { - assert_eq!(path, source_path); - } - } - } - #[test] fn skips_session_that_was_already_imported() { let root = TempDir::new().expect("tempdir"); let codex_home = root.path().join("codex-home"); let source_path = root.path().join("session.jsonl"); std::fs::write(&source_path, "{}\n").expect("session"); - record_imported_session(&codex_home, &source_path, ThreadId::new()).expect("record import"); + ledger::record_imported_session(&codex_home, &source_path, ThreadId::new()) + .expect("record import"); - let pending = prepare_pending_session_imports( - &codex_home, - vec![session_migration(&source_path)], - Vec::new(), - ) - .expect("already imported session should be skipped"); + let pending = + prepare_validated_session_import(&codex_home, session_migration(&source_path)) + .expect("already imported session should be skipped"); - assert!(pending.is_empty()); + assert!(pending.is_none()); + } + + #[test] + fn reports_session_preparation_errors() { + let root = TempDir::new().expect("tempdir"); + let source_path = root.path().join("missing-session.jsonl"); + + let err = prepare_validated_session_import(root.path(), session_migration(&source_path)) + .expect_err("missing session should fail preparation"); + + assert_eq!(err.kind(), io::ErrorKind::NotFound); + } + + #[test] + fn prepares_one_validated_session_import_with_content_hash() { + let root = TempDir::new().expect("tempdir"); + let source_path = root.path().join("session.jsonl"); + let contents = serde_json::json!({ + "type": "user", + "cwd": root.path(), + "timestamp": "2026-06-03T12:00:00Z", + "message": { "content": "first request" }, + }) + .to_string(); + std::fs::write(&source_path, &contents).expect("session"); + + let pending = + prepare_validated_session_import(root.path(), session_migration(&source_path)) + .expect("prepare session") + .expect("pending import"); + + assert_eq!( + pending.source_content_sha256, + format!("{:x}", Sha256::digest(contents)) + ); } fn session_migration(path: &Path) -> ExternalAgentSessionMigration { diff --git a/codex-rs/external-agent-sessions/src/records.rs b/codex-rs/external-agent-sessions/src/records.rs index 52f053545..00307fa1d 100644 --- a/codex-rs/external-agent-sessions/src/records.rs +++ b/codex-rs/external-agent-sessions/src/records.rs @@ -4,6 +4,8 @@ use crate::MessageRole; use crate::summarize_for_label; use crate::truncate; use serde_json::Value as JsonValue; +use sha2::Digest; +use sha2::Sha256; use std::fs::File; use std::io; use std::io::BufRead; @@ -21,6 +23,13 @@ pub struct SessionSummary { pub migration: ExternalAgentSessionMigration, } +pub(super) struct ParsedSessionImport { + pub cwd: Option, + pub source_title: Option, + pub messages: Vec, + pub content_sha256: String, +} + pub fn summarize_session(path: &Path) -> io::Result> { let file = File::open(path)?; let reader = BufReader::new(file); @@ -37,7 +46,7 @@ pub fn summarize_session(path: &Path) -> io::Result> { if trimmed.is_empty() { continue; } - let Ok(record) = serde_json::from_str::(trimmed) else { + let Ok(mut record) = serde_json::from_str::(trimmed) else { continue; }; if cwd.is_none() { @@ -52,7 +61,7 @@ pub fn summarize_session(path: &Path) -> io::Result> { if let Some(title) = ai_title_from_record(&record) { ai_title = Some(title.to_string()); } - let Some(message) = conversation_message_from_record(&record) else { + let Some(message) = conversation_message_from_owned_record(&mut record) else { continue; }; saw_message = true; @@ -84,54 +93,50 @@ pub fn summarize_session(path: &Path) -> io::Result> { })) } -pub(super) fn source_title_from_records(records: &[JsonValue]) -> Option { - latest_title_from_records(records, custom_title_from_record) - .or_else(|| latest_title_from_records(records, ai_title_from_record)) -} - -pub(super) fn read_records(path: &Path) -> io::Result> { +pub(super) fn read_session_import(path: &Path) -> io::Result { let file = File::open(path)?; - let reader = BufReader::new(file); - let mut records = Vec::new(); - for line in reader.lines() { - let line = line?; + let mut reader = BufReader::new(file); + let mut cwd = None; + let mut custom_title = None; + let mut ai_title = None; + let mut messages = Vec::new(); + let mut line = String::new(); + let mut hasher = Sha256::new(); + loop { + line.clear(); + if reader.read_line(&mut line)? == 0 { + break; + } + hasher.update(line.as_bytes()); let trimmed = line.trim(); if trimmed.is_empty() { continue; } - let Ok(value) = serde_json::from_str::(trimmed) else { + let Ok(mut record) = serde_json::from_str::(trimmed) else { continue; }; - if value.is_object() { - records.push(value); + if cwd.is_none() { + cwd = record + .get("cwd") + .and_then(JsonValue::as_str) + .map(PathBuf::from); + } + if let Some(title) = custom_title_from_record(&record) { + custom_title = Some(title.to_string()); + } + if let Some(title) = ai_title_from_record(&record) { + ai_title = Some(title.to_string()); + } + if let Some(message) = conversation_message_from_owned_record(&mut record) { + messages.push(message); } } - Ok(records) -} - -pub(super) fn project_root_from_records(records: &[JsonValue]) -> Option { - records - .iter() - .find_map(|record| record.get("cwd").and_then(JsonValue::as_str)) - .map(PathBuf::from) -} - -pub(super) fn conversation_messages(records: &[JsonValue]) -> Vec { - records - .iter() - .filter_map(conversation_message_from_record) - .collect() -} - -fn latest_title_from_records<'a>( - records: &'a [JsonValue], - title_from_record: impl Fn(&'a JsonValue) -> Option<&'a str>, -) -> Option { - records - .iter() - .filter_map(title_from_record) - .next_back() - .map(ToOwned::to_owned) + Ok(ParsedSessionImport { + cwd, + source_title: custom_title.or(ai_title), + messages, + content_sha256: format!("{:x}", hasher.finalize()), + }) } fn custom_title_from_record(record: &JsonValue) -> Option<&str> { @@ -150,7 +155,7 @@ fn title_from_record<'a>(record: &'a JsonValue, record_type: &str, field: &str) .filter(|title| !title.is_empty()) } -fn conversation_message_from_record(record: &JsonValue) -> Option { +fn conversation_message_from_owned_record(record: &mut JsonValue) -> Option { let record_type = record.get("type")?.as_str()?; if record_type != "assistant" && record_type != "user" { return None; @@ -161,18 +166,30 @@ fn conversation_message_from_record(record: &JsonValue) -> Option { + if text.trim().is_empty() { + return None; + } + ExtractedMessage { + text, + only_tool_result: false, + } + } + content => extract_message_text(&content)?, + }; Some(ConversationMessage { - role, + role: if is_assistant || extracted.only_tool_result { + MessageRole::Assistant + } else { + MessageRole::User + }, text: extracted.text, timestamp, }) @@ -324,6 +341,46 @@ fn parse_timestamp(timestamp: &str) -> Option { #[cfg(test)] mod tests { use super::*; + use tempfile::TempDir; + + #[test] + fn reads_session_import_in_one_pass() { + let root = TempDir::new().expect("tempdir"); + let path = root.path().join("session.jsonl"); + let contents = [ + serde_json::json!({ + "type": "user", + "cwd": root.path(), + "timestamp": "2026-06-03T12:00:00Z", + "message": { "content": "first request" }, + }) + .to_string(), + "not json".to_string(), + serde_json::json!({ + "type": "ai-title", + "aiTitle": "generated title", + }) + .to_string(), + serde_json::json!({ + "type": "custom-title", + "customTitle": "custom title", + }) + .to_string(), + ] + .join("\n"); + std::fs::write(&path, &contents).expect("session"); + + let parsed = read_session_import(&path).expect("parse session"); + + assert_eq!(parsed.cwd.as_deref(), Some(root.path())); + assert_eq!(parsed.source_title.as_deref(), Some("custom title")); + assert_eq!(parsed.messages.len(), 1); + assert_eq!(parsed.messages[0].text, "first request"); + assert_eq!( + parsed.content_sha256, + format!("{:x}", Sha256::digest(contents)) + ); + } #[test] fn converts_tool_use_blocks_to_bounded_external_agent_tags() {