From 498205b7ecf1ef11e9dba9a6f3553c321c0378e7 Mon Sep 17 00:00:00 2001 From: houseme Date: Sun, 30 Aug 2026 03:39:51 +0800 Subject: [PATCH] fix(ecstore): keep 1MiB GET off mid-size reader (#6861) Co-authored-by: heihutu --- crates/ecstore/src/set_disk/mod.rs | 1 + crates/ecstore/src/set_disk/ops/object.rs | 78 +++++++++++++++++++++-- 2 files changed, 74 insertions(+), 5 deletions(-) diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 7fd95925a..4403d0db8 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -2525,6 +2525,7 @@ fn record_get_object_reader_path_observation( GET_OBJECT_PATH_CODEC_STREAMING => 5, GET_OBJECT_PATH_REMOTE_TRANSITION => 6, GET_OBJECT_PATH_EMPTY => 7, + GET_OBJECT_PATH_LEGACY_DUPLEX => 8, _ => 255, }, Ordering::Relaxed, diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 7e0d86b19..c92ad68de 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -89,7 +89,10 @@ use tokio::io::AsyncWriteExt; const ENV_RUSTFS_GET_MID_SIZE_STREAMING_ENABLE: &str = "RUSTFS_GET_MID_SIZE_STREAMING_ENABLE"; const DEFAULT_RUSTFS_GET_MID_SIZE_STREAMING_ENABLE: bool = true; const GET_MID_SIZE_STREAMING_MIN_SIZE: usize = 128 * 1024 + 1; -const GET_MID_SIZE_STREAMING_MAX_SIZE: usize = 1024 * 1024; +// Exclude 1 MiB from the bounded mid-size reader until it has a demonstrated +// high-concurrency performance envelope; existing codec/legacy gates decide +// which established reader handles the object. +const GET_MID_SIZE_STREAMING_MAX_SIZE: usize = 512 * 1024; fn is_get_mid_size_streaming_enabled() -> bool { #[cfg(test)] @@ -8408,14 +8411,14 @@ mod mid_size_streaming_gate_tests { #[test] #[serial] - fn mid_size_streaming_includes_one_mib_and_rejects_larger_objects() { - let (object_info, fi) = plain_metadata(1024 * 1024); + fn mid_size_streaming_stops_at_512kib_and_rejects_one_mib() { + let (object_info, fi) = plain_metadata(512 * 1024); assert_eq!( get_mid_size_streaming_object_size_with_flags(&None, &object_info, &fi, &ObjectOptions::default(), true, true, true), - Some(1024 * 1024) + Some(512 * 1024) ); - let (large_info, large_fi) = plain_metadata(1024 * 1024 + 1); + let (large_info, large_fi) = plain_metadata(512 * 1024 + 1); assert_eq!( get_mid_size_streaming_object_size_with_flags( &None, @@ -8428,6 +8431,20 @@ mod mid_size_streaming_gate_tests { ), None ); + + let (one_mib_info, one_mib_fi) = plain_metadata(1024 * 1024); + assert_eq!( + get_mid_size_streaming_object_size_with_flags( + &None, + &one_mib_info, + &one_mib_fi, + &ObjectOptions::default(), + true, + true, + true + ), + None + ); } #[test] @@ -9801,6 +9818,57 @@ mod inline_put_commit_path_tests { .await; } + #[tokio::test] + #[serial] + async fn get_object_reader_routes_one_mib_away_from_mid_size_reader() { + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; + let bucket = "one-mib-legacy-reader"; + let object = "object.bin"; + let payload = vec![0x5a; 1024 * 1024]; + make_bucket(&disk_stores, bucket).await; + let storage_class = temp_env::with_var(INLINE_BLOCK_ENV, Some("1KiB"), || lookup_config_for_pools(&KVS::new(), &[4])) + .expect("test storage class should resolve"); + set_disks.set_test_storage_class_config(storage_class); + + let mut writer = PutObjReader::from_vec(payload.clone()); + temp_env::async_with_vars( + [ + (ENV_RUSTFS_GET_MID_SIZE_STREAMING_ENABLE, Some("true")), + (crate::set_disk::ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, Some("true")), + (crate::set_disk::ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, Some("true")), + (crate::set_disk::ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, Some("true")), + (crate::set_disk::ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, Some("off")), + (rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("true")), + ], + async { + set_disks + .put_object(bucket, object, &mut writer, &ObjectOptions::default()) + .await + .expect("1 MiB fixture should commit"); + + crate::set_disk::reset_test_get_object_reader_path(); + let mut reader = set_disks + .get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("1 MiB legacy GET should succeed"); + let mut restored = Vec::new(); + reader + .stream + .read_to_end(&mut restored) + .await + .expect("1 MiB legacy reader should stream"); + + assert_eq!(restored, payload); + assert_eq!( + crate::set_disk::test_get_object_reader_path_id(), + 8, + "1 MiB must bypass mid-size and use legacy duplex when codec rollout is off" + ); + }, + ) + .await; + } + #[tokio::test] async fn repeated_gets_reuse_the_set_erasure_shell() { let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;