fix(storage): cover inline reader fallback controls (#5169)

* test(ecstore): cover inline fast path boundaries

Co-Authored-By: heihutu <heihutu@gmail.com>

* fix(storage): cover inline reader fallback controls

Co-Authored-By: heihutu <heihutu@gmail.com>

* perf(tier): keep commit fanout concurrent

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-07-24 14:08:17 +08:00
committed by GitHub
parent 6f6d8a4d3e
commit 4133fbe0fc
11 changed files with 1383 additions and 36 deletions
+9
View File
@@ -39,6 +39,7 @@ ecstore-serial-flaky = { max-threads = 1 }
# servers never run at once. ci-7's nightly picks these up via the e2e suite;
# they are deliberately NOT in the fast PR `e2e-smoke` filter.
e2e-reliability = { max-threads = 1 }
e2e-inline-boundaries = { max-threads = 1 }
# --- default profile (local): serialize the flaky groups, never retry --------
[[profile.default.overrides]]
@@ -61,6 +62,10 @@ test-group = 'ecstore-serial-flaky'
filter = 'package(e2e_test) & test(/^(reliability_disk_fault|degraded_read_eof_regression)_test::/)'
test-group = 'e2e-reliability'
[[profile.default.overrides]]
filter = 'package(e2e_test) & test(/^inline_fast_path_cluster_test::/)'
test-group = 'e2e-inline-boundaries'
# ---------------------------------------------------------------------------
# ci profile — the strict CI gate (ci.yml `cargo nextest run --profile ci`)
# ---------------------------------------------------------------------------
@@ -314,3 +319,7 @@ path = "junit.xml"
[[profile.e2e-full.overrides]]
filter = 'package(e2e_test) & test(/^(reliability_disk_fault|degraded_read_eof_regression)_test::/)'
test-group = 'e2e-reliability'
[[profile.e2e-full.overrides]]
filter = 'package(e2e_test) & test(/^inline_fast_path_cluster_test::/)'
test-group = 'e2e-inline-boundaries'
Generated
+2
View File
@@ -3676,6 +3676,8 @@ dependencies = [
"hyper-util",
"local-ip-address",
"md5",
"opentelemetry-proto",
"prost 0.14.4",
"rand 0.10.2",
"rcgen",
"reqwest",
+1
View File
@@ -323,6 +323,7 @@ dial9-tokio-telemetry = "0.3"
opentelemetry = { version = "0.32.0" }
opentelemetry-appender-tracing = { version = "0.32.0" }
opentelemetry-otlp = { version = "0.32.0" }
opentelemetry-proto = { version = "0.32.0", default-features = false, features = ["metrics", "gen-tonic-messages"] }
opentelemetry_sdk = { version = "0.32.1" }
opentelemetry-semantic-conventions = { version = "0.32.1" }
opentelemetry-stdout = { version = "0.32.0" }
+2
View File
@@ -69,6 +69,8 @@ base64 = { workspace = true }
rand = { workspace = true, features = ["serde"] }
chrono = { workspace = true, features = ["serde"] }
md5 = { workspace = true }
opentelemetry-proto = { workspace = true }
prost.workspace = true
sha2 = { workspace = true }
astral-tokio-tar = { workspace = true }
s3s = { workspace = true, features = ["minio"] }
File diff suppressed because it is too large Load Diff
+4
View File
@@ -167,6 +167,10 @@ mod cluster_concurrency_test;
#[cfg(test)]
mod cluster_multidrive_pool_test;
// backlog#1433: real 4-node EC boundary gate for inline storage and GET paths.
#[cfg(test)]
mod inline_fast_path_cluster_test;
// PutObject / MultipartUpload with checksum (Content-MD5, x-amz-checksum-*)
#[cfg(test)]
mod checksum_upload_test;
+44
View File
@@ -512,6 +512,50 @@ mod tests {
StorageClassEnvOverrides::default()
}
#[test]
fn should_inline_preserves_exact_default_shard_boundaries() {
let config = Config::default();
for (case, shard_size, versioned, expected) in [
("unversioned below", 128 * 1024 - 1, false, true),
("unversioned exact", 128 * 1024, false, true),
("unversioned above", 128 * 1024 + 1, false, false),
("versioned below", 16 * 1024 - 1, true, true),
("versioned exact", 16 * 1024, true, true),
("versioned above", 16 * 1024 + 1, true, false),
("negative", -1, false, false),
] {
assert_eq!(
config.should_inline(shard_size, versioned),
expected,
"{case}: shard_size={shard_size}, versioned={versioned}"
);
}
}
#[test]
fn should_inline_preserves_exact_default_ec_2_2_object_boundaries() {
let config = Config::default();
let erasure = crate::erasure::coding::Erasure::new(2, 2, 1024 * 1024);
for (case, object_size, versioned, expected_shard_size, expected) in [
("unversioned below", 256 * 1024 - 1, false, 128 * 1024, true),
("unversioned exact", 256 * 1024, false, 128 * 1024, true),
("unversioned above", 256 * 1024 + 1, false, 128 * 1024 + 1, false),
("versioned below", 32 * 1024 - 1, true, 16 * 1024, true),
("versioned exact", 32 * 1024, true, 16 * 1024, true),
("versioned above", 32 * 1024 + 1, true, 16 * 1024 + 1, false),
] {
let shard_size = erasure.shard_file_size(object_size);
assert_eq!(shard_size, expected_shard_size, "{case}: object_size={object_size}");
assert_eq!(
config.should_inline(shard_size, versioned),
expected,
"{case}: object_size={object_size}, shard_size={shard_size}, versioned={versioned}"
);
}
}
#[test]
fn automatic_parity_is_resolved_per_pool() {
let cfg = lookup_config_for_pools_with_env(&KVS::new(), &[4, 2], no_env_overrides())
+54
View File
@@ -800,6 +800,60 @@ mod tests {
use super::*;
use rustfs_filemeta::{FileInfo, FileMeta, MetaCacheEntry, TRANSITION_COMPLETE};
fn inline_fast_path_object(size: i64, versioned: bool) -> ObjectInfo {
ObjectInfo {
size,
inlined: true,
version_id: versioned.then(|| Uuid::from_u128(1)),
parts: Arc::new(vec![ObjectPartInfo::default()]),
..Default::default()
}
}
#[test]
fn inline_fast_path_eligibility_preserves_exact_versioned_boundaries() {
for (case, size, versioned, expected) in [
("unversioned below", 128 * 1024 - 1, false, true),
("unversioned exact", 128 * 1024, false, true),
("unversioned above", 128 * 1024 + 1, false, false),
("versioned below", 16 * 1024 - 1, true, true),
("versioned exact", 16 * 1024, true, true),
("versioned above", 16 * 1024 + 1, true, false),
] {
assert_eq!(
inline_fast_path_object(size, versioned).is_inline_fast_path_eligible(),
expected,
"{case}: object_size={size}, versioned={versioned}"
);
}
}
#[test]
fn inline_fast_path_eligibility_rejects_incompatible_object_shapes() {
let mut object = inline_fast_path_object(ObjectInfo::INLINE_MAX_SIZE, false);
object.inlined = false;
assert!(!object.is_inline_fast_path_eligible(), "non-inline objects must fall back");
object.inlined = true;
object.parts = Arc::new(vec![ObjectPartInfo::default(), ObjectPartInfo::default()]);
assert!(!object.is_inline_fast_path_eligible(), "multipart objects must fall back");
object.parts = Arc::new(vec![ObjectPartInfo::default()]);
object.user_defined = Arc::new(HashMap::from([("x-amz-server-side-encryption".to_string(), "AES256".to_string())]));
assert!(!object.is_inline_fast_path_eligible(), "encrypted objects must fall back");
object.user_defined = Arc::new(HashMap::from([(
rustfs_utils::http::internal_key_rustfs(rustfs_utils::http::SUFFIX_COMPRESSION),
"zstd".to_string(),
)]));
assert!(!object.is_inline_fast_path_eligible(), "compressed objects must fall back");
object.user_defined = Arc::default();
object.transitioned_object.tier = "remote-tier".to_string();
assert!(!object.is_inline_fast_path_eligible(), "transitioned objects must fall back");
}
#[test]
fn versions_after_marker_handles_null_version_marker() {
let first_version = Uuid::parse_str("11111111-2222-3333-4444-555555555555").unwrap();
+137 -24
View File
@@ -923,37 +923,30 @@ async fn prepare_tier_mutation_peers(
error: tier_mutation_fanout_admin_error("prepare", err),
prepared_peers: Vec::new(),
})?);
let results = join_all(peers.into_iter().map(|peer| {
let payload = payload.clone();
async move {
let label = peer.peer_label();
let result = peer.prepare_tier_mutation(mutation_id, payload).await;
(peer, label, result)
}
}))
.await;
let mut prepared = Vec::with_capacity(results.len());
let mut failure = None;
for (peer, label, result) in results {
let mut prepared = Vec::with_capacity(peers.len());
for peer in peers {
let label = peer.peer_label();
let result = peer.prepare_tier_mutation(mutation_id, payload.clone()).await;
match result {
Ok(PeerTierMutationState::Prepared) => prepared.push(peer),
Ok(state) => {
failure = Some(tier_mutation_fanout_admin_error(
"prepare",
format!("peer {label} returned unexpected state {state:?}"),
));
return Err(TierMutationPrepareFailure {
error: tier_mutation_fanout_admin_error(
"prepare",
format!("peer {label} returned unexpected state {state:?}"),
),
prepared_peers: prepared,
});
}
Err(err) => {
return Err(TierMutationPrepareFailure {
error: tier_mutation_fanout_admin_error("prepare", format!("peer {label}: {err}")),
prepared_peers: prepared,
});
}
Err(err) => failure = Some(tier_mutation_fanout_admin_error("prepare", format!("peer {label}: {err}"))),
}
}
if let Some(err) = failure {
return Err(TierMutationPrepareFailure {
error: err,
prepared_peers: prepared,
});
}
Ok(prepared)
}
@@ -5905,6 +5898,126 @@ mod tests {
}
}
struct ConcurrencyTrackingTierMutationPeer {
label: &'static str,
calls: Arc<Mutex<Vec<String>>>,
active: Arc<AtomicUsize>,
max_active: Arc<AtomicUsize>,
}
impl ConcurrencyTrackingTierMutationPeer {
fn boxed(
label: &'static str,
calls: Arc<Mutex<Vec<String>>>,
active: Arc<AtomicUsize>,
max_active: Arc<AtomicUsize>,
) -> Arc<dyn TierMutationPeer> {
Arc::new(Self {
label,
calls,
active,
max_active,
})
}
async fn track(&self, action: &str) {
let current = self.active.fetch_add(1, Ordering::SeqCst) + 1;
self.max_active.fetch_max(current, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(20)).await;
self.active.fetch_sub(1, Ordering::SeqCst);
lock_unpoisoned(&self.calls).push(format!("{}:{action}", self.label));
}
}
#[async_trait::async_trait]
impl TierMutationPeer for ConcurrencyTrackingTierMutationPeer {
fn peer_label(&self) -> String {
self.label.to_string()
}
async fn prepare_tier_mutation(
&self,
_mutation_id: uuid::Uuid,
_canonical_payload: Bytes,
) -> Result<PeerTierMutationState> {
self.track("prepare").await;
Ok(PeerTierMutationState::Prepared)
}
async fn commit_tier_mutation(
&self,
_mutation_id: uuid::Uuid,
_canonical_payload: Bytes,
) -> Result<PeerTierMutationState> {
self.track("commit").await;
Ok(PeerTierMutationState::Committed)
}
async fn abort_tier_mutation(&self, _mutation_id: uuid::Uuid) -> Result<PeerTierMutationState> {
Ok(PeerTierMutationState::Aborted)
}
}
#[tokio::test]
async fn prepare_tier_mutation_peers_serializes_peer_prepare_writes() {
let mutation_id = uuid::Uuid::from_u128(34);
let intent = prepared_remove_intent("COLD-A", mutation_id);
let calls = Arc::new(Mutex::new(Vec::new()));
let active = Arc::new(AtomicUsize::new(0));
let max_active = Arc::new(AtomicUsize::new(0));
let prepared = prepare_tier_mutation_peers(
mutation_id,
vec![
ConcurrencyTrackingTierMutationPeer::boxed("peer-a", calls.clone(), active.clone(), max_active.clone()),
ConcurrencyTrackingTierMutationPeer::boxed("peer-b", calls.clone(), active.clone(), max_active.clone()),
ConcurrencyTrackingTierMutationPeer::boxed("peer-c", calls.clone(), active.clone(), max_active.clone()),
],
&intent,
)
.await
.unwrap_or_else(|failure| panic!("successful prepare fanout should prepare every peer: {}", failure.error.message));
assert_eq!(prepared.len(), 3);
assert_eq!(
max_active.load(Ordering::SeqCst),
1,
"peer prepare fanout must not write the same intent concurrently"
);
assert_eq!(
lock_unpoisoned(&calls).as_slice(),
&["peer-a:prepare", "peer-b:prepare", "peer-c:prepare"]
);
}
#[tokio::test]
async fn commit_tier_mutation_peers_keeps_peer_commits_concurrent() {
let mutation_id = uuid::Uuid::from_u128(35);
let calls = Arc::new(Mutex::new(Vec::new()));
let active = Arc::new(AtomicUsize::new(0));
let max_active = Arc::new(AtomicUsize::new(0));
commit_tier_mutation_peers(
mutation_id,
vec![
ConcurrencyTrackingTierMutationPeer::boxed("peer-a", calls.clone(), active.clone(), max_active.clone()),
ConcurrencyTrackingTierMutationPeer::boxed("peer-b", calls.clone(), active.clone(), max_active.clone()),
ConcurrencyTrackingTierMutationPeer::boxed("peer-c", calls.clone(), active.clone(), max_active.clone()),
],
"etag-new",
)
.await
.expect("successful commit fanout should commit every peer");
assert!(
max_active.load(Ordering::SeqCst) > 1,
"peer commit fanout should remain concurrent after serializing prepare"
);
let mut calls = lock_unpoisoned(&calls).clone();
calls.sort();
assert_eq!(calls.as_slice(), &["peer-a:commit", "peer-b:commit", "peer-c:commit"]);
}
#[tokio::test]
async fn committed_mutation_recovery_replays_peer_commit() {
let mutation_id = uuid::Uuid::from_u128(13);
+3 -9
View File
@@ -24,7 +24,6 @@ use super::storage_api::multipart_usecase::bucket::{
replication::{must_replicate_object, schedule_object_replication},
versioning_sys::BucketVersioningSys,
};
use super::storage_api::multipart_usecase::compression::is_disk_compressible;
#[cfg(test)]
use super::storage_api::multipart_usecase::contract::http::HTTPPreconditions;
use super::storage_api::multipart_usecase::contract::multipart::{CompletePart, MultipartOperations as _, MultipartUploadResult};
@@ -37,7 +36,7 @@ use super::storage_api::multipart_usecase::error::{StorageError, is_err_object_n
use super::storage_api::multipart_usecase::helper::OperationHelper;
#[cfg(test)]
use super::storage_api::multipart_usecase::io::{DecryptReader, EncryptReader, HardLimitReader, boxed_reader, wrap_reader};
use super::storage_api::multipart_usecase::io::{HashReader, WriteEncryption, WritePlan, compression_metadata_value};
use super::storage_api::multipart_usecase::io::{HashReader, WriteEncryption, WritePlan};
use super::storage_api::multipart_usecase::object_utils::to_s3s_etag;
use super::storage_api::multipart_usecase::options::{
copy_src_opts, extract_metadata_from_mime, get_complete_multipart_upload_opts, get_content_sha256_with_query, get_opts,
@@ -704,13 +703,8 @@ impl DefaultMultipartUsecase {
None => (None, None),
};
if is_disk_compressible(&req.headers, &key) {
rustfs_utils::http::insert_str(
&mut metadata,
rustfs_utils::http::SUFFIX_COMPRESSION,
compression_metadata_value(CompressionAlgorithm::default()),
);
}
// Multipart parts are independent physical streams. Advertising object-level
// compression here would make GET decode the completed object as one stream.
let mt2 = metadata.clone();
let mut opts: ObjectOptions = put_opts(&bucket, &key, version_id, &req.headers, metadata)
+1 -3
View File
@@ -1027,9 +1027,7 @@ pub(crate) mod multipart_usecase {
}
}
pub(crate) use super::{
access, bucket, compression, data_usage, error, helper, io, object_utils, options, s3_api, set_disk, sse,
};
pub(crate) use super::{access, bucket, data_usage, error, helper, io, object_utils, options, s3_api, set_disk, sse};
pub(crate) use crate::storage::storage_api::{ECStore, StorageObjectOptions, StoragePutObjReader};
}