From da18000cae9884ab45f83b2d07fbd5a220a1de39 Mon Sep 17 00:00:00 2001 From: jif Date: Wed, 16 Sep 2026 15:25:34 +0000 Subject: [PATCH] Measure rollout read and materialization durations (#45966) ## What changed - Emit one read counter and duration observation per rollout reader, tagged by format, outcome, failure stage, and error kind. Distinguish EOF, partial reads, and failures, including failed opens. - Accumulate time spent in completed open/retry and read calls, excluding caller processing and canceled calls, without exporting a metric for each JSONL record. - Record materialization duration when decompressing a rollout for append, for both successful decompression and failures. GitOrigin-RevId: 3d1e257863ffa696a63040172df54d8234d56ed5 --- codex-rs/rollout/src/compression.rs | 106 ++++++++++++------ .../rollout/src/compression/error_metrics.rs | 39 ++++--- .../rollout/src/compression/read_metrics.rs | 58 ++++++++++ 3 files changed, 153 insertions(+), 50 deletions(-) create mode 100644 codex-rs/rollout/src/compression/read_metrics.rs diff --git a/codex-rs/rollout/src/compression.rs b/codex-rs/rollout/src/compression.rs index da4f0bfcc3..de9fb65b6e 100644 --- a/codex-rs/rollout/src/compression.rs +++ b/codex-rs/rollout/src/compression.rs @@ -9,6 +9,7 @@ use std::path::PathBuf; use std::sync::atomic::AtomicU64; use std::sync::atomic::Ordering; use std::time::Duration; +use std::time::Instant; #[cfg(unix)] use std::os::unix::fs::OpenOptionsExt; @@ -16,8 +17,10 @@ use std::os::unix::fs::OpenOptionsExt; use std::os::unix::fs::PermissionsExt; mod error_metrics; +mod read_metrics; use error_metrics::FailureMetric; +use read_metrics::ReadMetrics; const COMPRESSED_SUFFIX: &str = ".zst"; const MAX_NOT_FOUND_RETRIES: usize = 3; @@ -47,16 +50,29 @@ pub(crate) async fn file_modified_time(path: &Path) -> io::Result io::Result { - for _ in 0..MAX_NOT_FOUND_RETRIES { - match reader::open_once(path).await { - Ok(reader) => return Ok(reader), - Err(err) if err.kind() == io::ErrorKind::NotFound => { - tokio::time::sleep(OPEN_ROLLOUT_LINE_READER_RETRY_DELAY).await; + let started_at = Instant::now(); + let mut metrics = ReadMetrics::default(); + let result = async { + for _ in 0..MAX_NOT_FOUND_RETRIES { + match reader::open_once(path, &mut metrics).await { + Ok(reader) => return Ok(reader), + Err(err) if err.kind() == io::ErrorKind::NotFound => { + tokio::time::sleep(OPEN_ROLLOUT_LINE_READER_RETRY_DELAY).await; + } + Err(err) => return Err(err), } - Err(err) => return Err(err), + } + reader::open_once(path, &mut metrics).await + } + .await; + metrics.duration = started_at.elapsed(); + match result { + Ok(inner) => Ok(RolloutLineReader { inner, metrics }), + Err(err) => { + metrics.failed("open", &err); + Err(err) } } - reader::open_once(path).await } /// Returns the compressed `.jsonl.zst` path for a rollout path. @@ -93,13 +109,14 @@ pub(crate) fn materialize_rollout_for_append_blocking(path: &Path) -> io::Result return Ok(plain_path); } + let started_at = Instant::now(); let temp_path = temp_path_for(plain_path.as_path(), "decompress"); - if let Some(parent) = plain_path.parent() { - std::fs::create_dir_all(parent) - .inspect_err(|err| FailureMetric::Materialize.record("prepare_directory", err))?; - } - let mut stage = "read_metadata"; + let mut stage = "prepare_directory"; let result: io::Result<()> = (|| { + if let Some(parent) = plain_path.parent() { + std::fs::create_dir_all(parent)?; + } + stage = "read_metadata"; let metadata = std::fs::metadata(compressed_path.as_path())?; let permissions = metadata.permissions(); stage = "create_temp"; @@ -138,9 +155,11 @@ pub(crate) fn materialize_rollout_for_append_blocking(path: &Path) -> io::Result if let Err(err) = &result { let _ = std::fs::remove_file(temp_path.as_path()); FailureMetric::Materialize.record(stage, err); + metrics::materialize_duration("failed", started_at.elapsed()); } result?; metrics::materialize("decompressed"); + metrics::materialize_duration("decompressed", started_at.elapsed()); Ok(plain_path) } @@ -217,6 +236,7 @@ impl RolloutFile { /// Line-oriented rollout reader returned by [`open_rollout_line_reader`]. pub struct RolloutLineReader { inner: RolloutLineReaderInner, + metrics: ReadMetrics, } enum RolloutLineReaderInner { @@ -227,20 +247,31 @@ enum RolloutLineReaderInner { impl RolloutLineReader { /// Reads the next JSONL record from the rollout. pub async fn next_line(&mut self) -> io::Result> { - match &mut self.inner { - RolloutLineReaderInner::Plain(lines) => lines.next_line().await, - RolloutLineReaderInner::Blocking(slot) => { - let Some(mut reader) = slot.take() else { - return Err(io::Error::other("compressed rollout reader is busy")); - }; - let (line, reader) = - tokio::task::spawn_blocking(move || (reader.next().transpose(), reader)) - .await - .map_err(io::Error::other)?; - *slot = Some(reader); - line + let started_at = Instant::now(); + self.metrics.reached_eof = false; + let result = async { + match &mut self.inner { + RolloutLineReaderInner::Plain(lines) => lines.next_line().await, + RolloutLineReaderInner::Blocking(slot) => { + let Some(mut reader) = slot.take() else { + return Err(io::Error::other("compressed rollout reader is busy")); + }; + let (line, reader) = + tokio::task::spawn_blocking(move || (reader.next().transpose(), reader)) + .await + .map_err(io::Error::other)?; + *slot = Some(reader); + line + } } } + .await; + self.metrics.duration = self.metrics.duration.saturating_add(started_at.elapsed()); + match &result { + Ok(line) => self.metrics.reached_eof = line.is_none(), + Err(err) => self.metrics.failed("read", err), + } + result } } @@ -1039,6 +1070,14 @@ mod metrics { counter(MATERIALIZE_COUNTER, &[("outcome", outcome)]); } + pub(super) fn materialize_duration(outcome: &'static str, duration: Duration) { + duration_histogram( + "codex.rollout_compression.materialize.duration_ms", + duration, + &[("outcome", outcome)], + ); + } + pub(super) fn run(status: &'static str) { counter(RUN_COUNTER, &[("status", status)]); } @@ -1169,16 +1208,20 @@ mod reader { use std::io::Read; use std::path::Path; - use super::RolloutLineReader; + use super::ReadMetrics; use super::RolloutLineReaderInner; use super::path; use tokio::io::AsyncBufReadExt; - pub(super) async fn open_once(path: &Path) -> io::Result { + pub(super) async fn open_once( + path: &Path, + metrics: &mut ReadMetrics, + ) -> io::Result { let path = path::existing_rollout_path(path) .await .unwrap_or_else(|| path.to_path_buf()); if path::is_compressed_rollout_path(path.as_path()) { + metrics.format = "zstd"; let reader = tokio::task::spawn_blocking(move || { let input = File::open(path.as_path())?; let decoder = zstd::stream::read::Decoder::new(input)?; @@ -1188,14 +1231,13 @@ mod reader { }) .await .map_err(io::Error::other)??; - return Ok(RolloutLineReader { - inner: RolloutLineReaderInner::Blocking(Some(reader)), - }); + return Ok(RolloutLineReaderInner::Blocking(Some(reader))); } + metrics.format = "plain"; let file = tokio::fs::File::open(path).await?; - Ok(RolloutLineReader { - inner: RolloutLineReaderInner::Plain(tokio::io::BufReader::new(file).lines()), - }) + Ok(RolloutLineReaderInner::Plain( + tokio::io::BufReader::new(file).lines(), + )) } } diff --git a/codex-rs/rollout/src/compression/error_metrics.rs b/codex-rs/rollout/src/compression/error_metrics.rs index 04e3f202c6..614e85ca9d 100644 --- a/codex-rs/rollout/src/compression/error_metrics.rs +++ b/codex-rs/rollout/src/compression/error_metrics.rs @@ -22,23 +22,6 @@ impl FailureMetric { Self::Scan => ("codex.rollout_compression.scan", "outcome"), Self::TempCleanup => (metrics::TEMP_CLEANUP_COUNTER, "outcome"), }; - let error_kind = match error.kind() { - io::ErrorKind::NotFound => "not_found", - io::ErrorKind::PermissionDenied => "permission_denied", - io::ErrorKind::AlreadyExists => "already_exists", - io::ErrorKind::InvalidInput => "invalid_input", - io::ErrorKind::InvalidData => "invalid_data", - io::ErrorKind::TimedOut => "timed_out", - io::ErrorKind::WriteZero => "write_zero", - io::ErrorKind::Interrupted => "interrupted", - io::ErrorKind::Unsupported => "unsupported", - io::ErrorKind::UnexpectedEof => "unexpected_eof", - io::ErrorKind::StorageFull => "storage_full", - io::ErrorKind::ReadOnlyFilesystem => "read_only_filesystem", - io::ErrorKind::NotADirectory => "not_a_directory", - io::ErrorKind::IsADirectory => "is_a_directory", - _ => "other", - }; let Some(metrics) = codex_otel::global() else { return; }; @@ -48,8 +31,28 @@ impl FailureMetric { &[ (outcome_key, "failed"), ("stage", stage), - ("error_kind", error_kind), + ("error_kind", error_kind(error)), ], ); } } + +pub(super) fn error_kind(error: &io::Error) -> &'static str { + match error.kind() { + io::ErrorKind::NotFound => "not_found", + io::ErrorKind::PermissionDenied => "permission_denied", + io::ErrorKind::AlreadyExists => "already_exists", + io::ErrorKind::InvalidInput => "invalid_input", + io::ErrorKind::InvalidData => "invalid_data", + io::ErrorKind::TimedOut => "timed_out", + io::ErrorKind::WriteZero => "write_zero", + io::ErrorKind::Interrupted => "interrupted", + io::ErrorKind::Unsupported => "unsupported", + io::ErrorKind::UnexpectedEof => "unexpected_eof", + io::ErrorKind::StorageFull => "storage_full", + io::ErrorKind::ReadOnlyFilesystem => "read_only_filesystem", + io::ErrorKind::NotADirectory => "not_a_directory", + io::ErrorKind::IsADirectory => "is_a_directory", + _ => "other", + } +} diff --git a/codex-rs/rollout/src/compression/read_metrics.rs b/codex-rs/rollout/src/compression/read_metrics.rs new file mode 100644 index 0000000000..3840f0215a --- /dev/null +++ b/codex-rs/rollout/src/compression/read_metrics.rs @@ -0,0 +1,58 @@ +//! One observation per rollout reader, including failed opens and partial reads. +//! Duration sums completed open/retry and read calls, excluding caller processing +//! and any canceled call. Partial readers must not be treated as successful EOFs. +//! Aggregate locally rather than exporting a metric for every JSONL record. + +use std::io; +use std::time::Duration; + +use super::error_metrics::error_kind; + +pub(super) struct ReadMetrics { + pub(super) format: &'static str, + pub(super) reached_eof: bool, + pub(super) duration: Duration, + failure: Option<(&'static str, &'static str)>, +} + +impl Default for ReadMetrics { + fn default() -> Self { + Self { + format: "unknown", + reached_eof: false, + duration: Duration::ZERO, + failure: None, + } + } +} + +impl ReadMetrics { + pub(super) fn failed(&mut self, stage: &'static str, error: &io::Error) { + self.failure.get_or_insert((stage, error_kind(error))); + } +} + +impl Drop for ReadMetrics { + fn drop(&mut self) { + let Some(metrics) = codex_otel::global() else { + return; + }; + let (outcome, stage, error_kind) = match self.failure { + Some((stage, error_kind)) => ("failed", stage, error_kind), + None if self.reached_eof => ("eof", "none", "none"), + None => ("partial", "none", "none"), + }; + let tags = [ + ("format", self.format), + ("outcome", outcome), + ("stage", stage), + ("error_kind", error_kind), + ]; + let _ = metrics.counter("codex.rollout_compression.read", /*inc*/ 1, &tags); + let _ = metrics.record_duration( + "codex.rollout_compression.read.io_duration_ms", + self.duration, + &tags, + ); + } +}