From 370b02ac96fa29f2d46e8ca61f299764819f6575 Mon Sep 17 00:00:00 2001 From: stefanstokic-oai Date: Fri, 31 Jul 2026 16:10:04 +0000 Subject: [PATCH] Sync updates to imported external agent sessions (#36356) ## Why External agent session files can gain messages after their initial import. Re-importing those files should extend the existing Codex thread instead of creating a duplicate. ## What changed - Map a changed source session back to its uniquely imported thread and append only the missing transcript suffix. - Update the import ledger after verifying that the source and destination transcripts match. - Defer the update when the target is active, archived, ambiguous, diverged, or otherwise unsafe to modify. ## Testing Added unit and app-server coverage for suffix planning, ledger checkpointing, concurrent updates, and unsafe targets that must be deferred. GitOrigin-RevId: 3d9e71cd66e8b31cf5128e8869063868bfb3eb05 --- codex-rs/Cargo.lock | 1 + .../session_importer.rs | 105 ++++- .../suite/v2/external_agent_import_sync.rs | 377 ++++++++++++++++++ codex-rs/app-server/tests/suite/v2/mod.rs | 1 + codex-rs/external-agent-migration/Cargo.toml | 2 + .../src/sessions/append.rs | 316 +++++++++++++++ .../src/sessions/append_tests.rs | 122 ++++++ .../src/sessions/export.rs | 4 +- .../src/sessions/ledger.rs | 75 ++++ .../src/sessions/ledger_tests.rs | 76 ++++ .../src/sessions/mod.rs | 110 ++++- 11 files changed, 1170 insertions(+), 19 deletions(-) create mode 100644 codex-rs/app-server/tests/suite/v2/external_agent_import_sync.rs create mode 100644 codex-rs/external-agent-migration/src/sessions/append.rs create mode 100644 codex-rs/external-agent-migration/src/sessions/append_tests.rs diff --git a/codex-rs/Cargo.lock b/codex-rs/Cargo.lock index 25875f5492..c77e37f013 100644 --- a/codex-rs/Cargo.lock +++ b/codex-rs/Cargo.lock @@ -3067,6 +3067,7 @@ dependencies = [ "codex-plugin", "codex-protocol", "codex-rollout", + "codex-thread-store", "codex-utils-output-truncation", "pretty_assertions", "serde", diff --git a/codex-rs/app-server/src/external_agent_migration/session_importer.rs b/codex-rs/app-server/src/external_agent_migration/session_importer.rs index 3f70821e57..3825d8e5da 100644 --- a/codex-rs/app-server/src/external_agent_migration/session_importer.rs +++ b/codex-rs/app-server/src/external_agent_migration/session_importer.rs @@ -1,3 +1,4 @@ +use std::collections::BTreeSet; use std::path::PathBuf; use std::sync::Arc; @@ -9,11 +10,14 @@ use codex_core::config::ConfigOverrides; use codex_external_agent_migration::ExternalAgentConfigImportItemResult; use codex_external_agent_migration::record_import_error; use codex_external_agent_migration::sessions::CompletedExternalAgentSessionImport; +use codex_external_agent_migration::sessions::ExistingSessionAppend; use codex_external_agent_migration::sessions::ExternalAgentSessionMigration; use codex_external_agent_migration::sessions::ImportedExternalAgentSession; use codex_external_agent_migration::sessions::ImportedSessionConnectorAttribution; use codex_external_agent_migration::sessions::PendingSessionImport; +use codex_external_agent_migration::sessions::SessionImportTarget; use codex_external_agent_migration::sessions::SessionMetadataMode; +use codex_external_agent_migration::sessions::append_existing_session; use codex_external_agent_migration::sessions::detect_imported_cla_session_connectors; use codex_external_agent_migration::sessions::prepare_validated_session_import_with_metadata_mode; use codex_external_agent_migration::sessions::record_completed_session_imports; @@ -44,11 +48,21 @@ struct CompletedSessionImport { connector_attribution: Option, } +enum SessionImportOutcome { + Created(CompletedSessionImport), + Appended { + source_path: PathBuf, + imported_thread_id: ThreadId, + title: Option, + }, +} + #[derive(Clone)] pub(super) struct ExternalAgentSessionImporter { codex_home: PathBuf, connector_metadata_roots: Vec, permits: Arc, + append_checkpoint_permits: Arc, thread_manager: Arc, thread_store: Arc, config_manager: ConfigManager, @@ -68,6 +82,7 @@ impl ExternalAgentSessionImporter { codex_home, connector_metadata_roots, permits: Arc::new(Semaphore::new(1)), + append_checkpoint_permits: Arc::new(Semaphore::new(1)), thread_manager, thread_store, config_manager, @@ -109,7 +124,7 @@ impl ExternalAgentSessionImporter { let mut completed_imports = Vec::new(); while let Some(result) = import_results.next().await { match result { - Ok(Some(completed_import)) => { + Ok(Some(SessionImportOutcome::Created(completed_import))) => { item_result.record_success( Some(completed_import.import.source_path.display().to_string()), Some(completed_import.import.imported_thread_id.to_string()), @@ -117,6 +132,17 @@ impl ExternalAgentSessionImporter { ); completed_imports.push(completed_import); } + Ok(Some(SessionImportOutcome::Appended { + source_path, + imported_thread_id, + title, + })) => { + item_result.record_success( + Some(source_path.display().to_string()), + Some(imported_thread_id.to_string()), + title, + ); + } Ok(None) => {} Err(failure) => { let SessionImportFailure { @@ -135,6 +161,9 @@ impl ExternalAgentSessionImporter { } } } + if completed_imports.is_empty() { + return item_result; + } let connector_attributions = completed_imports .iter() .filter_map(|completed_import| completed_import.connector_attribution.clone()) @@ -188,7 +217,7 @@ impl ExternalAgentSessionImporter { &self, session: ExternalAgentSessionMigration, metadata_mode: SessionMetadataMode, - ) -> Result, SessionImportFailure> { + ) -> Result, SessionImportFailure> { let source_path = session.path.clone(); let Some(pending_import) = self .prepare_session_import(session, metadata_mode) @@ -202,36 +231,88 @@ impl ExternalAgentSessionImporter { else { return Ok(None); }; - let connector_attribution = pending_import - .source_path + let PendingSessionImport { + source_path, + source_content_sha256, + target, + attributed_mcp_server_ids, + session, + } = pending_import; + match target { + SessionImportTarget::New => self + .create_session_import( + source_path, + source_content_sha256, + attributed_mcp_server_ids, + session, + ) + .await + .map(SessionImportOutcome::Created) + .map(Some), + SessionImportTarget::Existing { + thread_id, + expected_source_content_sha256, + } => { + let title = session.title.clone(); + let appended = append_existing_session( + &self.codex_home, + self.append_checkpoint_permits.as_ref(), + self.thread_manager.as_ref(), + self.thread_store.as_ref(), + ExistingSessionAppend { + source_path: &source_path, + source_content_sha256: &source_content_sha256, + expected_source_content_sha256: &expected_source_content_sha256, + thread_id, + source_items: &session.rollout_items, + }, + ) + .await; + Ok(appended.then_some(SessionImportOutcome::Appended { + source_path, + imported_thread_id: thread_id, + title, + })) + } + } + } + + async fn create_session_import( + &self, + source_path: PathBuf, + source_content_sha256: String, + attributed_mcp_server_ids: BTreeSet, + session: ImportedExternalAgentSession, + ) -> Result { + let connector_attribution = source_path .file_stem() .and_then(|stem| stem.to_str()) .map(str::trim) .filter(|session_id| !session_id.is_empty()) .map(|session_id| ImportedSessionConnectorAttribution { session_id: session_id.to_string(), - server_ids: pending_import.attributed_mcp_server_ids, + server_ids: attributed_mcp_server_ids, }); - let title = pending_import.session.title.clone(); + let title = session.title.clone(); let imported_thread_id = - self.persist_session(pending_import.session) + self.persist_session(session) .await .map_err(|failure| SessionImportFailure { - source_path: pending_import.source_path.clone(), + source_path: source_path.clone(), message: failure.message, stage: "session_persist", sub_error_type: failure.sub_error_type, })?; - Ok(Some(CompletedSessionImport { + Ok(CompletedSessionImport { import: CompletedExternalAgentSessionImport { - source_path: pending_import.source_path, - source_content_sha256: pending_import.source_content_sha256, + source_path, + source_content_sha256, imported_thread_id, connector_names: Vec::new(), title, }, connector_attribution, - })) + }) } async fn prepare_session_import( diff --git a/codex-rs/app-server/tests/suite/v2/external_agent_import_sync.rs b/codex-rs/app-server/tests/suite/v2/external_agent_import_sync.rs new file mode 100644 index 0000000000..4142127303 --- /dev/null +++ b/codex-rs/app-server/tests/suite/v2/external_agent_import_sync.rs @@ -0,0 +1,377 @@ +use std::fs::FileTimes; +use std::fs::OpenOptions; +use std::io::Write; +use std::path::Path; +use std::path::PathBuf; +use std::time::Duration; + +use anyhow::Context; +use anyhow::Result; +use app_test_support::MockResponsesConfig; +use app_test_support::TestAppServer; +use app_test_support::create_mock_responses_server_repeating_assistant; +use codex_app_server_protocol::ExternalAgentConfigDetectResponse; +use codex_app_server_protocol::ExternalAgentConfigImportCompletedNotification; +use codex_app_server_protocol::ExternalAgentConfigImportResponse; +use codex_app_server_protocol::ExternalAgentConfigMigrationItemType; +use codex_app_server_protocol::ThreadItem; +use codex_app_server_protocol::ThreadListResponse; +use codex_app_server_protocol::ThreadReadResponse; +use codex_app_server_protocol::TurnStartParams; +use codex_app_server_protocol::UserInput; +use pretty_assertions::assert_eq; +use serde_json::Value; +use serde_json::json; +use tempfile::TempDir; +use tokio::time::timeout; + +const TIMEOUT: Duration = Duration::from_secs(60); +const IMPORT_MARKER: &str = ""; +const FIRST_USER: &str = "original external message"; +const LATE_ASSISTANT: &str = "late external assistant reply"; +const SECOND_USER: &str = "second original external message"; +const SECOND_LATE_ASSISTANT: &str = "second late external assistant reply"; +const LATER_USER: &str = "later external user message"; +const NATIVE_USER: &str = "native Codex message"; +const NATIVE_ASSISTANT: &str = "native Codex answer"; + +fn source_record(cwd: &Path, role: &str, text: &str) -> Value { + json!({ + "timestamp": chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true), + "type": role, + "cwd": cwd, + "message": { "content": text }, + }) +} + +struct ImportFixture { + _codex_home: TempDir, + project_root: PathBuf, + session_path: PathBuf, + app_server: TestAppServer, +} + +impl ImportFixture { + async fn new() -> Result { + let codex_home = TempDir::new()?; + let project_root = codex_home.path().join("repo"); + std::fs::create_dir_all(&project_root)?; + let session_path = codex_home + .path() + .join(concat!(".", "cla", "ude")) + .join("projects/repo/session.jsonl"); + std::fs::create_dir_all(session_path.parent().context("session parent")?)?; + let initial_record = source_record(&project_root, "user", FIRST_USER); + std::fs::write(&session_path, format!("{initial_record}\n"))?; + + let app_server = start_app_server(codex_home.path()).await?; + Ok(Self { + _codex_home: codex_home, + project_root, + session_path, + app_server, + }) + } + + fn ledger_path(&self) -> PathBuf { + self._codex_home + .path() + .join("external_agent_session_imports.json") + } + + fn append(&self, records: &[(&str, &str)]) -> Result<()> { + self.append_to(&self.session_path, records) + } + + fn append_to(&self, session_path: &Path, records: &[(&str, &str)]) -> Result<()> { + let raw = records + .iter() + .map(|(role, text)| source_record(&self.project_root, role, text).to_string()) + .collect::>() + .join("\n"); + self.append_raw_to(session_path, &format!("{raw}\n")) + } + + fn append_raw(&self, raw: &str) -> Result<()> { + self.append_raw_to(&self.session_path, raw) + } + + fn append_raw_to(&self, session_path: &Path, raw: &str) -> Result<()> { + let modified_at = std::fs::metadata(session_path)?.modified()?; + let mut session = OpenOptions::new().append(true).open(session_path)?; + session.write_all(raw.as_bytes())?; + session.set_times(FileTimes::new().set_modified(modified_at + Duration::from_secs(1)))?; + Ok(()) + } + + async fn rpc( + &mut self, + method: &str, + params: Value, + ) -> Result { + let request_id = self + .app_server + .send_raw_request(method, Some(params)) + .await?; + timeout(TIMEOUT, self.app_server.read_response(request_id)).await? + } + + async fn detect(&mut self) -> Result { + self.rpc("externalAgentConfig/detect", json!({ "includeHome": true })) + .await + } + + async fn import_detected(&mut self) -> Result> { + let detected = self.detect().await?; + assert!( + !detected.items.is_empty(), + "changed session must be detected" + ); + let imported: ExternalAgentConfigImportResponse = self + .rpc( + "externalAgentConfig/import", + json!({ "migrationItems": detected.items }), + ) + .await?; + let completed: ExternalAgentConfigImportCompletedNotification = timeout( + TIMEOUT, + self.app_server + .read_notification("externalAgentConfig/import/completed"), + ) + .await??; + assert_eq!(completed.import_id, imported.import_id); + let sessions = completed + .item_type_results + .iter() + .find(|result| result.item_type == ExternalAgentConfigMigrationItemType::Sessions) + .context("sessions import result")?; + assert!( + sessions.failures.is_empty(), + "session import must not report failures: {:?}", + sessions.failures + ); + sessions + .successes + .iter() + .map(|success| success.target.clone().context("imported session target")) + .collect() + } + + async fn import_one(&mut self) -> Result { + let targets = self.import_detected().await?; + let [target] = targets.as_slice() else { + anyhow::bail!("session import must produce exactly one task"); + }; + Ok(target.clone()) + } + + async fn read(&mut self, thread_id: &str) -> Result { + self.rpc( + "thread/read", + json!({ "threadId": thread_id, "includeTurns": true }), + ) + .await + } + + async fn thread_count(&mut self) -> Result { + let response: ThreadListResponse = self.rpc("thread/list", json!({})).await?; + Ok(response.data.len()) + } + + async fn resume(&mut self, thread_id: &str) -> Result<()> { + let _: Value = self + .rpc("thread/resume", json!({ "threadId": thread_id })) + .await?; + Ok(()) + } + + async fn restart(&mut self) -> Result<()> { + timeout(TIMEOUT, self.app_server.shutdown_gracefully()).await??; + self.app_server = start_app_server(self._codex_home.path()).await?; + Ok(()) + } + + async fn assert_deferred( + &mut self, + ledger_before: &[u8], + threads: &[(&str, &[(&str, &str)])], + ) -> Result<()> { + assert!(self.import_detected().await?.is_empty()); + assert_eq!(self.thread_count().await?, threads.len()); + for &(thread_id, expected_history) in threads { + assert_history(&self.read(thread_id).await?, expected_history); + } + assert_eq!(std::fs::read(self.ledger_path())?, ledger_before); + assert!(!self.detect().await?.items.is_empty()); + Ok(()) + } + + async fn assert_one_deferred( + &mut self, + ledger_before: &[u8], + thread_id: &str, + expected_history: &[(&str, &str)], + ) -> Result<()> { + self.assert_deferred(ledger_before, &[(thread_id, expected_history)]) + .await + } +} + +async fn start_app_server(codex_home: &Path) -> Result { + let home = codex_home.display().to_string(); + TestAppServer::builder() + .with_codex_home(codex_home) + .with_env_overrides(&[("HOME", Some(home.as_str()))]) + .build_initialized_with_timeout(TIMEOUT) + .await +} + +fn assert_history(response: &ThreadReadResponse, expected: &[(&str, &str)]) { + let mut actual = Vec::new(); + let mut markers = 0; + for item in response.thread.turns.iter().flat_map(|turn| &turn.items) { + match item { + ThreadItem::UserMessage { content, .. } => { + actual.extend(content.iter().filter_map(|input| match input { + UserInput::Text { text, .. } => Some(("user", text.as_str())), + _ => None, + })); + } + ThreadItem::AgentMessage { text, .. } if text == IMPORT_MARKER => markers += 1, + ThreadItem::AgentMessage { text, .. } => actual.push(("assistant", text.as_str())), + _ => {} + } + } + assert_eq!(actual.as_slice(), expected); + assert_eq!(markers, 1, "the import marker must appear exactly once"); +} + +#[tokio::test] +async fn cold_exact_prefix_appends_suffix_to_same_task_and_checkpoints() -> Result<()> { + let mut fixture = ImportFixture::new().await?; + let original = fixture.import_one().await?; + + fixture.append(&[("assistant", LATE_ASSISTANT)])?; + assert_eq!(fixture.import_detected().await?, vec![original.clone()]); + + assert_eq!(fixture.thread_count().await?, 1); + assert_history( + &fixture.read(&original).await?, + &[("user", FIRST_USER), ("assistant", LATE_ASSISTANT)], + ); + assert!(fixture.detect().await?.items.is_empty()); + Ok(()) +} + +#[tokio::test] +async fn concurrent_changed_sessions_checkpoint_without_lost_update() -> Result<()> { + let mut fixture = ImportFixture::new().await?; + let second_session_path = fixture.session_path.with_file_name("second.jsonl"); + let second_initial = source_record(&fixture.project_root, "user", SECOND_USER); + std::fs::write(&second_session_path, format!("{second_initial}\n"))?; + + let mut original_threads = fixture.import_detected().await?; + assert_eq!(original_threads.len(), 2); + original_threads.sort(); + + fixture.append(&[("assistant", LATE_ASSISTANT)])?; + fixture.append_to( + &second_session_path, + &[("assistant", SECOND_LATE_ASSISTANT)], + )?; + let mut appended_threads = fixture.import_detected().await?; + appended_threads.sort(); + + assert_eq!(appended_threads, original_threads); + assert_eq!(fixture.thread_count().await?, 2); + assert!(fixture.detect().await?.items.is_empty()); + Ok(()) +} + +#[tokio::test] +async fn active_target_is_deferred_without_checkpoint() -> Result<()> { + let mut fixture = ImportFixture::new().await?; + let original = fixture.import_one().await?; + let ledger_before = std::fs::read(fixture.ledger_path())?; + fixture.resume(&original).await?; + fixture.append(&[("user", LATER_USER)])?; + + fixture + .assert_one_deferred(&ledger_before, &original, &[("user", FIRST_USER)]) + .await +} + +#[tokio::test] +async fn cold_diverged_target_is_deferred_without_checkpoint() -> Result<()> { + let model = create_mock_responses_server_repeating_assistant(NATIVE_ASSISTANT).await; + let mut fixture = ImportFixture::new().await?; + MockResponsesConfig::new(&model.uri()).write(fixture._codex_home.path())?; + fixture.restart().await?; + let original = fixture.import_one().await?; + fixture.resume(&original).await?; + let environment = fixture.app_server.auto_env_params()?; + timeout( + TIMEOUT, + fixture + .app_server + .start_turn_and_wait_for_completion(TurnStartParams { + thread_id: original.clone(), + input: vec![UserInput::Text { + text: NATIVE_USER.to_string(), + text_elements: Vec::new(), + }], + environments: Some(vec![environment]), + ..Default::default() + }), + ) + .await??; + + fixture.restart().await?; + let ledger_before = std::fs::read(fixture.ledger_path())?; + fixture.append(&[("user", LATER_USER)])?; + fixture + .assert_one_deferred( + &ledger_before, + &original, + &[ + ("user", FIRST_USER), + ("user", NATIVE_USER), + ("assistant", NATIVE_ASSISTANT), + ], + ) + .await +} + +#[tokio::test] +async fn changed_hash_with_equal_transcript_is_deferred_without_checkpoint() -> Result<()> { + let mut fixture = ImportFixture::new().await?; + let original = fixture.import_one().await?; + let ledger_before = std::fs::read(fixture.ledger_path())?; + fixture.append_raw("\n")?; + + fixture + .assert_one_deferred(&ledger_before, &original, &[("user", FIRST_USER)]) + .await +} + +#[tokio::test] +async fn ambiguous_legacy_targets_are_deferred_without_checkpoint() -> Result<()> { + let mut fixture = ImportFixture::new().await?; + let original = fixture.import_one().await?; + let mut ambiguous_ledger: Value = + serde_json::from_slice(&std::fs::read(fixture.ledger_path())?)?; + let mut second_target = ambiguous_ledger["records"][0].clone(); + second_target["imported_thread_id"] = json!("01800000-0001-7000-8000-000000000001"); + ambiguous_ledger["records"] + .as_array_mut() + .context("ledger records")? + .push(second_target); + let ambiguous_ledger = serde_json::to_vec(&ambiguous_ledger)?; + std::fs::write(fixture.ledger_path(), ambiguous_ledger)?; + + fixture.append(&[("user", LATER_USER)])?; + let ledger_before = std::fs::read(fixture.ledger_path())?; + fixture + .assert_one_deferred(&ledger_before, &original, &[("user", FIRST_USER)]) + .await +} diff --git a/codex-rs/app-server/tests/suite/v2/mod.rs b/codex-rs/app-server/tests/suite/v2/mod.rs index 767e4ca401..505393ea1c 100644 --- a/codex-rs/app-server/tests/suite/v2/mod.rs +++ b/codex-rs/app-server/tests/suite/v2/mod.rs @@ -29,6 +29,7 @@ mod executor_skills; mod experimental_api; mod experimental_feature_list; mod external_agent_config; +mod external_agent_import_sync; mod fs; mod git_attribution; mod hooks_list; diff --git a/codex-rs/external-agent-migration/Cargo.toml b/codex-rs/external-agent-migration/Cargo.toml index e26b94286c..28d20770f6 100644 --- a/codex-rs/external-agent-migration/Cargo.toml +++ b/codex-rs/external-agent-migration/Cargo.toml @@ -24,12 +24,14 @@ codex-otel = { workspace = true } codex-plugin = { workspace = true } codex-protocol = { workspace = true } codex-rollout = { workspace = true } +codex-thread-store = { workspace = true } codex-utils-output-truncation = { workspace = true } serde = { workspace = true, features = ["derive"] } serde_json = { workspace = true } serde_yaml = { workspace = true } sha2 = { workspace = true } toml = { workspace = true } +tokio = { workspace = true, features = ["rt", "sync"] } tracing = { workspace = true } [dev-dependencies] diff --git a/codex-rs/external-agent-migration/src/sessions/append.rs b/codex-rs/external-agent-migration/src/sessions/append.rs new file mode 100644 index 0000000000..4437dd29eb --- /dev/null +++ b/codex-rs/external-agent-migration/src/sessions/append.rs @@ -0,0 +1,316 @@ +use std::path::Path; + +use super::export::EXTERNAL_SESSION_IMPORTED_MARKER; +use super::ledger::checkpoint_existing_session_import; +use codex_core::ThreadManager; +use codex_protocol::ThreadId; +use codex_protocol::models::ResponseItem; +use codex_protocol::protocol::EventMsg; +use codex_protocol::protocol::RolloutItem; +use codex_protocol::protocol::ThreadMemoryMode; +use codex_thread_store::AppendThreadItemsParams; +use codex_thread_store::ReadThreadParams; +use codex_thread_store::ResumeThreadParams; +use codex_thread_store::ThreadPersistenceMetadata; +use codex_thread_store::ThreadStore; +use tokio::sync::Semaphore; + +/// A changed external session and the existing native thread it may extend. +pub struct ExistingSessionAppend<'a> { + pub source_path: &'a Path, + pub source_content_sha256: &'a str, + pub expected_source_content_sha256: &'a str, + pub thread_id: ThreadId, + pub source_items: &'a [RolloutItem], +} + +/// Appends an exact missing suffix to one cold imported thread and checkpoints the source. +/// +/// Any unavailable, active, archived, malformed, or diverged destination fails closed. +pub async fn append_existing_session( + codex_home: &Path, + checkpoint_permits: &Semaphore, + thread_manager: &ThreadManager, + thread_store: &dyn ThreadStore, + request: ExistingSessionAppend<'_>, +) -> bool { + let ExistingSessionAppend { + source_path, + source_content_sha256, + expected_source_content_sha256, + thread_id, + source_items, + } = request; + let Ok(mut thread) = thread_store + .read_thread(ReadThreadParams { + thread_id, + include_archived: true, + include_history: true, + }) + .await + else { + return false; + }; + if thread.thread_id != thread_id || thread.archived_at.is_some() { + return false; + } + let Some(history) = thread.history.take() else { + return false; + }; + if history.thread_id != thread_id || thread_manager.get_thread(thread_id).await.is_ok() { + return false; + } + let Some(metadata) = persistence_metadata(&history.items, thread_id) else { + return false; + }; + let Some(rollout_path) = thread.rollout_path else { + return false; + }; + if plan_append(source_items, &history.items).is_none() { + return false; + } + if thread_store + .resume_thread(ResumeThreadParams { + thread_id, + rollout_path: Some(rollout_path), + history: None, + include_archived: false, + metadata, + }) + .await + .is_err() + { + return false; + } + + let fresh_items = match thread_store + .read_thread(ReadThreadParams { + thread_id, + include_archived: true, + include_history: true, + }) + .await + { + Ok(mut thread) if thread.thread_id == thread_id && thread.archived_at.is_none() => thread + .history + .take() + .filter(|history| history.thread_id == thread_id) + .filter(|history| persistence_metadata(&history.items, thread_id).is_some()) + .and_then(|history| plan_append(source_items, &history.items)), + Ok(_) | Err(_) => None, + }; + let Some(items) = fresh_items else { + let _ = thread_store.discard_thread(thread_id).await; + return false; + }; + if thread_store + .append_items(AppendThreadItemsParams { thread_id, items }) + .await + .is_err() + { + let _ = thread_store.discard_thread(thread_id).await; + return false; + } + if thread_store.shutdown_thread(thread_id).await.is_err() { + let _ = thread_store.discard_thread(thread_id).await; + return false; + } + + let Ok(_checkpoint_permit) = checkpoint_permits.acquire().await else { + return false; + }; + let target_matches_source = match thread_store + .read_thread(ReadThreadParams { + thread_id, + include_archived: true, + include_history: true, + }) + .await + { + Ok(mut thread) if thread.thread_id == thread_id && thread.archived_at.is_none() => thread + .history + .take() + .filter(|history| history.thread_id == thread_id) + .filter(|history| persistence_metadata(&history.items, thread_id).is_some()) + .is_some_and(|history| model_transcripts_match(source_items, &history.items)), + Ok(_) | Err(_) => false, + }; + if !target_matches_source { + return false; + } + + let codex_home = codex_home.to_path_buf(); + let source_path = source_path.to_path_buf(); + let expected_source_content_sha256 = expected_source_content_sha256.to_string(); + let source_content_sha256 = source_content_sha256.to_string(); + matches!( + tokio::task::spawn_blocking(move || { + checkpoint_existing_session_import( + &codex_home, + &source_path, + thread_id, + &expected_source_content_sha256, + &source_content_sha256, + ) + }) + .await, + Ok(Ok(true)) + ) +} + +struct SourceModelItem<'a> { + response_item: &'a ResponseItem, + append_start_index: usize, +} + +/// Returns the nonempty source suffix that can be safely appended to `history_items`. +/// +/// The destination's complete model-visible transcript must be an exact prefix of the +/// source transcript. Rollout metadata does not participate in transcript identity. +fn plan_append( + source_items: &[RolloutItem], + history_items: &[RolloutItem], +) -> Option> { + let source = source_model_items(source_items)?; + let history = history_model_items(history_items)?; + if history.len() >= source.len() + || history + .iter() + .zip(&source) + .any(|(history, source)| *history != source.response_item) + { + return None; + } + + let append_start_index = source.get(history.len())?.append_start_index; + let suffix = source_items[append_start_index..] + .iter() + .filter(|item| !is_import_marker(item)) + .cloned() + .collect::>(); + suffix + .iter() + .any(|item| matches!(item, RolloutItem::ResponseItem(_))) + .then_some(suffix) +} + +/// Returns whether source and destination contain the same complete model-visible transcript. +/// +/// Rollout metadata does not participate in transcript identity. Unsupported history shapes fail +/// closed. +fn model_transcripts_match(source_items: &[RolloutItem], history_items: &[RolloutItem]) -> bool { + let Some(source) = source_model_items(source_items) else { + return false; + }; + let Some(history) = history_model_items(history_items) else { + return false; + }; + history.len() == source.len() + && history + .iter() + .zip(&source) + .all(|(history, source)| *history == source.response_item) +} + +fn source_model_items(items: &[RolloutItem]) -> Option>> { + let mut model_items = Vec::new(); + let mut append_start_index = None; + for (index, item) in items.iter().enumerate() { + match item { + RolloutItem::SessionMeta(_) | RolloutItem::InterAgentCommunicationMetadata { .. } => {} + RolloutItem::ResponseItem(response_item) => { + model_items.push(SourceModelItem { + response_item, + append_start_index: append_start_index.take().unwrap_or(index), + }); + } + RolloutItem::EventMsg(EventMsg::TurnStarted(_)) => { + append_start_index = Some(index); + } + RolloutItem::EventMsg(EventMsg::UserMessage(_)) => { + append_start_index.get_or_insert(index); + } + RolloutItem::EventMsg(EventMsg::AgentMessage(event)) + if event.message != EXTERNAL_SESSION_IMPORTED_MARKER => + { + append_start_index = Some(index); + } + RolloutItem::EventMsg( + EventMsg::ContextCompacted(_) | EventMsg::ThreadRolledBack(_), + ) + | RolloutItem::InterAgentCommunication(_) + | RolloutItem::Compacted(_) + | RolloutItem::TurnContext(_) + | RolloutItem::WorldState(_) => return None, + RolloutItem::EventMsg(_) => {} + } + } + Some(model_items) +} + +fn history_model_items(items: &[RolloutItem]) -> Option> { + let mut model_items = Vec::new(); + for item in items { + match item { + RolloutItem::SessionMeta(_) | RolloutItem::InterAgentCommunicationMetadata { .. } => {} + RolloutItem::ResponseItem(response_item) => model_items.push(response_item), + RolloutItem::EventMsg( + EventMsg::ContextCompacted(_) | EventMsg::ThreadRolledBack(_), + ) + | RolloutItem::InterAgentCommunication(_) + | RolloutItem::Compacted(_) + | RolloutItem::TurnContext(_) + | RolloutItem::WorldState(_) => return None, + RolloutItem::EventMsg(_) => {} + } + } + Some(model_items) +} + +fn is_import_marker(item: &RolloutItem) -> bool { + matches!( + item, + RolloutItem::EventMsg(EventMsg::AgentMessage(event)) + if event.message == EXTERNAL_SESSION_IMPORTED_MARKER + ) +} + +fn persistence_metadata( + history_items: &[RolloutItem], + thread_id: ThreadId, +) -> Option { + let RolloutItem::SessionMeta(first_meta_line) = history_items.first()? else { + return None; + }; + if first_meta_line.meta.id != thread_id { + return None; + } + let meta = history_items.iter().rev().find_map(|item| match item { + RolloutItem::SessionMeta(meta_line) if meta_line.meta.id == thread_id => { + Some(&meta_line.meta) + } + _ => None, + })?; + if meta.cwd.as_os_str().is_empty() + || meta + .model_provider + .as_deref() + .is_none_or(|provider| provider.trim().is_empty()) + { + return None; + } + let memory_mode = match meta.memory_mode.as_deref() { + None | Some("enabled") => ThreadMemoryMode::Enabled, + Some("disabled") => ThreadMemoryMode::Disabled, + Some(_) => return None, + }; + Some(ThreadPersistenceMetadata { + cwd: Some(meta.cwd.clone()), + model_provider: meta.model_provider.clone()?, + memory_mode, + }) +} + +#[cfg(test)] +#[path = "append_tests.rs"] +mod tests; diff --git a/codex-rs/external-agent-migration/src/sessions/append_tests.rs b/codex-rs/external-agent-migration/src/sessions/append_tests.rs new file mode 100644 index 0000000000..40f765f217 --- /dev/null +++ b/codex-rs/external-agent-migration/src/sessions/append_tests.rs @@ -0,0 +1,122 @@ +use super::super::ConversationMessage; +use super::super::MessageRole; +use super::super::export::EXTERNAL_SESSION_IMPORTED_MARKER; +use super::super::export::rollout_items_from_messages; +use super::*; +use codex_protocol::models::ContentItem; +use codex_protocol::models::ResponseItem; +use codex_protocol::protocol::ContextCompactedEvent; +use codex_protocol::protocol::ThreadRolledBackEvent; +use pretty_assertions::assert_eq; + +#[test] +fn returns_the_missing_suffix_from_its_visible_boundary() { + let history = rollout(&[(MessageRole::User, "first request")]); + let source = rollout(&[ + (MessageRole::User, "first request"), + (MessageRole::Assistant, "late answer"), + (MessageRole::User, "follow-up request"), + ]); + + let suffix = plan_append(&source, &history).expect("exact prefix should append"); + + assert!(matches!( + suffix.first(), + Some(RolloutItem::EventMsg(EventMsg::AgentMessage(event))) + if event.message == "late answer" + )); + assert_eq!( + model_messages(&suffix), + vec![ + (MessageRole::Assistant, "late answer"), + (MessageRole::User, "follow-up request"), + ] + ); + assert!(!suffix.iter().any(|item| matches!( + item, + RolloutItem::EventMsg(EventMsg::AgentMessage(event)) + if event.message == EXTERNAL_SESSION_IMPORTED_MARKER + ))); +} + +#[test] +fn requires_a_strict_nonempty_model_prefix() { + let history = rollout(&[(MessageRole::User, "first request")]); + let source = rollout(&[ + (MessageRole::User, "first request"), + (MessageRole::User, "follow-up request"), + ]); + assert!(plan_append(&history, &history).is_none()); + assert!( + plan_append( + &source, + &rollout(&[(MessageRole::User, "rewritten request")]) + ) + .is_none() + ); + + let mut metadata_changed = history.clone(); + for item in &mut metadata_changed { + if let RolloutItem::EventMsg(EventMsg::TurnStarted(event)) = item { + event.turn_id = "different-turn-id".to_string(); + event.started_at = Some(9_999); + } + } + assert!(model_transcripts_match(&history, &metadata_changed)); + assert!(!model_transcripts_match(&source, &history)); + assert!(plan_append(&source, &metadata_changed).is_some()); + + for event in [ + EventMsg::ContextCompacted(ContextCompactedEvent), + EventMsg::ThreadRolledBack(ThreadRolledBackEvent { num_turns: 1 }), + ] { + let mut rewritten = history.clone(); + rewritten.push(RolloutItem::EventMsg(event)); + assert!(plan_append(&source, &rewritten).is_none()); + } + let mut with_tool_call = history; + with_tool_call.push(RolloutItem::ResponseItem(ResponseItem::FunctionCall { + id: None, + name: "native_tool".to_string(), + namespace: None, + arguments: "{}".to_string(), + encrypted_function_args: None, + call_id: "native-call".to_string(), + internal_chat_message_metadata_passthrough: None, + })); + assert!(plan_append(&source, &with_tool_call).is_none()); +} + +fn rollout(messages: &[(MessageRole, &str)]) -> Vec { + rollout_items_from_messages( + messages + .iter() + .enumerate() + .map(|(index, &(role, text))| ConversationMessage { + role, + text: text.to_string(), + timestamp: Some(index as i64), + }) + .collect(), + ) +} + +fn model_messages(items: &[RolloutItem]) -> Vec<(MessageRole, &str)> { + items + .iter() + .filter_map(|item| match item { + RolloutItem::ResponseItem(ResponseItem::Message { role, content, .. }) => { + match (role.as_str(), content.as_slice()) { + ("user", [ContentItem::InputText { text }]) => { + Some((MessageRole::User, text.as_str())) + } + ("assistant", [ContentItem::OutputText { text }]) => { + Some((MessageRole::Assistant, text.as_str())) + } + _ => None, + } + } + _ => None, + }) + .collect() +} diff --git a/codex-rs/external-agent-migration/src/sessions/export.rs b/codex-rs/external-agent-migration/src/sessions/export.rs index 89a66e368c..c993664783 100644 --- a/codex-rs/external-agent-migration/src/sessions/export.rs +++ b/codex-rs/external-agent-migration/src/sessions/export.rs @@ -24,7 +24,7 @@ use std::collections::BTreeSet; use std::io; use std::path::Path; -const EXTERNAL_SESSION_IMPORTED_MARKER: &str = ""; +pub(super) const EXTERNAL_SESSION_IMPORTED_MARKER: &str = ""; #[cfg(test)] fn load_session_for_import(path: &Path) -> io::Result> { @@ -82,7 +82,7 @@ pub(crate) fn load_session_for_import_with_content_sha256( ))) } -fn rollout_items_from_messages(messages: Vec) -> Vec { +pub(super) fn rollout_items_from_messages(messages: Vec) -> Vec { let mut items = Vec::new(); let mut current_turn = None; let mut response_item_bytes = 0i64; diff --git a/codex-rs/external-agent-migration/src/sessions/ledger.rs b/codex-rs/external-agent-migration/src/sessions/ledger.rs index 86bcef58c4..9e90cf7a0b 100644 --- a/codex-rs/external-agent-migration/src/sessions/ledger.rs +++ b/codex-rs/external-agent-migration/src/sessions/ledger.rs @@ -44,6 +44,16 @@ pub struct CompletedExternalAgentSessionImport { pub title: Option, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) enum SessionImportSourceMapping { + None, + Unique { + source_content_sha256: String, + imported_thread_id: ThreadId, + }, + Ambiguous, +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct ImportedConnectorCandidate { pub name: String, @@ -82,6 +92,71 @@ pub(crate) fn record_imported_session( ) } +pub(crate) fn find_existing_session_import( + codex_home: &Path, + source_path: &Path, +) -> io::Result { + let source_path = canonical_source_path(source_path)?; + let ledger = load_import_ledger(codex_home)?; + let mut matching_records = ledger + .records + .iter() + .filter(|record| record.source_path == source_path); + let Some(record) = matching_records.next() else { + return Ok(SessionImportSourceMapping::None); + }; + if matching_records.next().is_some() { + return Ok(SessionImportSourceMapping::Ambiguous); + } + Ok(SessionImportSourceMapping::Unique { + source_content_sha256: record.content_sha256.clone(), + imported_thread_id: record.imported_thread_id, + }) +} + +pub(super) fn checkpoint_existing_session_import( + codex_home: &Path, + source_path: &Path, + thread_id: ThreadId, + expected_source_content_sha256: &str, + source_content_sha256: &str, +) -> io::Result { + if expected_source_content_sha256 == source_content_sha256 { + return Ok(false); + } + let source_path = canonical_source_path(source_path)?; + let source_modified_at = session_modified_at(&source_path)?; + if session_content_sha256(&source_path)? != source_content_sha256 + || source_modified_at != session_modified_at(&source_path)? + { + return Ok(false); + } + + let mut ledger = load_import_ledger(codex_home)?; + let mut matching_indices = ledger + .records + .iter() + .enumerate() + .filter_map(|(index, record)| (record.source_path == source_path).then_some(index)); + let Some(index) = matching_indices.next() else { + return Ok(false); + }; + if matching_indices.next().is_some() { + return Ok(false); + } + let record = &mut ledger.records[index]; + if record.imported_thread_id != thread_id + || record.content_sha256 != expected_source_content_sha256 + { + return Ok(false); + } + record.content_sha256 = source_content_sha256.to_string(); + record.imported_at = now_unix_seconds(); + record.source_modified_at = source_modified_at; + save_import_ledger(codex_home, &ledger)?; + Ok(true) +} + pub fn record_completed_session_imports( codex_home: &Path, imports: Vec, diff --git a/codex-rs/external-agent-migration/src/sessions/ledger_tests.rs b/codex-rs/external-agent-migration/src/sessions/ledger_tests.rs index 1ca425a660..1e93c72e7d 100644 --- a/codex-rs/external-agent-migration/src/sessions/ledger_tests.rs +++ b/codex-rs/external-agent-migration/src/sessions/ledger_tests.rs @@ -1,11 +1,14 @@ use super::CompletedExternalAgentSessionImport; use super::ImportedConnectorCandidate; use super::ImportedExternalAgentSessionLedger; +use super::checkpoint_existing_session_import; use super::read_imported_connector_candidates; use super::record_completed_session_imports; use codex_protocol::ThreadId; +use pretty_assertions::assert_eq; use sha2::Digest; use sha2::Sha256; +use std::path::Path; use tempfile::TempDir; #[test] @@ -95,6 +98,75 @@ fn completed_import_refreshes_existing_record_metadata() { assert_eq!(ledger.records[0].title.as_deref(), Some("Second title")); } +#[test] +fn checkpoint_is_compare_and_swap_and_preserves_record_metadata() { + 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, "old").expect("initial source"); + let source_path = std::fs::canonicalize(source_path).expect("canonical source"); + let thread_id = ThreadId::new(); + let old_hash = format!("{:x}", Sha256::digest("old")); + let new_hash = format!("{:x}", Sha256::digest("new")); + record_completed_session_imports( + &codex_home, + vec![CompletedExternalAgentSessionImport { + source_path: source_path.clone(), + source_content_sha256: old_hash.clone(), + imported_thread_id: thread_id, + connector_names: vec!["Gmail".to_string()], + title: Some("Original title".to_string()), + }], + ) + .expect("record initial import"); + std::fs::write(&source_path, "new").expect("updated source"); + let ledger_before = ledger_bytes(&codex_home); + for (candidate_thread_id, expected_hash, candidate_hash) in [ + (ThreadId::new(), old_hash.as_str(), new_hash.as_str()), + (thread_id, "stale", new_hash.as_str()), + (thread_id, old_hash.as_str(), "not-the-source-hash"), + ] { + assert!( + !checkpoint_existing_session_import( + &codex_home, + &source_path, + candidate_thread_id, + expected_hash, + candidate_hash, + ) + .expect("rejected checkpoint") + ); + } + assert_eq!(ledger_bytes(&codex_home), ledger_before); + + assert!( + checkpoint_existing_session_import( + &codex_home, + &source_path, + thread_id, + &old_hash, + &new_hash, + ) + .expect("checkpoint") + ); + let ledger = super::load_import_ledger(&codex_home).expect("updated ledger"); + let [record] = ledger.records.as_slice() else { + panic!("checkpoint must preserve one record"); + }; + assert_eq!( + record, + &super::ImportedExternalAgentSessionRecord { + source_path: source_path.clone(), + content_sha256: new_hash, + imported_thread_id: thread_id, + imported_at: record.imported_at, + source_modified_at: super::session_modified_at(&source_path).expect("source mtime"), + connector_names: vec!["Gmail".to_string()], + title: Some("Original title".to_string()), + } + ); +} + #[test] fn connector_candidates_use_latest_import_for_each_source() { let root = TempDir::new().expect("tempdir"); @@ -144,3 +216,7 @@ fn connector_candidates_use_latest_import_for_each_source() { ] ); } + +fn ledger_bytes(codex_home: &Path) -> Vec { + std::fs::read(super::import_ledger_path(codex_home)).expect("ledger") +} diff --git a/codex-rs/external-agent-migration/src/sessions/mod.rs b/codex-rs/external-agent-migration/src/sessions/mod.rs index 0d449b5f4a..c39fa8ce72 100644 --- a/codex-rs/external-agent-migration/src/sessions/mod.rs +++ b/codex-rs/external-agent-migration/src/sessions/mod.rs @@ -1,5 +1,6 @@ //! Parsing and export helpers for external-agent session histories. +mod append; mod export; pub(crate) mod ledger; pub(crate) mod records_cla; @@ -7,6 +8,7 @@ mod records_common; pub(crate) mod records_cur; mod title; +use codex_protocol::ThreadId; use codex_protocol::protocol::RolloutItem; use std::collections::BTreeSet; use std::io; @@ -17,6 +19,8 @@ pub use crate::detect::sessions::ImportedSessionConnectorAttribution; pub use crate::detect::sessions::detect_imported_cla_session_connectors; pub use crate::detect::sessions::detect_recent_cla_sessions; pub use crate::detect::sessions::detect_recent_cur_sessions; +pub use append::ExistingSessionAppend; +pub use append::append_existing_session; use export::load_session_for_import_with_content_sha256; pub use ledger::CompletedExternalAgentSessionImport; pub use ledger::ImportedConnectorCandidate; @@ -77,10 +81,20 @@ pub struct ImportedExternalAgentSession { pub rollout_items: Vec, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum SessionImportTarget { + New, + Existing { + thread_id: ThreadId, + expected_source_content_sha256: String, + }, +} + #[derive(Debug, Clone)] pub struct PendingSessionImport { pub source_path: PathBuf, pub source_content_sha256: String, + pub target: SessionImportTarget, pub attributed_mcp_server_ids: BTreeSet, pub session: ImportedExternalAgentSession, } @@ -101,11 +115,27 @@ pub fn prepare_validated_session_import_with_metadata_mode( session: ExternalAgentSessionMigration, metadata_mode: SessionMetadataMode, ) -> io::Result> { - let has_been_imported = has_current_session_been_imported(codex_home, &session.path)?; - if has_been_imported { + let Some(mut pending) = load_importable_session(&session.path, &session.cwd, metadata_mode)? + else { return Ok(None); - } - load_importable_session(&session.path, &session.cwd, metadata_mode) + }; + pending.target = match ledger::find_existing_session_import(codex_home, &pending.source_path)? { + ledger::SessionImportSourceMapping::None => SessionImportTarget::New, + ledger::SessionImportSourceMapping::Unique { + source_content_sha256, + imported_thread_id, + } => { + if source_content_sha256 == pending.source_content_sha256 { + return Ok(None); + } + SessionImportTarget::Existing { + thread_id: imported_thread_id, + expected_source_content_sha256: source_content_sha256, + } + } + ledger::SessionImportSourceMapping::Ambiguous => return Ok(None), + }; + Ok(Some(pending)) } fn load_importable_session( @@ -129,6 +159,7 @@ fn load_importable_session( .then_some(PendingSessionImport { source_path, source_content_sha256, + target: SessionImportTarget::New, attributed_mcp_server_ids, session: imported_session, })) @@ -174,6 +205,7 @@ pub(crate) fn now_unix_seconds() -> i64 { mod tests { use super::*; use codex_protocol::ThreadId; + use pretty_assertions::assert_eq; use sha2::Digest; use sha2::Sha256; use tempfile::TempDir; @@ -183,7 +215,8 @@ mod tests { 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"); + std::fs::write(&source_path, session_record(root.path(), "first request")) + .expect("session"); ledger::record_imported_session(&codex_home, &source_path, ThreadId::new()) .expect("record import"); @@ -194,6 +227,62 @@ mod tests { assert!(pending.is_none()); } + #[test] + fn prepares_changed_session_for_its_unique_import() { + let root = TempDir::new().expect("tempdir"); + let codex_home = root.path().join("codex-home"); + let source_path = root.path().join("session.jsonl"); + let initial = session_record(root.path(), "first request"); + let expected_source_content_sha256 = format!("{:x}", Sha256::digest(&initial)); + std::fs::write(&source_path, initial).expect("initial session"); + let imported_thread_id = ThreadId::new(); + ledger::record_imported_session(&codex_home, &source_path, imported_thread_id) + .expect("record import"); + std::fs::write(&source_path, session_record(root.path(), "changed request")) + .expect("changed session"); + + let pending = + prepare_validated_session_import(&codex_home, session_migration(&source_path)) + .expect("prepare changed session") + .expect("changed session should be eligible"); + + assert_eq!( + pending.target, + SessionImportTarget::Existing { + thread_id: imported_thread_id, + expected_source_content_sha256, + } + ); + } + + #[test] + fn skips_changed_session_with_ambiguous_imports() { + 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, session_record(root.path(), "changed request")) + .expect("session"); + let source_path = std::fs::canonicalize(source_path).expect("canonical source"); + let import = |source_content_sha256: &str| ledger::CompletedExternalAgentSessionImport { + source_path: source_path.clone(), + source_content_sha256: source_content_sha256.to_string(), + imported_thread_id: ThreadId::new(), + connector_names: Vec::new(), + title: None, + }; + ledger::record_completed_session_imports( + &codex_home, + vec![import("first hash"), import("second hash")], + ) + .expect("record ambiguous imports"); + + let pending = + prepare_validated_session_import(&codex_home, session_migration(&source_path)) + .expect("prepare ambiguous session"); + + assert!(pending.is_none()); + } + #[test] fn reports_session_preparation_errors() { let root = TempDir::new().expect("tempdir"); @@ -227,6 +316,7 @@ mod tests { pending.source_content_sha256, format!("{:x}", Sha256::digest(contents)) ); + assert_eq!(pending.target, SessionImportTarget::New); } #[test] @@ -281,4 +371,14 @@ mod tests { title: None, } } + + fn session_record(cwd: &Path, text: &str) -> String { + serde_json::json!({ + "type": "user", + "cwd": cwd, + "timestamp": "2026-06-03T12:00:00Z", + "message": { "content": text }, + }) + .to_string() + } }