perf(storage): converge Wave 2 hot-path optimizations (#6065)

* perf(get): share inline shards and lock clients

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

* perf(ecstore): converge PUT encoding on contiguous blocks

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

* perf(get): cache codec streaming gate config

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

* fix(sse): redact projected customer headers

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

* perf(ecstore): collapse GET metadata snapshots

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

* perf(ecstore): reuse decode stripe scratch

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

* refactor(ecstore): trim decode scratch adapters

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

* test(ecstore): adapt transition checks to metadata snapshots

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

* perf(get): release metadata snapshots at ownership boundary

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

* refactor(ecstore): close cumulative fast-path findings

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

* fix(storage): preserve lock and header invariants

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

* test(ecstore): adapt cumulative paths after rebase

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

* fix(rio-v2): adapt generated metadata fixture

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

---------

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-08-13 16:34:28 +08:00
committed by GitHub
parent 36deab8670
commit d2b1003612
24 changed files with 1233 additions and 697 deletions
@@ -54,8 +54,9 @@ use crate::disk::{
use crate::erasure::coding::BitrotReader;
use crate::io_support::bitrot::ShardReader;
use crate::io_support::bitrot::{
BitrotReaderStageMetrics, DeferredReaderStripeHandle, adjust_shard_read_params, create_bitrot_reader_with_stage_metrics,
create_deferred_bitrot_reader_with_stripe_handle, object_mmap_read_enabled, object_mmap_read_max_length,
BitrotReaderStageMetrics, DeferredReaderStripeHandle, adjust_shard_read_params,
create_bitrot_reader_from_bytes_with_stage_metrics, create_deferred_bitrot_reader_with_stripe_handle,
object_mmap_read_enabled, object_mmap_read_max_length,
};
use crate::set_disk::shard_source::ShardReadCost;
use futures::FutureExt as _;
@@ -1262,13 +1263,13 @@ pub(in crate::set_disk) fn schedule_bitrot_reader_task<'a>(
return;
}
let inline_data = files[idx].data.as_deref();
let inline_data = files[idx].data.clone();
let data_dir = files[idx].data_dir.unwrap_or_default();
let disk = disks[idx].as_ref();
let path = format!("{object}/{data_dir}/part.{part_number}");
reader_tasks.push(Box::pin(async move {
let result = create_bitrot_reader_with_stage_metrics(
let result = create_bitrot_reader_from_bytes_with_stage_metrics(
inline_data,
disk,
bucket,
@@ -1560,14 +1561,14 @@ pub(in crate::set_disk) async fn create_bitrot_readers_until_quorum_all_shards(
let schedule_stage_start = stage_metrics.map(|_| Instant::now());
for (idx, disk_op) in disks.iter().enumerate() {
setup.mark_scheduled(idx);
let inline_data = files[idx].data.as_deref();
let inline_data = files[idx].data.clone();
let data_dir = files[idx].data_dir.unwrap_or_default();
let disk = disk_op.as_ref();
let path = format!("{object}/{data_dir}/part.{part_number}");
let checksum_algo = checksum_algo.clone();
reader_tasks.push(async move {
let result = create_bitrot_reader_with_stage_metrics(
let result = create_bitrot_reader_from_bytes_with_stage_metrics(
inline_data,
disk,
bucket,
+284 -114
View File
@@ -247,8 +247,8 @@ impl SetDisks {
no_lock: true,
..Default::default()
};
let (current, _, _) = self.get_object_fileinfo(bucket, object, &read_opts, true, false).await?;
restore_operation_id_from_metadata(&current.metadata)?
let current = self.get_object_fileinfo(bucket, object, &read_opts, true, false).await?;
restore_operation_id_from_metadata(&current.fi().metadata)?
.filter(|actual| *actual == expected)
.ok_or_else(|| Error::other(format!("restore operation id changed before {mode}: expected {expected}")))?;
Ok(())
@@ -736,41 +736,91 @@ mod transition_matrix_tests;
pub use ops::heal_walk::HealWalkVersion;
pub(in crate::set_disk) enum GetObjectMetadata<T> {
Owned(T),
Shared(Arc<T>),
pub(in crate::set_disk) struct GetObjectFileInfo {
owned: Option<OwnedGetObjectFileInfo>,
shared: Option<Arc<GetObjectMetadataCacheEntry>>,
}
impl<T> std::ops::Deref for GetObjectMetadata<T> {
type Target = T;
struct OwnedGetObjectFileInfo {
fi: FileInfo,
parts_metadata: Vec<FileInfo>,
online_disks: Vec<Option<DiskStore>>,
}
fn deref(&self) -> &Self::Target {
match self {
Self::Owned(value) => value,
Self::Shared(value) => value,
impl GetObjectFileInfo {
fn owned(fi: FileInfo, parts_metadata: Vec<FileInfo>, online_disks: Vec<Option<DiskStore>>) -> Self {
Self {
owned: Some(OwnedGetObjectFileInfo {
fi,
parts_metadata,
online_disks,
}),
shared: None,
}
}
}
impl<T: Clone> GetObjectMetadata<T> {
fn into_owned(self) -> T {
match self {
Self::Owned(value) => value,
Self::Shared(value) => Arc::try_unwrap(value).unwrap_or_else(|value| (*value).clone()),
fn shared(entry: Arc<GetObjectMetadataCacheEntry>) -> Self {
Self {
owned: None,
shared: Some(entry),
}
}
}
type GetObjectFileInfo = (
GetObjectMetadata<FileInfo>,
GetObjectMetadata<Vec<FileInfo>>,
GetObjectMetadata<Vec<Option<DiskStore>>>,
);
fn fi(&self) -> &FileInfo {
match (&self.owned, &self.shared) {
(Some(snapshot), None) => &snapshot.fi,
(None, Some(entry)) => &entry.fi,
_ => unreachable!("GET metadata snapshot representation must be exclusive"),
}
}
fn parts_metadata(&self) -> &[FileInfo] {
match (&self.owned, &self.shared) {
(Some(snapshot), None) => &snapshot.parts_metadata,
(None, Some(entry)) => &entry.parts_metadata,
_ => unreachable!("GET metadata snapshot representation must be exclusive"),
}
}
fn online_disks(&self) -> &[Option<DiskStore>] {
match (&self.owned, &self.shared) {
(Some(snapshot), None) => &snapshot.online_disks,
(None, Some(entry)) => &entry.online_disks,
_ => unreachable!("GET metadata snapshot representation must be exclusive"),
}
}
fn into_owned(self) -> (FileInfo, Vec<FileInfo>, Vec<Option<DiskStore>>) {
match (self.owned, self.shared) {
(Some(snapshot), None) => {
let OwnedGetObjectFileInfo {
fi,
parts_metadata,
online_disks,
} = snapshot;
(fi, parts_metadata, online_disks)
}
(None, Some(entry)) => match Arc::try_unwrap(entry) {
Ok(entry) => (entry.fi, entry.parts_metadata, entry.online_disks),
Err(entry) => (entry.fi.clone(), entry.parts_metadata.clone(), entry.online_disks.clone()),
},
_ => unreachable!("GET metadata snapshot representation must be exclusive"),
}
}
#[cfg(test)]
fn has_valid_representation(&self) -> bool {
self.owned.is_some() ^ self.shared.is_some()
}
#[cfg(test)]
fn shared_entry(&self) -> Option<&Arc<GetObjectMetadataCacheEntry>> {
self.shared.as_ref()
}
}
pub(crate) struct PreparedGetObjectMetadata {
fi: GetObjectMetadata<FileInfo>,
files: GetObjectMetadata<Vec<FileInfo>>,
disks: GetObjectMetadata<Vec<Option<DiskStore>>>,
snapshot: GetObjectFileInfo,
object_info: Option<ObjectInfo>,
}
@@ -788,7 +838,7 @@ impl PreparedGetObjectMetadata {
}
pub(crate) fn read_semantics_identity(&self) -> [u8; 32] {
SetDisks::file_info_quorum_hash(&self.fi)
SetDisks::file_info_quorum_hash(self.snapshot.fi())
}
}
@@ -838,10 +888,11 @@ mod prepared_get_object_metadata_tests {
#[tokio::test]
async fn prepared_metadata_is_consumed_exactly_once() {
let snapshot = GetObjectFileInfo::owned(FileInfo::default(), Vec::new(), Vec::new());
assert!(snapshot.has_valid_representation());
assert!(snapshot.shared_entry().is_none());
let metadata = PreparedGetObjectMetadata {
fi: GetObjectMetadata::Owned(FileInfo::default()),
files: GetObjectMetadata::Owned(Vec::new()),
disks: GetObjectMetadata::Owned(Vec::new()),
snapshot,
object_info: None,
};
@@ -853,6 +904,44 @@ mod prepared_get_object_metadata_tests {
assert!(take_prepared_get_object_metadata().is_none());
}
#[test]
fn cache_hit_consumers_release_snapshot_at_legacy_ownership() {
let fi = FileInfo {
name: "object".to_owned(),
..Default::default()
};
let cached = Arc::new(GetObjectMetadataCacheEntry {
created_at: Instant::now(),
parts_metadata: vec![fi.clone()],
fi,
online_disks: vec![None],
read_quorum: 0,
});
let snapshot = GetObjectFileInfo::shared(Arc::clone(&cached));
assert!(snapshot.has_valid_representation());
assert!(std::mem::size_of::<GetObjectFileInfo>() >= std::mem::size_of::<OwnedGetObjectFileInfo>());
assert!(
std::mem::size_of::<GetObjectFileInfo>()
<= std::mem::size_of::<OwnedGetObjectFileInfo>() + 2 * std::mem::size_of::<usize>()
);
assert_eq!(Arc::strong_count(&cached), 2, "a cache hit must add one snapshot reference");
assert_eq!(snapshot.fi().name, "object");
assert_eq!(snapshot.parts_metadata().len(), 1);
assert_eq!(snapshot.online_disks().len(), 1);
assert_eq!(Arc::strong_count(&cached), 2, "borrowing consumers must not clone the snapshot");
let (owned_fi, owned_parts, disks) = snapshot.into_owned();
assert_eq!(owned_fi.name, "object");
assert_eq!(owned_parts.len(), 1);
assert_eq!(disks.len(), 1);
assert_eq!(
Arc::strong_count(&cached),
1,
"legacy ownership must release the cache snapshot after cloning its owned inputs"
);
}
#[tokio::test]
#[serial_test::serial(body_cache_hook)]
async fn prepared_reader_reuses_metadata_fanout_exactly_once() {
@@ -1016,12 +1105,10 @@ impl SetDisks {
object: &str,
opts: &ObjectOptions,
) -> Result<PreparedGetObjectMetadata> {
let (fi, files, disks) = self.get_object_fileinfo(bucket, object, opts, true, true).await?;
let object_info = build_get_object_info(&fi, bucket, object, opts.versioned || opts.version_suspended);
let snapshot = self.get_object_fileinfo(bucket, object, opts, true, true).await?;
let object_info = build_get_object_info(snapshot.fi(), bucket, object, opts.versioned || opts.version_suspended);
Ok(PreparedGetObjectMetadata {
fi,
files,
disks,
snapshot,
object_info: Some(object_info),
})
}
@@ -1137,24 +1224,6 @@ pub fn is_deadlock_detection_enabled() -> bool {
// require process restart to take effect.
// ============================================================================
/// Check if codec streaming is enabled (base flag).
///
/// **Note**: Cached via `OnceLock` — env var changes require process restart.
/// In test mode, bypasses cache to allow per-test env var overrides.
fn is_get_codec_streaming_enabled() -> bool {
#[cfg(test)]
{
rustfs_utils::get_env_bool(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, DEFAULT_RUSTFS_GET_CODEC_STREAMING_ENABLE)
}
#[cfg(not(test))]
{
static CACHED: OnceLock<bool> = OnceLock::new();
*CACHED.get_or_init(|| {
rustfs_utils::get_env_bool(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, DEFAULT_RUSTFS_GET_CODEC_STREAMING_ENABLE)
})
}
}
/// Check if multipart codec streaming is enabled.
///
/// When enabled, multipart objects use per-part codec streaming
@@ -1309,22 +1378,6 @@ fn is_multipart_reader_setup_prefetch_enabled() -> bool {
}
}
// --- Rollout Percentage Functions ---
fn get_codec_streaming_rollout_pct() -> u32 {
#[cfg(test)]
{
rustfs_utils::get_env_u32(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT_PCT, DEFAULT_RUSTFS_GET_CODEC_STREAMING_ROLLOUT_PCT)
}
#[cfg(not(test))]
{
static CACHED: OnceLock<u32> = OnceLock::new();
*CACHED.get_or_init(|| {
rustfs_utils::get_env_u32(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT_PCT, DEFAULT_RUSTFS_GET_CODEC_STREAMING_ROLLOUT_PCT)
})
}
}
fn get_metadata_early_stop_rollout_pct() -> u32 {
static CACHED: OnceLock<u32> = OnceLock::new();
*CACHED.get_or_init(|| {
@@ -1359,10 +1412,8 @@ fn is_optimization_enabled_for_request(base_enabled: bool, rollout_pct: u32, buc
(hash as u32) < rollout_pct
}
/// Should this specific request use codec streaming?
pub fn should_use_codec_streaming(bucket: &str, object: &str) -> bool {
let base = is_get_codec_streaming_enabled();
let pct = get_codec_streaming_rollout_pct();
is_optimization_enabled_for_request(base, pct, bucket, object)
fn should_use_codec_streaming(config: GetCodecStreamingConfig, bucket: &str, object: &str) -> bool {
is_optimization_enabled_for_request(config.enabled, config.rollout_pct, bucket, object)
}
/// Should this specific request use metadata early-stop?
@@ -1372,20 +1423,6 @@ pub fn should_use_metadata_early_stop(bucket: &str, object: &str) -> bool {
is_optimization_enabled_for_request(base, pct, bucket, object)
}
fn get_codec_streaming_min_size() -> usize {
if std::env::var_os(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE).is_some() {
return rustfs_utils::get_env_usize(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, DEFAULT_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE);
}
match get_codec_streaming_engine() {
GetCodecStreamingEngine::Rustfs => rustfs_utils::get_env_usize(
ENV_RUSTFS_GET_CODEC_STREAMING_RUSTFS_MIN_SIZE,
DEFAULT_RUSTFS_GET_CODEC_STREAMING_RUSTFS_MIN_SIZE,
),
GetCodecStreamingEngine::Legacy => DEFAULT_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE,
}
}
fn is_get_codec_streaming_data_blocks_first_enabled() -> bool {
#[cfg(test)]
{
@@ -1488,8 +1525,18 @@ enum GetCodecStreamingEngine {
Rustfs,
}
fn get_codec_streaming_engine() -> GetCodecStreamingEngine {
let engine = rustfs_utils::get_env_str(ENV_RUSTFS_GET_CODEC_STREAMING_ENGINE, DEFAULT_RUSTFS_GET_CODEC_STREAMING_ENGINE);
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
struct GetCodecStreamingConfig {
enabled: bool,
rollout: GetCodecStreamingRollout,
rollout_pct: u32,
body_compat_confirmed: bool,
header_compat_confirmed: bool,
engine: GetCodecStreamingEngine,
min_size: usize,
}
fn parse_get_codec_streaming_engine(engine: &str) -> GetCodecStreamingEngine {
match engine.trim() {
value if value.eq_ignore_ascii_case(GET_CODEC_STREAMING_ENGINE_RUSTFS) => GetCodecStreamingEngine::Rustfs,
value if value.eq_ignore_ascii_case(GET_CODEC_STREAMING_ENGINE_LEGACY) => GetCodecStreamingEngine::Legacy,
@@ -1497,8 +1544,7 @@ fn get_codec_streaming_engine() -> GetCodecStreamingEngine {
}
}
fn get_codec_streaming_rollout() -> GetCodecStreamingRollout {
let rollout = rustfs_utils::get_env_str(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, DEFAULT_RUSTFS_GET_CODEC_STREAMING_ROLLOUT);
fn parse_get_codec_streaming_rollout(rollout: &str) -> GetCodecStreamingRollout {
match rollout.trim() {
// Clean production token. `internal`/`benchmark` remain accepted aliases
// for backward compatibility; all three opt the fast path in.
@@ -1515,20 +1561,58 @@ fn get_codec_streaming_rollout() -> GetCodecStreamingRollout {
}
}
/// Emergency kill-switch (defaults to `true`). Set
/// `RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED=false` to force the fast path
/// off. Body compatibility is confirmed by the parity e2e net + bench (backlog#1183),
/// so this no longer gates enablement — the `..._ROLLOUT` switch does.
fn is_get_codec_streaming_body_compat_confirmed() -> bool {
rustfs_utils::get_env_bool(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, true)
fn load_get_codec_streaming_config() -> GetCodecStreamingConfig {
let engine = parse_get_codec_streaming_engine(&rustfs_utils::get_env_str(
ENV_RUSTFS_GET_CODEC_STREAMING_ENGINE,
DEFAULT_RUSTFS_GET_CODEC_STREAMING_ENGINE,
));
let min_size = if std::env::var_os(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE).is_some() {
rustfs_utils::get_env_usize(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, DEFAULT_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE)
} else {
match engine {
GetCodecStreamingEngine::Rustfs => rustfs_utils::get_env_usize(
ENV_RUSTFS_GET_CODEC_STREAMING_RUSTFS_MIN_SIZE,
DEFAULT_RUSTFS_GET_CODEC_STREAMING_RUSTFS_MIN_SIZE,
),
GetCodecStreamingEngine::Legacy => DEFAULT_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE,
}
};
GetCodecStreamingConfig {
enabled: rustfs_utils::get_env_bool(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, DEFAULT_RUSTFS_GET_CODEC_STREAMING_ENABLE),
rollout: parse_get_codec_streaming_rollout(&rustfs_utils::get_env_str(
ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT,
DEFAULT_RUSTFS_GET_CODEC_STREAMING_ROLLOUT,
)),
rollout_pct: rustfs_utils::get_env_u32(
ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT_PCT,
DEFAULT_RUSTFS_GET_CODEC_STREAMING_ROLLOUT_PCT,
),
body_compat_confirmed: rustfs_utils::get_env_bool(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, true),
header_compat_confirmed: rustfs_utils::get_env_bool(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, true),
engine,
min_size,
}
}
/// Emergency kill-switch (defaults to `true`). Set
/// `RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED=false` to force the fast path
/// off. Header compatibility is confirmed by the parity e2e net + bench (backlog#1183),
/// so this no longer gates enablement — the `..._ROLLOUT` switch does.
fn is_get_codec_streaming_header_compat_confirmed() -> bool {
rustfs_utils::get_env_bool(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, true)
fn get_codec_streaming_config_cached_core(load: impl FnOnce() -> GetCodecStreamingConfig) -> GetCodecStreamingConfig {
static CACHED: OnceLock<GetCodecStreamingConfig> = OnceLock::new();
*CACHED.get_or_init(load)
}
fn get_codec_streaming_config() -> GetCodecStreamingConfig {
#[cfg(test)]
{
load_get_codec_streaming_config()
}
#[cfg(not(test))]
{
get_codec_streaming_config_cached_core(load_get_codec_streaming_config)
}
}
fn get_codec_streaming_engine() -> GetCodecStreamingEngine {
get_codec_streaming_config().engine
}
fn build_get_codec_streaming_decode_engine(erasure: coding::Erasure) -> std::io::Result<CodecStreamingDecodeEngine> {
@@ -1897,43 +1981,43 @@ fn should_prefer_codec_streaming_data_blocks_first_reader_setup(
fn get_codec_streaming_reader_gate(
bucket: &str,
object: &str,
range: &Option<HTTPRangeSpec>,
part_number: Option<usize>,
object_class: GetCodecStreamingObjectClass,
object_info: &ObjectInfo,
fi: &FileInfo,
lock_optimization_enabled: bool,
) -> GetCodecStreamingGate {
let object_class = classify_get_codec_streaming_object_class(range, object_info, fi);
let config = get_codec_streaming_config();
if !is_get_codec_streaming_enabled() {
if !config.enabled {
return GetCodecStreamingGate {
object_class,
decision: GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::Disabled),
prefer_data_blocks_first_reader_setup: false,
};
}
if !get_codec_streaming_rollout().is_opted_in() {
if !config.rollout.is_opted_in() {
return GetCodecStreamingGate {
object_class,
decision: GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::RolloutNotOptedIn),
prefer_data_blocks_first_reader_setup: false,
};
}
if !should_use_codec_streaming(bucket, object) {
if !should_use_codec_streaming(config, bucket, object) {
return GetCodecStreamingGate {
object_class,
decision: GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::RolloutPctNotSelected),
prefer_data_blocks_first_reader_setup: false,
};
}
if !is_get_codec_streaming_body_compat_confirmed() {
if !config.body_compat_confirmed {
return GetCodecStreamingGate {
object_class,
decision: GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::BodyCompatibilityUnconfirmed),
prefer_data_blocks_first_reader_setup: false,
};
}
if !is_get_codec_streaming_header_compat_confirmed() {
if !config.header_compat_confirmed {
return GetCodecStreamingGate {
object_class,
decision: GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::HeaderCompatibilityUnconfirmed),
@@ -2006,7 +2090,7 @@ fn get_codec_streaming_reader_gate(
};
}
}
let Ok(min_size) = i64::try_from(get_codec_streaming_min_size()) else {
let Ok(min_size) = i64::try_from(config.min_size) else {
return GetCodecStreamingGate {
object_class,
decision: GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::InvalidMinSize),
@@ -2376,6 +2460,7 @@ pub struct SetDisks {
get_object_metadata_cache_hash_builder: std::collections::hash_map::RandomState,
get_object_metadata_cache_generations: Arc<[AtomicU64]>,
pub lockers: Vec<Arc<dyn LockClient>>,
shared_lockers: Arc<[Arc<dyn LockClient>]>,
local_lock_manager: Arc<rustfs_lock::GlobalLockManager>,
/// Per-instance runtime context (Phase 5, backlog#939).
///
@@ -2503,9 +2588,9 @@ impl Hash for GetObjectMetadataCacheKey {
struct GetObjectMetadataCacheEntry {
#[allow(dead_code)] // Kept for debugging; moka handles TTL internally
created_at: Instant,
fi: Arc<FileInfo>,
parts_metadata: Arc<Vec<FileInfo>>,
online_disks: Arc<Vec<Option<DiskStore>>>,
fi: FileInfo,
parts_metadata: Vec<FileInfo>,
online_disks: Vec<Option<DiskStore>>,
read_quorum: usize,
}
@@ -2772,6 +2857,7 @@ impl SetDisks {
) -> Arc<Self> {
let ctx = instance_ctx;
let set_lock_namespace: Arc<str> = format!("set-{pool_index}-{set_index}").into();
let shared_lockers = Arc::from(lockers.to_vec());
Arc::new(SetDisks {
locker_owner,
disks,
@@ -2794,6 +2880,7 @@ impl SetDisks {
.collect::<Vec<_>>(),
),
lockers,
shared_lockers,
// Sourced from the instance context so each instance owns its lock
// namespace (Phase 5 Slice 3). Single-instance: ctx aliases the
// process lock-manager singleton, so this is unchanged.
@@ -4999,6 +5086,89 @@ mod tests {
assert_eq!(Arc::strong_count(&set.set_lock_namespace), before);
}
#[tokio::test]
async fn new_ns_lock_shares_clients_without_changing_quorum() {
let healthy: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(Arc::new(rustfs_lock::GlobalLockManager::new())));
let failing: Arc<dyn LockClient> = Arc::new(FailingClient);
let ctx = Arc::new(InstanceContext::new());
ctx.update_erasure_type(SetupType::DistErasure).await;
let set = make_test_set_disks_with_ctx(vec![healthy.clone(), failing.clone()], ctx).await;
assert!(Arc::ptr_eq(&set.lockers[0], &healthy));
assert!(Arc::ptr_eq(&set.lockers[1], &failing));
let clients_before = Arc::strong_count(&set.shared_lockers);
let healthy_before = Arc::strong_count(&healthy);
let failing_before = Arc::strong_count(&failing);
let write_lock = set
.new_ns_lock("bucket", "write-object")
.await
.expect("namespace lock should be created");
assert_eq!(
Arc::strong_count(&set.shared_lockers),
clients_before + 1,
"each object lock should share one client slice allocation"
);
assert_eq!(
Arc::strong_count(&healthy),
healthy_before,
"constructing an object lock must not clone each client Arc"
);
assert_eq!(
Arc::strong_count(&failing),
failing_before,
"constructing an object lock must not clone each client Arc"
);
let write_error = write_lock
.get_write_lock(Duration::from_millis(500))
.await
.expect_err("one healthy client must not satisfy the two-client write quorum");
assert!(
matches!(
write_error,
LockError::QuorumNotReached {
required: 2,
achieved: 1
}
),
"the shared client representation must preserve the exact write quorum result: {write_error}"
);
let read_lock = set
.new_ns_lock("bucket", "read-object")
.await
.expect("second namespace lock should be created");
assert_eq!(Arc::strong_count(&set.shared_lockers), clients_before + 2);
let read_guard = read_lock
.get_read_lock(Duration::from_millis(500))
.await
.expect("one healthy client should satisfy the two-client read quorum");
assert!(matches!(read_guard, NamespaceLockGuard::Standard(_)));
}
#[tokio::test]
async fn new_ns_lock_uses_the_current_public_client_domain() {
let stale_a: Arc<dyn LockClient> = Arc::new(FailingClient);
let stale_b: Arc<dyn LockClient> = Arc::new(FailingClient);
let healthy_a: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(Arc::new(rustfs_lock::GlobalLockManager::new())));
let healthy_b: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(Arc::new(rustfs_lock::GlobalLockManager::new())));
let ctx = Arc::new(InstanceContext::new());
ctx.update_erasure_type(SetupType::DistErasure).await;
let set = make_test_set_disks_with_ctx(vec![stale_a, stale_b], ctx).await;
let mut set = (*set).clone();
set.lockers = vec![healthy_a, healthy_b];
let lock = set
.new_ns_lock("bucket", "object")
.await
.expect("namespace lock should use the current public clients");
let guard = lock
.get_write_lock(Duration::from_millis(500))
.await
.expect("the current healthy clients should satisfy the two-client quorum");
assert!(matches!(guard, NamespaceLockGuard::Standard(_)));
}
struct SetupTypeGuard {
previous: SetupType,
}
+2 -1
View File
@@ -3171,10 +3171,11 @@ mod heal_result_report_tests {
.await
.expect("object should be written");
let (fi, _, _) = set
let snapshot = set
.get_object_fileinfo(bucket, object, &opts, true, false)
.await
.expect("object metadata should resolve");
let fi = snapshot.fi();
assert_eq!(fi.erasure.parity_blocks, 0);
let data_dir = fi.data_dir.expect("non-inline object should have a data directory");
let part_path = dir.path().join(bucket).join(object).join(data_dir.to_string()).join("part.1");
+14 -3
View File
@@ -36,10 +36,21 @@ impl crate::storage_api_contracts::namespace::NamespaceLocking for SetDisks {
// test's transient DistErasure window) would push this set's namespace
// locking onto its own — possibly empty — dist locker list.
let set_lock = if self.ctx.is_dist_erasure().await {
// Calculate quorum based on lockers count (majority)
let lockers_count = self.lockers.len();
let lockers = if self.lockers.len() == self.shared_lockers.len()
&& self
.lockers
.iter()
.zip(self.shared_lockers.iter())
.all(|(current, shared)| Arc::ptr_eq(current, shared))
{
self.shared_lockers.clone()
} else {
Arc::from(self.lockers.clone())
};
// Calculate quorum from the exact client domain used by this lock.
let lockers_count = lockers.len();
let write_quorum = if lockers_count > 1 { (lockers_count / 2) + 1 } else { 1 };
NamespaceLock::with_clients_and_quorum_shared(self.set_lock_namespace.clone(), self.lockers.clone(), write_quorum)
NamespaceLock::with_clients_and_quorum_shared(self.set_lock_namespace.clone(), lockers, write_quorum)
} else {
NamespaceLock::with_local_manager_shared(self.set_lock_namespace.clone(), self.local_lock_manager.clone())
};
+145 -71
View File
@@ -379,7 +379,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
type GetObjectReader = GetObjectReader;
type PutObjectReader = PutObjReader;
#[tracing::instrument(level = "debug", skip(self))]
#[tracing::instrument(level = "debug", skip(self, h))]
async fn get_object_reader(
&self,
bucket: &str,
@@ -428,11 +428,11 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
};
let metadata_stage_start = Instant::now();
let (fi, files, disks, prepared_object_info) = if let Some(prepared) = take_prepared_get_object_metadata() {
(prepared.fi, prepared.files, prepared.disks, prepared.object_info)
let (snapshot, prepared_object_info) = if let Some(prepared) = take_prepared_get_object_metadata() {
(prepared.snapshot, prepared.object_info)
} else {
match self.get_object_fileinfo(bucket, object, opts, true, true).await {
Ok((fi, files, disks)) => (fi, files, disks, None),
Ok(snapshot) => (snapshot, None),
Err(err) => {
rustfs_io_metrics::record_get_object_metadata_phase_duration(metadata_stage_start.elapsed().as_secs_f64());
let failure_path = if is_meta_bucketname(bucket) {
@@ -445,10 +445,13 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
}
}
};
let fi = snapshot.fi();
let files = snapshot.parts_metadata();
let disks = snapshot.online_disks();
let object_info_stage_start = get_stage_timer_if_enabled(stage_metrics_enabled);
let object_info = prepared_object_info
.unwrap_or_else(|| build_get_object_info(&fi, bucket, object, opts.versioned || opts.version_suspended));
let object_class = classify_get_codec_streaming_object_class(&range, &object_info, &fi);
.unwrap_or_else(|| build_get_object_info(fi, bucket, object, opts.versioned || opts.version_suspended));
let object_class = classify_get_codec_streaming_object_class(&range, &object_info, fi);
let size_bucket = rustfs_io_metrics::get_object_size_bucket(object_info.size);
record_get_stage_duration_if_enabled(GET_OBJECT_PATH_SET_DISK, GET_STAGE_OBJECT_INFO, object_info_stage_start);
let metadata_elapsed = metadata_stage_start.elapsed().as_secs_f64();
@@ -496,7 +499,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
// Uses the shared predicate from ObjectInfo; additionally checks that
// inline data is actually present and neither range nor partNumber is
// in flight.
if should_use_inline_fast_path(&range, &object_info, &fi, opts) {
if should_use_inline_fast_path(&range, &object_info, fi, opts) {
let mut inline_prepare_stage_start = get_stage_timer_if_enabled(stage_metrics_enabled);
let data_shards = fi.erasure.data_blocks;
@@ -512,7 +515,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
};
if can_try_inline_data_shards_direct(object_size, fi.erasure.block_size)
&& let Some(data_files) = collect_inline_data_shard_fileinfos_by_index(&files, &fi, data_shards, |index| {
&& let Some(data_files) = collect_inline_data_shard_fileinfos_by_index(files, fi, data_shards, |index| {
disks.get(index).is_some_and(Option::is_some)
})
{
@@ -580,10 +583,10 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
}
}
let erasure = erasure_from_file_info(&fi, fi.uses_legacy_checksum)?;
let erasure = erasure_from_file_info(fi, fi.uses_legacy_checksum)?;
let read_length = erasure.shard_file_offset(0, object_size, object_size);
let total_shards = data_shards + fi.erasure.parity_blocks;
let (_disks, files) = Self::shuffle_disks_and_parts_metadata_by_index(&disks, &files, &fi);
let (_disks, files) = Self::shuffle_disks_and_parts_metadata_by_index(disks, files, fi);
// Check if we have enough inline data shards
let inline_count = files
@@ -662,10 +665,10 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
let codec_streaming_gate = get_codec_streaming_reader_gate(
bucket,
object,
&range,
opts.part_number,
object_class,
&object_info,
&fi,
fi,
lock_optimization_enabled,
);
record_get_stage_duration_if_enabled(GET_OBJECT_PATH_SET_DISK, GET_STAGE_PATH_DECISION, path_decision_stage_start);
@@ -743,15 +746,15 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
}
}
let direct_memory_decision = get_small_object_direct_memory_decision(&range, &object_info, &fi, opts);
let direct_memory_decision = get_small_object_direct_memory_decision(&range, &object_info, fi, opts);
record_get_direct_memory_decision(object_class, direct_memory_decision, size_bucket);
if let GetDirectMemoryDecision::Use { object_size } = direct_memory_decision {
if let Some(body) = Self::try_get_object_direct_data_shards_with_fileinfo(
bucket,
object,
&fi,
&files,
&disks,
fi,
files,
disks,
opts.skip_verify_bitrot,
object_class.as_str(),
size_bucket,
@@ -780,14 +783,15 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
}
let mut output = Vec::with_capacity(object_size);
let (fi, files, disks) = snapshot.into_owned();
Self::get_object_with_fileinfo(
bucket,
object,
0,
object_info.size,
&mut output,
fi.into_owned(),
files.into_owned(),
fi,
files,
&disks,
self.set_index,
self.pool_index,
@@ -826,9 +830,9 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
match Self::get_object_decode_reader_with_fileinfo(
bucket,
object,
&fi,
&files,
&disks,
fi,
files,
disks,
self.set_index,
self.pool_index,
opts.skip_verify_bitrot,
@@ -890,6 +894,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
let set_index = self.set_index;
let pool_index = self.pool_index;
let skip_verify = opts.skip_verify_bitrot;
let (fi, files, disks) = snapshot.into_owned();
tokio::spawn(async move {
let _guard = read_lock_guard;
let mut writer = GetObjectDownstreamWriter::new(wd);
@@ -903,8 +908,8 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
offset,
length,
&mut writer,
fi.into_owned(),
files.into_owned(),
fi,
files,
&disks,
set_index,
pool_index,
@@ -3294,10 +3299,10 @@ impl SetDisks {
// quorum, failing write quorum on update_object_meta (backlog#872).
let mut read_opts = opts.clone();
read_opts.include_part_checksums = true;
let (fi, _, disks) = self
let (mut fi, _, disks) = self
.get_object_fileinfo_gated(bucket, object, &read_opts, false, false)
.await?;
let mut fi = fi.into_owned();
.await?
.into_owned();
fi.metadata.insert(AMZ_OBJECT_TAGGING.to_owned(), tags.to_owned());
if let Some(eval_metadata) = &opts.eval_metadata {
@@ -4575,12 +4580,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
// Use the same full xl.meta read path as GetObject metadata resolution.
// This avoids HEAD/GetObject metadata visibility skew immediately after
// PutObject/CompleteMultipartUpload.
let (fi, _, _) = self
let snapshot = self
.get_object_fileinfo(bucket, object, opts, true, false)
.await
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
let oi = ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended);
let oi = ObjectInfo::from_file_info(snapshot.fi(), bucket, object, opts.versioned || opts.version_suspended);
Ok(oi)
}
@@ -4748,10 +4753,10 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
let mut transition_read_opts = opts.clone();
transition_read_opts.include_part_checksums = true;
let (fi, meta_arr, online_disks) = self
let (mut fi, meta_arr, online_disks) = self
.get_object_fileinfo(bucket, object, &transition_read_opts, true, false)
.await?;
let mut fi = fi.into_owned();
.await?
.into_owned();
/*if err != nil {
return Err(to_object_err(err, vec![bucket, object]));
}*/
@@ -4866,7 +4871,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
cloned_fi.size,
&mut writer,
cloned_fi,
meta_arr.into_owned(),
meta_arr,
&online_disks,
set_index,
pool_index,
@@ -4991,7 +4996,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
};
self.invalidate_get_object_metadata_cache(bucket, object).await;
let current = self.get_object_fileinfo(bucket, object, &commit_opts, true, false).await;
let (current_fi, _, _) = match current {
let current = match current {
Ok(current) => current,
Err(err) => {
drop(transition_lock_guard);
@@ -5002,7 +5007,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
return Err(err);
}
};
let mut current_fi = current_fi.into_owned();
let (mut current_fi, _, _) = current.into_owned();
let source_matches = current_fi.version_id == fi.version_id
&& current_fi.data_dir == fi.data_dir
&& current_fi.mod_time == fi.mod_time
@@ -5233,9 +5238,10 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
if let Err(err) = fi {
return set_restore_header_fn(&mut oi, Some(to_object_err(err, vec![bucket, object]))).await;
}
let (actual_fi, _, _) = fi?;
let actual = fi?;
let actual_fi = actual.fi();
oi = ObjectInfo::from_file_info(&actual_fi, bucket, object, opts.versioned || opts.version_suspended);
oi = ObjectInfo::from_file_info(actual_fi, bucket, object, opts.versioned || opts.version_suspended);
let expected_operation_id = restore_operation_id_from_metadata(&opts.user_defined)?;
if let Some(expected_operation_id) = expected_operation_id {
require_restore_operation_id(oi.user_defined.as_ref(), expected_operation_id)?;
@@ -5513,8 +5519,40 @@ mod object_encryption_resolver_wiring_tests {
use super::*;
use crate::object_api::{EncryptionResolutionError, ObjectEncryptionResolver, ReadEncryptionMaterial, ReadEncryptionRequest};
use std::io::Cursor;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Clone, Default)]
struct CapturedLogs(Arc<Mutex<Vec<u8>>>);
struct CapturedLogWriter(Arc<Mutex<Vec<u8>>>);
impl CapturedLogs {
fn contents(&self) -> String {
String::from_utf8(self.0.lock().expect("captured logs mutex should not poison").clone())
.expect("captured logs should be valid UTF-8")
}
}
impl std::io::Write for CapturedLogWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
let mut captured = self.0.lock().expect("captured logs mutex should not poison");
std::io::Write::write(&mut *captured, buf)
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
impl<'writer> tracing_subscriber::fmt::MakeWriter<'writer> for CapturedLogs {
type Writer = CapturedLogWriter;
fn make_writer(&'writer self) -> Self::Writer {
CapturedLogWriter(Arc::clone(&self.0))
}
}
struct CountingResolver {
calls: AtomicUsize,
}
@@ -5561,6 +5599,34 @@ mod object_encryption_resolver_wiring_tests {
assert!(result.is_err(), "resolver returning no material must fail closed");
assert_eq!(resolver.calls.load(Ordering::Relaxed), 1);
}
#[tokio::test(flavor = "current_thread")]
async fn get_object_reader_span_never_records_transport_headers() {
use super::hermetic_set_disks_support::hermetic_set_disks_isolated;
use crate::storage_api_contracts::object::ObjectIO as _;
use rustfs_utils::http::headers::SSEC_KEY_HEADER;
let logs = CapturedLogs::default();
let subscriber = tracing_subscriber::fmt()
.with_max_level(tracing::Level::DEBUG)
.with_writer(logs.clone())
.with_ansi(false)
.without_time()
.finish();
let _guard = tracing::subscriber::set_default(subscriber);
let (_temp_dirs, _disks, set_disks) = hermetic_set_disks_isolated(4).await;
let mut headers = HeaderMap::new();
headers.insert(http::header::AUTHORIZATION, HeaderValue::from_static("credential-must-not-be-logged"));
headers.insert(SSEC_KEY_HEADER, HeaderValue::from_static("customer-key-must-not-be-logged"));
let _ = set_disks
.get_object_reader("missing-bucket", "missing-object", None, headers, &ObjectOptions::default())
.await;
let captured = logs.contents();
assert!(!captured.contains("credential-must-not-be-logged"));
assert!(!captured.contains("customer-key-must-not-be-logged"));
}
}
#[cfg(test)]
@@ -6200,9 +6266,9 @@ mod metadata_mutation_generation_tests {
false,
)
.await
.expect("object metadata should be readable before adding the checksum sidecar");
let mut fi = fi.into_owned();
let disks = disks.into_owned();
.expect("object metadata should be readable before adding the checksum sidecar")
.into_owned();
let mut fi = fi;
rustfs_utils::http::insert_str(&mut fi.metadata, rustfs_utils::http::SUFFIX_PART_CHECKSUMS, value.to_string());
set_disks
.update_object_meta(bucket, object, fi, &disks)
@@ -6359,9 +6425,9 @@ mod metadata_mutation_generation_tests {
false,
)
.await
.expect("conflicting object metadata should be readable before corruption is injected");
let mut fi = fi.into_owned();
let disks = disks.into_owned();
.expect("conflicting object metadata should be readable before corruption is injected")
.into_owned();
let mut fi = fi;
let rustfs_key = format!(
"{}{}",
rustfs_utils::http::RUSTFS_INTERNAL_PREFIX,
@@ -6524,9 +6590,9 @@ mod transition_commit_failure_tests {
false,
)
.await
.expect("source metadata should be readable before adding the checksum sidecar");
let mut source_fi = source_fi.into_owned();
let online_disks = online_disks.into_owned();
.expect("source metadata should be readable before adding the checksum sidecar")
.into_owned();
let mut source_fi = source_fi;
rustfs_utils::http::insert_str(
&mut source_fi.metadata,
rustfs_utils::http::SUFFIX_PART_CHECKSUMS,
@@ -6800,10 +6866,11 @@ mod transition_commit_failure_tests {
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.await
.expect("source object should be written");
let (fi, parts_metadata, online_disks) = set_disks
let snapshot = set_disks
.get_object_fileinfo(bucket, object, &ObjectOptions::default(), true, false)
.await
.expect("source metadata should resolve");
let (fi, parts_metadata, online_disks) = snapshot.into_owned();
let generation = set_disks
.get_object_metadata_cache_generation(bucket, object)
.expect("metadata cache generation should be active");
@@ -6814,9 +6881,9 @@ mod transition_commit_failure_tests {
cache_key.clone(),
Arc::new(GetObjectMetadataCacheEntry {
created_at: Instant::now(),
fi: Arc::new((*fi).clone()),
parts_metadata: Arc::new(parts_metadata.into_owned()),
online_disks: Arc::new(online_disks.into_owned()),
fi,
parts_metadata,
online_disks,
read_quorum: 2,
}),
)
@@ -6943,7 +7010,7 @@ mod transition_commit_failure_tests {
let result = transition.await.expect("transition task should not panic");
result.expect_err("partial local commit must fail when only two of four disks are writable");
let fi = set_disks
let snapshot = set_disks
.get_object_fileinfo(
bucket,
object,
@@ -6955,10 +7022,10 @@ mod transition_commit_failure_tests {
false,
)
.await
.expect("rollback should keep the source metadata readable on applied disks")
.0;
.expect("rollback should keep the source metadata readable on applied disks");
assert_ne!(
fi.transition_status, TRANSITION_COMPLETE,
snapshot.fi().transition_status,
TRANSITION_COMPLETE,
"rollback must not leave the applied disks marked as transitioned"
);
let mut restored = Vec::new();
@@ -8093,7 +8160,7 @@ mod transition_upload_integrity_tests {
.await
.expect("local source should drain after failed transition");
assert_eq!(restored, payload);
let (fi, _, _) = set_disks
let snapshot = set_disks
.get_object_fileinfo(
bucket,
object,
@@ -8107,7 +8174,7 @@ mod transition_upload_integrity_tests {
)
.await
.expect("local source metadata should remain available");
assert_ne!(fi.transition_status, TRANSITION_COMPLETE);
assert_ne!(snapshot.fi().transition_status, TRANSITION_COMPLETE);
}
async fn write_source(
@@ -8204,7 +8271,8 @@ mod transition_upload_integrity_tests {
false,
)
.await
.expect("the existing target should remain readable");
.expect("the existing target should remain readable")
.into_owned();
assert_ne!(stored.transition_status, TRANSITION_COMPLETE);
assert!(stored.transition_version.is_none());
}
@@ -8276,9 +8344,9 @@ mod transition_upload_integrity_tests {
false,
)
.await
.expect("source metadata should be readable");
let mut source_fi = source_fi.into_owned();
let online_disks = online_disks.into_owned();
.expect("source metadata should be readable")
.into_owned();
let mut source_fi = source_fi;
rustfs_utils::http::insert_str(
&mut source_fi.metadata,
rustfs_utils::http::SUFFIX_PART_CHECKSUMS,
@@ -8305,7 +8373,7 @@ mod transition_upload_integrity_tests {
remote_object.starts_with(crate::bucket::lifecycle::transition_transaction::TRANSITION_TRANSACTION_PREFIX),
"remote object should be transaction-scoped: {remote_object}"
);
let (fi, _, _) = set_disks
let snapshot = set_disks
.get_object_fileinfo(
bucket,
object,
@@ -8320,6 +8388,7 @@ mod transition_upload_integrity_tests {
)
.await
.expect("committed transition metadata should be readable");
let fi = snapshot.fi();
assert_eq!(fi.transition_status, TRANSITION_COMPLETE);
assert_eq!(fi.transitioned_objname, *remote_object);
assert_eq!(
@@ -8754,7 +8823,7 @@ mod transition_upload_integrity_tests {
let object = format!("{}-corrupt.bin", position.label());
let payload = vec![0x41; 2 * 1024 * 1024];
let original = write_source(&set_disks, &disk_stores, &bucket, &object, &payload).await;
let (source, _, _) = set_disks
let source = set_disks
.get_object_fileinfo(
&bucket,
&object,
@@ -8768,6 +8837,7 @@ mod transition_upload_integrity_tests {
)
.await
.expect("source metadata should be available before shard corruption");
let source = source.fi();
let data_dir = source.data_dir.expect("source object should have a data directory");
corrupt_beyond_read_quorum(&temp_dirs, &bucket, &object, data_dir, source.erasure.parity_blocks, position).await;
@@ -8787,7 +8857,7 @@ mod transition_upload_integrity_tests {
),
"{position:?}: unexpected transition producer error: {error:?}"
);
let (after, _, _) = set_disks
let after = set_disks
.get_object_fileinfo(
&bucket,
&object,
@@ -8801,6 +8871,7 @@ mod transition_upload_integrity_tests {
)
.await
.expect("failed transition must leave metadata readable");
let after = after.fi();
assert_eq!(
after.data_dir,
Some(data_dir),
@@ -8847,7 +8918,7 @@ mod transition_upload_integrity_tests {
.transition_object(bucket, object, &transition_options(&original, tier_name))
.await
.expect("an unversioned remote version must commit");
let (fi, _, _) = set_disks
let snapshot = set_disks
.get_object_fileinfo(
bucket,
object,
@@ -8861,6 +8932,7 @@ mod transition_upload_integrity_tests {
)
.await
.expect("committed unversioned transition metadata should be readable");
let fi = snapshot.fi();
assert_eq!(fi.transition_version_id, None);
assert_eq!(fi.transition_version, None);
assert_eq!(fi.transition_version_state, rustfs_filemeta::TransitionVersionState::KnownDisabled);
@@ -9068,7 +9140,7 @@ mod transition_upload_integrity_tests {
matches!(error, StorageError::NamespaceLockQuorumUnavailable { .. }),
"unexpected tagging lock-lost error: {error:?}"
);
let (fi, _, _) = set_disks
let snapshot = set_disks
.get_object_fileinfo(
bucket,
object,
@@ -9083,7 +9155,7 @@ mod transition_upload_integrity_tests {
.await
.expect("source metadata should remain readable");
assert!(
!fi.metadata.contains_key(AMZ_OBJECT_TAGGING),
!snapshot.fi().metadata.contains_key(AMZ_OBJECT_TAGGING),
"a stale tagging writer must not write metadata after refresh-quorum loss"
);
}
@@ -9279,13 +9351,14 @@ mod transition_source_identity_matrix_tests {
.put_object(bucket, &object, &mut reader, &source_opts)
.await
.expect("source object should be written");
let (source, _, _) = set_disks
let source = set_disks
.get_object_fileinfo(bucket, &object, &source_opts, true, false)
.await
.expect("source metadata should resolve");
let source = source.fi();
assert_eq!(source.version_id, Some(source_version_id));
assert_eq!(
transition_source_identity(bucket, &object, &source, &source_opts, &get_raw_etag(&source.metadata))
transition_source_identity(bucket, &object, source, &source_opts, &get_raw_etag(&source.metadata))
.expect("persisted versioned source identity should build")
.version_mode,
TransitionSourceVersionMode::Versioned
@@ -9333,10 +9406,11 @@ mod transition_source_identity_matrix_tests {
versioned: true,
..Default::default()
};
let (persisted, _, _) = set_disks
let persisted = set_disks
.get_object_fileinfo(bucket, &object, &persisted_opts, true, false)
.await
.expect("drifted source metadata should resolve");
let persisted = persisted.fi();
put_barrier.release();
let result = transition.await.expect("transition task should not panic");
@@ -10356,9 +10430,9 @@ mod put_object_tags_early_stop_regression_tests {
false,
)
.await
.expect("object metadata should be readable before adding the checksum sidecar");
let mut fi = fi.into_owned();
let disks = disks.into_owned();
.expect("object metadata should be readable before adding the checksum sidecar")
.into_owned();
let mut fi = fi;
rustfs_utils::http::insert_str(
&mut fi.metadata,
rustfs_utils::http::SUFFIX_PART_CHECKSUMS,
+191 -42
View File
@@ -180,9 +180,9 @@ impl SetDisks {
let key = GetObjectMetadataCacheKey::new(bucket, object, generation);
let entry = Arc::new(GetObjectMetadataCacheEntry {
created_at: Instant::now(),
fi: Arc::new(fi.clone()),
parts_metadata: Arc::new(parts_metadata.to_vec()),
online_disks: Arc::new(online_disks.to_vec()),
fi: fi.clone(),
parts_metadata: parts_metadata.to_vec(),
online_disks: online_disks.to_vec(),
read_quorum,
});
self.insert_get_object_metadata_cache_entry_after_insert(key, generation, entry, || {})
@@ -300,11 +300,7 @@ impl SetDisks {
GET_STAGE_METADATA_CACHE_LOOKUP,
metadata_cache_lookup_start,
);
return Ok((
GetObjectMetadata::Shared(Arc::clone(&cached.fi)),
GetObjectMetadata::Shared(Arc::clone(&cached.parts_metadata)),
GetObjectMetadata::Shared(Arc::clone(&cached.online_disks)),
));
return Ok(GetObjectFileInfo::shared(cached));
}
MetadataCacheLookup::Miss => {
rustfs_io_metrics::record_get_object_metadata_cache_decision(
@@ -436,11 +432,7 @@ impl SetDisks {
// let online_disks: Vec<Option<DiskStore>> = op_online_disks.iter().filter(|v| v.is_some()).cloned().collect();
Ok((
GetObjectMetadata::Owned(fi),
GetObjectMetadata::Owned(parts_metadata),
GetObjectMetadata::Owned(op_online_disks),
))
Ok(GetObjectFileInfo::owned(fi, parts_metadata, op_online_disks))
}
#[hotpath::measure(impl_type = "SetDisks")]
@@ -450,14 +442,15 @@ impl SetDisks {
object: &str,
opts: &ObjectOptions,
) -> (ObjectInfo, usize, Option<StorageError>) {
let fi = match self.get_object_fileinfo(bucket, object, opts, false, false).await {
Ok((fi, _, _)) => fi,
let snapshot = match self.get_object_fileinfo(bucket, object, opts, false, false).await {
Ok(snapshot) => snapshot,
Err(e) => return (ObjectInfo::default(), 0, Some(e)),
};
let fi = snapshot.fi();
let write_quorum = fi.write_quorum(self.default_write_quorum());
let oi = ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended);
let oi = ObjectInfo::from_file_info(fi, bucket, object, opts.versioned || opts.version_suspended);
if !fi.version_purge_status().is_empty() && opts.version_id.is_some() {
return (
@@ -2723,22 +2716,14 @@ mod metadata_cache_tests {
.await
.expect("fresh cache entry should be returned");
let (returned_fi, returned_parts_metadata, returned_online_disks) = set
let returned = set
.get_object_fileinfo("bucket", "object", &ObjectOptions::default(), true, false)
.await
.expect("cache-backed metadata lookup should succeed");
assert!(
matches!(returned_fi, GetObjectMetadata::Shared(ref value) if Arc::ptr_eq(value, &cached.fi)),
"cache hits must share FileInfo ownership"
);
assert!(
matches!(returned_parts_metadata, GetObjectMetadata::Shared(ref value) if Arc::ptr_eq(value, &cached.parts_metadata)),
"cache hits must share the metadata vector"
);
assert!(
matches!(returned_online_disks, GetObjectMetadata::Shared(ref value) if Arc::ptr_eq(value, &cached.online_disks)),
"cache hits must share the online-disk vector"
returned.shared_entry().is_some_and(|value| Arc::ptr_eq(value, &cached)),
"cache hits must share the complete metadata snapshot"
);
}
@@ -2781,9 +2766,9 @@ mod metadata_cache_tests {
),
Arc::new(GetObjectMetadataCacheEntry {
created_at: Instant::now(),
fi: Arc::new(fi.clone()),
parts_metadata: Arc::new(vec![fi]),
online_disks: Arc::new(vec![None]),
fi: fi.clone(),
parts_metadata: vec![fi],
online_disks: vec![None],
read_quorum: 1,
}),
)
@@ -2877,13 +2862,12 @@ mod metadata_cache_tests {
barrier.wait_until_paused().await;
set.invalidate_get_object_metadata_cache(bucket, object).await;
barrier.release();
let (fi, parts_metadata, online_disks) = read
let snapshot = read
.await
.expect("metadata read task should not panic")
.expect("metadata fanout should still return its selected FileInfo");
assert!(matches!(fi, GetObjectMetadata::Owned(_)));
assert!(matches!(parts_metadata, GetObjectMetadata::Owned(_)));
assert!(matches!(online_disks, GetObjectMetadata::Owned(_)));
assert!(snapshot.owned.is_some());
assert!(snapshot.has_valid_representation());
assert!(
set.get_object_metadata_cache
@@ -2930,9 +2914,9 @@ mod metadata_cache_tests {
let key = GetObjectMetadataCacheKey::new("bucket", "object", generation);
let entry = Arc::new(GetObjectMetadataCacheEntry {
created_at: Instant::now(),
fi: Arc::new(fi.clone()),
parts_metadata: Arc::new(vec![fi]),
online_disks: Arc::new(Vec::new()),
fi: fi.clone(),
parts_metadata: vec![fi],
online_disks: Vec::new(),
read_quorum: 0,
});
@@ -3040,9 +3024,9 @@ mod metadata_cache_tests {
let entry = |fi: FileInfo| {
Arc::new(GetObjectMetadataCacheEntry {
created_at: Instant::now(),
parts_metadata: Arc::new(vec![fi.clone()]),
fi: Arc::new(fi),
online_disks: Arc::new(Vec::new()),
parts_metadata: vec![fi.clone()],
fi,
online_disks: Vec::new(),
read_quorum: 0,
})
};
@@ -4105,8 +4089,8 @@ mod tests {
get_codec_streaming_reader_gate(
CODEC_STREAMING_TEST_BUCKET,
CODEC_STREAMING_TEST_OBJECT,
range,
None,
classify_get_codec_streaming_object_class(range, object_info, fi),
object_info,
fi,
lock_optimization_enabled,
@@ -4123,8 +4107,8 @@ mod tests {
get_codec_streaming_reader_gate(
CODEC_STREAMING_TEST_BUCKET,
CODEC_STREAMING_TEST_OBJECT,
range,
part_number,
classify_get_codec_streaming_object_class(range, object_info, fi),
object_info,
fi,
lock_optimization_enabled,
@@ -4850,6 +4834,114 @@ mod tests {
.await
}
async fn encoded_inline_blocks(blocks: &[&[u8]], shard_size: usize, hash_algo: HashAlgorithm) -> Bytes {
let mut writer = BitrotWriter::new(Cursor::new(Vec::new()), shard_size, hash_algo);
for block in blocks {
writer.write(block).await.expect("test block should be encoded");
}
Bytes::from(writer.into_inner().into_inner())
}
fn assert_reader_shares_inline_allocation(reader: &ObjectBitrotReader, source: &Bytes) {
let reader_bytes = reader
.inner_ref()
.inline_bytes()
.expect("inline scheduler should retain an in-memory Bytes source");
assert_eq!(
reader_bytes.as_ptr(),
source.as_ptr(),
"the scheduler must clone Bytes ownership instead of copying the inline shard payload"
);
}
#[tokio::test]
async fn inline_range_scheduler_shares_bytes_and_rejects_bitrot_mismatch() {
const SHARD_SIZE: usize = 16;
let hash_algo = HashAlgorithm::HighwayHash256S;
let first = [b'a'; SHARD_SIZE];
let second = [b'b'; SHARD_SIZE];
let mut source = encoded_inline_blocks(&[&first, &second], SHARD_SIZE, hash_algo.clone()).await;
let second_payload = hash_algo.size() * 2 + SHARD_SIZE;
source = {
let mut corrupt = source.to_vec();
corrupt[second_payload] ^= 0xff;
Bytes::from(corrupt)
};
let files = vec![encoded_reader_setup_fileinfo(Some(source.to_vec()))];
let source = files[0].data.clone().expect("inline shard should exist");
let disks = vec![None];
let mut setup = create_bitrot_readers_until_quorum_with_preference(
&files,
&disks,
"bucket",
"object",
1,
SHARD_SIZE,
SHARD_SIZE,
SHARD_SIZE,
hash_algo,
false,
false,
1,
0,
BitrotReaderSetupMode::ReadQuorum,
true,
None,
None,
)
.await;
let mut reader = setup.readers[0].take().expect("range reader should be ready");
assert_reader_shares_inline_allocation(&reader, &source);
let err = reader
.read(&mut [0; SHARD_SIZE])
.await
.expect_err("corrupt ranged inline block must fail bitrot verification");
assert_eq!(err.kind(), ErrorKind::InvalidData);
}
#[tokio::test]
async fn inline_part_scheduler_shares_bytes_and_rejects_bitrot_mismatch() {
const SHARD_SIZE: usize = 16;
let hash_algo = HashAlgorithm::HighwayHash256S;
let block = [b'p'; SHARD_SIZE];
let encoded = encoded_inline_blocks(&[&block], SHARD_SIZE, hash_algo.clone()).await;
let mut corrupt = encoded.to_vec();
corrupt[hash_algo.size()] ^= 0xff;
let files = vec![encoded_reader_setup_fileinfo(Some(corrupt))];
let source = files[0].data.clone().expect("inline shard should exist");
let disks = vec![None];
let mut setup = create_bitrot_readers_until_quorum_all_shards(
&files,
&disks,
"bucket",
"object",
7,
0,
SHARD_SIZE,
SHARD_SIZE,
hash_algo,
false,
false,
1,
0,
BitrotReaderSetupMode::VerifyReconstruction,
None,
None,
)
.await;
let mut reader = setup.readers[0].take().expect("part reader should be ready");
assert_reader_shares_inline_allocation(&reader, &source);
let err = reader
.read(&mut [0; SHARD_SIZE])
.await
.expect_err("corrupt inline part must fail bitrot verification");
assert_eq!(err.kind(), ErrorKind::InvalidData);
}
async fn decode_codec_data_blocks_first_setup(
erasure: coding::Erasure,
data: &[u8],
@@ -5550,6 +5642,63 @@ mod tests {
});
}
#[test]
fn codec_streaming_config_cache_loads_once() {
use std::cell::Cell;
let loads = Cell::new(0);
let expected = GetCodecStreamingConfig {
enabled: true,
rollout: GetCodecStreamingRollout::Off,
rollout_pct: 100,
body_compat_confirmed: true,
header_compat_confirmed: true,
engine: GetCodecStreamingEngine::Legacy,
min_size: DEFAULT_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE,
};
for _ in 0..3 {
assert_eq!(
get_codec_streaming_config_cached_core(|| {
loads.set(loads.get() + 1);
expected
}),
expected
);
}
assert_eq!(loads.get(), 1, "production config cache must not reload env per GET");
}
#[test]
fn codec_streaming_config_loader_preserves_all_gate_env_overrides() {
temp_env::with_vars(
[
(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, Some("false")),
(ENV_RUSTFS_GET_CODEC_STREAMING_ENGINE, Some(GET_CODEC_STREAMING_ENGINE_RUSTFS)),
(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, Some("production")),
(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT_PCT, Some("37")),
(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, Some("false")),
(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, Some("false")),
(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, None::<&str>),
(ENV_RUSTFS_GET_CODEC_STREAMING_RUSTFS_MIN_SIZE, Some("262144")),
],
|| {
assert_eq!(
load_get_codec_streaming_config(),
GetCodecStreamingConfig {
enabled: false,
rollout: GetCodecStreamingRollout::On,
rollout_pct: 37,
body_compat_confirmed: false,
header_compat_confirmed: false,
engine: GetCodecStreamingEngine::Rustfs,
min_size: 262144,
}
);
},
);
}
#[test]
fn codec_streaming_default_min_size_meets_direct_memory_ceiling() {
for engine in [None, Some(GET_CODEC_STREAMING_ENGINE_RUSTFS)] {
+6 -6
View File
@@ -78,10 +78,10 @@ impl SetDisks {
include_part_checksums: true,
..Default::default()
};
let (fi, _, disks) = self
let (mut fi, _, disks) = self
.get_object_fileinfo_gated(bucket, object, &read_opts, false, false)
.await?;
let mut fi = fi.into_owned();
.await?
.into_owned();
if let Some(expected_operation_id) = expected_operation_id {
require_restore_operation_id(&fi.metadata, expected_operation_id)?;
}
@@ -146,10 +146,10 @@ impl SetDisks {
include_part_checksums: true,
..Default::default()
};
let (fi, _, disks) = self
let (mut fi, _, disks) = self
.get_object_fileinfo_gated(bucket, object, &read_opts, false, false)
.await?;
let mut fi = fi.into_owned();
.await?
.into_owned();
if let Some(expected_operation_id) = expected_operation_id {
match restore_operation_id_from_metadata(&fi.metadata)? {
Some(actual_operation_id) if actual_operation_id == expected_operation_id => {}
+94 -159
View File
@@ -16,6 +16,14 @@ use crate::diagnostics::get::{
GET_SHARD_READ_COST_LOCAL, GET_SHARD_READ_COST_REMOTE, GET_SHARD_READ_COST_SAME_NODE, GET_SHARD_READ_COST_UNKNOWN,
};
use crate::disk::error::Error;
use crate::layout::disks_layout::MAX_ERASURE_SET_DRIVE_COUNT;
use smallvec::SmallVec;
/// Generic codec callers may exceed the production set limit; `SmallVec` then
/// spills without changing slot semantics.
pub(crate) const INLINE_SHARD_SLOTS: usize = MAX_ERASURE_SET_DRIVE_COUNT;
pub(crate) type ShardBuffers = SmallVec<[Option<Vec<u8>>; INLINE_SHARD_SLOTS]>;
pub(crate) type ShardErrors = SmallVec<[Option<Error>; INLINE_SHARD_SLOTS]>;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ShardReadCost {
@@ -43,202 +51,137 @@ impl ShardReadCost {
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ShardSlot {
index: usize,
read_cost: ShardReadCost,
data: Option<Vec<u8>>,
error: Option<Error>,
}
impl ShardSlot {
pub(crate) fn new(index: usize, data: Option<Vec<u8>>, error: Option<Error>) -> Self {
Self::with_read_cost(index, ShardReadCost::Unknown, data, error)
}
pub(crate) fn with_read_cost(index: usize, read_cost: ShardReadCost, data: Option<Vec<u8>>, error: Option<Error>) -> Self {
Self {
index,
read_cost,
data,
error,
}
}
pub(crate) fn data(index: usize, data: Vec<u8>) -> Self {
Self::new(index, Some(data), None)
}
pub(crate) fn data_with_read_cost(index: usize, read_cost: ShardReadCost, data: Vec<u8>) -> Self {
Self::with_read_cost(index, read_cost, Some(data), None)
}
pub(crate) fn missing(index: usize, error: Error) -> Self {
Self::new(index, None, Some(error))
}
pub(crate) fn missing_with_read_cost(index: usize, read_cost: ShardReadCost, error: Error) -> Self {
Self::with_read_cost(index, read_cost, None, Some(error))
}
pub(crate) fn index(&self) -> usize {
self.index
}
pub(crate) fn read_cost(&self) -> ShardReadCost {
self.read_cost
}
pub(crate) fn has_data(&self) -> bool {
self.data.is_some()
}
pub(crate) fn data_bytes(&self) -> Option<&[u8]> {
self.data.as_deref()
}
pub(crate) fn error(&self) -> Option<&Error> {
self.error.as_ref()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct StripeReadState {
slots: Vec<ShardSlot>,
shards: ShardBuffers,
errors: ShardErrors,
read_quorum: usize,
}
impl StripeReadState {
pub(crate) fn new(slots: Vec<ShardSlot>, read_quorum: usize) -> Self {
Self { slots, read_quorum }
}
#[cfg(test)]
pub(crate) fn from_parts(shards: Vec<Option<Vec<u8>>>, errors: Vec<Option<Error>>, read_quorum: usize) -> Self {
Self::from_parts_with_read_costs(shards, errors, &[], read_quorum)
let mut shards = SmallVec::from_vec(shards);
let mut errors = SmallVec::from_vec(errors);
let slot_count = shards.len().max(errors.len());
shards.resize_with(slot_count, || None);
errors.resize_with(slot_count, || None);
Self {
shards,
errors,
read_quorum,
}
}
pub(crate) fn from_parts_with_read_costs<S, E>(shards: S, errors: E, read_costs: &[ShardReadCost], read_quorum: usize) -> Self
where
S: IntoIterator<Item = Option<Vec<u8>>>,
S::IntoIter: ExactSizeIterator,
E: IntoIterator<Item = Option<Error>>,
E::IntoIter: ExactSizeIterator,
{
let mut shards = shards.into_iter();
let mut errors = errors.into_iter();
let slot_count = shards.len().max(errors.len());
let mut slots = Vec::with_capacity(slot_count);
for index in 0..slot_count {
let read_cost = read_costs.get(index).copied().unwrap_or(ShardReadCost::Unknown);
slots.push(ShardSlot::with_read_cost(
index,
read_cost,
shards.next().flatten(),
errors.next().flatten(),
));
}
Self::new(slots, read_quorum)
pub(crate) fn with_slot_count(slot_count: usize, read_quorum: usize) -> Self {
let mut state = Self {
shards: SmallVec::new(),
errors: SmallVec::new(),
read_quorum,
};
state.reset(slot_count, read_quorum);
state
}
pub(crate) fn reset(&mut self, slot_count: usize, read_quorum: usize) {
self.shards.clear();
self.shards.resize_with(slot_count, || None);
self.errors.clear();
self.errors.resize_with(slot_count, || None);
self.read_quorum = read_quorum;
}
pub(crate) fn available_shards(&self) -> usize {
self.slots.iter().filter(|slot| slot.has_data()).count()
self.shards.iter().filter(|shard| shard.is_some()).count()
}
pub(crate) fn can_decode(&self) -> bool {
self.available_shards() >= self.read_quorum
}
pub(crate) fn slots(&self) -> &[ShardSlot] {
&self.slots
pub(crate) fn is_empty(&self) -> bool {
self.shards.is_empty()
}
pub(crate) fn slot_by_index(&self, index: usize) -> Option<&ShardSlot> {
if let Some(slot) = self.slots.get(index)
&& slot.index == index
{
return Some(slot);
}
self.slots.iter().find(|slot| slot.index == index)
pub(crate) fn data_bytes(&self, index: usize) -> Option<&[u8]> {
self.shards.get(index).and_then(Option::as_deref)
}
#[cfg(test)]
pub(crate) fn error(&self, index: usize) -> Option<&Error> {
self.errors.get(index).and_then(Option::as_ref)
}
pub(crate) fn data_shards_complete(&self, data_shards: usize) -> bool {
(0..data_shards).all(|index| self.slot_by_index(index).is_some_and(ShardSlot::has_data))
self.shards.len() >= data_shards && self.shards.iter().take(data_shards).all(Option::is_some)
}
pub(crate) fn into_parts(self) -> (Vec<Option<Vec<u8>>>, Vec<Option<Error>>) {
let part_count = self.slots.iter().map(|slot| slot.index).max().map_or(0, |index| index + 1);
let mut shards = Vec::with_capacity(part_count);
shards.resize_with(part_count, || None);
let mut errors = Vec::with_capacity(part_count);
errors.resize_with(part_count, || None);
for slot in self.slots {
shards[slot.index] = slot.data;
errors[slot.index] = slot.error;
}
(shards, errors)
pub(crate) fn parts_mut(&mut self) -> (&mut ShardBuffers, &mut ShardErrors) {
(&mut self.shards, &mut self.errors)
}
pub(crate) fn shards_mut(&mut self) -> &mut ShardBuffers {
&mut self.shards
}
pub(crate) fn into_parts(self) -> (ShardBuffers, ShardErrors) {
(self.shards, self.errors)
}
#[cfg(test)]
pub(crate) fn scratch_storage(&self) -> (*const Option<Vec<u8>>, *const Option<Error>, bool, bool) {
(self.shards.as_ptr(), self.errors.as_ptr(), self.shards.spilled(), self.errors.spilled())
}
#[cfg(test)]
pub(crate) fn shard_allocation(&self, index: usize) -> Option<(*const u8, usize)> {
self.shards
.get(index)
.and_then(|shard| shard.as_ref().map(|shard| (shard.as_ptr(), shard.capacity())))
}
}
#[async_trait::async_trait]
pub(crate) trait ShardStripeSource: Send {
async fn read_next_stripe(&mut self) -> StripeReadState;
async fn read_next_stripe(&mut self) -> Box<StripeReadState>;
fn recycle_stripe(&mut self, _state: Box<StripeReadState>) {}
}
#[cfg(test)]
mod tests {
use super::*;
use std::mem::size_of;
#[test]
fn stripe_read_state_tracks_decode_quorum() {
let state = StripeReadState::new(
vec![
ShardSlot::data_with_read_cost(0, ShardReadCost::Local, vec![1]),
ShardSlot::missing_with_read_cost(1, ShardReadCost::Remote, Error::FileNotFound),
ShardSlot::data_with_read_cost(2, ShardReadCost::SameNode, vec![2]),
],
2,
);
fn stripe_scratch_capacity_matches_the_production_set_limit() {
type OversizedShardBuffers = SmallVec<[Option<Vec<u8>>; 32]>;
type OversizedShardErrors = SmallVec<[Option<Error>; 32]>;
assert_eq!(INLINE_SHARD_SLOTS, MAX_ERASURE_SET_DRIVE_COUNT);
assert!(size_of::<ShardBuffers>() < size_of::<OversizedShardBuffers>());
assert!(size_of::<ShardErrors>() < size_of::<OversizedShardErrors>());
}
#[test]
fn stripe_read_state_tracks_decode_quorum_and_slot_access() {
let state =
StripeReadState::from_parts(vec![Some(vec![1]), None, Some(vec![2])], vec![None, Some(Error::FileNotFound), None], 2);
assert_eq!(state.available_shards(), 2);
assert!(state.can_decode());
assert_eq!(state.slots()[1].index(), 1);
assert_eq!(state.slots()[0].read_cost(), ShardReadCost::Local);
assert!(state.slots()[2].read_cost().is_low_cost());
assert_eq!(state.data_bytes(0), Some(&[1][..]));
assert_eq!(state.error(1), Some(&Error::FileNotFound));
}
#[test]
fn stripe_read_state_preserves_shards_and_errors() {
let state = StripeReadState::new(vec![ShardSlot::missing(1, Error::FileCorrupt), ShardSlot::data(0, vec![1, 2, 3])], 2);
let state = StripeReadState::from_parts(vec![Some(vec![1, 2, 3]), None], vec![None, Some(Error::FileCorrupt)], 2);
assert!(!state.can_decode());
let (shards, errors) = state.into_parts();
assert_eq!(shards, vec![Some(vec![1, 2, 3]), None]);
assert_eq!(errors, vec![None, Some(Error::FileCorrupt)]);
}
#[test]
fn stripe_read_state_builds_slots_from_parallel_reader_parts() {
let state =
StripeReadState::from_parts(vec![Some(vec![1]), None, Some(vec![3])], vec![None, Some(Error::FileNotFound)], 2);
assert!(state.can_decode());
assert_eq!(state.slots()[1].index(), 1);
assert_eq!(state.slots()[1].error(), Some(&Error::FileNotFound));
}
#[test]
fn stripe_read_state_preserves_read_cost_hints() {
let state = StripeReadState::from_parts_with_read_costs(
vec![Some(vec![1]), None, Some(vec![3])],
vec![None, Some(Error::FileNotFound)],
&[ShardReadCost::Local, ShardReadCost::Remote, ShardReadCost::Unknown],
2,
);
assert_eq!(state.slots()[0].read_cost(), ShardReadCost::Local);
assert_eq!(state.slots()[1].read_cost(), ShardReadCost::Remote);
assert_eq!(state.slots()[2].read_cost(), ShardReadCost::Unknown);
assert_eq!(ShardReadCost::SameNode.as_str(), GET_SHARD_READ_COST_SAME_NODE);
assert_eq!(shards.as_slice(), &[Some(vec![1, 2, 3]), None]);
assert_eq!(errors.as_slice(), &[None, Some(Error::FileCorrupt)]);
}
#[test]
@@ -258,8 +201,8 @@ mod tests {
let state = StripeReadState::from_parts(vec![Some(vec![1]), Some(vec![2]), None], Vec::new(), 2);
assert!(state.data_shards_complete(2));
assert_eq!(state.slots()[0].data_bytes(), Some(&[1][..]));
assert_eq!(state.slot_by_index(1).and_then(ShardSlot::data_bytes), Some(&[2][..]));
assert_eq!(state.data_bytes(0), Some(&[1][..]));
assert_eq!(state.data_bytes(1), Some(&[2][..]));
}
#[test]
@@ -268,12 +211,4 @@ mod tests {
assert!(!state.data_shards_complete(2));
}
#[test]
fn stripe_read_state_finds_out_of_order_slots_by_index() {
let state = StripeReadState::new(vec![ShardSlot::data(2, vec![3]), ShardSlot::data(0, vec![1])], 2);
assert_eq!(state.slot_by_index(0).and_then(ShardSlot::data_bytes), Some(&[1][..]));
assert!(state.slot_by_index(1).is_none());
}
}