mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 20:59:05 +00:00
Merge remote-tracking branch 'origin/main' into overtrue/fix-1911-ilm-metadata-decommission
# Conflicts: # crates/ecstore/src/core/pools.rs
This commit is contained in:
@@ -53,11 +53,12 @@ use crate::diagnostics::get::{
|
||||
GetObjectFailureReason, classify_disk_error, get_stage_timer_if_enabled, record_get_object_pipeline_failure,
|
||||
record_get_object_pipeline_failure_for_path, record_get_stage_duration_if_enabled,
|
||||
};
|
||||
use crate::disk::disk_store::DiskStoreRenameDataExt;
|
||||
use crate::disk::disk_store::{DiskStoreRenameDataExt, get_drive_metadata_timeout};
|
||||
use crate::disk::local::DELETE_DATA_DIR_MARKER_PREFIX;
|
||||
use crate::disk::{
|
||||
DataDirDeleteStatus, OldCurrentSize, PART_TRANSACTION_NEW_META, PART_TRANSACTION_OLD_META, PART_TRANSACTION_ROLLBACK,
|
||||
PartTransactionAction, STORAGE_FORMAT_FILE_BACKUP, part_transaction_path,
|
||||
BATCH_READ_VERSION_MAX_ITEMS, BatchReadVersionItem, BatchReadVersionReq, BatchReadVersionResp, DataDirDeleteStatus, Disk,
|
||||
OldCurrentSize, PART_TRANSACTION_NEW_META, PART_TRANSACTION_OLD_META, PART_TRANSACTION_ROLLBACK, PartTransactionAction,
|
||||
STORAGE_FORMAT_FILE_BACKUP, part_transaction_path,
|
||||
};
|
||||
use crate::erasure::coding::BitrotReader;
|
||||
use crate::io_support::bitrot::ShardReader;
|
||||
@@ -75,7 +76,7 @@ use std::{
|
||||
future::Future,
|
||||
pin::Pin,
|
||||
sync::{
|
||||
OnceLock,
|
||||
Arc, OnceLock,
|
||||
atomic::{AtomicUsize, Ordering},
|
||||
},
|
||||
task::{Context, Poll},
|
||||
@@ -94,6 +95,242 @@ fn metadata_distribution_key(bucket: &str, object: &str) -> String {
|
||||
[bucket, object].join("/")
|
||||
}
|
||||
|
||||
fn read_version_coalescing_enabled() -> bool {
|
||||
let enabled = || {
|
||||
rustfs_utils::get_env_opt_str(ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE)
|
||||
.is_some_and(|value| value.eq_ignore_ascii_case("auto") || value.eq_ignore_ascii_case("on"))
|
||||
};
|
||||
|
||||
#[cfg(test)]
|
||||
{
|
||||
enabled()
|
||||
}
|
||||
|
||||
#[cfg(not(test))]
|
||||
{
|
||||
static ENABLED: OnceLock<bool> = OnceLock::new();
|
||||
*ENABLED.get_or_init(enabled)
|
||||
}
|
||||
}
|
||||
|
||||
fn read_version_coalescing_delay() -> Duration {
|
||||
#[cfg(test)]
|
||||
{
|
||||
let micros = rustfs_utils::get_env_u64(
|
||||
ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS,
|
||||
DEFAULT_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS,
|
||||
);
|
||||
Duration::from_micros(micros)
|
||||
}
|
||||
|
||||
#[cfg(not(test))]
|
||||
{
|
||||
static DELAY: OnceLock<Duration> = OnceLock::new();
|
||||
*DELAY.get_or_init(|| {
|
||||
Duration::from_micros(rustfs_utils::get_env_u64(
|
||||
ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS,
|
||||
DEFAULT_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS,
|
||||
))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
struct CoalescedReadVersionRequest {
|
||||
item: BatchReadVersionItem,
|
||||
tx: oneshot::Sender<disk::error::Result<FileInfo>>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
|
||||
struct ReadVersionCoalescerKey {
|
||||
disk: usize,
|
||||
incl_free_versions: bool,
|
||||
read_data: bool,
|
||||
healing: bool,
|
||||
}
|
||||
|
||||
impl ReadVersionCoalescerKey {
|
||||
fn new(disk: &DiskStore, opts: &ReadOptions) -> Self {
|
||||
Self {
|
||||
disk: Arc::as_ptr(disk) as usize,
|
||||
incl_free_versions: opts.incl_free_versions,
|
||||
read_data: opts.read_data,
|
||||
healing: opts.healing,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct ReadVersionCoalescer {
|
||||
lanes: HashMap<ReadVersionCoalescerKey, Vec<CoalescedReadVersionRequest>>,
|
||||
}
|
||||
|
||||
fn read_version_coalescer() -> &'static Mutex<ReadVersionCoalescer> {
|
||||
static COALESCER: OnceLock<Mutex<ReadVersionCoalescer>> = OnceLock::new();
|
||||
COALESCER.get_or_init(|| Mutex::new(ReadVersionCoalescer::default()))
|
||||
}
|
||||
|
||||
fn record_read_version_coalescer_event(event: &'static str, item_count: usize) {
|
||||
counter!(
|
||||
METRIC_GET_METADATA_READ_VERSION_COALESCER_TOTAL,
|
||||
"event" => event,
|
||||
"item_count" => item_count.to_string()
|
||||
)
|
||||
.increment(1);
|
||||
}
|
||||
|
||||
async fn read_version_via_coalescer(
|
||||
disk: DiskStore,
|
||||
org_bucket: &str,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
version_id: &str,
|
||||
opts: &ReadOptions,
|
||||
allow_coalescing: bool,
|
||||
) -> disk::error::Result<FileInfo> {
|
||||
if !allow_coalescing || !read_version_coalescing_enabled() {
|
||||
return disk.read_version(org_bucket, bucket, object, version_id, opts).await;
|
||||
}
|
||||
if !matches!(disk.as_ref(), Disk::Remote(_)) {
|
||||
record_read_version_coalescer_event("bypass_non_remote", 1);
|
||||
return disk.read_version(org_bucket, bucket, object, version_id, opts).await;
|
||||
}
|
||||
|
||||
let (tx, rx) = oneshot::channel();
|
||||
let item = BatchReadVersionItem {
|
||||
org_volume: org_bucket.to_string(),
|
||||
volume: bucket.to_string(),
|
||||
path: object.to_string(),
|
||||
version_id: version_id.to_string(),
|
||||
};
|
||||
let lane_key = ReadVersionCoalescerKey::new(&disk, opts);
|
||||
let pending = {
|
||||
let mut coalescer = read_version_coalescer().lock().await;
|
||||
let lane = coalescer.lanes.entry(lane_key).or_default();
|
||||
let schedule_delayed_flush = lane.is_empty();
|
||||
lane.push(CoalescedReadVersionRequest { item, tx });
|
||||
if lane.len() >= BATCH_READ_VERSION_MAX_ITEMS {
|
||||
coalescer.lanes.remove(&lane_key)
|
||||
} else if schedule_delayed_flush {
|
||||
let disk = disk.clone();
|
||||
let task_opts = *opts;
|
||||
tokio::spawn(async move {
|
||||
tokio::time::sleep(read_version_coalescing_delay()).await;
|
||||
flush_read_version_coalescer_lane(lane_key, disk, task_opts).await;
|
||||
});
|
||||
None
|
||||
} else {
|
||||
None
|
||||
}
|
||||
};
|
||||
|
||||
if let Some(pending) = pending {
|
||||
flush_read_version_coalescer_pending(lane_key, disk, *opts, pending).await;
|
||||
}
|
||||
|
||||
rx.await
|
||||
.unwrap_or_else(|_| Err(DiskError::other("coalesced read_version response channel closed")))
|
||||
}
|
||||
|
||||
async fn flush_read_version_coalescer_lane(lane_key: ReadVersionCoalescerKey, disk: DiskStore, opts: ReadOptions) {
|
||||
let pending = {
|
||||
let mut coalescer = read_version_coalescer().lock().await;
|
||||
coalescer.lanes.remove(&lane_key).unwrap_or_default()
|
||||
};
|
||||
flush_read_version_coalescer_pending(lane_key, disk, opts, pending).await;
|
||||
}
|
||||
|
||||
async fn flush_read_version_coalescer_pending(
|
||||
lane_key: ReadVersionCoalescerKey,
|
||||
disk: DiskStore,
|
||||
opts: ReadOptions,
|
||||
pending: Vec<CoalescedReadVersionRequest>,
|
||||
) {
|
||||
if pending.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
{
|
||||
let mut observed_paths = HashSet::new();
|
||||
for request in &pending {
|
||||
if observed_paths.insert(request.item.path.as_str()) {
|
||||
disk_call_counters::record(&request.item.path, disk_call_counters::KIND_BATCH_READ_VERSION, lane_key.disk);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let mut senders = Vec::with_capacity(pending.len());
|
||||
let mut items = Vec::with_capacity(pending.len());
|
||||
for request in pending {
|
||||
senders.push(request.tx);
|
||||
items.push(request.item);
|
||||
}
|
||||
|
||||
let expected_items = items.clone();
|
||||
record_read_version_coalescer_event("attempted_batch", items.len());
|
||||
let result =
|
||||
match tokio::time::timeout(get_drive_metadata_timeout(), disk.batch_read_version(BatchReadVersionReq { items, opts }))
|
||||
.await
|
||||
{
|
||||
Ok(result) => result,
|
||||
Err(_) => Err(DiskError::Timeout),
|
||||
};
|
||||
match result {
|
||||
Ok(responses) => {
|
||||
let results = map_batch_read_version_responses(&expected_items, responses);
|
||||
for (tx, result) in senders.into_iter().zip(results) {
|
||||
let _ = tx.send(result);
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
let message = err.to_string();
|
||||
for tx in senders {
|
||||
let _ = tx.send(Err(DiskError::other(message.clone())));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn map_batch_read_version_responses(
|
||||
expected_items: &[BatchReadVersionItem],
|
||||
responses: Vec<BatchReadVersionResp>,
|
||||
) -> Vec<crate::disk::error::Result<FileInfo>> {
|
||||
let mut results = (0..expected_items.len())
|
||||
.map(|_| Err(DiskError::other("coalesced read_version response missing")))
|
||||
.collect::<Vec<_>>();
|
||||
let mut seen = vec![false; expected_items.len()];
|
||||
for response in responses {
|
||||
let Some(expected) = expected_items.get(response.index) else {
|
||||
continue;
|
||||
};
|
||||
let Some(slot) = results.get_mut(response.index) else {
|
||||
continue;
|
||||
};
|
||||
if seen[response.index] {
|
||||
*slot = Err(DiskError::other("coalesced read_version response duplicate index"));
|
||||
continue;
|
||||
}
|
||||
seen[response.index] = true;
|
||||
if response.path != expected.path || response.version_id != expected.version_id {
|
||||
*slot = Err(DiskError::other("coalesced read_version response identity mismatch"));
|
||||
} else {
|
||||
*slot = if response.success {
|
||||
Ok(response.file_info)
|
||||
} else {
|
||||
Err(batch_read_version_response_error(response.error_code, response.error))
|
||||
};
|
||||
}
|
||||
}
|
||||
results
|
||||
}
|
||||
|
||||
fn batch_read_version_response_error(error_code: u32, error: String) -> DiskError {
|
||||
match DiskError::from_u32(error_code) {
|
||||
Some(DiskError::Io(_)) | None => DiskError::other(error),
|
||||
Some(error) => error,
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn bounded_metadata_fanout_order(
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
@@ -133,11 +370,15 @@ pub(in crate::set_disk) fn bounded_metadata_fanout_order(
|
||||
order
|
||||
}
|
||||
use tokio::io::{AsyncRead, ReadBuf};
|
||||
use tokio::sync::RwLock;
|
||||
use tokio::sync::{Mutex, RwLock, oneshot};
|
||||
use tokio::task::JoinSet;
|
||||
|
||||
pub(in crate::set_disk) const EVENT_SET_DISK_READ: &str = "set_disk_read";
|
||||
pub(in crate::set_disk) const ENV_RUSTFS_GET_DATA_BLOCKS_FIRST_READER_SETUP: &str = "RUSTFS_GET_DATA_BLOCKS_FIRST_READER_SETUP";
|
||||
const ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE: &str = "RUSTFS_GET_METADATA_READ_VERSION_COALESCE";
|
||||
const ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS: &str = "RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS";
|
||||
const DEFAULT_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS: u64 = 200;
|
||||
const METRIC_GET_METADATA_READ_VERSION_COALESCER_TOTAL: &str = "rustfs_get_metadata_read_version_coalescer_total";
|
||||
pub(in crate::set_disk) const ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE: &str = "RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE";
|
||||
/// Default reader-setup strategy for the GET read path (rustfs/backlog#1215,
|
||||
/// #1159, #923).
|
||||
@@ -2356,6 +2597,7 @@ impl SetDisks {
|
||||
false,
|
||||
true,
|
||||
0,
|
||||
false,
|
||||
)
|
||||
.await?;
|
||||
Ok((ress, errors))
|
||||
@@ -2386,6 +2628,36 @@ impl SetDisks {
|
||||
true,
|
||||
caller_allows_early_stop,
|
||||
default_parity_count,
|
||||
false,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub(in crate::set_disk) async fn read_all_fileinfo_observed_for_get_object(
|
||||
disks: &[Option<DiskStore>],
|
||||
org_bucket: &str,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
version_id: &str,
|
||||
read_data: bool,
|
||||
incl_free_versions: bool,
|
||||
caller_allows_early_stop: 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,
|
||||
false,
|
||||
incl_free_versions,
|
||||
true,
|
||||
caller_allows_early_stop,
|
||||
default_parity_count,
|
||||
true,
|
||||
)
|
||||
.await
|
||||
}
|
||||
@@ -2408,6 +2680,7 @@ impl SetDisks {
|
||||
// subset would fail write quorum (backlog#872 regression).
|
||||
caller_allows_early_stop: bool,
|
||||
default_parity_count: usize,
|
||||
allow_coalescing: bool,
|
||||
) -> disk::error::Result<(Vec<FileInfo>, Vec<Option<DiskError>>, MetadataFanoutDiagnostics)> {
|
||||
let early_stop_enabled =
|
||||
caller_allows_early_stop && observe && (is_get_metadata_early_stop_enabled() || is_version_early_stop_enabled());
|
||||
@@ -2424,6 +2697,7 @@ impl SetDisks {
|
||||
healing,
|
||||
incl_free_versions,
|
||||
default_parity_count,
|
||||
allow_coalescing,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
@@ -2446,6 +2720,7 @@ impl SetDisks {
|
||||
healing,
|
||||
incl_free_versions,
|
||||
observe,
|
||||
allow_coalescing,
|
||||
)
|
||||
.await
|
||||
}
|
||||
@@ -2461,6 +2736,7 @@ impl SetDisks {
|
||||
healing: bool,
|
||||
incl_free_versions: bool,
|
||||
observe: bool,
|
||||
allow_coalescing: 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());
|
||||
@@ -2492,7 +2768,7 @@ impl SetDisks {
|
||||
if let Some(delay) = slowtail_fault.as_ref().and_then(|fault| fault.delay_for_disk(disk_index)) {
|
||||
tokio::time::sleep(delay).await;
|
||||
}
|
||||
disk.read_version(&org_bucket, &bucket, &object, &version_id, &task_opts)
|
||||
read_version_via_coalescer(disk, &org_bucket, &bucket, &object, &version_id, &task_opts, allow_coalescing)
|
||||
.await
|
||||
} else {
|
||||
Err(DiskError::DiskNotFound)
|
||||
@@ -2559,6 +2835,7 @@ impl SetDisks {
|
||||
healing: bool,
|
||||
incl_free_versions: bool,
|
||||
default_parity_count: usize,
|
||||
allow_coalescing: bool,
|
||||
) -> disk::error::Result<(Vec<FileInfo>, Vec<Option<DiskError>>, MetadataFanoutDiagnostics)> {
|
||||
let fanout_start = Instant::now();
|
||||
let mut ress = vec![FileInfo::default(); disks.len()];
|
||||
@@ -2607,7 +2884,7 @@ impl SetDisks {
|
||||
if let Some(delay) = slowtail_fault.as_ref().and_then(|fault| fault.delay_for_disk(index)) {
|
||||
tokio::time::sleep(delay).await;
|
||||
}
|
||||
disk.read_version(&org_bucket, &bucket, &object, &version_id, &task_opts)
|
||||
read_version_via_coalescer(disk, &org_bucket, &bucket, &object, &version_id, &task_opts, allow_coalescing)
|
||||
.await
|
||||
} else {
|
||||
Err(DiskError::DiskNotFound)
|
||||
@@ -5737,6 +6014,7 @@ pub(crate) mod disk_call_counters {
|
||||
|
||||
/// Kind label for the per-disk `read_version` metadata RPC.
|
||||
pub const KIND_READ_VERSION: &str = "read_version";
|
||||
pub const KIND_BATCH_READ_VERSION: &str = "batch_read_version";
|
||||
|
||||
/// Registry key: (object, kind, disk_index).
|
||||
type CountKey = (String, String, usize);
|
||||
@@ -6460,6 +6738,286 @@ mod tests {
|
||||
drop(dirs);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
async fn metadata_read_version_coalescer_bypasses_local_disks() {
|
||||
const DISKS: usize = 4;
|
||||
let bucket = "coalesced-read-version-local-bypass-bucket";
|
||||
let object_a = "coalesced-local-object-a";
|
||||
let object_b = "coalesced-local-object-b";
|
||||
let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await;
|
||||
install_metadata_fanout_fileinfo(&disks, bucket, object_a, None).await;
|
||||
install_metadata_fanout_fileinfo(&disks, bucket, object_b, None).await;
|
||||
|
||||
temp_env::async_with_vars(
|
||||
[
|
||||
(ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE, Some("auto")),
|
||||
(ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE_DELAY_MICROS, Some("5000")),
|
||||
],
|
||||
async {
|
||||
let calls = disk_call_counters::observe(object_a);
|
||||
let disks_a = disks.clone();
|
||||
let disks_b = disks.clone();
|
||||
let read_a = tokio::spawn(async move {
|
||||
SetDisks::read_all_fileinfo_observed_for_get_object(
|
||||
&disks_a, "", bucket, object_a, "", false, false, false, 2,
|
||||
)
|
||||
.await
|
||||
.map(|(file_infos, errors, _)| (file_infos, errors))
|
||||
});
|
||||
tokio::task::yield_now().await;
|
||||
let read_b = tokio::spawn(async move {
|
||||
SetDisks::read_all_fileinfo_observed_for_get_object(
|
||||
&disks_b, "", bucket, object_b, "", false, false, false, 2,
|
||||
)
|
||||
.await
|
||||
.map(|(file_infos, errors, _)| (file_infos, errors))
|
||||
});
|
||||
|
||||
let (metadata_a, errs_a) = read_a
|
||||
.await
|
||||
.expect("first read task should not panic")
|
||||
.expect("first coalesced read should resolve");
|
||||
let (metadata_b, errs_b) = read_b
|
||||
.await
|
||||
.expect("second read task should not panic")
|
||||
.expect("second coalesced read should resolve");
|
||||
|
||||
assert_eq!(metadata_a.iter().filter(|fi| fi.name == object_a).count(), DISKS);
|
||||
assert_eq!(metadata_b.iter().filter(|fi| fi.name == object_b).count(), DISKS);
|
||||
assert!(errs_a.iter().all(Option::is_none));
|
||||
assert!(errs_b.iter().all(Option::is_none));
|
||||
assert_eq!(
|
||||
calls.total(disk_call_counters::KIND_READ_VERSION),
|
||||
DISKS as u64,
|
||||
"local disks still execute the ordinary per-disk read_version path"
|
||||
);
|
||||
assert_eq!(
|
||||
calls.total(disk_call_counters::KIND_BATCH_READ_VERSION),
|
||||
0,
|
||||
"GET coalescing targets internode RPC count only and must not batch local disk reads"
|
||||
);
|
||||
},
|
||||
)
|
||||
.await;
|
||||
|
||||
drop(dirs);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn metadata_read_version_coalescer_requires_get_object_intent() {
|
||||
const DISKS: usize = 4;
|
||||
let bucket = "coalesced-read-version-default-bypass-bucket";
|
||||
let object = "default-bypass-object";
|
||||
let (dirs, disks) = call_counter_local_disks(bucket, DISKS).await;
|
||||
install_metadata_fanout_fileinfo(&disks, bucket, object, None).await;
|
||||
|
||||
temp_env::async_with_vars([(ENV_RUSTFS_GET_METADATA_READ_VERSION_COALESCE, Some("auto"))], async {
|
||||
let calls = disk_call_counters::observe(object);
|
||||
let (metadata, errs) = SetDisks::read_all_fileinfo(&disks, "", bucket, object, "", false, false, false)
|
||||
.await
|
||||
.expect("default metadata read should resolve");
|
||||
|
||||
assert_eq!(metadata.iter().filter(|fi| fi.name == object).count(), DISKS);
|
||||
assert!(errs.iter().all(Option::is_none));
|
||||
assert_eq!(calls.total(disk_call_counters::KIND_READ_VERSION), DISKS as u64);
|
||||
assert_eq!(
|
||||
calls.total(disk_call_counters::KIND_BATCH_READ_VERSION),
|
||||
0,
|
||||
"non-GET metadata paths must bypass coalescer even when the env gate is enabled"
|
||||
);
|
||||
})
|
||||
.await;
|
||||
|
||||
drop(dirs);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn batch_read_version_response_mapping_preserves_index_and_errors() {
|
||||
let expected_items = vec![
|
||||
BatchReadVersionItem {
|
||||
org_volume: String::new(),
|
||||
volume: "bucket".to_string(),
|
||||
path: "object-a".to_string(),
|
||||
version_id: "v-a".to_string(),
|
||||
},
|
||||
BatchReadVersionItem {
|
||||
org_volume: String::new(),
|
||||
volume: "bucket".to_string(),
|
||||
path: "object-b".to_string(),
|
||||
version_id: "v-b".to_string(),
|
||||
},
|
||||
BatchReadVersionItem {
|
||||
org_volume: String::new(),
|
||||
volume: "bucket".to_string(),
|
||||
path: "object-c".to_string(),
|
||||
version_id: "v-c".to_string(),
|
||||
},
|
||||
];
|
||||
let ok_file_info = FileInfo {
|
||||
name: "object-a".to_string(),
|
||||
..Default::default()
|
||||
};
|
||||
let responses = vec![
|
||||
BatchReadVersionResp {
|
||||
index: 2,
|
||||
path: "object-c".to_string(),
|
||||
version_id: "v-c".to_string(),
|
||||
success: false,
|
||||
file_info: FileInfo::default(),
|
||||
error: "disk read failed".to_string(),
|
||||
error_code: 0,
|
||||
},
|
||||
BatchReadVersionResp {
|
||||
index: 0,
|
||||
path: "object-a".to_string(),
|
||||
version_id: "v-a".to_string(),
|
||||
success: true,
|
||||
file_info: ok_file_info,
|
||||
error: String::new(),
|
||||
error_code: 0,
|
||||
},
|
||||
];
|
||||
|
||||
let mut results = map_batch_read_version_responses(&expected_items, responses).into_iter();
|
||||
let first = results
|
||||
.next()
|
||||
.expect("slot 0 should exist")
|
||||
.expect("slot 0 should map the success response by index");
|
||||
assert_eq!(first.name, "object-a");
|
||||
|
||||
let missing = results
|
||||
.next()
|
||||
.expect("slot 1 should exist")
|
||||
.expect_err("slot 1 should stay missing");
|
||||
assert!(
|
||||
missing.to_string().contains("response missing"),
|
||||
"unexpected missing response error: {missing}"
|
||||
);
|
||||
|
||||
let failed = results
|
||||
.next()
|
||||
.expect("slot 2 should exist")
|
||||
.expect_err("slot 2 should map the response error");
|
||||
assert!(failed.to_string().contains("disk read failed"), "unexpected per-item error: {failed}");
|
||||
assert!(results.next().is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn batch_read_version_response_mapping_preserves_typed_not_found_errors() {
|
||||
let expected_items = vec![
|
||||
BatchReadVersionItem {
|
||||
org_volume: String::new(),
|
||||
volume: "bucket".to_string(),
|
||||
path: "object-a".to_string(),
|
||||
version_id: "v-a".to_string(),
|
||||
},
|
||||
BatchReadVersionItem {
|
||||
org_volume: String::new(),
|
||||
volume: "bucket".to_string(),
|
||||
path: "object-b".to_string(),
|
||||
version_id: "v-b".to_string(),
|
||||
},
|
||||
];
|
||||
let results = map_batch_read_version_responses(
|
||||
&expected_items,
|
||||
vec![
|
||||
BatchReadVersionResp {
|
||||
index: 0,
|
||||
path: "object-a".to_string(),
|
||||
version_id: "v-a".to_string(),
|
||||
success: false,
|
||||
file_info: FileInfo::default(),
|
||||
error: DiskError::FileNotFound.to_string(),
|
||||
error_code: DiskError::FileNotFound.to_u32(),
|
||||
},
|
||||
BatchReadVersionResp {
|
||||
index: 1,
|
||||
path: "object-b".to_string(),
|
||||
version_id: "v-b".to_string(),
|
||||
success: false,
|
||||
file_info: FileInfo::default(),
|
||||
error: DiskError::FileVersionNotFound.to_string(),
|
||||
error_code: DiskError::FileVersionNotFound.to_u32(),
|
||||
},
|
||||
],
|
||||
);
|
||||
|
||||
assert!(matches!(results.first().expect("slot 0 should exist"), Err(DiskError::FileNotFound)));
|
||||
assert!(matches!(
|
||||
results.get(1).expect("slot 1 should exist"),
|
||||
Err(DiskError::FileVersionNotFound)
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn batch_read_version_response_mapping_rejects_identity_mismatch_and_duplicate_index() {
|
||||
let expected_items = vec![BatchReadVersionItem {
|
||||
org_volume: String::new(),
|
||||
volume: "bucket".to_string(),
|
||||
path: "object-a".to_string(),
|
||||
version_id: "v-a".to_string(),
|
||||
}];
|
||||
let mismatched = map_batch_read_version_responses(
|
||||
&expected_items,
|
||||
vec![BatchReadVersionResp {
|
||||
index: 0,
|
||||
path: "object-b".to_string(),
|
||||
version_id: "v-a".to_string(),
|
||||
success: true,
|
||||
file_info: FileInfo {
|
||||
name: "object-b".to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
error: String::new(),
|
||||
error_code: 0,
|
||||
}],
|
||||
)
|
||||
.pop()
|
||||
.expect("slot 0 should exist")
|
||||
.expect_err("identity mismatch should fail closed");
|
||||
assert!(
|
||||
mismatched.to_string().contains("identity mismatch"),
|
||||
"unexpected mismatch error: {mismatched}"
|
||||
);
|
||||
|
||||
let duplicate = map_batch_read_version_responses(
|
||||
&expected_items,
|
||||
vec![
|
||||
BatchReadVersionResp {
|
||||
index: 0,
|
||||
path: "object-a".to_string(),
|
||||
version_id: "v-a".to_string(),
|
||||
success: true,
|
||||
file_info: FileInfo {
|
||||
name: "object-a".to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
error: String::new(),
|
||||
error_code: 0,
|
||||
},
|
||||
BatchReadVersionResp {
|
||||
index: 0,
|
||||
path: "object-a".to_string(),
|
||||
version_id: "v-a".to_string(),
|
||||
success: true,
|
||||
file_info: FileInfo {
|
||||
name: "object-a".to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
error: String::new(),
|
||||
error_code: 0,
|
||||
},
|
||||
],
|
||||
)
|
||||
.pop()
|
||||
.expect("slot 0 should exist")
|
||||
.expect_err("duplicate response index should fail closed");
|
||||
assert!(
|
||||
duplicate.to_string().contains("duplicate index"),
|
||||
"unexpected duplicate error: {duplicate}"
|
||||
);
|
||||
}
|
||||
|
||||
/// Isolation guard: unobserved objects record nothing (so parallel tests do
|
||||
/// not inflate one another), and a scope clears its own counts on drop.
|
||||
#[tokio::test]
|
||||
|
||||
@@ -735,8 +735,12 @@ pub(crate) use core::io_primitives::disk_call_counters;
|
||||
mod ctx;
|
||||
mod metadata;
|
||||
mod ops;
|
||||
#[cfg(test)]
|
||||
pub(crate) use ops::multipart::NewMultipartUploadCommitObservation;
|
||||
#[cfg(any(test, feature = "test-util"))]
|
||||
pub use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause};
|
||||
#[cfg(test)]
|
||||
pub(crate) use ops::object::DeleteObjectCommitBarrier;
|
||||
#[cfg(feature = "test-util")]
|
||||
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
|
||||
pub(crate) use ops::object::body_cache_plaintext_len;
|
||||
@@ -3025,6 +3029,16 @@ pub struct SetDisks {
|
||||
storage_class_config_override: Arc<std::sync::RwLock<Option<Arc<storageclass::Config>>>>,
|
||||
}
|
||||
|
||||
// DistributedLock sends the raw ObjectKey to its clients; LockRegistry clones
|
||||
// each endpoint's canonical Arc, so an exact Arc set identifies the lock domain.
|
||||
pub(crate) fn same_distributed_lock_domain(left: &[Arc<dyn LockClient>], right: &[Arc<dyn LockClient>]) -> bool {
|
||||
left.iter()
|
||||
.all(|left_client| right.iter().any(|right_client| Arc::ptr_eq(left_client, right_client)))
|
||||
&& right
|
||||
.iter()
|
||||
.all(|right_client| left.iter().any(|left_client| Arc::ptr_eq(left_client, right_client)))
|
||||
}
|
||||
|
||||
const ERASURE_CACHE_MAX_ENTRIES: usize = 32;
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
|
||||
@@ -3600,6 +3614,15 @@ impl SetDisks {
|
||||
&self.ctx
|
||||
}
|
||||
|
||||
/// Whether both sets' namespace-lock implementations cover the same object key.
|
||||
pub(crate) async fn shares_namespace_lock_domain(&self, other: &Self) -> bool {
|
||||
match (self.ctx.is_dist_erasure().await, other.ctx.is_dist_erasure().await) {
|
||||
(false, false) => Arc::ptr_eq(&self.local_lock_manager, &other.local_lock_manager),
|
||||
(true, true) => same_distributed_lock_domain(&self.lockers, &other.lockers),
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
/// The lock manager this set actually uses (test-only; Phase 5 Slice 3).
|
||||
#[cfg(test)]
|
||||
pub(crate) fn local_lock_manager_for_test(&self) -> &Arc<rustfs_lock::GlobalLockManager> {
|
||||
@@ -4584,11 +4607,11 @@ fn should_preserve_delete_replication_state(opts: &ObjectOptions) -> bool {
|
||||
}
|
||||
|
||||
fn should_force_delete_marker_for_missing_version(opts: &ObjectOptions) -> bool {
|
||||
opts.delete_marker || (opts.versioned && opts.version_id.is_none() && !opts.data_movement)
|
||||
opts.delete_marker || ((opts.versioned || opts.version_suspended) && opts.version_id.is_none() && !opts.data_movement)
|
||||
}
|
||||
|
||||
fn resolve_delete_version_state(opts: &ObjectOptions, goi: &ObjectInfo, version_found: bool) -> (bool, bool) {
|
||||
let mut mark_delete = goi.version_id.is_some() || (opts.versioned && opts.version_id.is_none());
|
||||
let mut mark_delete = goi.version_id.is_some() || ((opts.versioned || opts.version_suspended) && opts.version_id.is_none());
|
||||
let mut delete_marker = opts.versioned;
|
||||
|
||||
if opts.version_id.is_some() {
|
||||
|
||||
@@ -32,6 +32,8 @@ use crate::crash_inject::{self, CrashPoint};
|
||||
use crate::multipart_listing::paginate_multipart_listing;
|
||||
use futures::{StreamExt, stream};
|
||||
use std::future::Future;
|
||||
#[cfg(test)]
|
||||
use std::sync::atomic::AtomicBool;
|
||||
#[cfg(any(test, feature = "test-util"))]
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
use std::time::Duration;
|
||||
@@ -65,6 +67,7 @@ impl StaleMultipartCleanupGuard {
|
||||
#[cfg(any(test, feature = "test-util"))]
|
||||
#[derive(Clone, Copy, PartialEq, Eq)]
|
||||
pub enum MultipartCommitPause {
|
||||
NewUploadBeforeLockLost,
|
||||
PutPartBeforeLockAcquire,
|
||||
PutPartBeforeLockLost,
|
||||
PutPartAfterRename,
|
||||
@@ -156,6 +159,72 @@ impl Drop for MultipartCommitBarrier {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
struct NewMultipartUploadCommitObservationState {
|
||||
bucket: String,
|
||||
object: String,
|
||||
committed: AtomicBool,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) struct NewMultipartUploadCommitObservation {
|
||||
state: Arc<NewMultipartUploadCommitObservationState>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
static NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION: std::sync::OnceLock<
|
||||
std::sync::Mutex<Option<Arc<NewMultipartUploadCommitObservationState>>>,
|
||||
> = std::sync::OnceLock::new();
|
||||
|
||||
#[cfg(test)]
|
||||
impl NewMultipartUploadCommitObservation {
|
||||
pub(crate) fn install(bucket: &str, object: &str) -> Self {
|
||||
let state = Arc::new(NewMultipartUploadCommitObservationState {
|
||||
bucket: bucket.to_string(),
|
||||
object: object.to_string(),
|
||||
committed: AtomicBool::new(false),
|
||||
});
|
||||
let mut slot = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
|
||||
.get_or_init(|| std::sync::Mutex::new(None))
|
||||
.lock()
|
||||
.expect("new multipart upload commit observation mutex should not poison");
|
||||
assert!(slot.is_none(), "new multipart upload commit observation must be unique");
|
||||
*slot = Some(Arc::clone(&state));
|
||||
Self { state }
|
||||
}
|
||||
|
||||
pub(crate) fn committed(&self) -> bool {
|
||||
self.state.committed.load(Ordering::Acquire)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl Drop for NewMultipartUploadCommitObservation {
|
||||
fn drop(&mut self) {
|
||||
let mut slot = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
|
||||
.get_or_init(|| std::sync::Mutex::new(None))
|
||||
.lock()
|
||||
.expect("new multipart upload commit observation mutex should not poison");
|
||||
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
|
||||
*slot = None;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn observe_new_multipart_upload_commit(bucket: &str, object: &str) {
|
||||
let state = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
|
||||
.get_or_init(|| std::sync::Mutex::new(None))
|
||||
.lock()
|
||||
.expect("new multipart upload commit observation mutex should not poison")
|
||||
.as_ref()
|
||||
.filter(|state| state.bucket == bucket && state.object == object)
|
||||
.cloned();
|
||||
if let Some(state) = state {
|
||||
state.committed.store(true, Ordering::Release);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(any(test, feature = "test-util"))]
|
||||
async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartCommitPause) {
|
||||
let barrier = {
|
||||
@@ -1615,6 +1684,30 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
||||
|
||||
let upload_path = Self::get_multipart_upload_dir(bucket, object, upload_uuid.as_str(), opts.data_movement);
|
||||
|
||||
#[cfg(any(test, feature = "test-util"))]
|
||||
pause_multipart_commit(bucket, object, MultipartCommitPause::NewUploadBeforeLockLost).await;
|
||||
if _object_lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) {
|
||||
return Err(StorageError::NamespaceLockQuorumUnavailable {
|
||||
mode: "new_multipart_upload_commit",
|
||||
bucket: bucket.to_string(),
|
||||
object: object.to_string(),
|
||||
required: 1,
|
||||
achieved: 0,
|
||||
});
|
||||
}
|
||||
if opts
|
||||
.namespace_lock_fence
|
||||
.as_ref()
|
||||
.is_some_and(NamespaceLockFence::is_lock_lost)
|
||||
{
|
||||
return Err(StorageError::NamespaceLockQuorumUnavailable {
|
||||
mode: "new_multipart_upload_outer_lock",
|
||||
bucket: bucket.to_string(),
|
||||
object: object.to_string(),
|
||||
required: 1,
|
||||
achieved: 0,
|
||||
});
|
||||
}
|
||||
ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?;
|
||||
Self::write_unique_file_info(
|
||||
&shuffle_disks,
|
||||
@@ -1626,6 +1719,8 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
||||
)
|
||||
.await
|
||||
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
|
||||
#[cfg(test)]
|
||||
observe_new_multipart_upload_commit(bucket, object);
|
||||
|
||||
// evalDisks
|
||||
|
||||
|
||||
@@ -1294,7 +1294,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
|
||||
(prepared.snapshot, prepared.object_info)
|
||||
} else {
|
||||
match self
|
||||
.get_object_fileinfo(
|
||||
.get_object_fileinfo_for_get_object_reader(
|
||||
bucket,
|
||||
object,
|
||||
opts,
|
||||
@@ -2497,6 +2497,7 @@ impl SetDisks {
|
||||
})
|
||||
.await?,
|
||||
);
|
||||
notify_put_object_commit_namespace_acquired(bucket, object);
|
||||
}
|
||||
#[cfg(not(any(test, feature = "test-util")))]
|
||||
{
|
||||
@@ -4649,6 +4650,7 @@ struct PutObjectCommitBarrierState {
|
||||
arrived: tokio::sync::Notify,
|
||||
release: tokio::sync::Notify,
|
||||
namespace_pending: tokio::sync::Notify,
|
||||
namespace_acquired: std::sync::atomic::AtomicBool,
|
||||
}
|
||||
|
||||
#[cfg(any(test, feature = "test-util"))]
|
||||
@@ -4670,6 +4672,7 @@ impl PutObjectCommitBarrier {
|
||||
arrived: tokio::sync::Notify::new(),
|
||||
release: tokio::sync::Notify::new(),
|
||||
namespace_pending: tokio::sync::Notify::new(),
|
||||
namespace_acquired: std::sync::atomic::AtomicBool::new(false),
|
||||
});
|
||||
let mut slot = PUT_OBJECT_COMMIT_BARRIER
|
||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||
@@ -4704,6 +4707,10 @@ impl PutObjectCommitBarrier {
|
||||
.await
|
||||
.expect("put object should wait for the namespace lock after leaving the commit barrier");
|
||||
}
|
||||
|
||||
pub fn namespace_acquired(&self) -> bool {
|
||||
self.state.namespace_acquired.load(std::sync::atomic::Ordering::Acquire)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(any(test, feature = "test-util"))]
|
||||
@@ -4760,6 +4767,22 @@ fn notify_put_object_commit_namespace_pending(bucket: &str, object: &str) {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(any(test, feature = "test-util"))]
|
||||
fn notify_put_object_commit_namespace_acquired(bucket: &str, object: &str) {
|
||||
let barrier = PUT_OBJECT_COMMIT_BARRIER
|
||||
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
|
||||
.lock()
|
||||
.expect("put object commit barrier mutex should not poison")
|
||||
.iter()
|
||||
.find(|barrier| {
|
||||
barrier.bucket == bucket && barrier.object == object && barrier.pause == PutObjectCommitPause::BeforeNamespace
|
||||
})
|
||||
.cloned();
|
||||
if let Some(barrier) = barrier {
|
||||
barrier.namespace_acquired.store(true, std::sync::atomic::Ordering::Release);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
struct DeleteObjectCommitBarrierState {
|
||||
bucket: String,
|
||||
@@ -4769,7 +4792,7 @@ struct DeleteObjectCommitBarrierState {
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
struct DeleteObjectCommitBarrier {
|
||||
pub(crate) struct DeleteObjectCommitBarrier {
|
||||
state: Arc<DeleteObjectCommitBarrierState>,
|
||||
}
|
||||
|
||||
@@ -4779,7 +4802,7 @@ static DELETE_OBJECT_COMMIT_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option
|
||||
|
||||
#[cfg(test)]
|
||||
impl DeleteObjectCommitBarrier {
|
||||
fn install(bucket: &str, object: &str) -> Self {
|
||||
pub(crate) fn install(bucket: &str, object: &str) -> Self {
|
||||
let state = Arc::new(DeleteObjectCommitBarrierState {
|
||||
bucket: bucket.to_string(),
|
||||
object: object.to_string(),
|
||||
@@ -4795,13 +4818,13 @@ impl DeleteObjectCommitBarrier {
|
||||
Self { state }
|
||||
}
|
||||
|
||||
async fn wait_until_paused(&self) {
|
||||
pub(crate) async fn wait_until_paused(&self) {
|
||||
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
|
||||
.await
|
||||
.expect("delete object should reach the deterministic commit barrier");
|
||||
}
|
||||
|
||||
fn release(&self) {
|
||||
pub(crate) fn release(&self) {
|
||||
self.state.release.notify_one();
|
||||
}
|
||||
}
|
||||
@@ -5923,6 +5946,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
||||
if dobj.version_id.is_none() && (version_suspended || versioned) {
|
||||
vr.mod_time = Some(OffsetDateTime::now_utc());
|
||||
vr.deleted = true;
|
||||
vr.mark_deleted = true;
|
||||
if versioned {
|
||||
vr.version_id = Some(Uuid::new_v4());
|
||||
}
|
||||
@@ -11809,6 +11833,7 @@ mod transition_upload_integrity_tests {
|
||||
crate::data_movement::SourceCleanupBucketFence {
|
||||
expected_incarnation_id: None,
|
||||
lifecycle_guard: Some(&bucket_guard),
|
||||
..Default::default()
|
||||
},
|
||||
"test_data_movement",
|
||||
)
|
||||
|
||||
@@ -259,10 +259,33 @@ impl SetDisks {
|
||||
read_data: bool,
|
||||
caller_allows_early_stop: bool,
|
||||
) -> Result<GetObjectFileInfo> {
|
||||
self.get_object_fileinfo_gated(bucket, object, opts, read_data, caller_allows_early_stop)
|
||||
self.get_object_fileinfo_gated_inner(bucket, object, opts, read_data, caller_allows_early_stop, false)
|
||||
.await
|
||||
}
|
||||
|
||||
#[tracing::instrument(level = "debug", skip(self))]
|
||||
#[hotpath::measure(impl_type = "SetDisks")]
|
||||
pub(super) async fn get_object_fileinfo_for_get_object_reader(
|
||||
&self,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
opts: &ObjectOptions,
|
||||
read_data: bool,
|
||||
caller_allows_early_stop: bool,
|
||||
) -> Result<GetObjectFileInfo> {
|
||||
let allow_read_version_coalescing = !crate::bucket::utils::is_meta_bucketname(bucket)
|
||||
&& crate::runtime::global::get_metadata_read_version_coalescing_service_ready();
|
||||
self.get_object_fileinfo_gated_inner(
|
||||
bucket,
|
||||
object,
|
||||
opts,
|
||||
read_data,
|
||||
caller_allows_early_stop,
|
||||
allow_read_version_coalescing,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Like `get_object_fileinfo`, but `allow_early_stop=false` forces the full
|
||||
/// quorum fanout. Read-before-write callers (object tagging) must use this:
|
||||
/// the returned online-disk set is the write target, and the early-stop
|
||||
@@ -275,6 +298,20 @@ impl SetDisks {
|
||||
opts: &ObjectOptions,
|
||||
read_data: bool,
|
||||
allow_early_stop: bool,
|
||||
) -> Result<GetObjectFileInfo> {
|
||||
self.get_object_fileinfo_gated_inner(bucket, object, opts, read_data, allow_early_stop, false)
|
||||
.await
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
async fn get_object_fileinfo_gated_inner(
|
||||
&self,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
opts: &ObjectOptions,
|
||||
read_data: bool,
|
||||
allow_early_stop: bool,
|
||||
allow_read_version_coalescing: bool,
|
||||
) -> Result<GetObjectFileInfo> {
|
||||
let vid = opts.version_id.clone().unwrap_or_default();
|
||||
let stage_metrics_enabled = rustfs_io_metrics::get_stage_metrics_enabled();
|
||||
@@ -337,19 +374,34 @@ impl SetDisks {
|
||||
// read_all_fileinfo_observed (see read_all_fileinfo_early_stop in
|
||||
// core/io_primitives.rs); unsafe requests and callers that opt out
|
||||
// (allow_early_stop=false) fall back to full-wait.
|
||||
let (mut parts_metadata, errs, metadata_fanout_diagnostics) = Self::read_all_fileinfo_observed(
|
||||
&disks,
|
||||
"",
|
||||
bucket,
|
||||
object,
|
||||
vid.as_str(),
|
||||
read_data,
|
||||
false,
|
||||
opts.incl_free_versions,
|
||||
allow_early_stop,
|
||||
self.default_parity_count,
|
||||
)
|
||||
.await?;
|
||||
let (mut parts_metadata, errs, metadata_fanout_diagnostics) = if allow_read_version_coalescing {
|
||||
Self::read_all_fileinfo_observed_for_get_object(
|
||||
&disks,
|
||||
"",
|
||||
bucket,
|
||||
object,
|
||||
vid.as_str(),
|
||||
read_data,
|
||||
opts.incl_free_versions,
|
||||
allow_early_stop,
|
||||
self.default_parity_count,
|
||||
)
|
||||
.await?
|
||||
} else {
|
||||
Self::read_all_fileinfo_observed(
|
||||
&disks,
|
||||
"",
|
||||
bucket,
|
||||
object,
|
||||
vid.as_str(),
|
||||
read_data,
|
||||
false,
|
||||
opts.incl_free_versions,
|
||||
allow_early_stop,
|
||||
self.default_parity_count,
|
||||
)
|
||||
.await?
|
||||
};
|
||||
let metadata_metrics_path = if crate::bucket::utils::is_meta_bucketname(bucket) {
|
||||
GET_OBJECT_PATH_INTERNAL_META
|
||||
} else {
|
||||
|
||||
Reference in New Issue
Block a user