Files
rustfs/crates/ecstore/src/set_disk/read.rs
T
houseme 27468ebfa9 feat(get): consolidate GET performance optimization (#3972)
* feat(get): consolidate GET performance optimization

Consolidated implementation of all GET performance optimizations into
a single, well-organized commit replacing the previous patch-on-patch
approach.

## Changes

### Configuration (set_disk/mod.rs)
- Consolidated all GET optimization flags into a single organized section
- Enabled by default: codec streaming, metadata early-stop, page cache reclaim
- Added codec streaming multipart flag (default: disabled)
- Added version-aware early-stop flag (default: disabled)
- Added adaptive duplex buffer sizing based on object size
- All flags use OnceLock caching with rollout percentage support

### Metadata Early-Stop (set_disk/read.rs)
- Delete marker early-stop when quorum agrees
- Version-aware early-stop for versioned GET requests
- MetadataQuorumAccumulator enhanced with:
  - delete_marker_votes tracking
  - requested_version_id and matching_version_votes tracking
  - version_early_stop_decision() method
- 6 new tests for version early-stop scenarios

### Codec Streaming (erasure/coding/decode_reader.rs)
- DualInFlight (2-stripe lookahead) enabled by default

### Decode Pipeline (erasure/coding/decode.rs)
- Stripe prefetch count configuration
- Bitrot-decode overlap configuration

### Disk Layer (disk/local.rs)
- O_DIRECT read configuration constants (preparation)

### Metrics (io-metrics/lib.rs)
- BytesPool acquisition/return metrics
- Metadata phase duration with early-stop label
- Total duration with reader_path label

### Diagnostics (diagnostics/)
- Early-stop reason constants
- Pool tier/outcome label constants

### Observability (.docker/observability/)
- 3 Grafana dashboards for GET optimization monitoring
- Prometheus alert rules (6 alerts: 3 critical, 3 warning)
- Updated README.md and README_ZH.md with usage docs

### Config (config/src/constants/runtime.rs)
- Page cache reclaim read enabled by default

## Environment Variables

| Variable | Default | Description |
|----------|---------|-------------|
| RUSTFS_GET_CODEC_STREAMING_ENABLE | true | Codec streaming base flag |
| RUSTFS_GET_CODEC_STREAMING_ROLLOUT_PCT | 100 | Codec streaming rollout % |
| RUSTFS_GET_CODEC_STREAMING_MULTIPART_ENABLE | false | Multipart codec streaming |
| RUSTFS_GET_METADATA_EARLY_STOP_ENABLE | true | Early-stop base flag |
| RUSTFS_GET_METADATA_EARLY_STOP_ROLLOUT_PCT | 100 | Early-stop rollout % |
| RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE | false | Version-aware early-stop |
| RUSTFS_OBJECT_FILE_CACHE_RECLAIM_READ_ENABLE | true | Page cache reclaim |
| RUSTFS_OBJECT_DIRECT_IO_READ_ENABLE | false | O_DIRECT (preparation) |
| RUSTFS_GET_DECODE_STRIPE_PREFETCH_COUNT | 1 | Stripe prefetch |
| RUSTFS_GET_BITROT_DECODE_OVERLAP_ENABLE | false | Bitrot-decode overlap |
| RUSTFS_GET_CODEC_STREAMING_MAX_INFLIGHT | 2 | DualInFlight stripes |

## Rollback

All optimizations can be disabled via environment variables:
RUSTFS_GET_CODEC_STREAMING_ENABLE=false
RUSTFS_GET_METADATA_EARLY_STOP_ENABLE=false
RUSTFS_OBJECT_FILE_CACHE_RECLAIM_READ_ENABLE=false

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

* test(get): add stress test scripts for GET optimization validation

- quick-validate-get-optimization.sh: Quick 5-minute validation
- stress-test-get-optimization.sh: Full 30+ minute stress test
- README-stress-test.md: Usage documentation

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

* test(ecstore): align file cache reclaim defaults

* chore(deps): update redis and erasure codec

* test(ecstore): align decode fill policy default

* test(ecstore): align metadata early-stop default

* fix(ecstore): keep metadata early stop opt-in

---------

Co-authored-by: heihutu <heihutu@gmail.com>
2026-06-28 07:14:07 +08:00

3870 lines
144 KiB
Rust

// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use super::*;
use crate::diagnostics::get::{
GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA, GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER,
GET_METADATA_EARLY_STOP_REASON_ERROR, GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM,
GET_METADATA_EARLY_STOP_REASON_NOT_FOUND, GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST,
GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM, GET_METADATA_EARLY_STOP_REASON_VERSION_MATCH_QUORUM,
GET_METADATA_EARLY_STOP_REASON_VERSION_NOT_FOUND, GET_METADATA_RESPONSE_CORRUPT, GET_METADATA_RESPONSE_DISK_NOT_FOUND,
GET_METADATA_RESPONSE_ERROR, GET_METADATA_RESPONSE_IGNORED, GET_METADATA_RESPONSE_NOT_FOUND, GET_METADATA_RESPONSE_TIMEOUT,
GET_METADATA_RESPONSE_VALID, GET_METADATA_RESPONSE_VERSION_NOT_FOUND, GET_OBJECT_PATH_CODEC_STREAMING,
GET_OBJECT_PATH_LEGACY_DUPLEX, GET_STAGE_DECODE, GET_STAGE_RANGE, GET_STAGE_READER_SETUP, GetObjectFailureReason,
classify_disk_error, record_get_object_pipeline_failure, record_get_object_pipeline_failure_for_path,
};
use crate::erasure::coding::BitrotReader;
use crate::io_support::bitrot::create_deferred_bitrot_reader;
use crate::set_disk::shard_source::ShardReadCost;
use futures::stream::{FuturesUnordered, StreamExt};
use metrics::counter;
use rustfs_config::{DEFAULT_OBJECT_ZERO_COPY_ENABLE, ENV_OBJECT_ZERO_COPY_ENABLE};
use std::{
collections::HashMap,
future::Future,
pin::Pin,
sync::OnceLock,
time::{Duration, Instant},
};
use tokio::io::AsyncRead;
use tokio::sync::RwLock;
use tokio::task::JoinSet;
const EVENT_SET_DISK_READ: &str = "set_disk_read";
const SLOW_OBJECT_READ_LOG_THRESHOLD: Duration = Duration::from_secs(5);
const READ_REPAIR_HEAL_DEDUP_TTL: Duration = Duration::from_secs(60);
const READ_REPAIR_HEAL_DEDUP_MAX_ENTRIES: usize = 4096;
static READ_REPAIR_HEAL_CACHE: OnceLock<RwLock<HashMap<ReadRepairHealCacheKey, Instant>>> = OnceLock::new();
pub(super) enum GetCodecStreamingReaderBuildOutcome {
Reader(Box<dyn AsyncRead + Unpin + Send + Sync>),
Fallback(GetCodecStreamingFallbackReason),
}
pub(super) fn codec_streaming_reader_setup_fallback_reason(missing_shards: usize) -> Option<GetCodecStreamingFallbackReason> {
(missing_shards > 0).then_some(GetCodecStreamingFallbackReason::ReadQuorumNotSafe)
}
#[derive(Clone, Copy, Debug)]
struct MetadataFanoutObservation {
outcome: &'static str,
elapsed: Duration,
valid: bool,
ignored: bool,
}
impl MetadataFanoutObservation {
fn from_file_info(file_info: &FileInfo, elapsed: Duration) -> Self {
if file_info.is_valid() {
Self {
outcome: GET_METADATA_RESPONSE_VALID,
elapsed,
valid: true,
ignored: false,
}
} else {
Self {
outcome: GET_METADATA_RESPONSE_ERROR,
elapsed,
valid: false,
ignored: false,
}
}
}
fn from_error(err: &DiskError, elapsed: Duration) -> Self {
Self {
outcome: classify_metadata_response_error(err),
elapsed,
valid: false,
ignored: is_metadata_fanout_ignored_error(err),
}
}
}
#[derive(Clone, Debug, Default)]
struct MetadataFanoutDiagnostics {
fanout_duration: Duration,
observations: Vec<MetadataFanoutObservation>,
}
impl MetadataFanoutDiagnostics {
fn new(fanout_duration: Duration, observations: Vec<MetadataFanoutObservation>) -> Self {
Self {
fanout_duration,
observations,
}
}
fn total_responses(&self) -> usize {
self.observations.len()
}
fn valid_responses(&self) -> usize {
self.observations.iter().filter(|observation| observation.valid).count()
}
fn ignored_responses(&self) -> usize {
self.observations.iter().filter(|observation| observation.ignored).count()
}
fn error_responses(&self) -> usize {
self.total_responses().saturating_sub(self.valid_responses())
}
fn first_response_latency(&self) -> Option<Duration> {
self.observations.iter().map(|observation| observation.elapsed).min()
}
fn first_valid_response_latency(&self) -> Option<Duration> {
self.observations
.iter()
.filter(|observation| observation.valid)
.map(|observation| observation.elapsed)
.min()
}
fn slowest_response_latency(&self) -> Option<Duration> {
self.observations.iter().map(|observation| observation.elapsed).max()
}
fn quorum_candidate_latency(&self, read_quorum: usize) -> Option<Duration> {
if read_quorum == 0 {
return Some(Duration::ZERO);
}
let mut valid_latencies = self
.observations
.iter()
.filter(|observation| observation.valid)
.map(|observation| observation.elapsed)
.collect::<Vec<_>>();
valid_latencies.sort_unstable();
valid_latencies.get(read_quorum.saturating_sub(1)).copied()
}
fn record(&self, path: &'static str) {
rustfs_io_metrics::record_get_object_metadata_fanout_duration(path, self.fanout_duration.as_secs_f64());
if let Some(latency) = self.first_response_latency() {
rustfs_io_metrics::record_get_object_first_metadata_response_latency(path, latency.as_secs_f64());
}
if let Some(latency) = self.first_valid_response_latency() {
rustfs_io_metrics::record_get_object_first_valid_metadata_response_latency(path, latency.as_secs_f64());
}
if let Some(latency) = self.slowest_response_latency() {
rustfs_io_metrics::record_get_object_slowest_metadata_response_latency(path, latency.as_secs_f64());
}
rustfs_io_metrics::record_get_object_metadata_fanout_shape(
path,
self.total_responses(),
self.valid_responses(),
self.ignored_responses(),
self.error_responses(),
);
for observation in &self.observations {
rustfs_io_metrics::record_get_object_metadata_response(path, observation.outcome);
}
}
fn record_quorum_candidate_latency(&self, path: &'static str, read_quorum: usize) {
if let Some(latency) = self.quorum_candidate_latency(read_quorum) {
rustfs_io_metrics::record_get_object_quorum_reached_latency(path, latency.as_secs_f64());
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
struct MetadataEarlyStopDecision {
reason: &'static str,
}
#[derive(Clone, Debug)]
struct MetadataQuorumAccumulator {
total_disks: usize,
default_parity_count: usize,
allow_early_stop: bool,
valid_responses: usize,
not_found_responses: usize,
version_not_found_responses: usize,
ignored_errors: usize,
hard_errors: usize,
candidate: Option<FileInfo>,
candidate_votes: usize,
conflicting_metadata: bool,
delete_marker_seen: bool,
delete_marker_votes: usize,
requested_version_id: String,
matching_version_votes: usize,
}
impl MetadataQuorumAccumulator {
fn new(total_disks: usize, default_parity_count: usize, allow_early_stop: bool) -> Self {
Self {
total_disks,
default_parity_count,
allow_early_stop,
valid_responses: 0,
not_found_responses: 0,
version_not_found_responses: 0,
ignored_errors: 0,
hard_errors: 0,
candidate: None,
candidate_votes: 0,
conflicting_metadata: false,
delete_marker_seen: false,
delete_marker_votes: 0,
requested_version_id: String::new(),
matching_version_votes: 0,
}
}
fn with_requested_version_id(mut self, version_id: &str) -> Self {
self.requested_version_id = version_id.to_string();
self
}
fn observe_file_info(&mut self, file_info: &FileInfo) {
if !file_info.is_valid() {
self.hard_errors = self.hard_errors.saturating_add(1);
return;
}
self.valid_responses = self.valid_responses.saturating_add(1);
// Track version match for versioned requests
if !self.requested_version_id.is_empty()
&& let Some(ref vid) = file_info.version_id
&& vid.to_string() == self.requested_version_id
{
self.matching_version_votes = self.matching_version_votes.saturating_add(1);
}
if file_info.deleted {
self.delete_marker_votes = self.delete_marker_votes.saturating_add(1);
self.delete_marker_seen = true;
return;
}
match &self.candidate {
Some(candidate) if metadata_early_stop_candidate_matches(candidate, file_info) => {
self.candidate_votes = self.candidate_votes.saturating_add(1);
}
Some(_) => {
self.conflicting_metadata = true;
}
None => {
self.candidate = Some(file_info.clone());
self.candidate_votes = 1;
}
}
}
fn observe_error(&mut self, err: &DiskError) {
match err {
DiskError::FileNotFound | DiskError::VolumeNotFound => {
self.not_found_responses = self.not_found_responses.saturating_add(1);
}
DiskError::FileVersionNotFound => {
self.version_not_found_responses = self.version_not_found_responses.saturating_add(1);
}
_ if is_metadata_fanout_ignored_error(err) => {
self.ignored_errors = self.ignored_errors.saturating_add(1);
}
_ => {
self.hard_errors = self.hard_errors.saturating_add(1);
}
}
}
fn early_stop_decision(&self) -> Option<MetadataEarlyStopDecision> {
if !self.allow_early_stop {
return None;
}
if self.delete_marker_votes >= self.missing_response_quorum() {
return Some(MetadataEarlyStopDecision {
reason: GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER,
});
}
if self.conflicting_metadata
|| self.delete_marker_seen
|| self.not_found_responses > 0
|| self.version_not_found_responses > 0
|| self.hard_errors > 0
{
return None;
}
if self
.candidate
.as_ref()
.and_then(|candidate| self.candidate_read_quorum(candidate))
.is_some_and(|read_quorum| self.candidate_votes >= read_quorum)
{
return Some(MetadataEarlyStopDecision {
reason: GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM,
});
}
None
}
/// Check if a versioned request can early-stop because the requested
/// version_id has reached quorum across disks.
fn version_early_stop_decision(&self) -> Option<MetadataEarlyStopDecision> {
if self.requested_version_id.is_empty() {
return None;
}
if self.matching_version_votes >= self.read_quorum_for_version() {
return Some(MetadataEarlyStopDecision {
reason: GET_METADATA_EARLY_STOP_REASON_VERSION_MATCH_QUORUM,
});
}
None
}
/// Compute the read quorum threshold for version-aware early-stop.
/// Uses `total_disks / 2` (like `missing_response_quorum`) when
/// `default_parity_count` is set, otherwise requires all disks.
fn read_quorum_for_version(&self) -> usize {
self.missing_response_quorum()
}
fn final_miss_reason(&self) -> &'static str {
if !self.allow_early_stop {
return GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST;
}
if self.conflicting_metadata {
return GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA;
}
if self.delete_marker_seen {
return GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER;
}
let missing_response_quorum = self.missing_response_quorum();
if self.version_not_found_responses >= missing_response_quorum {
return GET_METADATA_EARLY_STOP_REASON_VERSION_NOT_FOUND;
}
if self.not_found_responses >= missing_response_quorum {
return GET_METADATA_EARLY_STOP_REASON_NOT_FOUND;
}
if self.hard_errors > 0 {
return GET_METADATA_EARLY_STOP_REASON_ERROR;
}
if self.ignored_errors > 0 {
return GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM;
}
GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM
}
fn candidate_read_quorum(&self, candidate: &FileInfo) -> Option<usize> {
if self.default_parity_count == 0 {
return Some(self.total_disks);
}
if candidate.deleted || candidate.size == 0 || candidate.erasure.parity_blocks >= self.total_disks {
return None;
}
Some(self.total_disks.saturating_sub(candidate.erasure.parity_blocks))
}
fn missing_response_quorum(&self) -> usize {
if self.default_parity_count == 0 {
self.total_disks
} else {
self.total_disks / 2
}
}
}
fn metadata_early_stop_candidate_matches(left: &FileInfo, right: &FileInfo) -> bool {
left.volume == right.volume
&& left.name == right.name
&& left.version_id == right.version_id
&& left.is_latest == right.is_latest
&& left.deleted == right.deleted
&& left.mark_deleted == right.mark_deleted
&& left.size == right.size
&& left.mod_time == right.mod_time
&& left.mode == right.mode
&& left.metadata == right.metadata
&& left.parts == right.parts
&& left.checksum == right.checksum
&& left.versioned == right.versioned
&& left.erasure.algorithm == right.erasure.algorithm
&& left.erasure.data_blocks == right.erasure.data_blocks
&& left.erasure.parity_blocks == right.erasure.parity_blocks
&& left.erasure.block_size == right.erasure.block_size
&& left.erasure.distribution == right.erasure.distribution
}
fn classify_metadata_response_error(err: &DiskError) -> &'static str {
match err {
DiskError::FileNotFound | DiskError::VolumeNotFound => GET_METADATA_RESPONSE_NOT_FOUND,
DiskError::FileVersionNotFound => GET_METADATA_RESPONSE_VERSION_NOT_FOUND,
DiskError::DiskNotFound => GET_METADATA_RESPONSE_DISK_NOT_FOUND,
DiskError::FileCorrupt | DiskError::CorruptedFormat | DiskError::CorruptedBackend | DiskError::OutdatedXLMeta => {
GET_METADATA_RESPONSE_CORRUPT
}
DiskError::Timeout => GET_METADATA_RESPONSE_TIMEOUT,
DiskError::FaultyDisk | DiskError::FaultyRemoteDisk => GET_METADATA_RESPONSE_IGNORED,
_ => GET_METADATA_RESPONSE_ERROR,
}
}
fn is_metadata_fanout_ignored_error(err: &DiskError) -> bool {
OBJECT_OP_IGNORED_ERRS.iter().any(|ignored| ignored == err)
}
fn is_confirmed_missing_part_error(err: Option<&str>) -> bool {
let Some(err) = err else {
return false;
};
err.contains("file not found")
|| err.contains("No such file or directory")
|| err.contains("Specified part could not be found")
|| (err.starts_with("part.") && err.ends_with(" not found"))
}
fn resolve_read_part_from_responses(
bucket: &str,
part_meta_path: &str,
part_number: usize,
part_idx: usize,
expected_part_count: usize,
responses: &[Option<Vec<ObjectPartInfo>>],
read_quorum: usize,
) -> disk::error::Result<ObjectPartInfo> {
let mut etag_quorum = HashMap::new();
let mut part_infos = Vec::new();
let mut present_count = 0usize;
let mut missing_count = 0usize;
let mut transient_error_count = 0usize;
let mut mismatched_response_count = 0usize;
for response in responses.iter() {
let Some(parts) = response else {
transient_error_count += 1;
continue;
};
if parts.len() != expected_part_count {
mismatched_response_count += 1;
continue;
}
if !parts[part_idx].etag.is_empty() {
present_count += 1;
*etag_quorum.entry(parts[part_idx].etag.clone()).or_insert(0) += 1;
part_infos.push(parts[part_idx].clone());
continue;
}
if is_confirmed_missing_part_error(parts[part_idx].error.as_deref()) {
missing_count += 1;
} else {
transient_error_count += 1;
}
}
let mut max_etag_quorum = 0;
let mut max_etag = None;
for (etag, quorum) in etag_quorum.iter() {
if quorum > &max_etag_quorum {
max_etag_quorum = *quorum;
max_etag = Some(etag);
}
}
let max_quorum = max_etag_quorum.max(missing_count);
let mut found = None;
for info in part_infos.iter() {
if let Some(etag) = max_etag
&& info.etag == *etag
{
found = Some(info.clone());
break;
}
}
if let (Some(found), Some(max_etag)) = (found, max_etag)
&& !found.etag.is_empty()
&& etag_quorum.get(max_etag).unwrap_or(&0) >= &read_quorum
{
return Ok(found);
}
if missing_count >= read_quorum {
return Ok(ObjectPartInfo {
number: part_number,
error: Some(format!("part.{part_number} not found")),
..Default::default()
});
}
if issue3031_diag_enabled() {
warn!(
target: "rustfs_ecstore::set_disk",
bucket = %bucket,
part_meta_path = %part_meta_path,
part_id = part_number,
read_quorum = read_quorum,
max_quorum = max_quorum,
disk_response_count = responses.len(),
present_count = present_count,
missing_count = missing_count,
transient_error_count = transient_error_count,
mismatched_response_count = mismatched_response_count,
"issue3031_read_parts_part_quorum"
);
}
Err(DiskError::ErasureReadQuorum)
}
fn shard_read_cost_for_disk(disk: Option<&DiskStore>) -> ShardReadCost {
match disk {
Some(disk) if disk.is_local() => ShardReadCost::Local,
Some(_) => ShardReadCost::Remote,
None => ShardReadCost::Unknown,
}
}
#[derive(Clone, Debug, Eq, Hash, PartialEq)]
struct ReadRepairHealCacheKey {
bucket: String,
object: String,
version_id: Option<String>,
pool_index: usize,
set_index: usize,
}
impl ReadRepairHealCacheKey {
fn new(bucket: &str, object: &str, version_id: Option<&str>, pool_index: usize, set_index: usize) -> Self {
Self {
bucket: bucket.to_string(),
object: object.to_string(),
version_id: version_id.filter(|value| !value.is_empty()).map(str::to_string),
pool_index,
set_index,
}
}
}
fn resolved_read_repair_version_id(fi: &FileInfo, requested_version_id: Option<&str>) -> Option<String> {
fi.version_id
.as_ref()
.map(ToString::to_string)
.or_else(|| requested_version_id.filter(|value| !value.is_empty()).map(str::to_string))
}
async fn reserve_read_repair_heal(
bucket: &str,
object: &str,
version_id: Option<&str>,
pool_index: usize,
set_index: usize,
) -> Option<ReadRepairHealCacheKey> {
let key = ReadRepairHealCacheKey::new(bucket, object, version_id, pool_index, set_index);
let now = Instant::now();
let cache = READ_REPAIR_HEAL_CACHE.get_or_init(|| RwLock::new(HashMap::new()));
{
let cache = cache.read().await;
if cache
.get(&key)
.is_some_and(|seen_at| now.saturating_duration_since(*seen_at) <= READ_REPAIR_HEAL_DEDUP_TTL)
{
return None;
}
}
let mut cache = cache.write().await;
if cache
.get(&key)
.is_some_and(|seen_at| now.saturating_duration_since(*seen_at) <= READ_REPAIR_HEAL_DEDUP_TTL)
{
return None;
}
if cache.len() >= READ_REPAIR_HEAL_DEDUP_MAX_ENTRIES {
cache.retain(|_, seen_at| now.saturating_duration_since(*seen_at) <= READ_REPAIR_HEAL_DEDUP_TTL);
}
if cache.len() >= READ_REPAIR_HEAL_DEDUP_MAX_ENTRIES
&& let Some(oldest_key) = cache.iter().min_by_key(|(_, seen_at)| **seen_at).map(|(key, _)| key.clone())
{
cache.remove(&oldest_key);
}
cache.insert(key.clone(), now);
Some(key)
}
async fn release_read_repair_heal_reservation(key: &ReadRepairHealCacheKey) {
if let Some(cache) = READ_REPAIR_HEAL_CACHE.get() {
cache.write().await.remove(key);
}
}
fn record_read_repair_dedup(reason: &'static str) {
counter!("rustfs_heal_read_repair_dedup_total", "reason" => reason.to_string()).increment(1);
}
enum ReadRepairAdmissionOutcome {
Response(HealAdmissionResult),
Failed(String),
}
type ReadRepairAdmissionFuture = Pin<Box<dyn Future<Output = ReadRepairAdmissionOutcome> + Send>>;
type ReadRepairAdmissionSubmitter = fn(rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture;
struct ReadRepairHealSubmission<'a> {
bucket: &'a str,
object: &'a str,
version_id: Option<&'a str>,
pool_index: usize,
set_index: usize,
part_number: Option<usize>,
reason: &'static str,
}
fn send_read_repair_heal_request(request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture {
Box::pin(async {
match send_heal_request_with_admission(request).await {
Ok(result) => ReadRepairAdmissionOutcome::Response(result),
Err(err) => ReadRepairAdmissionOutcome::Failed(err),
}
})
}
async fn submit_read_repair_heal(
bucket: &str,
object: &str,
version_id: Option<&str>,
pool_index: usize,
set_index: usize,
part_number: Option<usize>,
reason: &'static str,
) {
submit_read_repair_heal_with_submitter(
ReadRepairHealSubmission {
bucket,
object,
version_id,
pool_index,
set_index,
part_number,
reason,
},
send_read_repair_heal_request,
)
.await;
}
async fn submit_read_repair_heal_with_submitter(
submission: ReadRepairHealSubmission<'_>,
submitter: ReadRepairAdmissionSubmitter,
) {
let ReadRepairHealSubmission {
bucket,
object,
version_id,
pool_index,
set_index,
part_number,
reason,
} = submission;
let Some(dedup_key) = reserve_read_repair_heal(bucket, object, version_id, pool_index, set_index).await else {
record_read_repair_dedup("duplicate");
debug!(
bucket,
object, part_number, pool_index, set_index, reason, "Skipped duplicate read-repair heal request"
);
return;
};
let mut request = rustfs_common::heal_channel::create_heal_request_with_options(
bucket.to_string(),
Some(object.to_string()),
false,
Some(HealChannelPriority::Normal),
Some(pool_index),
Some(set_index),
);
request.source = HealRequestSource::ReadRepair;
request.object_version_id = version_id.filter(|value| !value.is_empty()).map(str::to_string);
request.recreate_missing = Some(true);
let request_id = request.id.clone();
let bucket = bucket.to_string();
let object = object.to_string();
tokio::spawn(async move {
match submitter(request).await {
ReadRepairAdmissionOutcome::Response(result) if result.is_admitted() => {
debug!(
bucket,
object,
part_number,
pool_index,
set_index,
request_id,
reason,
admission = result.result_label(),
"Read-repair heal request admitted"
);
}
ReadRepairAdmissionOutcome::Response(result) => {
release_read_repair_heal_reservation(&dedup_key).await;
debug!(
bucket,
object,
part_number,
pool_index,
set_index,
request_id,
reason,
admission = result.result_label(),
drop_reason = result.reason_label(),
"Read-repair heal request not admitted"
);
}
ReadRepairAdmissionOutcome::Failed(err) => {
release_read_repair_heal_reservation(&dedup_key).await;
debug!(
bucket,
object,
part_number,
pool_index,
set_index,
request_id,
reason,
error = %err,
"Read-repair heal request could not be submitted"
);
}
}
});
}
type ObjectBitrotReader = BitrotReader<Box<dyn AsyncRead + Send + Sync + Unpin>>;
struct BitrotReaderSetup {
readers: Vec<Option<ObjectBitrotReader>>,
errors: Vec<Option<DiskError>>,
attempted: Vec<bool>,
ready: Vec<bool>,
}
#[derive(Clone, Copy)]
enum BitrotReaderSetupMode {
ReadQuorum,
VerifyReconstruction,
}
impl BitrotReaderSetup {
fn available_shards(&self) -> usize {
self.ready.iter().filter(|ready| **ready).count()
}
fn available_data_shards(&self, data_shards: usize) -> usize {
self.ready.iter().take(data_shards).filter(|ready| **ready).count()
}
fn completed_failed_shards(&self) -> usize {
self.attempted
.iter()
.zip(&self.errors)
.filter(|(attempted, err)| **attempted && err.is_some())
.count()
}
fn data_shards_attempted(&self, data_shards: usize) -> bool {
self.attempted.iter().take(data_shards).all(|attempted| *attempted)
}
fn reconstruction_verification_target(&self, data_shards: usize, parity_shards: usize) -> usize {
let missing_data_sources = data_shards.saturating_sub(self.available_data_shards(data_shards));
if missing_data_sources > 0 && missing_data_sources < parity_shards {
data_shards.saturating_add(1).min(data_shards.saturating_add(parity_shards))
} else {
data_shards
}
}
fn has_setup_quorum(&self, data_shards: usize, parity_shards: usize, mode: BitrotReaderSetupMode) -> bool {
let target = match mode {
BitrotReaderSetupMode::ReadQuorum => data_shards,
BitrotReaderSetupMode::VerifyReconstruction => self.reconstruction_verification_target(data_shards, parity_shards),
};
self.available_shards() >= target
}
}
#[allow(clippy::too_many_arguments)]
async fn create_bitrot_readers_until_quorum(
files: &[FileInfo],
disks: &[Option<DiskStore>],
bucket: &str,
object: &str,
part_number: usize,
read_offset: usize,
read_length: usize,
shard_size: usize,
checksum_algo: HashAlgorithm,
skip_verify_bitrot: bool,
use_zero_copy: bool,
data_shards: usize,
parity_shards: usize,
mode: BitrotReaderSetupMode,
) -> BitrotReaderSetup {
let mut setup = BitrotReaderSetup {
readers: (0..disks.len()).map(|_| None).collect(),
errors: vec![Some(DiskError::DiskNotFound); disks.len()],
attempted: vec![false; disks.len()],
ready: vec![false; disks.len()],
};
let mut reader_tasks = FuturesUnordered::new();
for (idx, disk_op) in disks.iter().enumerate() {
let inline_data = files[idx].data.as_deref();
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(
inline_data,
disk,
bucket,
&path,
read_offset,
read_length,
shard_size,
checksum_algo,
skip_verify_bitrot,
use_zero_copy,
)
.await;
(idx, result)
});
}
while let Some((idx, result)) = reader_tasks.next().await {
setup.attempted[idx] = true;
match result {
Ok(Some(reader)) => {
setup.readers[idx] = Some(reader);
setup.errors[idx] = None;
setup.ready[idx] = true;
}
Ok(None) => {
setup.readers[idx] = None;
setup.errors[idx] = Some(DiskError::DiskNotFound);
setup.ready[idx] = false;
}
Err(e) => {
setup.readers[idx] = None;
setup.errors[idx] = Some(e);
setup.ready[idx] = false;
}
}
if setup.has_setup_quorum(data_shards, parity_shards, mode) {
break;
}
}
if setup.has_setup_quorum(data_shards, parity_shards, mode) {
for idx in 0..disks.len() {
if setup.attempted[idx] {
continue;
}
let inline_data = files[idx].data.clone();
let disk = disks[idx].clone();
let data_dir = files[idx].data_dir.unwrap_or_default();
let path = format!("{object}/{data_dir}/part.{part_number}");
setup.readers[idx] = Some(create_deferred_bitrot_reader(
inline_data,
disk,
bucket,
&path,
read_offset,
read_length,
shard_size,
checksum_algo.clone(),
skip_verify_bitrot,
use_zero_copy,
));
setup.errors[idx] = None;
}
}
drop(reader_tasks);
setup
}
async fn collect_read_multiple_results<F>(
tasks: Vec<F>,
read_quorum: usize,
) -> std::result::Result<(Vec<Option<Vec<ReadMultipleResp>>>, Vec<Option<DiskError>>), ()>
where
F: Future<Output = disk::error::Result<Vec<ReadMultipleResp>>> + Send + 'static,
{
let mut responses = vec![None; tasks.len()];
let mut errors = vec![Some(DiskError::DiskNotFound); tasks.len()];
let mut successful_responses = 0usize;
let mut pending = tasks.len();
let mut join_set = JoinSet::new();
for (index, task) in tasks.into_iter().enumerate() {
join_set.spawn(async move { (index, task.await) });
}
while let Some(join_result) = join_set.join_next().await {
pending = pending.saturating_sub(1);
match join_result {
Ok((index, Ok(resp))) => {
responses[index] = Some(resp);
errors[index] = None;
successful_responses += 1;
}
Ok((index, Err(err))) => {
errors[index] = Some(err);
}
Err(_) => {}
}
if successful_responses + pending < read_quorum {
return Err(());
}
}
Ok((responses, errors))
}
async fn collect_read_parts_results<F>(
tasks: Vec<F>,
read_quorum: usize,
) -> std::result::Result<(Vec<Option<Vec<ObjectPartInfo>>>, Vec<Option<DiskError>>), ()>
where
F: Future<Output = disk::error::Result<Vec<ObjectPartInfo>>> + Send + 'static,
{
let mut responses = vec![None; tasks.len()];
let mut errors = vec![Some(DiskError::DiskNotFound); tasks.len()];
let mut successful_responses = 0usize;
let mut pending = tasks.len();
let mut join_set = JoinSet::new();
for (index, task) in tasks.into_iter().enumerate() {
join_set.spawn(async move { (index, task.await) });
}
while let Some(join_result) = join_set.join_next().await {
pending = pending.saturating_sub(1);
match join_result {
Ok((index, Ok(resp))) => {
responses[index] = Some(resp);
errors[index] = None;
successful_responses += 1;
}
Ok((index, Err(err))) => {
errors[index] = Some(err);
}
Err(_) => {}
}
if successful_responses + pending < read_quorum {
return Err(());
}
}
Ok((responses, errors))
}
impl SetDisks {
async fn is_get_object_metadata_cache_enabled(&self, bucket: &str, opts: &ObjectOptions, read_data: bool) -> bool {
is_get_object_metadata_cache_request_eligible(bucket, opts, read_data) && !runtime_sources::setup_is_dist_erasure().await
}
async fn cached_get_object_fileinfo(&self, bucket: &str, object: &str) -> Option<GetObjectMetadataCacheEntry> {
let key = GetObjectMetadataCacheKey::new(bucket, object);
// moka handles TTL expiry automatically; no is_fresh() check needed
self.get_object_metadata_cache
.get(&key)
.await
.filter(|entry| entry.online_disks.iter().filter(|disk| disk.is_some()).count() >= entry.read_quorum)
}
async fn cache_get_object_fileinfo(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
parts_metadata: &[FileInfo],
online_disks: &[Option<DiskStore>],
read_quorum: usize,
) {
if fi.deleted || !fi.is_valid() {
return;
}
let key = GetObjectMetadataCacheKey::new(bucket, object);
// moka handles capacity eviction (LRU) automatically
self.get_object_metadata_cache
.insert(
key,
GetObjectMetadataCacheEntry {
created_at: Instant::now(),
fi: fi.clone(),
parts_metadata: parts_metadata.to_vec(),
online_disks: online_disks.to_vec(),
read_quorum,
},
)
.await;
}
pub(super) async fn read_parts(
disks: &[Option<DiskStore>],
bucket: &str,
part_meta_paths: &[String],
part_numbers: &[usize],
read_quorum: usize,
) -> disk::error::Result<Vec<ObjectPartInfo>> {
let bucket = bucket.to_string();
let part_meta_paths = part_meta_paths.to_vec();
let tasks: Vec<_> = disks
.iter()
.map(|disk| {
let disk = disk.clone();
let bucket = bucket.clone();
let part_meta_paths = part_meta_paths.clone();
async move {
if let Some(disk) = disk {
disk.read_parts(&bucket, &part_meta_paths).await
} else {
Err(DiskError::DiskNotFound)
}
}
})
.collect();
let (responses, collected_errors) = match collect_read_parts_results(tasks, read_quorum).await {
Ok(collected) => collected,
Err(()) => return Err(DiskError::ErasureReadQuorum),
};
if let Some(err) = reduce_read_quorum_errs(&collected_errors, OBJECT_OP_IGNORED_ERRS, read_quorum) {
return Err(err);
}
let mut ret = vec![ObjectPartInfo::default(); part_meta_paths.len()];
for (part_idx, part_info) in part_meta_paths.iter().enumerate() {
ret[part_idx] = resolve_read_part_from_responses(
&bucket,
part_info,
part_numbers[part_idx],
part_idx,
part_meta_paths.len(),
&responses,
read_quorum,
)?;
}
Ok(ret)
}
#[allow(clippy::too_many_arguments)]
#[tracing::instrument(level = "debug", skip(disks))]
pub(super) async fn read_all_fileinfo(
disks: &[Option<DiskStore>],
org_bucket: &str,
bucket: &str,
object: &str,
version_id: &str,
read_data: bool,
healing: bool,
incl_free_versions: bool,
) -> disk::error::Result<(Vec<FileInfo>, Vec<Option<DiskError>>)> {
let (ress, errors, _) = Self::read_all_fileinfo_inner(
disks,
org_bucket,
bucket,
object,
version_id,
read_data,
healing,
incl_free_versions,
false,
0,
)
.await?;
Ok((ress, errors))
}
#[allow(clippy::too_many_arguments)]
async fn read_all_fileinfo_observed(
disks: &[Option<DiskStore>],
org_bucket: &str,
bucket: &str,
object: &str,
version_id: &str,
read_data: bool,
healing: bool,
incl_free_versions: bool,
default_parity_count: usize,
) -> disk::error::Result<(Vec<FileInfo>, Vec<Option<DiskError>>, MetadataFanoutDiagnostics)> {
Self::read_all_fileinfo_inner(
disks,
org_bucket,
bucket,
object,
version_id,
read_data,
healing,
incl_free_versions,
true,
default_parity_count,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn read_all_fileinfo_inner(
disks: &[Option<DiskStore>],
org_bucket: &str,
bucket: &str,
object: &str,
version_id: &str,
read_data: bool,
healing: bool,
incl_free_versions: bool,
observe: bool,
default_parity_count: usize,
) -> disk::error::Result<(Vec<FileInfo>, Vec<Option<DiskError>>, MetadataFanoutDiagnostics)> {
let early_stop_enabled = observe && is_get_metadata_early_stop_enabled();
let allow_early_stop = observe
&& ((is_get_metadata_early_stop_enabled() && version_id.is_empty() && !healing && !incl_free_versions)
|| (is_version_early_stop_enabled() && !version_id.is_empty() && !healing));
if allow_early_stop {
return Self::read_all_fileinfo_early_stop(
disks,
org_bucket,
bucket,
object,
version_id,
read_data,
healing,
incl_free_versions,
default_parity_count,
)
.await;
}
if early_stop_enabled {
rustfs_io_metrics::record_get_object_metadata_early_stop_miss(
GET_OBJECT_PATH_LEGACY_DUPLEX,
GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST,
);
rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(GET_OBJECT_PATH_LEGACY_DUPLEX, 0);
}
Self::read_all_fileinfo_full_wait(
disks,
org_bucket,
bucket,
object,
version_id,
read_data,
healing,
incl_free_versions,
observe,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn read_all_fileinfo_full_wait(
disks: &[Option<DiskStore>],
org_bucket: &str,
bucket: &str,
object: &str,
version_id: &str,
read_data: bool,
healing: bool,
incl_free_versions: bool,
observe: bool,
) -> disk::error::Result<(Vec<FileInfo>, Vec<Option<DiskError>>, MetadataFanoutDiagnostics)> {
let fanout_start = observe.then(Instant::now);
let mut ress = Vec::with_capacity(disks.len());
let mut errors = Vec::with_capacity(disks.len());
let mut observations = observe.then(|| Vec::with_capacity(disks.len()));
let opts = Arc::new(ReadOptions {
incl_free_versions,
read_data,
healing,
});
let org_bucket = Arc::new(org_bucket.to_string());
let bucket = Arc::new(bucket.to_string());
let object = Arc::new(object.to_string());
let version_id = Arc::new(version_id.to_string());
let futures = disks.iter().map(|disk| {
let disk = disk.clone();
let opts = opts.clone();
let org_bucket = org_bucket.clone();
let bucket = bucket.clone();
let object = object.clone();
let version_id = version_id.clone();
tokio::spawn(async move {
let response_start = observe.then(Instant::now);
let result = if let Some(disk) = disk {
disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).await
} else {
Err(DiskError::DiskNotFound)
};
let elapsed = response_start.map(|start| start.elapsed());
(result, elapsed)
})
});
// Wait for all tasks to complete
let results = join_all(futures).await;
for result in results {
match result {
Ok((res, elapsed)) => match res {
Ok(file_info) => {
if let (Some(observations), Some(elapsed)) = (&mut observations, elapsed) {
observations.push(MetadataFanoutObservation::from_file_info(&file_info, elapsed));
}
ress.push(file_info);
errors.push(None);
}
Err(e) => {
if let (Some(observations), Some(elapsed)) = (&mut observations, elapsed) {
observations.push(MetadataFanoutObservation::from_error(&e, elapsed));
}
ress.push(FileInfo::default());
errors.push(Some(e));
}
},
Err(_) => {
let err = DiskError::Unexpected;
if let (Some(observations), Some(fanout_start)) = (&mut observations, fanout_start) {
observations.push(MetadataFanoutObservation::from_error(&err, fanout_start.elapsed()));
}
ress.push(FileInfo::default());
errors.push(Some(err));
}
}
}
let diagnostics = match (fanout_start, observations) {
(Some(fanout_start), Some(observations)) => MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations),
_ => MetadataFanoutDiagnostics::default(),
};
Ok((ress, errors, diagnostics))
}
#[allow(clippy::too_many_arguments)]
async fn read_all_fileinfo_early_stop(
disks: &[Option<DiskStore>],
org_bucket: &str,
bucket: &str,
object: &str,
version_id: &str,
read_data: bool,
healing: bool,
incl_free_versions: bool,
default_parity_count: usize,
) -> disk::error::Result<(Vec<FileInfo>, Vec<Option<DiskError>>, MetadataFanoutDiagnostics)> {
let fanout_start = Instant::now();
let mut ress = vec![FileInfo::default(); disks.len()];
let mut errors = vec![None; disks.len()];
let mut observations = Vec::with_capacity(disks.len());
let mut accumulator =
MetadataQuorumAccumulator::new(disks.len(), default_parity_count, true).with_requested_version_id(version_id);
let opts = Arc::new(ReadOptions {
incl_free_versions,
read_data,
healing,
});
let org_bucket = Arc::new(org_bucket.to_string());
let bucket = Arc::new(bucket.to_string());
let object = Arc::new(object.to_string());
let version_id = Arc::new(version_id.to_string());
let mut join_set = JoinSet::new();
for (index, disk) in disks.iter().cloned().enumerate() {
let opts = opts.clone();
let org_bucket = org_bucket.clone();
let bucket = bucket.clone();
let object = object.clone();
let version_id = version_id.clone();
join_set.spawn(async move {
let response_start = Instant::now();
let result = if let Some(disk) = disk {
disk.read_version(&org_bucket, &bucket, &object, &version_id, &opts).await
} else {
Err(DiskError::DiskNotFound)
};
(index, result, response_start.elapsed())
});
}
while let Some(result) = join_set.join_next().await {
match result {
Ok((index, res, elapsed)) => match res {
Ok(file_info) => {
observations.push(MetadataFanoutObservation::from_file_info(&file_info, elapsed));
accumulator.observe_file_info(&file_info);
if let Some(slot) = ress.get_mut(index) {
*slot = file_info;
}
}
Err(err) => {
observations.push(MetadataFanoutObservation::from_error(&err, elapsed));
accumulator.observe_error(&err);
if let Some(slot) = errors.get_mut(index) {
*slot = Some(err);
}
}
},
Err(_) => {
let err = DiskError::Unexpected;
observations.push(MetadataFanoutObservation::from_error(&err, fanout_start.elapsed()));
accumulator.observe_error(&err);
}
}
if let Some(decision) = accumulator
.early_stop_decision()
.or_else(|| accumulator.version_early_stop_decision())
{
let saved_responses = join_set.len();
join_set.abort_all();
rustfs_io_metrics::record_get_object_metadata_early_stop_hit(GET_OBJECT_PATH_LEGACY_DUPLEX, decision.reason);
rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(
GET_OBJECT_PATH_LEGACY_DUPLEX,
saved_responses,
);
let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations);
return Ok((ress, errors, diagnostics));
}
}
rustfs_io_metrics::record_get_object_metadata_early_stop_miss(
GET_OBJECT_PATH_LEGACY_DUPLEX,
accumulator.final_miss_reason(),
);
rustfs_io_metrics::record_get_object_metadata_early_stop_saved_responses(GET_OBJECT_PATH_LEGACY_DUPLEX, 0);
let diagnostics = MetadataFanoutDiagnostics::new(fanout_start.elapsed(), observations);
Ok((ress, errors, diagnostics))
}
pub async fn read_version_optimized(
&self,
bucket: &str,
object: &str,
version_id: &str,
opts: &ReadOptions,
) -> Result<Vec<FileInfo>> {
// Use existing disk selection logic
let disks = self.disks.read().await;
let required_reads = self.format.erasure.sets.len();
// Clone parameters outside the closure to avoid lifetime issues
let bucket = bucket.to_string();
let object = object.to_string();
let version_id = version_id.to_string();
let opts = opts.clone();
let processor = runtime_sources::batch_processors().read_processor();
let tasks: Vec<_> = disks
.iter()
.take(required_reads + 2) // Read a few extra for reliability
.filter_map(|disk| {
disk.as_ref().map(|d| {
let disk = d.clone();
let bucket = bucket.clone();
let object = object.clone();
let version_id = version_id.clone();
let opts = opts.clone();
async move { disk.read_version(&bucket, &bucket, &object, &version_id, &opts).await }
})
})
.collect();
match processor.execute_batch_with_quorum(tasks, required_reads).await {
Ok(results) => Ok(results),
Err(_) => Err(DiskError::FileNotFound.into()), // Use existing error type
}
}
pub(super) async fn read_all_xl(
disks: &[Option<DiskStore>],
bucket: &str,
object: &str,
read_data: bool,
incl_free_vers: bool,
) -> (Vec<FileInfo>, Vec<Option<DiskError>>) {
let (fileinfos, errs) = Self::read_all_raw_file_info(disks, bucket, object, read_data).await;
Self::pick_latest_quorum_files_info(fileinfos, errs, bucket, object, read_data, incl_free_vers).await
}
pub(crate) async fn load_file_info_versions_exact(
&self,
bucket: &str,
object: &str,
) -> Result<Option<rustfs_filemeta::FileInfoVersions>> {
let disks = self.get_disks_internal().await;
if disks.is_empty() {
return Err(to_object_err(StorageError::ErasureReadQuorum, vec![bucket, object]));
}
let read_quorum = disks.len().div_ceil(2).max(1);
let (raw_fileinfos, errs) = Self::read_all_raw_file_info(&disks, bucket, object, false).await;
if let Some(err) = reduce_read_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, read_quorum) {
let object_err = to_object_err(err.into(), vec![bucket, object]);
if is_err_object_not_found(&object_err) || is_err_version_not_found(&object_err) {
return Ok(None);
}
return Err(object_err);
}
let mut shallow_versions = Vec::with_capacity(raw_fileinfos.len());
for raw_fileinfo in raw_fileinfos.into_iter().flatten() {
let meta = FileMeta::load(&raw_fileinfo.buf)
.map_err(|err| Error::other(format!("exact object metadata decode failed for {bucket}/{object}: {err}")))?;
shallow_versions.push(meta.versions);
}
if shallow_versions.len() < read_quorum {
return Err(to_object_err(StorageError::ErasureReadQuorum, vec![bucket, object]));
}
let versions = merge_file_meta_versions(read_quorum, true, 0, &shallow_versions);
if versions.is_empty() {
return Err(Error::other(format!(
"exact object metadata read returned no quorum versions for {bucket}/{object}"
)));
}
FileMeta {
versions,
..Default::default()
}
.get_all_file_info_versions(bucket, object, true)
.map(Some)
.map_err(|err| Error::other(format!("exact object versions decode failed for {bucket}/{object}: {err}")))
}
pub(super) async fn read_all_raw_file_info(
disks: &[Option<DiskStore>],
bucket: &str,
object: &str,
read_data: bool,
) -> (Vec<Option<RawFileInfo>>, Vec<Option<DiskError>>) {
let mut ress = Vec::with_capacity(disks.len());
let mut errors = Vec::with_capacity(disks.len());
let mut futures = Vec::with_capacity(disks.len());
for disk in disks.iter() {
futures.push(async move {
if let Some(disk) = disk {
disk.read_xl(bucket, object, read_data).await
} else {
Err(DiskError::DiskNotFound)
}
});
}
let results = join_all(futures).await;
for result in results {
match result {
Ok(res) => {
ress.push(Some(res));
errors.push(None);
}
Err(e) => {
ress.push(None);
errors.push(Some(e));
}
}
}
(ress, errors)
}
pub(super) async fn pick_latest_quorum_files_info(
fileinfos: Vec<Option<RawFileInfo>>,
errs: Vec<Option<DiskError>>,
bucket: &str,
object: &str,
read_data: bool,
incl_free_vers: bool,
) -> (Vec<FileInfo>, Vec<Option<DiskError>>) {
let mut metadata_array = vec![None; fileinfos.len()];
let mut meta_file_infos = vec![FileInfo::default(); fileinfos.len()];
let mut metadata_shallow_versions = vec![None; fileinfos.len()];
let mut v2_bufs = {
if !read_data {
vec![Vec::new(); fileinfos.len()]
} else {
Vec::new()
}
};
let mut errs = errs;
for (idx, info_op) in fileinfos.iter().enumerate() {
if let Some(info) = info_op {
if !read_data {
v2_bufs[idx] = info.buf.clone();
}
let xlmeta = match FileMeta::load(&info.buf) {
Ok(res) => res,
Err(err) => {
errs[idx] = Some(err.into());
continue;
}
};
metadata_array[idx] = Some(xlmeta);
meta_file_infos[idx] = FileInfo::default();
}
}
for (idx, info_op) in metadata_array.iter().enumerate() {
if let Some(info) = info_op {
metadata_shallow_versions[idx] = Some(info.versions.clone());
}
}
let shallow_versions: Vec<Vec<FileMetaShallowVersion>> = metadata_shallow_versions.iter().flatten().cloned().collect();
let read_quorum = fileinfos.len().div_ceil(2);
let versions = merge_file_meta_versions(read_quorum, false, 1, &shallow_versions);
let meta = FileMeta {
versions,
..Default::default()
};
let finfo = match meta.into_fileinfo(bucket, object, "", true, incl_free_vers, true) {
Ok(res) => res,
Err(err) => {
for item in errs.iter_mut() {
if item.is_none() {
*item = Some(err.clone().into());
}
}
return (meta_file_infos, errs);
}
};
if !finfo.is_valid() {
for item in errs.iter_mut() {
if item.is_none() {
*item = Some(DiskError::FileCorrupt);
}
}
return (meta_file_infos, errs);
}
let vid = finfo.version_id.unwrap_or(Uuid::nil());
for (idx, meta_op) in metadata_array.iter().enumerate() {
if let Some(meta) = meta_op {
match meta.into_fileinfo(bucket, object, vid.to_string().as_str(), read_data, incl_free_vers, true) {
Ok(res) => meta_file_infos[idx] = res,
Err(err) => errs[idx] = Some(err.into()),
}
}
}
(meta_file_infos, errs)
}
pub(super) async fn read_multiple_files(
disks: &[Option<DiskStore>],
req: ReadMultipleReq,
read_quorum: usize,
) -> Vec<ReadMultipleResp> {
let mut futures = Vec::with_capacity(disks.len());
let empty_quorum_result = || {
req.files
.iter()
.map(|want| ReadMultipleResp {
bucket: req.bucket.clone(),
prefix: req.prefix.clone(),
file: want.clone(),
exists: false,
error: Error::ErasureReadQuorum.to_string(),
data: Vec::new(),
mod_time: None,
})
.collect::<Vec<_>>()
};
for disk in disks.iter() {
let disk = disk.clone();
let req = req.clone();
futures.push(async move {
if let Some(disk) = disk {
disk.read_multiple(req).await
} else {
Err(DiskError::DiskNotFound)
}
});
}
let (ress, errors) = match collect_read_multiple_results(futures, read_quorum).await {
Ok(collected) => collected,
Err(()) => return empty_quorum_result(),
};
// debug!("ReadMultipleResp ress {:?}", ress);
// debug!("ReadMultipleResp errors {:?}", errors);
let mut ret = Vec::with_capacity(req.files.len());
for want in req.files.iter() {
let mut quorum = 0;
let mut get_res = ReadMultipleResp::default();
for res in ress.iter() {
if res.is_none() {
continue;
}
let disk_res = res.as_ref().unwrap();
for resp in disk_res.iter() {
if !resp.error.is_empty() || !resp.exists {
continue;
}
if &resp.file != want || resp.bucket != req.bucket || resp.prefix != req.prefix {
continue;
}
quorum += 1;
if get_res.mod_time > resp.mod_time || get_res.data.len() > resp.data.len() {
continue;
}
get_res = resp.clone();
}
}
if quorum < read_quorum {
// debug!("quorum < read_quorum: {} < {}", quorum, read_quorum);
get_res.exists = false;
get_res.error = Error::ErasureReadQuorum.to_string();
get_res.data = Vec::new();
}
ret.push(get_res);
}
// log err
ret
}
#[tracing::instrument(level = "debug", skip(self))]
pub(super) async fn get_object_fileinfo(
&self,
bucket: &str,
object: &str,
opts: &ObjectOptions,
read_data: bool,
) -> Result<(FileInfo, Vec<FileInfo>, Vec<Option<DiskStore>>)> {
let vid = opts.version_id.clone().unwrap_or_default();
let use_metadata_cache = self.is_get_object_metadata_cache_enabled(bucket, opts, read_data).await;
if use_metadata_cache
&& vid.is_empty()
&& let Some(cached) = self.cached_get_object_fileinfo(bucket, object).await
{
return Ok((cached.fi, cached.parts_metadata, cached.online_disks));
}
let disks = self.disks.read().await;
let disks = disks.clone();
// TODO: optimize concurrency and break once enough slots are available
let (parts_metadata, errs, metadata_fanout_diagnostics) = Self::read_all_fileinfo_observed(
&disks,
"",
bucket,
object,
vid.as_str(),
read_data,
false,
opts.incl_free_versions,
self.default_parity_count,
)
.await?;
metadata_fanout_diagnostics.record(GET_OBJECT_PATH_LEGACY_DUPLEX);
// warn!("get_object_fileinfo parts_metadata {:?}", &parts_metadata);
// warn!("get_object_fileinfo {}/{} errs {:?}", bucket, object, &errs);
let _min_disks = self.set_drive_count - self.default_parity_count;
let (read_quorum, _) = match Self::object_quorum_from_meta(&parts_metadata, &errs, self.default_parity_count)
.map_err(|err| to_object_err(err.into(), vec![bucket, object]))
{
Ok(v) => v,
Err(e) => {
// error!("Self::object_quorum_from_meta: {:?}, bucket: {}, object: {}", &e, bucket, object);
return Err(e);
}
};
let read_quorum =
usize::try_from(read_quorum).map_err(|_| to_object_err(DiskError::ErasureReadQuorum.into(), vec![bucket, object]))?;
metadata_fanout_diagnostics.record_quorum_candidate_latency(GET_OBJECT_PATH_LEGACY_DUPLEX, read_quorum);
if let Some(err) = reduce_read_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, read_quorum) {
error!("reduce_read_quorum_errs: {:?}, bucket: {}, object: {}", &err, bucket, object);
return Err(to_object_err(err.into(), vec![bucket, object]));
}
let (op_online_disks, mot_time, etag) = Self::list_online_disks(&disks, &parts_metadata, &errs, read_quorum);
let fi = Self::pick_valid_fileinfo(&parts_metadata, mot_time, etag, read_quorum)?;
if errs.iter().any(|err| err.is_some()) {
let version_id = resolved_read_repair_version_id(&fi, opts.version_id.as_deref());
submit_read_repair_heal(
&fi.volume,
&fi.name,
version_id.as_deref(),
self.pool_index,
self.set_index,
None,
"metadata_read_error",
)
.await;
} else if use_metadata_cache {
self.cache_get_object_fileinfo(bucket, object, &fi, &parts_metadata, &op_online_disks, read_quorum)
.await;
}
// debug!("get_object_fileinfo pick fi {:?}", &fi);
// let online_disks: Vec<Option<DiskStore>> = op_online_disks.iter().filter(|v| v.is_some()).cloned().collect();
Ok((fi, parts_metadata, op_online_disks))
}
pub(super) async fn get_object_info_and_quorum(
&self,
bucket: &str,
object: &str,
opts: &ObjectOptions,
) -> (ObjectInfo, usize, Option<StorageError>) {
let fi = match self.get_object_fileinfo(bucket, object, opts, false).await {
Ok((fi, _, _)) => fi,
Err(e) => return (ObjectInfo::default(), 0, Some(e)),
};
let write_quorum = fi.write_quorum(self.default_write_quorum());
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 (
oi,
write_quorum,
Some(to_object_err(StorageError::MethodNotAllowed, vec![bucket, object])),
);
}
if fi.deleted {
if opts.incl_free_versions && fi.tier_free_version() && opts.version_id.is_some() {
return (oi, write_quorum, None);
}
return if opts.version_id.is_none() || opts.delete_marker {
(oi, write_quorum, Some(to_object_err(StorageError::FileNotFound, vec![bucket, object])))
} else {
(
oi,
write_quorum,
Some(to_object_err(StorageError::MethodNotAllowed, vec![bucket, object])),
)
};
}
(oi, write_quorum, None)
}
#[allow(clippy::too_many_arguments)]
pub(super) async fn get_object_with_fileinfo<W>(
// &self,
bucket: &str,
object: &str,
offset: usize,
length: i64,
writer: &mut W,
fi: FileInfo,
files: Vec<FileInfo>,
disks: &[Option<DiskStore>],
set_index: usize,
pool_index: usize,
skip_verify_bitrot: bool,
) -> Result<()>
where
W: AsyncWrite + Send + Sync + Unpin + 'static,
{
let pipeline_started = Instant::now();
debug!(bucket, object, requested_length = length, offset, "get_object_with_fileinfo start");
let (disks, files) = Self::shuffle_disks_and_parts_metadata_by_index(disks, &files, &fi);
let total_size = fi.size as usize;
if offset > total_size {
let reason = GetObjectFailureReason::RangeOrLengthInvalid;
record_get_object_pipeline_failure(GET_STAGE_RANGE, reason);
error!(
bucket,
object,
offset,
total_size,
requested_length = length,
stage = GET_STAGE_RANGE,
reason = reason.as_str(),
state = "range_or_length_invalid",
"GetObject range validation failed"
);
return Err(Error::other("offset out of range"));
}
let length = if length < 0 { total_size - offset } else { length as usize };
let Some(end_offset_exclusive) = offset.checked_add(length) else {
let reason = GetObjectFailureReason::RangeOrLengthInvalid;
record_get_object_pipeline_failure(GET_STAGE_RANGE, reason);
error!(
bucket,
object,
offset,
total_size,
requested_length = length,
stage = GET_STAGE_RANGE,
reason = reason.as_str(),
state = "range_or_length_invalid",
"GetObject range validation overflow"
);
return Err(Error::other("offset out of range"));
};
if end_offset_exclusive > total_size {
let reason = GetObjectFailureReason::RangeOrLengthInvalid;
record_get_object_pipeline_failure(GET_STAGE_RANGE, reason);
error!(
bucket,
object,
offset,
total_size,
requested_length = length,
end_offset_exclusive,
stage = GET_STAGE_RANGE,
reason = reason.as_str(),
state = "range_or_length_invalid",
"GetObject range validation failed"
);
return Err(Error::other("offset out of range"));
}
let (part_index, mut part_offset) = fi.to_part_offset(offset)?;
let mut end_offset = offset;
if length > 0 {
end_offset += length - 1
}
let (last_part_index, last_part_relative_offset) = fi.to_part_offset(end_offset)?;
debug!(
bucket,
object, offset, length, end_offset, part_index, last_part_index, last_part_relative_offset, "Multipart read bounds"
);
let erasure = crate::erasure::coding::Erasure::new_with_options(
fi.erasure.data_blocks,
fi.erasure.parity_blocks,
fi.erasure.block_size,
fi.uses_legacy_checksum,
);
let part_indices: Vec<usize> = (part_index..=last_part_index).collect();
debug!(bucket, object, ?part_indices, "Multipart part indices to stream");
let mut total_read = 0;
for current_part in part_indices {
if total_read == length {
debug!(
bucket,
object,
total_read,
requested_length = length,
part_index = current_part,
"Stopping multipart stream early because accumulated bytes match request"
);
break;
}
let part_number = fi.parts[current_part].number;
let part_size = fi.parts[current_part].size;
let mut part_length = part_size - part_offset;
if part_length > (length - total_read) {
part_length = length - total_read
}
let till_offset = erasure.shard_file_offset(part_offset, part_length, part_size);
let read_offset = (part_offset / erasure.block_size) * erasure.shard_size();
debug!(
bucket,
object,
part_index = current_part,
part_number,
part_offset,
part_size,
part_length,
read_offset,
till_offset,
total_read_before = total_read,
requested_length = length,
"Streaming multipart part"
);
let checksum_info = fi.erasure.get_checksum_info(part_number);
let checksum_algo =
if fi.uses_legacy_checksum && checksum_info.algorithm == rustfs_utils::HashAlgorithm::HighwayHash256S {
rustfs_utils::HashAlgorithm::HighwayHash256SLegacy
} else {
checksum_info.algorithm
};
let read_length = till_offset.saturating_sub(read_offset);
// Read zero-copy configuration from environment variable
// Default: enabled (true) for performance
let use_zero_copy = rustfs_utils::get_env_bool(ENV_OBJECT_ZERO_COPY_ENABLE, DEFAULT_OBJECT_ZERO_COPY_ENABLE);
let reader_setup_stage_start = Instant::now();
let read_costs = disks
.iter()
.map(|disk| shard_read_cost_for_disk(disk.as_ref()))
.collect::<Vec<_>>();
let reader_setup = create_bitrot_readers_until_quorum(
&files,
&disks,
bucket,
object,
part_number,
read_offset,
read_length,
erasure.shard_size(),
checksum_algo,
skip_verify_bitrot,
use_zero_copy,
erasure.data_shards,
erasure.parity_shards,
BitrotReaderSetupMode::ReadQuorum,
)
.await;
let reader_setup_elapsed = reader_setup_stage_start.elapsed();
rustfs_io_metrics::record_get_object_shard_reader_setup_duration(reader_setup_elapsed.as_secs_f64());
let setup_available_readers = reader_setup.available_shards();
if reader_setup_elapsed >= SLOW_OBJECT_READ_LOG_THRESHOLD {
warn!(
event = EVENT_SET_DISK_READ,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
bucket,
object,
part_index = current_part,
part_number,
read_offset,
read_length,
available_shards = setup_available_readers,
total_shards = reader_setup.errors.len(),
data_shards = erasure.data_shards,
parity_shards = erasure.parity_shards,
elapsed_ms = reader_setup_elapsed.as_millis(),
errors = ?reader_setup.errors,
state = "reader_setup_slow",
"Set disk object reader setup is slow"
);
}
let nil_count = reader_setup.available_shards();
if nil_count < erasure.data_shards {
if let Some(read_err) = reduce_read_quorum_errs(&reader_setup.errors, OBJECT_OP_IGNORED_ERRS, erasure.data_shards)
{
let reason = classify_disk_error(&read_err);
error!(
event = EVENT_SET_DISK_READ,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
bucket,
object,
part_index = current_part,
part_number,
part_offset,
part_length,
read_offset,
till_offset,
total_shards = erasure.data_shards + erasure.parity_shards,
available_shards = nil_count,
data_shards = erasure.data_shards,
stage = GET_STAGE_READER_SETUP,
reason = reason.as_str(),
errors = ?reader_setup.errors,
state = "read_quorum_unavailable",
"Create bitrot reader failed read quorum"
);
record_get_object_pipeline_failure(GET_STAGE_READER_SETUP, reason);
return Err(to_object_err(read_err.into(), vec![bucket, object]));
}
let reason = GetObjectFailureReason::ReadQuorum;
error!(
event = EVENT_SET_DISK_READ,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
bucket,
object,
part_index = current_part,
part_number,
part_offset,
part_length,
read_offset,
till_offset,
total_shards = erasure.data_shards + erasure.parity_shards,
available_shards = nil_count,
data_shards = erasure.data_shards,
stage = GET_STAGE_READER_SETUP,
reason = reason.as_str(),
errors = ?reader_setup.errors,
state = "not_enough_readers",
"Create bitrot reader did not have enough disks"
);
record_get_object_pipeline_failure(GET_STAGE_READER_SETUP, reason);
return Err(Error::other(format!("not enough disks to read: {:?}", reader_setup.errors)));
}
// Check if we have missing shards even though we can read successfully
// This happens when a node was offline during write and comes back online
let total_shards = erasure.data_shards + erasure.parity_shards;
let available_shards = nil_count;
let missing_shards = reader_setup.completed_failed_shards();
debug!(
bucket,
object,
part_number,
total_shards,
available_shards,
attempted_shards = reader_setup.attempted.iter().filter(|attempted| **attempted).count(),
missing_shards,
data_shards = erasure.data_shards,
parity_shards = erasure.parity_shards,
"Shard availability check"
);
if missing_shards > 0 && available_shards >= erasure.data_shards {
// We have missing shards but enough to read - trigger background heal
debug!(
bucket,
object,
part_number,
missing_shards,
available_shards,
pool_index,
set_index,
"Detected missing shards during read, triggering background heal"
);
let version_id = fi.version_id.as_ref().map(ToString::to_string);
submit_read_repair_heal(
bucket,
object,
version_id.as_deref(),
pool_index,
set_index,
Some(part_number),
"missing_shards",
)
.await;
}
// debug!(
// "read part {} part_offset {},part_length {},part_size {} ",
// part_number, part_offset, part_length, part_size
// );
let decode_stage_start = Instant::now();
let unattempted_data_shards = !reader_setup.data_shards_attempted(erasure.data_shards);
let readers = reader_setup.readers;
let (written, err) = erasure
.decode_with_read_costs(writer, readers, part_offset, part_length, part_size, read_costs)
.await;
let decode_elapsed = decode_stage_start.elapsed();
rustfs_io_metrics::record_get_object_decode_duration(decode_elapsed.as_secs_f64());
if decode_elapsed >= SLOW_OBJECT_READ_LOG_THRESHOLD || err.is_some() {
warn!(
event = EVENT_SET_DISK_READ,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
bucket,
object,
part_index = current_part,
part_number,
part_offset,
part_size,
part_length,
bytes_written = written,
available_shards,
missing_shards,
elapsed_ms = decode_elapsed.as_millis(),
error = ?err,
state = if err.is_some() { "decode_failed_or_slow" } else { "decode_slow" },
"Set disk object decode stage is slow or failed"
);
}
debug!(
bucket,
object,
part_index = current_part,
part_number,
part_length,
bytes_written = written,
"Finished decoding multipart part"
);
if let Some(e) = err {
let de_err: DiskError = e.into();
let mut has_err = true;
if written == part_length {
let should_enqueue_heal = matches!(de_err, DiskError::FileCorrupt)
|| (matches!(de_err, DiskError::FileNotFound) && !unattempted_data_shards);
if should_enqueue_heal {
debug!(
bucket,
object,
part_number,
error = ?de_err,
"Recoverable decode error triggered read repair"
);
let version_id = fi.version_id.as_ref().map(ToString::to_string);
submit_read_repair_heal(
bucket,
object,
version_id.as_deref(),
pool_index,
set_index,
Some(part_number),
"decode_error",
)
.await;
has_err = false;
}
}
if has_err {
let reason = classify_disk_error(&de_err);
error!(
bucket,
object,
part_index = current_part,
part_number,
part_offset,
part_length,
bytes_written = written,
stage = GET_STAGE_DECODE,
reason = reason.as_str(),
error = ?de_err,
"Erasure decode failed during GetObject"
);
record_get_object_pipeline_failure(GET_STAGE_DECODE, reason);
return Err(de_err.into());
}
}
// debug!("ec decode {} written size {}", part_number, n);
total_read += part_length;
part_offset = 0;
}
// debug!("read end");
debug!(bucket, object, total_read, expected_length = length, "Multipart read finished");
let pipeline_elapsed = pipeline_started.elapsed();
if pipeline_elapsed >= SLOW_OBJECT_READ_LOG_THRESHOLD {
warn!(
event = EVENT_SET_DISK_READ,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
bucket,
object,
offset,
total_read,
expected_length = length,
elapsed_ms = pipeline_elapsed.as_millis(),
state = "pipeline_slow",
"Set disk object read pipeline is slow"
);
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub(super) async fn get_object_decode_reader_with_fileinfo(
bucket: &str,
object: &str,
fi: &FileInfo,
files: &[FileInfo],
disks: &[Option<DiskStore>],
_set_index: usize,
_pool_index: usize,
skip_verify_bitrot: bool,
) -> Result<GetCodecStreamingReaderBuildOutcome> {
let (disks, files) = Self::shuffle_disks_and_parts_metadata_by_index(disks, files, fi);
if fi.parts.len() != 1 {
return Err(Error::other("codec streaming reader only supports single-part plain objects"));
}
let erasure = crate::erasure::coding::Erasure::new_with_options(
fi.erasure.data_blocks,
fi.erasure.parity_blocks,
fi.erasure.block_size,
fi.uses_legacy_checksum,
);
let part = &fi.parts[0];
let part_number = part.number;
let part_size = part.size;
let part_length = usize::try_from(fi.size).map_err(|_| Error::other("codec streaming reader object size is invalid"))?;
if part_length > part_size {
return Err(Error::other("codec streaming reader part length exceeds part size"));
}
let checksum_info = fi.erasure.get_checksum_info(part_number);
let checksum_algo = if fi.uses_legacy_checksum && checksum_info.algorithm == rustfs_utils::HashAlgorithm::HighwayHash256S
{
rustfs_utils::HashAlgorithm::HighwayHash256SLegacy
} else {
checksum_info.algorithm
};
let use_zero_copy = rustfs_utils::get_env_bool(ENV_OBJECT_ZERO_COPY_ENABLE, DEFAULT_OBJECT_ZERO_COPY_ENABLE);
let till_offset = erasure.shard_file_offset(0, part_length, part_size);
let reader_setup_stage_start = Instant::now();
let read_costs = disks
.iter()
.map(|disk| shard_read_cost_for_disk(disk.as_ref()))
.collect::<Vec<_>>();
let reader_setup = create_bitrot_readers_until_quorum(
&files,
&disks,
bucket,
object,
part_number,
0,
till_offset,
erasure.shard_size(),
checksum_algo,
skip_verify_bitrot,
use_zero_copy,
erasure.data_shards,
erasure.parity_shards,
BitrotReaderSetupMode::VerifyReconstruction,
)
.await;
let metrics_path = get_codec_streaming_metrics_path();
rustfs_io_metrics::record_get_object_stage_duration(
metrics_path,
GET_STAGE_READER_SETUP,
reader_setup_stage_start.elapsed().as_secs_f64(),
);
let available_shards = reader_setup.available_shards();
if available_shards < erasure.data_shards {
if let Some(read_err) = reduce_read_quorum_errs(&reader_setup.errors, OBJECT_OP_IGNORED_ERRS, erasure.data_shards) {
let reason = classify_disk_error(&read_err);
record_get_object_pipeline_failure_for_path(metrics_path, GET_STAGE_READER_SETUP, reason);
return Err(to_object_err(read_err.into(), vec![bucket, object]));
}
record_get_object_pipeline_failure_for_path(metrics_path, GET_STAGE_READER_SETUP, GetObjectFailureReason::ReadQuorum);
return Err(Error::other(format!("not enough disks to read: {:?}", reader_setup.errors)));
}
let missing_shards = reader_setup.completed_failed_shards();
if let Some(reason) = codec_streaming_reader_setup_fallback_reason(missing_shards) {
return Ok(GetCodecStreamingReaderBuildOutcome::Fallback(reason));
}
let readers = reader_setup.readers;
let source =
crate::erasure::coding::decode::ParallelReader::new_with_metrics_path_read_costs_and_reconstruction_verification(
readers,
erasure.clone(),
0,
part_size,
Some(metrics_path),
read_costs,
);
let engine = build_get_codec_streaming_decode_engine(erasure)?;
let reader = crate::erasure::coding::decode_reader::ErasureDecodeReader::new_with_metrics_path(
source,
engine,
part_length,
metrics_path,
)?;
Ok(GetCodecStreamingReaderBuildOutcome::Reader(Box::new(
crate::erasure::coding::decode_reader::SyncErasureDecodeReader::new_with_metrics_path(reader, metrics_path),
)))
}
}
fn is_get_object_metadata_cache_request_eligible(bucket: &str, opts: &ObjectOptions, read_data: bool) -> bool {
read_data
&& !opts.no_lock
&& opts.version_id.is_none()
&& !opts.versioned
&& !opts.version_suspended
&& !opts.incl_free_versions
&& !opts.delete_marker
&& opts.part_number.is_none()
&& !opts.data_movement
&& !opts.raw_data_movement_read
&& !bucket.starts_with(RUSTFS_META_BUCKET)
}
#[cfg(test)]
mod metadata_cache_tests {
use super::*;
use rustfs_common::heal_channel::HealAdmissionDropReason;
use serial_test::serial;
use std::sync::atomic::{AtomicUsize, Ordering};
static SLOW_READ_REPAIR_SUBMITTER_CALLS: AtomicUsize = AtomicUsize::new(0);
static DROPPED_READ_REPAIR_SUBMITTER_CALLS: AtomicUsize = AtomicUsize::new(0);
fn slow_read_repair_submitter(_request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture {
SLOW_READ_REPAIR_SUBMITTER_CALLS.fetch_add(1, Ordering::Relaxed);
Box::pin(async {
tokio::time::sleep(Duration::from_millis(250)).await;
ReadRepairAdmissionOutcome::Response(HealAdmissionResult::Accepted)
})
}
fn dropped_read_repair_submitter(_request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture {
DROPPED_READ_REPAIR_SUBMITTER_CALLS.fetch_add(1, Ordering::Relaxed);
Box::pin(async {
ReadRepairAdmissionOutcome::Response(HealAdmissionResult::Dropped(HealAdmissionDropReason::PolicyDropped))
})
}
async fn new_metadata_cache_test_set() -> Arc<SetDisks> {
SetDisks::new(
"metadata-cache-test".to_string(),
Arc::new(RwLock::new(Vec::new())),
4,
2,
0,
0,
Vec::new(),
FormatV3::new(1, 4),
Vec::new(),
)
.await
}
fn valid_test_fileinfo(object: &str) -> FileInfo {
let mut fi = FileInfo::new(object, 2, 2);
fi.volume = "bucket".to_string();
fi.name = object.to_string();
fi.size = 1;
fi.erasure.index = 1;
fi.metadata.insert("etag".to_string(), "etag-1".to_string());
fi
}
#[test]
fn get_object_metadata_cache_request_eligibility_is_conservative() {
let opts = ObjectOptions::default();
assert!(is_get_object_metadata_cache_request_eligible("bucket", &opts, true));
assert!(!is_get_object_metadata_cache_request_eligible("bucket", &opts, false));
assert!(!is_get_object_metadata_cache_request_eligible(RUSTFS_META_BUCKET, &opts, true));
let mut opts = ObjectOptions {
no_lock: true,
..Default::default()
};
assert!(!is_get_object_metadata_cache_request_eligible("bucket", &opts, true));
opts = ObjectOptions {
version_id: Some("version".to_string()),
..Default::default()
};
assert!(!is_get_object_metadata_cache_request_eligible("bucket", &opts, true));
opts = ObjectOptions {
versioned: true,
..Default::default()
};
assert!(!is_get_object_metadata_cache_request_eligible("bucket", &opts, true));
opts = ObjectOptions {
version_suspended: true,
..Default::default()
};
assert!(!is_get_object_metadata_cache_request_eligible("bucket", &opts, true));
opts = ObjectOptions {
incl_free_versions: true,
..Default::default()
};
assert!(!is_get_object_metadata_cache_request_eligible("bucket", &opts, true));
opts = ObjectOptions {
delete_marker: true,
..Default::default()
};
assert!(!is_get_object_metadata_cache_request_eligible("bucket", &opts, true));
opts = ObjectOptions {
part_number: Some(1),
..Default::default()
};
assert!(!is_get_object_metadata_cache_request_eligible("bucket", &opts, true));
opts = ObjectOptions {
data_movement: true,
..Default::default()
};
assert!(!is_get_object_metadata_cache_request_eligible("bucket", &opts, true));
opts = ObjectOptions {
raw_data_movement_read: true,
..Default::default()
};
assert!(!is_get_object_metadata_cache_request_eligible("bucket", &opts, true));
}
#[tokio::test]
async fn read_repair_heal_dedupes_same_object_version() {
let bucket = format!("bucket-{}", Uuid::new_v4());
assert!(reserve_read_repair_heal(&bucket, "object", None, 0, 0).await.is_some());
assert!(reserve_read_repair_heal(&bucket, "object", None, 0, 0).await.is_none());
assert!(
reserve_read_repair_heal(&bucket, "object", Some("version-2"), 0, 0)
.await
.is_some()
);
assert!(reserve_read_repair_heal(&bucket, "object", None, 0, 1).await.is_some());
}
#[tokio::test]
async fn read_repair_heal_reservation_can_be_released_after_rejection() {
let bucket = format!("bucket-{}", Uuid::new_v4());
let key = reserve_read_repair_heal(&bucket, "object", None, 0, 0)
.await
.expect("first reservation should be accepted");
assert!(reserve_read_repair_heal(&bucket, "object", None, 0, 0).await.is_none());
release_read_repair_heal_reservation(&key).await;
assert!(reserve_read_repair_heal(&bucket, "object", None, 0, 0).await.is_some());
}
#[tokio::test]
#[serial]
async fn submit_read_repair_heal_does_not_wait_for_slow_admission() {
SLOW_READ_REPAIR_SUBMITTER_CALLS.store(0, Ordering::Relaxed);
let bucket = format!("bucket-{}", Uuid::new_v4());
let started = Instant::now();
submit_read_repair_heal_with_submitter(
ReadRepairHealSubmission {
bucket: &bucket,
object: "object",
version_id: None,
pool_index: 0,
set_index: 0,
part_number: Some(1),
reason: "missing_shards",
},
slow_read_repair_submitter,
)
.await;
assert!(
started.elapsed() < Duration::from_millis(50),
"read-repair submission should not wait for admission response"
);
tokio::time::timeout(Duration::from_secs(1), async {
while SLOW_READ_REPAIR_SUBMITTER_CALLS.load(Ordering::Relaxed) == 0 {
tokio::time::sleep(Duration::from_millis(5)).await;
}
})
.await
.expect("background read-repair submitter should be called");
}
#[tokio::test]
#[serial]
async fn submit_read_repair_heal_releases_reservation_after_policy_drop() {
DROPPED_READ_REPAIR_SUBMITTER_CALLS.store(0, Ordering::Relaxed);
let bucket = format!("bucket-{}", Uuid::new_v4());
submit_read_repair_heal_with_submitter(
ReadRepairHealSubmission {
bucket: &bucket,
object: "object",
version_id: None,
pool_index: 0,
set_index: 0,
part_number: Some(1),
reason: "missing_shards",
},
dropped_read_repair_submitter,
)
.await;
let released_key = tokio::time::timeout(Duration::from_secs(1), async {
loop {
if let Some(key) = reserve_read_repair_heal(&bucket, "object", None, 0, 0).await {
break key;
}
tokio::time::sleep(Duration::from_millis(5)).await;
}
})
.await
.expect("read-repair reservation should be released after policy drop");
assert_eq!(DROPPED_READ_REPAIR_SUBMITTER_CALLS.load(Ordering::Relaxed), 1);
release_read_repair_heal_reservation(&released_key).await;
}
#[test]
fn resolved_read_repair_version_prefers_selected_fileinfo_version() {
let mut fi = valid_test_fileinfo("object");
fi.version_id = Some(Uuid::parse_str("00000000-0000-0000-0000-000000000001").unwrap());
assert_eq!(
resolved_read_repair_version_id(&fi, Some("00000000-0000-0000-0000-000000000002")).as_deref(),
Some("00000000-0000-0000-0000-000000000001")
);
fi.version_id = None;
assert_eq!(
resolved_read_repair_version_id(&fi, Some("00000000-0000-0000-0000-000000000002")).as_deref(),
Some("00000000-0000-0000-0000-000000000002")
);
}
#[tokio::test]
async fn get_object_metadata_cache_hit_returns_stored_metadata() {
let set = new_metadata_cache_test_set().await;
let fi = valid_test_fileinfo("object");
let parts_metadata = vec![fi.clone()];
let online_disks = Vec::new();
set.cache_get_object_fileinfo("bucket", "object", &fi, &parts_metadata, &online_disks, 0)
.await;
let cached = set
.cached_get_object_fileinfo("bucket", "object")
.await
.expect("fresh cache entry should be returned");
assert_eq!(cached.fi.name, "object");
assert_eq!(cached.parts_metadata.len(), 1);
assert_eq!(cached.online_disks.len(), 0);
assert_eq!(cached.read_quorum, 0);
}
#[tokio::test]
async fn get_object_metadata_cache_rejects_deleted_and_invalid_fileinfo() {
let set = new_metadata_cache_test_set().await;
let mut deleted = valid_test_fileinfo("deleted-object");
deleted.deleted = true;
set.cache_get_object_fileinfo("bucket", "deleted-object", &deleted, &[deleted.clone()], &[], 0)
.await;
assert!(
set.cached_get_object_fileinfo("bucket", "deleted-object").await.is_none(),
"deleted metadata must not be cached"
);
let invalid = FileInfo::default();
set.cache_get_object_fileinfo("bucket", "invalid-object", &invalid, std::slice::from_ref(&invalid), &[], 0)
.await;
assert!(
set.cached_get_object_fileinfo("bucket", "invalid-object").await.is_none(),
"invalid metadata must not be cached"
);
}
#[tokio::test]
async fn get_object_metadata_cache_requires_cached_read_quorum() {
let set = new_metadata_cache_test_set().await;
let fi = valid_test_fileinfo("object");
set.get_object_metadata_cache
.insert(
GetObjectMetadataCacheKey::new("bucket", "object"),
GetObjectMetadataCacheEntry {
created_at: Instant::now(),
fi: fi.clone(),
parts_metadata: vec![fi],
online_disks: vec![None],
read_quorum: 1,
},
)
.await;
assert!(
set.cached_get_object_fileinfo("bucket", "object").await.is_none(),
"cache hit must be rejected when cached online disks cannot satisfy read quorum"
);
}
#[tokio::test]
async fn get_object_metadata_cache_rejects_stale_entries() {
// moka handles TTL expiry automatically via time_to_live(250ms).
// This test verifies that entries inserted with the cache API are retrievable
// while fresh, and that the cache API works correctly.
let set = new_metadata_cache_test_set().await;
let fi = valid_test_fileinfo("object");
set.cache_get_object_fileinfo("bucket", "object", &fi, std::slice::from_ref(&fi), &[], 0)
.await;
assert!(
set.cached_get_object_fileinfo("bucket", "object").await.is_some(),
"freshly inserted entry should be returned"
);
}
#[tokio::test]
async fn get_object_metadata_cache_invalidation_removes_object_entry() {
let set = new_metadata_cache_test_set().await;
let fi = valid_test_fileinfo("object");
set.cache_get_object_fileinfo("bucket", "object", &fi, std::slice::from_ref(&fi), &[], 0)
.await;
assert!(set.cached_get_object_fileinfo("bucket", "object").await.is_some());
set.invalidate_get_object_metadata_cache("bucket", "object").await;
assert!(
set.cached_get_object_fileinfo("bucket", "object").await.is_none(),
"explicit invalidation must remove the cached object metadata"
);
}
#[tokio::test]
async fn get_object_metadata_cache_prunes_when_capacity_is_reached() {
// moka handles capacity eviction automatically via max_capacity(1024).
// This test verifies that the cache can hold entries and that insertion works.
let set = new_metadata_cache_test_set().await;
let fresh_fi = valid_test_fileinfo("fresh-object");
set.cache_get_object_fileinfo("bucket", "fresh-object", &fresh_fi, std::slice::from_ref(&fresh_fi), &[], 0)
.await;
assert!(
set.cached_get_object_fileinfo("bucket", "fresh-object").await.is_some(),
"freshly inserted entry should be retrievable"
);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn metadata_fanout_test_fileinfo(object: &str) -> FileInfo {
let mut fi = FileInfo::new(object, 2, 2);
fi.volume = "bucket".to_string();
fi.name = object.to_string();
fi.size = 1;
fi.erasure.index = 1;
fi.metadata.insert("etag".to_string(), "etag-1".to_string());
fi
}
fn codec_streaming_test_fileinfo(size: i64, part_count: usize) -> FileInfo {
let mut fi = FileInfo::new("object", 4, 2);
fi.volume = "bucket".to_string();
fi.name = "object".to_string();
fi.size = size;
fi.metadata.insert("etag".to_string(), "etag-1".to_string());
let size = usize::try_from(size).expect("test object size should fit usize");
let part_count = part_count.max(1);
let part_size = size.div_ceil(part_count);
for index in 0..part_count {
let remaining = size.saturating_sub(index * part_size);
let current_part_size = remaining.min(part_size);
fi.add_object_part(
index + 1,
format!("etag-{index}"),
current_part_size,
None,
i64::try_from(current_part_size).expect("test part size should fit i64"),
None,
None,
);
}
fi
}
fn read_part_test_part(number: usize, etag: &str) -> ObjectPartInfo {
ObjectPartInfo {
number,
etag: etag.to_string(),
..Default::default()
}
}
fn read_part_test_missing(number: usize) -> ObjectPartInfo {
ObjectPartInfo {
number,
error: Some("file not found".to_string()),
..Default::default()
}
}
#[test]
fn resolve_read_part_returns_part_when_etag_reaches_quorum() {
let part = read_part_test_part(1, "etag-1");
let responses = vec![
Some(vec![part.clone()]),
Some(vec![part]),
Some(vec![read_part_test_missing(1)]),
None,
];
let resolved = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2)
.expect("etag quorum should resolve part");
assert_eq!(resolved.etag, "etag-1");
assert!(resolved.error.is_none());
}
#[test]
fn resolve_read_part_returns_missing_only_when_missing_reaches_quorum() {
let responses = vec![
Some(vec![read_part_test_missing(1)]),
Some(vec![read_part_test_missing(1)]),
None,
];
let resolved = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2)
.expect("missing quorum should resolve as a confirmed missing part");
assert_eq!(resolved.number, 1);
assert_eq!(resolved.error.as_deref(), Some("part.1 not found"));
}
#[test]
fn resolve_read_part_treats_os_not_found_as_confirmed_missing() {
let responses = vec![
Some(vec![ObjectPartInfo {
number: 9999,
error: Some("No such file or directory (os error 2)".to_string()),
..Default::default()
}]),
Some(vec![ObjectPartInfo {
number: 9999,
error: Some("No such file or directory (os error 2)".to_string()),
..Default::default()
}]),
Some(vec![read_part_test_part(1, "stale-etag")]),
None,
];
let resolved = resolve_read_part_from_responses("bucket", "upload/part.9999.meta", 9999, 0, 1, &responses, 2)
.expect("OS not-found quorum should resolve as a missing part");
assert_eq!(resolved.number, 9999);
assert_eq!(resolved.error.as_deref(), Some("part.9999 not found"));
}
#[test]
fn resolve_read_part_returns_missing_when_missing_quorum_beats_stale_present_part() {
let responses = vec![
Some(vec![read_part_test_missing(1)]),
Some(vec![read_part_test_missing(1)]),
Some(vec![read_part_test_part(1, "stale-etag")]),
None,
];
let resolved = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2)
.expect("confirmed missing quorum should resolve as a missing part despite stale metadata");
assert_eq!(resolved.number, 1);
assert_eq!(resolved.error.as_deref(), Some("part.1 not found"));
}
#[test]
fn resolve_read_part_preserves_read_quorum_when_present_part_lacks_quorum() {
let responses = vec![
Some(vec![read_part_test_part(1, "etag-1")]),
Some(vec![read_part_test_missing(1)]),
None,
];
let err = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2)
.expect_err("mixed present and missing observations should not become InvalidPart");
assert_eq!(err, DiskError::ErasureReadQuorum);
}
#[test]
fn resolve_read_part_does_not_count_mismatched_response_as_missing() {
let responses = vec![Some(Vec::new()), Some(vec![read_part_test_missing(1)]), None];
let err = resolve_read_part_from_responses("bucket", "upload/part.1.meta", 1, 0, 1, &responses, 2)
.expect_err("invalid disk responses must not vote for a missing part");
assert_eq!(err, DiskError::ErasureReadQuorum);
}
#[test]
fn metadata_fanout_diagnostics_classifies_response_outcomes() {
let valid = metadata_fanout_test_fileinfo("object");
let invalid = FileInfo::default();
let observations = vec![
MetadataFanoutObservation::from_file_info(&valid, Duration::from_millis(3)),
MetadataFanoutObservation::from_file_info(&invalid, Duration::from_millis(4)),
MetadataFanoutObservation::from_error(&DiskError::FileNotFound, Duration::from_millis(5)),
MetadataFanoutObservation::from_error(&DiskError::FileVersionNotFound, Duration::from_millis(6)),
MetadataFanoutObservation::from_error(&DiskError::DiskNotFound, Duration::from_millis(7)),
MetadataFanoutObservation::from_error(&DiskError::FileCorrupt, Duration::from_millis(8)),
MetadataFanoutObservation::from_error(&DiskError::Timeout, Duration::from_millis(9)),
MetadataFanoutObservation::from_error(&DiskError::FaultyDisk, Duration::from_millis(10)),
MetadataFanoutObservation::from_error(&DiskError::Unexpected, Duration::from_millis(11)),
];
let diagnostics = MetadataFanoutDiagnostics::new(Duration::from_millis(12), observations);
let outcomes = diagnostics
.observations
.iter()
.map(|observation| observation.outcome)
.collect::<Vec<_>>();
assert_eq!(
outcomes,
vec![
GET_METADATA_RESPONSE_VALID,
GET_METADATA_RESPONSE_ERROR,
GET_METADATA_RESPONSE_NOT_FOUND,
GET_METADATA_RESPONSE_VERSION_NOT_FOUND,
GET_METADATA_RESPONSE_DISK_NOT_FOUND,
GET_METADATA_RESPONSE_CORRUPT,
GET_METADATA_RESPONSE_TIMEOUT,
GET_METADATA_RESPONSE_IGNORED,
GET_METADATA_RESPONSE_ERROR,
]
);
assert_eq!(diagnostics.total_responses(), 9);
assert_eq!(diagnostics.valid_responses(), 1);
assert_eq!(diagnostics.error_responses(), 8);
}
#[test]
fn metadata_fanout_diagnostics_tracks_ignored_errors_without_hiding_outcomes() {
let diagnostics = MetadataFanoutDiagnostics::new(
Duration::from_millis(15),
vec![
MetadataFanoutObservation::from_error(&DiskError::DiskNotFound, Duration::from_millis(1)),
MetadataFanoutObservation::from_error(&DiskError::FaultyRemoteDisk, Duration::from_millis(2)),
MetadataFanoutObservation::from_error(&DiskError::FileNotFound, Duration::from_millis(3)),
],
);
assert_eq!(diagnostics.ignored_responses(), 2);
assert_eq!(diagnostics.error_responses(), 3);
assert_eq!(diagnostics.observations[0].outcome, GET_METADATA_RESPONSE_DISK_NOT_FOUND);
assert_eq!(diagnostics.observations[1].outcome, GET_METADATA_RESPONSE_IGNORED);
assert_eq!(diagnostics.observations[2].outcome, GET_METADATA_RESPONSE_NOT_FOUND);
}
#[test]
fn metadata_fanout_quorum_candidate_latency_ignores_slow_trailing_valid_response() {
let valid = metadata_fanout_test_fileinfo("object");
let diagnostics = MetadataFanoutDiagnostics::new(
Duration::from_millis(250),
vec![
MetadataFanoutObservation::from_error(&DiskError::FileNotFound, Duration::from_millis(2)),
MetadataFanoutObservation::from_file_info(&valid, Duration::from_millis(5)),
MetadataFanoutObservation::from_file_info(&valid, Duration::from_millis(7)),
MetadataFanoutObservation::from_file_info(&valid, Duration::from_millis(250)),
],
);
assert_eq!(diagnostics.first_response_latency(), Some(Duration::from_millis(2)));
assert_eq!(diagnostics.first_valid_response_latency(), Some(Duration::from_millis(5)));
assert_eq!(diagnostics.slowest_response_latency(), Some(Duration::from_millis(250)));
assert_eq!(diagnostics.quorum_candidate_latency(2), Some(Duration::from_millis(7)));
}
#[test]
fn metadata_fanout_quorum_candidate_latency_requires_enough_valid_responses() {
let valid = metadata_fanout_test_fileinfo("object");
let diagnostics = MetadataFanoutDiagnostics::new(
Duration::from_millis(8),
vec![
MetadataFanoutObservation::from_file_info(&valid, Duration::from_millis(5)),
MetadataFanoutObservation::from_error(&DiskError::FileNotFound, Duration::from_millis(6)),
],
);
assert_eq!(diagnostics.quorum_candidate_latency(0), Some(Duration::ZERO));
assert_eq!(diagnostics.quorum_candidate_latency(1), Some(Duration::from_millis(5)));
assert_eq!(diagnostics.quorum_candidate_latency(2), None);
}
#[test]
fn metadata_fanout_counts_delete_marker_and_conflicting_versions_as_valid_observations_only() {
let mut latest = metadata_fanout_test_fileinfo("object");
latest.version_id = Some(Uuid::parse_str("00000000-0000-0000-0000-000000000001").expect("static uuid should parse"));
let mut historical = metadata_fanout_test_fileinfo("object");
historical.version_id = Some(Uuid::parse_str("00000000-0000-0000-0000-000000000002").expect("static uuid should parse"));
let mut delete_marker = metadata_fanout_test_fileinfo("object");
delete_marker.deleted = true;
delete_marker.version_id =
Some(Uuid::parse_str("00000000-0000-0000-0000-000000000003").expect("static uuid should parse"));
let diagnostics = MetadataFanoutDiagnostics::new(
Duration::from_millis(9),
vec![
MetadataFanoutObservation::from_file_info(&latest, Duration::from_millis(3)),
MetadataFanoutObservation::from_file_info(&historical, Duration::from_millis(4)),
MetadataFanoutObservation::from_file_info(&delete_marker, Duration::from_millis(5)),
],
);
assert_eq!(diagnostics.total_responses(), 3);
assert_eq!(diagnostics.valid_responses(), 3);
assert_eq!(diagnostics.error_responses(), 0);
assert!(
diagnostics
.observations
.iter()
.all(|observation| observation.outcome == GET_METADATA_RESPONSE_VALID)
);
}
fn metadata_early_stop_accumulator() -> MetadataQuorumAccumulator {
MetadataQuorumAccumulator::new(4, 2, true)
}
fn metadata_early_stop_candidate(object: &str, disk_index: usize) -> FileInfo {
let mut fi = metadata_fanout_test_fileinfo(object);
fi.erasure.index = disk_index;
fi.add_object_part(1, "part-etag-1".to_string(), 1, None, 1, None, None);
fi
}
#[test]
fn metadata_quorum_accumulator_hits_all_valid_same_version() {
let mut accumulator = metadata_early_stop_accumulator();
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 1));
assert!(accumulator.early_stop_decision().is_none());
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 2));
assert_eq!(
accumulator.early_stop_decision(),
Some(MetadataEarlyStopDecision {
reason: GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM
})
);
assert_eq!(accumulator.valid_responses, 2);
}
#[test]
fn metadata_quorum_accumulator_hits_quorum_valid_with_slow_trailing_disk() {
let mut accumulator = metadata_early_stop_accumulator();
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 1));
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 2));
assert!(accumulator.early_stop_decision().is_some());
assert_eq!(accumulator.candidate_votes, 2);
}
#[test]
fn metadata_quorum_accumulator_partial_result_remains_quorum_compatible() {
let first = metadata_early_stop_candidate("object", 1);
let second = metadata_early_stop_candidate("object", 2);
let parts_metadata = vec![first.clone(), second, FileInfo::default(), FileInfo::default()];
let errs = vec![None, None, None, None];
let (read_quorum, _) = SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 2)
.expect("partial early-stop metadata should preserve read quorum");
let read_quorum = usize::try_from(read_quorum).expect("read quorum should be non-negative");
let selected = SetDisks::pick_valid_fileinfo(&parts_metadata, None, Some("etag-1".to_string()), read_quorum)
.expect("partial early-stop metadata should preserve selected FileInfo");
assert_eq!(read_quorum, 2);
assert_eq!(selected.name, first.name);
assert_eq!(selected.get_etag(), first.get_etag());
}
#[test]
fn metadata_quorum_accumulator_falls_back_on_conflicting_versions() {
let mut accumulator = metadata_early_stop_accumulator();
let first = metadata_early_stop_candidate("object", 1);
let mut second = metadata_early_stop_candidate("object", 2);
second.version_id = Some(Uuid::parse_str("00000000-0000-0000-0000-000000000002").expect("static uuid should parse"));
accumulator.observe_file_info(&first);
accumulator.observe_file_info(&second);
assert!(accumulator.early_stop_decision().is_none());
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA);
}
#[test]
fn metadata_quorum_accumulator_falls_back_on_explicit_version_quorum() {
let mut accumulator = MetadataQuorumAccumulator::new(4, 2, false);
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 1));
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 2));
assert!(accumulator.early_stop_decision().is_none());
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_UNSAFE_REQUEST);
}
#[test]
fn metadata_quorum_accumulator_hits_delete_marker_quorum_early_stop() {
let mut accumulator = metadata_early_stop_accumulator();
let mut deleted = metadata_early_stop_candidate("object", 1);
deleted.deleted = true;
accumulator.observe_file_info(&deleted);
accumulator.observe_file_info(&deleted);
assert_eq!(
accumulator.early_stop_decision(),
Some(MetadataEarlyStopDecision {
reason: GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER
})
);
}
#[test]
fn metadata_quorum_accumulator_falls_back_on_delete_marker_below_quorum() {
let mut accumulator = metadata_early_stop_accumulator();
let mut deleted = metadata_early_stop_candidate("object", 1);
deleted.deleted = true;
accumulator.observe_file_info(&deleted);
assert!(accumulator.early_stop_decision().is_none());
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER);
}
#[test]
fn metadata_quorum_accumulator_falls_back_on_object_not_found_quorum() {
let mut accumulator = metadata_early_stop_accumulator();
accumulator.observe_error(&DiskError::FileNotFound);
accumulator.observe_error(&DiskError::VolumeNotFound);
assert!(accumulator.early_stop_decision().is_none());
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_NOT_FOUND);
}
#[test]
fn metadata_quorum_accumulator_falls_back_on_version_not_found_quorum() {
let mut accumulator = metadata_early_stop_accumulator();
accumulator.observe_error(&DiskError::FileVersionNotFound);
accumulator.observe_error(&DiskError::FileVersionNotFound);
assert!(accumulator.early_stop_decision().is_none());
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_VERSION_NOT_FOUND);
}
#[test]
fn metadata_quorum_accumulator_falls_back_on_insufficient_quorum() {
let mut accumulator = metadata_early_stop_accumulator();
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 1));
assert!(accumulator.early_stop_decision().is_none());
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_INSUFFICIENT_QUORUM);
}
#[test]
fn metadata_quorum_accumulator_falls_back_on_mixed_corrupt_and_valid() {
let mut accumulator = metadata_early_stop_accumulator();
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 1));
accumulator.observe_error(&DiskError::FileCorrupt);
accumulator.observe_file_info(&metadata_early_stop_candidate("object", 2));
assert!(accumulator.early_stop_decision().is_none());
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_ERROR);
}
#[test]
fn metadata_quorum_accumulator_uses_same_gate_for_head_and_get() {
let mut get_accumulator = metadata_early_stop_accumulator();
let mut head_accumulator = metadata_early_stop_accumulator();
for disk_index in [1, 2] {
let fi = metadata_early_stop_candidate("object", disk_index);
get_accumulator.observe_file_info(&fi);
head_accumulator.observe_file_info(&fi);
}
assert_eq!(get_accumulator.early_stop_decision(), head_accumulator.early_stop_decision());
assert_eq!(get_accumulator.final_miss_reason(), head_accumulator.final_miss_reason());
}
#[test]
fn metadata_early_stop_gate_defaults_to_disabled() {
temp_env::with_var(ENV_RUSTFS_GET_METADATA_EARLY_STOP_ENABLE, None::<&str>, || {
assert!(!is_get_metadata_early_stop_enabled());
});
}
#[test]
fn version_early_stop_gate_defaults_to_disabled() {
temp_env::with_var(ENV_RUSTFS_GET_METADATA_VERSION_EARLY_STOP_ENABLE, None::<&str>, || {
assert!(!is_version_early_stop_enabled());
});
}
fn version_early_stop_candidate(object: &str, disk_index: usize, version_id: Uuid) -> FileInfo {
let mut fi = metadata_early_stop_candidate(object, disk_index);
fi.version_id = Some(version_id);
fi
}
fn version_early_stop_accumulator(requested_version_id: &str) -> MetadataQuorumAccumulator {
MetadataQuorumAccumulator::new(4, 2, true).with_requested_version_id(requested_version_id)
}
#[test]
fn version_early_stop_hits_quorum_with_matching_versions() {
let vid = Uuid::new_v4();
let mut accumulator = version_early_stop_accumulator(&vid.to_string());
accumulator.observe_file_info(&version_early_stop_candidate("object", 1, vid));
assert!(accumulator.version_early_stop_decision().is_none());
accumulator.observe_file_info(&version_early_stop_candidate("object", 2, vid));
assert_eq!(
accumulator.version_early_stop_decision(),
Some(MetadataEarlyStopDecision {
reason: GET_METADATA_EARLY_STOP_REASON_VERSION_MATCH_QUORUM
})
);
assert_eq!(accumulator.matching_version_votes, 2);
}
#[test]
fn version_early_stop_does_not_fire_without_requested_version() {
let vid = Uuid::new_v4();
let mut accumulator = version_early_stop_accumulator("");
accumulator.observe_file_info(&version_early_stop_candidate("object", 1, vid));
accumulator.observe_file_info(&version_early_stop_candidate("object", 2, vid));
assert!(accumulator.version_early_stop_decision().is_none());
}
#[test]
fn version_early_stop_does_not_fire_with_mismatched_versions() {
let requested_vid = Uuid::new_v4();
let other_vid = Uuid::new_v4();
let mut accumulator = version_early_stop_accumulator(&requested_vid.to_string());
accumulator.observe_file_info(&version_early_stop_candidate("object", 1, requested_vid));
accumulator.observe_file_info(&version_early_stop_candidate("object", 2, other_vid));
assert!(accumulator.version_early_stop_decision().is_none());
assert_eq!(accumulator.matching_version_votes, 1);
}
#[test]
fn version_early_stop_does_not_fire_below_quorum() {
let vid = Uuid::new_v4();
let mut accumulator = version_early_stop_accumulator(&vid.to_string());
accumulator.observe_file_info(&version_early_stop_candidate("object", 1, vid));
assert!(accumulator.version_early_stop_decision().is_none());
assert_eq!(accumulator.matching_version_votes, 1);
}
#[test]
fn version_early_stop_tracks_matching_votes_independently_of_candidate() {
let requested_vid = Uuid::new_v4();
let mut accumulator = version_early_stop_accumulator(&requested_vid.to_string());
// Two valid responses with matching version_id but different erasure.index
// (so they conflict on the candidate path but still count for version votes)
let mut fi1 = version_early_stop_candidate("object", 1, requested_vid);
fi1.size = 100;
let mut fi2 = version_early_stop_candidate("object", 2, requested_vid);
fi2.size = 200; // different size → conflicting metadata on candidate path
accumulator.observe_file_info(&fi1);
accumulator.observe_file_info(&fi2);
// Candidate path sees conflict, but version path sees quorum
assert!(accumulator.early_stop_decision().is_none());
assert_eq!(
accumulator.version_early_stop_decision(),
Some(MetadataEarlyStopDecision {
reason: GET_METADATA_EARLY_STOP_REASON_VERSION_MATCH_QUORUM
})
);
}
fn codec_streaming_test_object_info(fi: &FileInfo) -> ObjectInfo {
ObjectInfo::from_file_info(fi, "bucket", "object", false)
}
fn inline_reader_setup_fileinfo(data: Option<&'static [u8]>) -> FileInfo {
let mut fi = FileInfo::new("object", 2, 2);
fi.volume = "bucket".to_string();
fi.name = "object".to_string();
fi.size = data.map_or(0, |data| i64::try_from(data.len()).expect("test data length should fit i64"));
fi.data = data.map(Bytes::from_static);
fi
}
async fn setup_inline_bitrot_readers(
data: Vec<Option<&'static [u8]>>,
data_shards: usize,
parity_shards: usize,
mode: BitrotReaderSetupMode,
) -> BitrotReaderSetup {
let files = data.into_iter().map(inline_reader_setup_fileinfo).collect::<Vec<_>>();
let disks = vec![None; files.len()];
create_bitrot_readers_until_quorum(
&files,
&disks,
"bucket",
"object",
1,
0,
4,
4,
HashAlgorithm::None,
false,
false,
data_shards,
parity_shards,
mode,
)
.await
}
#[tokio::test]
async fn bitrot_reader_setup_stops_at_read_quorum() {
let setup = setup_inline_bitrot_readers(
vec![Some(b"aaaa"), Some(b"bbbb"), Some(b"cccc"), Some(b"dddd")],
2,
2,
BitrotReaderSetupMode::ReadQuorum,
)
.await;
assert_eq!(setup.available_shards(), 2);
assert!(setup.data_shards_attempted(2));
assert_eq!(setup.completed_failed_shards(), 0);
}
#[tokio::test]
async fn bitrot_reader_setup_retains_unattempted_fallback_readers() {
let mut setup = setup_inline_bitrot_readers(
vec![Some(b"aaaa"), Some(b"bbbb"), Some(b"cccc"), Some(b"dddd")],
2,
2,
BitrotReaderSetupMode::ReadQuorum,
)
.await;
assert_eq!(setup.available_shards(), 2);
assert_eq!(setup.readers.iter().filter(|reader| reader.is_some()).count(), 4);
let fallback_index = setup
.attempted
.iter()
.position(|attempted| !*attempted)
.expect("read quorum should leave at least one deferred fallback");
assert!(!setup.ready[fallback_index]);
let mut fallback = setup.readers[fallback_index]
.take()
.expect("deferred fallback reader should be retained");
let mut out = [0u8; 4];
let n = fallback
.read(&mut out)
.await
.expect("deferred fallback reader should open on read");
assert_eq!(n, 4);
assert_eq!(&out[..n], [b"aaaa", b"bbbb", b"cccc", b"dddd"][fallback_index]);
}
#[tokio::test]
async fn bitrot_reader_setup_verify_mode_stops_when_data_quorum_is_available() {
let setup = setup_inline_bitrot_readers(
vec![Some(b"aaaa"), Some(b"bbbb"), Some(b"cccc"), Some(b"dddd")],
2,
2,
BitrotReaderSetupMode::VerifyReconstruction,
)
.await;
assert_eq!(setup.available_shards(), 2);
assert_eq!(setup.available_data_shards(2), 2);
assert_eq!(setup.completed_failed_shards(), 0);
}
#[tokio::test]
async fn bitrot_reader_setup_verify_mode_collects_extra_source_for_reconstruction() {
let setup = setup_inline_bitrot_readers(
vec![None, Some(b"bbbb"), Some(b"cccc"), Some(b"dddd")],
2,
2,
BitrotReaderSetupMode::VerifyReconstruction,
)
.await;
assert_eq!(setup.available_shards(), 3);
assert_eq!(setup.available_data_shards(2), 1);
assert!(setup.data_shards_attempted(2));
assert_eq!(setup.completed_failed_shards(), 1);
}
#[tokio::test]
async fn bitrot_reader_setup_verify_mode_does_not_wait_for_impossible_extra_source() {
let setup = setup_inline_bitrot_readers(
vec![None, Some(b"bbbb"), Some(b"cccc")],
2,
1,
BitrotReaderSetupMode::VerifyReconstruction,
)
.await;
assert_eq!(setup.available_shards(), 2);
assert_eq!(setup.available_data_shards(2), 1);
assert_eq!(setup.completed_failed_shards(), 1);
}
#[test]
fn codec_streaming_reader_gate_is_conservative() {
temp_env::with_vars(
[
(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, Some("benchmark")),
(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, Some("1")),
],
|| {
let fi = codec_streaming_test_fileinfo(1024, 1);
let object_info = codec_streaming_test_object_info(&fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &object_info, &fi, true).decision,
GetCodecStreamingDecision::Use
);
assert_eq!(
get_codec_streaming_reader_gate(&None, &object_info, &fi, false).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::LockOptimizationDisabled)
);
let range = Some(HTTPRangeSpec {
is_suffix_length: false,
start: 0,
end: 1,
});
assert_eq!(
get_codec_streaming_reader_gate(&range, &object_info, &fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::Range)
);
let multipart_fi = codec_streaming_test_fileinfo(1024, 2);
let multipart_object_info = codec_streaming_test_object_info(&multipart_fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &multipart_object_info, &multipart_fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::Multipart)
);
let mut encrypted_fi = fi.clone();
encrypted_fi
.metadata
.insert("x-amz-server-side-encryption".to_string(), "AES256".to_string());
let encrypted = codec_streaming_test_object_info(&encrypted_fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &encrypted, &encrypted_fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::Encrypted)
);
let mut compressed_fi = fi.clone();
insert_str(&mut compressed_fi.metadata, SUFFIX_COMPRESSION, "lz4".to_string());
let compressed = codec_streaming_test_object_info(&compressed_fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &compressed, &compressed_fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::Compressed)
);
let small_fi = codec_streaming_test_fileinfo(0, 1);
let small_object_info = codec_streaming_test_object_info(&small_fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &small_object_info, &small_fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::BelowMinSize)
);
let mut remote_fi = fi;
remote_fi.transition_status = crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string();
let remote = codec_streaming_test_object_info(&remote_fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &remote, &remote_fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::Remote)
);
},
);
}
#[test]
fn codec_streaming_engine_defaults_to_legacy_and_parses_rustfs() {
temp_env::with_var(ENV_RUSTFS_GET_CODEC_STREAMING_ENGINE, None::<&str>, || {
assert_eq!(get_codec_streaming_engine(), GetCodecStreamingEngine::Legacy);
});
temp_env::with_var(ENV_RUSTFS_GET_CODEC_STREAMING_ENGINE, Some(GET_CODEC_STREAMING_ENGINE_RUSTFS), || {
assert_eq!(get_codec_streaming_engine(), GetCodecStreamingEngine::Rustfs);
});
temp_env::with_var(ENV_RUSTFS_GET_CODEC_STREAMING_ENGINE, Some("unknown"), || {
assert_eq!(get_codec_streaming_engine(), GetCodecStreamingEngine::Legacy);
});
}
#[test]
fn codec_streaming_engine_env_is_ignored_when_streaming_is_disabled() {
temp_env::with_vars(
[
(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, Some("false")),
(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, None::<&str>),
(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, None::<&str>),
(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, None::<&str>),
(ENV_RUSTFS_GET_CODEC_STREAMING_ENGINE, Some(GET_CODEC_STREAMING_ENGINE_RUSTFS)),
(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, Some("1")),
],
|| {
let fi = codec_streaming_test_fileinfo(1024, 1);
let object_info = codec_streaming_test_object_info(&fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &object_info, &fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::Disabled)
);
},
);
}
#[test]
fn codec_streaming_decode_engine_builder_selects_rustfs() {
temp_env::with_var(ENV_RUSTFS_GET_CODEC_STREAMING_ENGINE, Some(GET_CODEC_STREAMING_ENGINE_RUSTFS), || {
let erasure = crate::erasure::coding::Erasure::new(4, 2, 32);
let engine = build_get_codec_streaming_decode_engine(erasure).expect("engine should be built");
assert!(matches!(engine, CodecStreamingDecodeEngine::Rustfs(_)));
});
}
#[test]
fn codec_streaming_metrics_path_matches_selected_engine() {
temp_env::with_var(ENV_RUSTFS_GET_CODEC_STREAMING_ENGINE, None::<&str>, || {
assert_eq!(
get_codec_streaming_metrics_path(),
crate::diagnostics::get::GET_OBJECT_PATH_CODEC_STREAMING_LEGACY_ENGINE
);
});
temp_env::with_var(ENV_RUSTFS_GET_CODEC_STREAMING_ENGINE, Some(GET_CODEC_STREAMING_ENGINE_RUSTFS), || {
assert_eq!(
get_codec_streaming_metrics_path(),
crate::diagnostics::get::GET_OBJECT_PATH_CODEC_STREAMING_RUSTFS_ENGINE
);
});
}
#[test]
fn codec_streaming_fallback_metric_labels_are_stable() {
assert_eq!(GetCodecStreamingFallbackReason::Disabled.as_str(), "disabled");
assert_eq!(GetCodecStreamingFallbackReason::RolloutNotOptedIn.as_str(), "rollout_not_opted_in");
assert_eq!(
GetCodecStreamingFallbackReason::BodyCompatibilityUnconfirmed.as_str(),
"body_compatibility_unconfirmed"
);
assert_eq!(
GetCodecStreamingFallbackReason::HeaderCompatibilityUnconfirmed.as_str(),
"header_compatibility_unconfirmed"
);
assert_eq!(
GetCodecStreamingFallbackReason::LockOptimizationDisabled.as_str(),
"lock_optimization_disabled"
);
assert_eq!(GetCodecStreamingFallbackReason::Range.as_str(), "range");
assert_eq!(GetCodecStreamingFallbackReason::BelowMinSize.as_str(), "below_min_size");
assert_eq!(GetCodecStreamingFallbackReason::Encrypted.as_str(), "encrypted");
assert_eq!(GetCodecStreamingFallbackReason::Compressed.as_str(), "compressed");
assert_eq!(GetCodecStreamingFallbackReason::Remote.as_str(), "remote");
assert_eq!(GetCodecStreamingFallbackReason::Multipart.as_str(), "multipart");
assert_eq!(GetCodecStreamingFallbackReason::InvalidMinSize.as_str(), "invalid_min_size");
assert_eq!(GetCodecStreamingFallbackReason::ReadQuorumNotSafe.as_str(), "read_quorum_not_safe");
assert_eq!(GetCodecStreamingObjectClass::PlainSinglePart.as_str(), "plain_single_part");
assert_eq!(GetCodecStreamingObjectClass::Range.as_str(), "range");
assert_eq!(GetCodecStreamingObjectClass::Encrypted.as_str(), "encrypted");
assert_eq!(GetCodecStreamingObjectClass::Compressed.as_str(), "compressed");
assert_eq!(GetCodecStreamingObjectClass::Remote.as_str(), "remote");
assert_eq!(GetCodecStreamingObjectClass::Multipart.as_str(), "multipart");
}
#[test]
fn codec_streaming_reader_gate_defaults_to_disabled() {
temp_env::with_vars(
[
(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, None::<&str>),
(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, None::<&str>),
(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, None::<&str>),
(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, None::<&str>),
(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, Some("1")),
],
|| {
let fi = codec_streaming_test_fileinfo(1024, 1);
let object_info = codec_streaming_test_object_info(&fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &object_info, &fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::Disabled)
);
},
);
}
#[test]
fn codec_streaming_reader_gate_requires_explicit_rollout_and_compat_confirmation() {
temp_env::with_vars(
[
(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, Some("1")),
(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, Some("off")),
(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, None::<&str>),
(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, None::<&str>),
],
|| {
let fi = codec_streaming_test_fileinfo(1024, 1);
let object_info = codec_streaming_test_object_info(&fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &object_info, &fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::RolloutNotOptedIn)
);
},
);
temp_env::with_vars(
[
(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, Some("1")),
(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, Some("benchmark")),
(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, Some("false")),
(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, Some("true")),
],
|| {
let fi = codec_streaming_test_fileinfo(1024, 1);
let object_info = codec_streaming_test_object_info(&fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &object_info, &fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::BodyCompatibilityUnconfirmed)
);
},
);
temp_env::with_vars(
[
(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, Some("1")),
(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, Some("benchmark")),
(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, Some("false")),
],
|| {
let fi = codec_streaming_test_fileinfo(1024, 1);
let object_info = codec_streaming_test_object_info(&fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &object_info, &fi, true).decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::HeaderCompatibilityUnconfirmed)
);
},
);
}
#[test]
fn codec_streaming_reader_gate_records_object_classes() {
temp_env::with_vars(
[
(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, Some("benchmark")),
(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, Some("1")),
],
|| {
let fi = codec_streaming_test_fileinfo(1024, 1);
let object_info = codec_streaming_test_object_info(&fi);
assert_eq!(
get_codec_streaming_reader_gate(&None, &object_info, &fi, true).object_class,
GetCodecStreamingObjectClass::PlainSinglePart
);
let range = Some(HTTPRangeSpec {
is_suffix_length: false,
start: 0,
end: 1,
});
assert_eq!(
get_codec_streaming_reader_gate(&range, &object_info, &fi, true).object_class,
GetCodecStreamingObjectClass::Range
);
},
);
}
#[tokio::test]
async fn codec_streaming_reader_build_falls_back_when_read_quorum_is_not_safe() {
let setup = setup_inline_bitrot_readers(
vec![None, Some(b"bbbb"), Some(b"cccc"), Some(b"dddd")],
2,
2,
BitrotReaderSetupMode::VerifyReconstruction,
)
.await;
assert_eq!(setup.completed_failed_shards(), 1);
assert_eq!(
codec_streaming_reader_setup_fallback_reason(setup.completed_failed_shards()),
Some(GetCodecStreamingFallbackReason::ReadQuorumNotSafe)
);
}
#[tokio::test]
async fn collect_read_multiple_results_fails_early_when_quorum_is_impossible() {
let started = std::time::Instant::now();
let resp = ReadMultipleResp {
bucket: "bucket".to_string(),
prefix: "prefix".to_string(),
file: "file".to_string(),
exists: true,
error: String::new(),
data: vec![1],
mod_time: None,
};
let tasks: Vec<_> = vec![
(10_u64, Err(DiskError::DiskNotFound)),
(15, Err(DiskError::DiskNotFound)),
(250, Ok::<Vec<ReadMultipleResp>, DiskError>(vec![resp])),
]
.into_iter()
.map(|(delay_ms, outcome)| async move {
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
outcome
})
.collect();
let result = collect_read_multiple_results(tasks, 2).await;
assert!(result.is_err(), "quorum should become impossible before slow tail completes");
assert!(started.elapsed() < std::time::Duration::from_millis(120));
}
#[tokio::test]
async fn collect_read_multiple_results_returns_collected_responses_on_quorum() {
let resp = ReadMultipleResp {
bucket: "bucket".to_string(),
prefix: "prefix".to_string(),
file: "file".to_string(),
exists: true,
error: String::new(),
data: vec![1, 2, 3],
mod_time: None,
};
let tasks: Vec<_> = vec![
(10_u64, Ok::<Vec<ReadMultipleResp>, DiskError>(vec![resp.clone()])),
(15, Ok::<Vec<ReadMultipleResp>, DiskError>(vec![resp.clone()])),
(250, Err(DiskError::DiskNotFound)),
]
.into_iter()
.map(|(delay_ms, outcome)| async move {
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
outcome
})
.collect();
let (responses, errors) = collect_read_multiple_results(tasks, 2).await.expect("quorum should succeed");
assert_eq!(responses.iter().filter(|item| item.is_some()).count(), 2);
assert_eq!(errors.iter().filter(|item| item.is_none()).count(), 2);
}
#[tokio::test]
async fn collect_read_multiple_results_tolerates_single_panicked_task_when_quorum_is_met() {
let resp = ReadMultipleResp {
bucket: "bucket".to_string(),
prefix: "prefix".to_string(),
file: "file".to_string(),
exists: true,
error: String::new(),
data: vec![1, 2, 3],
mod_time: None,
};
let tasks: Vec<_> = vec![(5_u64, true), (10, false), (12, false)]
.into_iter()
.map(|(delay_ms, should_panic)| {
let resp = resp.clone();
async move {
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
if should_panic {
panic!("simulated task panic");
}
Ok::<Vec<ReadMultipleResp>, DiskError>(vec![resp])
}
})
.collect();
let (responses, errors) = collect_read_multiple_results(tasks, 2)
.await
.expect("quorum should still succeed");
assert_eq!(responses.iter().filter(|item| item.is_some()).count(), 2);
assert_eq!(errors.iter().filter(|item| item.is_none()).count(), 2);
}
#[tokio::test]
async fn collect_read_parts_results_fails_early_when_quorum_is_impossible() {
let started = std::time::Instant::now();
let part = ObjectPartInfo {
number: 1,
etag: "etag".to_string(),
..Default::default()
};
let tasks: Vec<_> = vec![
(10_u64, Err(DiskError::DiskNotFound)),
(15, Err(DiskError::DiskNotFound)),
(250, Ok::<Vec<ObjectPartInfo>, DiskError>(vec![part])),
]
.into_iter()
.map(|(delay_ms, outcome)| async move {
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
outcome
})
.collect();
let result = collect_read_parts_results(tasks, 2).await;
assert!(result.is_err(), "quorum should become impossible before slow tail completes");
assert!(started.elapsed() < std::time::Duration::from_millis(120));
}
#[tokio::test]
async fn collect_read_parts_results_returns_collected_responses_on_quorum() {
let part = ObjectPartInfo {
number: 1,
etag: "etag".to_string(),
..Default::default()
};
let tasks: Vec<_> = vec![
(10_u64, Ok::<Vec<ObjectPartInfo>, DiskError>(vec![part.clone()])),
(15, Ok::<Vec<ObjectPartInfo>, DiskError>(vec![part.clone()])),
(250, Err(DiskError::DiskNotFound)),
]
.into_iter()
.map(|(delay_ms, outcome)| async move {
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
outcome
})
.collect();
let (responses, errors) = collect_read_parts_results(tasks, 2).await.expect("quorum should succeed");
assert_eq!(responses.iter().filter(|item| item.is_some()).count(), 2);
assert_eq!(errors.iter().filter(|item| item.is_none()).count(), 2);
}
#[tokio::test]
async fn collect_read_parts_results_tolerates_single_panicked_task_when_quorum_is_met() {
let part = ObjectPartInfo {
number: 1,
etag: "etag".to_string(),
..Default::default()
};
let tasks: Vec<_> = vec![(5_u64, true), (10, false), (12, false)]
.into_iter()
.map(|(delay_ms, should_panic)| {
let part = part.clone();
async move {
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
if should_panic {
panic!("simulated task panic");
}
Ok::<Vec<ObjectPartInfo>, DiskError>(vec![part])
}
})
.collect();
let (responses, errors) = collect_read_parts_results(tasks, 2)
.await
.expect("quorum should still succeed");
assert_eq!(responses.iter().filter(|item| item.is_some()).count(), 2);
assert_eq!(errors.iter().filter(|item| item.is_none()).count(), 2);
}
}