diff --git a/crates/ecstore/src/erasure_coding/encode.rs b/crates/ecstore/src/erasure_coding/encode.rs index 2310702fc..3b459a7d6 100644 --- a/crates/ecstore/src/erasure_coding/encode.rs +++ b/crates/ecstore/src/erasure_coding/encode.rs @@ -23,6 +23,7 @@ use bytes::Bytes; use futures::StreamExt; use futures::stream::FuturesUnordered; use std::sync::Arc; +use std::time::Instant; use std::vec; use tokio::io::AsyncRead; use tokio::runtime::RuntimeFlavor; @@ -40,6 +41,18 @@ const DEFAULT_RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS: usize = 4; static CACHED_MAX_INFLIGHT_BYTES: std::sync::OnceLock = std::sync::OnceLock::new(); static CACHED_BATCH_BLOCKS: std::sync::OnceLock = std::sync::OnceLock::new(); +#[inline(always)] +fn stage_timer_if_enabled() -> Option { + rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now) +} + +#[inline(always)] +fn record_internal_stage_if_enabled(stage: &'static str, started_at: Option) { + if let Some(started_at) = started_at { + rustfs_io_metrics::record_stage_duration(stage, started_at.elapsed().as_secs_f64() * 1000.0); + } +} + fn encode_channel_capacity(expanded_block_bytes: usize, max_inflight_bytes: usize) -> usize { if expanded_block_bytes == 0 { return 1; @@ -57,6 +70,15 @@ fn encode_batch_block_count() -> usize { }) } +fn erasure_encode_max_inflight_bytes() -> usize { + *CACHED_MAX_INFLIGHT_BYTES.get_or_init(|| { + rustfs_utils::get_env_usize( + ENV_RUSTFS_ERASURE_ENCODE_MAX_INFLIGHT_BYTES, + DEFAULT_RUSTFS_ERASURE_ENCODE_MAX_INFLIGHT_BYTES, + ) + }) +} + fn queued_block_bytes(block: &[Bytes]) -> usize { block.iter().map(Bytes::len).sum() } @@ -237,6 +259,7 @@ impl<'a> MultiWriter<'a> { impl Erasure { async fn encode_block(self: Arc, encode_buf: Vec, len: usize) -> std::io::Result<(Vec, Vec)> { + let encode_stage_start = stage_timer_if_enabled(); let encode_once = move || { let res = self.encode_data(&encode_buf[..len]); (res, encode_buf) @@ -252,6 +275,7 @@ impl Erasure { .map_err(|err| std::io::Error::other(format!("EC encode task failed: {err}")))?, }; + record_internal_stage_if_enabled("erasure_encode_cpu", encode_stage_start); Ok((res?, returned_buf)) } @@ -316,12 +340,7 @@ impl Erasure { // Bound queued encoded blocks by memory budget to avoid per-request spikes. let expanded_block_bytes = self.shard_size().saturating_mul(self.total_shard_count()); - let max_inflight_bytes = *CACHED_MAX_INFLIGHT_BYTES.get_or_init(|| { - rustfs_utils::get_env_usize( - ENV_RUSTFS_ERASURE_ENCODE_MAX_INFLIGHT_BYTES, - DEFAULT_RUSTFS_ERASURE_ENCODE_MAX_INFLIGHT_BYTES, - ) - }); + let max_inflight_bytes = erasure_encode_max_inflight_bytes(); let inflight_blocks = encode_channel_capacity(expanded_block_bytes, max_inflight_bytes); let (tx, mut rx) = mpsc::channel::>(inflight_blocks); @@ -339,10 +358,12 @@ impl Erasure { buf = returned_buf; let queued_bytes = queued_block_bytes(&res); rustfs_io_metrics::add_ec_encode_inflight_bytes(queued_bytes); + let send_wait_stage_start = stage_timer_if_enabled(); if let Err(err) = tx.send(res).await { rustfs_io_metrics::remove_ec_encode_inflight_bytes(queued_bytes); return Err(std::io::Error::other(format!("Failed to send encoded data : {err}"))); } + record_internal_stage_if_enabled("erasure_encode_send_wait", send_wait_stage_start); } Ok(None) => { break; @@ -369,30 +390,41 @@ impl Erasure { let mut write_err = None; - while let Some(block) = rx.recv().await { + loop { + let recv_wait_stage_start = stage_timer_if_enabled(); + let Some(block) = rx.recv().await else { + break; + }; + record_internal_stage_if_enabled("erasure_encode_recv_wait", recv_wait_stage_start); if block.is_empty() { break; } let queued_bytes = queued_block_bytes(&block); rustfs_io_metrics::remove_ec_encode_inflight_bytes(queued_bytes); + let write_stage_start = stage_timer_if_enabled(); if let Err(err) = writers.write(block).await { write_err = Some(err); break; } + record_internal_stage_if_enabled("erasure_encode_write", write_stage_start); } if let Some(err) = write_err { task.abort(); let _ = task.await; drain_queued_inflight_bytes(&mut rx).await; + let shutdown_stage_start = stage_timer_if_enabled(); if let Err(shutdown_err) = writers.shutdown().await { error!("failed to shutdown erasure writers after write error: {:?}", shutdown_err); } + record_internal_stage_if_enabled("erasure_encode_shutdown", shutdown_stage_start); return Err(err); } let (reader, total) = task.await??; + let shutdown_stage_start = stage_timer_if_enabled(); writers.shutdown().await?; + record_internal_stage_if_enabled("erasure_encode_shutdown", shutdown_stage_start); Ok((reader, total)) } @@ -413,12 +445,7 @@ impl Erasure { } let expanded_block_bytes = self.shard_size().saturating_mul(self.total_shard_count()); - let max_inflight_bytes = *CACHED_MAX_INFLIGHT_BYTES.get_or_init(|| { - rustfs_utils::get_env_usize( - ENV_RUSTFS_ERASURE_ENCODE_MAX_INFLIGHT_BYTES, - DEFAULT_RUSTFS_ERASURE_ENCODE_MAX_INFLIGHT_BYTES, - ) - }); + let max_inflight_bytes = erasure_encode_max_inflight_bytes(); let inflight_blocks = encode_channel_capacity(expanded_block_bytes, max_inflight_bytes); let batch_blocks = encode_batch_block_count().min(inflight_blocks); let channel_capacity = inflight_blocks.div_ceil(batch_blocks).max(1); @@ -444,10 +471,12 @@ impl Erasure { if pending_batch.len() >= batch_blocks { rustfs_io_metrics::add_ec_encode_inflight_bytes(pending_batch_bytes); + let send_wait_stage_start = stage_timer_if_enabled(); if let Err(err) = tx.send(pending_batch).await { rustfs_io_metrics::remove_ec_encode_inflight_bytes(pending_batch_bytes); 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); pending_batch = Vec::with_capacity(batch_blocks); pending_batch_bytes = 0; } @@ -471,10 +500,12 @@ impl Erasure { if !pending_batch.is_empty() { rustfs_io_metrics::add_ec_encode_inflight_bytes(pending_batch_bytes); + let send_wait_stage_start = stage_timer_if_enabled(); if let Err(err) = tx.send(pending_batch).await { rustfs_io_metrics::remove_ec_encode_inflight_bytes(pending_batch_bytes); 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); } Ok((reader, total)) @@ -483,14 +514,21 @@ impl Erasure { let mut writers = MultiWriter::new(writers, quorum); let mut write_err = None; - while let Some(batch) = rx.recv().await { + loop { + let recv_wait_stage_start = stage_timer_if_enabled(); + let Some(batch) = rx.recv().await else { + break; + }; + record_internal_stage_if_enabled("erasure_encode_batched_recv_wait", recv_wait_stage_start); rustfs_io_metrics::remove_ec_encode_inflight_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 { write_err = Some(err); break; } } + record_internal_stage_if_enabled("erasure_encode_batched_write", write_stage_start); if write_err.is_some() { break; } @@ -500,14 +538,18 @@ impl Erasure { task.abort(); let _ = task.await; drain_queued_batched_inflight_bytes(&mut rx).await; + let shutdown_stage_start = stage_timer_if_enabled(); if let Err(shutdown_err) = writers.shutdown().await { error!("failed to shutdown erasure writers after write error: {:?}", shutdown_err); } + record_internal_stage_if_enabled("erasure_encode_batched_shutdown", shutdown_stage_start); return Err(err); } let (reader, total) = task.await??; + let shutdown_stage_start = stage_timer_if_enabled(); writers.shutdown().await?; + record_internal_stage_if_enabled("erasure_encode_batched_shutdown", shutdown_stage_start); Ok((reader, total)) } diff --git a/crates/ecstore/src/rpc/remote_disk.rs b/crates/ecstore/src/rpc/remote_disk.rs index 6a9a38020..cf4f14a0e 100644 --- a/crates/ecstore/src/rpc/remote_disk.rs +++ b/crates/ecstore/src/rpc/remote_disk.rs @@ -751,6 +751,12 @@ fn encode_msgpack(value: &T) -> Result> { Ok(serializer.into_inner()) } +fn encode_msgpack_named(value: &T) -> Result> { + let mut serializer = rmp_serde::Serializer::new(Vec::new()).with_struct_map(); + value.serialize(&mut serializer)?; + Ok(serializer.into_inner()) +} + fn decode_msgpack_or_json(binary: &[u8], json: &str) -> Result { if !binary.is_empty() { let mut deserializer = rmp_serde::Deserializer::new(Cursor::new(binary)); @@ -1524,6 +1530,7 @@ impl DiskAPI for RemoteDisk { "rename_data", || async { let file_info = serde_json::to_string(&fi)?; + let file_info_bin = encode_msgpack_named(&fi)?; let mut client = self .get_client() .await @@ -1535,6 +1542,7 @@ impl DiskAPI for RemoteDisk { file_info, dst_volume: dst_volume.to_string(), dst_path: dst_path.to_string(), + file_info_bin: file_info_bin.into(), }); let response = client.rename_data(request).await?.into_inner(); @@ -1543,7 +1551,8 @@ impl DiskAPI for RemoteDisk { return Err(response.error.unwrap_or_default().into()); } - let rename_data_resp = serde_json::from_str::(&response.rename_data_resp)?; + let rename_data_resp = + decode_msgpack_or_json::(&response.rename_data_resp_bin, &response.rename_data_resp)?; Ok(rename_data_resp) }, @@ -2373,6 +2382,64 @@ mod tests { } } + fn sample_rename_data_file_info() -> FileInfo { + FileInfo { + volume: "bucket".to_string(), + name: "object".to_string(), + version_id: Some(Uuid::new_v4()), + data_dir: Some(Uuid::new_v4()), + size: 64 * 1024, + mod_time: Some(::time::OffsetDateTime::UNIX_EPOCH + ::time::Duration::seconds(1)), + metadata: [ + ("etag".to_string(), "etag-value".to_string()), + ("content-type".to_string(), "application/octet-stream".to_string()), + ] + .into_iter() + .collect(), + erasure: rustfs_filemeta::ErasureInfo { + algorithm: rustfs_filemeta::ERASURE_ALGORITHM.to_string(), + data_blocks: 4, + parity_blocks: 2, + block_size: 1024 * 1024, + index: 1, + distribution: vec![1, 2, 3, 4, 5, 6], + ..Default::default() + }, + ..Default::default() + } + } + + #[test] + fn rename_data_file_info_named_msgpack_is_smaller_than_json() { + let file_info = sample_rename_data_file_info(); + let json = serde_json::to_vec(&file_info).expect("file info json should encode"); + let named_msgpack = encode_msgpack_named(&file_info).expect("file info named msgpack should encode"); + + assert!( + named_msgpack.len() < json.len(), + "expected named msgpack payload to be smaller than json (msgpack={}, json={})", + named_msgpack.len(), + json.len() + ); + } + + #[test] + fn rename_data_resp_named_msgpack_is_smaller_than_json() { + let response = RenameDataResp { + old_data_dir: Some(Uuid::new_v4()), + sign: Some(vec![1_u8; 32]), + }; + let json = serde_json::to_vec(&response).expect("rename data response json should encode"); + let named_msgpack = encode_msgpack_named(&response).expect("rename data response named msgpack should encode"); + + assert!( + named_msgpack.len() < json.len(), + "expected named msgpack payload to be smaller than json (msgpack={}, json={})", + named_msgpack.len(), + json.len() + ); + } + #[derive(Debug, Default)] struct SinkTestWriter; diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index 4f66f6e8f..1756ab3f9 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -152,6 +152,9 @@ const SET_DISK_COMMIT_TAIL_WARN_THRESHOLD_MS: u128 = 5_000; const ENV_RUSTFS_PUT_LARGE_BATCH_MIN_SIZE_BYTES: &str = "RUSTFS_PUT_LARGE_BATCH_MIN_SIZE_BYTES"; const DEFAULT_RUSTFS_PUT_LARGE_BATCH_MIN_SIZE_BYTES: usize = 64 * 1024 * 1024; static CACHED_PUT_LARGE_BATCH_MIN_SIZE_BYTES: std::sync::OnceLock = std::sync::OnceLock::new(); +const ENV_RUSTFS_MULTIPART_PUT_LARGE_BATCH_MIN_SIZE_BYTES: &str = "RUSTFS_MULTIPART_PUT_LARGE_BATCH_MIN_SIZE_BYTES"; +const DEFAULT_RUSTFS_MULTIPART_PUT_LARGE_BATCH_MIN_SIZE_BYTES: usize = 128 * 1024 * 1024; +static CACHED_MULTIPART_PUT_LARGE_BATCH_MIN_SIZE_BYTES: std::sync::OnceLock = std::sync::OnceLock::new(); use crate::rio::{EtagResolvable, HashReader, HashReaderMut, TryGetIndex as _}; @@ -827,6 +830,15 @@ impl SmallWritePath { SmallWritePath::PipelineBatchedLarge => "write_pipeline_batched_large", } } + + fn multipart_metric_label(&self) -> &'static str { + match self { + SmallWritePath::Inline => "multipart_write_inline", + SmallWritePath::SingleBlockNonInline => "multipart_write_single_block_non_inline", + SmallWritePath::Pipeline => "multipart_write_pipeline", + SmallWritePath::PipelineBatchedLarge => "multipart_write_pipeline_batched_large", + } + } } fn put_large_batch_min_size_bytes() -> usize { @@ -835,6 +847,15 @@ fn put_large_batch_min_size_bytes() -> usize { }) } +fn multipart_put_large_batch_min_size_bytes() -> usize { + *CACHED_MULTIPART_PUT_LARGE_BATCH_MIN_SIZE_BYTES.get_or_init(|| { + rustfs_utils::get_env_usize( + ENV_RUSTFS_MULTIPART_PUT_LARGE_BATCH_MIN_SIZE_BYTES, + DEFAULT_RUSTFS_MULTIPART_PUT_LARGE_BATCH_MIN_SIZE_BYTES, + ) + }) +} + fn classify_small_write_path(is_inline_buffer: bool, object_size: i64, block_size: usize) -> SmallWritePath { if should_use_inline_small_fast_path(is_inline_buffer, object_size, block_size) { SmallWritePath::Inline @@ -859,6 +880,17 @@ fn classify_put_write_path(is_inline_buffer: bool, object_size: i64, block_size: } } +fn classify_multipart_part_write_path(object_size: i64, block_size: usize) -> SmallWritePath { + if should_use_single_block_non_inline_fast_path(false, object_size, block_size) { + return SmallWritePath::SingleBlockNonInline; + } + + match usize::try_from(object_size) { + Ok(size) if size >= multipart_put_large_batch_min_size_bytes() => SmallWritePath::PipelineBatchedLarge, + _ => SmallWritePath::Pipeline, + } +} + #[async_trait::async_trait] impl rustfs_storage_api::ObjectIO for SetDisks { type Error = Error; @@ -3392,6 +3424,7 @@ impl rustfs_storage_api::MultipartOperations for SetDisks { let tmp_part_path = Arc::new(format!("{tmp_part}/{part_suffix}")); let erasure = erasure_coding::Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); + let writer_setup_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now); let mut writers = Vec::with_capacity(shuffle_disks.len()); let mut errors = Vec::with_capacity(shuffle_disks.len()); @@ -3433,6 +3466,13 @@ impl rustfs_storage_api::MultipartOperations for SetDisks { } } + if let Some(stage_start) = writer_setup_stage_start { + rustfs_io_metrics::record_put_object_stage_duration( + "multipart_set_disk_writer_setup", + stage_start.elapsed().as_secs_f64() * 1000.0, + ); + } + let nil_count = errors.iter().filter(|&e| e.is_none()).count(); if nil_count < write_quorum { if let Some(write_err) = reduce_write_quorum_errs(&errors, OBJECT_OP_IGNORED_ERRS, write_quorum) { @@ -3442,12 +3482,16 @@ impl rustfs_storage_api::MultipartOperations for SetDisks { return Err(Error::other(format!("not enough disks to write: {errors:?}"))); } + // Capture the original part size before swapping the stream out for encoding. + let multipart_part_size = data.size(); let stream = mem::replace( &mut data.stream, HashReader::from_stream(Cursor::new(Vec::new()), 0, 0, None, None, false)?, ); - let write_path = classify_small_write_path(false, data.size(), fi.erasure.block_size); + let write_path = classify_multipart_part_write_path(multipart_part_size, fi.erasure.block_size); + rustfs_io_metrics::record_put_object_path(write_path.multipart_metric_label()); + let encode_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now); let (reader, w_size) = match write_path { SmallWritePath::SingleBlockNonInline => { @@ -3455,11 +3499,19 @@ impl rustfs_storage_api::MultipartOperations for SetDisks { .encode_single_block_non_inline(stream, &mut writers, write_quorum) .await? } - SmallWritePath::Inline | SmallWritePath::Pipeline | SmallWritePath::PipelineBatchedLarge => { + SmallWritePath::PipelineBatchedLarge => Arc::new(erasure).encode_batched(stream, &mut writers, write_quorum).await?, + SmallWritePath::Inline | SmallWritePath::Pipeline => { Arc::new(erasure).encode(stream, &mut writers, write_quorum).await? } }; // TODO: delete temporary directory on error + if let Some(stage_start) = encode_stage_start { + rustfs_io_metrics::record_put_object_stage_duration( + "multipart_set_disk_encode", + stage_start.elapsed().as_secs_f64() * 1000.0, + ); + } + let _ = mem::replace(&mut data.stream, reader); if (w_size as i64) < data.size() { @@ -4337,6 +4389,7 @@ impl rustfs_storage_api::MultipartOperations for SetDisks { ); } + let complete_tail_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(Instant::now); self.cleanup_multipart_path(&parts).await; let (online_disks, versions, op_old_dir, cleanup_disks) = Self::rename_data( @@ -4355,6 +4408,13 @@ impl rustfs_storage_api::MultipartOperations for SetDisks { .await?; } + if let Some(stage_start) = complete_tail_stage_start { + rustfs_io_metrics::record_put_object_stage_duration( + "multipart_complete_tail", + stage_start.elapsed().as_secs_f64() * 1000.0, + ); + } + drop(object_lock_guard); // drop object lock guard to release the lock if let Some(versions) = versions { @@ -7492,6 +7552,35 @@ mod tests { )); } + #[test] + fn multipart_put_large_batch_path_only_applies_at_128m_and_above() { + assert!(matches!( + classify_multipart_part_write_path(128 * 1024 * 1024, 1024 * 1024), + SmallWritePath::PipelineBatchedLarge + )); + assert!(matches!( + classify_multipart_part_write_path(64 * 1024 * 1024, 1024 * 1024), + SmallWritePath::Pipeline + )); + assert!(matches!( + classify_multipart_part_write_path(1024 * 1024, 1024 * 1024), + SmallWritePath::SingleBlockNonInline + )); + } + + #[test] + fn multipart_write_paths_use_distinct_metric_labels() { + assert_eq!(SmallWritePath::Pipeline.multipart_metric_label(), "multipart_write_pipeline"); + assert_eq!( + SmallWritePath::PipelineBatchedLarge.multipart_metric_label(), + "multipart_write_pipeline_batched_large" + ); + assert_eq!( + SmallWritePath::SingleBlockNonInline.multipart_metric_label(), + "multipart_write_single_block_non_inline" + ); + } + #[test] fn test_is_cold_storage_class() { // Test cold storage classes diff --git a/crates/ecstore/tests/protobuf_bytes_regression_test.rs b/crates/ecstore/tests/protobuf_bytes_regression_test.rs index a09dedd9b..d7cf76e32 100644 --- a/crates/ecstore/tests/protobuf_bytes_regression_test.rs +++ b/crates/ecstore/tests/protobuf_bytes_regression_test.rs @@ -3,7 +3,8 @@ use bytes::Bytes; use rustfs_protos::proto_gen::node_service::{ - ReadMultipleRequest, ReadMultipleResponse, ReadVersionResponse, ReadXlResponse, UpdateMetadataRequest, WriteMetadataRequest, + ReadMultipleRequest, ReadMultipleResponse, ReadVersionResponse, ReadXlResponse, RenameDataRequest, RenameDataResponse, + UpdateMetadataRequest, WriteMetadataRequest, }; fn expect_bytes(_: &Bytes) {} @@ -17,6 +18,12 @@ fn protobuf_bytes_fields_use_bytes_consistently() { let write = WriteMetadataRequest::default(); expect_bytes(&write.file_info_bin); + let rename_data = RenameDataRequest::default(); + expect_bytes(&rename_data.file_info_bin); + + let rename_data_response = RenameDataResponse::default(); + expect_bytes(&rename_data_response.rename_data_resp_bin); + let version = ReadVersionResponse::default(); expect_bytes(&version.file_info_bin); diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index e3b123893..d5e0021ae 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -349,6 +349,8 @@ pub struct RenameDataRequest { pub dst_volume: ::prost::alloc::string::String, #[prost(string, tag = "6")] pub dst_path: ::prost::alloc::string::String, + #[prost(bytes = "bytes", tag = "7")] + pub file_info_bin: ::prost::bytes::Bytes, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct RenameDataResponse { @@ -358,6 +360,8 @@ pub struct RenameDataResponse { pub rename_data_resp: ::prost::alloc::string::String, #[prost(message, optional, tag = "3")] pub error: ::core::option::Option, + #[prost(bytes = "bytes", tag = "4")] + pub rename_data_resp_bin: ::prost::bytes::Bytes, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct MakeVolumesRequest { diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index 42bdbcacc..83a46d048 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -253,12 +253,14 @@ message RenameDataRequest { string file_info = 4; string dst_volume = 5; string dst_path = 6; + bytes file_info_bin = 7; } message RenameDataResponse { bool success = 1; string rename_data_resp = 2; optional Error error = 3; + bytes rename_data_resp_bin = 4; } message MakeVolumesRequest { diff --git a/docs/observability/issue-712-local-queryable-metrics-backend-en.md b/docs/observability/issue-712-local-queryable-metrics-backend-en.md new file mode 100644 index 000000000..537c09460 --- /dev/null +++ b/docs/observability/issue-712-local-queryable-metrics-backend-en.md @@ -0,0 +1,232 @@ +# Issue #712 Local Queryable Metrics Backend Guide + +## 1. Purpose + +This guide is intended to unblock the third validation batch for `#712` by making multipart PUT stage metrics queryable from a local backend. + +The current problem is not that multipart stage metrics are missing from the code path. The actual problem is that the local environment does not expose a queryable metrics backend: + +1. `rustfs/admin/v3/metrics` is available, but it returns an admin-side JSON snapshot rather than the `metrics` crate histogram series +2. there is no local Prometheus or equivalent queryable metrics endpoint listening by default +3. therefore the following stage labels cannot be queried directly yet: + - `multipart_ingress_prepare` + - `multipart_set_disk_writer_setup` + - `multipart_set_disk_encode` + - `multipart_complete_tail` + +The goal of this guide is to close that gap. + +## 2. Recommended approach + +Reuse the repository's existing observability stack: + +1. `.docker/observability/docker-compose.yml` + +Why this is the preferred path: + +1. it is already maintained in-repo +2. it includes OTEL Collector, Prometheus, and Grafana +3. it can receive telemetry from RustFS through `RUSTFS_OBS_ENDPOINT` + +## 3. Expected data flow + +After startup, the intended flow is: + +1. RustFS + - `RUSTFS_OBS_ENDPOINT=http://host.docker.internal:4318` +2. OTEL Collector + - receives OTLP/HTTP telemetry +3. Prometheus + - scrapes collector-exported metrics +4. Query surface + - `http://127.0.0.1:9090` + +## 4. Startup steps + +### 4.1 Start the observability stack + +From the repository root: + +```bash +cd .docker/observability +docker compose up -d +``` + +### 4.2 Wait for core services + +Recommended checks: + +```bash +curl -fsS http://127.0.0.1:9090/-/ready +curl -fsS http://127.0.0.1:3000/api/health +``` + +If you need container status: + +```bash +docker compose ps +``` + +### 4.3 Point RustFS to the OTEL Collector + +For local single-node multi-disk validation: + +```bash +export RUSTFS_OBS_ENDPOINT=http://host.docker.internal:4318 +``` + +If you use the repository-local restart helper: + +```bash +bash scripts/restart_local_single_node_multidisk_rustfs.sh +``` + +Make sure the final runtime environment really contains: + +```bash +RUSTFS_OBS_ENDPOINT=http://host.docker.internal:4318 +``` + +## 5. Minimal query validation + +### 5.1 Confirm Prometheus can see RustFS metrics + +```bash +curl -fsS 'http://127.0.0.1:9090/api/v1/query?query=rustfs_s3_put_object_total' +``` + +### 5.2 Confirm stage labels exist + +```bash +curl -fsS 'http://127.0.0.1:9090/api/v1/label/stage/values' +``` + +If the pipeline is working, the result should include: + +1. `multipart_ingress_prepare` +2. `multipart_set_disk_writer_setup` +3. `multipart_set_disk_encode` +4. `multipart_complete_tail` + +### 5.3 Direct P95 query + +```bash +curl -fsS 'http://127.0.0.1:9090/api/v1/query?query=histogram_quantile(0.95,sum by(stage,le)(rate(rustfs_s3_put_object_stage_duration_ms_bucket{stage=~"multipart_.*"}[5m])))' +``` + +## 6. Recommended silent validation order + +Once the backend is queryable, use this order for the third multipart validation batch: + +1. start the observability stack +2. restart RustFS and confirm `RUSTFS_OBS_ENDPOINT` is active +3. run one multipart baseline: + - `1g-64m-pc4` + - `2g-128m-pc4` +4. ignore streaming benchmark logs and only keep: + - `summary.csv` + - Prometheus query outputs +5. summarize: + - throughput / reqps / average latency + - P95 / P99 for the four multipart stages + +## 7. Recommended query set + +### 7.1 Multipart stage P95 + +```promql +histogram_quantile( + 0.95, + sum by (stage, le) ( + rate(rustfs_s3_put_object_stage_duration_ms_bucket{stage=~"multipart_.*"}[5m]) + ) +) +``` + +### 7.2 Multipart stage P99 + +```promql +histogram_quantile( + 0.99, + sum by (stage, le) ( + rate(rustfs_s3_put_object_stage_duration_ms_bucket{stage=~"multipart_.*"}[5m]) + ) +) +``` + +### 7.3 Complete tail focus + +```promql +histogram_quantile( + 0.95, + sum by (instance, le) ( + rate(rustfs_s3_put_object_stage_duration_ms_bucket{stage="multipart_complete_tail"}[5m]) + ) +) +``` + +### 7.4 Encode focus + +```promql +histogram_quantile( + 0.95, + sum by (instance, le) ( + rate(rustfs_s3_put_object_stage_duration_ms_bucket{stage="multipart_set_disk_encode"}[5m]) + ) +) +``` + +## 8. Suggested result layout + +Recommended layout: + +```text +target/bench/ + issue712-multipart-server-path-focus/ + summary.csv + metrics-query.txt + promql/ + multipart-stage-p95.txt + multipart-stage-p99.txt + multipart-complete-tail-p95.txt + multipart-encode-p95.txt +``` + +## 9. Troubleshooting + +### 9.1 Prometheus does not start + +Check: + +```bash +cd .docker/observability +docker compose logs prometheus +``` + +### 9.2 RustFS does not export telemetry + +Check: + +1. `RUSTFS_OBS_ENDPOINT` really points to `http://host.docker.internal:4318` +2. the OTEL collector container is running +3. RustFS was restarted after the environment variable changed + +### 9.3 Stage labels are missing + +Check: + +1. a multipart PUT workload really ran +2. `put_stage_metrics_enabled()` was enabled at runtime +3. the query window is not too short + +## 10. Recommendation + +When `#712` third-batch validation resumes, do not run the benchmark first and hunt for metrics later. + +Use this order instead: + +1. bring up a queryable backend +2. run the multipart baseline +3. read only: + - `summary.csv` + - Prometheus stage query outputs diff --git a/docs/observability/issue-712-local-queryable-metrics-backend-zh.md b/docs/observability/issue-712-local-queryable-metrics-backend-zh.md new file mode 100644 index 000000000..83c1b9b60 --- /dev/null +++ b/docs/observability/issue-712-local-queryable-metrics-backend-zh.md @@ -0,0 +1,232 @@ +# Issue #712 本地可查询 metrics backend 启动手册 + +## 1. 目的 + +本文用于打通 `#712` 第三批验证所需的本地可查询 metrics backend。 + +当前问题不是 multipart stage 指标没有打点,而是本地环境里没有可查询后端: + +1. `rustfs/admin/v3/metrics` 可用,但返回的是管理侧 JSON 快照,不包含 `metrics` crate 的 histogram 指标 +2. 本地默认没有 Prometheus / 可查询 metrics endpoint 在监听 +3. 因此无法直接查询: + - `multipart_ingress_prepare` + - `multipart_set_disk_writer_setup` + - `multipart_set_disk_encode` + - `multipart_complete_tail` + +本文的目标就是把这条链打通。 + +## 2. 推荐方案 + +推荐直接复用仓库内已有的 observability stack: + +1. `.docker/observability/docker-compose.yml` + +这套 stack 的优点: + +1. 已经是仓库现成维护的方案 +2. 包含 OTEL Collector、Prometheus、Grafana +3. 可以直接承接 RustFS 的 `RUSTFS_OBS_ENDPOINT` + +## 3. 核心链路 + +启动后,链路应该是: + +1. RustFS + - `RUSTFS_OBS_ENDPOINT=http://host.docker.internal:4318` +2. OTEL Collector + - 接收 OTLP/HTTP +3. Prometheus + - 抓取 collector 暴露的 metrics +4. 查询 + - `http://127.0.0.1:9090` + +## 4. 启动步骤 + +### 4.1 启动 observability stack + +在仓库根目录执行: + +```bash +cd .docker/observability +docker compose up -d +``` + +### 4.2 等待组件就绪 + +建议检查: + +```bash +curl -fsS http://127.0.0.1:9090/-/ready +curl -fsS http://127.0.0.1:3000/api/health +``` + +若需要查看容器: + +```bash +docker compose ps +``` + +### 4.3 启动 RustFS 时显式指向 OTEL Collector + +本地单机多盘场景建议: + +```bash +export RUSTFS_OBS_ENDPOINT=http://host.docker.internal:4318 +``` + +如果使用我们现成的本地重启脚本: + +```bash +bash scripts/restart_local_single_node_multidisk_rustfs.sh +``` + +请确保脚本最终生效的环境里包含: + +```bash +RUSTFS_OBS_ENDPOINT=http://host.docker.internal:4318 +``` + +## 5. 最小查询验证 + +### 5.1 确认 Prometheus 能查询到 RustFS 指标 + +```bash +curl -fsS 'http://127.0.0.1:9090/api/v1/query?query=rustfs_s3_put_object_total' +``` + +### 5.2 查询 stage 指标是否存在 + +```bash +curl -fsS 'http://127.0.0.1:9090/api/v1/label/stage/values' +``` + +如果链路打通,返回结果中应能看到: + +1. `multipart_ingress_prepare` +2. `multipart_set_disk_writer_setup` +3. `multipart_set_disk_encode` +4. `multipart_complete_tail` + +### 5.3 直接查询 P95 + +```bash +curl -fsS 'http://127.0.0.1:9090/api/v1/query?query=histogram_quantile(0.95,sum by(stage,le)(rate(rustfs_s3_put_object_stage_duration_ms_bucket{stage=~"multipart_.*"}[5m])))' +``` + +## 6. 建议的静默验证顺序 + +打通 backend 后,建议按如下顺序做第三批: + +1. 启动 observability stack +2. 重启 RustFS 并确认 `RUSTFS_OBS_ENDPOINT` 生效 +3. 先跑一轮 multipart baseline: + - `1g-64m-pc4` + - `2g-128m-pc4` +4. 不看过程日志,只保留: + - `summary.csv` + - Prometheus 查询结果 +5. 最后整理: + - throughput / reqps / avg latency + - 4 个 multipart stage 的 P95 / P99 + +## 7. 推荐查询集合 + +### 7.1 multipart 阶段 P95 + +```promql +histogram_quantile( + 0.95, + sum by (stage, le) ( + rate(rustfs_s3_put_object_stage_duration_ms_bucket{stage=~"multipart_.*"}[5m]) + ) +) +``` + +### 7.2 multipart 阶段 P99 + +```promql +histogram_quantile( + 0.99, + sum by (stage, le) ( + rate(rustfs_s3_put_object_stage_duration_ms_bucket{stage=~"multipart_.*"}[5m]) + ) +) +``` + +### 7.3 complete tail 单独看 + +```promql +histogram_quantile( + 0.95, + sum by (instance, le) ( + rate(rustfs_s3_put_object_stage_duration_ms_bucket{stage="multipart_complete_tail"}[5m]) + ) +) +``` + +### 7.4 encode 单独看 + +```promql +histogram_quantile( + 0.95, + sum by (instance, le) ( + rate(rustfs_s3_put_object_stage_duration_ms_bucket{stage="multipart_set_disk_encode"}[5m]) + ) +) +``` + +## 8. 结果目录建议 + +推荐目录: + +```text +target/bench/ + issue712-multipart-server-path-focus/ + summary.csv + metrics-query.txt + promql/ + multipart-stage-p95.txt + multipart-stage-p99.txt + multipart-complete-tail-p95.txt + multipart-encode-p95.txt +``` + +## 9. 失败排查 + +### 9.1 Prometheus 起不来 + +检查: + +```bash +cd .docker/observability +docker compose logs prometheus +``` + +### 9.2 RustFS 没有上报 + +检查: + +1. `RUSTFS_OBS_ENDPOINT` 是否真的是 `http://host.docker.internal:4318` +2. OTEL collector 是否运行 +3. RustFS 重启后是否带上了新环境变量 + +### 9.3 查不到 stage label + +检查: + +1. 是否真的跑过 multipart PUT 请求 +2. `put_stage_metrics_enabled()` 是否在运行期被开启 +3. 查询窗口是否太短 + +## 10. 建议 + +下次推进 `#712` 第三批时,不要先跑 benchmark,再临时找 metrics。 + +正确顺序应是: + +1. 先按本文把可查询 backend 起好 +2. 再跑 baseline +3. 最后只读: + - `summary.csv` + - Prometheus stage 查询结果 diff --git a/docs/observability/issue-712-multipart-put-stage-metrics-dashboard-guide-zh.md b/docs/observability/issue-712-multipart-put-stage-metrics-dashboard-guide-zh.md new file mode 100644 index 000000000..75bf40b59 --- /dev/null +++ b/docs/observability/issue-712-multipart-put-stage-metrics-dashboard-guide-zh.md @@ -0,0 +1,263 @@ +# Issue #712 multipart PUT 分阶段指标 Dashboard / PromQL 指南 + +## 1. 目的 + +本文给 `#712` 的第一批 server-path 观测增强配套一份可执行的 Dashboard / PromQL 指南。 + +本批次新增的 multipart 阶段指标依然复用现有指标名: + +1. `rustfs_s3_put_object_stage_duration_ms` + +但新增了四个 stage label: + +1. `multipart_ingress_prepare` +2. `multipart_set_disk_writer_setup` +3. `multipart_set_disk_encode` +4. `multipart_complete_tail` + +## 2. 使用前提 + +这些阶段指标严格受全局开关控制: + +1. `rustfs_io_metrics::put_stage_metrics_enabled() == true` + +如果该开关没有开启: + +1. 不会上报这些阶段指标 +2. 也不会额外做阶段计时 + +## 3. 推荐直接复用现有 Grafana Row + +当前 Dashboard 中已经有: + +1. `Large PUT Stage Breakdown` + +这意味着: + +1. 不需要重新设计一套全新 row +2. 只需要在现有 row / stage 变量里选新的 multipart stage label 即可 + +## 4. 推荐 PromQL + +### 4.1 multipart 阶段 P95 + +```promql +histogram_quantile( + 0.95, + sum by (stage, le) ( + rate( + rustfs_s3_put_object_stage_duration_ms_bucket{ + job=~"$job", + stage=~"multipart_.*" + }[$__rate_interval] + ) + ) +) +``` + +### 4.2 multipart 阶段 P99 + +```promql +histogram_quantile( + 0.99, + sum by (stage, le) ( + rate( + rustfs_s3_put_object_stage_duration_ms_bucket{ + job=~"$job", + stage=~"multipart_.*" + }[$__rate_interval] + ) + ) +) +``` + +### 4.3 单实例 multipart 阶段 P95 + +```promql +histogram_quantile( + 0.95, + sum by (instance, stage, le) ( + rate( + rustfs_s3_put_object_stage_duration_ms_bucket{ + job=~"$job", + instance=~"$instance", + stage=~"multipart_.*" + }[$__rate_interval] + ) + ) +) +``` + +### 4.4 multipart 与 ordinary PUT encode 对比 + +```promql +histogram_quantile( + 0.95, + sum by (stage, le) ( + rate( + rustfs_s3_put_object_stage_duration_ms_bucket{ + job=~"$job", + stage=~"set_disk_encode|multipart_set_disk_encode" + }[$__rate_interval] + ) + ) +) +``` + +### 4.5 multipart complete tail 重点盯盘 + +```promql +histogram_quantile( + 0.95, + sum by (instance, le) ( + rate( + rustfs_s3_put_object_stage_duration_ms_bucket{ + job=~"$job", + stage="multipart_complete_tail", + instance=~"$instance" + }[$__rate_interval] + ) + ) +) +``` + +### 4.6 multipart path 命中计数 + +```promql +sum by (path) ( + rustfs_s3_put_object_path_total{ + path=~"multipart_.*" + } +) +``` + +用于回答: + +1. 当前 run 是否真的命中了 `multipart_write_pipeline_batched_large` +2. batched gate 是否只是“代码存在”,还是“运行时实际生效” + +### 4.7 erasure encode 内部阶段均值 + +```promql +sum by (stage) ( + increase( + rustfs_internal_stage_duration_ms_sum{ + stage=~"erasure_encode.*" + }[$__rate_interval] + ) +) +/ +sum by (stage) ( + increase( + rustfs_internal_stage_duration_ms_count{ + stage=~"erasure_encode.*" + }[$__rate_interval] + ) +) +``` + +用于回答: + +1. `multipart_set_disk_encode` 内部到底更偏 CPU encode,还是更偏 writer write +2. producer / consumer 之间是 encoder 在等 writer,还是 writer 在等 encoder + +### 4.8 erasure encode 当前累计 counters + +当窗口查询容易受到 scrape 周期影响时,可以直接看当前累计值: + +```promql +rustfs_internal_stage_duration_ms_count{ + stage=~"erasure_encode.*" +} +``` + +```promql +rustfs_internal_stage_duration_ms_sum{ + stage=~"erasure_encode.*" +} +``` + +这在 focused 单 profile 验证里很有用,尤其适合“重启实例后只跑一轮”的场景。 + +## 5. 推荐看板顺序 + +当你在看 `>1GiB multipart PUT` 时,建议按下面顺序看: + +1. `multipart_ingress_prepare` +2. `multipart_set_disk_writer_setup` +3. `multipart_set_disk_encode` +4. `multipart_complete_tail` +5. `multipart_write_pipeline` vs `multipart_write_pipeline_batched_large` +6. `erasure_encode_*` / `erasure_encode_batched_*` + +解释顺序: + +1. 如果 ingress 先高,先看 part ingress buffer / request-body handling +2. 如果 writer setup 高,先看 bitrot writer / disk availability / shard_file_size path +3. 如果 encode 高,先看 multipart 是否需要独立 encode strategy +4. 如果 complete tail 高,优先看 `complete_multipart_upload()` 的 metadata / checksum / rename tail +5. 如果 batched 预期已打开,但 path 仍然只有 `multipart_write_pipeline`,优先检查 size gate 是否真正被命中 +6. 如果 `erasure_encode_batched_send_wait` 很低、但 `erasure_encode_batched_recv_wait` 明显更高,优先怀疑当前 batch barrier 让 writer 侧在等下一批 encode 完成 + +## 6. 推荐结合看的辅助指标 + +建议和上面四个阶段一起看: + +1. `rustfs_io_put_object_concurrent_requests` +2. `rustfs_ec_encode_inflight_bytes_current` +3. host CPU +4. per-instance disk write throughput +5. readiness / write quorum 异常计数 + +## 7. 典型解释模板 + +### 7.1 ingress 高 + +可能原因: + +1. part body stream buffering 不合适 +2. `part.size` 与 ingress buffer 不匹配 + +### 7.2 writer setup 高 + +可能原因: + +1. bitrot writer 构建成本偏高 +2. online disk / writer init 慢 +3. shard_file_size 相关路径有额外成本 + +### 7.3 encode 高 + +可能原因: + +1. multipart part 仍然借用了 ordinary PUT encode 行为 +2. `part.size` 太大,单 part encode CPU 时间过长 +3. batching / inflight 参数不合适 + +### 7.4 complete tail 高 + +可能原因: + +1. complete 阶段 part metadata 处理放大 +2. checksum combine 成本高 +3. rename / cleanup / commit tail 成本高 + +## 8. 建议的截图 / 归档内容 + +每次 `>1GiB multipart PUT` 复测,建议固定归档: + +1. multipart stage P95 截图 +2. multipart stage P99 截图 +3. `multipart_complete_tail` 单实例截图 +4. CPU / disk write 辅助图 + +## 9. 当前阶段建议 + +下一次进入 `#712` 继续推进时: + +1. 先开 `put_stage_metrics_enabled` +2. 先跑推荐 baseline: + - `1GiB -> 64MiB / pc4` + - `2GiB -> 128MiB / pc4` +3. 先看 `multipart_complete_tail` 是否明显高于其他阶段 +4. 再决定是先改 ingress / encode / writer setup / complete tail diff --git a/docs/operations/issue-712-batchblocks-candidate-low-noise-matrix-zh.md b/docs/operations/issue-712-batchblocks-candidate-low-noise-matrix-zh.md new file mode 100644 index 000000000..bed7a926d --- /dev/null +++ b/docs/operations/issue-712-batchblocks-candidate-low-noise-matrix-zh.md @@ -0,0 +1,201 @@ +# Issue #712 `RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS` 候选值低噪声复测矩阵 + +## 1. 目标 + +本文用于下一轮更小面、更低噪声的复测。 + +本轮只回答一个问题: + +1. `RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=2` 是否比 `4` 更稳地改善 `2g-128m-pc4` 的 multipart batched 路径 + +当前不回答: + +1. 是否直接改代码默认值 +2. 是否扩展到 ordinary PUT +3. 是否继续引入新的 encode 调度语义 + +## 2. 固定范围 + +只保留一个 profile: + +1. `2g-128m-pc4` + +只保留两个候选值: + +1. `RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=4` +2. `RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=2` + +不再混入: + +1. `1g-64m-pc4` +2. `2g-256m-pc4` +3. 其他 batching / inflight 配置 + +## 3. 低噪声原则 + +本轮强制执行以下原则: + +1. 每一轮都重启 RustFS 进程,避免 counters/sums 混轮 +2. 每一轮之间固定冷却 `30s` +3. 每一轮都固定只跑 `2m` +4. 每一轮都只抓: + - `summary.csv` + - `path.json` + - `internal_count.json` + - `internal_sum.json` +5. 不看过程日志,不在中间临时扩项 + +## 4. 推荐矩阵 + +推荐使用交叉顺序,避免单边连续运行: + +1. `b4-r1` +2. `b2-r1` +3. `b2-r2` +4. `b4-r2` +5. `b4-r3` +6. `b2-r3` + +这是一个偏保守的 `ABBAAB` 顺序。 + +原因: + +1. 不让 `2` 总是出现在后面 +2. 不让 `4` 总是出现在后面 +3. 能更快看出“只是后跑更快”还是“候选值真的更优” + +## 5. 目录约定 + +建议统一使用: + +```text +target/bench/issue712-batchblocks-candidate-runs/ + b4-r1/ + b2-r1/ + b2-r2/ + b4-r2/ + b4-r3/ + b2-r3/ +``` + +每轮目录固定包含: + +1. `summary.csv` +2. `path.json` +3. `internal_count.json` +4. `internal_sum.json` + +## 6. 单轮执行模板 + +### 6.1 `batch_blocks=4` + +启动: + +```bash +env \ + RUSTFS_UNSAFE_BYPASS_DISK_CHECK=true \ + RUSTFS_ADDRESS=127.0.0.1:9000 \ + RUSTFS_ACCESS_KEY=rustfsadmin \ + RUSTFS_SECRET_KEY=rustfsadmin \ + RUSTFS_RPC_SECRET=rustfs-rpc-secret \ + RUSTFS_REGION=us-east-1 \ + RUSTFS_CONSOLE_ENABLE=false \ + RUSTFS_OBS_ENDPOINT=http://127.0.0.1:4318 \ + RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=4 \ + target/debug/rustfs server \ + /private/tmp/issue708-single-node-multidisk/d1 \ + /private/tmp/issue708-single-node-multidisk/d2 \ + /private/tmp/issue708-single-node-multidisk/d3 \ + /private/tmp/issue708-single-node-multidisk/d4 +``` + +运行: + +```bash +bash scripts/run_gt1g_multipart_put_server_path_focus.sh \ + --host 127.0.0.1:9000 \ + --access-key rustfsadmin \ + --secret-key rustfsadmin \ + --profiles 2g-128m-pc4 \ + --duration 2m \ + --out-dir target/bench/issue712-batchblocks-candidate-runs/b4-r1 +``` + +采集: + +```bash +curl -fsS 'http://127.0.0.1:9090/api/v1/query?query=rustfs_s3_put_object_path_total{path=~"multipart_.*"}' \ + > target/bench/issue712-batchblocks-candidate-runs/b4-r1/path.json + +curl -fsS 'http://127.0.0.1:9090/api/v1/query?query=rustfs_internal_stage_duration_ms_count{stage=~"erasure_encode.*"}' \ + > target/bench/issue712-batchblocks-candidate-runs/b4-r1/internal_count.json + +curl -fsS 'http://127.0.0.1:9090/api/v1/query?query=rustfs_internal_stage_duration_ms_sum{stage=~"erasure_encode.*"}' \ + > target/bench/issue712-batchblocks-candidate-runs/b4-r1/internal_sum.json +``` + +### 6.2 `batch_blocks=2` + +只改一个变量: + +```bash +RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=2 +``` + +其余命令完全保持一致。 + +## 7. 每轮必须核对的点 + +### 7.1 path 命中 + +必须确认: + +1. `multipart_write_pipeline_batched_large` 已命中 + +如果没有命中: + +1. 该轮结果无效 +2. 不要纳入比较 + +### 7.2 internal stages + +至少比较以下 3 组: + +1. `erasure_encode_batched_recv_wait` +2. `erasure_encode_batched_write` +3. `erasure_encode_cpu` + +## 8. 推荐的人工比较方式 + +对每一轮,记录: + +1. throughput +2. avg latency +3. `recv_wait_avg = internal_sum / internal_count` +4. `write_avg = internal_sum / internal_count` +5. `cpu_avg = internal_sum / internal_count` + +重点不是单看 throughput,而是看: + +1. `2` 是否更稳定地降低 `recv_wait_avg` +2. `2` 是否没有把 `write_avg` 或 `cpu_avg` 反向放大 + +## 9. 建议的判定门槛 + +只有同时满足下面两条,才建议继续往“默认值候选”方向推进: + +1. 在至少 `3` 轮对照中,`batch_blocks=2` 的 throughput / latency 有 `>= 2` 轮优于 `4` +2. `recv_wait_avg` 的下降方向在大多数轮次都成立 + +如果只满足其中一条: + +1. 保持它为候选 env 配置 +2. 暂不改代码默认值 + +## 10. 当前建议 + +当前阶段最稳妥的动作顺序是: + +1. 先按这份矩阵补足 `6` 轮 +2. 再做一次汇总表 +3. 最后才决定是否让 multipart batched path 默认使用 `2` diff --git a/docs/operations/issue-712-deeper-zero-copy-next-steps-zh.md b/docs/operations/issue-712-deeper-zero-copy-next-steps-zh.md new file mode 100644 index 000000000..ccf678622 --- /dev/null +++ b/docs/operations/issue-712-deeper-zero-copy-next-steps-zh.md @@ -0,0 +1,72 @@ +# Issue #712 更深一层 zero-copy 下一阶段说明 + +## 1. 当前结论 + +当前分支已经具备: + +1. `rename_data` 的 `msgpack named-map + JSON fallback` +2. ordinary PUT 的 `zero_copy_eager` 实验路径 + +其中 `zero_copy_eager` 已经是实际命中的业务路径,但它当前仍然不是端到端严格意义上的 zero-copy write path。 + +## 2. 当前 copy 还存在的层 + +当前 plain PUT 即使命中了 `zero_copy_eager`,仍然会在后续链路里发生复制,主要位置在: + +1. `HashReader::from_stream(...)` +2. `Erasure::encode(...)` 的 block ingest + +从当前代码 review 看,优先级更高的是: + +1. `Erasure::encode` + +而不是: + +1. `HashReader` + +## 3. 为什么先看 `Erasure::encode` + +原因: + +1. `HashReader` 主要负责包装 `AsyncRead` 与 checksum 语义,不是最直接的热点 +2. `Erasure::encode` 目前仍然是每个 block 读入 `Vec` 后再进入编码 +3. 这更像是 plain PUT 的下一层实际 copy 热点 + +## 4. 这轮已经尝试过但不保留的方向 + +本轮已经试过一个很小的局部改动: + +1. 对 full block 优先走 `encode_data_owned(...)` + +结果: + +1. 对某些 ordinary PUT 面有正向迹象 +2. 但不同 size 下收益不稳定 +3. 因此当前不建议直接基于这个 patch 继续往前推 + +## 5. 下一阶段更稳妥的技术方向 + +如果要继续做更深一层 zero-copy,建议优先考虑: + +1. 为 `Erasure::encode` 设计更稳定的 block buffer 生命周期 +2. 让 full block ingest 更接近 `Bytes` / owned buffer 复用 +3. 避免在当前函数内继续做更多“局部替换一个调用”的小补丁 + +换句话说,下一阶段更适合做: + +1. block buffer 生命周期设计 +2. owned block ring / reusable block pool +3. encode/write 之间更明确的 buffer ownership + +而不是先做: + +1. `HashReader` 级别的大改 +2. 更多无设计托底的局部 `Vec` 调整 + +## 6. 当前建议 + +当前最稳妥的推进顺序: + +1. 先把当前 `zero_copy_eager` 路径继续作为实验性 ordinary PUT 路径保留 +2. 如果要继续做 deeper zero-copy,单独开一个新阶段 +3. 该阶段优先聚焦 `Erasure::encode` ingest / buffer lifecycle,而不是先改 `HashReader` diff --git a/docs/operations/issue-712-encode-write-overlap-observability-summary-zh.md b/docs/operations/issue-712-encode-write-overlap-observability-summary-zh.md new file mode 100644 index 000000000..88e44d478 --- /dev/null +++ b/docs/operations/issue-712-encode-write-overlap-observability-summary-zh.md @@ -0,0 +1,284 @@ +# Issue #712 encode / write overlap 观测小结 + +## 1. 目的 + +本文记录 `#712` 新一轮更细主线的第一步结果: + +1. 不先冒进改 `encode.rs` 调度语义 +2. 先把 `multipart_set_disk_encode` 内部拆成更细阶段 +3. 用 focused benchmark 判断当前 overlap 更像卡在 encode 侧、write 侧,还是 batch barrier + +## 2. 新增内部阶段 + +本轮在 `crates/ecstore/src/erasure_coding/encode.rs` 中补充了内部阶段观测,使用指标: + +1. `rustfs_internal_stage_duration_ms` + +新增阶段: + +1. `erasure_encode_cpu` +2. `erasure_encode_send_wait` +3. `erasure_encode_recv_wait` +4. `erasure_encode_write` +5. `erasure_encode_shutdown` +6. `erasure_encode_batched_send_wait` +7. `erasure_encode_batched_recv_wait` +8. `erasure_encode_batched_write` +9. `erasure_encode_batched_shutdown` + +## 3. focused baseline + +profile: + +1. `2g-128m-pc4` + +默认 batched 配置(`RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=4`)下,本轮 observed run 结果为: + +1. Throughput: `351.96 MiB/s` +2. Avg latency: `5842.1ms` +3. path: `multipart_write_pipeline_batched_large` + +从 raw sum / count 推出的内部均值近似为: + +1. `erasure_encode_cpu`: `~2.95ms` +2. `erasure_encode_batched_send_wait`: `~0.01ms` +3. `erasure_encode_batched_recv_wait`: `~105.3ms` +4. `erasure_encode_batched_write`: `~75.4ms` +5. `erasure_encode_batched_shutdown`: `~0.90ms` + +## 4. 第一轮判断 + +这组数据说明: + +1. `send_wait` 几乎可以忽略,说明 encode producer 基本没有被 queue backpressure 卡住 +2. `recv_wait` 明显高于 `write`,说明 consumer 侧更常见的是在等下一批 encode 结果,而不是 writer 太慢导致队列打满 +3. `cpu` 本身并不大,真正放大的是 batched producer/consumer 之间的批次屏障 + +换句话说,这一轮更像是: + +1. writer 在等 encoder / 等 batch 集齐 +2. 而不是 encoder 在等 writer + +## 5. 配置性验证 + +为了验证 batch barrier 假设,本轮只改一个变量: + +1. `RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=2` + +同样只跑: + +1. `2g-128m-pc4` + +结果: + +1. Throughput: `368.57 MiB/s` +2. Avg latency: `5537.3ms` +3. path: `multipart_write_pipeline_batched_large` + +raw sum / count 推导的内部均值近似为: + +1. `erasure_encode_cpu`: `~2.94ms` +2. `erasure_encode_batched_send_wait`: `~0.01ms` +3. `erasure_encode_batched_recv_wait`: `~48.9ms` +4. `erasure_encode_batched_write`: `~35.9ms` +5. `erasure_encode_batched_shutdown`: `~0.36ms` + +## 6. 当前结论 + +这轮可以先收敛出一个相对明确的方向: + +1. 当前 batched 路径的第一优先问题不像是 encode CPU 绝对太重 +2. 更像是 batch size 偏大,导致 writer 侧等待下一批 encode 结果的时间被放大 +3. 把 `RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS` 从 `4` 降到 `2` 后,结果明显好于默认值 + +## 7. 当前建议 + +下一步如果继续沿着 overlap 这条线推进,建议优先级如下: + +1. 先把 `RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=2` 作为 batched 路径的候选值继续复测 +2. 再决定是否要把 multipart batched path 的 batch size 做成与 ordinary PUT 分离 +3. 在没有更多证据前,不要继续尝试更激进的 batch 内单次 blocking encode 调度改法 + +## 8. 第二轮更稳妥复测 + +为了减少热态偏差,本轮又按 `4 -> 2 -> 4 -> 2` 做了四轮受控复测。 + +profile 固定: + +1. `2g-128m-pc4` + +结果: + +1. `b4-r1`: `317.20 MiB/s`, `6515.2ms` +2. `b2-r1`: `300.38 MiB/s`, `6842.8ms` +3. `b4-r2`: `334.58 MiB/s`, `5966.5ms` +4. `b2-r2`: `358.77 MiB/s`, `5748.2ms` + +从这组数据看: + +1. `batch_blocks=2` 不是每一轮都赢 +2. 但 `b2-r2` 明显优于同组前后的 `b4-r2` +3. 结果仍然存在不小波动,因此还不足以直接改代码默认值 + +## 9. 第二轮内部阶段对比 + +对 `b4-r2` 与 `b2-r2` 的 raw sum / count 做近似均值后,可以看到: + +### `b4-r2` + +1. `erasure_encode_batched_recv_wait`: `~107.4ms` +2. `erasure_encode_batched_write`: `~78.2ms` +3. `erasure_encode_cpu`: `~3.19ms` + +### `b2-r2` + +1. `erasure_encode_batched_recv_wait`: `~50.9ms` +2. `erasure_encode_batched_write`: `~37.7ms` +3. `erasure_encode_cpu`: `~3.06ms` + +这说明: + +1. `batch_blocks=2` 的主要收益方向仍然是降低 batched consumer 侧等待时间 +2. `cpu` 本身没有发生决定性变化 +3. 当前更像是在改善 batch barrier,而不是改变编码计算本体 + +## 10. 当前收敛结论 + +到这一轮为止,更稳妥的结论是: + +1. `RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=2` 仍然值得保留为候选配置 +2. 它对 `erasure_encode_batched_recv_wait` 的改善方向是清晰的 +3. 但吞吐/延迟收益还不够稳定,当前不建议直接改默认值 +4. 下一步更适合继续以环境变量方式复测,而不是马上把默认值写死到代码里 + +## 11. 第三轮补齐后的小结 + +按低噪声矩阵继续补到 `6` 轮之后,结果如下: + +1. `b4-r1`: `317.20 MiB/s`, `6515.2ms` +2. `b2-r1`: `300.38 MiB/s`, `6842.8ms` +3. `b4-r2`: `334.58 MiB/s`, `5966.5ms` +4. `b2-r2`: `358.77 MiB/s`, `5748.2ms` +5. `b4-r3`: `323.35 MiB/s`, `6412.4ms` +6. `b2-r3`: `350.87 MiB/s`, `5827.6ms` + +按组汇总: + +### `batch_blocks=4` + +1. Avg throughput: `325.04 MiB/s` +2. Median throughput: `323.35 MiB/s` +3. Avg latency: `6298.0ms` +4. Median latency: `6412.4ms` + +### `batch_blocks=2` + +1. Avg throughput: `336.67 MiB/s` +2. Median throughput: `350.87 MiB/s` +3. Avg latency: `6139.5ms` +4. Median latency: `5827.6ms` + +这说明: + +1. `batch_blocks=2` 经过 `6` 轮汇总后,组均值已经优于 `4` +2. 但单轮结果仍然存在明显波动,因此还不能把它视为“完全稳定结论” + +## 12. 第三轮内部阶段补充 + +从 `b4-r3` 与 `b2-r3` 的 raw sum / count 看: + +### `b4-r3` + +1. `erasure_encode_batched_recv_wait`: `~115.55ms` +2. `erasure_encode_batched_write`: `~80.15ms` +3. `erasure_encode_cpu`: `~3.55ms` + +### `b2-r3` + +1. `erasure_encode_batched_recv_wait`: `~52.46ms` +2. `erasure_encode_batched_write`: `~37.99ms` +3. `erasure_encode_cpu`: `~3.12ms` + +这与前一组 `b4-r2` / `b2-r2` 的方向一致,说明: + +1. `batch_blocks=2` 主要还是在改善 batch barrier +2. 下降最明显的仍然是 consumer 侧等待时间 +3. `cpu` 本体没有决定性变化 + +## 13. 当前更新后的建议 + +到这一步,建议可以进一步收敛为: + +1. `RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=2` 仍然保留为 multipart batched 路径的强候选配置 +2. 它已经具备“方向明确、组均值更优”的证据 +3. 但由于单轮波动仍在,当前更合适的动作仍然是继续以 env 方式复测,而不是直接改代码默认值 + +## 14. 第四轮补齐后的更新判断 + +继续按同一矩阵补到 `8` 轮之后,新增结果为: + +1. `b4-r4`: `296.34 MiB/s`, `6977.8ms` +2. `b2-r4`: `248.68 MiB/s`, `8310.8ms` + +这样 `8` 轮完整结果为: + +1. `b4-r1`: `317.20 MiB/s`, `6515.2ms` +2. `b2-r1`: `300.38 MiB/s`, `6842.8ms` +3. `b4-r2`: `334.58 MiB/s`, `5966.5ms` +4. `b2-r2`: `358.77 MiB/s`, `5748.2ms` +5. `b4-r3`: `323.35 MiB/s`, `6412.4ms` +6. `b2-r3`: `350.87 MiB/s`, `5827.6ms` +7. `b4-r4`: `296.34 MiB/s`, `6977.8ms` +8. `b2-r4`: `248.68 MiB/s`, `8310.8ms` + +按组重新汇总: + +### `batch_blocks=4` + +1. Avg throughput: `317.87 MiB/s` +2. Median throughput: `320.27 MiB/s` +3. Avg latency: `6468.0ms` +4. Median latency: `6463.8ms` + +### `batch_blocks=2` + +1. Avg throughput: `314.68 MiB/s` +2. Median throughput: `325.62 MiB/s` +3. Avg latency: `6682.4ms` +4. Median latency: `6335.2ms` + +## 15. 第四轮后的结论修正 + +补齐到 `8` 轮之后,需要把结论进一步收紧: + +1. `batch_blocks=2` 仍然能稳定改善 `erasure_encode_batched_recv_wait` +2. 但吞吐/延迟层面的总收益并没有收敛成稳定优势 +3. `b2-r4` 明显把组均值重新拉回,说明它目前还只是“有潜力的候选配置”,不是“已经证实优于 4 的配置” + +从 `b4-r4` 与 `b2-r4` 的 raw sum / count 看: + +### `b4-r4` + +1. `erasure_encode_batched_recv_wait`: `~131.18ms` +2. `erasure_encode_batched_write`: `~82.04ms` +3. `erasure_encode_cpu`: `~3.95ms` + +### `b2-r4` + +1. `erasure_encode_batched_recv_wait`: `~81.49ms` +2. `erasure_encode_batched_write`: `~45.22ms` +3. `erasure_encode_cpu`: `~5.38ms` + +这说明: + +1. `batch_blocks=2` 对 wait/write 的改善方向依旧成立 +2. 但这次同时伴随更差的整体 throughput/latency +3. 当前还不能把局部内部阶段改善直接视为端到端收益 + +## 16. 当前最稳妥的建议 + +到这一步,最稳妥的建议是: + +1. `RUSTFS_ERASURE_ENCODE_BATCH_BLOCKS=2` 继续保留为 env-only 候选配置 +2. 当前不要改代码默认值 +3. 后续如果继续验证,应优先排查为什么 `b2-r4` 会出现这种明显反向波动,而不是继续机械追加更多轮次 diff --git a/docs/operations/issue-712-multipart-put-server-path-validation-runbook-zh.md b/docs/operations/issue-712-multipart-put-server-path-validation-runbook-zh.md new file mode 100644 index 000000000..16be354d3 --- /dev/null +++ b/docs/operations/issue-712-multipart-put-server-path-validation-runbook-zh.md @@ -0,0 +1,213 @@ +# Issue #712 multipart PUT server-path 静默验证 Runbook + +## 1. 目的 + +本文给 `#712` 第二批工作提供一个只关注 multipart server path 的静默验证 runbook。 + +目标: + +1. 不再扩散客户端参数矩阵 +2. 固定当前推荐 baseline +3. 把注意力集中到 server-path 观测增强后的阶段结果 + +## 2. 固定 baseline + +当前固定 baseline: + +1. `1GiB -> 64MiB part / pc4` +2. `2GiB -> 128MiB part / pc4` + +可选补充: + +1. `2GiB -> 256MiB part / pc4` + +但默认不作为首选 baseline。 + +## 3. 推荐脚本 + +直接使用: + +1. `scripts/run_gt1g_multipart_put_server_path_focus.sh` + +该脚本默认只跑: + +1. `1g-64m-pc4` +2. `2g-128m-pc4` + +可选补充: + +1. `2g-256m-pc4` + +## 4. 静默执行命令 + +### 4.1 默认两组 + +```bash +bash scripts/run_gt1g_multipart_put_server_path_focus.sh \ + --host 127.0.0.1:9000 \ + --access-key rustfsadmin \ + --secret-key rustfsadmin \ + --bucket-prefix issue712-multipart-focus \ + --duration 10m \ + --out-dir target/bench/issue712-multipart-server-path-focus +``` + +### 4.2 加上 `2g-256m-pc4` + +```bash +bash scripts/run_gt1g_multipart_put_server_path_focus.sh \ + --host 127.0.0.1:9000 \ + --access-key rustfsadmin \ + --secret-key rustfsadmin \ + --bucket-prefix issue712-multipart-focus \ + --duration 10m \ + --profiles 1g-64m-pc4,2g-128m-pc4,2g-256m-pc4 \ + --out-dir target/bench/issue712-multipart-server-path-focus-wide +``` + +## 5. 强制要求 + +这轮 runbook 的要求是: + +1. 静默跑 +2. 只读 `summary.csv` +3. 如需解释异常,再去看 dashboard / 日志 + +## 6. 结果目录 + +建议统一: + +```text +target/bench/ + issue712-multipart-server-path-focus/ + run_manifest.txt + commands.txt + summary.csv + logs/ + benchdata/ +``` + +## 7. 需要记录的指标 + +最终结果表之外,强制记录以下阶段: + +1. `multipart_ingress_prepare` +2. `multipart_set_disk_writer_setup` +3. `multipart_set_disk_encode` +4. `multipart_complete_tail` + +同时建议固定记录 multipart path 计数: + +1. `multipart_write_pipeline` +2. `multipart_write_pipeline_batched_large` +3. `multipart_write_single_block_non_inline` + +## 8. 本轮的判断顺序 + +先看: + +1. `summary.csv` + +再看: + +1. `multipart_complete_tail` +2. `multipart_set_disk_encode` +3. `multipart_set_disk_writer_setup` +4. `multipart_ingress_prepare` + +## 9. 结果解释 + +### 9.1 `summary.csv` 先分出好坏组合 + +先回答: + +1. `1GiB / 64MiB / pc4` 是否仍是最稳 baseline +2. `2GiB / 128MiB / pc4` 是否仍是最稳 baseline + +### 9.2 再用 Dashboard 回答热点层 + +再回答: + +1. `multipart_complete_tail` 是否最高 +2. `multipart_set_disk_encode` 是否主导 +3. `multipart_set_disk_writer_setup` 是否异常高 +4. `multipart_ingress_prepare` 是否已经被 body buffering 放大 + +## 10. 下一步动作判定 + +### 如果 `multipart_set_disk_encode` 最高 + +下一步优先: + +1. multipart part 专用 batching gate +2. multipart encode path 单独策略 + +### 如果 `multipart_complete_tail` 最高 + +下一步优先: + +1. complete path metadata / checksum / rename tail 优化 + +### 如果 `multipart_set_disk_writer_setup` 最高 + +下一步优先: + +1. writer init / bitrot writer path 优化 + +### 如果 `multipart_ingress_prepare` 最高 + +下一步优先: + +1. part ingress buffer 分层 +2. body read / HashReader 前的缓冲调整 + +## 11. `multipart_*` 指标为空时的排查顺序 + +如果本轮跑的是 multipart PUT,但 Prometheus 里查不到任何 `multipart_*` stage: + +1. 先不要直接判定“新打点无效” +2. 先查当前 `rustfs_s3_put_object_stage_duration_ms` 里到底有哪些 stage +3. 如果只看到了 ordinary PUT 的 `ingress_prepare` / `set_disk_writer_setup` / `set_disk_encode` / `set_disk_rename`,要优先怀疑当前 `127.0.0.1:9000` 上跑的不是预期的新二进制 + +推荐先查: + +```promql +topk( + 40, + count by (__name__, stage) ( + {__name__=~"rustfs_s3_put_object_stage_duration_ms.*"} + ) +) +``` + +如果结果里只有 ordinary stage,而没有: + +1. `multipart_ingress_prepare` +2. `multipart_set_disk_writer_setup` +3. `multipart_set_disk_encode` +4. `multipart_complete_tail` + +则应优先检查: + +1. 本轮 RustFS 进程是否确实来自当前 worktree 的 `target/debug/rustfs` +2. 重启脚本是否真的清掉了旧进程 +3. `RUSTFS_OBS_ENDPOINT` 是否仍然指向当前可查询 backend + +## 12. 已确认的假阴性根因样例 + +在 `2026-06-24` 的第三批继续验证中,出现过一次典型假阴性: + +1. multipart benchmark 已经成功跑完 +2. Prometheus 里却只有 ordinary PUT stage,没有任何 `multipart_*` +3. 根因并不是打点代码失效,而是 `127.0.0.1:9000` 上仍然挂着更早启动的旧 RustFS 进程 + +纠偏方式: + +1. 用当前 worktree 的 `target/debug/rustfs` 前台直接拉起服务 +2. 再跑最小化 focused smoke +3. 立刻查询 `multipart_*` stage + +这次纠偏后的前台复测已经确认: + +1. `multipart_*` 四个阶段可以正常上报 +2. 当前热点仍然稳定落在 `multipart_set_disk_encode` diff --git a/docs/operations/issue-712-multipart-put-size-gated-ab-summary-zh.md b/docs/operations/issue-712-multipart-put-size-gated-ab-summary-zh.md new file mode 100644 index 000000000..32e5b38d7 --- /dev/null +++ b/docs/operations/issue-712-multipart-put-size-gated-ab-summary-zh.md @@ -0,0 +1,109 @@ +# Issue #712 multipart size-gated A/B 结果总结 + +## 1. 本轮目的 + +本轮不是继续扩面 benchmark,而是回答一个更窄的问题: + +1. `set_disk.rs` 上的 multipart size-gated batching 是否真的在运行时生效 +2. 如果生效,`2g-128m-pc4` 下是否优于纯 pipeline + +## 2. 先修正的运行时问题 + +在继续验证中发现一个关键问题: + +1. multipart 路径的 batching 分类使用了 `data.size()` +2. 但原逻辑是在 `data.stream` 被替换为空 reader 之后才读取这个 size +3. 这会导致 multipart size gate 在运行时无法按预期命中 + +本轮已在 multipart 路径上修正为: + +1. 先捕获原始 `multipart_part_size` +2. 再执行 stream swap +3. 再基于原始 size 做 `classify_multipart_part_write_path(...)` + +同时补充了 multipart path 计数标签: + +1. `multipart_write_pipeline` +2. `multipart_write_pipeline_batched_large` +3. `multipart_write_single_block_non_inline` + +## 3. focused 对照矩阵 + +本轮只保留受影响的 profile: + +1. `2g-128m-pc4` + +并做两组对照: + +1. pipeline 基线: + - `RUSTFS_MULTIPART_PUT_LARGE_BATCH_MIN_SIZE_BYTES=10737418240` +2. 默认门槛实验: + - 使用默认 `128MiB` 门槛 + +附加做了一组强制 batched 校验: + +1. `RUSTFS_MULTIPART_PUT_LARGE_BATCH_MIN_SIZE_BYTES=1` + +## 4. 结果 + +### 4.1 修正后的 pipeline 基线 + +结果目录: + +1. `target/bench/issue712-multipart-server-path-fixed-pipeline/summary.csv` + +结果: + +1. Throughput: `364.83 MiB/s` +2. Avg latency: `5604.5ms` +3. Path counter: 仅出现 `multipart_write_pipeline` +4. `multipart_set_disk_encode` P95: `7348.88ms` + +### 4.2 修正后的默认门槛实验 + +结果目录: + +1. `target/bench/issue712-multipart-server-path-fixed-batched/summary.csv` + +结果: + +1. Throughput: `359.00 MiB/s` +2. Avg latency: `5734.9ms` +3. Path counter: 已出现 `multipart_write_pipeline_batched_large` +4. `multipart_set_disk_encode` P95: `7375ms` + +说明: + +1. batched 路径在修正后已经真正开始命中 +2. 但当前结果并没有优于纯 pipeline + +### 4.3 强制 batched 校验 + +结果目录: + +1. `target/bench/issue712-multipart-server-path-forced-batched/summary.csv` + +结果: + +1. Throughput: `362.00 MiB/s` +2. Avg latency: `5664.3ms` + +这组结果同样没有优于修正后的 pipeline 基线。 + +## 5. 本轮结论 + +本轮可以明确收敛出三点: + +1. multipart size-gated batching 的运行时命中问题已经被定位并修正 +2. batched path 修正后确实可以被 Prometheus path counter 观察到 +3. 在当前单机多盘 `2g-128m-pc4` focused 验证下,默认 `128MiB` 门槛没有带来正收益,反而略逊于纯 pipeline + +## 6. 当前建议 + +在当前证据下,不建议把 multipart batched path 作为默认推荐优化结论直接推进。 + +更稳妥的后续方向是: + +1. 先保留这条路径为可控实验能力 +2. 继续围绕 `multipart_set_disk_encode` 本身做更细粒度优化 +3. 如果还要继续试 batching,应先解释为什么 `128MiB` part 在当前实现下没有优于 pipeline,而不是直接继续扩大默认启用范围 diff --git a/docs/operations/issue-712-multipart-put-stage-metrics-dashboard-guide-zh.md b/docs/operations/issue-712-multipart-put-stage-metrics-dashboard-guide-zh.md new file mode 100644 index 000000000..18f524912 --- /dev/null +++ b/docs/operations/issue-712-multipart-put-stage-metrics-dashboard-guide-zh.md @@ -0,0 +1,201 @@ +# Issue #712 multipart PUT 分阶段指标 Dashboard / PromQL 指南 + +## 1. 目的 + +本文给 `#712` 的第一批 server-path 观测增强配套一份可执行的 Dashboard / PromQL 指南。 + +本批次新增的 multipart 阶段指标依然复用现有指标名: + +1. `rustfs_s3_put_object_stage_duration_ms` + +但新增了四个 stage label: + +1. `multipart_ingress_prepare` +2. `multipart_set_disk_writer_setup` +3. `multipart_set_disk_encode` +4. `multipart_complete_tail` + +## 2. 使用前提 + +这些阶段指标严格受全局开关控制: + +1. `rustfs_io_metrics::put_stage_metrics_enabled() == true` + +如果该开关没有开启: + +1. 不会上报这些阶段指标 +2. 也不会额外做阶段计时 + +## 3. 推荐直接复用现有 Grafana Row + +当前 Dashboard 中已经有: + +1. `Large PUT Stage Breakdown` + +这意味着: + +1. 不需要重新设计一套全新 row +2. 只需要在现有 row / stage 变量里选新的 multipart stage label 即可 + +## 4. 推荐 PromQL + +### 4.1 multipart 阶段 P95 + +```promql +histogram_quantile( + 0.95, + sum by (stage, le) ( + rate( + rustfs_s3_put_object_stage_duration_ms_bucket{ + job=~"$job", + stage=~"multipart_.*" + }[$__rate_interval] + ) + ) +) +``` + +### 4.2 multipart 阶段 P99 + +```promql +histogram_quantile( + 0.99, + sum by (stage, le) ( + rate( + rustfs_s3_put_object_stage_duration_ms_bucket{ + job=~"$job", + stage=~"multipart_.*" + }[$__rate_interval] + ) + ) +) +``` + +### 4.3 单实例 multipart 阶段 P95 + +```promql +histogram_quantile( + 0.95, + sum by (instance, stage, le) ( + rate( + rustfs_s3_put_object_stage_duration_ms_bucket{ + job=~"$job", + instance=~"$instance", + stage=~"multipart_.*" + }[$__rate_interval] + ) + ) +) +``` + +### 4.4 multipart 与 ordinary PUT encode 对比 + +```promql +histogram_quantile( + 0.95, + sum by (stage, le) ( + rate( + rustfs_s3_put_object_stage_duration_ms_bucket{ + job=~"$job", + stage=~"set_disk_encode|multipart_set_disk_encode" + }[$__rate_interval] + ) + ) +) +``` + +### 4.5 multipart complete tail 重点盯盘 + +```promql +histogram_quantile( + 0.95, + sum by (instance, le) ( + rate( + rustfs_s3_put_object_stage_duration_ms_bucket{ + job=~"$job", + stage="multipart_complete_tail", + instance=~"$instance" + }[$__rate_interval] + ) + ) +) +``` + +## 5. 推荐看板顺序 + +当你在看 `>1GiB multipart PUT` 时,建议按下面顺序看: + +1. `multipart_ingress_prepare` +2. `multipart_set_disk_writer_setup` +3. `multipart_set_disk_encode` +4. `multipart_complete_tail` + +解释顺序: + +1. 如果 ingress 先高,先看 part ingress buffer / request-body handling +2. 如果 writer setup 高,先看 bitrot writer / disk availability / shard_file_size path +3. 如果 encode 高,先看 multipart 是否需要独立 encode strategy +4. 如果 complete tail 高,优先看 `complete_multipart_upload()` 的 metadata / checksum / rename tail + +## 6. 推荐结合看的辅助指标 + +建议和上面四个阶段一起看: + +1. `rustfs_io_put_object_concurrent_requests` +2. `rustfs_ec_encode_inflight_bytes_current` +3. host CPU +4. per-instance disk write throughput +5. readiness / write quorum 异常计数 + +## 7. 典型解释模板 + +### 7.1 ingress 高 + +可能原因: + +1. part body stream buffering 不合适 +2. `part.size` 与 ingress buffer 不匹配 + +### 7.2 writer setup 高 + +可能原因: + +1. bitrot writer 构建成本偏高 +2. online disk / writer init 慢 +3. shard_file_size 相关路径有额外成本 + +### 7.3 encode 高 + +可能原因: + +1. multipart part 仍然借用了 ordinary PUT encode 行为 +2. `part.size` 太大,单 part encode CPU 时间过长 +3. batching / inflight 参数不合适 + +### 7.4 complete tail 高 + +可能原因: + +1. complete 阶段 part metadata 处理放大 +2. checksum combine 成本高 +3. rename / cleanup / commit tail 成本高 + +## 8. 建议的截图 / 归档内容 + +每次 `>1GiB multipart PUT` 复测,建议固定归档: + +1. multipart stage P95 截图 +2. multipart stage P99 截图 +3. `multipart_complete_tail` 单实例截图 +4. CPU / disk write 辅助图 + +## 9. 当前阶段建议 + +下一次进入 `#712` 继续推进时: + +1. 先开 `put_stage_metrics_enabled` +2. 先跑推荐 baseline: + - `1GiB -> 64MiB / pc4` + - `2GiB -> 128MiB / pc4` +3. 先看 `multipart_complete_tail` 是否明显高于其他阶段 +4. 再决定是先改 ingress / encode / writer setup / complete tail diff --git a/docs/operations/issue-712-third-batch-false-negative-followup-zh.md b/docs/operations/issue-712-third-batch-false-negative-followup-zh.md new file mode 100644 index 000000000..79080112b --- /dev/null +++ b/docs/operations/issue-712-third-batch-false-negative-followup-zh.md @@ -0,0 +1,107 @@ +# Issue #712 第三批继续验证结果总结 + +## 1. 背景 + +`#712` 第三批的第一轮静默 baseline 已经确认: + +1. `multipart_set_disk_encode` 是当前最主要热点 +2. `multipart_complete_tail` 可见但不是第一瓶颈 + +在继续推进 multipart encode 优化线时,又出现了一次“指标空结果”的假阴性,需要单独记录,避免后续团队误判为打点无效。 + +## 2. 假阴性根因 + +现象: + +1. multipart benchmark 成功完成 +2. `summary.csv` 正常生成 +3. Prometheus 中却只出现 ordinary PUT 的 stage: + - `ingress_prepare` + - `set_disk_writer_setup` + - `set_disk_encode` + - `set_disk_rename` +4. 没有任何 `multipart_*` stage + +根因: + +1. 本轮 benchmark 实际命中了 `127.0.0.1:9000` 上残留的旧 RustFS 进程 +2. 该旧进程不是当前 worktree 的新二进制 +3. 因此即使 benchmark 是 multipart PUT,也不会产出本批新增的 multipart stage 标签 + +这次问题的本质不是“新打点代码无效”,而是“验证目标进程不对”。 + +## 3. 如何识别这类假阴性 + +推荐直接查询: + +```promql +topk( + 40, + count by (__name__, stage) ( + {__name__=~"rustfs_s3_put_object_stage_duration_ms.*"} + ) +) +``` + +如果你跑的是 multipart PUT,但结果里只有 ordinary PUT stage,而没有: + +1. `multipart_ingress_prepare` +2. `multipart_set_disk_writer_setup` +3. `multipart_set_disk_encode` +4. `multipart_complete_tail` + +则优先判断为“验证目标进程可能不对”,而不要先判断为“metrics 打点无效”。 + +## 4. 前台纠偏复测 + +为了排除旧进程干扰,本轮直接使用当前 worktree 的 `target/debug/rustfs` 前台拉起服务,再做最小化 focused smoke。 + +复测 profile: + +1. `2g-128m-pc4` + +复测结果: + +1. Throughput: `373.78 MiB/s` +2. Request rate: `2.92 obj/s` +3. Avg latency: `5471.4ms` + +结果目录: + +1. `target/bench/issue712-multipart-server-path-foreground-smoke/summary.csv` + +## 5. 前台复测阶段指标 + +Prometheus 近窗口查询已经确认 multipart 四阶段都正常出现。 + +P95: + +1. `multipart_ingress_prepare`: `4.75ms` +2. `multipart_set_disk_writer_setup`: `4.86ms` +3. `multipart_set_disk_encode`: `7375ms` +4. `multipart_complete_tail`: `13ms` + +同时还能看到 ordinary PUT 的阶段指标,但这不影响判断;关键是: + +1. multipart stage 已经实际落盘到 TSDB +2. encode 依旧是最显著热点 + +## 6. 当前直接结论 + +这轮继续验证后的结论没有变化,反而更稳: + +1. `multipart_set_disk_encode` 仍然是 `>1GiB multipart PUT` 当前最主要优化目标 +2. `multipart_complete_tail` 仍然是次级问题,不应抢在 encode 之前 +3. `multipart_ingress_prepare` 和 `multipart_set_disk_writer_setup` 暂时不是第一优先级 +4. 后续优化应继续围绕 multipart encode path,而不是先转去 complete tail / ingress + +## 7. 对后续验证的要求 + +后续再做 multipart server-path 对照验证时,建议强制执行: + +1. 先确认服务进程确实来自当前 worktree 二进制 +2. 先确认 `multipart_*` stage 能被 Prometheus 查询到 +3. 再读取 `summary.csv` +4. 最后再判断 encode 优化是否真的有效 + +否则很容易再次出现“summary 正常,但阶段指标是旧进程数据”的假阴性。 diff --git a/docs/operations/issue-712-zero-copy-eager-put-ops-guide-zh.md b/docs/operations/issue-712-zero-copy-eager-put-ops-guide-zh.md new file mode 100644 index 000000000..026d444f4 --- /dev/null +++ b/docs/operations/issue-712-zero-copy-eager-put-ops-guide-zh.md @@ -0,0 +1,153 @@ +# Issue #712 `zero_copy_eager` plain PUT 验证手册 + +## 1. 目的 + +本文用于记录当前 `zero_copy_eager` plain PUT 路径的实际使用方式、验证命令和当前边界。 + +这条路径当前的目标不是“端到端完全零拷贝”,而是先把 ordinary PUT 从“只有 zero-copy eligibility / metrics”推进到: + +1. 真实存在的业务路径 +2. 真实命中可观测的 `put_path` +3. 先减少请求体聚合阶段的额外复制 + +## 2. 当前启用条件 + +当前 `zero_copy_eager` 只在下面条件同时满足时才会命中: + +1. 非加密 +2. 非压缩 +3. 非 extract 请求 +4. 对象大小 `> 1MiB` +5. 对象大小 `<= 32MiB` +6. 请求长度已知 +7. 不属于需要特殊处理的 aws-chunked 未知长度场景 + +## 3. 当前实现边界 + +这条路径已经是真实业务路径,但当前还不是端到端严格意义上的 zero-copy。 + +已经做到: + +1. 请求体 chunk 以 `Bytes` 形式进入 `zero_copy_eager` +2. 不再先把整个对象拼成一个大 `Vec` 再进入后续写路径 +3. 运行时会记录 `put_path=zero_copy_eager` +4. 会真实记录 `rustfs_zero_copy_write_total` + +尚未做到: + +1. `HashReader` 之后完全无复制 +2. `Erasure::encode` 的 block ingest 完全无复制 +3. shard write 全链路严格 zero-copy + +因此当前最准确的描述是: + +1. 这是 ordinary PUT 的真实 zero-copy eager ingress 路径 +2. 不是完整的 end-to-end zero-copy write path + +## 4. 推荐验证命令 + +### 4.1 启动本地单机多盘 RustFS + +```bash +env \ + RUSTFS_UNSAFE_BYPASS_DISK_CHECK=true \ + RUSTFS_ADDRESS=127.0.0.1:9000 \ + RUSTFS_ACCESS_KEY=rustfsadmin \ + RUSTFS_SECRET_KEY=rustfsadmin \ + RUSTFS_RPC_SECRET=rustfs-rpc-secret \ + RUSTFS_REGION=us-east-1 \ + RUSTFS_CONSOLE_ENABLE=false \ + RUSTFS_OBS_ENDPOINT=http://127.0.0.1:4318 \ + target/debug/rustfs server \ + /private/tmp/issue708-single-node-multidisk/d1 \ + /private/tmp/issue708-single-node-multidisk/d2 \ + /private/tmp/issue708-single-node-multidisk/d3 \ + /private/tmp/issue708-single-node-multidisk/d4 +``` + +### 4.2 小面 ordinary PUT 验证 + +```bash +bash scripts/run_put_large_stage_breakdown.sh \ + --endpoint http://127.0.0.1:9000 \ + --access-key rustfsadmin \ + --secret-key rustfsadmin \ + --sizes 16MiB,32MiB \ + --concurrencies 16 \ + --duration 60s \ + --rounds 1 \ + --retry-per-round 1 \ + --retry-sleep-secs 2 \ + --cooldown-secs 15 \ + --out-dir target/bench/issue712-zero-copy-eager-put-verify +``` + +## 5. 运行时必须核对的指标 + +### 5.1 PUT path 命中 + +必须确认: + +```promql +rustfs_s3_put_object_path_total{ + path=~"zero_copy_eager|small_eager|streaming|stream_compressed" +} +``` + +如果看不到 `zero_copy_eager`: + +1. 说明本轮没有真正命中这条路径 +2. 不能据此评价它的收益 + +### 5.2 zero-copy write 指标 + +建议同时看: + +```promql +rustfs_zero_copy_write_total +``` + +```promql +rustfs_zero_copy_write_size_bytes_sum +``` + +### 5.3 普通 PUT summary + +最终固定读取: + +1. `aggregate_median_summary.csv` + +## 6. 当前已知结果 + +当前分支上已经验证过一次小面 ordinary PUT: + +1. `16MiB, c16`: `318.77 MiB/s`, `851.3ms` +2. `32MiB, c16`: `276.38 MiB/s`, `1846.5ms` + +并确认运行时命中了: + +1. `put_path=zero_copy_eager` + +这说明: + +1. 这条路径已经不只是 eligibility / metrics +2. 它已经是真实参与 ordinary PUT 的运行时路径 + +## 7. 当前最稳妥的使用建议 + +当前阶段建议把这条路径当作: + +1. 实验性 ordinary PUT 优化路径 +2. 需要继续压测验证的真实实现 + +当前不建议把它描述成: + +1. 已完成的端到端 zero-copy write path + +## 8. 下一步方向 + +如果要继续把收益往下穿透,下一阶段优先顺序建议是: + +1. 先继续验证 `zero_copy_eager` 的普通 PUT 命中率和端到端收益 +2. 再评估是否要继续下钻 `Erasure::encode` 的 block ingest +3. `HashReader` 层暂时不是第一优先级 diff --git a/rustfs/src/app/multipart_usecase.rs b/rustfs/src/app/multipart_usecase.rs index 6829ff479..8a1051ab3 100644 --- a/rustfs/src/app/multipart_usecase.rs +++ b/rustfs/src/app/multipart_usecase.rs @@ -800,6 +800,7 @@ impl DefaultMultipartUsecase { .map_err(ApiError::from)?; let mut size = size.ok_or_else(|| s3_error!(UnexpectedContent))?; + let ingress_stage_start = rustfs_io_metrics::put_stage_metrics_enabled().then(std::time::Instant::now); // Apply adaptive buffer sizing based on part size for optimal streaming performance. // Uses workload profile configuration (enabled by default) to select appropriate buffer size. @@ -931,6 +932,13 @@ impl DefaultMultipartUsecase { let mut reader = PutObjReader::new(reader); + if let Some(stage_start) = ingress_stage_start { + rustfs_io_metrics::record_put_object_stage_duration( + "multipart_ingress_prepare", + stage_start.elapsed().as_secs_f64() * 1000.0, + ); + } + let info = store .put_object_part(&bucket, &key, &upload_id, part_id, &mut reader, &opts) .await diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 2e16fbc84..07ece3f83 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -430,6 +430,43 @@ impl AsyncRead for PooledBufferReader { } } +struct ChunkedBytesReader { + chunks: Vec, + chunk_index: usize, + chunk_offset: usize, +} + +impl ChunkedBytesReader { + fn new(chunks: Vec) -> Self { + Self { + chunks, + chunk_index: 0, + chunk_offset: 0, + } + } +} + +impl AsyncRead for ChunkedBytesReader { + fn poll_read(mut self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + while self.chunk_index < self.chunks.len() { + let chunk = &self.chunks[self.chunk_index]; + if self.chunk_offset >= chunk.len() { + self.chunk_index += 1; + self.chunk_offset = 0; + continue; + } + + let remaining = &chunk[self.chunk_offset..]; + let to_read = remaining.len().min(buf.remaining()); + buf.put_slice(&remaining[..to_read]); + self.chunk_offset += to_read; + return Poll::Ready(Ok(())); + } + + Poll::Ready(Ok(())) + } +} + /// Determine if zero-copy write should be used for this PutObject operation. /// /// Zero-copy is beneficial for large objects without encryption or compression. @@ -485,6 +522,34 @@ fn should_use_zero_copy(size: i64, headers: &HeaderMap) -> bool { true } +fn should_use_zero_copy_eager_put_path( + size: i64, + headers: &HeaderMap, + server_side_encryption_requested: bool, + should_compress: bool, + is_extract: bool, +) -> bool { + const ZERO_COPY_EAGER_PUT_MAX_SIZE: i64 = 32 * 1024 * 1024; + + if is_extract || should_compress || server_side_encryption_requested { + return false; + } + + if size <= 0 || size > ZERO_COPY_EAGER_PUT_MAX_SIZE { + return false; + } + + if !should_use_zero_copy(size, headers) { + return false; + } + + if request_uses_aws_chunked(headers) && decoded_content_length_from_headers(headers).ok().flatten().is_none() { + return false; + } + + true +} + fn has_put_sse_request_headers(headers: &HeaderMap) -> bool { headers.get(AMZ_SERVER_SIDE_ENCRYPTION).is_some() || headers.get(AMZ_SERVER_SIDE_ENCRYPTION_CUSTOMER_ALGORITHM).is_some() @@ -584,6 +649,41 @@ where Ok(std::io::Cursor::new(buf)) } +async fn read_zero_copy_put_body_exact(mut body: S, size: usize) -> S3Result +where + S: futures::Stream> + Unpin, + E: std::fmt::Display, +{ + let mut chunks = Vec::new(); + let mut filled = 0usize; + + while filled < size { + let Some(chunk) = body.next().await else { + return Err(s3_error!(IncompleteBody)); + }; + let chunk = chunk.map_err(|err| ApiError::from(StorageError::other(err.to_string())))?; + if chunk.is_empty() { + continue; + } + if filled.saturating_add(chunk.len()) > size { + return Err(s3_error!(UnexpectedContent)); + } + + rustfs_io_metrics::record_zero_copy_buffer_operation("put_chunk", chunk.len()); + filled += chunk.len(); + chunks.push(chunk); + } + + while let Some(chunk) = body.next().await { + let chunk = chunk.map_err(|err| ApiError::from(StorageError::other(err.to_string())))?; + if !chunk.is_empty() { + return Err(s3_error!(UnexpectedContent)); + } + } + + Ok(ChunkedBytesReader::new(chunks)) +} + fn object_seek_support_threshold() -> usize { static OBJECT_SEEK_SUPPORT_THRESHOLD: OnceLock = OnceLock::new(); *OBJECT_SEEK_SUPPORT_THRESHOLD.get_or_init(|| { @@ -2076,8 +2176,12 @@ impl DefaultObjectUsecase { let use_small_eager_put_path = should_use_small_eager_put_path(size, &req.headers, server_side_encryption_requested, should_compress, false); + let use_zero_copy_eager_put_path = + should_use_zero_copy_eager_put_path(size, &req.headers, server_side_encryption_requested, should_compress, false); let put_path = if should_compress { "stream_compressed" + } else if use_zero_copy_eager_put_path { + "zero_copy_eager" } else if use_small_eager_put_path { "small_eager" } else { @@ -2229,7 +2333,12 @@ impl DefaultObjectUsecase { write_plan = write_plan.with_compression(algorithm); hrd } else { - if use_small_eager_put_path { + if use_zero_copy_eager_put_path { + let zero_copy_start = std::time::Instant::now(); + let eager_body = read_zero_copy_put_body_exact(body, actual_size as usize).await?; + rustfs_io_metrics::record_zero_copy_write(actual_size as usize, zero_copy_start.elapsed().as_secs_f64() * 1000.0); + HashReader::from_stream(eager_body, size, actual_size, md5hex, sha256hex, false).map_err(ApiError::from)? + } else if use_small_eager_put_path { if (actual_size as usize) <= POOL_BYPASS_MAX_SIZE { // Bypass BytesPool for very small objects to avoid Small-tier // Mutex contention under high concurrency. Direct allocation @@ -5486,6 +5595,24 @@ mod tests { assert!(!should_use_small_eager_put_path(1024 * 1024 + 1, &headers, false, false, false)); } + #[test] + fn should_use_zero_copy_eager_put_path_allows_large_plain_objects_within_cap() { + let headers = HeaderMap::new(); + + assert!(should_use_zero_copy_eager_put_path(2 * 1024 * 1024, &headers, false, false, false)); + assert!(should_use_zero_copy_eager_put_path(32 * 1024 * 1024, &headers, false, false, false)); + assert!(!should_use_zero_copy_eager_put_path(32 * 1024 * 1024 + 1, &headers, false, false, false)); + } + + #[test] + fn should_use_zero_copy_eager_put_path_rejects_compression_sse_and_extract() { + let headers = HeaderMap::new(); + + assert!(!should_use_zero_copy_eager_put_path(2 * 1024 * 1024, &headers, true, false, false)); + assert!(!should_use_zero_copy_eager_put_path(2 * 1024 * 1024, &headers, false, true, false)); + assert!(!should_use_zero_copy_eager_put_path(2 * 1024 * 1024, &headers, false, false, true)); + } + #[tokio::test] async fn read_small_put_body_exact_pooled_reads_exact_bytes() { let pool = get_concurrency_manager().bytes_pool(); @@ -5511,6 +5638,42 @@ mod tests { assert_eq!(err.code(), &S3ErrorCode::IncompleteBody); } + #[tokio::test] + async fn read_zero_copy_put_body_exact_reads_chunked_body() { + use tokio::io::AsyncReadExt; + + let body = futures::stream::iter(vec![ + Ok::(Bytes::from_static(b"hello ")), + Ok::(Bytes::from_static(b"world")), + ]); + + let mut reader = read_zero_copy_put_body_exact(body, 11) + .await + .expect("zero-copy eager body read should succeed"); + let mut out = Vec::new(); + reader + .read_to_end(&mut out) + .await + .expect("chunked bytes reader should be readable"); + + assert_eq!(out, b"hello world"); + } + + #[tokio::test] + async fn read_zero_copy_put_body_exact_rejects_extra_bytes() { + let body = futures::stream::iter(vec![ + Ok::(Bytes::from_static(b"hello")), + Ok::(Bytes::from_static(b"!")), + ]); + + let err = match read_zero_copy_put_body_exact(body, 5).await { + Ok(_) => panic!("extra bytes should fail"), + Err(err) => err, + }; + + assert_eq!(err.code(), &S3ErrorCode::UnexpectedContent); + } + #[tokio::test] async fn pooled_buffer_reader_keeps_buffer_alive_until_consumed() { use tokio::io::AsyncReadExt; diff --git a/rustfs/src/storage/rpc/disk.rs b/rustfs/src/storage/rpc/disk.rs index 13c09f9d6..fa39af177 100644 --- a/rustfs/src/storage/rpc/disk.rs +++ b/rustfs/src/storage/rpc/disk.rs @@ -38,6 +38,14 @@ fn encode_msgpack(value: &T, value_name: &str) -> std::resu Ok(serializer.into_inner()) } +fn encode_msgpack_named(value: &T, value_name: &str) -> std::result::Result, DiskError> { + let mut serializer = rmp_serde::Serializer::new(Vec::new()).with_struct_map(); + value + .serialize(&mut serializer) + .map_err(|err| DiskError::other(format!("encode {value_name} named msgpack failed: {err}")))?; + Ok(serializer.into_inner()) +} + impl NodeService { pub(super) async fn handle_disk_info(&self, request: Request) -> Result, Status> { let request = request.into_inner(); @@ -629,12 +637,13 @@ impl NodeService { ) -> Result, Status> { let request = request.into_inner(); if let Some(disk) = self.find_disk(&request.disk).await { - let file_info = match serde_json::from_str::(&request.file_info) { + let file_info = match decode_msgpack_or_json::(&request.file_info_bin, &request.file_info, "FileInfo") { Ok(file_info) => file_info, Err(err) => { return Ok(Response::new(RenameDataResponse { success: false, rename_data_resp: String::new(), + rename_data_resp_bin: Vec::new().into(), error: Some(DiskError::other(format!("decode FileInfo failed: {err}")).into()), })); } @@ -644,25 +653,33 @@ impl NodeService { .await { Ok(rename_data_resp) => { - let rename_data_resp = match serde_json::to_string(&rename_data_resp) { - Ok(file_info) => file_info, - Err(err) => { - return Ok(Response::new(RenameDataResponse { - success: false, - rename_data_resp: String::new(), - error: Some(DiskError::other(format!("encode data failed: {err}")).into()), - })); - } - }; - Ok(Response::new(RenameDataResponse { - success: true, - rename_data_resp, - error: None, - })) + let rename_data_resp_json = serde_json::to_string(&rename_data_resp); + let rename_data_resp_bin = encode_msgpack_named(&rename_data_resp, "RenameDataResp"); + match (rename_data_resp_json, rename_data_resp_bin) { + (Ok(rename_data_resp), Ok(rename_data_resp_bin)) => Ok(Response::new(RenameDataResponse { + success: true, + rename_data_resp, + rename_data_resp_bin: rename_data_resp_bin.into(), + error: None, + })), + (Err(err), _) => Ok(Response::new(RenameDataResponse { + success: false, + rename_data_resp: String::new(), + rename_data_resp_bin: Vec::new().into(), + error: Some(DiskError::other(format!("encode data failed: {err}")).into()), + })), + (_, Err(err)) => Ok(Response::new(RenameDataResponse { + success: false, + rename_data_resp: String::new(), + rename_data_resp_bin: Vec::new().into(), + error: Some(DiskError::other(format!("encode data failed: {err}")).into()), + })), + } } Err(err) => Ok(Response::new(RenameDataResponse { success: false, rename_data_resp: String::new(), + rename_data_resp_bin: Vec::new().into(), error: Some(err.into()), })), } @@ -670,6 +687,7 @@ impl NodeService { Ok(Response::new(RenameDataResponse { success: false, rename_data_resp: String::new(), + rename_data_resp_bin: Vec::new().into(), error: Some(DiskError::other("can not find disk".to_string()).into()), })) } diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 44bcf8716..9d686b141 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -1598,6 +1598,7 @@ mod tests { dst_volume: "dst-volume".to_string(), dst_path: "dst-path".to_string(), file_info: "{}".to_string(), + file_info_bin: Vec::new().into(), }); let response = service.rename_data(request).await; @@ -1619,6 +1620,7 @@ mod tests { dst_volume: "dst-volume".to_string(), dst_path: "dst-path".to_string(), file_info: "invalid json".to_string(), + file_info_bin: Vec::new().into(), }); let response = service.rename_data(request).await; diff --git a/scripts/restart_local_single_node_multidisk_rustfs.sh b/scripts/restart_local_single_node_multidisk_rustfs.sh new file mode 100644 index 000000000..e6ce8e099 --- /dev/null +++ b/scripts/restart_local_single_node_multidisk_rustfs.sh @@ -0,0 +1,81 @@ +#!/usr/bin/env bash +set -euo pipefail + +PIDFILE=/private/tmp/issue708-single-node-multidisk/rustfs.pid +LOGFILE=/private/tmp/issue708-single-node-multidisk/rustfs.log +RUSTFS_CMD_PATTERN='target/debug/rustfs server /private/tmp/issue708-single-node-multidisk/d1 /private/tmp/issue708-single-node-multidisk/d2 /private/tmp/issue708-single-node-multidisk/d3 /private/tmp/issue708-single-node-multidisk/d4' + +mkdir -p /private/tmp/issue708-single-node-multidisk/{d1,d2,d3,d4} + +if [[ -f "$PIDFILE" ]]; then + old_pid="$(cat "$PIDFILE")" + if kill -0 "$old_pid" >/dev/null 2>&1; then + kill "$old_pid" >/dev/null 2>&1 || true + sleep 2 + fi +fi + +lingering_pids="$(pgrep -f "$RUSTFS_CMD_PATTERN" || true)" +if [[ -n "$lingering_pids" ]]; then + while IFS= read -r lingering_pid; do + [[ -z "$lingering_pid" ]] && continue + kill "$lingering_pid" >/dev/null 2>&1 || true + done <<< "$lingering_pids" + sleep 2 +fi + +if [[ -n "${RUSTFS_TUNING_ENV_FILE:-}" && -f "${RUSTFS_TUNING_ENV_FILE:-}" ]]; then + set -a + source "$RUSTFS_TUNING_ENV_FILE" + set +a +fi + +launch_rustfs() { + if command -v setsid >/dev/null 2>&1; then + setsid env \ + RUSTFS_UNSAFE_BYPASS_DISK_CHECK=true \ + RUSTFS_ADDRESS=127.0.0.1:9000 \ + RUSTFS_ACCESS_KEY=rustfsadmin \ + RUSTFS_SECRET_KEY=rustfsadmin \ + RUSTFS_RPC_SECRET=rustfs-rpc-secret \ + RUSTFS_REGION=us-east-1 \ + RUSTFS_CONSOLE_ENABLE=false \ + RUSTFS_OBS_ENDPOINT=http://127.0.0.1:4318 \ + target/debug/rustfs server \ + /private/tmp/issue708-single-node-multidisk/d1 \ + /private/tmp/issue708-single-node-multidisk/d2 \ + /private/tmp/issue708-single-node-multidisk/d3 \ + /private/tmp/issue708-single-node-multidisk/d4 \ + "$LOGFILE" 2>&1 & + else + nohup env \ + RUSTFS_UNSAFE_BYPASS_DISK_CHECK=true \ + RUSTFS_ADDRESS=127.0.0.1:9000 \ + RUSTFS_ACCESS_KEY=rustfsadmin \ + RUSTFS_SECRET_KEY=rustfsadmin \ + RUSTFS_RPC_SECRET=rustfs-rpc-secret \ + RUSTFS_REGION=us-east-1 \ + RUSTFS_CONSOLE_ENABLE=false \ + RUSTFS_OBS_ENDPOINT=http://127.0.0.1:4318 \ + target/debug/rustfs server \ + /private/tmp/issue708-single-node-multidisk/d1 \ + /private/tmp/issue708-single-node-multidisk/d2 \ + /private/tmp/issue708-single-node-multidisk/d3 \ + /private/tmp/issue708-single-node-multidisk/d4 \ + "$LOGFILE" 2>&1 & + fi +} + +launch_rustfs + +new_pid=$! +sleep 5 + +actual_pid="$(pgrep -f "$RUSTFS_CMD_PATTERN" | head -n 1 || true)" +if [[ -z "$actual_pid" ]]; then + echo "failed to locate restarted rustfs process; launcher pid was $new_pid" >&2 + exit 1 +fi + +echo "$actual_pid" > "$PIDFILE" +curl -fsS http://127.0.0.1:9000/health >/dev/null diff --git a/scripts/run_gt1g_multipart_put_matrix.sh b/scripts/run_gt1g_multipart_put_matrix.sh new file mode 100644 index 000000000..2e5dc1dbd --- /dev/null +++ b/scripts/run_gt1g_multipart_put_matrix.sh @@ -0,0 +1,230 @@ +#!/usr/bin/env bash +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +PROJECT_ROOT="$(cd "${SCRIPT_DIR}/.." && pwd)" + +HOST="${HOST:-127.0.0.1:9000}" +ACCESS_KEY="${ACCESS_KEY:-}" +SECRET_KEY="${SECRET_KEY:-}" +BUCKET_PREFIX="${BUCKET_PREFIX:-rustfs-gt1g-multipart}" +LOOKUP="${LOOKUP:-path}" +DURATION="${DURATION:-10m}" +WARP_BIN="${WARP_BIN:-warp}" +OUT_DIR="${OUT_DIR:-${PROJECT_ROOT}/target/bench/gt1g-multipart-put-$(date +%Y%m%d-%H%M%S)}" +DRY_RUN=false + +DEFAULT_PROFILES="1g-64m-pc4,1g-128m-pc4,2g-128m-pc4,2g-256m-pc4" +PROFILES="${PROFILES:-$DEFAULT_PROFILES}" + +usage() { + cat <<'USAGE' +Usage: + scripts/run_gt1g_multipart_put_matrix.sh --access-key --secret-key [options] + +Required: + --access-key + --secret-key + +Options: + --host Default: 127.0.0.1:9000 + --bucket-prefix Default: rustfs-gt1g-multipart + --lookup Default: path + --duration Default: 10m + --profiles Default: 1g-64m-pc4,1g-128m-pc4,2g-128m-pc4,2g-256m-pc4 + --warp-bin Default: warp + --out-dir Default: target/bench/gt1g-multipart-put- + --dry-run + -h, --help + +Built-in profiles: + 1g-64m-pc4 concurrent=4 parts=16 part.size=64MiB part.concurrent=4 + 1g-128m-pc4 concurrent=4 parts=8 part.size=128MiB part.concurrent=4 + 2g-128m-pc4 concurrent=4 parts=16 part.size=128MiB part.concurrent=4 + 2g-256m-pc4 concurrent=4 parts=8 part.size=256MiB part.concurrent=4 + 1g-64m-pc8 concurrent=4 parts=16 part.size=64MiB part.concurrent=8 + 2g-128m-pc8 concurrent=4 parts=16 part.size=128MiB part.concurrent=8 + 2g-256m-pc8 concurrent=4 parts=8 part.size=256MiB part.concurrent=8 +USAGE +} + +require_cmd() { + if ! command -v "$1" >/dev/null 2>&1; then + echo "ERROR: command not found: $1" >&2 + exit 1 + fi +} + +trim() { + echo "$1" | awk '{$1=$1;print}' +} + +parse_args() { + while [[ $# -gt 0 ]]; do + case "$1" in + --host) HOST="$2"; shift 2 ;; + --access-key) ACCESS_KEY="$2"; shift 2 ;; + --secret-key) SECRET_KEY="$2"; shift 2 ;; + --bucket-prefix) BUCKET_PREFIX="$2"; shift 2 ;; + --lookup) LOOKUP="$2"; shift 2 ;; + --duration) DURATION="$2"; shift 2 ;; + --profiles) PROFILES="$2"; shift 2 ;; + --warp-bin) WARP_BIN="$2"; shift 2 ;; + --out-dir) OUT_DIR="$2"; shift 2 ;; + --dry-run) DRY_RUN=true; shift ;; + -h|--help) usage; exit 0 ;; + *) + echo "ERROR: unknown arg: $1" >&2 + usage + exit 1 + ;; + esac + done +} + +validate_args() { + if [[ -z "$ACCESS_KEY" || -z "$SECRET_KEY" ]]; then + echo "ERROR: --access-key and --secret-key are required" >&2 + exit 1 + fi +} + +setup_output() { + mkdir -p "$OUT_DIR/logs" "$OUT_DIR/benchdata" + SUMMARY_CSV="$OUT_DIR/summary.csv" + COMMANDS_TXT="$OUT_DIR/commands.txt" + MANIFEST_TXT="$OUT_DIR/run_manifest.txt" + + echo "profile,host,bucket,duration,concurrent,parts,part_size,part_concurrent,throughput_human,reqps,latency_human,log_file,benchdata_file,status" > "$SUMMARY_CSV" + : > "$COMMANDS_TXT" +} + +write_manifest() { + { + echo "created_at=$(date +%Y-%m-%dT%H:%M:%S%z)" + echo "host=${HOST}" + echo "bucket_prefix=${BUCKET_PREFIX}" + echo "lookup=${LOOKUP}" + echo "duration=${DURATION}" + echo "profiles=${PROFILES}" + echo "warp_bin=${WARP_BIN}" + echo "dry_run=${DRY_RUN}" + echo "access_key=REDACTED" + echo "secret_key=REDACTED" + } > "$MANIFEST_TXT" +} + +profile_spec() { + case "$1" in + 1g-64m-pc4) echo "4|16|64MiB|4" ;; + 1g-128m-pc4) echo "4|8|128MiB|4" ;; + 2g-128m-pc4) echo "4|16|128MiB|4" ;; + 2g-256m-pc4) echo "4|8|256MiB|4" ;; + 1g-64m-pc8) echo "4|16|64MiB|8" ;; + 2g-128m-pc8) echo "4|16|128MiB|8" ;; + 2g-256m-pc8) echo "4|8|256MiB|8" ;; + *) + echo "ERROR: unknown profile: $1" >&2 + exit 1 + ;; + esac +} + +extract_first() { + local regex="$1" + local file="$2" + rg -o "$regex" "$file" | head -n1 || true +} + +extract_metrics() { + local log_file="$1" + local throughput reqps latency + + throughput="$(extract_first '[0-9]+(\.[0-9]+)?[[:space:]]*(GiB/s|MiB/s|KiB/s|GB/s|MB/s|KB/s|B/s)' "$log_file")" + reqps="$(extract_first '[0-9]+(\.[0-9]+)?[[:space:]]*(obj/s|req/s|ops/s|requests/s)' "$log_file")" + latency="$(rg -o 'Reqs:[[:space:]]+Avg:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' "$log_file" | head -n1 | sed -E 's/^Reqs:[[:space:]]+Avg:[[:space:]]+//')" + + throughput="$(trim "${throughput:-N/A}")" + reqps="$(trim "${reqps:-N/A}")" + latency="$(trim "${latency:-N/A}")" + reqps="$(echo "$reqps" | awk '{print $1}')" + + echo "${throughput},${reqps:-N/A},${latency}" +} + +run_profile() { + local profile="$1" + local spec concurrent parts part_size part_concurrent + local bucket log_file benchdata_file status metrics throughput reqps latency + + spec="$(profile_spec "$profile")" + IFS='|' read -r concurrent parts part_size part_concurrent <<< "$spec" + + bucket="$(echo "${BUCKET_PREFIX}-${profile}-$(date +%Y%m%d%H%M%S)" | tr '[:upper:]' '[:lower:]' | sed -E 's/[^a-z0-9-]+/-/g; s/-{2,}/-/g; s/^-+//; s/-+$//')" + bucket="${bucket:0:63}" + log_file="$OUT_DIR/logs/${profile}.log" + benchdata_file="$OUT_DIR/benchdata/${profile}.csv.zst" + + local -a cmd=( + "$WARP_BIN" multipart-put + --host "$HOST" + --access-key "$ACCESS_KEY" + --secret-key "$SECRET_KEY" + --bucket "$bucket" + --lookup "$LOOKUP" + --duration "$DURATION" + --concurrent "$concurrent" + --parts "$parts" + --part.size "$part_size" + --part.concurrent "$part_concurrent" + --benchdata "$benchdata_file" + --analyze.v + --no-color + ) + + printf '%q ' "${cmd[@]}" >> "$COMMANDS_TXT" + printf '\n' >> "$COMMANDS_TXT" + + status="ok" + if [[ "$DRY_RUN" == "true" ]]; then + printf '[DRY-RUN] %q ' "${cmd[@]}" + printf '\n' + : > "$log_file" + else + if ! "${cmd[@]}" >"$log_file" 2>&1; then + status="failed" + fi + fi + + metrics="$(extract_metrics "$log_file")" + throughput="$(echo "$metrics" | cut -d',' -f1)" + reqps="$(echo "$metrics" | cut -d',' -f2)" + latency="$(echo "$metrics" | cut -d',' -f3)" + + echo "${profile},${HOST},${bucket},${DURATION},${concurrent},${parts},${part_size},${part_concurrent},${throughput},${reqps},${latency},${log_file},${benchdata_file},${status}" >> "$SUMMARY_CSV" +} + +main() { + parse_args "$@" + validate_args + require_cmd "$WARP_BIN" + require_cmd rg + setup_output + write_manifest + + echo "Output dir: $OUT_DIR" + echo "Profiles: $PROFILES" + + IFS=',' read -r -a profiles <<< "$PROFILES" + for raw in "${profiles[@]}"; do + profile="$(trim "$raw")" + [[ -z "$profile" ]] && continue + run_profile "$profile" + done + + echo + echo "Summary:" + cat "$SUMMARY_CSV" +} + +main "$@" diff --git a/scripts/run_gt1g_multipart_put_server_path_focus.sh b/scripts/run_gt1g_multipart_put_server_path_focus.sh new file mode 100644 index 000000000..d48199ddf --- /dev/null +++ b/scripts/run_gt1g_multipart_put_server_path_focus.sh @@ -0,0 +1,222 @@ +#!/usr/bin/env bash +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +PROJECT_ROOT="$(cd "${SCRIPT_DIR}/.." && pwd)" + +HOST="${HOST:-127.0.0.1:9000}" +ACCESS_KEY="${ACCESS_KEY:-}" +SECRET_KEY="${SECRET_KEY:-}" +BUCKET_PREFIX="${BUCKET_PREFIX:-issue712-gt1g-multipart-focus}" +LOOKUP="${LOOKUP:-path}" +DURATION="${DURATION:-10m}" +WARP_BIN="${WARP_BIN:-warp}" +OUT_DIR="${OUT_DIR:-${PROJECT_ROOT}/target/bench/gt1g-multipart-server-path-$(date +%Y%m%d-%H%M%S)}" +DRY_RUN=false + +DEFAULT_PROFILES="1g-64m-pc4,2g-128m-pc4" +PROFILES="${PROFILES:-$DEFAULT_PROFILES}" + +usage() { + cat <<'USAGE' +Usage: + scripts/run_gt1g_multipart_put_server_path_focus.sh --access-key --secret-key [options] + +Required: + --access-key + --secret-key + +Options: + --host Default: 127.0.0.1:9000 + --bucket-prefix Default: issue712-gt1g-multipart-focus + --lookup Default: path + --duration Default: 10m + --profiles Default: 1g-64m-pc4,2g-128m-pc4 + --warp-bin Default: warp + --out-dir Default: target/bench/gt1g-multipart-server-path- + --dry-run + -h, --help + +Profiles: + 1g-64m-pc4 concurrent=4 parts=16 part.size=64MiB part.concurrent=4 + 2g-128m-pc4 concurrent=4 parts=16 part.size=128MiB part.concurrent=4 + 2g-256m-pc4 concurrent=4 parts=8 part.size=256MiB part.concurrent=4 +USAGE +} + +require_cmd() { + if ! command -v "$1" >/dev/null 2>&1; then + echo "ERROR: command not found: $1" >&2 + exit 1 + fi +} + +trim() { + echo "$1" | awk '{$1=$1;print}' +} + +parse_args() { + while [[ $# -gt 0 ]]; do + case "$1" in + --host) HOST="$2"; shift 2 ;; + --access-key) ACCESS_KEY="$2"; shift 2 ;; + --secret-key) SECRET_KEY="$2"; shift 2 ;; + --bucket-prefix) BUCKET_PREFIX="$2"; shift 2 ;; + --lookup) LOOKUP="$2"; shift 2 ;; + --duration) DURATION="$2"; shift 2 ;; + --profiles) PROFILES="$2"; shift 2 ;; + --warp-bin) WARP_BIN="$2"; shift 2 ;; + --out-dir) OUT_DIR="$2"; shift 2 ;; + --dry-run) DRY_RUN=true; shift ;; + -h|--help) usage; exit 0 ;; + *) + echo "ERROR: unknown arg: $1" >&2 + usage + exit 1 + ;; + esac + done +} + +validate_args() { + if [[ -z "$ACCESS_KEY" || -z "$SECRET_KEY" ]]; then + echo "ERROR: --access-key and --secret-key are required" >&2 + exit 1 + fi +} + +setup_output() { + mkdir -p "$OUT_DIR/logs" "$OUT_DIR/benchdata" + SUMMARY_CSV="$OUT_DIR/summary.csv" + COMMANDS_TXT="$OUT_DIR/commands.txt" + MANIFEST_TXT="$OUT_DIR/run_manifest.txt" + + echo "profile,host,bucket,duration,concurrent,parts,part_size,part_concurrent,throughput_human,reqps,latency_human,log_file,benchdata_file,status" > "$SUMMARY_CSV" + : > "$COMMANDS_TXT" +} + +write_manifest() { + { + echo "created_at=$(date +%Y-%m-%dT%H:%M:%S%z)" + echo "host=${HOST}" + echo "bucket_prefix=${BUCKET_PREFIX}" + echo "lookup=${LOOKUP}" + echo "duration=${DURATION}" + echo "profiles=${PROFILES}" + echo "warp_bin=${WARP_BIN}" + echo "dry_run=${DRY_RUN}" + echo "access_key=REDACTED" + echo "secret_key=REDACTED" + } > "$MANIFEST_TXT" +} + +profile_spec() { + case "$1" in + 1g-64m-pc4) echo "4|16|64MiB|4" ;; + 2g-128m-pc4) echo "4|16|128MiB|4" ;; + 2g-256m-pc4) echo "4|8|256MiB|4" ;; + *) + echo "ERROR: unknown profile: $1" >&2 + exit 1 + ;; + esac +} + +extract_first() { + local regex="$1" + local file="$2" + rg -o "$regex" "$file" | head -n1 || true +} + +extract_metrics() { + local log_file="$1" + local throughput reqps latency + + throughput="$(extract_first '[0-9]+(\.[0-9]+)?[[:space:]]*(GiB/s|MiB/s|KiB/s|GB/s|MB/s|KB/s|B/s)' "$log_file")" + reqps="$(extract_first '[0-9]+(\.[0-9]+)?[[:space:]]*(obj/s|req/s|ops/s|requests/s)' "$log_file")" + latency="$(rg -o 'Reqs:[[:space:]]+Avg:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' "$log_file" | head -n1 | sed -E 's/^Reqs:[[:space:]]+Avg:[[:space:]]+//')" + + throughput="$(trim "${throughput:-N/A}")" + reqps="$(trim "${reqps:-N/A}")" + latency="$(trim "${latency:-N/A}")" + reqps="$(echo "$reqps" | awk '{print $1}')" + + echo "${throughput},${reqps:-N/A},${latency}" +} + +run_profile() { + local profile="$1" + local spec concurrent parts part_size part_concurrent + local bucket log_file benchdata_file status metrics throughput reqps latency + + spec="$(profile_spec "$profile")" + IFS='|' read -r concurrent parts part_size part_concurrent <<< "$spec" + + bucket="$(echo "${BUCKET_PREFIX}-${profile}-$(date +%Y%m%d%H%M%S)" | tr '[:upper:]' '[:lower:]' | sed -E 's/[^a-z0-9-]+/-/g; s/-{2,}/-/g; s/^-+//; s/-+$//')" + bucket="${bucket:0:63}" + log_file="$OUT_DIR/logs/${profile}.log" + benchdata_file="$OUT_DIR/benchdata/${profile}.csv.zst" + + local -a cmd=( + "$WARP_BIN" multipart-put + --host "$HOST" + --access-key "$ACCESS_KEY" + --secret-key "$SECRET_KEY" + --bucket "$bucket" + --lookup "$LOOKUP" + --duration "$DURATION" + --concurrent "$concurrent" + --parts "$parts" + --part.size "$part_size" + --part.concurrent "$part_concurrent" + --benchdata "$benchdata_file" + --analyze.v + --no-color + ) + + printf '%q ' "${cmd[@]}" >> "$COMMANDS_TXT" + printf '\n' >> "$COMMANDS_TXT" + + status="ok" + if [[ "$DRY_RUN" == "true" ]]; then + printf '[DRY-RUN] %q ' "${cmd[@]}" + printf '\n' + : > "$log_file" + else + if ! "${cmd[@]}" >"$log_file" 2>&1; then + status="failed" + fi + fi + + metrics="$(extract_metrics "$log_file")" + throughput="$(echo "$metrics" | cut -d',' -f1)" + reqps="$(echo "$metrics" | cut -d',' -f2)" + latency="$(echo "$metrics" | cut -d',' -f3)" + + echo "${profile},${HOST},${bucket},${DURATION},${concurrent},${parts},${part_size},${part_concurrent},${throughput},${reqps},${latency},${log_file},${benchdata_file},${status}" >> "$SUMMARY_CSV" +} + +main() { + parse_args "$@" + validate_args + require_cmd "$WARP_BIN" + require_cmd rg + setup_output + write_manifest + + echo "Output dir: $OUT_DIR" + echo "Profiles: $PROFILES" + + IFS=',' read -r -a profiles <<< "$PROFILES" + for raw in "${profiles[@]}"; do + profile="$(trim "$raw")" + [[ -z "$profile" ]] && continue + run_profile "$profile" + done + + echo + echo "Summary:" + cat "$SUMMARY_CSV" +} + +main "$@"