diff --git a/codex-rs/rollout/src/compression.rs b/codex-rs/rollout/src/compression.rs index 996ac4a7e0..da4f0bfcc3 100644 --- a/codex-rs/rollout/src/compression.rs +++ b/codex-rs/rollout/src/compression.rs @@ -289,6 +289,9 @@ mod worker { compressed: usize, skipped: usize, failed: usize, + scan_errors: bool, + cleanup_errors: bool, + time_budget_exhausted: bool, } pub(super) struct CompressionRunMarker { @@ -399,14 +402,17 @@ mod worker { let writer_locks = Arc::new(crate::WriterLockCoordinator::new(&codex_home)); let mut stage = "temp_cleanup"; let result = async { - cleanup_stale_temps(codex_home.as_path()).await?; + let mut stats = CompressionStats { + cleanup_errors: cleanup_stale_temps(codex_home.as_path()).await?, + ..Default::default() + }; stage = "scan"; - let mut stats = CompressionStats::default(); for root in [ codex_home.join(ARCHIVED_SESSIONS_SUBDIR), codex_home.join(SESSIONS_SUBDIR), ] { if started_at.elapsed() >= WORKER_MAX_RUNTIME { + stats.time_budget_exhausted = true; break; } compress_rollouts_in_root(root.as_path(), started_at, &mut stats, &writer_locks) @@ -427,8 +433,35 @@ mod worker { "rollout compression worker finished: scanned={}, compressed={}, skipped={}, failed={}", stats.scanned, stats.compressed, stats.skipped, stats.failed ); - metrics::run("completed"); - metrics::run_duration("completed", started_at.elapsed()); + // Keep the existing completed outcome: it means the pass returned, not + // that every directory was scanned or every file was compressed. + let completion = if stats.time_budget_exhausted { + "time_budget" + } else { + "scan_finished" + }; + let tags = [ + ("status", "completed"), + ("completion_reason", completion), + ( + "file_errors", + if stats.failed > 0 { "true" } else { "false" }, + ), + ( + "scan_errors", + if stats.scan_errors { "true" } else { "false" }, + ), + ( + "cleanup_errors", + if stats.cleanup_errors { + "true" + } else { + "false" + }, + ), + ]; + metrics::counter(metrics::RUN_COUNTER, &tags); + metrics::duration_histogram(metrics::RUN_DURATION_HISTOGRAM, started_at.elapsed(), &tags); marker.persist(); Ok(()) } @@ -453,18 +486,28 @@ mod worker { stats: &mut CompressionStats, writer_locks: &Arc, ) -> io::Result<()> { - if !tokio::fs::try_exists(root).await.unwrap_or(false) { + if !tokio::fs::try_exists(root) + .await + .inspect_err(|err| { + stats.scan_errors = true; + FailureMetric::Scan.record("check_root", err); + }) + .unwrap_or(false) + { return Ok(()); } let mut stack = vec![root.to_path_buf()]; let mut jobs = JoinSet::new(); while let Some(dir) = stack.pop() { if started_at.elapsed() >= WORKER_MAX_RUNTIME { + stats.time_budget_exhausted = true; break; } let mut read_dir = match tokio::fs::read_dir(dir.as_path()).await { Ok(read_dir) => read_dir, Err(err) => { + stats.scan_errors = true; + FailureMetric::Scan.record("read_directory", &err); warn!( "failed to read rollout compression directory {}: {err}", dir.display() @@ -482,12 +525,15 @@ mod worker { } }; if started_at.elapsed() >= WORKER_MAX_RUNTIME { + stats.time_budget_exhausted = true; break; } let path = entry.path(); let file_type = match entry.file_type().await { Ok(file_type) => file_type, Err(err) => { + stats.scan_errors = true; + FailureMetric::Scan.record("read_file_type", &err); warn!( "failed to read rollout compression file type {}: {err}", path.display() @@ -510,13 +556,16 @@ mod worker { } let path = rollout_file.into_path(); if crate::rollout_id_from_path(path.as_path()).is_none() { + stats.scan_errors = true; stats.skipped = stats.skipped.saturating_add(1); metrics::file("skipped_unreadable_meta"); continue; } let thread_id = match crate::read_session_meta_line(path.as_path()).await { Ok(metadata) => metadata.meta.id, - Err(_) => { + Err(err) => { + stats.scan_errors = true; + FailureMetric::Scan.record("read_metadata", &err); stats.skipped = stats.skipped.saturating_add(1); metrics::file("skipped_unreadable_meta"); continue; @@ -836,25 +885,36 @@ mod worker { file.set_permissions(permissions.clone()) } - async fn cleanup_stale_temps(codex_home: &Path) -> io::Result<()> { + async fn cleanup_stale_temps(codex_home: &Path) -> io::Result { + let mut errors = false; for root in [ codex_home.join(SESSIONS_SUBDIR), codex_home.join(ARCHIVED_SESSIONS_SUBDIR), ] { - cleanup_stale_temps_in_root(root.as_path()).await?; + errors |= cleanup_stale_temps_in_root(root.as_path()).await?; } - Ok(()) + Ok(errors) } - async fn cleanup_stale_temps_in_root(root: &Path) -> io::Result<()> { - if !tokio::fs::try_exists(root).await.unwrap_or(false) { - return Ok(()); + async fn cleanup_stale_temps_in_root(root: &Path) -> io::Result { + let mut errors = false; + if !tokio::fs::try_exists(root) + .await + .inspect_err(|err| { + errors = true; + FailureMetric::TempCleanup.record("check_root", err); + }) + .unwrap_or(false) + { + return Ok(errors); } let mut stack = vec![root.to_path_buf()]; while let Some(dir) = stack.pop() { let mut read_dir = match tokio::fs::read_dir(dir.as_path()).await { Ok(read_dir) => read_dir, Err(err) => { + errors = true; + FailureMetric::TempCleanup.record("read_directory", &err); warn!( "failed to read rollout temp cleanup directory {}: {err}", dir.display() @@ -867,6 +927,8 @@ mod worker { let file_type = match entry.file_type().await { Ok(file_type) => file_type, Err(err) => { + errors = true; + FailureMetric::TempCleanup.record("read_file_type", &err); warn!( "failed to read rollout temp cleanup file type {}: {err}", path.display() @@ -887,8 +949,12 @@ mod worker { let stale = entry .metadata() .await + .and_then(|metadata| metadata.modified()) + .inspect_err(|err| { + errors = true; + FailureMetric::TempCleanup.record("read_metadata", err); + }) .ok() - .and_then(|metadata| metadata.modified().ok()) .and_then(|modified| SystemTime::now().duration_since(modified).ok()) .is_some_and(|age| age >= TEMP_FILE_STALE_AFTER); if !stale { @@ -898,6 +964,7 @@ mod worker { Ok(()) => metrics::temp_cleanup("removed"), Err(err) if err.kind() == io::ErrorKind::NotFound => {} Err(err) => { + errors = true; FailureMetric::TempCleanup.record("remove_temp", &err); warn!( "failed to remove stale rollout temp {}: {err}", @@ -908,7 +975,7 @@ mod worker { } } } - Ok(()) + Ok(errors) } } @@ -923,7 +990,7 @@ mod metrics { "codex.rollout_compression.file.compression_ratio"; pub(super) const MATERIALIZE_COUNTER: &str = "codex.rollout_compression.materialize"; pub(super) const RUN_COUNTER: &str = "codex.rollout_compression.run"; - const RUN_DURATION_HISTOGRAM: &str = "codex.rollout_compression.run.duration_ms"; + pub(super) const RUN_DURATION_HISTOGRAM: &str = "codex.rollout_compression.run.duration_ms"; const RATIO_BASIS_POINTS: u128 = 10_000; pub(super) const TEMP_CLEANUP_COUNTER: &str = "codex.rollout_compression.temp_cleanup"; @@ -984,7 +1051,7 @@ mod metrics { counter(TEMP_CLEANUP_COUNTER, &[("outcome", outcome)]); } - fn counter(name: &str, tags: &[(&str, &str)]) { + pub(super) fn counter(name: &str, tags: &[(&str, &str)]) { let Some(metrics) = codex_otel::global() else { return; }; @@ -998,7 +1065,7 @@ mod metrics { let _ = metrics.histogram(name, value, tags); } - fn duration_histogram(name: &str, duration: Duration, tags: &[(&str, &str)]) { + pub(super) fn duration_histogram(name: &str, duration: Duration, tags: &[(&str, &str)]) { let Some(metrics) = codex_otel::global() else { return; }; diff --git a/codex-rs/rollout/src/compression/error_metrics.rs b/codex-rs/rollout/src/compression/error_metrics.rs index 3c75c3bd91..04e3f202c6 100644 --- a/codex-rs/rollout/src/compression/error_metrics.rs +++ b/codex-rs/rollout/src/compression/error_metrics.rs @@ -9,6 +9,7 @@ pub(super) enum FailureMetric { File, Materialize, Run, + Scan, TempCleanup, } @@ -18,6 +19,7 @@ impl FailureMetric { Self::File => (metrics::FILE_COUNTER, "outcome"), Self::Materialize => (metrics::MATERIALIZE_COUNTER, "outcome"), Self::Run => (metrics::RUN_COUNTER, "status"), + Self::Scan => ("codex.rollout_compression.scan", "outcome"), Self::TempCleanup => (metrics::TEMP_CLEANUP_COUNTER, "outcome"), }; let error_kind = match error.kind() {