diff --git a/codex-rs/thread-store/src/lib.rs b/codex-rs/thread-store/src/lib.rs index 75253e0541..70fc994d1a 100644 --- a/codex-rs/thread-store/src/lib.rs +++ b/codex-rs/thread-store/src/lib.rs @@ -25,6 +25,7 @@ pub use live_thread::LiveThread; pub use live_thread::LiveThreadInitGuard; pub use local::LocalThreadStore; pub use local::LocalThreadStoreConfig; +pub use local::RolloutMigrationFailureReason; pub use local::RolloutMigrationMode; pub use local::RolloutMigrationOptions; pub use local::RolloutMigrationOutcome; diff --git a/codex-rs/thread-store/src/local/mod.rs b/codex-rs/thread-store/src/local/mod.rs index 2b3488bba0..b08a17d261 100644 --- a/codex-rs/thread-store/src/local/mod.rs +++ b/codex-rs/thread-store/src/local/mod.rs @@ -99,6 +99,7 @@ use crate::UpdatedProject; use crate::local::writer_lock::WriterLockCoordinator; use crate::local::writer_lock::WriterLockGuard; +pub use rollout_migration::RolloutMigrationFailureReason; pub use rollout_migration::RolloutMigrationMode; pub use rollout_migration::RolloutMigrationOptions; pub use rollout_migration::RolloutMigrationOutcome; diff --git a/codex-rs/thread-store/src/local/rollout_migration.rs b/codex-rs/thread-store/src/local/rollout_migration.rs index 666e626f77..028b6d0094 100644 --- a/codex-rs/thread-store/src/local/rollout_migration.rs +++ b/codex-rs/thread-store/src/local/rollout_migration.rs @@ -9,6 +9,7 @@ //! journal must be enough for a later migration run to finish SQLite recovery safely. use std::collections::HashMap; +use std::io; use std::path::Path; use std::path::PathBuf; use std::time::Duration; @@ -131,12 +132,28 @@ pub enum RolloutMigrationStatus { Failed, } +/// A bounded explanation for why one rollout migration failed. +#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum RolloutMigrationFailureReason { + MissingSqliteMetadata, + InvalidSessionMetadata, + RolloutReadFailed, + LegacyRolloutConversionFailed, + SqliteMaterializationFailed, + RolloutPublishFailed, + InterruptedMigrationRecoveryFailed, + Unknown, +} + /// The per-thread result of a rollout migration run. #[derive(Clone, Debug, Eq, PartialEq, Serialize)] pub struct RolloutMigrationOutcome { pub thread_id: Option, pub rollout_path: PathBuf, pub status: RolloutMigrationStatus, + #[serde(skip_serializing_if = "Option::is_none")] + pub failure_reason: Option, pub bytes_processed: u64, pub message: Option, } @@ -174,6 +191,26 @@ enum RolloutMigrationKind { Subagent, } +struct RolloutMigrationFailure { + reason: RolloutMigrationFailureReason, + error: ThreadStoreError, +} + +impl RolloutMigrationFailure { + fn new(reason: RolloutMigrationFailureReason, error: ThreadStoreError) -> Self { + Self { reason, error } + } +} + +type ClassifiedMigrationResult = Result; + +fn with_failure_reason( + result: ThreadStoreResult, + reason: RolloutMigrationFailureReason, +) -> ClassifiedMigrationResult { + result.map_err(|error| RolloutMigrationFailure::new(reason, error)) +} + impl RolloutMigrationRateLimiter { fn new(max_mib_per_second: Option) -> ThreadStoreResult { let bytes_per_second = max_mib_per_second @@ -349,6 +386,13 @@ impl LocalThreadStore { let empty = tokio::fs::metadata(&path) .await .is_ok_and(|metadata| metadata.len() == 0); + let failure_reason = if empty { + None + } else if error.kind() != io::ErrorKind::Other { + Some(RolloutMigrationFailureReason::RolloutReadFailed) + } else { + Some(RolloutMigrationFailureReason::InvalidSessionMetadata) + }; return Ok(Some(RolloutMigrationOutcome { thread_id, rollout_path: path, @@ -357,6 +401,7 @@ impl LocalThreadStore { } else { RolloutMigrationStatus::Failed }, + failure_reason, bytes_processed: 0, message: (!empty).then(|| error.to_string()), })); @@ -417,7 +462,10 @@ impl LocalThreadStore { bytes_processed, ))); } - Err(error) => Err(error), + Err(error) => Err(RolloutMigrationFailure::new( + RolloutMigrationFailureReason::InterruptedMigrationRecoveryFailed, + error, + )), } } else { if options.mode == RolloutMigrationMode::Apply { @@ -455,7 +503,10 @@ impl LocalThreadStore { return Ok(Some(migration_outcome( thread_id, path, - Err(error), + Err(RolloutMigrationFailure::new( + RolloutMigrationFailureReason::Unknown, + error, + )), /*bytes_processed*/ 0, ))); } @@ -466,24 +517,29 @@ impl LocalThreadStore { .await { Ok(()) => Ok(RolloutMigrationStatus::Migrated), - Err(error) => { + Err(failure) => { if let Err(cleanup_error) = self .cleanup_failed_unpublished_migration(thread_id, &path, &journal_path) .await { - Err(migration_error(format!( - "{error}; failed to clean up unpublished migration: {cleanup_error}" - ))) + Err(RolloutMigrationFailure::new( + failure.reason, + migration_error(format!( + "{}; failed to clean up unpublished migration: {cleanup_error}", + failure.error + )), + )) } else { - Err(error) + Err(failure) } } }; let bytes_processed = limiter.bytes_processed.saturating_sub(bytes_before); Ok(Some(match result { - Err(ThreadStoreError::Conflict { message }) => { - skipped_busy_outcome(thread_id, path, message, bytes_processed) - } + Err(RolloutMigrationFailure { + error: ThreadStoreError::Conflict { message }, + .. + }) => skipped_busy_outcome(thread_id, path, message, bytes_processed), result => migration_outcome(thread_id, path, result, bytes_processed), })) } @@ -496,38 +552,63 @@ impl LocalThreadStore { kind: RolloutMigrationKind, legacy_names: &HashMap, limiter: &mut RolloutMigrationRateLimiter, - ) -> ThreadStoreResult<()> { + ) -> ClassifiedMigrationResult<()> { if let Some(state_db) = &self.state_db - && state_db - .get_thread(thread_id) - .await - .map_err(migration_error)? - .is_none() + && with_failure_reason( + state_db + .get_thread(thread_id) + .await + .map_err(migration_error), + RolloutMigrationFailureReason::Unknown, + )? + .is_none() { - return Err(migration_error(format!( - "thread {thread_id} is missing its SQLite metadata" - ))); + return Err(RolloutMigrationFailure::new( + RolloutMigrationFailureReason::MissingSqliteMetadata, + migration_error(format!("thread {thread_id} is missing its SQLite metadata")), + )); } let compressed = rollout_path_is_compressed(rollout_path); - let staged_path = staged_rollout_path(rollout_path)?; - let decompressed_path = compressed - .then(|| decompressed_staged_rollout_path(rollout_path)) - .transpose()?; - thread_history::delete_thread(self, thread_id).await?; - write_migration_journal(journal_path).await?; + let staged_path = with_failure_reason( + staged_rollout_path(rollout_path), + RolloutMigrationFailureReason::Unknown, + )?; + let decompressed_path = with_failure_reason( + compressed + .then(|| decompressed_staged_rollout_path(rollout_path)) + .transpose(), + RolloutMigrationFailureReason::Unknown, + )?; + with_failure_reason( + thread_history::delete_thread(self, thread_id).await, + RolloutMigrationFailureReason::SqliteMaterializationFailed, + )?; + with_failure_reason( + write_migration_journal(journal_path).await, + RolloutMigrationFailureReason::RolloutPublishFailed, + )?; - let source_metadata = tokio::fs::metadata(rollout_path) - .await - .map_err(migration_error)?; + let source_metadata = with_failure_reason( + tokio::fs::metadata(rollout_path) + .await + .map_err(migration_error), + RolloutMigrationFailureReason::RolloutReadFailed, + )?; let source_modified = source_metadata.modified().ok(); let source_permissions = source_metadata.permissions(); let source_path = if let Some(decompressed_path) = decompressed_path.as_ref() { - decompress_rollout_to_path(rollout_path, decompressed_path).await?; - let decompressed_bytes = tokio::fs::metadata(decompressed_path) - .await - .map_err(migration_error)? - .len(); + with_failure_reason( + decompress_rollout_to_path(rollout_path, decompressed_path).await, + RolloutMigrationFailureReason::RolloutReadFailed, + )?; + let decompressed_bytes = with_failure_reason( + tokio::fs::metadata(decompressed_path) + .await + .map_err(migration_error), + RolloutMigrationFailureReason::RolloutReadFailed, + )? + .len(); limiter .account(source_metadata.len().saturating_add(decompressed_bytes)) .await; @@ -535,7 +616,10 @@ impl LocalThreadStore { } else { rollout_path }; - let source_file = File::open(source_path).await.map_err(migration_error)?; + let source_file = with_failure_reason( + File::open(source_path).await.map_err(migration_error), + RolloutMigrationFailureReason::RolloutReadFailed, + )?; let mut source = BufReader::with_capacity(PROJECTION_BATCH_BYTES as usize, source_file); let mut bytes = Vec::new(); @@ -543,9 +627,16 @@ impl LocalThreadStore { // readers tolerate a pre-header prefix, so find that metadata before replaying the source // instead of buffering the prefix in memory. let canonical_session_meta = loop { - let record = read_rollout_record(&mut source, &mut bytes) - .await? - .ok_or_else(|| migration_error("rollout contains no session metadata"))?; + let record = with_failure_reason( + read_rollout_record(&mut source, &mut bytes).await, + RolloutMigrationFailureReason::RolloutReadFailed, + )? + .ok_or_else(|| { + RolloutMigrationFailure::new( + RolloutMigrationFailureReason::InvalidSessionMetadata, + migration_error("rollout contains no session metadata"), + ) + })?; limiter.account(record.byte_count).await; let Some(line) = record.line else { continue; @@ -563,97 +654,130 @@ impl LocalThreadStore { source_permissions: &source_permissions, canonical_session_meta: &canonical_session_meta, }; - let bounded_subagent_context = if kind == RolloutMigrationKind::Subagent { - let RolloutItem::SessionMeta(session_meta) = &canonical_session_meta.item else { - return Err(migration_error("canonical session metadata is missing")); - }; - let context = - subagent::select_bounded_context(source_path.to_path_buf(), session_meta.clone()) - .await?; - limiter.account(source_metadata.len()).await; - context - } else { - None - }; - let (_, expected_ordinal) = if let Some(items) = bounded_subagent_context { - Self::write_bounded_subagent_rollout(&canonicalization_source, items, limiter).await? - } else { - Self::write_rollout_with_rollback_plan(&canonicalization_source, limiter).await? - }; - if kind == RolloutMigrationKind::Subagent { - rewrite_subagent_history_boundary(&staged_path, expected_ordinal).await?; - } - let expected_length = tokio::fs::metadata(&staged_path) + // Everything up through a durable staged file is one legacy-to-paginated conversion + // phase. Keep the individual operations readable and tag the phase once if it fails. + let conversion_result = async { + let bounded_subagent_context = if kind == RolloutMigrationKind::Subagent { + let RolloutItem::SessionMeta(session_meta) = &canonical_session_meta.item else { + return Err(migration_error("canonical session metadata is missing")); + }; + let context = subagent::select_bounded_context( + source_path.to_path_buf(), + session_meta.clone(), + ) + .await?; + limiter.account(source_metadata.len()).await; + context + } else { + None + }; + let (_, expected_ordinal) = if let Some(items) = bounded_subagent_context { + Self::write_bounded_subagent_rollout(&canonicalization_source, items, limiter) + .await? + } else { + Self::write_rollout_with_rollback_plan(&canonicalization_source, limiter).await? + }; + + if kind == RolloutMigrationKind::Subagent { + rewrite_subagent_history_boundary(&staged_path, expected_ordinal).await?; + } + let expected_length = tokio::fs::metadata(&staged_path) + .await + .map_err(migration_error)? + .len(); + + // SQLite projection only starts after every staged-file mutation is durable. + let modified_at = source_modified; + let path = staged_path.clone(); + tokio::task::spawn_blocking(move || { + let file = std::fs::OpenOptions::new().write(true).open(path)?; + if let Some(modified_at) = modified_at { + file.set_times(std::fs::FileTimes::new().set_modified(modified_at))?; + } + file.sync_all() + }) .await .map_err(migration_error)? - .len(); - - // SQLite projection only starts after every staged-file mutation is durable. - let modified_at = source_modified; - let path = staged_path.clone(); - tokio::task::spawn_blocking(move || { - let file = std::fs::OpenOptions::new().write(true).open(path)?; - if let Some(modified_at) = modified_at { - file.set_times(std::fs::FileTimes::new().set_modified(modified_at))?; - } - file.sync_all() - }) - .await - .map_err(migration_error)? - .map_err(migration_error)?; - - self.project_rollout_in_batches(thread_id, &staged_path, limiter) - .await?; - let projection = thread_history::projection_state(self, thread_id) - .await? - .ok_or_else(|| migration_error("completed rollout has no SQLite projection"))?; - if projection.next_byte_offset != expected_length - || projection.next_ordinal != expected_ordinal - { - return Err(migration_error( - "SQLite projection does not cover the complete staged rollout", - )); - } - - let compressed_staged_path = if compressed { - let path = compressed_staged_rollout_path(rollout_path)?; - compress_rollout_to_path(&staged_path, &path, source_permissions, source_modified) - .await?; - Some(path) - } else { - None - }; - - // Older writers do not know about migration locks. Keep the legacy path visible if its - // append-only source changed while the replacement rollout was being staged. - let current_source_metadata = tokio::fs::metadata(rollout_path) - .await .map_err(migration_error)?; - if current_source_metadata.len() != source_metadata.len() - || current_source_metadata.modified().ok() != source_modified - { - return Err(ThreadStoreError::Conflict { - message: "rollout changed while migration was staging it; close older Codex processes and retry".to_string(), - }); - } - if let Some(compressed_staged_path) = compressed_staged_path { - tokio::fs::rename(compressed_staged_path, rollout_path) - .await - .map_err(migration_error)?; - remove_file_if_present(&staged_path).await?; - if let Some(decompressed_path) = decompressed_path.as_ref() { - remove_file_if_present(decompressed_path).await?; - } - } else { - tokio::fs::rename(&staged_path, rollout_path) - .await - .map_err(migration_error)?; + Ok::<_, ThreadStoreError>((expected_length, expected_ordinal)) } - sync_parent_directory(rollout_path).await?; - self.finish_published_migration(thread_id, journal_path, legacy_names) - .await + .await; + let (expected_length, expected_ordinal) = with_failure_reason( + conversion_result, + RolloutMigrationFailureReason::LegacyRolloutConversionFailed, + )?; + + // SQLite sees only the final staged bytes, so all projection failures share one reason. + let projection_result = async { + self.project_rollout_in_batches(thread_id, &staged_path, limiter) + .await?; + let projection = thread_history::projection_state(self, thread_id) + .await? + .ok_or_else(|| migration_error("completed rollout has no SQLite projection"))?; + if projection.next_byte_offset != expected_length + || projection.next_ordinal != expected_ordinal + { + return Err(migration_error( + "SQLite projection does not cover the complete staged rollout", + )); + } + Ok::<_, ThreadStoreError>(()) + } + .await; + with_failure_reason( + projection_result, + RolloutMigrationFailureReason::SqliteMaterializationFailed, + )?; + + // Once projection is verified, the remaining work publishes the replacement and clears + // the durable pending journal. + let publish_result = async { + let compressed_staged_path = if compressed { + let path = compressed_staged_rollout_path(rollout_path)?; + compress_rollout_to_path(&staged_path, &path, source_permissions, source_modified) + .await?; + Some(path) + } else { + None + }; + + // Older writers do not know about migration locks. Keep the legacy path visible if + // its append-only source changed while the replacement rollout was being staged. + let current_source_metadata = tokio::fs::metadata(rollout_path) + .await + .map_err(migration_error)?; + if current_source_metadata.len() != source_metadata.len() + || current_source_metadata.modified().ok() != source_modified + { + return Err(ThreadStoreError::Conflict { + message: "rollout changed while migration was staging it; close older Codex processes and retry".to_string(), + }); + } + + if let Some(compressed_staged_path) = compressed_staged_path { + tokio::fs::rename(compressed_staged_path, rollout_path) + .await + .map_err(migration_error)?; + remove_file_if_present(&staged_path).await?; + if let Some(decompressed_path) = decompressed_path.as_ref() { + remove_file_if_present(decompressed_path).await?; + } + } else { + tokio::fs::rename(&staged_path, rollout_path) + .await + .map_err(migration_error)?; + } + sync_parent_directory(rollout_path).await?; + self.finish_published_migration(thread_id, journal_path, legacy_names) + .await + } + .await; + with_failure_reason( + publish_result, + RolloutMigrationFailureReason::RolloutPublishFailed, + ) } async fn build_rollback_plan( @@ -1126,7 +1250,7 @@ fn thread_id_from_rollout_filename(path: &Path) -> Option { fn migration_outcome( thread_id: ThreadId, rollout_path: PathBuf, - result: ThreadStoreResult, + result: ClassifiedMigrationResult, bytes_processed: u64, ) -> RolloutMigrationOutcome { match result { @@ -1134,15 +1258,17 @@ fn migration_outcome( thread_id: Some(thread_id), rollout_path, status, + failure_reason: None, bytes_processed, message: None, }, - Err(error) => RolloutMigrationOutcome { + Err(failure) => RolloutMigrationOutcome { thread_id: Some(thread_id), rollout_path, status: RolloutMigrationStatus::Failed, + failure_reason: Some(failure.reason), bytes_processed, - message: Some(error.to_string()), + message: Some(failure.error.to_string()), }, } } @@ -1157,6 +1283,7 @@ fn skipped_busy_outcome( thread_id: Some(thread_id), rollout_path, status: RolloutMigrationStatus::SkippedBusy, + failure_reason: None, bytes_processed, message: Some(message), } diff --git a/codex-rs/thread-store/src/local/rollout_migration/telemetry.rs b/codex-rs/thread-store/src/local/rollout_migration/telemetry.rs index 8a8d672b30..e27c66b577 100644 --- a/codex-rs/thread-store/src/local/rollout_migration/telemetry.rs +++ b/codex-rs/thread-store/src/local/rollout_migration/telemetry.rs @@ -5,6 +5,7 @@ use std::time::Instant; +use super::RolloutMigrationFailureReason; use super::RolloutMigrationMode; use super::RolloutMigrationOptions; use super::RolloutMigrationReport; @@ -56,12 +57,15 @@ impl RolloutMigrationTelemetry { let mut io_bytes = 0_u64; if let Ok(report) = result { for outcome in &report.outcomes { - let tags = [ + let mut tags = vec![ ("trigger", self.trigger.tag()), ("mode", self.mode.tag()), ("scope", self.scope.tag()), ("status", outcome.status.tag()), ]; + if let Some(failure_reason) = outcome.failure_reason { + tags.push(("failure_reason", failure_reason.tag())); + } let _ = metrics.counter(THREAD_METRIC, /*inc*/ 1, &tags); io_bytes = io_bytes.saturating_add(outcome.bytes_processed); } @@ -131,3 +135,18 @@ impl RolloutMigrationStatus { } } } + +impl RolloutMigrationFailureReason { + fn tag(self) -> &'static str { + match self { + Self::MissingSqliteMetadata => "missing_sqlite_metadata", + Self::InvalidSessionMetadata => "invalid_session_metadata", + Self::RolloutReadFailed => "rollout_read_failed", + Self::LegacyRolloutConversionFailed => "legacy_rollout_conversion_failed", + Self::SqliteMaterializationFailed => "sqlite_materialization_failed", + Self::RolloutPublishFailed => "rollout_publish_failed", + Self::InterruptedMigrationRecoveryFailed => "interrupted_migration_recovery_failed", + Self::Unknown => "unknown", + } + } +} diff --git a/codex-rs/thread-store/src/local/rollout_migration_tests.rs b/codex-rs/thread-store/src/local/rollout_migration_tests.rs index ab7cbf6307..96ee8cb0c2 100644 --- a/codex-rs/thread-store/src/local/rollout_migration_tests.rs +++ b/codex-rs/thread-store/src/local/rollout_migration_tests.rs @@ -45,6 +45,7 @@ use serde_json::json; use tempfile::TempDir; use super::LocalThreadStore; +use super::RolloutMigrationFailureReason; use super::RolloutMigrationMode; use super::RolloutMigrationOptions; use super::RolloutMigrationProgress; @@ -263,6 +264,16 @@ fn apply_options() -> RolloutMigrationOptions { } } +fn assert_failed_with_reason( + outcome: &super::RolloutMigrationOutcome, + failure_reason: RolloutMigrationFailureReason, +) { + assert_eq!( + (outcome.status, outcome.failure_reason), + (RolloutMigrationStatus::Failed, Some(failure_reason)) + ); +} + async fn indexed_store(home: &Path) -> LocalThreadStore { let config = test_config(home); let rollout_config = RolloutConfig { @@ -2252,6 +2263,57 @@ async fn migration_skips_empty_rollout_files() { ); } +#[tokio::test] +async fn migration_reports_missing_sqlite_metadata() { + let home = TempDir::new().expect("create Codex home"); + let thread_id = ThreadId::new(); + write_rollout( + home.path(), + thread_id, + SessionSource::Cli, + vec![user_message("question")], + ); + let store = indexed_store(home.path()).await; + store + .state_db + .as_ref() + .expect("state db") + .delete_thread(thread_id) + .await + .expect("remove thread metadata"); + + let report = store + .migrate_rollouts(apply_options()) + .await + .expect("inspect rollout with missing metadata"); + + assert_failed_with_reason( + &report.outcomes[0], + RolloutMigrationFailureReason::MissingSqliteMetadata, + ); +} + +#[tokio::test] +async fn migration_reports_invalid_session_metadata() { + let home = TempDir::new().expect("create Codex home"); + let thread_id = ThreadId::new(); + let directory = home.path().join("sessions/2025/01/03"); + fs::create_dir_all(&directory).expect("create rollout directory"); + let path = directory.join(format!("rollout-2025-01-03T12-00-00-{thread_id}.jsonl")); + fs::write(path, "not a rollout record\n").expect("write malformed rollout"); + let store = indexed_store(home.path()).await; + + let report = store + .migrate_rollouts(apply_options()) + .await + .expect("inspect rollout with invalid metadata"); + + assert_failed_with_reason( + &report.outcomes[0], + RolloutMigrationFailureReason::InvalidSessionMetadata, + ); +} + #[tokio::test] async fn migration_skips_malformed_lines_and_trailing_partial_tail() { let home = TempDir::new().expect("create Codex home");