mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-13 08:36:54 +00:00
perf: harden HotPath closure evidence (#5573)
* test(ecstore): cover raw shard write errors Co-Authored-By: heihutu <heihutu@gmail.com> * perf(ecstore): track encode payload stage peaks Co-Authored-By: heihutu <heihutu@gmail.com> * test(perf): harden formal warp ABBA evidence Co-Authored-By: heihutu <heihutu@gmail.com> --------- Co-authored-by: heihutu <heihutu@gmail.com> Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -641,6 +641,7 @@ impl Erasure {
|
||||
let res = self.clone().encode_block_bytes_mut(encode_buf, n).await?;
|
||||
buf = BytesMut::with_capacity(ingest_capacity);
|
||||
let queued_bytes = queued_block_bytes(&res);
|
||||
let _producer_stage = rustfs_io_metrics::track_ec_encode_producer_bytes(queued_bytes);
|
||||
let send_wait_stage_start = stage_timer_if_enabled();
|
||||
if let Err(err) = send_queued(&tx, res, queued_bytes).await {
|
||||
return Err(std::io::Error::other(format!("Failed to send encoded data : {err}")));
|
||||
@@ -670,6 +671,7 @@ impl Erasure {
|
||||
let (res, returned_buf) = self.clone().encode_block(encode_buf, n).await?;
|
||||
buf = returned_buf;
|
||||
let queued_bytes = queued_block_bytes(&res);
|
||||
let _producer_stage = rustfs_io_metrics::track_ec_encode_producer_bytes(queued_bytes);
|
||||
let send_wait_stage_start = stage_timer_if_enabled();
|
||||
if let Err(err) = send_queued(&tx, res, queued_bytes).await {
|
||||
return Err(std::io::Error::other(format!("Failed to send encoded data : {err}")));
|
||||
@@ -712,6 +714,7 @@ impl Erasure {
|
||||
if block.is_empty() {
|
||||
break;
|
||||
}
|
||||
let _writer_stage = rustfs_io_metrics::track_ec_encode_writer_bytes(queued_block_bytes(&block));
|
||||
let write_stage_start = stage_timer_if_enabled();
|
||||
if let Err(err) = writers.write(block).await {
|
||||
write_err = Some(err);
|
||||
@@ -768,6 +771,7 @@ impl Erasure {
|
||||
let mut buf = vec![0u8; block_size];
|
||||
let mut pending_batch = Vec::with_capacity(batch_blocks);
|
||||
let mut pending_batch_bytes = 0usize;
|
||||
let mut pending_batch_stage = None;
|
||||
loop {
|
||||
match rustfs_utils::read_full_or_eof(&mut reader, &mut buf).await {
|
||||
Ok(Some(n)) => {
|
||||
@@ -779,6 +783,8 @@ impl Erasure {
|
||||
let queued_bytes = queued_block_bytes(&res);
|
||||
pending_batch_bytes = pending_batch_bytes.saturating_add(queued_bytes);
|
||||
pending_batch.push(res);
|
||||
drop(pending_batch_stage.take());
|
||||
pending_batch_stage = Some(rustfs_io_metrics::track_ec_encode_producer_bytes(pending_batch_bytes));
|
||||
|
||||
if pending_batch.len() >= batch_blocks {
|
||||
let send_wait_stage_start = stage_timer_if_enabled();
|
||||
@@ -786,6 +792,7 @@ impl Erasure {
|
||||
return Err(std::io::Error::other(format!("Failed to send encoded data : {err}")));
|
||||
}
|
||||
record_internal_stage_if_enabled("erasure_encode_batched_send_wait", send_wait_stage_start);
|
||||
drop(pending_batch_stage.take());
|
||||
pending_batch = Vec::with_capacity(batch_blocks);
|
||||
pending_batch_bytes = 0;
|
||||
}
|
||||
@@ -813,6 +820,7 @@ impl Erasure {
|
||||
return Err(std::io::Error::other(format!("Failed to send encoded data : {err}")));
|
||||
}
|
||||
record_internal_stage_if_enabled("erasure_encode_batched_send_wait", send_wait_stage_start);
|
||||
drop(pending_batch_stage);
|
||||
}
|
||||
|
||||
Ok((reader, total))
|
||||
@@ -828,6 +836,7 @@ impl Erasure {
|
||||
};
|
||||
record_internal_stage_if_enabled("erasure_encode_batched_recv_wait", recv_wait_stage_start);
|
||||
let batch = batch.into_inner();
|
||||
let _writer_stage = rustfs_io_metrics::track_ec_encode_writer_bytes(queued_batch_bytes(&batch));
|
||||
let write_stage_start = stage_timer_if_enabled();
|
||||
for block in batch {
|
||||
if let Err(err) = writers.write(block).await {
|
||||
@@ -1302,24 +1311,36 @@ mod tests {
|
||||
}
|
||||
|
||||
async fn writer_error_aborts_blocked_producer(pipeline: EncodePipeline) {
|
||||
const BLOCK_SIZE: usize = 16;
|
||||
|
||||
let gauge_baseline = rustfs_io_metrics::current_ec_encode_inflight_bytes();
|
||||
let producer_baseline = rustfs_io_metrics::current_ec_encode_producer_bytes();
|
||||
let writer_baseline = rustfs_io_metrics::current_ec_encode_writer_bytes();
|
||||
let producer_peak_before = rustfs_io_metrics::current_ec_encode_producer_bytes_peak();
|
||||
let writer_peak_before = rustfs_io_metrics::current_ec_encode_writer_bytes_peak();
|
||||
let block_size = match pipeline {
|
||||
EncodePipeline::Batched => usize::try_from(producer_peak_before.max(writer_peak_before).saturating_add(1))
|
||||
.expect("stage peak fits the test address space"),
|
||||
EncodePipeline::Vec | EncodePipeline::BytesMut => 16,
|
||||
};
|
||||
rustfs_io_metrics::set_put_stage_metrics_enabled(true);
|
||||
let erasure = Arc::new(Erasure::new(1, 0, block_size));
|
||||
let batch_blocks = encode_batch_block_count().min(encode_channel_capacity(
|
||||
erasure.shard_size().saturating_mul(erasure.total_shard_count()),
|
||||
erasure_encode_max_inflight_bytes(),
|
||||
));
|
||||
let blocks_before_pending = match pipeline {
|
||||
EncodePipeline::Batched => encode_batch_block_count(),
|
||||
EncodePipeline::Batched => batch_blocks,
|
||||
EncodePipeline::Vec | EncodePipeline::BytesMut => 1,
|
||||
};
|
||||
let (reader, reader_blocked, reader_dropped, _final_block) =
|
||||
BlocksThenPendingReader::new(blocks_before_pending, BLOCK_SIZE);
|
||||
BlocksThenPendingReader::new(blocks_before_pending, block_size);
|
||||
let writes = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let mut writers = vec![Some(bitrot_writer(
|
||||
FailAfterReaderBlocksWriter {
|
||||
reader_blocked,
|
||||
writes: writes.clone(),
|
||||
},
|
||||
BLOCK_SIZE,
|
||||
block_size,
|
||||
))];
|
||||
let erasure = Arc::new(Erasure::new(1, 0, BLOCK_SIZE));
|
||||
|
||||
let result = match pipeline {
|
||||
EncodePipeline::Vec => erasure.encode_with_ingest_mode(reader, &mut writers, 1, false).await,
|
||||
@@ -1346,12 +1367,40 @@ mod tests {
|
||||
gauge_baseline,
|
||||
"writer failure must settle all queued and pending encoded bytes"
|
||||
);
|
||||
assert_eq!(
|
||||
rustfs_io_metrics::current_ec_encode_producer_bytes(),
|
||||
producer_baseline,
|
||||
"writer failure must settle producer stage bytes for every ingest pipeline"
|
||||
);
|
||||
assert_eq!(
|
||||
rustfs_io_metrics::current_ec_encode_writer_bytes(),
|
||||
writer_baseline,
|
||||
"writer failure must settle writer stage bytes for every ingest pipeline"
|
||||
);
|
||||
if matches!(pipeline, EncodePipeline::Batched) {
|
||||
let expected_batch_bytes = u64::try_from(block_size)
|
||||
.expect("block size fits the stage gauge")
|
||||
.checked_mul(u64::try_from(batch_blocks).expect("batch block count fits the stage gauge"))
|
||||
.expect("test batch payload fits the stage gauge");
|
||||
assert!(
|
||||
rustfs_io_metrics::current_ec_encode_producer_bytes_peak() >= producer_peak_before.max(expected_batch_bytes),
|
||||
"batched producer must expose its full pending batch before writer failure"
|
||||
);
|
||||
assert!(
|
||||
rustfs_io_metrics::current_ec_encode_writer_bytes_peak() >= writer_peak_before.max(expected_batch_bytes),
|
||||
"batched writer must expose its full batch before writer failure"
|
||||
);
|
||||
}
|
||||
rustfs_io_metrics::set_put_stage_metrics_enabled(false);
|
||||
}
|
||||
|
||||
async fn aborting_full_queue_settles_pending_send() {
|
||||
const BLOCK_SIZE: usize = 16;
|
||||
|
||||
let gauge_baseline = rustfs_io_metrics::current_ec_encode_inflight_bytes();
|
||||
let producer_baseline = rustfs_io_metrics::current_ec_encode_producer_bytes();
|
||||
let writer_baseline = rustfs_io_metrics::current_ec_encode_writer_bytes();
|
||||
rustfs_io_metrics::set_put_stage_metrics_enabled(true);
|
||||
let erasure = Arc::new(Erasure::new(1, 0, BLOCK_SIZE));
|
||||
let inflight_blocks = encode_channel_capacity(
|
||||
erasure.shard_size().saturating_mul(erasure.total_shard_count()),
|
||||
@@ -1389,6 +1438,15 @@ mod tests {
|
||||
})
|
||||
.await
|
||||
.expect("producer should account for the pending send after the queue fills");
|
||||
tokio::time::timeout(Duration::from_secs(1), async {
|
||||
while rustfs_io_metrics::current_ec_encode_producer_bytes() == producer_baseline
|
||||
|| rustfs_io_metrics::current_ec_encode_writer_bytes() == writer_baseline
|
||||
{
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("full queue must retain producer and writer stage ownership before cancellation");
|
||||
|
||||
encode.abort();
|
||||
assert!(matches!(encode.await, Err(err) if err.is_cancelled()), "encode task should be cancelled");
|
||||
@@ -1406,6 +1464,17 @@ mod tests {
|
||||
gauge_baseline,
|
||||
"cancelling a full queue must settle queued and pending bytes"
|
||||
);
|
||||
assert_eq!(
|
||||
rustfs_io_metrics::current_ec_encode_producer_bytes(),
|
||||
producer_baseline,
|
||||
"cancelling a full queue must settle the pending producer stage"
|
||||
);
|
||||
assert_eq!(
|
||||
rustfs_io_metrics::current_ec_encode_writer_bytes(),
|
||||
writer_baseline,
|
||||
"cancelling a full queue must settle the stalled writer stage"
|
||||
);
|
||||
rustfs_io_metrics::set_put_stage_metrics_enabled(false);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -1872,6 +1941,61 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn bytesmut_streaming_encode_observes_all_payload_stage_peaks() {
|
||||
let queue_baseline = rustfs_io_metrics::current_ec_encode_inflight_bytes();
|
||||
let producer_peak_before = rustfs_io_metrics::current_ec_encode_producer_bytes_peak();
|
||||
let queue_peak_before = rustfs_io_metrics::current_ec_encode_queue_bytes_peak();
|
||||
let writer_peak_before = rustfs_io_metrics::current_ec_encode_writer_bytes_peak();
|
||||
let prior_peak = producer_peak_before.max(queue_peak_before).max(writer_peak_before);
|
||||
let block_size = usize::try_from(prior_peak.saturating_add(1)).expect("stage peak fits the test address space");
|
||||
let erasure = Arc::new(Erasure::new(1, 0, block_size));
|
||||
let encoded_block_bytes = erasure.shard_size() * erasure.total_shard_count();
|
||||
let committed = Arc::new(Mutex::new(Vec::new()));
|
||||
let mut writers = vec![Some(bitrot_writer(DeferredCommitWriter::new(committed.clone()), block_size))];
|
||||
let reader = tokio::io::BufReader::new(Cursor::new(vec![0x5a; block_size * 2]));
|
||||
|
||||
rustfs_io_metrics::set_put_stage_metrics_enabled(true);
|
||||
let result = erasure.encode_with_ingest_mode(reader, &mut writers, 1, true).await;
|
||||
rustfs_io_metrics::set_put_stage_metrics_enabled(false);
|
||||
|
||||
let (_reader, written) = result.expect("bytesmut streaming encode should complete");
|
||||
assert_eq!(written, block_size * 2);
|
||||
assert_eq!(
|
||||
rustfs_io_metrics::current_ec_encode_inflight_bytes(),
|
||||
queue_baseline,
|
||||
"completed streaming encode must not retain queue bytes"
|
||||
);
|
||||
assert_eq!(
|
||||
rustfs_io_metrics::current_ec_encode_producer_bytes(),
|
||||
0,
|
||||
"completed streaming encode must not retain producer bytes"
|
||||
);
|
||||
assert_eq!(
|
||||
rustfs_io_metrics::current_ec_encode_writer_bytes(),
|
||||
0,
|
||||
"completed streaming encode must not retain writer bytes"
|
||||
);
|
||||
let expected_peak = u64::try_from(encoded_block_bytes).expect("encoded block bytes fit the gauge");
|
||||
assert!(
|
||||
rustfs_io_metrics::current_ec_encode_producer_bytes_peak() >= producer_peak_before.max(expected_peak),
|
||||
"producer peak must observe encoded bytes before queue hand-off"
|
||||
);
|
||||
assert!(
|
||||
rustfs_io_metrics::current_ec_encode_queue_bytes_peak() >= queue_peak_before.max(expected_peak),
|
||||
"queue peak must observe encoded bytes pending shard writers"
|
||||
);
|
||||
assert!(
|
||||
rustfs_io_metrics::current_ec_encode_writer_bytes_peak() >= writer_peak_before.max(expected_peak),
|
||||
"writer peak must observe encoded bytes after queue hand-off"
|
||||
);
|
||||
assert!(
|
||||
!committed.lock().expect("committed buffer should be lockable").is_empty(),
|
||||
"stage peak observation must not change writer commit behavior"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn encode_streaming_write_quorum_failure_aborts_and_reports_error() {
|
||||
const DATA_SHARDS: usize = 2;
|
||||
|
||||
@@ -768,10 +768,18 @@ mod tests {
|
||||
#[derive(Debug, Clone, Default)]
|
||||
struct TestRemoteDataTransport {
|
||||
bytes: Arc<Mutex<Vec<u8>>>,
|
||||
write_error: Option<io::ErrorKind>,
|
||||
}
|
||||
|
||||
#[cfg(feature = "hotpath")]
|
||||
impl TestRemoteDataTransport {
|
||||
fn with_write_error(kind: io::ErrorKind) -> Self {
|
||||
Self {
|
||||
bytes: Arc::default(),
|
||||
write_error: Some(kind),
|
||||
}
|
||||
}
|
||||
|
||||
fn bytes(&self) -> Vec<u8> {
|
||||
self.bytes
|
||||
.lock()
|
||||
@@ -784,11 +792,15 @@ mod tests {
|
||||
#[derive(Debug)]
|
||||
struct TestRemoteWriter {
|
||||
bytes: Arc<Mutex<Vec<u8>>>,
|
||||
write_error: Option<io::ErrorKind>,
|
||||
}
|
||||
|
||||
#[cfg(feature = "hotpath")]
|
||||
impl tokio::io::AsyncWrite for TestRemoteWriter {
|
||||
fn poll_write(self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &[u8]) -> Poll<io::Result<usize>> {
|
||||
if let Some(kind) = self.write_error {
|
||||
return Poll::Ready(Err(io::Error::from(kind)));
|
||||
}
|
||||
self.bytes
|
||||
.lock()
|
||||
.expect("test remote transport bytes lock should not be poisoned")
|
||||
@@ -815,6 +827,7 @@ mod tests {
|
||||
async fn open_write(&self, _request: WriteStreamRequest) -> Result<FileWriter> {
|
||||
Ok(Box::new(TestRemoteWriter {
|
||||
bytes: Arc::clone(&self.bytes),
|
||||
write_error: self.write_error,
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -855,7 +868,7 @@ mod tests {
|
||||
}
|
||||
|
||||
#[cfg(feature = "hotpath")]
|
||||
async fn remote_test_disk() -> (DiskStore, TestRemoteDataTransport) {
|
||||
async fn remote_test_disk(transport: TestRemoteDataTransport) -> (DiskStore, TestRemoteDataTransport) {
|
||||
use crate::disk::endpoint::Endpoint;
|
||||
|
||||
let endpoint = Endpoint {
|
||||
@@ -865,7 +878,6 @@ mod tests {
|
||||
set_idx: 0,
|
||||
disk_idx: 0,
|
||||
};
|
||||
let transport = TestRemoteDataTransport::default();
|
||||
let remote = RemoteDisk::new(
|
||||
&endpoint,
|
||||
&DiskOption {
|
||||
@@ -1021,7 +1033,7 @@ mod tests {
|
||||
.expect("mmap-copy fallback stream should preserve bytes");
|
||||
assert_eq!(fallback_read, fallback_payload);
|
||||
|
||||
let (remote_disk, remote_transport) = remote_test_disk().await;
|
||||
let (remote_disk, remote_transport) = remote_test_disk(TestRemoteDataTransport::default()).await;
|
||||
let remote_payload = b"remote shard bytes";
|
||||
let mut remote_writer = create_bitrot_writer(
|
||||
false,
|
||||
@@ -1074,6 +1086,43 @@ mod tests {
|
||||
}
|
||||
assert_eq!(remote_read, remote_payload);
|
||||
|
||||
let (failing_remote_disk, _) = remote_test_disk(TestRemoteDataTransport::with_write_error(io::ErrorKind::Other)).await;
|
||||
let mut failing_writer = create_bitrot_writer(
|
||||
false,
|
||||
Some(&failing_remote_disk),
|
||||
bucket,
|
||||
"obj/hotpath-remote-write-error-part.1",
|
||||
4,
|
||||
shard_size,
|
||||
HashAlgorithm::None,
|
||||
)
|
||||
.await
|
||||
.expect("failing remote bitrot writer should open before its first write");
|
||||
let write_error = failing_writer
|
||||
.write(b"fail")
|
||||
.await
|
||||
.expect_err("raw shard writer failures must remain visible through the bitrot writer");
|
||||
assert_eq!(write_error.kind(), io::ErrorKind::Other);
|
||||
|
||||
let (would_block_remote_disk, _) =
|
||||
remote_test_disk(TestRemoteDataTransport::with_write_error(io::ErrorKind::WouldBlock)).await;
|
||||
let mut would_block_writer = create_bitrot_writer(
|
||||
false,
|
||||
Some(&would_block_remote_disk),
|
||||
bucket,
|
||||
"obj/hotpath-remote-write-would-block-part.1",
|
||||
4,
|
||||
shard_size,
|
||||
HashAlgorithm::None,
|
||||
)
|
||||
.await
|
||||
.expect("would-block remote bitrot writer should open before its first write");
|
||||
let would_block = would_block_writer
|
||||
.write(b"wait")
|
||||
.await
|
||||
.expect_err("would-block must remain visible to the caller");
|
||||
assert_eq!(would_block.kind(), io::ErrorKind::WouldBlock);
|
||||
|
||||
drop(guard);
|
||||
let report = std::fs::read_to_string(&report_path).expect("HotPath I/O report should be written");
|
||||
let report: serde_json::Value = serde_json::from_str(&report).expect("HotPath I/O report should be valid JSON");
|
||||
@@ -1089,6 +1138,13 @@ mod tests {
|
||||
assert!(byte_count > 0, "report must include fixed label {label}");
|
||||
byte_count
|
||||
};
|
||||
let io_errors = |label: &str, direction: &str| {
|
||||
entries
|
||||
.iter()
|
||||
.filter(|entry| entry["label"].as_str().is_some_and(|entry_label| entry_label == label))
|
||||
.filter_map(|entry| entry[direction]["errors"].as_u64())
|
||||
.sum::<u64>()
|
||||
};
|
||||
let payload_len = u64::try_from(payload.len()).expect("test payload length should fit u64");
|
||||
assert_eq!(
|
||||
io_bytes(RAW_SHARD_READ_LOCAL_LABEL, "read"),
|
||||
@@ -1103,6 +1159,11 @@ mod tests {
|
||||
io_bytes(RAW_SHARD_WRITE_REMOTE_LABEL, "write"),
|
||||
u64::try_from(remote_payload.len()).expect("remote payload length should fit u64")
|
||||
);
|
||||
assert_eq!(
|
||||
io_errors(RAW_SHARD_WRITE_REMOTE_LABEL, "write"),
|
||||
1,
|
||||
"the wrapper must record the real remote writer failure without changing it, while excluding WouldBlock"
|
||||
);
|
||||
for label in [
|
||||
RAW_SHARD_READ_LOCAL_LABEL,
|
||||
RAW_SHARD_READ_REMOTE_LABEL,
|
||||
|
||||
@@ -49,7 +49,10 @@
|
||||
#[macro_use]
|
||||
extern crate metrics;
|
||||
|
||||
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
||||
use std::sync::{
|
||||
Mutex,
|
||||
atomic::{AtomicBool, AtomicU64, Ordering},
|
||||
};
|
||||
|
||||
/// Global switch for detailed per-stage PUT metrics (path label, stage durations).
|
||||
/// When `false`, `record_put_object_path` and `record_put_object_stage_duration`
|
||||
@@ -269,6 +272,12 @@ pub use collector::MetricsCollector;
|
||||
pub use performance::PerformanceMetrics;
|
||||
|
||||
static EC_ENCODE_INFLIGHT_BYTES: AtomicU64 = AtomicU64::new(0);
|
||||
static EC_ENCODE_PRODUCER_BYTES_CURRENT: AtomicU64 = AtomicU64::new(0);
|
||||
static EC_ENCODE_PRODUCER_BYTES_PEAK: AtomicU64 = AtomicU64::new(0);
|
||||
static EC_ENCODE_QUEUE_BYTES_PEAK: AtomicU64 = AtomicU64::new(0);
|
||||
static EC_ENCODE_WRITER_BYTES_CURRENT: AtomicU64 = AtomicU64::new(0);
|
||||
static EC_ENCODE_WRITER_BYTES_PEAK: AtomicU64 = AtomicU64::new(0);
|
||||
static EC_ENCODE_PEAK_PUBLISH_LOCK: Mutex<()> = Mutex::new(());
|
||||
static GET_OBJECT_BUFFERED_BYTES: AtomicU64 = AtomicU64::new(0);
|
||||
const SHARD_READ_COST_LOCAL: &str = "local";
|
||||
const SHARD_READ_COST_REMOTE: &str = "remote";
|
||||
@@ -313,6 +322,49 @@ fn saturating_sub_atomic(counter: &AtomicU64, bytes: u64) -> u64 {
|
||||
}
|
||||
}
|
||||
|
||||
#[inline(always)]
|
||||
fn update_peak_atomic(counter: &AtomicU64, value: u64) -> Option<u64> {
|
||||
let mut peak = counter.load(Ordering::Relaxed);
|
||||
loop {
|
||||
if value <= peak {
|
||||
return None;
|
||||
}
|
||||
match counter.compare_exchange_weak(peak, value, Ordering::Relaxed, Ordering::Relaxed) {
|
||||
Ok(_) => return Some(value),
|
||||
Err(actual) => peak = actual,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
enum EcEncodePeakMetric {
|
||||
Producer,
|
||||
Queue,
|
||||
Writer,
|
||||
}
|
||||
|
||||
fn publish_ec_encode_peak_with(counter: &AtomicU64, value: u64, before_publish_lock: impl FnOnce(), set_gauge: impl FnOnce(f64)) {
|
||||
if update_peak_atomic(counter, value).is_none() {
|
||||
return;
|
||||
}
|
||||
before_publish_lock();
|
||||
let _guard = EC_ENCODE_PEAK_PUBLISH_LOCK.lock().unwrap_or_else(|error| error.into_inner());
|
||||
set_gauge(counter.load(Ordering::Relaxed) as f64);
|
||||
}
|
||||
|
||||
fn publish_ec_encode_peak(counter: &AtomicU64, metric: EcEncodePeakMetric, value: u64) {
|
||||
publish_ec_encode_peak_with(
|
||||
counter,
|
||||
value,
|
||||
|| {},
|
||||
|peak| match metric {
|
||||
EcEncodePeakMetric::Producer => gauge_set_cached!("rustfs_ec_encode_producer_bytes_peak", peak),
|
||||
EcEncodePeakMetric::Queue => gauge_set_cached!("rustfs_ec_encode_queue_bytes_peak", peak),
|
||||
EcEncodePeakMetric::Writer => gauge_set_cached!("rustfs_ec_encode_writer_bytes_peak", peak),
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
#[inline(always)]
|
||||
fn usize_to_f64(value: usize) -> f64 {
|
||||
value as f64
|
||||
@@ -2204,6 +2256,9 @@ pub fn record_allocator_memory_observation(backend: &'static str, observation: A
|
||||
pub fn add_ec_encode_inflight_bytes(bytes: usize) {
|
||||
let next = EC_ENCODE_INFLIGHT_BYTES.fetch_add(bytes as u64, Ordering::Relaxed) + bytes as u64;
|
||||
gauge_set_cached!("rustfs_ec_encode_inflight_bytes_current", next as f64);
|
||||
if put_stage_metrics_enabled() {
|
||||
publish_ec_encode_peak(&EC_ENCODE_QUEUE_BYTES_PEAK, EcEncodePeakMetric::Queue, next);
|
||||
}
|
||||
}
|
||||
|
||||
/// Remove encoded bytes from the tracked erasure encode in-flight gauge.
|
||||
@@ -2219,6 +2274,109 @@ pub fn current_ec_encode_inflight_bytes() -> u64 {
|
||||
EC_ENCODE_INFLIGHT_BYTES.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Tracks encoded payload bytes held before queue hand-off or during shard writes.
|
||||
///
|
||||
/// Each guard contributes to a process-wide stage total until it is dropped. The
|
||||
/// reported peak therefore includes concurrent PUTs, but excludes reader,
|
||||
/// allocator, and transport buffers; it is not a per-PUT or process-RSS limit.
|
||||
pub struct EcEncodePayloadStageGuard {
|
||||
counter: &'static AtomicU64,
|
||||
bytes: u64,
|
||||
enabled: bool,
|
||||
current_metric: EcEncodePeakMetric,
|
||||
}
|
||||
|
||||
impl Drop for EcEncodePayloadStageGuard {
|
||||
fn drop(&mut self) {
|
||||
if !self.enabled {
|
||||
return;
|
||||
}
|
||||
let next = saturating_sub_atomic(self.counter, self.bytes);
|
||||
match self.current_metric {
|
||||
EcEncodePeakMetric::Producer => gauge_set_cached!("rustfs_ec_encode_producer_bytes_current", next as f64),
|
||||
EcEncodePeakMetric::Queue => unreachable!("queue bytes use their own ownership guard"),
|
||||
EcEncodePeakMetric::Writer => gauge_set_cached!("rustfs_ec_encode_writer_bytes_current", next as f64),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn track_ec_encode_payload_stage(
|
||||
bytes: usize,
|
||||
counter: &'static AtomicU64,
|
||||
peak: &'static AtomicU64,
|
||||
metric: EcEncodePeakMetric,
|
||||
) -> EcEncodePayloadStageGuard {
|
||||
let enabled = put_stage_metrics_enabled();
|
||||
let bytes = bytes as u64;
|
||||
if enabled {
|
||||
let next = counter.fetch_add(bytes, Ordering::Relaxed) + bytes;
|
||||
match metric {
|
||||
EcEncodePeakMetric::Producer => gauge_set_cached!("rustfs_ec_encode_producer_bytes_current", next as f64),
|
||||
EcEncodePeakMetric::Queue => unreachable!("queue bytes use their own ownership guard"),
|
||||
EcEncodePeakMetric::Writer => gauge_set_cached!("rustfs_ec_encode_writer_bytes_current", next as f64),
|
||||
}
|
||||
publish_ec_encode_peak(peak, metric, next);
|
||||
}
|
||||
EcEncodePayloadStageGuard {
|
||||
counter,
|
||||
bytes,
|
||||
enabled,
|
||||
current_metric: metric,
|
||||
}
|
||||
}
|
||||
|
||||
/// Track encoded producer payload bytes until queue hand-off completes.
|
||||
#[inline(always)]
|
||||
pub fn track_ec_encode_producer_bytes(bytes: usize) -> EcEncodePayloadStageGuard {
|
||||
track_ec_encode_payload_stage(
|
||||
bytes,
|
||||
&EC_ENCODE_PRODUCER_BYTES_CURRENT,
|
||||
&EC_ENCODE_PRODUCER_BYTES_PEAK,
|
||||
EcEncodePeakMetric::Producer,
|
||||
)
|
||||
}
|
||||
|
||||
/// Track encoded payload bytes while shard writers own the batch.
|
||||
#[inline(always)]
|
||||
pub fn track_ec_encode_writer_bytes(bytes: usize) -> EcEncodePayloadStageGuard {
|
||||
track_ec_encode_payload_stage(
|
||||
bytes,
|
||||
&EC_ENCODE_WRITER_BYTES_CURRENT,
|
||||
&EC_ENCODE_WRITER_BYTES_PEAK,
|
||||
EcEncodePeakMetric::Writer,
|
||||
)
|
||||
}
|
||||
|
||||
/// Return the process-lifetime high-water mark of encoded producer payload bytes.
|
||||
#[inline(always)]
|
||||
pub fn current_ec_encode_producer_bytes_peak() -> u64 {
|
||||
EC_ENCODE_PRODUCER_BYTES_PEAK.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Return the current process-wide encoded producer payload bytes.
|
||||
#[inline(always)]
|
||||
pub fn current_ec_encode_producer_bytes() -> u64 {
|
||||
EC_ENCODE_PRODUCER_BYTES_CURRENT.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Return the process-lifetime high-water mark of encoded queue payload bytes.
|
||||
#[inline(always)]
|
||||
pub fn current_ec_encode_queue_bytes_peak() -> u64 {
|
||||
EC_ENCODE_QUEUE_BYTES_PEAK.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Return the process-lifetime high-water mark of encoded writer payload bytes.
|
||||
#[inline(always)]
|
||||
pub fn current_ec_encode_writer_bytes_peak() -> u64 {
|
||||
EC_ENCODE_WRITER_BYTES_PEAK.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Return the current process-wide encoded writer payload bytes.
|
||||
#[inline(always)]
|
||||
pub fn current_ec_encode_writer_bytes() -> u64 {
|
||||
EC_ENCODE_WRITER_BYTES_CURRENT.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Track whole-object buffering on the GET path.
|
||||
#[inline(always)]
|
||||
pub fn track_get_object_buffered_bytes(bytes: usize) -> Option<MemoryGaugeGuard> {
|
||||
@@ -2385,7 +2543,9 @@ pub fn record_io_latency_p99(latency_ms: f64) {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use std::sync::Mutex;
|
||||
use metrics_util::MetricKind;
|
||||
use metrics_util::debugging::{DebugValue, DebuggingRecorder};
|
||||
use std::sync::{Arc, Barrier, Mutex};
|
||||
|
||||
// Serialize tests that mutate the process-global PUT_STAGE_METRICS_ENABLED flag.
|
||||
static METRICS_FLAG_LOCK: Mutex<()> = Mutex::new(());
|
||||
@@ -2838,6 +2998,9 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn test_ec_encode_inflight_bytes_tracking() {
|
||||
let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner());
|
||||
set_put_stage_metrics_enabled(true);
|
||||
let queue_peak_before = current_ec_encode_queue_bytes_peak();
|
||||
EC_ENCODE_INFLIGHT_BYTES.store(0, Ordering::Relaxed);
|
||||
add_ec_encode_inflight_bytes(1024);
|
||||
add_ec_encode_inflight_bytes(2048);
|
||||
@@ -2845,6 +3008,104 @@ mod tests {
|
||||
remove_ec_encode_inflight_bytes(2048);
|
||||
remove_ec_encode_inflight_bytes(4096);
|
||||
assert_eq!(current_ec_encode_inflight_bytes(), 0);
|
||||
assert!(
|
||||
current_ec_encode_queue_bytes_peak() >= queue_peak_before.max(3072),
|
||||
"queue peak must retain the largest observed queue occupancy"
|
||||
);
|
||||
set_put_stage_metrics_enabled(false);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_ec_encode_producer_and_writer_stage_guards_aggregate_and_settle() {
|
||||
let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner());
|
||||
set_put_stage_metrics_enabled(true);
|
||||
|
||||
let producer_bytes = 1024;
|
||||
let producer_first = track_ec_encode_producer_bytes(producer_bytes);
|
||||
let producer_second = track_ec_encode_producer_bytes(producer_bytes);
|
||||
assert_eq!(current_ec_encode_producer_bytes(), 2048);
|
||||
assert!(
|
||||
current_ec_encode_producer_bytes_peak() >= 2048,
|
||||
"producer peak must include simultaneous stage ownership"
|
||||
);
|
||||
drop((producer_first, producer_second));
|
||||
assert_eq!(current_ec_encode_producer_bytes(), 0);
|
||||
|
||||
let writer_bytes = 2048;
|
||||
let writer_first = track_ec_encode_writer_bytes(writer_bytes);
|
||||
let writer_second = track_ec_encode_writer_bytes(writer_bytes);
|
||||
assert_eq!(current_ec_encode_writer_bytes(), 4096);
|
||||
assert!(
|
||||
current_ec_encode_writer_bytes_peak() >= 4096,
|
||||
"writer peak must include simultaneous stage ownership"
|
||||
);
|
||||
drop((writer_first, writer_second));
|
||||
assert_eq!(current_ec_encode_writer_bytes(), 0);
|
||||
|
||||
set_put_stage_metrics_enabled(false);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_ec_encode_producer_peak_exports_the_high_water_mark() {
|
||||
let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner());
|
||||
assert_eq!(current_ec_encode_producer_bytes(), 0, "test must start without producer stage ownership");
|
||||
let previous_peak = EC_ENCODE_PRODUCER_BYTES_PEAK.swap(0, Ordering::Relaxed);
|
||||
let recorder = DebuggingRecorder::new();
|
||||
let snapshotter = recorder.snapshotter();
|
||||
|
||||
metrics::with_local_recorder(&recorder, || {
|
||||
set_put_stage_metrics_enabled(true);
|
||||
let first = track_ec_encode_producer_bytes(1024);
|
||||
let second = track_ec_encode_producer_bytes(2048);
|
||||
drop((first, second));
|
||||
set_put_stage_metrics_enabled(false);
|
||||
});
|
||||
|
||||
let exported_peak = snapshotter
|
||||
.snapshot()
|
||||
.into_vec()
|
||||
.into_iter()
|
||||
.find_map(|(composite, _, _, value)| {
|
||||
(composite.kind() == MetricKind::Gauge && composite.key().name() == "rustfs_ec_encode_producer_bytes_peak")
|
||||
.then_some(value)
|
||||
});
|
||||
assert!(
|
||||
matches!(exported_peak, Some(DebugValue::Gauge(value)) if value.0 == 3072.0),
|
||||
"exported producer peak must retain the aggregate high-water mark"
|
||||
);
|
||||
EC_ENCODE_PRODUCER_BYTES_PEAK.fetch_max(previous_peak, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_ec_encode_peak_publish_does_not_regress_after_out_of_order_cas() {
|
||||
let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner());
|
||||
let peak = Arc::new(AtomicU64::new(0));
|
||||
let exported = Arc::new(AtomicU64::new(0));
|
||||
let first_ready = Arc::new(Barrier::new(2));
|
||||
let release_first = Arc::new(Barrier::new(2));
|
||||
let first_peak = Arc::clone(&peak);
|
||||
let first_exported = Arc::clone(&exported);
|
||||
let first_ready_for_thread = Arc::clone(&first_ready);
|
||||
let release_first_for_thread = Arc::clone(&release_first);
|
||||
|
||||
let first = std::thread::spawn(move || {
|
||||
publish_ec_encode_peak_with(
|
||||
&first_peak,
|
||||
10,
|
||||
|| {
|
||||
first_ready_for_thread.wait();
|
||||
release_first_for_thread.wait();
|
||||
},
|
||||
|value| first_exported.store(value as u64, Ordering::Relaxed),
|
||||
);
|
||||
});
|
||||
first_ready.wait();
|
||||
publish_ec_encode_peak_with(&peak, 20, || {}, |value| exported.store(value as u64, Ordering::Relaxed));
|
||||
release_first.wait();
|
||||
first.join().expect("first publisher should not panic");
|
||||
|
||||
assert_eq!(peak.load(Ordering::Relaxed), 20);
|
||||
assert_eq!(exported.load(Ordering::Relaxed), 20, "late publisher must reload the high-water mark");
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
Reference in New Issue
Block a user