From 3a0dbccc2e7baab1ce9bde5fdb37ccdf65f295d5 Mon Sep 17 00:00:00 2001 From: houseme Date: Thu, 13 Aug 2026 10:04:20 +0800 Subject: [PATCH] perf(ecstore): reduce inline PUT commit overhead (#6033) * perf(metrics): attribute PUT stage costs Co-Authored-By: heihutu * perf(ecstore): move PUT metadata during shuffle Co-Authored-By: heihutu * perf(s3): reuse PUT object lock state Co-Authored-By: heihutu * perf(ecstore): trim PUT metadata fanout clones Build per-disk PUT metadata only for committed writer slots, move the response metadata out of the fanout vector, and preserve fresh FileInfo shuffle semantics. Co-Authored-By: heihutu * perf(metrics): make PUT stage attribution opt-in Co-Authored-By: heihutu * perf(ecstore): commit inline PUT shards directly Co-Authored-By: heihutu * perf(ecstore): streamline rename staging cleanup Use the directory-specific removal operation for rename_data staging parents. This avoids a guaranteed failed file-removal probe on Unix-like hosts and lets Windows remove the empty directory directly while preserving best-effort non-empty handling. Co-Authored-By: heihutu * test(ecstore): cover inline PUT rename failures Cache the detailed stage metrics gate once per PUT and exercise exact-quorum and quorum-minus-one failures after inline shard encoding. Co-Authored-By: heihutu --------- Co-authored-by: heihutu --- crates/config/src/constants/app.rs | 5 + crates/config/src/observability/mod.rs | 5 + crates/ecstore/src/disk/local.rs | 49 +- crates/ecstore/src/erasure/coding/encode.rs | 63 +++ crates/ecstore/src/set_disk/metadata.rs | 52 ++ crates/ecstore/src/set_disk/ops/object.rs | 527 ++++++++++++++++---- crates/io-metrics/src/lib.rs | 73 ++- rustfs/src/app/object_usecase.rs | 62 ++- rustfs/src/startup_observability.rs | 57 ++- 9 files changed, 785 insertions(+), 108 deletions(-) diff --git a/crates/config/src/constants/app.rs b/crates/config/src/constants/app.rs index 6dd3cb6a5..36ea95d5a 100644 --- a/crates/config/src/constants/app.rs +++ b/crates/config/src/constants/app.rs @@ -353,6 +353,11 @@ pub const DEFAULT_OBS_TRACES_EXPORT_ENABLED: bool = true; /// Environment variable: RUSTFS_OBS_METRICS_EXPORT_ENABLED pub const DEFAULT_OBS_METRICS_EXPORT_ENABLED: bool = true; +/// Default detailed PUT stage metrics enabled +/// Default value: false +/// Environment variable: RUSTFS_OBS_PUT_STAGE_METRICS_ENABLED +pub const DEFAULT_OBS_PUT_STAGE_METRICS_ENABLED: bool = false; + /// Default logs export enabled /// It is used to enable or disable exporting logs /// Default value: true diff --git a/crates/config/src/observability/mod.rs b/crates/config/src/observability/mod.rs index 9dd2ba654..b23e5f2c6 100644 --- a/crates/config/src/observability/mod.rs +++ b/crates/config/src/observability/mod.rs @@ -44,6 +44,10 @@ pub const ENV_OBS_METRICS_EXPORT_ENABLED: &str = "RUSTFS_OBS_METRICS_EXPORT_ENAB pub const ENV_OBS_LOGS_EXPORT_ENABLED: &str = "RUSTFS_OBS_LOGS_EXPORT_ENABLED"; pub const ENV_OBS_PROFILING_EXPORT_ENABLED: &str = "RUSTFS_OBS_PROFILING_EXPORT_ENABLED"; +/// Enables detailed per-stage PUT metrics. Disabled by default because each +/// PUT records multiple timers and histograms when attribution is active. +pub const ENV_OBS_PUT_STAGE_METRICS_ENABLED: &str = "RUSTFS_OBS_PUT_STAGE_METRICS_ENABLED"; + pub const ENV_OBS_LOGGER_LEVEL: &str = "RUSTFS_OBS_LOGGER_LEVEL"; pub const ENV_OBS_LOG_STDOUT_ENABLED: &str = "RUSTFS_OBS_LOG_STDOUT_ENABLED"; pub const ENV_OBS_LOG_DIRECTORY: &str = "RUSTFS_OBS_LOG_DIRECTORY"; @@ -141,6 +145,7 @@ mod tests { assert_eq!(ENV_OBS_METRICS_EXPORT_ENABLED, "RUSTFS_OBS_METRICS_EXPORT_ENABLED"); assert_eq!(ENV_OBS_LOGS_EXPORT_ENABLED, "RUSTFS_OBS_LOGS_EXPORT_ENABLED"); assert_eq!(ENV_OBS_PROFILING_EXPORT_ENABLED, "RUSTFS_OBS_PROFILING_EXPORT_ENABLED"); + assert_eq!(ENV_OBS_PUT_STAGE_METRICS_ENABLED, "RUSTFS_OBS_PUT_STAGE_METRICS_ENABLED"); // Test log cleanup related env keys assert_eq!(ENV_OBS_LOG_MAX_TOTAL_SIZE_BYTES, "RUSTFS_OBS_LOG_MAX_TOTAL_SIZE_BYTES"); assert_eq!(ENV_OBS_LOG_MAX_SINGLE_FILE_SIZE_BYTES, "RUSTFS_OBS_LOG_MAX_SINGLE_FILE_SIZE_BYTES"); diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 5a42245d8..3778d837e 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -9174,7 +9174,7 @@ impl DiskAPI for LocalDisk { if let Some(src_file_path_parent) = src_file_path.parent() { if src_volume != super::RUSTFS_META_MULTIPART_BUCKET { - let _ = remove_std(src_file_path_parent); + let _ = std::fs::remove_dir(src_file_path_parent); } else { let _ = self .delete_file(&dst_volume_dir, &src_file_path_parent.to_path_buf(), true, false) @@ -9499,7 +9499,7 @@ impl DiskAPI for LocalDisk { if let Some(ref cleanup) = cleanup_path { let _ = self.delete_file(&dst_volume_dir, cleanup, true, false).await; } else if let Some(parent) = src_file_path.parent() { - let _ = remove_std(parent); + let _ = std::fs::remove_dir(parent); } // Heal reuses a version's `data_dir` and lands the rebuilt shard on @@ -12439,6 +12439,10 @@ mod test { .join(RUSTFS_META_TMP_BUCKET) .join(tmp_object) .join(new_data_dir.to_string()); + let tmp_parent = tmp_data_dir + .parent() + .expect("tmp data dir should have a parent") + .to_path_buf(); fs::create_dir_all(&tmp_data_dir) .await .expect("new tmp data dir should be created"); @@ -12450,6 +12454,10 @@ mod test { disk.rename_data(RUSTFS_META_TMP_BUCKET, tmp_object, new_fi, bucket, object) .await .expect("rename_data should commit"); + assert!( + !tmp_parent.exists(), + "successful non-inline commit should remove the empty staging parent" + ); // The tmp xl.meta write point uses SyncMode::FileOnly: its parent dir // ({tmp}/{tmp_object}) must not be fsynced. @@ -12654,6 +12662,9 @@ mod test { let tmp_object = "tmp-new-inline"; ensure_test_volume(&disk, bucket).await; ensure_test_volume(&disk, RUSTFS_META_TMP_BUCKET).await; + let tmp_parent = disk + .get_object_path(RUSTFS_META_TMP_BUCKET, tmp_object) + .expect("tmp parent should resolve"); let _mode = durability_mode_override::set(DurabilityMode::Strict); let version_id = Uuid::parse_str("99999999-9999-9999-9999-999999999999").expect("version id should parse"); @@ -12662,6 +12673,7 @@ mod test { disk.rename_data(RUSTFS_META_TMP_BUCKET, tmp_object, new_fi, bucket, object) .await .expect("inline rename_data should commit the new object"); + assert!(!tmp_parent.exists(), "successful inline commit should remove the empty staging parent"); let bucket_dir = disk.get_bucket_path(bucket).expect("bucket path should resolve"); let prefix_dir = disk.get_object_path(bucket, "prefix").expect("prefix path should resolve"); @@ -12685,6 +12697,34 @@ mod test { ); } + #[tokio::test] + async fn rename_data_inline_preserves_non_empty_staging_parent() { + use tempfile::tempdir; + + let dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + let bucket = "inline-staging-sentinel-bucket"; + let object = "inline-object"; + let tmp_object = "inline-stage-with-sentinel"; + ensure_test_volume(&disk, bucket).await; + ensure_test_volume(&disk, RUSTFS_META_TMP_BUCKET).await; + + let tmp_parent = disk + .get_object_path(RUSTFS_META_TMP_BUCKET, tmp_object) + .expect("tmp parent should resolve"); + fs::create_dir_all(&tmp_parent).await.expect("tmp parent should be created"); + let sentinel = tmp_parent.join("sentinel"); + fs::write(&sentinel, b"keep").await.expect("sentinel should be written"); + + let fi = test_file_info(object, Uuid::new_v4(), None, Some(Bytes::from_static(b"inline-payload"))); + disk.rename_data(RUSTFS_META_TMP_BUCKET, tmp_object, fi, bucket, object) + .await + .expect("non-empty staging cleanup must not negate the committed object"); + + assert_eq!(fs::read(&sentinel).await.expect("sentinel should remain"), b"keep"); + } + #[cfg(unix)] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[allow(clippy::await_holding_lock)] @@ -12859,7 +12899,10 @@ mod test { .expect("non-inline rename_data should commit"); assert!(!replacement_dir.exists(), "the destination object directory must not be replaced"); - assert!(staging_parent.exists(), "the guarded staging parent must retain its identity"); + assert!( + !staging_parent.exists(), + "successful commit should remove the empty staging parent after releasing its guard" + ); assert!( !replacement_staging_parent.exists(), "the staging parent must not be replaced between data and metadata publication" diff --git a/crates/ecstore/src/erasure/coding/encode.rs b/crates/ecstore/src/erasure/coding/encode.rs index e600603d8..39b792485 100644 --- a/crates/ecstore/src/erasure/coding/encode.rs +++ b/crates/ecstore/src/erasure/coding/encode.rs @@ -23,6 +23,7 @@ use crate::runtime::sources as runtime_sources; use bytes::{Bytes, BytesMut}; use futures::StreamExt; use futures::stream::FuturesUnordered; +use rustfs_utils::HashAlgorithm; use std::sync::Arc; use std::time::Instant; use std::vec; @@ -592,6 +593,39 @@ impl Erasure { Ok((reader, total)) } + /// Encode a small inline object directly into its per-disk bitrot payloads. + /// The returned bytes are the same `[hash][shard]` representation produced + /// by `BitrotWriter`, ready to be embedded in each disk's staged `xl.meta`. + #[hotpath::measure(impl_type = "Erasure")] + pub(crate) async fn encode_inline_shards_with_size_hint( + self: Arc, + mut reader: R, + size_hint: usize, + ) -> std::io::Result<(R, usize, Vec)> + where + R: AsyncRead + Send + Sync + Unpin, + { + use tokio::io::AsyncReadExt; + + let mut buf = Vec::with_capacity(small_ingest_capacity(&self, size_hint)); + let total = reader.read_to_end(&mut buf).await?; + if total == 0 { + return Ok((reader, 0, Vec::new())); + } + + let shards = self.encode_data_owned(buf)?; + let mut inline_shards = Vec::with_capacity(shards.len()); + for shard in shards { + let hash = HashAlgorithm::HighwayHash256S.hash_encode(&shard); + let mut encoded = BytesMut::with_capacity(hash.as_ref().len() + shard.len()); + encoded.extend_from_slice(hash.as_ref()); + encoded.extend_from_slice(&shard); + inline_shards.push(encoded.freeze()); + } + + Ok((reader, total, inline_shards)) + } + #[hotpath::measure(impl_type = "Erasure")] pub async fn encode( self: Arc, @@ -2356,6 +2390,35 @@ mod tests { assert!(committed.lock().unwrap().is_empty()); } + #[tokio::test] + async fn encode_inline_shards_matches_writer_bitrot_layout() { + const DATA_SHARDS: usize = 2; + const PARITY_SHARDS: usize = 2; + const BLOCK_SIZE: usize = 64; + let payload = b"inline commit payload".to_vec(); + let checksum_algo = HashAlgorithm::HighwayHash256S; + let erasure = Arc::new(Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE)); + let reader = tokio::io::BufReader::new(Cursor::new(payload.clone())); + + let (_reader, total, inline_shards) = erasure + .clone() + .encode_inline_shards_with_size_hint(reader, payload.len()) + .await + .expect("inline shards should encode"); + let raw_shards = erasure + .encode_data_owned(payload.clone()) + .expect("reference shards should encode"); + + assert_eq!(total, payload.len()); + assert_eq!(inline_shards.len(), DATA_SHARDS + PARITY_SHARDS); + for (inline, raw) in inline_shards.iter().zip(raw_shards) { + let mut writer = BitrotWriterWrapper::new(CustomWriter::new_inline_buffer(), raw.len(), checksum_algo.clone()); + writer.write(&raw).await.expect("reference writer should accept shard"); + writer.shutdown().await.expect("reference writer should shutdown"); + assert_eq!(inline.as_ref(), writer.into_inline_data().expect("reference writer should retain bytes")); + } + } + /// encode_inline_small: small payload is encoded into the correct number of shards /// and each writer receives data after shutdown. #[tokio::test] diff --git a/crates/ecstore/src/set_disk/metadata.rs b/crates/ecstore/src/set_disk/metadata.rs index a4c4dfa66..b07ed428c 100644 --- a/crates/ecstore/src/set_disk/metadata.rs +++ b/crates/ecstore/src/set_disk/metadata.rs @@ -1079,6 +1079,25 @@ impl SetDisks { shuffled_disks } + pub(super) fn shuffle_disks_owned(mut disks: Vec>, distribution: &[usize]) -> Vec> { + if distribution.is_empty() { + return disks; + } + + let mut shuffled_disks = vec![None; disks.len()]; + for (index, disk) in disks.iter_mut().enumerate() { + let Some(slot) = distribution + .get(index) + .and_then(|block_index| block_index.checked_sub(1)) + .filter(|slot| *slot < shuffled_disks.len()) + else { + continue; + }; + shuffled_disks[slot] = disk.take(); + } + shuffled_disks + } + pub(super) fn shuffle_check_parts(parts_errs: &[usize], distribution: &[usize]) -> Vec { if distribution.is_empty() { return parts_errs.to_vec(); @@ -1390,6 +1409,23 @@ mod tests { assert_eq!(owned_slots, expected_slots, "fallback disk slots must match the borrowing variant"); } + #[tokio::test] + async fn owned_shuffle_preserves_fresh_put_metadata() { + let tempdir = tempfile::tempdir().expect("tempdir should be created"); + let fi = FileInfo::new("bucket/object", 2, 1); + let parts = vec![fi.clone(); fi.erasure.distribution.len()]; + let disks = shuffle_test_disks(&tempdir, parts.len()).await; + + let (owned_disks, owned_parts) = SetDisks::shuffle_disks_and_parts_metadata_by_index_owned(disks, parts, &fi); + + assert!(owned_disks.iter().all(Option::is_some), "fresh PUT must retain every online disk"); + assert_eq!( + owned_parts, + vec![fi; owned_disks.len()], + "fresh PUT metadata with pending shard indexes must survive init fallback" + ); + } + // backlog#949: corrupt/adversarial distribution values (0 or > N) must not // trigger a `usize` underflow / out-of-bounds panic in the shuffle helpers. #[test] @@ -1419,6 +1455,22 @@ mod tests { assert_eq!(result.len(), disks.len(), "output length must be preserved"); } + #[tokio::test] + async fn owned_disk_shuffle_matches_borrowing_variant() { + let tempdir = tempfile::tempdir().expect("tempdir should be created"); + let mut disks = shuffle_test_disks(&tempdir, 4).await; + disks[1] = None; + disks[3] = None; + let distribution = [3, 1, 4, 2]; + + let expected = SetDisks::shuffle_disks(&disks, &distribution); + let actual = SetDisks::shuffle_disks_owned(disks, &distribution); + + let expected_slots = expected.iter().map(Option::is_some).collect::>(); + let actual_slots = actual.iter().map(Option::is_some).collect::>(); + assert_eq!(actual_slots, expected_slots, "owned shuffle must preserve disk placement"); + } + #[tokio::test] async fn shuffle_disks_and_parts_metadata_survives_corrupt_distribution() { let tempdir = tempfile::tempdir().expect("tempdir should be created"); diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 20d271cd9..df11384e3 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -61,6 +61,10 @@ fn duration_millis_f64(duration: std::time::Duration) -> f64 { duration.as_secs_f64() * 1000.0 } +fn committed_response_metadata_slot(committed_disks: &[Option], fallback_slot: usize) -> usize { + committed_disks.iter().position(Option::is_some).unwrap_or(fallback_slot) +} + #[cfg(test)] mod duration_metrics_tests { use super::duration_millis_f64; @@ -72,6 +76,38 @@ mod duration_metrics_tests { } } +#[cfg(test)] +mod put_metadata_tests { + use super::*; + + #[test] + fn committed_file_info_follows_exact_quorum_success_slot() { + let mut first_success = FileInfo::new("bucket/object", 2, 2); + first_success.name = "first-success".to_string(); + let mut second_success = first_success.clone(); + second_success.name = "second-success".to_string(); + let mut parts_metadata = [FileInfo::default(), first_success, second_success, FileInfo::default()]; + let committed_disks = [None, Some(()), Some(()), None]; + + assert_eq!( + committed_disks.iter().filter(|disk| disk.is_some()).count(), + 2, + "fixture must meet exact quorum" + ); + let selected_slot = committed_response_metadata_slot(&committed_disks, 3); + let selected = std::mem::take(&mut parts_metadata[selected_slot]); + + assert_eq!(selected.name, "first-success"); + assert_eq!(parts_metadata[1], FileInfo::default(), "selected metadata should move without cloning"); + assert_eq!(parts_metadata[2].name, "second-success", "other committed metadata must remain available"); + assert_eq!( + committed_response_metadata_slot::<()>(&[None, None, None, None], 3), + 3, + "a violated post-commit success-mask invariant must not turn a durable PUT into an error" + ); + } +} + fn is_restore_control_metadata(key: &str) -> bool { key.eq_ignore_ascii_case(X_AMZ_RESTORE.as_str()) || key.eq_ignore_ascii_case(rustfs_utils::http::headers::AMZ_RESTORE_EXPIRY_DAYS) @@ -1059,10 +1095,7 @@ impl SetDisks { } fi.data_dir = Some(Uuid::new_v4()); - - let parts_metadata = vec![fi.clone(); disks.len()]; - - let (mut shuffle_disks, mut parts_metadatas) = Self::shuffle_disks_and_parts_metadata(&disks, &parts_metadata, &fi); + let mut shuffle_disks = Self::shuffle_disks_owned(disks, &fi.erasure.distribution); let tmp_dir = Uuid::new_v4().to_string(); @@ -1074,61 +1107,91 @@ impl SetDisks { let put_object_size = known_put_object_storage_size(data.size()); let is_inline_buffer = storage_class_config.should_inline(erasure.shard_file_size(put_object_size), opts.versioned); + let collect_stage_timing = rustfs_io_metrics::put_stage_metrics_enabled() || issue3031_diag_enabled(); let shard_file_size = erasure.shard_file_size(put_object_size); let shard_size = erasure.shard_size(); - let writer_setup_stage_start = Instant::now(); - let writer_futs: Vec<_> = shuffle_disks - .iter() - .map(|disk_op| { - let tmp_obj = tmp_object.clone(); - async move { - if let Some(disk) = disk_op - && disk.is_online().await - { - match create_bitrot_writer( - is_inline_buffer, - Some(disk), - RUSTFS_META_TMP_BUCKET, - &tmp_obj, - shard_file_size, - shard_size, - HashAlgorithm::HighwayHash256S, - ) - .await - { - Ok(writer) => (Some(writer), None), - Err(err) => { - warn!( - event = EVENT_SET_DISK_WRITE, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_SET_DISK, - disk = ?disk, - state = "bitrot_writer_skipped", - error = ?err, - "Set disk bitrot writer skipped" - ); - (None, Some(err)) - } - } - } else { - (None, Some(DiskError::DiskNotFound)) - } + let write_path = classify_put_write_path(is_inline_buffer, put_object_size, fi.erasure.block_size); + let direct_inline_commit = matches!(write_path, SmallWritePath::Inline); + rustfs_io_metrics::record_put_object_path(write_path.metric_label()); + let writer_setup_stage_start = collect_stage_timing.then(Instant::now); + let (mut writers, errors) = if direct_inline_commit { + let online = join_all(shuffle_disks.iter().map(|disk| async move { + if let Some(disk) = disk { + disk.is_online().await + } else { + false } - }) - .collect(); - let writer_results = join_all(writer_futs).await; - let mut writers = Vec::with_capacity(writer_results.len()); - let mut errors = Vec::with_capacity(writer_results.len()); - for (w, e) in writer_results { - writers.push(w); - errors.push(e); + })) + .await; + let mut errors = Vec::with_capacity(online.len()); + for (disk, is_online) in shuffle_disks.iter_mut().zip(online) { + if is_online { + errors.push(None); + } else { + *disk = None; + errors.push(Some(DiskError::DiskNotFound)); + } + } + (std::iter::repeat_with(|| None).take(shuffle_disks.len()).collect(), errors) + } else { + let writer_futs: Vec<_> = shuffle_disks + .iter() + .map(|disk_op| { + let tmp_obj = tmp_object.clone(); + async move { + if let Some(disk) = disk_op + && disk.is_online().await + { + match create_bitrot_writer( + is_inline_buffer, + Some(disk), + RUSTFS_META_TMP_BUCKET, + &tmp_obj, + shard_file_size, + shard_size, + HashAlgorithm::HighwayHash256S, + ) + .await + { + Ok(writer) => (Some(writer), None), + Err(err) => { + warn!( + event = EVENT_SET_DISK_WRITE, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + disk = ?disk, + state = "bitrot_writer_skipped", + error = ?err, + "Set disk bitrot writer skipped" + ); + (None, Some(err)) + } + } + } else { + (None, Some(DiskError::DiskNotFound)) + } + } + }) + .collect(); + let writer_results = join_all(writer_futs).await; + let mut writers = Vec::with_capacity(writer_results.len()); + let mut errors = Vec::with_capacity(writer_results.len()); + for (writer, error) in writer_results { + writers.push(writer); + errors.push(error); + } + (writers, errors) + }; + let writer_setup_elapsed = writer_setup_stage_start.map(|stage_start| stage_start.elapsed()); + let writer_setup_ms = writer_setup_elapsed + .map(|elapsed| elapsed.as_millis() as u64) + .unwrap_or_default(); + if let Some(writer_setup_elapsed) = writer_setup_elapsed { + rustfs_io_metrics::record_put_object_stage_duration( + "set_disk_writer_setup", + duration_millis_f64(writer_setup_elapsed), + ); } - let writer_setup_elapsed = writer_setup_stage_start.elapsed(); - let writer_setup_ms = writer_setup_elapsed.as_millis() as u64; - rustfs_io_metrics::record_put_object_stage_duration( - "set_disk_writer_setup", - duration_millis_f64(writer_setup_elapsed), - ); let nil_count = errors.iter().filter(|&e| e.is_none()).count(); if nil_count < write_quorum { @@ -1156,21 +1219,23 @@ impl SetDisks { HashReader::from_stream(Cursor::new(Vec::new()), 0, 0, None, None, false)?, ); - let write_path = classify_put_write_path(is_inline_buffer, put_object_size, fi.erasure.block_size); - rustfs_io_metrics::record_put_object_path(write_path.metric_label()); let small_size_hint = if matches!(write_path, SmallWritePath::Inline | SmallWritePath::SingleBlockNonInline) { usize::try_from(put_object_size).map_err(Error::other)? } else { 0 }; - let encode_stage_start = Instant::now(); + let encode_stage_start = collect_stage_timing.then(Instant::now); + let mut inline_shards = None; let (reader, w_size) = match write_path { SmallWritePath::Inline => match Arc::clone(&erasure) - .encode_inline_small_with_size_hint(stream, &mut writers, write_quorum, small_size_hint) + .encode_inline_shards_with_size_hint(stream, small_size_hint) .await { - Ok((r, w)) => (r, w), + Ok((r, w, shards)) => { + inline_shards = Some(shards); + (r, w) + } Err(e) => { error!("encode_inline_small err {:?}", e); return Err(e.into()); @@ -1203,9 +1268,11 @@ impl SetDisks { } }, }; - let encode_elapsed = encode_stage_start.elapsed(); - let encode_ms = encode_elapsed.as_millis() as u64; - rustfs_io_metrics::record_put_object_stage_duration("set_disk_encode", duration_millis_f64(encode_elapsed)); + let encode_elapsed = encode_stage_start.map(|stage_start| stage_start.elapsed()); + let encode_ms = encode_elapsed.map(|elapsed| elapsed.as_millis() as u64).unwrap_or_default(); + if let Some(encode_elapsed) = encode_elapsed { + rustfs_io_metrics::record_put_object_stage_duration("set_disk_encode", duration_millis_f64(encode_elapsed)); + } let _ = mem::replace(&mut data.stream, reader); // if let Err(err) = close_bitrot_writers(&mut writers).await { @@ -1291,37 +1358,67 @@ impl SetDisks { // drop it below reconstructable quorum (backlog#852 / #799 B3). // `rename_data` re-checks write quorum over the surviving disks and // rolls back if too few remain. - let committed_shards = drop_failed_writer_disks(&mut shuffle_disks, &writers); + let committed_shards = if matches!(write_path, SmallWritePath::Inline) { + shuffle_disks.iter().filter(|disk| disk.is_some()).count() + } else { + drop_failed_writer_disks(&mut shuffle_disks, &writers) + }; if committed_shards < write_quorum { return Err(Error::other(format!( "put_object write quorum unavailable after encode: {committed_shards} shard(s) committed, need {write_quorum}" ))); } - for (i, pfi) in parts_metadatas.iter_mut().enumerate() { - pfi.metadata = user_defined.clone(); + fi.metadata = user_defined; + fi.mod_time = mod_time; + fi.size = w_size as i64; + fi.versioned = opts.versioned || opts.version_suspended; + fi.add_object_part(1, etag, w_size, mod_time, actual_size, index_op, None); + if opts.data_movement { + fi.set_data_moved(); + } + let parity_blocks = fi.erasure.parity_blocks; + + let response_metadata_slot = shuffle_disks + .iter() + .rposition(Option::is_some) + .ok_or_else(|| Error::other("put_object write quorum unavailable after encode"))?; + let mut base_file_info = fi; + let mut parts_metadatas = Vec::with_capacity(shuffle_disks.len()); + for (i, disk) in shuffle_disks.iter().enumerate() { + if disk.is_none() { + parts_metadatas.push(FileInfo::default()); + continue; + } + + let mut pfi = if i == response_metadata_slot { + std::mem::take(&mut base_file_info) + } else { + base_file_info.clone() + }; if is_inline_buffer { - if let Some(writer) = writers[i].take() { + if let Some(shards) = inline_shards.as_ref() { + pfi.data = Some( + shards + .get(i) + .cloned() + .ok_or_else(|| Error::other(format!("inline encoder omitted disk shard {i}")))?, + ); + } else if let Some(writer) = writers[i].take() { pfi.data = Some(writer.into_inline_data().map(Bytes::from).unwrap_or_default()); } pfi.set_inline_data(); } - - pfi.mod_time = mod_time; - pfi.size = w_size as i64; - pfi.versioned = opts.versioned || opts.version_suspended; - pfi.add_object_part(1, etag.clone(), w_size, mod_time, actual_size, index_op.clone(), None); - pfi.checksum = fi.checksum.clone(); - - if opts.data_movement { - pfi.set_data_moved(); - } + parts_metadatas.push(pfi); } + let committed_version_id = parts_metadatas[response_metadata_slot].version_id; + let committed_data_dir = parts_metadatas[response_metadata_slot].data_dir; + let is_compressed = parts_metadatas[response_metadata_slot].is_compressed(); drop(writers); // drop writers to close all files, this is to prevent FileAccessDenied errors when renaming data - if fi.erasure.parity_blocks == 0 { + if parity_blocks == 0 { let written_size = i64::try_from(w_size).map_err(|_| Error::other("put_object written size overflows i64"))?; let logical_shard_size = usize::try_from(erasure.shard_file_size(written_size)) .map_err(|_| Error::other("put_object shard size overflows usize"))?; @@ -1518,7 +1615,7 @@ impl SetDisks { Some(self.pool_index), Some(self.set_index), ); - request.object_version_id = fi.version_id.map(|version_id| version_id.to_string()); + request.object_version_id = committed_version_id.map(|version_id| version_id.to_string()); tokio::spawn(async move { let _ = rustfs_common::heal_channel::send_heal_request(request).await; }); @@ -1553,7 +1650,7 @@ impl SetDisks { let mut cleanup_stage_ms: Option = None; if let Some(old_dir) = op_old_dir { - let committed_dir = fi.data_dir.unwrap_or_default().to_string(); + let committed_dir = committed_data_dir.unwrap_or_default().to_string(); let cleanup_stage_start = Instant::now(); // backlog#898: reclaiming the dereferenced old data dir is // best-effort and returns a receipt (never `Err`). A failed GC @@ -1590,16 +1687,10 @@ impl SetDisks { } } - for (i, op_disk) in online_disks.iter().enumerate() { - if let Some(disk) = op_disk - && disk.is_online().await - { - fi = parts_metadatas[i].clone(); - break; - } - } + let committed_metadata_slot = committed_response_metadata_slot(&online_disks, response_metadata_slot); + let mut fi = std::mem::take(&mut parts_metadatas[committed_metadata_slot]); - if fi.is_compressed() { + if is_compressed { record_compression_total_memory(actual_size as u64, w_size as u64).await; } self.record_capacity_scope_if_needed(opts.capacity_scope_token, &online_disks); @@ -5568,6 +5659,264 @@ pub(in crate::set_disk::ops) mod hermetic_set_disks_support { } } +#[cfg(test)] +mod inline_put_commit_path_tests { + use super::hermetic_set_disks_support::hermetic_set_disks_isolated as hermetic_set_disks; + use super::*; + use crate::disk::{DiskAPI as _, ReadOptions}; + use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _}; + use tokio::io::AsyncReadExt; + + async fn make_bucket(disks: &[DiskStore], bucket: &str) { + for disk in disks { + disk.make_volume(bucket).await.expect("bucket volume should be created"); + } + } + + #[tokio::test] + async fn inline_put_direct_commit_round_trips_verified_bitrot_shards() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "inline-direct-commit"; + let object = "object.bin"; + let payload: Vec = (0..16 * 1024).map(|index| (index % 251) as u8).collect(); + make_bucket(&disk_stores, bucket).await; + + let mut reader = PutObjReader::from_vec(payload.clone()); + set_disks + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + .expect("inline PUT should commit"); + + let read_data = ReadOptions { + read_data: true, + ..Default::default() + }; + for (disk_index, disk) in disk_stores.iter().enumerate() { + let file_info = disk + .read_version("", bucket, object, "", &read_data) + .await + .unwrap_or_else(|err| panic!("disk {disk_index} should persist inline metadata: {err}")); + assert!(file_info.inline_data(), "disk {disk_index} should mark the shard inline"); + let inline_data = file_info + .data + .as_ref() + .unwrap_or_else(|| panic!("disk {disk_index} should persist inline bitrot bytes")); + let erasure = erasure_from_file_info(&file_info, false).expect("persisted erasure layout should be valid"); + let logical_shard_size = + usize::try_from(erasure.shard_file_size(payload.len() as i64)).expect("logical shard size should fit usize"); + coding::bitrot_verify( + Cursor::new(inline_data.clone()), + inline_data.len(), + logical_shard_size, + HashAlgorithm::HighwayHash256S, + erasure.shard_size(), + ) + .await + .unwrap_or_else(|err| panic!("disk {disk_index} inline shard should pass bitrot verification: {err}")); + } + + let mut object_reader = set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("committed inline object should be readable"); + let mut restored = Vec::new(); + object_reader + .stream + .read_to_end(&mut restored) + .await + .expect("inline object should stream"); + assert_eq!(restored, payload); + } + + #[tokio::test] + async fn inline_put_direct_commit_accepts_exact_quorum_and_rejects_quorum_minus_one() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "inline-direct-quorum"; + let exact_quorum_object = "exact-quorum.bin"; + let below_quorum_object = "below-quorum.bin"; + make_bucket(&disk_stores, bucket).await; + { + let mut disks = set_disks.disks.write().await; + disks[3] = None; + } + + let mut reader = PutObjReader::from_vec(vec![0x5a; 4 * 1024]); + set_disks + .put_object(bucket, exact_quorum_object, &mut reader, &ObjectOptions::default()) + .await + .expect("three online disks should satisfy the four-disk write quorum"); + for (disk_index, disk) in disk_stores.iter().enumerate() { + let persisted = disk + .read_version("", bucket, exact_quorum_object, "", &ReadOptions::default()) + .await; + assert_eq!( + persisted.is_ok(), + disk_index < 3, + "exact-quorum commit should publish only on the three online disks" + ); + } + + set_disks.disks.write().await[2] = None; + let mut reader = PutObjReader::from_vec(vec![0xa5; 4 * 1024]); + set_disks + .put_object(bucket, below_quorum_object, &mut reader, &ObjectOptions::default()) + .await + .expect_err("two online disks are one below the four-disk write quorum"); + + for (disk_index, disk) in disk_stores.iter().enumerate() { + assert!( + disk.read_version("", bucket, below_quorum_object, "", &ReadOptions::default()) + .await + .is_err(), + "disk {disk_index} must not expose an object after pre-commit quorum failure" + ); + } + } + + #[tokio::test] + async fn inline_put_direct_commit_handles_post_encode_rename_failures() { + use crate::disk::health_state::RuntimeDriveHealthState; + + let payload = vec![0x5a; 4 * 1024]; + + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "inline-direct-post-encode-quorum"; + let object = "exact-quorum.bin"; + make_bucket(&disk_stores, bucket).await; + let barrier = PutObjectCommitBarrier::install(bucket, object, PutObjectCommitPause::AfterNamespace); + let put = { + let set_disks = Arc::clone(&set_disks); + let payload = payload.clone(); + tokio::spawn(async move { + let mut reader = PutObjReader::from_vec(payload); + set_disks + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + }) + }; + barrier.wait_until_paused().await; + disk_stores[3].force_runtime_state_for_test(RuntimeDriveHealthState::Offline); + barrier.release(); + put.await + .expect("exact-quorum PUT task should complete") + .expect("one post-encode rename failure should preserve write quorum"); + disk_stores[3].force_runtime_state_for_test(RuntimeDriveHealthState::Online); + for (disk_index, disk) in disk_stores.iter().enumerate() { + let persisted = disk.read_version("", bucket, object, "", &ReadOptions::default()).await; + assert_eq!( + persisted.is_ok(), + disk_index < 3, + "only disks that completed rename_data may publish the exact-quorum object" + ); + } + + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "inline-direct-post-encode-rollback"; + let object = "rollback.bin"; + let old_payload = vec![0x31; 4 * 1024]; + make_bucket(&disk_stores, bucket).await; + let mut old_reader = PutObjReader::from_vec(old_payload.clone()); + set_disks + .put_object(bucket, object, &mut old_reader, &ObjectOptions::default()) + .await + .expect("old inline object should commit"); + let read_data = ReadOptions { + read_data: true, + ..Default::default() + }; + let mut old_disk_data = Vec::with_capacity(disk_stores.len()); + for disk in &disk_stores { + old_disk_data.push( + disk.read_version("", bucket, object, "", &read_data) + .await + .expect("old inline shard should be readable before overwrite") + .data, + ); + } + + let barrier = PutObjectCommitBarrier::install(bucket, object, PutObjectCommitPause::AfterNamespace); + let put = { + let set_disks = Arc::clone(&set_disks); + let payload = payload.clone(); + tokio::spawn(async move { + let mut reader = PutObjReader::from_vec(payload); + set_disks + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + }) + }; + barrier.wait_until_paused().await; + for disk in &disk_stores[2..] { + disk.force_runtime_state_for_test(RuntimeDriveHealthState::Offline); + } + barrier.release(); + put.await + .expect("quorum-minus-one PUT task should complete") + .expect_err("two post-encode rename failures must fail write quorum"); + for disk in &disk_stores[2..] { + disk.force_runtime_state_for_test(RuntimeDriveHealthState::Online); + } + + for (disk_index, disk) in disk_stores.iter().enumerate() { + let restored = disk + .read_version("", bucket, object, "", &read_data) + .await + .unwrap_or_else(|err| panic!("disk {disk_index} should retain the old inline object: {err}")); + assert_eq!(restored.data, old_disk_data[disk_index]); + } + let mut object_reader = set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("old object should remain readable after quorum rollback"); + let mut restored = Vec::new(); + object_reader + .stream + .read_to_end(&mut restored) + .await + .expect("old object should stream after quorum rollback"); + assert_eq!(restored, old_payload); + } + + #[tokio::test] + async fn zero_length_put_keeps_existing_pipeline_layout_and_round_trips() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "zero-length-put"; + let object = "empty.bin"; + make_bucket(&disk_stores, bucket).await; + + let mut reader = PutObjReader::from_vec(Vec::new()); + set_disks + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + .expect("zero-length PUT should commit through the existing pipeline"); + + let read_data = ReadOptions { + read_data: true, + ..Default::default() + }; + for (disk_index, disk) in disk_stores.iter().enumerate() { + let file_info = disk + .read_version("", bucket, object, "", &read_data) + .await + .unwrap_or_else(|err| panic!("disk {disk_index} should persist empty-object metadata: {err}")); + assert_eq!(file_info.size, 0); + assert_eq!(file_info.data.as_deref(), Some(&[][..])); + } + + let mut object_reader = set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("empty object should be readable"); + let mut restored = Vec::new(); + object_reader + .stream + .read_to_end(&mut restored) + .await + .expect("empty object should stream"); + assert!(restored.is_empty()); + } +} + #[cfg(test)] mod get_object_downstream_close_accounting_tests { use super::hermetic_set_disks_support::hermetic_set_disks; diff --git a/crates/io-metrics/src/lib.rs b/crates/io-metrics/src/lib.rs index 6d2df2b9c..658dee400 100644 --- a/crates/io-metrics/src/lib.rs +++ b/crates/io-metrics/src/lib.rs @@ -58,7 +58,7 @@ use std::sync::{ /// When `false`, `record_put_object_path` and `record_put_object_stage_duration` /// become no-ops, and callers can skip the `Instant::now()` syscalls entirely. /// -/// Set to `true` during startup when OTEL metric export is enabled. +/// Enabled only through an explicit runtime opt-in. static PUT_STAGE_METRICS_ENABLED: AtomicBool = AtomicBool::new(false); static GET_STAGE_METRICS_ENABLED: AtomicBool = AtomicBool::new(false); @@ -78,7 +78,7 @@ static METRICS_ENABLED: AtomicBool = AtomicBool::new(false); /// Enable or disable detailed per-stage PUT metrics. /// -/// Called once during startup, typically gated by `rustfs_obs::observability_metric_enabled()`. +/// Called once during startup after applying the detailed PUT attribution opt-in. pub fn set_put_stage_metrics_enabled(enabled: bool) { PUT_STAGE_METRICS_ENABLED.store(enabled, Ordering::Relaxed); } @@ -103,6 +103,12 @@ pub fn put_stage_metrics_enabled() -> bool { PUT_STAGE_METRICS_ENABLED.load(Ordering::Relaxed) } +/// Start a PUT-stage timer only when detailed PUT attribution is enabled. +#[inline(always)] +pub fn put_stage_timer() -> Option { + put_stage_metrics_enabled().then(std::time::Instant::now) +} + #[inline(always)] pub fn get_stage_metrics_enabled() -> bool { GET_STAGE_METRICS_ENABLED.load(Ordering::Relaxed) @@ -434,7 +440,7 @@ pub fn record_get_object_request_result(status: &str, duration_secs: f64) { /// Record PutObject request start. #[inline(always)] pub fn record_put_object_request_start(concurrent_requests: usize) { - if !put_stage_metrics_enabled() { + if !metrics_enabled() { return; } counter!("rustfs_io_put_object_requests_total").increment(1); @@ -444,7 +450,7 @@ pub fn record_put_object_request_start(concurrent_requests: usize) { /// Record PutObject request result. #[inline(always)] pub fn record_put_object_request_result(status: &str, duration_secs: f64) { - if !put_stage_metrics_enabled() { + if !metrics_enabled() { return; } counter!("rustfs_io_put_object_request_results_total", "status" => status.to_string()).increment(1); @@ -1905,7 +1911,7 @@ pub fn record_get_object(duration_ms: f64, size_bytes: i64) { /// * `zero_copy_eligible` - Whether the request was eligible for a zero-copy path #[inline(always)] pub fn record_put_object(duration_ms: f64, size_bytes: i64, zero_copy_eligible: bool) { - if !put_stage_metrics_enabled() { + if !metrics_enabled() { return; } counter!("rustfs_s3_put_object_total").increment(1); @@ -2004,6 +2010,13 @@ pub fn record_put_object_stage_duration(stage: &'static str, duration_ms: f64) { histogram!("rustfs_s3_put_object_stage_duration_ms", "stage" => stage).record(duration_ms); } +#[inline(always)] +pub fn record_put_object_stage_duration_from(stage: &'static str, started_at: Option) { + if let Some(started_at) = started_at { + record_put_object_stage_duration(stage, started_at.elapsed().as_secs_f64() * 1000.0); + } +} + /// Record generic internal operation stage duration (non-PUT paths). /// Use this for metacache walks, listing, lifecycle, and other background /// operations that are NOT part of the PUT object hot path. @@ -2819,6 +2832,56 @@ mod tests { assert!(!put_stage_metrics_enabled()); } + #[test] + fn put_stage_gate_does_not_disable_basic_put_metrics() { + let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner()); + let recorder = DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + + metrics::with_local_recorder(&recorder, || { + set_metrics_enabled(true); + set_put_stage_metrics_enabled(false); + record_put_object_request_start(1); + record_put_object_request_result("ok", 0.001); + record_put_object(1.0, 1024, false); + record_put_object_stage_duration("disabled_stage", 0.5); + + set_put_stage_metrics_enabled(true); + record_put_object_stage_duration("enabled_stage", 0.5); + + set_put_stage_metrics_enabled(false); + set_metrics_enabled(false); + }); + + let metrics = snapshotter.snapshot().into_vec(); + assert!(metrics.iter().any(|(composite, _, _, _)| { + composite.kind() == MetricKind::Counter && composite.key().name() == "rustfs_s3_put_object_total" + })); + assert!(metrics.iter().any(|(composite, _, _, _)| { + composite.kind() == MetricKind::Counter && composite.key().name() == "rustfs_io_put_object_requests_total" + })); + + let stages = metrics + .iter() + .filter(|(composite, _, _, _)| { + composite.kind() == MetricKind::Histogram && composite.key().name() == "rustfs_s3_put_object_stage_duration_ms" + }) + .flat_map(|(composite, _, _, _)| composite.key().labels().map(|label| label.value().to_string())) + .collect::>(); + assert_eq!(stages, ["enabled_stage"]); + } + + #[test] + fn test_put_stage_timer_follows_metrics_switch() { + let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner()); + set_put_stage_metrics_enabled(false); + assert!(put_stage_timer().is_none()); + + set_put_stage_metrics_enabled(true); + assert!(put_stage_timer().is_some()); + set_put_stage_metrics_enabled(false); + } + #[test] fn test_record_get_object_path_and_stage() { let _guard = METRICS_FLAG_LOCK.lock().unwrap_or_else(|e| e.into_inner()); diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 8e420e97e..a33d5ff8a 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -194,7 +194,7 @@ use std::str::FromStr; use std::sync::atomic::AtomicUsize; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Mutex, OnceLock}; -use std::time::Duration; +use std::time::{Duration, Instant}; use time::{OffsetDateTime, format_description::well_known::Rfc3339}; use tokio::io::{AsyncRead, ReadBuf}; use tokio::sync::{OwnedSemaphorePermit, RwLock}; @@ -5596,7 +5596,8 @@ impl DefaultObjectUsecase { // Bucket-quota admission runs exactly once, and only now that the authoritative object length is known. `size` is the same basis the settle phase records via ObjectInfo.size (actual, pre-compression/pre-encryption logical size), NOT the aws-chunked wire Content-Length. When no quota is configured this stays a zero-extra-I/O fast path; once a hard quota is set, checker/config/usage faults fail closed with a retryable error. self.check_bucket_quota(&bucket, quota_operation, size as u64).await?; - let ingress_stage_start = std::time::Instant::now(); + let put_stage_metrics_enabled = rustfs_io_metrics::put_stage_metrics_enabled(); + let ingress_stage_start = put_stage_metrics_enabled.then(Instant::now); let should_compress = is_disk_compressible(&req.headers, &key) && size > MIN_DISK_COMPRESSIBLE_SIZE as i64 && !ciphertext_passthrough; let server_side_encryption_requested = @@ -5660,9 +5661,13 @@ impl DefaultObjectUsecase { let Some(store) = self.object_store() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; + let bucket_validate_stage_start = put_stage_metrics_enabled.then(Instant::now); validate_bucket_exists(&store, &bucket).await?; + rustfs_io_metrics::record_put_object_stage_duration_from("app_bucket_validate", bucket_validate_stage_start); + let sse_config_stage_start = put_stage_metrics_enabled.then(Instant::now); let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok(); + rustfs_io_metrics::record_put_object_stage_duration_from("app_sse_config_lookup", sse_config_stage_start); debug!( target: "rustfs::app::object_usecase", component = "app", @@ -5728,7 +5733,9 @@ impl DefaultObjectUsecase { let mut metadata = metadata.unwrap_or_default(); let has_explicit_object_lock_retention = object_lock_mode.is_some() || object_lock_retain_until_date.is_some(); + let object_lock_config_stage_start = put_stage_metrics_enabled.then(Instant::now); let object_lock_config_state = load_bucket_object_lock_config_state(&bucket).await?; + rustfs_io_metrics::record_put_object_stage_duration_from("app_object_lock_config_lookup", object_lock_config_stage_start); apply_put_request_metadata( &mut metadata, &req.headers, @@ -5750,6 +5757,7 @@ impl DefaultObjectUsecase { has_explicit_object_lock_retention, )?; + let put_opts_stage_start = put_stage_metrics_enabled.then(Instant::now); let mut opts: ObjectOptions = put_opts_with_replication_authorization( &bucket, &key, @@ -5760,6 +5768,7 @@ impl DefaultObjectUsecase { ) .await .map_err(ApiError::from)?; + rustfs_io_metrics::record_put_object_stage_duration_from("app_put_opts_build", put_opts_stage_start); apply_bucket_generation_guard(&req, &bucket, &mut opts)?; apply_put_request_object_lock_opts( &bucket, @@ -5778,10 +5787,11 @@ impl DefaultObjectUsecase { // replication), the lookup is skipped and accounting is backfilled from // the dst xl.meta that rename_data already reads, saving a full-disk // metadata fanout per PUT. - let prelookup_required = version_id.is_some() || object_lock_checks_required(&bucket).await; + let prelookup_required = version_id.is_some() || object_lock_checks_required_for_state(&object_lock_config_state); // Outer None = prelookup skipped (accounting comes from the commit // backfill); Some(inner) = the previous current size as observed by the // lookup, with the pre-#1009 semantics kept bit-for-bit. + let prelookup_stage_start = (prelookup_required && put_stage_metrics_enabled).then(Instant::now); let prelookup_previous_current_size: Option> = if prelookup_required { let current_opts: ObjectOptions = internal_object_info_lookup_opts( get_opts(&bucket, &key, version_id.clone(), None, &req.headers) @@ -5807,6 +5817,7 @@ impl DefaultObjectUsecase { } else { None }; + rustfs_io_metrics::record_put_object_stage_duration_from("app_prelookup", prelookup_stage_start); let actual_size = size; @@ -5895,15 +5906,13 @@ impl DefaultObjectUsecase { put_extra_checksum_headers = additional_checksum_echo_pairs(&opts.want_checksum); } rustfs_io_metrics::record_put_object_path(put_path); - rustfs_io_metrics::record_put_object_stage_duration( - "ingress_prepare", - ingress_stage_start.elapsed().as_secs_f64() * 1000.0, - ); + rustfs_io_metrics::record_put_object_stage_duration_from("ingress_prepare", ingress_stage_start); let mut helper = OperationHelper::new(&req, event_name, S3Operation::PutObject); let ssekms_context = extract_ssekms_context_from_headers(&req.headers)?; // Apply encryption using unified SSE API. + let encryption_stage_start = put_stage_metrics_enabled.then(Instant::now); let write_principal = SseKmsPrincipal::from_request(&req); let encryption_request = EncryptionRequest { bucket: &bucket, @@ -5947,6 +5956,7 @@ impl DefaultObjectUsecase { } reader = write_plan.apply(reader, actual_size).map_err(ApiError::from)?; + rustfs_io_metrics::record_put_object_stage_duration_from("app_encryption_prepare", encryption_stage_start); let mut reader = PutObjReader::new(reader); @@ -5963,9 +5973,11 @@ impl DefaultObjectUsecase { // post-commit schedule (see the reuse site further down), so a // replication-config hot update can no longer split the two phases // (https://github.com/rustfs/backlog/issues/1320). + let replication_decision_stage_start = put_stage_metrics_enabled.then(Instant::now); let dsc = must_replicate_object(&bucket, &key, &mt2, "".to_string(), opts.delete_marker_replication_status(), opts.clone()) .await; + rustfs_io_metrics::record_put_object_stage_duration_from("app_replication_decision", replication_decision_stage_start); if dsc.replicate_any() { insert_str(&mut opts.user_defined, SUFFIX_REPLICATION_TIMESTAMP, jiff::Zoned::now().to_string()); @@ -5977,7 +5989,12 @@ impl DefaultObjectUsecase { } let cache_adapter = self.object_data_cache(); + let cache_invalidate_before_stage_start = put_stage_metrics_enabled.then(Instant::now); let _ = invalidate_object_data_cache_before_mutation(&cache_adapter, &bucket, &key).await; + rustfs_io_metrics::record_put_object_stage_duration_from( + "app_cache_invalidate_before", + cache_invalidate_before_stage_start, + ); let store_put_watchdog = tokio_util::sync::CancellationToken::new(); spawn_traced({ @@ -6017,6 +6034,7 @@ impl DefaultObjectUsecase { let object_traffic_progress = object_traffic_health .as_deref() .and_then(ObjectTrafficHealth::track_write_storage); + let store_put_stage_start = put_stage_metrics_enabled.then(Instant::now); let (obj_info, backfilled_old_current_size) = match store .put_object_with_old_current_size(&bucket, &key, &mut reader, &opts) .await @@ -6042,6 +6060,7 @@ impl DefaultObjectUsecase { } Err(err) => { store_put_watchdog.cancel(); + rustfs_io_metrics::record_put_object_stage_duration_from("app_store_put", store_put_stage_start); warn!( target: "rustfs::app::object_usecase", event = EVENT_PUT_OBJECT_STORE_RETURNED, @@ -6063,10 +6082,12 @@ impl DefaultObjectUsecase { return result; } }; + rustfs_io_metrics::record_put_object_stage_duration_from("app_store_put", store_put_stage_start); drop(object_traffic_progress); #[cfg(test)] wait_for_put_post_store_test_hook(&bucket).await; + let post_store_stage_start = put_stage_metrics_enabled.then(Instant::now); maybe_enqueue_transition_immediate(&obj_info, LcEventSrc::S3PutObject).await; let _ = invalidate_object_data_cache_after_put_success(&cache_adapter, &bucket, &key).await; @@ -6165,10 +6186,13 @@ impl DefaultObjectUsecase { let result = Ok(response); let _ = helper.complete(&result); rustfs_scanner::record_dirty_usage_bucket(&bucket); + rustfs_io_metrics::record_put_object_stage_duration_from("app_post_store_bookkeeping", post_store_stage_start); // Record write operation for capacity management (inline to avoid per-request tokio::spawn overhead) + let capacity_update_stage_start = put_stage_metrics_enabled.then(Instant::now); let manager = get_capacity_manager(); manager.record_write_operation().await; + rustfs_io_metrics::record_put_object_stage_duration_from("app_capacity_update", capacity_update_stage_start); // Record PutObject metrics via zero-copy-metrics { @@ -9572,6 +9596,13 @@ pub(super) async fn object_lock_checks_required(bucket: &str) -> bool { .map_or(true, |metadata| metadata.object_locking()) } +fn object_lock_checks_required_for_state(state: &metadata_sys::ObjectLockConfigState) -> bool { + match state { + metadata_sys::ObjectLockConfigState::Configured { .. } | metadata_sys::ObjectLockConfigState::Fabricated => true, + metadata_sys::ObjectLockConfigState::ConfirmedAbsent => false, + } +} + /// rustfs/backlog#1009: map the rename_data old-size backfill onto the /// `previous_current_size` value the usage-accounting helpers expect. Outer /// `None` = unknown (no quorum agreement, or a peer predates the field) — the @@ -10198,6 +10229,23 @@ mod tests { assert_eq!(err.message(), Some(ERR_OBJECT_LOCK_RETENTION_HEADERS_MUST_BE_PAIRED)); } + #[test] + fn object_lock_checks_required_reuses_authoritative_state() { + assert!(!object_lock_checks_required_for_state( + &metadata_sys::ObjectLockConfigState::ConfirmedAbsent + )); + + let configured = metadata_sys::ObjectLockConfigState::Configured { + config: ObjectLockConfiguration { + object_lock_enabled: Some(ObjectLockEnabled::from_static(ObjectLockEnabled::ENABLED)), + rule: None, + }, + updated_at: OffsetDateTime::now_utc(), + }; + assert!(object_lock_checks_required_for_state(&configured)); + assert!(object_lock_checks_required_for_state(&metadata_sys::ObjectLockConfigState::Fabricated)); + } + #[test] fn build_put_like_object_lock_metadata_rejects_retain_until_date_without_mode() { let retain_until = Timestamp::from(OffsetDateTime::now_utc().add(time::Duration::days(1))); diff --git a/rustfs/src/startup_observability.rs b/rustfs/src/startup_observability.rs index 0dd66b875..01ec050b8 100644 --- a/rustfs/src/startup_observability.rs +++ b/rustfs/src/startup_observability.rs @@ -24,14 +24,63 @@ pub(crate) async fn init_observability_runtime(store: Arc, ctx: Cancell init_update_check(); crate::allocator_reclaim::init_allocator_reclaim(ctx.clone()); - if startup_runtime_sources::observability_metric_enabled() { + let metrics_enabled = startup_runtime_sources::observability_metric_enabled(); + configure_metric_gates(metrics_enabled); + + if metrics_enabled { // Load persisted compression stats into memory early, before any PUTs can occur. init_compression_total_memory_from_backend(store).await; - startup_runtime_sources::set_put_stage_metrics_enabled(true); - startup_runtime_sources::set_get_stage_metrics_enabled(true); - startup_runtime_sources::set_metrics_enabled(true); startup_runtime_sources::init_metrics_runtime(ctx.clone()); crate::memory_observability::init_memory_observability(ctx.clone()); init_auto_tuner(ctx).await; } } + +fn configure_metric_gates(metrics_enabled: bool) { + let put_stage_metrics_enabled = metrics_enabled + && rustfs_utils::get_env_bool( + rustfs_config::observability::ENV_OBS_PUT_STAGE_METRICS_ENABLED, + rustfs_config::DEFAULT_OBS_PUT_STAGE_METRICS_ENABLED, + ); + startup_runtime_sources::set_put_stage_metrics_enabled(put_stage_metrics_enabled); + startup_runtime_sources::set_get_stage_metrics_enabled(metrics_enabled); + startup_runtime_sources::set_metrics_enabled(metrics_enabled); +} + +#[cfg(test)] +mod tests { + use super::*; + + const PUT_STAGE_ENV: &str = rustfs_config::observability::ENV_OBS_PUT_STAGE_METRICS_ENABLED; + + #[test] + #[serial_test::serial] + fn put_stage_metrics_require_explicit_opt_in() { + let previous_metrics = rustfs_io_metrics::metrics_enabled(); + let previous_get_stages = rustfs_io_metrics::get_stage_metrics_enabled(); + let previous_put_stages = rustfs_io_metrics::put_stage_metrics_enabled(); + + temp_env::with_var(PUT_STAGE_ENV, None::<&str>, || { + configure_metric_gates(true); + assert!(rustfs_io_metrics::metrics_enabled()); + assert!(rustfs_io_metrics::get_stage_metrics_enabled()); + assert!(!rustfs_io_metrics::put_stage_metrics_enabled()); + }); + + temp_env::with_var(PUT_STAGE_ENV, Some("true"), || { + configure_metric_gates(true); + assert!(rustfs_io_metrics::metrics_enabled()); + assert!(rustfs_io_metrics::get_stage_metrics_enabled()); + assert!(rustfs_io_metrics::put_stage_metrics_enabled()); + + configure_metric_gates(false); + assert!(!rustfs_io_metrics::metrics_enabled()); + assert!(!rustfs_io_metrics::get_stage_metrics_enabled()); + assert!(!rustfs_io_metrics::put_stage_metrics_enabled()); + }); + + startup_runtime_sources::set_metrics_enabled(previous_metrics); + startup_runtime_sources::set_get_stage_metrics_enabled(previous_get_stages); + startup_runtime_sources::set_put_stage_metrics_enabled(previous_put_stages); + } +}