mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-02 05:14:35 +08:00
fix(ecstore): retain read quota and control reserve test hedging (#8211)
This commit is contained in:
@@ -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();
|
||||
|
||||
@@ -41,6 +41,8 @@ pub struct MonitoredReader<R> {
|
||||
m: Arc<Monitor>,
|
||||
opts: MonitorReaderOptions,
|
||||
wait_state: std::sync::Mutex<Option<WaitState>>,
|
||||
// Payload limit and unpaid tokens for the read retained across Pending polls.
|
||||
pending_read: Option<(usize, u64)>,
|
||||
temp_buf: Vec<u8>,
|
||||
}
|
||||
|
||||
@@ -63,6 +65,7 @@ impl<R> MonitoredReader<R> {
|
||||
m,
|
||||
opts,
|
||||
wait_state: std::sync::Mutex::new(None),
|
||||
pending_read: None,
|
||||
temp_buf: Vec::new(),
|
||||
}
|
||||
}
|
||||
@@ -90,17 +93,30 @@ impl<R: AsyncRead + Unpin> AsyncRead for MonitoredReader<R> {
|
||||
|
||||
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<R: AsyncRead + Unpin> AsyncRead for MonitoredReader<R> {
|
||||
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 {
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user