mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-13 16:46:55 +00:00
perf(ecstore): reduce inline PUT commit overhead (#6033)
* perf(metrics): attribute PUT stage costs Co-Authored-By: heihutu <heihutu@gmail.com> * perf(ecstore): move PUT metadata during shuffle Co-Authored-By: heihutu <heihutu@gmail.com> * perf(s3): reuse PUT object lock state Co-Authored-By: heihutu <heihutu@gmail.com> * 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 <heihutu@gmail.com> * perf(metrics): make PUT stage attribution opt-in Co-Authored-By: heihutu <heihutu@gmail.com> * perf(ecstore): commit inline PUT shards directly Co-Authored-By: heihutu <heihutu@gmail.com> * 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 <heihutu@gmail.com> * 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 <heihutu@gmail.com> --------- Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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<R>(
|
||||
self: Arc<Self>,
|
||||
mut reader: R,
|
||||
size_hint: usize,
|
||||
) -> std::io::Result<(R, usize, Vec<Bytes>)>
|
||||
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<R>(
|
||||
self: Arc<Self>,
|
||||
@@ -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]
|
||||
|
||||
@@ -1079,6 +1079,25 @@ impl SetDisks {
|
||||
shuffled_disks
|
||||
}
|
||||
|
||||
pub(super) fn shuffle_disks_owned(mut disks: Vec<Option<DiskStore>>, distribution: &[usize]) -> Vec<Option<DiskStore>> {
|
||||
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<usize> {
|
||||
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::<Vec<_>>();
|
||||
let actual_slots = actual.iter().map(Option::is_some).collect::<Vec<_>>();
|
||||
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");
|
||||
|
||||
@@ -61,6 +61,10 @@ fn duration_millis_f64(duration: std::time::Duration) -> f64 {
|
||||
duration.as_secs_f64() * 1000.0
|
||||
}
|
||||
|
||||
fn committed_response_metadata_slot<D>(committed_disks: &[Option<D>], 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<u64> = 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<u8> = (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;
|
||||
|
||||
@@ -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<std::time::Instant> {
|
||||
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<std::time::Instant>) {
|
||||
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::<Vec<_>>();
|
||||
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());
|
||||
|
||||
Reference in New Issue
Block a user