refactor(ecstore): sink write + shared-read primitives into set_disk::core::io_primitives (backlog#820) (#4288)

* refactor(ecstore): sink write/rename/delete primitives into set_disk::core::io_primitives (backlog#820)

P5 step 2 of the SetDisks God-Object split (tracking backlog#815, issue
backlog#820; follows step 1 #4285). Relocate the entire set_disk/write.rs
primitive family into the core/io_primitives.rs module home established by
step 1:

- Module items: dangling_delete_grace() + its env consts, and the
  OrphanDirScan scan-result enum.
- impl SetDisks write/rename/delete primitives shared across mod.rs and the
  ops/ operation families: rename_data, commit_rename_data_dir,
  cleanup_multipart_path, rename_part, eval_disks, write_unique_file_info,
  update_object_meta(_with_opts), delete_if_dangling, delete_prefix,
  scan_orphan_dir, purge_orphan_dir_object, check_write_precondition,
  default_read_quorum, default_write_quorum.
- The dangling_delete_grace unit tests move with their subject.

set_disk/write.rs is removed and 'mod write;' dropped from mod.rs.

Pure move + visibility adjustment, zero logic change. Method bodies are moved
verbatim; pub(super) items are widened to pub(in crate::set_disk) so mod.rs /
ops still reach them (write.rs's super was set_disk; the deeper module needs the
explicit path to preserve identical reach). Callers use inherent self./Self::
calls, so no call sites change.

Verification:
- cargo check / clippy -D warnings -p rustfs-ecstore --all-targets: clean
- cargo test -p rustfs-ecstore --lib: 1841 passed, 0 failed (moved
  dangling_delete_grace tests run as set_disk::core::io_primitives::tests::*)
- all five arch guard scripts: pass
- token-stream diff of moved block vs original write.rs: identical modulo
  visibility tokens (2713 tokens each)

* refactor(ecstore): sink shared read primitives into set_disk::core::io_primitives (backlog#820)

P5 step 3 (final) of the SetDisks God-Object split (tracking backlog#815, issue
backlog#820; follows steps 1 #4285 and 2). Relocate the shared, low-level
metadata/erasure READ PRIMITIVE methods out of set_disk/read.rs into the
core/io_primitives.rs module home, leaving the object-read operation itself in
read.rs.

Moved into a new impl SetDisks block in io_primitives.rs (verbatim bodies):
- read_parts, read_all_fileinfo (+ _observed / _inner / _full_wait /
  _early_stop variants), read_all_xl, load_file_info_versions_exact,
  read_all_raw_file_info, pick_latest_quorum_files_info, read_multiple_files.
- The should_allow_metadata_early_stop free helper (called by the moved
  read_all_fileinfo_observed).

Kept in read.rs: the object-read operation and its private helpers
(read_version_optimized, get_object_fileinfo, get_object_info_and_quorum,
try_get_object_direct_data_shards_with_fileinfo, get_object_with_fileinfo,
get_object_decode_reader_with_fileinfo, build_codec_streaming_part_reader) and
the metadata-cache helpers/tests. read.rs reaches the moved primitives through
the SetDisks core (inherent self./Self:: calls; the sole cross-boundary edge,
get_object_fileinfo -> read_all_fileinfo_observed, is handled by widening that
one method to pub(in crate::set_disk)).

Pure move + visibility adjustment, zero logic change. pub(super) primitives are
widened to pub(in crate::set_disk); the three internal-only fanout variants
(_inner/_full_wait/_early_stop) stay private; load_file_info_versions_exact
stays pub(crate). Method bodies are moved verbatim.

Verification:
- cargo check / clippy -D warnings -p rustfs-ecstore --all-targets: clean
- cargo test -p rustfs-ecstore --lib: 1841 passed, 0 failed
- all five arch guard scripts: pass
- token-stream diff of the three moved regions vs original read.rs: identical
  modulo visibility tokens, the new impl-block scaffolding, and rustfmt
  signature re-wrapping (read.rs diff is pure deletion, zero added lines)
This commit is contained in:
Zhengchao An
2026-07-05 22:35:05 +08:00
committed by GitHub
parent 23075518c2
commit 166679a723
4 changed files with 1571 additions and 1568 deletions
File diff suppressed because it is too large Load Diff
-1
View File
@@ -451,7 +451,6 @@ mod ops;
mod read;
mod replication;
pub(crate) mod shard_source;
mod write;
/// Get lock acquire timeout from environment variable RUSTFS_LOCK_ACQUIRE_TIMEOUT (in seconds)
/// Defaults to 30 seconds if not set or invalid
-626
View File
@@ -120,344 +120,6 @@ impl SetDisks {
.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() || is_version_early_stop_enabled());
let allow_early_stop = observe && should_allow_metadata_early_stop(read_data, version_id, healing, incl_free_versions);
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 futures to complete
let results = join_all(futures).await;
for join_result in results {
match join_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(_join_err) => {
// A spawned task panicked — treat as unexpected disk error
if let Some(observations) = &mut observations {
observations.push(MetadataFanoutObservation::from_error(&DiskError::Unexpected, Duration::ZERO));
}
ress.push(FileInfo::default());
errors.push(Some(DiskError::Unexpected));
}
}
}
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,
);
while join_set.join_next().await.is_some() {}
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,
@@ -498,285 +160,6 @@ impl SetDisks {
}
}
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,
@@ -1772,15 +1155,6 @@ impl SetDisks {
}
}
fn should_allow_metadata_early_stop(read_data: bool, version_id: &str, healing: bool, incl_free_versions: bool) -> bool {
if read_data {
return false;
}
(is_get_metadata_early_stop_enabled() && version_id.is_empty() && !healing && !incl_free_versions)
|| (is_version_early_stop_enabled() && !version_id.is_empty() && !healing)
}
fn get_object_metadata_cache_request_bypass_reason(bucket: &str, opts: &ObjectOptions, read_data: bool) -> Option<&'static str> {
if !read_data {
return Some(GET_METADATA_CACHE_REASON_NOT_READ_DATA);
-941
View File
@@ -1,941 +0,0 @@
// 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::*;
/// Grace window during which a recently modified object is never deleted as
/// dangling. 0 disables the grace window.
const ENV_HEAL_DANGLING_DELETE_GRACE_SECS: &str = "RUSTFS_HEAL_DANGLING_DELETE_GRACE_SECS";
const DEFAULT_HEAL_DANGLING_DELETE_GRACE_SECS: u64 = 3600;
fn dangling_delete_grace() -> time::Duration {
let secs = rustfs_utils::get_env_u64(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, DEFAULT_HEAL_DANGLING_DELETE_GRACE_SECS);
time::Duration::seconds(i64::try_from(secs).unwrap_or(i64::MAX))
}
/// Result of scanning one disk's copy of a directory prefix while deciding
/// whether an orphan (metadata-less) directory tree can be safely purged.
enum OrphanDirScan {
/// The subtree holds at least one regular file (object metadata or data), so
/// it is a real object and must not be purged.
HasData,
/// The prefix exists on this disk and contains only nested empty directories.
/// Carries every directory path in pre-order (parents before children).
Empty(Vec<String>),
/// The prefix does not exist on this disk.
Missing,
}
impl SetDisks {
pub(super) fn default_read_quorum(&self) -> usize {
self.set_drive_count - self.default_parity_count
}
pub(super) fn default_write_quorum(&self) -> usize {
let mut data_count = self.set_drive_count - self.default_parity_count;
if data_count == self.default_parity_count {
data_count += 1
}
data_count
}
#[tracing::instrument(level = "debug", skip(disks, file_infos))]
#[allow(clippy::type_complexity)]
pub(super) async fn rename_data(
disks: &[Option<DiskStore>],
src_bucket: &str,
src_object: &str,
file_infos: &[FileInfo],
dst_bucket: &str,
dst_object: &str,
write_quorum: usize,
) -> disk::error::Result<(Vec<Option<DiskStore>>, Option<Vec<u8>>, Option<Uuid>, Vec<Option<DiskStore>>)> {
let mut futures = Vec::with_capacity(disks.len());
let mut errs = Vec::with_capacity(disks.len());
let src_bucket = Arc::new(src_bucket.to_string());
let src_object = Arc::new(src_object.to_string());
let dst_bucket = Arc::new(dst_bucket.to_string());
let dst_object = Arc::new(dst_object.to_string());
for (i, (disk, file_info)) in disks.iter().zip(file_infos.iter()).enumerate() {
let mut file_info = file_info.clone();
let disk = disk.clone();
let src_bucket = src_bucket.clone();
let src_object = src_object.clone();
let dst_object = dst_object.clone();
let dst_bucket = dst_bucket.clone();
futures.push(tokio::spawn(async move {
if file_info.erasure.index == 0 {
file_info.erasure.index = i + 1;
}
if !file_info.is_valid() {
return Err(DiskError::FileCorrupt);
}
if let Some(disk) = disk {
disk.rename_data(&src_bucket, &src_object, file_info, &dst_bucket, &dst_object)
.await
} else {
Err(DiskError::DiskNotFound)
}
}));
}
let mut disk_versions = vec![None; disks.len()];
let mut data_dirs = vec![None; disks.len()];
let results = join_all(futures).await;
for (idx, result) in results.iter().enumerate() {
match result.as_ref().map_err(|_| DiskError::Unexpected)? {
Ok(res) => {
data_dirs[idx] = res.old_data_dir;
disk_versions[idx].clone_from(&res.sign);
errs.push(None);
}
Err(e) => {
errs.push(Some(e.clone()));
}
}
}
if issue3031_diag_enabled() {
let success_count = errs.iter().filter(|err| err.is_none()).count();
let failure_count = errs.len().saturating_sub(success_count);
let ignored_failure_count = errs
.iter()
.filter(|err| err.as_ref().is_some_and(|err| OBJECT_OP_IGNORED_ERRS.contains(err)))
.count();
let data_dir_vote_count = data_dirs.iter().filter(|data_dir| data_dir.is_some()).count();
let reduced_data_dir = Self::reduce_common_data_dir(&data_dirs, write_quorum);
warn!(
target: "rustfs_ecstore::set_disk",
src_bucket = %src_bucket,
src_object = %src_object,
dst_bucket = %dst_bucket,
dst_object = %dst_object,
write_quorum,
disk_count = errs.len(),
success_count,
failure_count,
ignored_failure_count,
data_dir_vote_count,
reduced_data_dir = ?reduced_data_dir,
errs = ?errs,
data_dirs = ?data_dirs,
"issue3031_rename_data_quorum_context"
);
}
let mut futures = Vec::with_capacity(disks.len());
if let Some(ret_err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) {
// TODO: add concurrency
for (i, err) in errs.iter().enumerate() {
if err.is_some() {
continue;
}
if let Some(disk) = disks[i].as_ref() {
let fi = file_infos[i].clone();
let old_data_dir = data_dirs[i];
let disk = disk.clone();
let dst_bucket = dst_bucket.clone();
let dst_object = dst_object.clone();
futures.push(tokio::spawn(async move {
disk.delete_version(
&dst_bucket,
&dst_object,
fi,
false,
DeleteOptions {
undo_write: true,
old_data_dir,
..Default::default()
},
)
.await
}));
}
}
if issue3031_diag_enabled() {
warn!(
target: "rustfs_ecstore::set_disk",
src_bucket = %src_bucket,
src_object = %src_object,
dst_bucket = %dst_bucket,
dst_object = %dst_object,
write_quorum,
ret_err = %ret_err,
errs = ?errs,
data_dirs = ?data_dirs,
"issue3031_rename_data_quorum_failed"
);
}
let undo_results = join_all(futures).await;
let undo_error_count = undo_results
.iter()
.filter(|result| match result {
Err(_) | Ok(Err(_)) => true,
Ok(Ok(_)) => false,
})
.count();
if undo_error_count > 0 {
warn!(
target: "rustfs_ecstore::set_disk",
dst_bucket = %dst_bucket,
dst_object = %dst_object,
undo_error_count,
"rename_data quorum rollback reported errors"
);
}
return Err(ret_err);
}
let versions = None;
// TODO: reduceCommonVersions
let data_dir = Self::reduce_common_data_dir(&data_dirs, write_quorum);
let online_disks = Self::eval_disks(disks, &errs);
let cleanup_disks = if let Some(data_dir) = data_dir {
disks
.iter()
.zip(errs.iter())
.zip(data_dirs.iter())
.map(|((disk, err), old_data_dir)| {
if err.is_none() && *old_data_dir == Some(data_dir) {
disk.clone()
} else {
None
}
})
.collect()
} else {
vec![None; disks.len()]
};
// // TODO: reduce_common_data_dir
// if let Some(old_dir) = rename_ress
// .iter()
// .filter_map(|v| if v.is_some() { v.as_ref().unwrap().old_data_dir } else { None })
// .map(|v| v.to_string())
// .next()
// {
// let cm_errs = self.commit_rename_data_dir(&shuffle_disks, &bucket, &object, &old_dir).await;
// warn!("put_object commit_rename_data_dir:{:?}", &cm_errs);
// }
// self.delete_all(RUSTFS_META_TMP_BUCKET, &tmp_dir).await?;
Ok((online_disks, versions, data_dir, cleanup_disks))
}
#[allow(dead_code)]
#[tracing::instrument(level = "debug", skip(self, disks))]
pub(super) async fn commit_rename_data_dir(
&self,
disks: &[Option<DiskStore>],
bucket: &str,
object: &str,
data_dir: &str,
write_quorum: usize,
) -> disk::error::Result<()> {
let file_path = Arc::new(format!("{object}/{data_dir}"));
let bucket = Arc::new(bucket.to_string());
let futures = disks.iter().map(|disk| {
let file_path = file_path.clone();
let bucket = bucket.clone();
let disk = disk.clone();
tokio::spawn(async move {
if let Some(disk) = disk {
(disk
.delete(
&bucket,
&file_path,
DeleteOptions {
recursive: true,
..Default::default()
},
)
.await)
.err()
} else {
Some(DiskError::DiskNotFound)
}
})
});
let errs: Vec<Option<DiskError>> = join_all(futures)
.await
.into_iter()
.map(|e| e.unwrap_or(Some(DiskError::Unexpected)))
.collect();
if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) {
return Err(err);
}
Ok(())
}
#[tracing::instrument(skip(self))]
pub(super) async fn cleanup_multipart_path(&self, paths: &[String]) {
if paths.is_empty() {
return;
}
let disks = self.get_disks_internal().await;
let mut errs = Vec::with_capacity(disks.len());
// Use improved simple batch processor instead of join_all for better performance
let processor = runtime_sources::batch_processors().write_processor();
let tasks: Vec<_> = disks
.iter()
.map(|disk| {
let disk = disk.clone();
let paths = paths.to_vec();
async move {
if let Some(disk) = disk {
disk.delete_paths(RUSTFS_META_MULTIPART_BUCKET, &paths).await
} else {
Err(DiskError::DiskNotFound)
}
}
})
.collect();
let results = processor.execute_batch(tasks).await;
for result in results {
match result {
Ok(_) => {
errs.push(None);
}
Err(e) => {
errs.push(Some(e));
}
}
}
if errs.iter().any(|e| e.is_some()) {
warn!("cleanup_multipart_path errs {:?}", &errs);
}
}
#[tracing::instrument(skip(disks, meta))]
#[allow(clippy::too_many_arguments)]
pub(super) async fn rename_part(
&self,
disks: &[Option<DiskStore>],
src_bucket: &str,
src_object: &str,
dst_bucket: &str,
dst_object: &str,
meta: Bytes,
write_quorum: usize,
quorum_context: Option<MultipartWriteQuorumContext<'_>>,
) -> disk::error::Result<Vec<Option<DiskStore>>> {
let src_bucket = Arc::new(src_bucket.to_string());
let src_object = Arc::new(src_object.to_string());
let dst_bucket = Arc::new(dst_bucket.to_string());
let dst_object = Arc::new(dst_object.to_string());
// Match MinIO's multipart overwrite semantics: clear any stale destination
// part payload and metadata before the new per-disk rename fan-out begins.
self.cleanup_multipart_path(&[dst_object.to_string(), format!("{dst_object}.meta")])
.await;
let mut errs = Vec::with_capacity(disks.len());
let futures = disks.iter().map(|disk| {
let disk = disk.clone();
let meta = meta.clone();
let src_bucket = src_bucket.clone();
let src_object = src_object.clone();
let dst_bucket = dst_bucket.clone();
let dst_object = dst_object.clone();
tokio::spawn(async move {
if let Some(disk) = disk {
disk.rename_part(&src_bucket, &src_object, &dst_bucket, &dst_object, meta)
.await
} else {
Err(DiskError::DiskNotFound)
}
})
});
let results = join_all(futures).await;
for result in results {
match result? {
Ok(_) => {
errs.push(None);
}
Err(e) => {
errs.push(Some(e));
}
}
}
if issue3031_diag_enabled() {
let success_count = errs.iter().filter(|err| err.is_none()).count();
let error_count = errs.len().saturating_sub(success_count);
let disk_not_found_count = errs.iter().filter(|err| matches!(err, Some(DiskError::DiskNotFound))).count();
let file_not_found_count = errs.iter().filter(|err| matches!(err, Some(DiskError::FileNotFound))).count();
warn!(
target: "rustfs_ecstore::set_disk",
src_bucket = %src_bucket,
src_object = %src_object,
dst_bucket = %dst_bucket,
dst_object = %dst_object,
write_quorum = write_quorum,
disk_count = errs.len(),
success_count = success_count,
error_count = error_count,
disk_not_found_count = disk_not_found_count,
file_not_found_count = file_not_found_count,
"issue3031_rename_part_context"
);
}
if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) {
if let Some(context) = quorum_context {
log_multipart_write_quorum_failure(context, &errs, write_quorum, &err);
} else {
warn!("rename_part errs {:?}", &errs);
}
self.cleanup_multipart_path(&[dst_object.to_string(), format!("{dst_object}.meta")])
.await;
return Err(err);
}
let disks = Self::eval_disks(disks, &errs);
Ok(disks)
}
pub(super) fn eval_disks(disks: &[Option<DiskStore>], errs: &[Option<DiskError>]) -> Vec<Option<DiskStore>> {
if disks.len() != errs.len() {
return Vec::new();
}
let mut online_disks = vec![None; disks.len()];
for (i, err_op) in errs.iter().enumerate() {
if err_op.is_none() {
online_disks[i].clone_from(&disks[i]);
}
}
online_disks
}
#[tracing::instrument(skip(disks, files))]
pub(super) async fn write_unique_file_info(
disks: &[Option<DiskStore>],
org_bucket: &str,
bucket: &str,
prefix: &str,
files: &[FileInfo],
write_quorum: usize,
) -> disk::error::Result<()> {
let mut futures = Vec::with_capacity(disks.len());
let mut errs = Vec::with_capacity(disks.len());
for (i, disk) in disks.iter().enumerate() {
let mut file_info = files[i].clone();
file_info.erasure.index = i + 1;
futures.push(async move {
if let Some(disk) = disk {
disk.write_metadata(org_bucket, bucket, prefix, file_info).await
} else {
Err(DiskError::DiskNotFound)
}
});
}
let results = join_all(futures).await;
for result in results {
match result {
Ok(_) => {
errs.push(None);
}
Err(e) => {
errs.push(Some(e));
}
}
}
if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) {
let mut revert_futures = Vec::with_capacity(disks.len());
for (i, err) in errs.iter().enumerate() {
if err.is_some() {
continue;
}
if let Some(disk) = disks[i].as_ref() {
let disk = disk.clone();
let bucket = bucket.to_string();
let path = path_join_buf(&[prefix, STORAGE_FORMAT_FILE]);
revert_futures.push(async move {
if let Err(err) = disk
.delete(
&bucket,
&path,
DeleteOptions {
recursive: true,
..Default::default()
},
)
.await
{
warn!("write meta revert err {:?}", err);
}
});
}
}
join_all(revert_futures).await;
return Err(err);
}
Ok(())
}
pub(super) async fn update_object_meta(
&self,
bucket: &str,
object: &str,
fi: FileInfo,
disks: &[Option<DiskStore>],
) -> disk::error::Result<()> {
self.update_object_meta_with_opts(bucket, object, fi, disks, &UpdateMetadataOpts::default())
.await
}
pub(super) async fn update_object_meta_with_opts(
&self,
bucket: &str,
object: &str,
fi: FileInfo,
disks: &[Option<DiskStore>],
opts: &UpdateMetadataOpts,
) -> disk::error::Result<()> {
if fi.metadata.is_empty() && !opts.replace_user_metadata {
return Ok(());
}
self.invalidate_get_object_metadata_cache(bucket, object).await;
let mut futures = Vec::with_capacity(disks.len());
let mut errs = Vec::with_capacity(disks.len());
for disk in disks.iter() {
let fi = fi.clone();
futures.push(async move {
if let Some(disk) = disk {
disk.update_metadata(bucket, object, fi, opts).await
} else {
Err(DiskError::DiskNotFound)
}
})
}
let results = join_all(futures).await;
for result in results {
match result {
Ok(_) => {
errs.push(None);
}
Err(e) => {
errs.push(Some(e));
}
}
}
if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, fi.write_quorum(self.default_write_quorum())) {
return Err(err);
}
self.invalidate_get_object_metadata_cache(bucket, object).await;
Ok(())
}
pub(super) async fn delete_if_dangling(
&self,
bucket: &str,
object: &str,
meta_arr: &[FileInfo],
errs: &[Option<DiskError>],
data_errs_by_part: &HashMap<usize, Vec<usize>>,
opts: ObjectOptions,
) -> disk::error::Result<FileInfo> {
let (m, can_heal) = is_object_dangling(meta_arr, errs, data_errs_by_part);
if !can_heal {
return Err(DiskError::ErasureReadQuorum);
}
// Recently written objects get a grace window before dangling cleanup: after
// an unclean shutdown some disks may still be catching up (or carry writes
// that were never made durable), and deleting the surviving shards right away
// turns a partial loss into a total one. Skip deletion and leave the object
// for a later heal/scanner pass to re-evaluate.
if m.is_valid()
&& let Some(mod_time) = m.mod_time
{
let grace = dangling_delete_grace();
if !grace.is_zero() && OffsetDateTime::now_utc() - mod_time < grace {
info!(
bucket = bucket,
object = object,
mod_time = %mod_time,
grace_secs = grace.whole_seconds(),
"skipping dangling-object deletion within grace window"
);
return Err(DiskError::ErasureReadQuorum);
}
}
let mut tags: HashMap<String, String> = HashMap::new();
tags.insert("set".to_string(), self.set_index.to_string());
tags.insert("pool".to_string(), self.pool_index.to_string());
tags.insert("merrs".to_string(), join_errs(errs));
tags.insert("derrs".to_string(), format!("{data_errs_by_part:?}"));
if m.is_valid() {
tags.insert("sz".to_string(), m.size.to_string());
tags.insert(
"mt".to_string(),
m.mod_time
.as_ref()
.map_or(String::new(), |mod_time| mod_time.unix_timestamp().to_string()),
);
tags.insert("d:p".to_string(), format!("{}:{}", m.erasure.data_blocks, m.erasure.parity_blocks));
} else {
tags.insert("invalid".to_string(), "1".to_string());
tags.insert(
"d:p".to_string(),
format!("{}:{}", self.set_drive_count - self.default_parity_count, self.default_parity_count),
);
}
let mut offline = 0;
for (i, err) in errs.iter().enumerate() {
let mut found = false;
if let Some(err) = err
&& err == &DiskError::DiskNotFound
{
found = true;
}
for p in data_errs_by_part {
if let Some(v) = p.1.get(i)
&& *v == CHECK_PART_DISK_NOT_FOUND
{
found = true;
break;
}
}
if found {
offline += 1;
}
}
if offline > 0 {
tags.insert("offline".to_string(), offline.to_string());
}
let mut fi = FileInfo::default();
if let Some(ref version_id) = opts.version_id {
fi.version_id = Uuid::parse_str(version_id).ok();
}
fi.set_tier_free_version_id(&Uuid::new_v4().to_string());
let disks = self.get_disks_internal().await;
let mut futures = Vec::with_capacity(disks.len());
for disk_op in disks.iter() {
let bucket = bucket.to_string();
let object = object.to_string();
let fi = fi.clone();
futures.push(async move {
if let Some(disk) = disk_op {
disk.delete_version(&bucket, &object, fi, false, DeleteOptions::default())
.await
} else {
Err(DiskError::DiskNotFound)
}
});
}
let results = join_all(futures).await;
let mut delete_errs = Vec::with_capacity(results.len());
for (index, result) in results.into_iter().enumerate() {
let key = format!("ddisk-{index}");
let already_absent = matches!(
errs.get(index).and_then(Option::as_ref),
Some(DiskError::FileNotFound | DiskError::FileVersionNotFound)
);
match result {
Ok(_) => {
tags.insert(key, "<nil>".to_string());
delete_errs.push(None);
}
Err(e) => {
tags.insert(key, e.to_string());
if already_absent || matches!(&e, DiskError::FileNotFound | DiskError::FileVersionNotFound) {
delete_errs.push(None);
} else {
delete_errs.push(Some(e));
}
}
}
}
let write_quorum = if m.is_valid() {
m.write_quorum(self.default_write_quorum())
} else {
self.default_write_quorum()
};
if let Some(err) = reduce_write_quorum_errs(&delete_errs, OBJECT_OP_IGNORED_ERRS, write_quorum) {
return Err(err);
}
Ok(m)
}
pub(super) async fn delete_prefix(&self, bucket: &str, prefix: &str) -> disk::error::Result<()> {
let disks = self.get_disks_internal().await;
let write_quorum = disks.len() / 2 + 1;
let mut futures = Vec::with_capacity(disks.len());
for disk_op in disks.iter() {
let bucket = bucket.to_string();
let prefix = prefix.to_string();
futures.push(async move {
if let Some(disk) = disk_op {
disk.delete(
&bucket,
&prefix,
DeleteOptions {
recursive: true,
immediate: true,
..Default::default()
},
)
.await
} else {
Ok(())
}
});
}
let errs = join_all(futures).await.into_iter().map(|v| v.err()).collect::<Vec<_>>();
if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) {
return Err(err);
}
Ok(())
}
/// Scan a single disk's copy of `prefix` and decide whether it is an orphan
/// (metadata-less) directory subtree. Walks the tree iteratively and returns
/// [`OrphanDirScan::HasData`] as soon as any regular file is found.
async fn scan_orphan_dir(disk: &DiskStore, bucket: &str, prefix: &str) -> OrphanDirScan {
let root = prefix.trim_end_matches(SLASH_SEPARATOR).to_string();
let mut stack = vec![root.clone()];
// Pre-order list of directories (a parent always precedes its descendants),
// so reversing it yields a safe children-first removal order.
let mut dirs: Vec<String> = Vec::new();
let mut existed = false;
while let Some(dir) = stack.pop() {
let entries = match disk.list_dir("", bucket, &dir, 0).await {
Ok(entries) => entries,
Err(_) => {
// The root missing (or never existing) means there is nothing to
// purge on this disk. A nested directory vanishing mid-scan is a
// benign race, so skip it and keep walking.
if dir == root {
return OrphanDirScan::Missing;
}
continue;
}
};
existed = true;
dirs.push(dir.clone());
for entry in entries {
match entry.strip_suffix(SLASH_SEPARATOR) {
// `read_dir` marks directories with a trailing slash; anything else
// is a regular file, which means real object data lives here.
Some(child) => stack.push(format!("{dir}{SLASH_SEPARATOR}{child}")),
None => return OrphanDirScan::HasData,
}
}
}
if existed {
OrphanDirScan::Empty(dirs)
} else {
OrphanDirScan::Missing
}
}
/// Purge an orphan directory prefix — a trailing-slash key that exists on disk
/// as an empty directory tree with no object metadata on any disk of this set.
/// Such prefixes are listable (see `scan_dir`) yet are not real objects, so the
/// normal delete path returns NotFound and leaves them stranded (issue #4189).
///
/// Callers pass the *decoded* directory name (`prefix/`), not the `__XLDIR__`
/// encoded object key — the orphan tree on disk uses the plain path.
///
/// Returns `Ok(true)` when the prefix was an orphan tree on this set and was
/// removed, `Ok(false)` when it holds real data or does not exist on any disk
/// of this set (the caller should surface the original NotFound), and `Err` on
/// a hard disk failure.
pub(crate) async fn purge_orphan_dir_object(&self, bucket: &str, object: &str) -> disk::error::Result<bool> {
let disks = self.get_disks_internal().await;
// Phase 1: classify every online disk. Refuse to purge if ANY disk holds
// object data under the prefix, so a degraded/healable object is never
// destroyed.
let mut per_disk_dirs: Vec<(usize, Vec<String>)> = Vec::new();
let mut existed = false;
for (i, disk) in disks.iter().enumerate() {
let Some(disk) = disk else { continue };
match Self::scan_orphan_dir(disk, bucket, object).await {
OrphanDirScan::HasData => return Ok(false),
OrphanDirScan::Empty(dirs) => {
existed = true;
per_disk_dirs.push((i, dirs));
}
OrphanDirScan::Missing => {}
}
}
if !existed {
return Ok(false);
}
// Phase 2: remove the empty directories children-first on each disk. A
// non-recursive delete performs an empty-only `rmdir`, so a directory that
// concurrently gained an object fails with DirectoryNotEmpty and is skipped —
// a racing PutObject is never clobbered.
for (i, mut dirs) in per_disk_dirs {
let Some(disk) = disks[i].as_ref() else { continue };
dirs.reverse();
for dir in dirs {
if let Err(err) = disk
.delete(
bucket,
&dir,
DeleteOptions {
recursive: false,
immediate: true,
..Default::default()
},
)
.await
{
// Best effort: a sibling removal may have already cleared a shared
// parent, or a concurrent writer repopulated the directory. Neither
// is fatal to purging the orphan tree.
debug!(bucket, object, dir, error = ?err, "purge_orphan_dir_object: skipped non-empty/absent directory");
}
}
}
Ok(true)
}
pub(super) async fn check_write_precondition(
&self,
bucket: &str,
object: &str,
opts: &ObjectOptions,
) -> Option<StorageError> {
let mut opts = opts.clone();
let http_preconditions = opts.http_preconditions?;
opts.http_preconditions = None;
// Never claim a lock here, to avoid deadlock
// - If no_lock is false, we must have obtained the lock out side of this function
// - If no_lock is true, we should not obtain locks
opts.no_lock = true;
let oi = self.get_object_info(bucket, object, &opts).await;
match oi {
Ok(oi) => {
// If top level is a delete marker proceed to upload.
if oi.delete_marker {
return None;
}
let if_none_match = http_preconditions.if_none_match_value().map(str::to_owned);
let if_match = http_preconditions.if_match_value().map(str::to_owned);
if should_prevent_write(&oi, if_none_match, if_match) {
return Some(StorageError::PreconditionFailed);
}
}
Err(StorageError::VersionNotFound(_, _, _))
| Err(StorageError::ObjectNotFound(_, _))
| Err(StorageError::ErasureReadQuorum) => {
// When the object is not found,
// - if If-Match is set, we should return 404 NotFound
// - if If-None-Match is set, we should be able to proceed with the request
if http_preconditions.if_match_value().is_some() {
return Some(StorageError::ObjectNotFound(bucket.to_string(), object.to_string()));
}
}
Err(e) => {
return Some(e);
}
}
None
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn dangling_delete_grace_defaults_to_one_hour() {
temp_env::with_var(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, None::<&str>, || {
assert_eq!(dangling_delete_grace(), time::Duration::seconds(3600));
});
}
#[test]
fn dangling_delete_grace_env_override_and_disable() {
temp_env::with_var(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, Some("120"), || {
assert_eq!(dangling_delete_grace(), time::Duration::seconds(120));
});
temp_env::with_var(ENV_HEAL_DANGLING_DELETE_GRACE_SECS, Some("0"), || {
assert!(dangling_delete_grace().is_zero());
});
}
}