From de86025f2c5128f4bedd36ff2443e216bc7463df Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 27 Jun 2026 15:14:44 +0800 Subject: [PATCH] feat(get): overlap read verify decode (#3945) --- .../src/erasure/coding/decode_reader.rs | 402 ++++++++++++++++-- crates/ecstore/src/set_disk/read.rs | 4 +- crates/io-metrics/src/lib.rs | 62 +++ scripts/run_get_codec_streaming_smoke.sh | 8 +- 4 files changed, 449 insertions(+), 27 deletions(-) diff --git a/crates/ecstore/src/erasure/coding/decode_reader.rs b/crates/ecstore/src/erasure/coding/decode_reader.rs index 576e9eef1..c197b9007 100644 --- a/crates/ecstore/src/erasure/coding/decode_reader.rs +++ b/crates/ecstore/src/erasure/coding/decode_reader.rs @@ -22,6 +22,7 @@ use crate::diagnostics::get::{ use crate::disk::error::Error as DiskError; use crate::erasure::codec::bridge::ErasureDecodeEngine; use crate::set_disk::shard_source::{ShardStripeSource, StripeReadState}; +use std::collections::VecDeque; use std::io; use std::io::ErrorKind; use std::pin::Pin; @@ -31,12 +32,55 @@ use std::time::Instant; use tokio::io::{AsyncRead, ReadBuf}; use tokio::task::JoinHandle; +const ENV_RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT: &str = "RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT"; +const DEFAULT_RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT: usize = 1; +const FILL_POLICY_SINGLE_INFLIGHT: &str = "single_inflight"; +const FILL_POLICY_DUAL_INFLIGHT: &str = "dual_inflight"; + type FillTask = JoinHandle>; struct FillResult { source: S, workspace: W, result: io::Result>>, + queued_buffers: VecDeque>, + deferred_error: Option, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum FillPolicy { + SingleInFlight, + DualInFlight, +} + +impl FillPolicy { + fn from_env() -> Self { + match rustfs_utils::get_env_usize( + ENV_RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT, + DEFAULT_RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT, + ) { + 2 => Self::DualInFlight, + _ => Self::SingleInFlight, + } + } + + const fn max_inflight(self) -> usize { + match self { + Self::SingleInFlight => 1, + Self::DualInFlight => 2, + } + } + + const fn additional_queued_buffers(self) -> usize { + self.max_inflight().saturating_sub(1) + } + + const fn as_str(self) -> &'static str { + match self { + Self::SingleInFlight => FILL_POLICY_SINGLE_INFLIGHT, + Self::DualInFlight => FILL_POLICY_DUAL_INFLIGHT, + } + } } pub(crate) struct ErasureDecodeReader @@ -44,16 +88,18 @@ where E: ErasureDecodeEngine, { metrics_path: &'static str, + fill_policy: FillPolicy, source: Option, engine: E, workspace: Option, output_buf: Vec, output_pos: usize, - prefetched_buf: Option>, + prefetched_bufs: VecDeque>, prefetch_error: Option, prefetch_wait_started_at: Option, + output_wait_started_at: Option, remaining: usize, - // Bounded lookahead: at most one background stripe read/decode is in flight. + // Bounded lookahead controlled by `FillPolicy`. fill: Option>, } @@ -71,6 +117,16 @@ where engine: E, total_length: usize, metrics_path: &'static str, + ) -> io::Result { + Self::new_with_fill_policy_inner(source, engine, total_length, metrics_path, FillPolicy::from_env()) + } + + fn new_with_fill_policy_inner( + source: S, + engine: E, + total_length: usize, + metrics_path: &'static str, + fill_policy: FillPolicy, ) -> io::Result { if engine.data_shards() == 0 { return Err(io::Error::new(ErrorKind::InvalidInput, "erasure reader requires data shards")); @@ -84,19 +140,32 @@ where Ok(Self { metrics_path, + fill_policy, source: Some(source), engine, workspace: Some(workspace), output_buf: Vec::new(), output_pos: 0, - prefetched_buf: None, + prefetched_bufs: VecDeque::new(), prefetch_error: None, prefetch_wait_started_at: None, + output_wait_started_at: None, remaining: total_length, fill: None, }) } + #[cfg(test)] + fn new_with_fill_policy( + source: S, + engine: E, + total_length: usize, + metrics_path: &'static str, + fill_policy: FillPolicy, + ) -> io::Result { + Self::new_with_fill_policy_inner(source, engine, total_length, metrics_path, fill_policy) + } + fn poll_fill_result(&mut self, cx: &mut Context<'_>) -> Poll>>> { if self.fill.is_none() { let Some(mut source) = self.source.take() else { @@ -109,8 +178,11 @@ where let engine = self.engine.clone(); let metrics_path = self.metrics_path; + let fill_policy = self.fill_policy; let remaining = self.remaining; self.fill = Some(tokio::spawn(async move { + let mut queued_buffers = VecDeque::new(); + let mut deferred_error = None; let fill_stage_start = Instant::now(); let stripe_read_stage_start = Instant::now(); let state = source.read_next_stripe().await; @@ -126,6 +198,44 @@ where GET_STAGE_DECODE, decode_stage_start.elapsed().as_secs_f64(), ); + if let Ok(Some(first_buf)) = result.as_ref() { + let mut remaining_after_first = remaining.saturating_sub(first_buf.len()); + for _ in 0..fill_policy.additional_queued_buffers() { + if remaining_after_first == 0 { + break; + } + let stripe_read_stage_start = Instant::now(); + let state = source.read_next_stripe().await; + rustfs_io_metrics::record_get_object_stage_duration( + metrics_path, + GET_STAGE_STRIPE_READ, + stripe_read_stage_start.elapsed().as_secs_f64(), + ); + let decode_stage_start = Instant::now(); + let queued_result = decode_stripe(metrics_path, &engine, &mut workspace, state, remaining_after_first); + rustfs_io_metrics::record_get_object_stage_duration( + metrics_path, + GET_STAGE_DECODE, + decode_stage_start.elapsed().as_secs_f64(), + ); + match queued_result { + Ok(Some(buf)) => { + remaining_after_first = remaining_after_first.saturating_sub(buf.len()); + queued_buffers.push_back(buf); + } + Ok(None) => { + if remaining_after_first > 0 { + deferred_error = Some(DiskError::LessData.into()); + } + break; + } + Err(err) => { + deferred_error = Some(err); + break; + } + } + } + } rustfs_io_metrics::record_get_object_stage_duration( metrics_path, GET_STAGE_FILL, @@ -135,6 +245,8 @@ where source, workspace, result, + queued_buffers, + deferred_error, } })); } @@ -148,6 +260,8 @@ where source, workspace, result, + queued_buffers, + deferred_error, } = match fill_result { Ok(result) => result, Err(err) => { @@ -159,6 +273,10 @@ where self.source = Some(source); self.workspace = Some(workspace); self.fill = None; + if let Some(deferred_error) = deferred_error { + self.prefetch_error = Some(deferred_error); + } + self.prefetched_bufs.extend(queued_buffers); match result { Ok(Some(buf)) => { @@ -182,7 +300,7 @@ where } fn poll_prefetch(&mut self, cx: &mut Context<'_>) -> Poll> { - if self.prefetched_buf.is_some() || self.prefetch_error.is_some() || self.remaining == 0 { + if self.prefetched_bufs.len() >= self.fill_policy.max_inflight() || self.prefetch_error.is_some() || self.remaining == 0 { return Poll::Ready(Ok(())); } @@ -205,16 +323,48 @@ where match fill { Ok(Some(buf)) => { - if self.output_pos < self.output_buf.len() { - rustfs_io_metrics::record_get_object_reader_prefetch(self.metrics_path, GET_READER_PREFETCH_STORED); - rustfs_io_metrics::record_get_object_reader_buffer(self.metrics_path, GET_READER_BUFFER_PREFETCH, buf.len()); - self.prefetched_buf = Some(buf); + let output_has_remaining = self.output_pos < self.output_buf.len(); + let mut queued_count = 0usize; + let mut queued_bytes = 0usize; + + if output_has_remaining { + rustfs_io_metrics::record_get_object_fill_completed_before_output_drained( + self.metrics_path, + self.fill_policy.as_str(), + ); + queued_bytes += buf.len(); + queued_count += 1; + self.prefetched_bufs.push_back(buf); } else { rustfs_io_metrics::record_get_object_reader_prefetch(self.metrics_path, GET_READER_PREFETCH_DIRECT); rustfs_io_metrics::record_get_object_reader_buffer(self.metrics_path, GET_READER_BUFFER_OUTPUT, buf.len()); self.output_buf = buf; self.output_pos = 0; } + + if !self.prefetched_bufs.is_empty() { + let total_prefetched_bytes = self.prefetched_bufs.iter().map(Vec::len).sum::(); + if total_prefetched_bytes > queued_bytes { + queued_bytes = total_prefetched_bytes; + } + if self.prefetched_bufs.len() > queued_count { + queued_count = self.prefetched_bufs.len(); + } + } + if queued_count > 0 { + rustfs_io_metrics::record_get_object_fill_queued(self.metrics_path, self.fill_policy.as_str(), queued_count); + rustfs_io_metrics::record_get_object_reader_prefetch_bytes( + self.metrics_path, + self.fill_policy.as_str(), + queued_bytes, + ); + rustfs_io_metrics::record_get_object_reader_prefetch(self.metrics_path, GET_READER_PREFETCH_STORED); + rustfs_io_metrics::record_get_object_reader_buffer( + self.metrics_path, + GET_READER_BUFFER_PREFETCH, + queued_bytes, + ); + } Poll::Ready(Ok(())) } Ok(None) => { @@ -241,6 +391,7 @@ where { fn drop(&mut self) { if let Some(fill) = self.fill.take() { + rustfs_io_metrics::record_get_object_fill_cancelled_on_drop(self.metrics_path, self.fill_policy.as_str()); fill.abort(); } } @@ -256,7 +407,7 @@ where fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { loop { if self.output_pos < self.output_buf.len() { - if self.prefetched_buf.is_none() + if self.prefetched_bufs.len() < self.fill_policy.max_inflight() && self.prefetch_error.is_none() && self.remaining > 0 && let Poll::Ready(result) = self.poll_prefetch(cx) @@ -283,7 +434,7 @@ where return Poll::Ready(Ok(())); } - if let Some(next_buf) = self.prefetched_buf.take() { + if let Some(next_buf) = self.prefetched_bufs.pop_front() { rustfs_io_metrics::record_get_object_reader_buffer(self.metrics_path, GET_READER_BUFFER_OUTPUT, next_buf.len()); self.output_buf = next_buf; self.output_pos = 0; @@ -298,19 +449,36 @@ where return Poll::Ready(Ok(())); } - ready!(self.poll_prefetch(cx))?; + if self.output_wait_started_at.is_none() { + self.output_wait_started_at = Some(Instant::now()); + } + let prefetch = ready!(self.poll_prefetch(cx)); + if let Some(started_at) = self.output_wait_started_at.take() { + rustfs_io_metrics::record_get_object_fill_waited_by_output( + self.metrics_path, + self.fill_policy.as_str(), + started_at.elapsed().as_secs_f64(), + ); + } + prefetch?; } } } pub(crate) struct SyncErasureDecodeReader { inner: Mutex, + metrics_path: &'static str, } impl SyncErasureDecodeReader { pub(crate) fn new(inner: R) -> Self { + Self::new_with_metrics_path(inner, GET_OBJECT_PATH_CODEC_STREAMING) + } + + pub(crate) fn new_with_metrics_path(inner: R, metrics_path: &'static str) -> Self { Self { inner: Mutex::new(inner), + metrics_path, } } } @@ -324,7 +492,7 @@ where let mut inner = match self.inner.lock() { Ok(inner) => { rustfs_io_metrics::record_get_object_stage_duration( - GET_OBJECT_PATH_CODEC_STREAMING, + self.metrics_path, GET_STAGE_OUTPUT_LOCK_WAIT, lock_wait_start.elapsed().as_secs_f64(), ); @@ -332,7 +500,7 @@ where } Err(_) => { rustfs_io_metrics::record_get_object_stage_duration( - GET_OBJECT_PATH_CODEC_STREAMING, + self.metrics_path, GET_STAGE_OUTPUT_LOCK_WAIT, lock_wait_start.elapsed().as_secs_f64(), ); @@ -351,13 +519,9 @@ where Poll::Ready(Err(_)) => GET_READER_POLL_READY_ERROR, Poll::Pending => GET_READER_POLL_PENDING, }; - rustfs_io_metrics::record_get_object_stage_duration( - GET_OBJECT_PATH_CODEC_STREAMING, - GET_STAGE_OUTPUT_POLL, - poll_duration, - ); + rustfs_io_metrics::record_get_object_stage_duration(self.metrics_path, GET_STAGE_OUTPUT_POLL, poll_duration); rustfs_io_metrics::record_get_object_reader_poll( - GET_OBJECT_PATH_CODEC_STREAMING, + self.metrics_path, poll_outcome, read_buf_remaining_before, filled_bytes, @@ -469,10 +633,13 @@ mod tests { use crate::set_disk::shard_source::{ShardSlot, StripeReadState}; use rustfs_utils::HashAlgorithm; use std::collections::VecDeque; + use std::future::pending; use std::io::Cursor; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; + use temp_env::with_var; use tokio::io::AsyncReadExt; + use tokio::sync::Notify; use tokio::task::yield_now; use tokio::time::{Duration, timeout}; @@ -482,6 +649,22 @@ mod tests { read_count: Option>, } + struct BlockingSource { + started: Arc, + dropped: Arc, + read_quorum: usize, + } + + struct BlockingSourceDropGuard { + dropped: Arc, + } + + impl Drop for BlockingSourceDropGuard { + fn drop(&mut self) { + self.dropped.fetch_add(1, Ordering::SeqCst); + } + } + #[async_trait::async_trait] impl ShardStripeSource for VecStripeSource { async fn read_next_stripe(&mut self) -> StripeReadState { @@ -494,6 +677,18 @@ mod tests { } } + #[async_trait::async_trait] + impl ShardStripeSource for BlockingSource { + async fn read_next_stripe(&mut self) -> StripeReadState { + let _guard = BlockingSourceDropGuard { + dropped: Arc::clone(&self.dropped), + }; + self.started.notify_one(); + pending::<()>().await; + StripeReadState::new(Vec::new(), self.read_quorum) + } + } + fn source_from_data(erasure: &Erasure, data: &[u8], missing_indexes: &[usize]) -> VecStripeSource { let read_quorum = erasure.data_shards; let stripes = data @@ -533,7 +728,13 @@ mod tests { E: ErasureDecodeEngine + Clone + Send + Sync + 'static, { let source = source_from_data(erasure, data, missing_indexes); - let mut reader = ErasureDecodeReader::new(source, engine, data.len())?; + let mut reader = ErasureDecodeReader::new_with_fill_policy( + source, + engine, + data.len(), + GET_OBJECT_PATH_CODEC_STREAMING, + FillPolicy::SingleInFlight, + )?; let mut decoded = Vec::new(); reader.read_to_end(&mut decoded).await?; Ok(decoded) @@ -544,6 +745,21 @@ mod tests { decode_all_with_engine(&erasure, engine, data, missing_indexes).await } + #[test] + fn fill_policy_defaults_to_single_inflight() { + with_var(ENV_RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT, None::<&str>, || { + assert_eq!(FillPolicy::from_env(), FillPolicy::SingleInFlight); + }); + + with_var(ENV_RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT, Some("2"), || { + assert_eq!(FillPolicy::from_env(), FillPolicy::DualInFlight); + }); + + with_var(ENV_RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT, Some("99"), || { + assert_eq!(FillPolicy::from_env(), FillPolicy::SingleInFlight); + }); + } + #[tokio::test] async fn erasure_decode_reader_reads_single_stripe() { let erasure = Erasure::new(4, 2, 64); @@ -573,7 +789,14 @@ mod tests { let erasure = Erasure::new(4, 2, 32); let source = source_from_data(&erasure, &[], &[]); let engine = LegacyEcDecodeEngine::new(erasure); - let mut reader = ErasureDecodeReader::new(source, engine, 0).expect("empty reader should be constructed"); + let mut reader = ErasureDecodeReader::new_with_fill_policy( + source, + engine, + 0, + GET_OBJECT_PATH_CODEC_STREAMING, + FillPolicy::SingleInFlight, + ) + .expect("empty reader should be constructed"); let mut decoded = Vec::new(); let read = reader @@ -599,6 +822,107 @@ mod tests { assert_eq!(decoded, data); } + #[tokio::test] + async fn erasure_decode_reader_dual_inflight_prefetches_an_extra_stripe() { + let erasure = Erasure::new(4, 2, 16); + let data = (0..48u8).collect::>(); + let read_count = Arc::new(AtomicUsize::new(0)); + let mut source = source_from_data(&erasure, &data, &[]); + source.read_count = Some(Arc::clone(&read_count)); + let engine = LegacyEcDecodeEngine::new(erasure); + let mut reader = ErasureDecodeReader::new_with_fill_policy( + source, + engine, + data.len(), + GET_OBJECT_PATH_CODEC_STREAMING, + FillPolicy::DualInFlight, + ) + .expect("reader should be constructed"); + let mut first_read = [0u8; 1]; + + let read = reader.read(&mut first_read).await.expect("first read should succeed"); + + assert_eq!(read, 1); + timeout(Duration::from_secs(1), async { + while read_count.load(Ordering::SeqCst) < 3 { + yield_now().await; + } + }) + .await + .expect("dual inflight policy should prefetch two future stripes before current output drains"); + } + + #[tokio::test] + async fn erasure_decode_reader_defers_short_read_error_until_buffer_drains() { + let erasure = Erasure::new(4, 2, 32); + let first_stripe = (0..32u8).collect::>(); + let source = VecStripeSource { + stripes: VecDeque::from([ + source_from_data(&erasure, &first_stripe, &[]) + .stripes + .pop_front() + .expect("first stripe should exist"), + StripeReadState::new(Vec::new(), erasure.data_shards), + ]), + read_quorum: erasure.data_shards, + read_count: None, + }; + let engine = LegacyEcDecodeEngine::new(erasure); + let mut reader = ErasureDecodeReader::new_with_fill_policy( + source, + engine, + first_stripe.len() + 1, + GET_OBJECT_PATH_CODEC_STREAMING, + FillPolicy::SingleInFlight, + ) + .expect("reader should be constructed"); + let mut decoded = Vec::new(); + + let err = reader + .read_to_end(&mut decoded) + .await + .expect_err("short-read error should surface after buffered output drains"); + + assert_eq!(err.kind(), ErrorKind::Other); + assert_eq!(decoded, first_stripe); + } + + #[tokio::test] + async fn erasure_decode_reader_drop_aborts_inflight_fill_task() { + let started = Arc::new(Notify::new()); + let dropped = Arc::new(AtomicUsize::new(0)); + let source = BlockingSource { + started: Arc::clone(&started), + dropped: Arc::clone(&dropped), + read_quorum: 1, + }; + let engine = LegacyEcDecodeEngine::new(Erasure::new(1, 0, 32)); + let task = tokio::spawn(async move { + let mut reader = ErasureDecodeReader::new_with_fill_policy( + source, + engine, + 1, + GET_OBJECT_PATH_CODEC_STREAMING, + FillPolicy::SingleInFlight, + ) + .expect("reader should be constructed"); + let mut first_read = [0u8; 1]; + let _ = reader.read(&mut first_read).await; + }); + + started.notified().await; + task.abort(); + let _ = task.await; + + timeout(Duration::from_secs(1), async { + while dropped.load(Ordering::SeqCst) == 0 { + yield_now().await; + } + }) + .await + .expect("reader drop should abort the in-flight fill task"); + } + #[tokio::test] async fn erasure_decode_reader_rejects_inconsistent_reconstruction_sources() { const DATA_SHARDS: usize = 2; @@ -626,7 +950,14 @@ mod tests { read_count: None, }; let engine = LegacyEcDecodeEngine::new(erasure); - let mut reader = ErasureDecodeReader::new(source, engine, data.len()).expect("reader should be constructed"); + let mut reader = ErasureDecodeReader::new_with_fill_policy( + source, + engine, + data.len(), + GET_OBJECT_PATH_CODEC_STREAMING, + FillPolicy::SingleInFlight, + ) + .expect("reader should be constructed"); let mut decoded = Vec::new(); let err = reader @@ -705,7 +1036,14 @@ mod tests { Some(GET_OBJECT_PATH_CODEC_STREAMING), ); let engine = LegacyEcDecodeEngine::new(erasure); - let mut reader = ErasureDecodeReader::new(source, engine, data.len()).expect("reader should be constructed"); + let mut reader = ErasureDecodeReader::new_with_fill_policy( + source, + engine, + data.len(), + GET_OBJECT_PATH_CODEC_STREAMING, + FillPolicy::SingleInFlight, + ) + .expect("reader should be constructed"); let mut decoded = Vec::new(); let err = reader @@ -775,7 +1113,14 @@ mod tests { read_count: None, }; let engine = LegacyEcDecodeEngine::new(erasure); - let mut reader = ErasureDecodeReader::new(source, engine, 1).expect("reader should be constructed"); + let mut reader = ErasureDecodeReader::new_with_fill_policy( + source, + engine, + 1, + GET_OBJECT_PATH_CODEC_STREAMING, + FillPolicy::SingleInFlight, + ) + .expect("reader should be constructed"); let mut decoded = Vec::new(); let err = reader @@ -797,7 +1142,14 @@ mod tests { let mut source = source_from_data(&erasure, &data, &[]); source.read_count = Some(Arc::clone(&read_count)); let engine = LegacyEcDecodeEngine::new(erasure); - let mut reader = ErasureDecodeReader::new(source, engine, data.len()).expect("reader should be constructed"); + let mut reader = ErasureDecodeReader::new_with_fill_policy( + source, + engine, + data.len(), + GET_OBJECT_PATH_CODEC_STREAMING, + FillPolicy::SingleInFlight, + ) + .expect("reader should be constructed"); let mut first_read = [0u8; 1]; let read = reader.read(&mut first_read).await.expect("first read should succeed"); diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index dddc8c898..d2286489c 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -2294,7 +2294,9 @@ impl SetDisks { part_length, metrics_path, )?; - Ok(Box::new(crate::erasure::coding::decode_reader::SyncErasureDecodeReader::new(reader))) + Ok(Box::new( + crate::erasure::coding::decode_reader::SyncErasureDecodeReader::new_with_metrics_path(reader, metrics_path), + )) } } diff --git a/crates/io-metrics/src/lib.rs b/crates/io-metrics/src/lib.rs index 1d6638547..efba48a27 100644 --- a/crates/io-metrics/src/lib.rs +++ b/crates/io-metrics/src/lib.rs @@ -624,6 +624,57 @@ pub fn record_get_object_reader_prefetch_wait(path: &'static str, duration_secs: histogram!("rustfs_io_get_object_reader_prefetch_wait_seconds", "path" => path).record(duration_secs); } +/// Record how many decoded fills were queued ahead of the current output. +#[inline(always)] +pub fn record_get_object_fill_queued(path: &'static str, policy: &'static str, queued: usize) { + if !get_stage_metrics_enabled() { + return; + } + counter!("rustfs_io_get_object_fill_queued_total", "path" => path, "policy" => policy) + .increment(u64::try_from(queued).unwrap_or(u64::MAX)); + histogram!("rustfs_io_get_object_fill_queued", "path" => path, "policy" => policy).record(usize_to_f64(queued)); +} + +/// Record that a fill completed while the current output buffer still had unread bytes. +#[inline(always)] +pub fn record_get_object_fill_completed_before_output_drained(path: &'static str, policy: &'static str) { + if !get_stage_metrics_enabled() { + return; + } + counter!("rustfs_io_get_object_fill_completed_before_output_drained_total", "path" => path, "policy" => policy).increment(1); +} + +/// Record how long output polling waited on the fill pipeline. +#[inline(always)] +pub fn record_get_object_fill_waited_by_output(path: &'static str, policy: &'static str, duration_secs: f64) { + if !get_stage_metrics_enabled() { + return; + } + counter!("rustfs_io_get_object_fill_waited_by_output_total", "path" => path, "policy" => policy).increment(1); + histogram!("rustfs_io_get_object_fill_waited_by_output_seconds", "path" => path, "policy" => policy).record(duration_secs); +} + +/// Record that a background fill task was cancelled during reader drop. +#[inline(always)] +pub fn record_get_object_fill_cancelled_on_drop(path: &'static str, policy: &'static str) { + if !get_stage_metrics_enabled() { + return; + } + counter!("rustfs_io_get_object_fill_cancelled_on_drop_total", "path" => path, "policy" => policy).increment(1); +} + +/// Record bytes staged into the prefetch queue. +#[inline(always)] +pub fn record_get_object_reader_prefetch_bytes(path: &'static str, policy: &'static str, bytes: usize) { + if !get_stage_metrics_enabled() { + return; + } + let bytes = u64::try_from(bytes).unwrap_or(u64::MAX); + counter!("rustfs_io_get_object_reader_prefetch_bytes_total", "path" => path, "policy" => policy).increment(bytes); + histogram!("rustfs_io_get_object_reader_prefetch_bytes", "path" => path, "policy" => policy) + .record(usize_to_f64(bytes as usize)); +} + /// Record one underlying shard read attempt for GetObject read-path attribution. #[inline(always)] pub fn record_get_object_shard_read( @@ -1490,6 +1541,17 @@ mod tests { assert!(0.005_f64.is_sign_positive()); } + #[test] + fn test_record_get_object_fill_metrics() { + record_get_object_fill_queued("codec_streaming", "single_inflight", 1); + record_get_object_fill_completed_before_output_drained("codec_streaming", "single_inflight"); + record_get_object_fill_waited_by_output("codec_streaming", "single_inflight", 0.0003); + record_get_object_fill_cancelled_on_drop("codec_streaming", "single_inflight"); + record_get_object_reader_prefetch_bytes("codec_streaming", "single_inflight", 4096); + + assert!(0.0003_f64.is_sign_positive()); + } + #[test] fn test_record_put_object() { record_put_object(200.0, 1024 * 1024, true); diff --git a/scripts/run_get_codec_streaming_smoke.sh b/scripts/run_get_codec_streaming_smoke.sh index 8cc134640..d67448e68 100755 --- a/scripts/run_get_codec_streaming_smoke.sh +++ b/scripts/run_get_codec_streaming_smoke.sh @@ -25,6 +25,7 @@ RETRY_PER_ROUND=1 ROUND_COOLDOWN_SECS=0 MODE="both" CODEC_ENGINES="legacy" +CODEC_MAX_INFLIGHT=1 OUT_DIR="" RUSTFS_BIN="${PROJECT_ROOT}/target/release/rustfs" WARP_BIN="warp" @@ -51,6 +52,7 @@ Purpose: Core options: --mode Which profile(s) to run (default: both) --codec-engine Codec engine(s) for codec profiles (default: legacy) + --codec-max-inflight RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT (default: 1) --address RustFS listen address (default: 127.0.0.1:19030) --bucket Benchmark bucket (default: rustfs-get-codec-smoke) --sizes Object sizes (default: 1MiB,4MiB,10MiB) @@ -138,6 +140,7 @@ parse_args() { case "$1" in --mode) MODE="$2"; shift 2 ;; --codec-engine) CODEC_ENGINES="$2"; shift 2 ;; + --codec-max-inflight) CODEC_MAX_INFLIGHT="$2"; shift 2 ;; --address) ADDRESS="$2"; shift 2 ;; --bucket) BUCKET="$2"; shift 2 ;; --sizes) SIZES="$2"; shift 2 ;; @@ -190,6 +193,7 @@ validate_args() { [[ -n "$BUCKET" ]] || die "--bucket must not be empty" [[ -n "$SIZES" ]] || die "--sizes must not be empty" validate_positive_int "$CONCURRENCY" "--concurrency" + validate_positive_int "$CODEC_MAX_INFLIGHT" "--codec-max-inflight" validate_positive_int "$ROUNDS" "--rounds" validate_positive_int "$RETRY_PER_ROUND" "--retry-per-round" validate_non_negative_int "$ROUND_COOLDOWN_SECS" "--round-cooldown-secs" @@ -301,6 +305,7 @@ compat_object_size=${COMPAT_OBJECT_SIZE} rustfs_volumes=${volumes} RUSTFS_GET_CODEC_STREAMING_ENABLE=${codec_enabled} RUSTFS_GET_CODEC_STREAMING_ENGINE=${codec_engine} +RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT=${CODEC_MAX_INFLIGHT} RUSTFS_GET_CODEC_STREAMING_MIN_SIZE=${CODEC_MIN_SIZE} metrics_path=${metrics_path} RUSTFS_SCANNER_ENABLED=false @@ -378,6 +383,7 @@ start_server() { export RUSTFS_CONSOLE_ENABLE=false export RUSTFS_GET_CODEC_STREAMING_ENABLE="$codec_enabled" export RUSTFS_GET_CODEC_STREAMING_ENGINE="$codec_engine" + export RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT="$CODEC_MAX_INFLIGHT" export RUSTFS_GET_CODEC_STREAMING_MIN_SIZE="$CODEC_MIN_SIZE" export RUSTFS_REGION="$REGION" export RUSTFS_RPC_SECRET="rustfs-get-codec-smoke-rpc-secret" @@ -661,7 +667,7 @@ write_root_compat_summary() { local codec_snapshots=() local profile snapshot - [[ -f "$legacy_snapshot" ]] || return + [[ -f "$legacy_snapshot" ]] || return 0 for profile in "$@"; do [[ "$profile" == legacy ]] && continue