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
This commit is contained in:
stefanstokic-oai
2026-07-31 16:10:04 +00:00
committed by jif-oai
parent 796faee10f
commit 370b02ac96
11 changed files with 1170 additions and 19 deletions

1
codex-rs/Cargo.lock generated
View File

@@ -3067,6 +3067,7 @@ dependencies = [
"codex-plugin",
"codex-protocol",
"codex-rollout",
"codex-thread-store",
"codex-utils-output-truncation",
"pretty_assertions",
"serde",

View File

@@ -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<ImportedSessionConnectorAttribution>,
}
enum SessionImportOutcome {
Created(CompletedSessionImport),
Appended {
source_path: PathBuf,
imported_thread_id: ThreadId,
title: Option<String>,
},
}
#[derive(Clone)]
pub(super) struct ExternalAgentSessionImporter {
codex_home: PathBuf,
connector_metadata_roots: Vec<PathBuf>,
permits: Arc<Semaphore>,
append_checkpoint_permits: Arc<Semaphore>,
thread_manager: Arc<ThreadManager>,
thread_store: Arc<dyn ThreadStore>,
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<Option<CompletedSessionImport>, SessionImportFailure> {
) -> Result<Option<SessionImportOutcome>, 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<String>,
session: ImportedExternalAgentSession,
) -> Result<CompletedSessionImport, SessionImportFailure> {
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(

View File

@@ -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 = "<EXTERNAL SESSION IMPORTED>";
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<Self> {
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::<Vec<_>>()
.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<T: serde::de::DeserializeOwned>(
&mut self,
method: &str,
params: Value,
) -> Result<T> {
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<ExternalAgentConfigDetectResponse> {
self.rpc("externalAgentConfig/detect", json!({ "includeHome": true }))
.await
}
async fn import_detected(&mut self) -> Result<Vec<String>> {
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<String> {
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<ThreadReadResponse> {
self.rpc(
"thread/read",
json!({ "threadId": thread_id, "includeTurns": true }),
)
.await
}
async fn thread_count(&mut self) -> Result<usize> {
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<TestAppServer> {
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
}

View File

@@ -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;

View File

@@ -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]

View File

@@ -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<Vec<RolloutItem>> {
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::<Vec<_>>();
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<Vec<SourceModelItem<'_>>> {
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<Vec<&ResponseItem>> {
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<ThreadPersistenceMetadata> {
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;

View File

@@ -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<RolloutItem> {
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()
}

View File

@@ -24,7 +24,7 @@ use std::collections::BTreeSet;
use std::io;
use std::path::Path;
const EXTERNAL_SESSION_IMPORTED_MARKER: &str = "<EXTERNAL SESSION IMPORTED>";
pub(super) const EXTERNAL_SESSION_IMPORTED_MARKER: &str = "<EXTERNAL SESSION IMPORTED>";
#[cfg(test)]
fn load_session_for_import(path: &Path) -> io::Result<Option<ImportedExternalAgentSession>> {
@@ -82,7 +82,7 @@ pub(crate) fn load_session_for_import_with_content_sha256(
)))
}
fn rollout_items_from_messages(messages: Vec<ConversationMessage>) -> Vec<RolloutItem> {
pub(super) fn rollout_items_from_messages(messages: Vec<ConversationMessage>) -> Vec<RolloutItem> {
let mut items = Vec::new();
let mut current_turn = None;
let mut response_item_bytes = 0i64;

View File

@@ -44,6 +44,16 @@ pub struct CompletedExternalAgentSessionImport {
pub title: Option<String>,
}
#[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<SessionImportSourceMapping> {
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<bool> {
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<CompletedExternalAgentSessionImport>,

View File

@@ -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<u8> {
std::fs::read(super::import_ledger_path(codex_home)).expect("ledger")
}

View File

@@ -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<RolloutItem>,
}
#[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<String>,
pub session: ImportedExternalAgentSession,
}
@@ -101,11 +115,27 @@ pub fn prepare_validated_session_import_with_metadata_mode(
session: ExternalAgentSessionMigration,
metadata_mode: SessionMetadataMode,
) -> io::Result<Option<PendingSessionImport>> {
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()
}
}