diff --git a/crates/ecstore/src/bucket/bandwidth/monitor.rs b/crates/ecstore/src/bucket/bandwidth/monitor.rs index e99c81385..0a7d7a07f 100644 --- a/crates/ecstore/src/bucket/bandwidth/monitor.rs +++ b/crates/ecstore/src/bucket/bandwidth/monitor.rs @@ -495,6 +495,193 @@ mod tests { assert!(rate > 0.0); } + // Tokio's paused clock does not advance ratelimit's StdClock. Keep natural + // refill below one whole token for years and supply each exact budget here. + fn supply_reader_test_tokens(throttle: &BucketThrottle, available: u64) { + let rate = u64::try_from(throttle.node_bandwidth_per_sec).expect("test throttle rate should be positive"); + let burst = throttle.burst(); + *throttle.limiter.lock().expect("test throttle should not poison") = Ratelimiter::builder(rate) + .max_tokens(burst) + .initial_available(available) + .period(Duration::from_nanos(u64::MAX)) + .build() + .expect("controlled token budget should be valid"); + } + + #[tokio::test(start_paused = true)] + async fn test_monitored_reader_retains_partial_tokens_across_waits() { + use crate::bucket::bandwidth::reader::{MonitorReaderOptions, MonitoredReader}; + use futures_util::task::noop_waker_ref; + use std::pin::Pin; + use std::task::{Context, Poll}; + use tokio::io::{AsyncRead, ReadBuf}; + + for header_size in [0, 25] { + let monitor = Monitor::new(1); + let opts = BucketOptions { + name: "reader-budget".to_string(), + replication_arn: "reader-budget-arn".to_string(), + }; + monitor.set_bandwidth_limit(&opts.name, &opts.replication_arn, 100); + let throttle = monitor.throttle(&opts).expect("test throttle should exist"); + supply_reader_test_tokens(&throttle, 25); + let data = [0xAB; 100]; + let mut reader = MonitoredReader::new( + monitor, + data.as_slice(), + MonitorReaderOptions { + bucket_options: opts, + header_size, + }, + ); + let mut output = [0; 100]; + let mut buf = ReadBuf::new(&mut output); + let mut cx = Context::from_waker(noop_waker_ref()); + + assert!(Pin::new(&mut reader).poll_read(&mut cx, &mut buf).is_pending()); + assert!(buf.filled().is_empty(), "partial payment must not authorize payload reads"); + assert_eq!(throttle.limiter.lock().expect("test throttle").available(), 0); + + supply_reader_test_tokens(&throttle, 75); + tokio::time::advance(Duration::from_secs(1)).await; + assert!( + matches!(Pin::new(&mut reader).poll_read(&mut cx, &mut buf), Poll::Ready(Ok(()))), + "25 paid tokens plus 75 later tokens must fund the original read with header_size={header_size}" + ); + assert_eq!(buf.filled(), &data[..100 - header_size]); + assert_eq!(throttle.limiter.lock().expect("test throttle").available(), 0); + } + } + + #[tokio::test(start_paused = true)] + async fn test_monitored_reader_preserves_paid_read_across_cancellation() { + use crate::bucket::bandwidth::reader::{MonitorReaderOptions, MonitoredReader}; + use futures_util::task::noop_waker_ref; + use std::future::Future; + use std::pin::Pin; + use std::task::{Context, Poll}; + use tokio::io::{AsyncRead, AsyncReadExt, AsyncWriteExt, ReadBuf}; + + let monitor = Monitor::new(1); + let opts = BucketOptions { + name: "reader-pending".to_string(), + replication_arn: "reader-pending-arn".to_string(), + }; + monitor.set_bandwidth_limit(&opts.name, &opts.replication_arn, 100); + let throttle = monitor.throttle(&opts).expect("test throttle should exist"); + supply_reader_test_tokens(&throttle, 100); + let (mut writer, inner) = tokio::io::duplex(100); + let mut reader = MonitoredReader::new( + monitor, + inner, + MonitorReaderOptions { + bucket_options: opts, + header_size: 0, + }, + ); + let mut cx = Context::from_waker(noop_waker_ref()); + let mut first_output = [0; 100]; + let mut first_read = Box::pin(reader.read(&mut first_output)); + assert!(first_read.as_mut().poll(&mut cx).is_pending()); + assert_eq!(throttle.limiter.lock().expect("test throttle").available(), 0); + drop(first_read); + writer + .write_all(&[0xAB; 100]) + .await + .expect("duplex writer should supply the pending read"); + + // A canceled read future keeps its reader alive. Resume in a smaller + // buffer without charging the already-paid allowance a second time. + let mut short_output = [0; 40]; + let mut short_buf = ReadBuf::new(&mut short_output); + assert!(matches!(Pin::new(&mut reader).poll_read(&mut cx, &mut short_buf), Poll::Ready(Ok(())))); + assert_eq!(short_buf.filled(), &[0xAB; 40]); + assert_eq!(throttle.limiter.lock().expect("test throttle").available(), 0); + + // Completing the short read discards unused allowance, matching the + // existing precharge policy. The next read must buy a fresh budget. + let mut remaining_output = [0; 100]; + let mut remaining_buf = ReadBuf::new(&mut remaining_output); + assert!(Pin::new(&mut reader).poll_read(&mut cx, &mut remaining_buf).is_pending()); + supply_reader_test_tokens(&throttle, 100); + tokio::time::advance(Duration::from_secs(1)).await; + assert!(matches!( + Pin::new(&mut reader).poll_read(&mut cx, &mut remaining_buf), + Poll::Ready(Ok(())) + )); + assert_eq!(remaining_buf.filled(), &[0xAB; 60]); + assert_eq!(throttle.limiter.lock().expect("test throttle").available(), 0); + + drop(writer); + let mut eof_output = [0; 100]; + let mut eof_buf = ReadBuf::new(&mut eof_output); + assert!(Pin::new(&mut reader).poll_read(&mut cx, &mut eof_buf).is_pending()); + supply_reader_test_tokens(&throttle, 100); + tokio::time::advance(Duration::from_secs(1)).await; + assert!(matches!(Pin::new(&mut reader).poll_read(&mut cx, &mut eof_buf), Poll::Ready(Ok(())))); + assert!(eof_buf.filled().is_empty()); + assert_eq!(throttle.limiter.lock().expect("test throttle").available(), 0); + assert!(Pin::new(&mut reader).poll_read(&mut cx, &mut eof_buf).is_pending()); + } + + #[tokio::test(start_paused = true)] + async fn test_monitored_reader_shrinks_unpaid_payload_after_bandwidth_reduction() { + use crate::bucket::bandwidth::reader::{MonitorReaderOptions, MonitoredReader}; + use futures_util::task::noop_waker_ref; + use std::pin::Pin; + use std::task::{Context, Poll}; + use tokio::io::{AsyncRead, ReadBuf}; + + for header_size in [0, 25] { + let monitor = Monitor::new(1); + let opts = BucketOptions { + name: "reader-reduced-budget".to_string(), + replication_arn: "reader-reduced-budget-arn".to_string(), + }; + monitor.set_bandwidth_limit(&opts.name, &opts.replication_arn, 100); + let initial_throttle = monitor.throttle(&opts).expect("initial throttle should exist"); + supply_reader_test_tokens(&initial_throttle, 25); + let data = [0xAB; 100]; + let mut reader = MonitoredReader::new( + monitor.clone(), + data.as_slice(), + MonitorReaderOptions { + bucket_options: opts.clone(), + header_size, + }, + ); + let mut output = [0; 100]; + let mut buf = ReadBuf::new(&mut output); + let mut cx = Context::from_waker(noop_waker_ref()); + assert!(Pin::new(&mut reader).poll_read(&mut cx, &mut buf).is_pending()); + assert!(buf.filled().is_empty()); + assert_eq!(initial_throttle.limiter.lock().expect("initial throttle").available(), 0); + + monitor.set_bandwidth_limit(&opts.name, &opts.replication_arn, 10); + let reduced_throttle = monitor.throttle(&opts).expect("reduced throttle should exist"); + assert_eq!(reduced_throttle.burst(), 10); + supply_reader_test_tokens(&reduced_throttle, 0); + tokio::time::advance(Duration::from_secs(1)).await; + let result = Pin::new(&mut reader).poll_read(&mut cx, &mut buf); + let expected_payload = if header_size == 0 { + assert!( + matches!(result, Poll::Ready(Ok(()))), + "already-paid payload must not wait for the old burst's unpaid remainder" + ); + 25 + } else { + assert!(result.is_pending(), "the initial 25 tokens paid only the header"); + assert!(buf.filled().is_empty(), "payload still requires its own reduced budget"); + supply_reader_test_tokens(&reduced_throttle, 10); + tokio::time::advance(Duration::from_secs(1)).await; + assert!(matches!(Pin::new(&mut reader).poll_read(&mut cx, &mut buf), Poll::Ready(Ok(())))); + 10 + }; + assert_eq!(buf.filled(), &data[..expected_payload]); + assert_eq!(reduced_throttle.limiter.lock().expect("reduced throttle").available(), 0); + } + } + #[test] fn test_consume_refills_continuously() { let clock = TestClock::new(); diff --git a/crates/ecstore/src/bucket/bandwidth/reader.rs b/crates/ecstore/src/bucket/bandwidth/reader.rs index 9bf27bfda..644dd1689 100644 --- a/crates/ecstore/src/bucket/bandwidth/reader.rs +++ b/crates/ecstore/src/bucket/bandwidth/reader.rs @@ -41,6 +41,8 @@ pub struct MonitoredReader { m: Arc, opts: MonitorReaderOptions, wait_state: std::sync::Mutex>, + // Payload limit and unpaid tokens for the read retained across Pending polls. + pending_read: Option<(usize, u64)>, temp_buf: Vec, } @@ -63,6 +65,7 @@ impl MonitoredReader { m, opts, wait_state: std::sync::Mutex::new(None), + pending_read: None, temp_buf: Vec::new(), } } @@ -90,17 +93,30 @@ impl AsyncRead for MonitoredReader { let throttle = match this.m.throttle(&this.opts.bucket_options) { Some(t) => t, - None => return Pin::new(&mut this.r).poll_read(cx, buf), + None => { + this.pending_read = None; + return Pin::new(&mut this.r).poll_read(cx, buf); + } }; let b = throttle.burst(); debug_assert!(b >= 1, "burst must be at least 1"); - let (need, tokens) = calc_need_and_tokens(b, buf.remaining(), &mut this.opts.header_size); - let (deficit, rate, consumed) = throttle.consume(tokens); - let need = need.min(consumed as usize); + let (need, tokens) = match this.pending_read { + Some(reservation) => reservation, + None => calc_need_and_tokens(b, buf.remaining(), &mut this.opts.header_size), + }; + // A lower limit may discard unpaid payload, but must retain paid + // payload and any header tokens still owed by this pending read. + let excess_payload = need.saturating_sub(usize::try_from(b).unwrap_or(usize::MAX)); + let reduction = tokens.min(u64::try_from(excess_payload).unwrap_or(u64::MAX)); + let need = need.saturating_sub(usize::try_from(reduction).unwrap_or(usize::MAX)); + let tokens = tokens - reduction; + let (deficit, rate, _) = throttle.consume(tokens); + this.pending_read = Some((need, deficit)); if deficit > 0 && rate > 0.0 { - let duration = std::time::Duration::from_secs_f64(deficit as f64 / rate); + // A changed limit may reduce the burst while this read is pending. + let duration = std::time::Duration::from_secs_f64(deficit.min(b) as f64 / rate); debug!( tokens = tokens, deficit = deficit, @@ -115,14 +131,17 @@ impl AsyncRead for MonitoredReader { warn!("MonitoredReader wait_state mutex poisoned, recovering"); e.into_inner() }) = Some(WaitState { sleep }); - return Poll::Pending; } - Poll::Ready(()) => {} + Poll::Ready(()) => cx.waker().wake_by_ref(), } + return Poll::Pending; } let filled_before = buf.filled().len(); let result = poll_limited_read(&mut this.r, cx, buf, need, &mut this.temp_buf); + if result.is_ready() { + this.pending_read = None; + } if let Poll::Ready(Ok(())) = result { let read_bytes = buf.filled().len().saturating_sub(filled_before) as u64; if read_bytes > 0 { diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index b9ea895b6..b4e8db506 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -2884,7 +2884,16 @@ impl SetDisks { let result = if let Some(deadline) = single_pending_hedge_deadline.take() { tokio::select! { result = join_set.join_next() => result, - _ = tokio::time::sleep_until(deadline) => { + _ = async { + tokio::time::sleep_until(deadline).await; + #[cfg(test)] + rename_fanout_barrier::checkpoint( + object.as_ref(), + 0, + rename_fanout_barrier::PHASE_NON_INLINE_HEDGE_TIMER, + ) + .await; + } => { if bounded_fanout && !force_full_wait && join_set.len() == 1 @@ -7322,6 +7331,9 @@ pub(crate) mod rename_fanout_barrier { CLEANUP as PHASE_CLEANUP, READ_VERSION as PHASE_READ_VERSION, RENAME as PHASE_RENAME, ROLLBACK as PHASE_ROLLBACK, }; + /// Object-scoped hedge timer checkpoint; slot zero identifies the timer, not a disk. + pub const PHASE_NON_INLINE_HEDGE_TIMER: &str = "non_inline_hedge_timer"; + /// One armed barrier: the fan-out task matching `(disk_index, phase)` pauses. struct Armed { disk_index: usize, diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index e6059f6cd..fb18239a8 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -1593,6 +1593,8 @@ mod prepared_get_object_metadata_tests { ("RUSTFS_GET_METADATA_EARLY_STOP_BOUNDED_FANOUT", Some("true")), ], async { + // Hold this object's hedge timer so real-disk latency cannot add speculative fanout. + let _hedge_timer = rename_fanout_barrier::arm(&object, 0, rename_fanout_barrier::PHASE_NON_INLINE_HEDGE_TIMER); let calls = disk_call_counters::observe(&object); let mut reader = set_disks .get_object_reader(bucket, &object, None, HeaderMap::new(), &opts)