mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-27 15:37:02 +00:00
fix(storage): resolve erasure parity per pool (#4977)
* fix(filemeta): add state-aware file info validation
* fix(filemeta): validate shard arithmetic and delete paths
* fix(ecstore): add fallible erasure construction
* fix(ecstore): resolve storage parity per pool
* fix(storage): report heterogeneous erasure layouts
* fix(admin): publish prepared storage config atomically
* fix(storage): harden per-pool parity boundaries
* fix(storage): address pre-PR validation findings
* test(ci): fix strict-topology validation fixtures
* fix(heal): preserve delete markers during repair
* refactor(filemeta): drop unused ValidatedFileInfo witness
ValidatedFileInfo wrapped an unread `_file_info` reference alongside an `Option<ValidatedErasureLayout>`, but only the layout was ever consumed. Return the layout directly from `FileInfo::validate` so the sole production consumer (`LocalDisk::check_parts`) and the two unit tests read it without the extra witness type and lifetime.
No behavior change.
* fix(filemeta): keep compressed and MinIO-migrated tiered objects readable
The new decode-path validation rejected several legitimate on-disk shapes that older RustFS and MinIO-migrated data carry, turning readable objects into FileCorrupt:
- Compressed objects written with an unknown upload size persist a negative per-part actual_size (the documented "unknown size" sentinel that ObjectInfo::get_actual_size already tolerates). validate_collection_contents rejected it via usize::try_from; now a negative actual_size skips shard validation and only real, non-negative sizes are checked.
- MinIO-migrated objects transitioned to a versioned remote tier store the tier version id as a UUID string, not 16 raw bytes. MetaObject::into_fileinfo returned FileCorrupt (main tolerated it as None), making all versions of the object unreadable; MetaDeleteMarker free-version records took a Some(nil) sentinel path with the same effect, which also breaks free-version expiry (remote-tier leak). Both now decode through a shared transitioned_version_id_from_meta_sys helper: 16 raw bytes or a UUID string are accepted, anything else is tolerated as None instead of failing the read.
Regression tests updated to assert the readable/compat behavior, with new tests covering MinIO string-form recovery.
* fix(scanner): build the delete-marker test fixture without erasure geometry
get_size_counts_delete_markers_separately_from_versions built its delete marker with `FileInfo::new(object, 1, 1)`, which attaches erasure geometry (data=1/parity=1/distribution). This PR classifies versions by shape via `is_storage_delete_marker()` (no geometry) rather than the raw `deleted` flag, so a geometry-bearing "delete marker" is correctly serialized as a purge-pending payload Object and counted as a version — CI saw summary.versions=3, expected 2.
Real delete markers carry no erasure geometry (delete paths build them as `FileInfo { deleted: true, ..Default::default() }`), so construct the fixture the same way. It then classifies as a storage delete marker and the counts (versions=2, delete_markers=1) hold. This keeps the PR's more-correct classification, which prevents a purge-pending object's geometry from being dropped when serialized as a bare delete marker.
* docs(changelog): note per-pool parity fix and storage-class startup upgrade caveat
Records the #4801 per-pool erasure parity fix under Fixed, and documents the upgrade behavior where a persisted storage class that a small or heterogeneous pool cannot satisfy now fails startup — with the RUSTFS_STORAGE_CLASS_STANDARD recovery steps. Docs-only; covers R4 from the on-disk compatibility audit.
* fix(heal): report parity from erasure geometry, not is_valid()
heal_object set HealResultItem.parity_blocks via `if lfi.is_valid()`, which was missed by the migration of the other quorum/metadata predicates. With the new `is_valid()` semantics (full payload validation; delete markers now return false), a delete marker or a geometry-bearing version with a benign collection quirk would misreport parity as the pool default instead of its own. Use `has_valid_erasure_geometry()` — the narrow "does this carry erasure geometry" predicate the rest of the migration uses — so reporting matches the object's actual layout. Reporting-only; no data-path change.
* fix(filemeta): do not silently serialize a non-canonical deleted FileInfo as an Object
`From<FileInfo> for FileMetaVersion` classifies by `is_storage_delete_marker()` (shape), which correctly routes canonical delete markers to Delete and purge-pending payloads (deleted=true with real erasure geometry) to Object. But a `deleted` FileInfo that is neither a canonical marker nor a valid erasure payload would silently serialize as a zero-geometry MetaObject that later fails `validate_for_metadata_read`. Write paths validate first (`validate_for_erasure_write` / `validate_for_metadata_read`), so this is a caller bug; `From` is infallible, so surface it with a structured `warn!` on the malformed branch instead of writing corrupt metadata silently. Legitimate purge-pending objects (valid geometry) are unaffected — the guard only fires for `deleted && !has_valid_erasure_geometry()`.
* test(filemeta): assert real historical xl.meta versions pass metadata-read validation
Empirical companion to the code-reasoned decode-tolerance invariants (docs/architecture/erasure-coding.md §11) and the rolling-upgrade / MinIO-migration compatibility concern: the tightened `validate_for_metadata_read` runs on every local disk read and peer-RPC-decoded FileInfo, so it must accept every version of real historically-written xl.meta, never reject it as FileCorrupt.
Loads five real fixtures — MinIO small-inline, MinIO versioned (two object versions + a delete marker), MinIO large multipart, a legacy V1 (xl.json-derived) object, and a legacy meta_ver 2 object — decodes every version with parts materialized, and asserts validate_for_metadata_read() is Ok for each. Reverting the tolerant handling (delete-marker shape, legacy per-part checksums, string/short transitioned-versionID, negative actual_size) turns this red.
* fix(ci): remove duplicate storage test re-exports
---------
Co-authored-by: overtrue <anzhengchao@gmail.com>
This commit is contained in:
@@ -164,7 +164,7 @@ pub(in crate::set_disk) struct MetadataFanoutObservation {
|
||||
|
||||
impl MetadataFanoutObservation {
|
||||
pub(in crate::set_disk) fn from_file_info(file_info: &FileInfo, elapsed: Duration) -> Self {
|
||||
if file_info.is_valid() {
|
||||
if file_info_is_valid_for_metadata(file_info) {
|
||||
Self {
|
||||
outcome: GET_METADATA_RESPONSE_VALID,
|
||||
elapsed,
|
||||
@@ -301,6 +301,7 @@ pub(in crate::set_disk) struct MetadataQuorumAccumulator {
|
||||
pub(in crate::set_disk) candidate_votes: usize,
|
||||
pub(in crate::set_disk) conflicting_metadata: bool,
|
||||
pub(in crate::set_disk) delete_marker_seen: bool,
|
||||
pub(in crate::set_disk) delete_marker_candidates: Vec<(FileInfo, usize)>,
|
||||
pub(in crate::set_disk) delete_marker_votes: usize,
|
||||
pub(in crate::set_disk) requested_version_id: String,
|
||||
pub(in crate::set_disk) matching_version_votes: usize,
|
||||
@@ -321,6 +322,7 @@ impl MetadataQuorumAccumulator {
|
||||
candidate_votes: 0,
|
||||
conflicting_metadata: false,
|
||||
delete_marker_seen: false,
|
||||
delete_marker_candidates: Vec::new(),
|
||||
delete_marker_votes: 0,
|
||||
requested_version_id: String::new(),
|
||||
matching_version_votes: 0,
|
||||
@@ -333,7 +335,7 @@ impl MetadataQuorumAccumulator {
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn observe_file_info(&mut self, file_info: &FileInfo) {
|
||||
if !file_info.is_valid() {
|
||||
if !file_info_is_valid_for_metadata(file_info) {
|
||||
self.hard_errors = self.hard_errors.saturating_add(1);
|
||||
return;
|
||||
}
|
||||
@@ -348,9 +350,24 @@ impl MetadataQuorumAccumulator {
|
||||
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);
|
||||
if file_info.is_canonical_delete_marker() {
|
||||
self.delete_marker_seen = true;
|
||||
if let Some((_, votes)) = self
|
||||
.delete_marker_candidates
|
||||
.iter_mut()
|
||||
.find(|(candidate, _)| metadata_early_stop_candidate_matches(candidate, file_info))
|
||||
{
|
||||
*votes = votes.saturating_add(1);
|
||||
} else {
|
||||
self.delete_marker_candidates.push((file_info.clone(), 1));
|
||||
}
|
||||
self.delete_marker_votes = self
|
||||
.delete_marker_candidates
|
||||
.iter()
|
||||
.map(|(_, votes)| *votes)
|
||||
.max()
|
||||
.unwrap_or_default();
|
||||
self.conflicting_metadata |= self.delete_marker_candidates.len() > 1;
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -477,10 +494,15 @@ impl MetadataQuorumAccumulator {
|
||||
if self.default_parity_count == 0 {
|
||||
return Some(self.total_disks);
|
||||
}
|
||||
if candidate.deleted || candidate.size == 0 || candidate.erasure.parity_blocks >= self.total_disks {
|
||||
if candidate.is_canonical_delete_marker() || candidate.size == 0 || candidate.erasure.parity_blocks >= self.total_disks {
|
||||
return None;
|
||||
}
|
||||
Some(candidate.write_quorum(self.default_write_quorum()))
|
||||
let data_blocks = candidate.erasure.data_blocks;
|
||||
Some(if data_blocks == candidate.erasure.parity_blocks {
|
||||
data_blocks.saturating_add(1)
|
||||
} else {
|
||||
data_blocks
|
||||
})
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn default_write_quorum(&self) -> usize {
|
||||
@@ -518,13 +540,22 @@ pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo
|
||||
&& left.is_latest == right.is_latest
|
||||
&& left.deleted == right.deleted
|
||||
&& left.mark_deleted == right.mark_deleted
|
||||
&& left.transition_status == right.transition_status
|
||||
&& left.transitioned_objname == right.transitioned_objname
|
||||
&& left.transition_tier == right.transition_tier
|
||||
&& left.transition_version_id == right.transition_version_id
|
||||
&& left.expire_restored == right.expire_restored
|
||||
&& left.size == right.size
|
||||
&& left.mod_time == right.mod_time
|
||||
&& left.mode == right.mode
|
||||
&& left.written_by_version == right.written_by_version
|
||||
&& left.metadata == right.metadata
|
||||
&& left.replication_state_internal == right.replication_state_internal
|
||||
&& left.parts == right.parts
|
||||
&& left.checksum == right.checksum
|
||||
&& left.versioned == right.versioned
|
||||
&& left.num_versions == right.num_versions
|
||||
&& left.successor_mod_time == right.successor_mod_time
|
||||
&& left.data_dir == right.data_dir
|
||||
&& left.erasure.algorithm == right.erasure.algorithm
|
||||
&& left.erasure.data_blocks == right.erasure.data_blocks
|
||||
@@ -2257,13 +2288,24 @@ impl SetDisks {
|
||||
)));
|
||||
}
|
||||
|
||||
FileMeta {
|
||||
let file_info_versions = 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}")))
|
||||
.map_err(|err| Error::other(format!("exact object versions decode failed for {bucket}/{object}: {err}")))?;
|
||||
|
||||
for file_info in file_info_versions
|
||||
.versions
|
||||
.iter()
|
||||
.chain(file_info_versions.free_versions.iter())
|
||||
{
|
||||
file_info
|
||||
.validate_for_metadata_read()
|
||||
.map_err(|err| Error::other(format!("exact object versions validation failed for {bucket}/{object}: {err}")))?;
|
||||
}
|
||||
|
||||
Ok(Some(file_info_versions))
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) async fn read_all_raw_file_info(
|
||||
@@ -2375,7 +2417,7 @@ impl SetDisks {
|
||||
// `into_fileinfo` with an empty version_id selects the first non-free version
|
||||
// (see FileMeta::into_fileinfo); replicate that selection from the header here.
|
||||
let vid = match meta.into_fileinfo(bucket, object, "", true, incl_free_vers, true) {
|
||||
Ok(finfo) if finfo.is_valid() => finfo.version_id.unwrap_or(Uuid::nil()),
|
||||
Ok(finfo) if file_info_is_valid_for_metadata(&finfo) => finfo.version_id.unwrap_or(Uuid::nil()),
|
||||
_ => match meta
|
||||
.versions
|
||||
.iter()
|
||||
@@ -2398,7 +2440,10 @@ impl SetDisks {
|
||||
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,
|
||||
Ok(res) => match res.validate_for_metadata_read() {
|
||||
Ok(_) => meta_file_infos[idx] = res,
|
||||
Err(err) => errs[idx] = Some(err.into()),
|
||||
},
|
||||
Err(err) => errs[idx] = Some(err.into()),
|
||||
}
|
||||
}
|
||||
@@ -2608,6 +2653,21 @@ impl SetDisks {
|
||||
Vec<Option<DiskStore>>,
|
||||
Option<OldCurrentSize>,
|
||||
)> {
|
||||
if let Some(file_info) = disks
|
||||
.iter()
|
||||
.zip(file_infos.iter())
|
||||
.find_map(|(disk, file_info)| disk.as_ref().map(|_| file_info))
|
||||
{
|
||||
// Newly encoded metadata does not acquire its per-disk shard index
|
||||
// until the fanout below. Validate the shared metadata shape once,
|
||||
// using an online slot because shuffled offline slots contain the
|
||||
// default placeholder, then validate each assigned geometry in its task.
|
||||
if file_info.is_canonical_delete_marker() {
|
||||
file_info.validate_for_metadata_read()?;
|
||||
} else {
|
||||
file_info.validate_for_erasure_write()?;
|
||||
}
|
||||
}
|
||||
let mut futures = Vec::with_capacity(disks.len());
|
||||
|
||||
let mut errs = Vec::with_capacity(disks.len());
|
||||
@@ -2631,11 +2691,16 @@ impl SetDisks {
|
||||
#[allow(clippy::let_unit_value)]
|
||||
let _fanout_task_guard = Self::rename_fanout_task_guard(&dst_object);
|
||||
|
||||
let Some(disk) = disk else {
|
||||
return Err(DiskError::DiskNotFound);
|
||||
};
|
||||
|
||||
let is_delete_marker = file_info.is_canonical_delete_marker();
|
||||
if file_info.erasure.index == 0 {
|
||||
file_info.erasure.index = i + 1;
|
||||
}
|
||||
|
||||
if !file_info.is_valid() {
|
||||
if !is_delete_marker && !file_info.has_valid_erasure_geometry() {
|
||||
return Err(DiskError::FileCorrupt);
|
||||
}
|
||||
|
||||
@@ -2643,12 +2708,8 @@ impl SetDisks {
|
||||
// A no-op immediately-ready future in production.
|
||||
Self::rename_fanout_barrier(&dst_object, i, rename_fanout_barrier_phase::RENAME).await;
|
||||
|
||||
if let Some(disk) = disk {
|
||||
disk.rename_data(&src_bucket, &src_object, file_info, &dst_bucket, &dst_object)
|
||||
.await
|
||||
} else {
|
||||
Err(DiskError::DiskNotFound)
|
||||
}
|
||||
disk.rename_data(&src_bucket, &src_object, file_info, &dst_bucket, &dst_object)
|
||||
.await
|
||||
}));
|
||||
}
|
||||
|
||||
@@ -3391,7 +3452,7 @@ impl SetDisks {
|
||||
// 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()
|
||||
if file_info_is_valid_for_metadata(&m)
|
||||
&& let Some(mod_time) = m.mod_time
|
||||
{
|
||||
let grace = dangling_delete_grace();
|
||||
@@ -3412,7 +3473,7 @@ impl SetDisks {
|
||||
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() {
|
||||
if file_info_is_valid_for_metadata(&m) {
|
||||
tags.insert("sz".to_string(), m.size.to_string());
|
||||
tags.insert(
|
||||
"mt".to_string(),
|
||||
@@ -3502,7 +3563,7 @@ impl SetDisks {
|
||||
}
|
||||
}
|
||||
|
||||
let write_quorum = if m.is_valid() {
|
||||
let write_quorum = if file_info_is_valid_for_metadata(&m) {
|
||||
m.write_quorum(self.default_write_quorum())
|
||||
} else {
|
||||
self.default_write_quorum()
|
||||
@@ -4375,6 +4436,17 @@ mod tests {
|
||||
fi
|
||||
}
|
||||
|
||||
fn metadata_test_delete_marker(object: &str, version_id: Uuid, mod_time: OffsetDateTime) -> FileInfo {
|
||||
FileInfo {
|
||||
volume: "bucket".to_string(),
|
||||
name: object.to_string(),
|
||||
version_id: Some(version_id),
|
||||
deleted: true,
|
||||
mod_time: Some(mod_time),
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
fn read_part_test_part(number: usize, etag: &str) -> ObjectPartInfo {
|
||||
ObjectPartInfo {
|
||||
number,
|
||||
@@ -4434,6 +4506,21 @@ mod tests {
|
||||
.await
|
||||
}
|
||||
|
||||
async fn write_raw_file_meta_unchecked(disk: &DiskStore, bucket: &str, object: &str, metadata: FileMeta) {
|
||||
let encoded = metadata.marshal_msg().expect("raw regression metadata should serialize");
|
||||
disk.write_all(bucket, &format!("{object}/{STORAGE_FORMAT_FILE}"), Bytes::from(encoded))
|
||||
.await
|
||||
.expect("raw regression metadata should be installed");
|
||||
}
|
||||
|
||||
async fn write_raw_file_info_unchecked(disk: &DiskStore, bucket: &str, object: &str, file_info: FileInfo) {
|
||||
let mut metadata = FileMeta::new();
|
||||
metadata
|
||||
.add_version(file_info)
|
||||
.expect("raw regression metadata should encode");
|
||||
write_raw_file_meta_unchecked(disk, bucket, object, metadata).await;
|
||||
}
|
||||
|
||||
fn failed_read_repair_submitter(_request: rustfs_common::heal_channel::HealChannelRequest) -> ReadRepairAdmissionFuture {
|
||||
Box::pin(async { ReadRepairAdmissionOutcome::Failed("injected submit failure".to_string()) })
|
||||
}
|
||||
@@ -4591,6 +4678,153 @@ mod tests {
|
||||
(0..count).map(|_| metadata_test_fileinfo(object)).collect()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rename_data_skips_offline_placeholder_when_validating_new_metadata() {
|
||||
let (_dirs, mut online_disks) = call_counter_local_disks("rename-validation-bucket", 1).await;
|
||||
let online_disk = online_disks.pop().expect("one test disk should be present");
|
||||
let mut file_info = metadata_test_fileinfo("rename-unassigned-index");
|
||||
file_info.erasure.index = 0;
|
||||
|
||||
let err = SetDisks::rename_data(
|
||||
&[None, online_disk],
|
||||
RUSTFS_META_TMP_BUCKET,
|
||||
"source",
|
||||
&[FileInfo::default(), file_info],
|
||||
"bucket",
|
||||
"object",
|
||||
1,
|
||||
)
|
||||
.await
|
||||
.expect_err("the missing staged source must fail after metadata validation");
|
||||
|
||||
assert_ne!(err, DiskError::FileCorrupt);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rename_data_accepts_canonical_delete_marker() {
|
||||
let bucket = "rename-delete-marker-bucket";
|
||||
let object = "object";
|
||||
let (_dirs, mut online_disks) = call_counter_local_disks(bucket, 1).await;
|
||||
let online_disk = online_disks.pop().expect("one test disk slot should be present");
|
||||
let disk = online_disk.as_ref().expect("test disk should be online");
|
||||
match disk.make_volume(RUSTFS_META_TMP_BUCKET).await {
|
||||
Ok(()) | Err(DiskError::VolumeExists) => {}
|
||||
Err(err) => panic!("temporary metadata volume should be available: {err:?}"),
|
||||
}
|
||||
let version_id = Uuid::new_v4();
|
||||
let mut marker = metadata_test_delete_marker(object, version_id, OffsetDateTime::now_utc());
|
||||
marker
|
||||
.metadata
|
||||
.insert("x-rustfs-internal-purgestatus".to_string(), "pending".to_string());
|
||||
|
||||
SetDisks::rename_data(
|
||||
&[None, online_disk.clone()],
|
||||
RUSTFS_META_TMP_BUCKET,
|
||||
"source",
|
||||
&[FileInfo::default(), marker],
|
||||
bucket,
|
||||
object,
|
||||
1,
|
||||
)
|
||||
.await
|
||||
.expect("canonical delete marker should commit without erasure payload");
|
||||
|
||||
let stored = disk
|
||||
.read_version("", bucket, object, &version_id.to_string(), &ReadOptions::default())
|
||||
.await
|
||||
.expect("committed delete marker should be readable");
|
||||
|
||||
assert!(stored.deleted);
|
||||
assert_eq!(stored.version_id, Some(version_id));
|
||||
assert_eq!(stored.erasure.index, 0);
|
||||
assert_eq!(stored.metadata.get("x-rustfs-internal-purgestatus").map(String::as_str), Some("pending"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rename_data_preserves_null_delete_marker_type() {
|
||||
let bucket = "rename-null-marker-bucket";
|
||||
let object = "object";
|
||||
let (_dirs, mut online_disks) = call_counter_local_disks(bucket, 1).await;
|
||||
let online_disk = online_disks.pop().expect("one test disk slot should be present");
|
||||
let disk = online_disk.as_ref().expect("test disk should be online");
|
||||
match disk.make_volume(RUSTFS_META_TMP_BUCKET).await {
|
||||
Ok(()) | Err(DiskError::VolumeExists) => {}
|
||||
Err(err) => panic!("temporary metadata volume should be available: {err:?}"),
|
||||
}
|
||||
let marker = metadata_test_delete_marker(object, Uuid::new_v4(), OffsetDateTime::now_utc());
|
||||
let marker = FileInfo {
|
||||
version_id: None,
|
||||
..marker
|
||||
};
|
||||
|
||||
SetDisks::rename_data(
|
||||
std::slice::from_ref(&online_disk),
|
||||
RUSTFS_META_TMP_BUCKET,
|
||||
"source",
|
||||
&[marker],
|
||||
bucket,
|
||||
object,
|
||||
1,
|
||||
)
|
||||
.await
|
||||
.expect("null delete marker should commit");
|
||||
|
||||
let stored = disk
|
||||
.read_version("", bucket, object, "", &ReadOptions::default())
|
||||
.await
|
||||
.expect("null delete marker should remain readable");
|
||||
assert!(stored.deleted);
|
||||
assert_eq!(stored.version_id, None);
|
||||
assert!(stored.is_canonical_delete_marker());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rename_delete_marker_quorum_failure_restores_existing_metadata() {
|
||||
let bucket = "rename-marker-quorum-bucket";
|
||||
let object = "object";
|
||||
let (_dirs, mut online_disks) = call_counter_local_disks(bucket, 1).await;
|
||||
let online_disk = online_disks.pop().expect("one test disk slot should be present");
|
||||
let disk = online_disk.as_ref().expect("test disk should be online");
|
||||
match disk.make_volume(RUSTFS_META_TMP_BUCKET).await {
|
||||
Ok(()) | Err(DiskError::VolumeExists) => {}
|
||||
Err(err) => panic!("temporary metadata volume should be available: {err:?}"),
|
||||
}
|
||||
let old_version_id = Uuid::new_v4();
|
||||
let mut old = metadata_test_fileinfo(object);
|
||||
old.version_id = Some(old_version_id);
|
||||
old.mod_time = Some(OffsetDateTime::now_utc());
|
||||
disk.write_metadata(bucket, bucket, object, old)
|
||||
.await
|
||||
.expect("old metadata should be written");
|
||||
let marker_version_id = Uuid::new_v4();
|
||||
let marker = metadata_test_delete_marker(object, marker_version_id, OffsetDateTime::now_utc());
|
||||
|
||||
let err = SetDisks::rename_data(
|
||||
&[online_disk.clone(), None],
|
||||
RUSTFS_META_TMP_BUCKET,
|
||||
"source",
|
||||
&[marker, FileInfo::default()],
|
||||
bucket,
|
||||
object,
|
||||
2,
|
||||
)
|
||||
.await
|
||||
.expect_err("quorum-minus-one marker commit should fail");
|
||||
|
||||
assert_eq!(err, DiskError::ErasureWriteQuorum);
|
||||
let restored = disk
|
||||
.read_version("", bucket, object, &old_version_id.to_string(), &ReadOptions::default())
|
||||
.await
|
||||
.expect("old metadata should remain after rollback");
|
||||
assert!(!restored.deleted);
|
||||
assert_eq!(restored.version_id, Some(old_version_id));
|
||||
assert!(matches!(
|
||||
disk.read_version("", bucket, object, &marker_version_id.to_string(), &ReadOptions::default())
|
||||
.await,
|
||||
Err(DiskError::FileVersionNotFound)
|
||||
));
|
||||
}
|
||||
|
||||
/// Demo / regression guard for the backlog#1325 rename fan-out pause barrier
|
||||
/// and background-task introspection. Serves the barrier-style acceptance of
|
||||
/// #1312 ("assert no background disk write remains after release").
|
||||
@@ -4793,6 +5027,47 @@ mod tests {
|
||||
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_ERROR);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metadata_quorum_accumulator_early_stops_on_one_delete_marker_majority() {
|
||||
let marker = metadata_test_delete_marker("object", Uuid::new_v4(), OffsetDateTime::now_utc());
|
||||
let mut accumulator = MetadataQuorumAccumulator::new(6, 3, true);
|
||||
|
||||
for _ in 0..4 {
|
||||
accumulator.observe_file_info(&marker);
|
||||
}
|
||||
|
||||
assert_eq!(accumulator.default_write_quorum(), 4);
|
||||
assert_eq!(accumulator.delete_marker_votes, 4);
|
||||
assert_eq!(
|
||||
accumulator.early_stop_decision(),
|
||||
Some(MetadataEarlyStopDecision {
|
||||
reason: GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER,
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metadata_quorum_accumulator_does_not_combine_distinct_delete_markers() {
|
||||
let now = OffsetDateTime::now_utc();
|
||||
let first = metadata_test_delete_marker("object", Uuid::new_v4(), now);
|
||||
let second = metadata_test_delete_marker("object", Uuid::new_v4(), now + time::Duration::seconds(1));
|
||||
let third = metadata_test_delete_marker("object", Uuid::new_v4(), now + time::Duration::seconds(2));
|
||||
let mut accumulator = MetadataQuorumAccumulator::new(4, 2, true);
|
||||
|
||||
accumulator.observe_file_info(&first);
|
||||
accumulator.observe_file_info(&second);
|
||||
accumulator.observe_file_info(&third);
|
||||
|
||||
assert_eq!(accumulator.delete_marker_votes, 1);
|
||||
assert_eq!(accumulator.delete_marker_candidates.len(), 3);
|
||||
assert_eq!(accumulator.early_stop_decision(), None);
|
||||
|
||||
accumulator.observe_file_info(&second);
|
||||
accumulator.observe_file_info(&second);
|
||||
assert_eq!(accumulator.delete_marker_votes, 3);
|
||||
assert!(accumulator.early_stop_decision().is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metadata_quorum_accumulator_candidate_latest_quorum_handles_zero_parity_and_invalid_candidates() {
|
||||
let accumulator = MetadataQuorumAccumulator::new(4, 0, true);
|
||||
@@ -4803,7 +5078,10 @@ mod tests {
|
||||
let accumulator = MetadataQuorumAccumulator::new(4, 2, true);
|
||||
let mut deleted = candidate.clone();
|
||||
deleted.deleted = true;
|
||||
assert_eq!(accumulator.candidate_latest_quorum(&deleted), None);
|
||||
assert_eq!(accumulator.candidate_latest_quorum(&deleted), Some(3));
|
||||
|
||||
let marker = metadata_test_delete_marker("object", Uuid::new_v4(), OffsetDateTime::now_utc());
|
||||
assert_eq!(accumulator.candidate_latest_quorum(&marker), None);
|
||||
|
||||
let mut empty = candidate.clone();
|
||||
empty.size = 0;
|
||||
@@ -5052,6 +5330,59 @@ mod tests {
|
||||
assert_eq!(versions.versions[0].name, object);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn load_file_info_versions_exact_rejects_transitioned_duplicate_parts() {
|
||||
let bucket = "exact-versions-bucket";
|
||||
let object = "poisoned-transitioned-object";
|
||||
let (_dir, disk) = read_multiple_test_disk(bucket, &[]).await;
|
||||
let mut file_info = metadata_test_fileinfo(object);
|
||||
file_info.version_id = Some(Uuid::new_v4());
|
||||
file_info.mod_time = Some(OffsetDateTime::now_utc());
|
||||
file_info.transition_status = TRANSITION_COMPLETE.to_string();
|
||||
file_info.transitioned_objname = "remote/object".to_string();
|
||||
file_info.transition_tier = "WARM".to_string();
|
||||
file_info.parts.push(file_info.parts[0].clone());
|
||||
write_raw_file_info_unchecked(&disk, bucket, object, file_info).await;
|
||||
let set = io_primitives_test_set(vec![Some(disk)], 0).await;
|
||||
|
||||
let err = set
|
||||
.load_file_info_versions_exact(bucket, object)
|
||||
.await
|
||||
.expect_err("exact loader must reject metadata that would poison decommission");
|
||||
|
||||
assert!(err.to_string().contains("validation failed"), "unexpected error: {err}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn load_file_info_versions_exact_rejects_default_like_delete_marker() {
|
||||
let bucket = "exact-versions-bucket";
|
||||
let object = "forged-delete-marker";
|
||||
let (_dir, disk) = read_multiple_test_disk(bucket, &[]).await;
|
||||
let forged_version = rustfs_filemeta::FileMetaVersion {
|
||||
version_type: rustfs_filemeta::VersionType::Delete,
|
||||
delete_marker: Some(rustfs_filemeta::MetaDeleteMarker {
|
||||
version_id: Some(Uuid::new_v4()),
|
||||
mod_time: None,
|
||||
..Default::default()
|
||||
}),
|
||||
write_version: 1,
|
||||
..Default::default()
|
||||
};
|
||||
let mut forged_meta = FileMeta::new();
|
||||
forged_meta
|
||||
.versions
|
||||
.push(FileMetaShallowVersion::try_from(forged_version).expect("forged marker body should encode"));
|
||||
write_raw_file_meta_unchecked(&disk, bucket, object, forged_meta).await;
|
||||
let set = io_primitives_test_set(vec![Some(disk)], 0).await;
|
||||
|
||||
let err = set
|
||||
.load_file_info_versions_exact(bucket, object)
|
||||
.await
|
||||
.expect_err("default-like delete marker must be rejected at the exact loader boundary");
|
||||
|
||||
assert!(err.to_string().contains("exact object versions decode failed"), "unexpected error: {err}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn commit_rename_data_dir_reclaims_old_data_dir_and_reports_receipt() {
|
||||
let bucket = "commit-rename-bucket";
|
||||
|
||||
@@ -245,12 +245,12 @@ impl SetDisks {
|
||||
continue;
|
||||
}
|
||||
|
||||
if !metadata.is_valid() {
|
||||
if !file_info_is_valid_for_metadata(metadata) {
|
||||
parities[index] = -1;
|
||||
continue;
|
||||
}
|
||||
|
||||
if metadata.deleted || metadata.size == 0 {
|
||||
if metadata.is_canonical_delete_marker() || metadata.size == 0 {
|
||||
parities[index] = half;
|
||||
} else if metadata.transition_status == TRANSITION_COMPLETE {
|
||||
let majority_metadata_parity = total_shards_i32 - (half + 1);
|
||||
@@ -327,7 +327,7 @@ impl SetDisks {
|
||||
for (i, etag_item) in etags.iter().enumerate() {
|
||||
if let Some(etag_item) = etag_item
|
||||
&& etag_item == &etag
|
||||
&& parts_metadata[i].is_valid()
|
||||
&& file_info_is_valid_for_metadata(&parts_metadata[i])
|
||||
{
|
||||
new_disk[i].clone_from(&disks[i]);
|
||||
}
|
||||
@@ -340,7 +340,7 @@ impl SetDisks {
|
||||
let mut new_disk = vec![None; disks.len()];
|
||||
|
||||
for (i, &t) in mod_times.iter().enumerate() {
|
||||
if parts_metadata[i].is_valid() && mod_time == t {
|
||||
if file_info_is_valid_for_metadata(&parts_metadata[i]) && mod_time == t {
|
||||
new_disk[i].clone_from(&disks[i]);
|
||||
}
|
||||
}
|
||||
@@ -357,7 +357,7 @@ impl SetDisks {
|
||||
continue;
|
||||
}
|
||||
|
||||
if meta.is_valid() {
|
||||
if file_info_is_valid_for_metadata(meta) {
|
||||
usable_metadata += 1;
|
||||
}
|
||||
}
|
||||
@@ -388,7 +388,7 @@ impl SetDisks {
|
||||
|
||||
let mut identity_counts = HashMap::with_capacity(usable_metadata);
|
||||
for (meta, err) in parts_metadata.iter().zip(errs.iter()) {
|
||||
if err.is_some() || !meta.is_valid() {
|
||||
if err.is_some() || !file_info_is_valid_for_metadata(meta) {
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -613,7 +613,7 @@ impl SetDisks {
|
||||
}
|
||||
}
|
||||
|
||||
if !meta.deleted && meta.size != 0 {
|
||||
if !meta.is_canonical_delete_marker() && meta.size != 0 {
|
||||
hasher.update(meta.erasure.data_blocks.to_le_bytes());
|
||||
hasher.update(meta.erasure.parity_blocks.to_le_bytes());
|
||||
hasher.update(meta.erasure.distribution.len().to_le_bytes());
|
||||
@@ -626,7 +626,7 @@ impl SetDisks {
|
||||
fn latest_fileinfo_identity_groups(parts_metadata: &[FileInfo], errs: &[Option<DiskError>]) -> Vec<FileInfoIdentityGroup> {
|
||||
let mut groups: Vec<FileInfoIdentityGroup> = Vec::with_capacity(parts_metadata.len());
|
||||
for (meta, err) in parts_metadata.iter().zip(errs.iter()) {
|
||||
if err.is_some() || !meta.is_valid() {
|
||||
if err.is_some() || !file_info_is_valid_for_metadata(meta) {
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -658,7 +658,7 @@ impl SetDisks {
|
||||
let mut count = 0;
|
||||
|
||||
for (i, ((meta, err), disk)) in parts_metadata.iter().zip(errs.iter()).zip(disks.iter()).enumerate() {
|
||||
if err.is_some() || !meta.is_valid() || Self::file_info_quorum_hash(meta) != hash {
|
||||
if err.is_some() || !file_info_is_valid_for_metadata(meta) || Self::file_info_quorum_hash(meta) != hash {
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -731,7 +731,7 @@ impl SetDisks {
|
||||
let mut meta_hashes = vec![None; metas.len()];
|
||||
|
||||
for (i, meta) in metas.iter().enumerate() {
|
||||
if !meta.is_valid() {
|
||||
if !file_info_is_valid_for_metadata(meta) {
|
||||
debug!(
|
||||
index = i,
|
||||
valid = false,
|
||||
@@ -807,7 +807,7 @@ impl SetDisks {
|
||||
if let Some(hash) = op_hash
|
||||
&& let Some(max_hash) = max_val
|
||||
&& *hash == max_hash
|
||||
&& metas[i].is_valid()
|
||||
&& file_info_is_valid_for_metadata(&metas[i])
|
||||
{
|
||||
if !found {
|
||||
found_fi = Some(metas[i].clone());
|
||||
@@ -861,7 +861,7 @@ impl SetDisks {
|
||||
|
||||
let mut inconsistent = 0;
|
||||
for (k, v) in parts_metadata.iter().enumerate() {
|
||||
if disks[k].is_none() || !v.is_valid() || distribution[k] != v.erasure.index {
|
||||
if disks[k].is_none() || !v.has_valid_erasure_geometry() || distribution[k] != v.erasure.index {
|
||||
inconsistent += 1;
|
||||
}
|
||||
}
|
||||
@@ -877,9 +877,9 @@ impl SetDisks {
|
||||
continue;
|
||||
}
|
||||
let eligible = if use_by_index {
|
||||
parts_metadata[k].is_valid() && distribution[k] == parts_metadata[k].erasure.index
|
||||
parts_metadata[k].has_valid_erasure_geometry() && distribution[k] == parts_metadata[k].erasure.index
|
||||
} else {
|
||||
init || parts_metadata[k].is_valid()
|
||||
init || parts_metadata[k].has_valid_erasure_geometry()
|
||||
};
|
||||
if !eligible {
|
||||
continue;
|
||||
@@ -917,7 +917,7 @@ impl SetDisks {
|
||||
continue;
|
||||
}
|
||||
|
||||
if !v.is_valid() {
|
||||
if !v.has_valid_erasure_geometry() {
|
||||
inconsistent += 1;
|
||||
continue;
|
||||
}
|
||||
@@ -963,7 +963,7 @@ impl SetDisks {
|
||||
continue;
|
||||
}
|
||||
|
||||
if !init && !parts_metadata[k].is_valid() {
|
||||
if !init && !parts_metadata[k].has_valid_erasure_geometry() {
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -1156,6 +1156,43 @@ mod tests {
|
||||
assert_ne!(SetDisks::file_info_quorum_hash(&left), SetDisks::file_info_quorum_hash(&right));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn purge_pending_quorum_hash_keeps_erasure_layouts_separate() {
|
||||
let mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
|
||||
let version_id = Uuid::new_v4();
|
||||
let data_dir = Uuid::new_v4();
|
||||
let mut honest = FileInfo::new("bucket/object", 5, 1);
|
||||
honest.name = "bucket/object".to_string();
|
||||
honest.version_id = Some(version_id);
|
||||
honest.data_dir = Some(data_dir);
|
||||
honest.mod_time = Some(mod_time);
|
||||
honest.size = 1;
|
||||
honest.deleted = true;
|
||||
honest.add_object_part(1, "part-etag".to_string(), 1, Some(mod_time), 1, None, None);
|
||||
|
||||
let mut parts_metadata = (1..=6)
|
||||
.map(|index| {
|
||||
let mut metadata = honest.clone();
|
||||
metadata.erasure.index = index;
|
||||
metadata
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
let mut tampered_layout = FileInfo::new("bucket/object", 3, 3).erasure;
|
||||
tampered_layout.index = 1;
|
||||
parts_metadata[0].erasure = tampered_layout;
|
||||
let errs = vec![None; 6];
|
||||
|
||||
assert_eq!(
|
||||
SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 3)
|
||||
.expect("five honest EC:1 payload copies should determine object quorum"),
|
||||
(5, 5)
|
||||
);
|
||||
let selected = SetDisks::find_file_info_in_quorum(&parts_metadata, &Some(mod_time), &None, 5)
|
||||
.expect("the five matching EC:1 payload copies should determine metadata identity");
|
||||
assert_eq!(selected.erasure.data_blocks, 5);
|
||||
assert_eq!(selected.erasure.parity_blocks, 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn quorum_helpers_reject_zero_quorum_and_shuffle_check_parts_by_distribution() {
|
||||
let err = SetDisks::find_file_info_in_quorum(&[], &None, &None, 0).expect_err("zero quorum cannot select metadata");
|
||||
|
||||
@@ -224,6 +224,14 @@ pub(crate) const RUSTFS_MULTIPART_BUCKET_KEY: &str = "x-rustfs-internal-multipar
|
||||
pub(crate) const RUSTFS_MULTIPART_OBJECT_KEY: &str = "x-rustfs-internal-multipart-object";
|
||||
const ENV_ISSUE3031_DIAG_ENABLE: &str = "RUSTFS_ISSUE3031_DIAG_ENABLE";
|
||||
|
||||
/// Validate disk metadata at a boundary that may legitimately return a delete
|
||||
/// marker. Disk/RPC decode boundaries perform the full collection validation
|
||||
/// once; repeated quorum passes use the cheap erasure-geometry predicate for
|
||||
/// payload entries and the canonical marker predicate for pure delete markers.
|
||||
pub(in crate::set_disk) fn file_info_is_valid_for_metadata(file_info: &FileInfo) -> bool {
|
||||
file_info.has_valid_metadata_shape()
|
||||
}
|
||||
|
||||
struct ObjectLockDiagGuard {
|
||||
guard: NamespaceLockGuard,
|
||||
enabled: bool,
|
||||
@@ -1971,25 +1979,209 @@ fn issue3031_diag_enabled() -> bool {
|
||||
rustfs_utils::get_env_bool(ENV_ISSUE3031_DIAG_ENABLE, false)
|
||||
}
|
||||
|
||||
fn build_tiered_decommission_file_info(
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
fi: &FileInfo,
|
||||
disk_count: usize,
|
||||
default_parity_count: usize,
|
||||
storage_class: Option<&str>,
|
||||
) -> (FileInfo, usize) {
|
||||
let parity_drives = runtime_sources::storage_class_parity(storage_class).unwrap_or(default_parity_count);
|
||||
let data_drives = disk_count - parity_drives;
|
||||
let mut write_quorum = data_drives;
|
||||
if data_drives == parity_drives {
|
||||
write_quorum += 1;
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub(super) struct WriteLayout {
|
||||
data_drives: usize,
|
||||
parity_drives: usize,
|
||||
write_quorum: usize,
|
||||
}
|
||||
|
||||
impl WriteLayout {
|
||||
fn from_parity(drive_count: usize, parity_drives: usize) -> Result<Self> {
|
||||
let max_parity = drive_count / 2;
|
||||
if parity_drives > max_parity {
|
||||
return Err(Error::other(format!(
|
||||
"write parity {parity_drives} exceeds the maximum {max_parity} for {drive_count} drives"
|
||||
)));
|
||||
}
|
||||
|
||||
let data_drives = drive_count
|
||||
.checked_sub(parity_drives)
|
||||
.filter(|&data_drives| data_drives > 0 && parity_drives <= data_drives)
|
||||
.ok_or_else(|| Error::other(format!("invalid write layout with {drive_count} drives and parity {parity_drives}")))?;
|
||||
let write_quorum = data_drives
|
||||
.checked_add(usize::from(data_drives == parity_drives))
|
||||
.filter(|&write_quorum| write_quorum <= drive_count)
|
||||
.ok_or_else(|| Error::other(format!("invalid write quorum for {drive_count} drives and parity {parity_drives}")))?;
|
||||
|
||||
Ok(Self {
|
||||
data_drives,
|
||||
parity_drives,
|
||||
write_quorum,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn resolve_write_layout(
|
||||
config: &storageclass::Config,
|
||||
pool_index: usize,
|
||||
drive_count: usize,
|
||||
fallback_parity: usize,
|
||||
storage_class: Option<&str>,
|
||||
max_parity: bool,
|
||||
) -> Result<WriteLayout> {
|
||||
let configured_parity = if config.is_initialized() {
|
||||
config
|
||||
.parity_for_pool(storage_class.unwrap_or_default(), pool_index, drive_count)
|
||||
.ok_or_else(|| {
|
||||
Error::other(format!("storage class layout does not match pool {pool_index} with {drive_count} drives"))
|
||||
})?
|
||||
} else {
|
||||
fallback_parity
|
||||
};
|
||||
let parity_drives = if max_parity { drive_count / 2 } else { configured_parity };
|
||||
|
||||
WriteLayout::from_parity(drive_count, parity_drives)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod write_layout_tests {
|
||||
use super::{WriteLayout, resolve_write_layout};
|
||||
use crate::config::storageclass::{
|
||||
CLASS_RRS, CLASS_STANDARD, INLINE_BLOCK_ENV, OPTIMIZE_ENV, RRS, RRS_ENV, STANDARD_ENV, lookup_config_for_pools,
|
||||
lookup_config_for_pools_without_env,
|
||||
};
|
||||
use arc_swap::ArcSwap;
|
||||
use rustfs_config::server_config::KVS;
|
||||
use std::sync::Arc;
|
||||
|
||||
#[test]
|
||||
fn automatic_standard_layout_is_resolved_per_pool() {
|
||||
let config = lookup_config_for_pools_without_env(&KVS::new(), &[4, 2])
|
||||
.expect("automatic storage class should resolve for both pools");
|
||||
|
||||
assert_eq!(
|
||||
resolve_write_layout(&config, 0, 4, 2, None, false).expect("first pool should resolve"),
|
||||
WriteLayout {
|
||||
data_drives: 2,
|
||||
parity_drives: 2,
|
||||
write_quorum: 3,
|
||||
}
|
||||
);
|
||||
assert_eq!(
|
||||
resolve_write_layout(&config, 1, 2, 1, None, false).expect("second pool should resolve"),
|
||||
WriteLayout {
|
||||
data_drives: 1,
|
||||
parity_drives: 1,
|
||||
write_quorum: 2,
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reduced_redundancy_layout_allows_single_disk_zero_parity() {
|
||||
let config =
|
||||
lookup_config_for_pools_without_env(&KVS::new(), &[4, 1]).expect("reduced redundancy should resolve for both pools");
|
||||
|
||||
assert_eq!(
|
||||
resolve_write_layout(&config, 0, 4, 2, Some(RRS), false).expect("four-drive RRS pool should resolve"),
|
||||
WriteLayout {
|
||||
data_drives: 3,
|
||||
parity_drives: 1,
|
||||
write_quorum: 3,
|
||||
}
|
||||
);
|
||||
assert_eq!(
|
||||
resolve_write_layout(&config, 1, 1, 0, Some(RRS), false).expect("single-drive RRS pool should resolve"),
|
||||
WriteLayout {
|
||||
data_drives: 1,
|
||||
parity_drives: 0,
|
||||
write_quorum: 1,
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_layout_rejects_unknown_topology_and_invalid_parity() {
|
||||
let config = lookup_config_for_pools_without_env(&KVS::new(), &[4, 2]).expect("test storage class should resolve");
|
||||
|
||||
assert!(resolve_write_layout(&config, 2, 2, 1, None, false).is_err());
|
||||
assert!(resolve_write_layout(&config, 1, 4, 2, None, false).is_err());
|
||||
assert!(WriteLayout::from_parity(4, 3).is_err());
|
||||
assert!(WriteLayout::from_parity(0, 0).is_err());
|
||||
|
||||
let mut zero_parity_kvs = KVS::new();
|
||||
zero_parity_kvs.insert(CLASS_STANDARD.to_string(), "EC:0".to_string());
|
||||
zero_parity_kvs.insert(CLASS_RRS.to_string(), "EC:0".to_string());
|
||||
let zero_parity = lookup_config_for_pools_without_env(&zero_parity_kvs, &[4]).expect("zero-parity config should resolve");
|
||||
assert_eq!(
|
||||
resolve_write_layout(&zero_parity, 0, 4, 2, None, true).expect("max parity should override configured parity"),
|
||||
WriteLayout {
|
||||
data_drives: 2,
|
||||
parity_drives: 2,
|
||||
write_quorum: 3,
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn only_uninitialized_config_falls_back_to_pool_startup_parity() {
|
||||
let uninitialized = crate::config::storageclass::Config::default();
|
||||
assert_eq!(
|
||||
resolve_write_layout(&uninitialized, 99, 2, 0, None, false)
|
||||
.expect("uninitialized config should preserve the pool's startup fallback"),
|
||||
WriteLayout {
|
||||
data_drives: 2,
|
||||
parity_drives: 0,
|
||||
write_quorum: 2,
|
||||
}
|
||||
);
|
||||
|
||||
let initialized = lookup_config_for_pools_without_env(&KVS::new(), &[4, 2]).expect("initialized config should resolve");
|
||||
assert!(resolve_write_layout(&initialized, 99, 2, 1, None, false).is_err());
|
||||
assert!(resolve_write_layout(&initialized, 1, 4, 2, None, false).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
fn held_snapshot_keeps_parity_and_inline_policy_consistent_across_reload() {
|
||||
let old = temp_env::with_vars(
|
||||
[
|
||||
(STANDARD_ENV, Some("")),
|
||||
(RRS_ENV, Some("")),
|
||||
(OPTIMIZE_ENV, None),
|
||||
(INLINE_BLOCK_ENV, Some("1KiB")),
|
||||
],
|
||||
|| lookup_config_for_pools(&KVS::new(), &[4, 2]),
|
||||
)
|
||||
.expect("old config should resolve");
|
||||
let new = temp_env::with_vars(
|
||||
[
|
||||
(STANDARD_ENV, Some("EC:1")),
|
||||
(RRS_ENV, Some("EC:1")),
|
||||
(OPTIMIZE_ENV, None),
|
||||
(INLINE_BLOCK_ENV, Some("0B")),
|
||||
],
|
||||
|| lookup_config_for_pools(&KVS::new(), &[4, 2]),
|
||||
)
|
||||
.expect("new config should resolve");
|
||||
|
||||
let published = ArcSwap::from_pointee(old);
|
||||
let held = published.load_full();
|
||||
published.store(Arc::new(new));
|
||||
|
||||
let held_layout = resolve_write_layout(&held, 0, 4, 2, None, false).expect("held snapshot should remain valid");
|
||||
assert_eq!(held_layout.parity_drives, 2);
|
||||
assert!(held.should_inline(512, false));
|
||||
|
||||
let current = published.load_full();
|
||||
let current_layout = resolve_write_layout(¤t, 0, 4, 2, None, false).expect("new snapshot should resolve");
|
||||
assert_eq!(current_layout.parity_drives, 1);
|
||||
assert!(!current.should_inline(512, false));
|
||||
}
|
||||
}
|
||||
|
||||
fn build_tiered_decommission_file_info(bucket: &str, object: &str, fi: &FileInfo, layout: WriteLayout) -> FileInfo {
|
||||
let WriteLayout {
|
||||
data_drives,
|
||||
parity_drives,
|
||||
..
|
||||
} = layout;
|
||||
|
||||
let mut updated = fi.clone();
|
||||
updated.erasure = FileInfo::new([bucket, object].join("/").as_str(), data_drives, parity_drives).erasure;
|
||||
|
||||
(updated, write_quorum)
|
||||
updated
|
||||
}
|
||||
|
||||
fn resolve_tiered_decommission_write_quorum_result(
|
||||
@@ -2036,6 +2228,8 @@ pub struct SetDisks {
|
||||
/// writes skip the global registry mutex (backlog#1315). `Arc` so clones of
|
||||
/// a set share one generation marker.
|
||||
capacity_dirty_generation: Arc<AtomicU64>,
|
||||
#[cfg(test)]
|
||||
storage_class_config_override: Arc<std::sync::RwLock<Option<Arc<storageclass::Config>>>>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Eq, PartialEq)]
|
||||
@@ -2168,6 +2362,28 @@ impl DiskHealthEntry {
|
||||
}
|
||||
|
||||
impl SetDisks {
|
||||
fn storage_class_config_snapshot(&self) -> Arc<storageclass::Config> {
|
||||
#[cfg(test)]
|
||||
if let Some(config) = self
|
||||
.storage_class_config_override
|
||||
.read()
|
||||
.expect("test storage class override lock should not be poisoned")
|
||||
.as_ref()
|
||||
{
|
||||
return config.clone();
|
||||
}
|
||||
|
||||
runtime_sources::storage_class_config_snapshot()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn set_test_storage_class_config(&self, config: storageclass::Config) {
|
||||
*self
|
||||
.storage_class_config_override
|
||||
.write()
|
||||
.expect("test storage class override lock should not be poisoned") = Some(Arc::new(config));
|
||||
}
|
||||
|
||||
fn get_object_metadata_cache_hash(&self, bucket: &str, object: &str) -> u64 {
|
||||
let mut hasher = self.get_object_metadata_cache_hash_builder.build_hasher();
|
||||
bucket.hash(&mut hasher);
|
||||
@@ -2372,6 +2588,8 @@ impl SetDisks {
|
||||
ctx,
|
||||
capacity_scope_cache: Arc::new(std::sync::RwLock::new(CapacityScopeCache::default())),
|
||||
capacity_dirty_generation: Arc::new(AtomicU64::new(u64::MAX)),
|
||||
#[cfg(test)]
|
||||
storage_class_config_override: Arc::new(std::sync::RwLock::new(None)),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -2888,7 +3106,7 @@ fn collect_inline_data_shard_fileinfos_by_index<'a>(
|
||||
if block_index == 0 || block_index > data_shards {
|
||||
continue;
|
||||
}
|
||||
if !file_info.is_valid() {
|
||||
if !file_info.has_valid_erasure_geometry() {
|
||||
continue;
|
||||
}
|
||||
if file_info.data.as_ref().is_none_or(|data| data.is_empty()) {
|
||||
@@ -3322,6 +3540,7 @@ impl SetDisks {
|
||||
fi: &FileInfo,
|
||||
opts: &ObjectOptions,
|
||||
) -> Result<()> {
|
||||
let storage_class_config = self.storage_class_config_snapshot();
|
||||
let _lock_guard = if !opts.no_lock {
|
||||
Some(
|
||||
self.new_ns_lock(bucket, object)
|
||||
@@ -3336,8 +3555,16 @@ impl SetDisks {
|
||||
|
||||
let disks = self.disks.read().await.clone();
|
||||
let storage_class = opts.user_defined.get(AMZ_STORAGE_CLASS).map(String::as_str);
|
||||
let (fi, write_quorum) =
|
||||
build_tiered_decommission_file_info(bucket, object, fi, disks.len(), self.default_parity_count, storage_class);
|
||||
let layout = resolve_write_layout(
|
||||
&storage_class_config,
|
||||
self.pool_index,
|
||||
disks.len(),
|
||||
self.default_parity_count,
|
||||
storage_class,
|
||||
opts.max_parity,
|
||||
)?;
|
||||
let fi = build_tiered_decommission_file_info(bucket, object, fi, layout);
|
||||
let write_quorum = layout.write_quorum;
|
||||
let parts_metadata = vec![fi.clone(); disks.len()];
|
||||
let (shuffle_disks, parts_metadata) = Self::shuffle_disks_and_parts_metadata(&disks, &parts_metadata, &fi);
|
||||
|
||||
@@ -3407,13 +3634,13 @@ fn is_object_dangling(
|
||||
let mut valid_meta = FileInfo::default();
|
||||
|
||||
for fi in meta_arr.iter() {
|
||||
if fi.is_valid() {
|
||||
if file_info_is_valid_for_metadata(fi) {
|
||||
valid_meta = fi.clone();
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if !valid_meta.is_valid() {
|
||||
if !file_info_is_valid_for_metadata(&valid_meta) {
|
||||
let data_blocks = meta_arr.len().div_ceil(2);
|
||||
if not_found_parts_errs > data_blocks {
|
||||
return (valid_meta, true);
|
||||
@@ -3426,7 +3653,7 @@ fn is_object_dangling(
|
||||
return (valid_meta, false);
|
||||
}
|
||||
|
||||
if valid_meta.deleted {
|
||||
if valid_meta.is_canonical_delete_marker() {
|
||||
let data_blocks = errs.len().div_ceil(2);
|
||||
return (valid_meta, not_found_meta_errs > data_blocks);
|
||||
}
|
||||
@@ -3539,12 +3766,12 @@ async fn disks_with_all_parts(
|
||||
// Check for inconsistent erasure distribution
|
||||
let mut inconsistent = 0;
|
||||
for (index, meta) in parts_metadata.iter().enumerate() {
|
||||
if !meta.is_valid() {
|
||||
if !file_info_is_valid_for_metadata(meta) {
|
||||
// Since for majority of the cases erasure.Index matches with erasure.Distribution we can
|
||||
// consider the offline disks as consistent.
|
||||
continue;
|
||||
}
|
||||
if !meta.deleted {
|
||||
if !meta.is_canonical_delete_marker() {
|
||||
if meta.erasure.distribution.len() != online_disks.len() {
|
||||
// Erasure distribution seems to have lesser
|
||||
// number of items than number of online disks.
|
||||
@@ -3604,7 +3831,7 @@ async fn disks_with_all_parts(
|
||||
}
|
||||
|
||||
if erasure_distribution_reliable {
|
||||
if !meta.is_valid() {
|
||||
if !file_info_is_valid_for_metadata(meta) {
|
||||
info!(
|
||||
"disks_with_all_partsv2: metadata is not valid, object_name={}, index: {index}",
|
||||
object_name
|
||||
@@ -3615,7 +3842,7 @@ async fn disks_with_all_parts(
|
||||
continue;
|
||||
}
|
||||
|
||||
if !meta.deleted && meta.erasure.distribution.len() != online_disks_len {
|
||||
if !meta.is_canonical_delete_marker() && meta.erasure.distribution.len() != online_disks_len {
|
||||
// Erasure distribution is not the same as onlineDisks
|
||||
// attempt a fix if possible, assuming other entries
|
||||
// might have the right erasure distribution.
|
||||
@@ -3658,7 +3885,7 @@ async fn disks_with_all_parts(
|
||||
};
|
||||
|
||||
let meta = &mut parts_metadata[index];
|
||||
if meta.deleted || meta.is_remote() {
|
||||
if meta.is_canonical_delete_marker() || meta.is_remote() {
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -3788,7 +4015,7 @@ pub fn should_heal_object_on_disk(
|
||||
return (true, true, Some(DiskError::OutdatedXLMeta));
|
||||
}
|
||||
|
||||
if !meta.deleted && !meta.is_remote() {
|
||||
if !meta.is_canonical_delete_marker() && !meta.is_remote() {
|
||||
let err_vec = [CHECK_PART_FILE_NOT_FOUND, CHECK_PART_FILE_CORRUPT];
|
||||
for part_err in parts_errs.iter() {
|
||||
if err_vec.contains(part_err) {
|
||||
@@ -6356,6 +6583,7 @@ mod tests {
|
||||
erasure: ErasureInfo {
|
||||
data_blocks: 4,
|
||||
parity_blocks: 2,
|
||||
block_size: 4,
|
||||
index: 1, // Must be > 0 for is_valid() to return true
|
||||
distribution: vec![1, 2, 3, 4, 5, 6], // Must match data_blocks + parity_blocks
|
||||
..Default::default()
|
||||
@@ -6368,6 +6596,7 @@ mod tests {
|
||||
erasure: ErasureInfo {
|
||||
data_blocks: 6,
|
||||
parity_blocks: 3,
|
||||
block_size: 4,
|
||||
index: 1, // Must be > 0 for is_valid() to return true
|
||||
distribution: vec![1, 2, 3, 4, 5, 6, 7, 8, 9], // Must match data_blocks + parity_blocks
|
||||
..Default::default()
|
||||
@@ -6380,6 +6609,7 @@ mod tests {
|
||||
erasure: ErasureInfo {
|
||||
data_blocks: 2,
|
||||
parity_blocks: 1,
|
||||
block_size: 4,
|
||||
index: 1, // Must be > 0 for is_valid() to return true
|
||||
distribution: vec![1, 2, 3], // Must match data_blocks + parity_blocks
|
||||
..Default::default()
|
||||
@@ -6399,6 +6629,102 @@ mod tests {
|
||||
assert_eq!(parities[2], 1); // half of total shards (3/2 = 1) for zero size file
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn delete_markers_participate_in_four_disk_metadata_quorum_without_erasure_geometry() {
|
||||
let marker = FileInfo {
|
||||
name: "bucket/deleted".to_string(),
|
||||
deleted: true,
|
||||
version_id: Some(Uuid::new_v4()),
|
||||
mod_time: Some(OffsetDateTime::now_utc()),
|
||||
..Default::default()
|
||||
};
|
||||
let parts_metadata = vec![marker; 4];
|
||||
let errs = vec![None; 4];
|
||||
|
||||
assert!(parts_metadata.iter().all(file_info_is_valid_for_metadata));
|
||||
assert!(parts_metadata.iter().all(|metadata| !metadata.is_valid()));
|
||||
assert_eq!(SetDisks::list_object_parities(&parts_metadata, &errs), vec![2; 4]);
|
||||
assert_eq!(
|
||||
SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 2)
|
||||
.expect("four matching delete markers must reach metadata quorum"),
|
||||
(2, 3)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metadata_boundary_does_not_relax_non_delete_or_malformed_delete_metadata() {
|
||||
assert!(!file_info_is_valid_for_metadata(&FileInfo::default()));
|
||||
|
||||
let mut transitioned = FileInfo::new("bucket/transitioned", 2, 2);
|
||||
transitioned.erasure.index = 1;
|
||||
transitioned.transition_status = TRANSITION_COMPLETE.to_string();
|
||||
assert!(file_info_is_valid_for_metadata(&transitioned));
|
||||
|
||||
transitioned.erasure = ErasureInfo::default();
|
||||
assert!(
|
||||
!file_info_is_valid_for_metadata(&transitioned),
|
||||
"transition state must not relax local erasure validation"
|
||||
);
|
||||
|
||||
let mut purge_pending = FileInfo::new("bucket/purge-pending", 2, 2);
|
||||
purge_pending.erasure.index = 1;
|
||||
purge_pending.deleted = true;
|
||||
purge_pending.parts.push(ObjectPartInfo {
|
||||
number: 1,
|
||||
..Default::default()
|
||||
});
|
||||
assert!(
|
||||
file_info_is_valid_for_metadata(&purge_pending),
|
||||
"purge-pending payload metadata must retain its valid erasure vote"
|
||||
);
|
||||
assert!(!purge_pending.is_canonical_delete_marker());
|
||||
|
||||
let mut malformed_marker = FileInfo {
|
||||
deleted: true,
|
||||
..Default::default()
|
||||
};
|
||||
malformed_marker.parts = vec![
|
||||
ObjectPartInfo {
|
||||
number: 1,
|
||||
..Default::default()
|
||||
},
|
||||
ObjectPartInfo {
|
||||
number: 1,
|
||||
..Default::default()
|
||||
},
|
||||
];
|
||||
assert!(!file_info_is_valid_for_metadata(&malformed_marker));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn purge_pending_payload_uses_its_erasure_parity_for_metadata_quorum() {
|
||||
let version_id = Uuid::new_v4();
|
||||
let mod_time = OffsetDateTime::now_utc();
|
||||
let parts_metadata = (1..=6)
|
||||
.map(|disk_index| {
|
||||
let mut purge_pending = FileInfo::new("bucket/purge-pending", 5, 1);
|
||||
purge_pending.name = "bucket/purge-pending".to_string();
|
||||
purge_pending.version_id = Some(version_id);
|
||||
purge_pending.mod_time = Some(mod_time);
|
||||
purge_pending.size = 1;
|
||||
purge_pending.deleted = true;
|
||||
purge_pending.erasure.index = disk_index;
|
||||
purge_pending.add_object_part(1, "part-etag-1".to_string(), 1, None, 1, None, None);
|
||||
purge_pending
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
let errs = vec![None; 6];
|
||||
|
||||
assert!(parts_metadata.iter().all(file_info_is_valid_for_metadata));
|
||||
assert!(parts_metadata.iter().all(|metadata| !metadata.is_canonical_delete_marker()));
|
||||
assert_eq!(SetDisks::list_object_parities(&parts_metadata, &errs), vec![1; 6]);
|
||||
assert_eq!(
|
||||
SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 3)
|
||||
.expect("purge-pending payload should retain its EC:1 quorum"),
|
||||
(5, 5)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_conv_part_err_to_int() {
|
||||
// Test error conversion to integer codes
|
||||
@@ -7141,7 +7467,8 @@ mod tests {
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let (updated, write_quorum) = build_tiered_decommission_file_info("bucket", "object", &original, 16, 4, None);
|
||||
let layout = WriteLayout::from_parity(16, 4).expect("tiered write layout should be valid");
|
||||
let updated = build_tiered_decommission_file_info("bucket", "object", &original, layout);
|
||||
|
||||
assert_eq!(updated.version_id, original.version_id);
|
||||
assert_eq!(updated.transition_status, original.transition_status);
|
||||
@@ -7150,7 +7477,7 @@ mod tests {
|
||||
assert_eq!(updated.transition_version_id, original.transition_version_id);
|
||||
assert_eq!(updated.erasure.data_blocks, 12);
|
||||
assert_eq!(updated.erasure.parity_blocks, 4);
|
||||
assert_eq!(write_quorum, 12);
|
||||
assert_eq!(layout.write_quorum, 12);
|
||||
assert_ne!(updated.erasure.distribution, original.erasure.distribution);
|
||||
}
|
||||
|
||||
@@ -8395,11 +8722,15 @@ mod tests {
|
||||
}
|
||||
|
||||
async fn make_local_bucket_test_set_disks() -> Arc<SetDisks> {
|
||||
let format = FormatV3::new(1, 2);
|
||||
make_local_bucket_test_set_disks_with_drive_count(2).await
|
||||
}
|
||||
|
||||
async fn make_local_bucket_test_set_disks_with_drive_count(drive_count: usize) -> Arc<SetDisks> {
|
||||
let format = FormatV3::new(1, drive_count);
|
||||
let mut endpoints = Vec::new();
|
||||
let mut disks = Vec::new();
|
||||
|
||||
for disk_idx in 0..2 {
|
||||
for disk_idx in 0..drive_count {
|
||||
let dir = tempfile::tempdir().expect("tempdir should be created");
|
||||
let mut endpoint =
|
||||
Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse");
|
||||
@@ -8428,18 +8759,23 @@ mod tests {
|
||||
disks.push(Some(disk));
|
||||
}
|
||||
|
||||
SetDisks::new(
|
||||
let set_disks = SetDisks::new(
|
||||
"test-owner".to_string(),
|
||||
Arc::new(RwLock::new(disks)),
|
||||
2,
|
||||
1,
|
||||
drive_count,
|
||||
drive_count / 2,
|
||||
0,
|
||||
0,
|
||||
endpoints,
|
||||
format,
|
||||
Vec::new(),
|
||||
)
|
||||
.await
|
||||
.await;
|
||||
set_disks.set_test_storage_class_config(
|
||||
storageclass::lookup_config_for_pools_without_env(&rustfs_config::server_config::KVS::new(), &[drive_count])
|
||||
.expect("test storage class should resolve for the local drive count"),
|
||||
);
|
||||
set_disks
|
||||
}
|
||||
|
||||
async fn make_local_bucket_test_set_disks_with_missing_format() -> Arc<SetDisks> {
|
||||
@@ -8909,7 +9245,7 @@ mod tests {
|
||||
|
||||
#[tokio::test]
|
||||
async fn set_level_versioned_delete_marker_hides_object_without_corrupting_version_metadata() {
|
||||
let set_disks = make_local_bucket_test_set_disks().await;
|
||||
let set_disks = make_local_bucket_test_set_disks_with_drive_count(4).await;
|
||||
let bucket = "bucket-versioned-delete";
|
||||
let object = "object.txt";
|
||||
let opts = ObjectOptions {
|
||||
|
||||
@@ -771,7 +771,7 @@ impl SetDisks {
|
||||
// A surviving valid, non-deleted, non-remote data FileInfo to rebuild from.
|
||||
let Some(surviving) = parts_metadata
|
||||
.iter()
|
||||
.find(|fi| fi.is_valid() && !fi.deleted && !fi.is_remote())
|
||||
.find(|fi| fi.has_valid_erasure_geometry() && !fi.deleted && !fi.is_remote())
|
||||
.cloned()
|
||||
else {
|
||||
return Ok(false);
|
||||
@@ -1013,7 +1013,13 @@ impl SetDisks {
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
if lfi.is_valid() {
|
||||
// Report the object's own parity only when it actually carries erasure
|
||||
// geometry; delete markers and geometry-less versions fall back to the
|
||||
// pool default. Uses `has_valid_erasure_geometry()` (not `is_valid()`)
|
||||
// to stay in step with the rest of the metadata-predicate migration —
|
||||
// `is_valid()` now requires full payload validation and returns `false`
|
||||
// for delete markers, which would misreport their parity here.
|
||||
if lfi.has_valid_erasure_geometry() {
|
||||
result.parity_blocks = lfi.erasure.parity_blocks;
|
||||
} else {
|
||||
result.parity_blocks = self.default_parity_count;
|
||||
|
||||
@@ -845,6 +845,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
||||
#[tracing::instrument(skip(self))]
|
||||
async fn new_multipart_upload(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<MultipartUploadResult> {
|
||||
crate::hp_guard!("SetDisks::new_multipart_upload");
|
||||
let storage_class_config = self.storage_class_config_snapshot();
|
||||
let mut _object_lock_guard = None;
|
||||
|
||||
if opts.http_preconditions.is_some() {
|
||||
@@ -876,18 +877,18 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
||||
let _ = user_defined.remove(AMZ_STORAGE_CLASS);
|
||||
}
|
||||
|
||||
let sc_parity_drives = runtime_sources::storage_class_parity(user_defined.get(AMZ_STORAGE_CLASS).map(String::as_str));
|
||||
|
||||
let mut parity_drives = sc_parity_drives.unwrap_or(self.default_parity_count);
|
||||
if opts.max_parity {
|
||||
parity_drives = disks.len() / 2;
|
||||
}
|
||||
|
||||
let data_drives = disks.len() - parity_drives;
|
||||
let mut write_quorum = data_drives;
|
||||
if data_drives == parity_drives {
|
||||
write_quorum += 1
|
||||
}
|
||||
let WriteLayout {
|
||||
data_drives,
|
||||
parity_drives,
|
||||
write_quorum,
|
||||
} = resolve_write_layout(
|
||||
&storage_class_config,
|
||||
self.pool_index,
|
||||
disks.len(),
|
||||
self.default_parity_count,
|
||||
user_defined.get(AMZ_STORAGE_CLASS).map(String::as_str),
|
||||
opts.max_parity,
|
||||
)?;
|
||||
|
||||
let mut fi = FileInfo::new([bucket, object].join("/").as_str(), data_drives, parity_drives);
|
||||
|
||||
@@ -1369,7 +1370,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
||||
}
|
||||
|
||||
for meta in parts_metadatas.iter_mut() {
|
||||
if meta.is_valid() {
|
||||
if meta.has_valid_erasure_geometry() {
|
||||
meta.size = fi.size;
|
||||
meta.mod_time = fi.mod_time;
|
||||
meta.parts.clone_from(&fi.parts);
|
||||
@@ -1556,10 +1557,14 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::config::storageclass::lookup_config_for_pools_without_env;
|
||||
use crate::disk::DiskAPI as _;
|
||||
use crate::disk::{endpoint::Endpoint, format::FormatV3};
|
||||
use crate::set_disk::ops::object::hermetic_set_disks_support::hermetic_set_disks;
|
||||
use crate::set_disk::ops::object::hermetic_set_disks_support::{
|
||||
hermetic_set_disks, hermetic_set_disks_for_pool_with_default_parity,
|
||||
};
|
||||
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
|
||||
use rustfs_config::server_config::KVS;
|
||||
use rustfs_lock::{LockClient, client::local::LocalClient};
|
||||
use serial_test::serial;
|
||||
use tempfile::TempDir;
|
||||
@@ -1890,6 +1895,65 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn second_pool_multipart_uses_its_own_layout_and_round_trips() {
|
||||
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_for_pool_with_default_parity(2, 1, 2).await;
|
||||
set_disks.set_test_storage_class_config(
|
||||
lookup_config_for_pools_without_env(&KVS::new(), &[4, 2]).expect("heterogeneous pool storage class should resolve"),
|
||||
);
|
||||
|
||||
let bucket = "multipart-second-pool-bucket";
|
||||
let object = "object";
|
||||
for disk in &disk_stores {
|
||||
disk.make_volume(bucket).await.expect("bucket volume should be created");
|
||||
}
|
||||
|
||||
let upload = set_disks
|
||||
.new_multipart_upload(bucket, object, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("second-pool multipart upload should be created");
|
||||
let (upload_info, _) = set_disks
|
||||
.check_upload_id_exists(bucket, object, &upload.upload_id, false)
|
||||
.await
|
||||
.expect("stored multipart layout should be readable");
|
||||
assert_eq!(upload_info.erasure.data_blocks, 1);
|
||||
assert_eq!(upload_info.erasure.parity_blocks, 1);
|
||||
|
||||
let payload = vec![0x5a; 4096];
|
||||
let mut reader = PutObjReader::from_vec(payload.clone());
|
||||
let part = set_disks
|
||||
.put_object_part(bucket, object, &upload.upload_id, 1, &mut reader, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("second-pool part should encode without zero data shards");
|
||||
set_disks
|
||||
.clone()
|
||||
.complete_multipart_upload(
|
||||
bucket,
|
||||
object,
|
||||
&upload.upload_id,
|
||||
vec![CompletePart {
|
||||
part_num: part.part_num,
|
||||
etag: part.etag,
|
||||
..Default::default()
|
||||
}],
|
||||
&ObjectOptions::default(),
|
||||
)
|
||||
.await
|
||||
.expect("second-pool multipart upload should complete");
|
||||
|
||||
let mut object_reader = set_disks
|
||||
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("completed second-pool object should be readable");
|
||||
let mut restored = Vec::new();
|
||||
object_reader
|
||||
.stream
|
||||
.read_to_end(&mut restored)
|
||||
.await
|
||||
.expect("completed second-pool object should stream");
|
||||
assert_eq!(restored, payload);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn list_multipart_uploads_caps_each_page_at_max_uploads_and_paginates_cleanly() {
|
||||
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
|
||||
|
||||
@@ -702,6 +702,7 @@ impl SetDisks {
|
||||
opts: &ObjectOptions,
|
||||
) -> Result<(ObjectInfo, Option<OldCurrentSize>)> {
|
||||
crate::hp_guard!("SetDisks::put_object");
|
||||
let storage_class_config = self.storage_class_config_snapshot();
|
||||
self.invalidate_get_object_metadata_cache(bucket, object).await;
|
||||
|
||||
let disks = self.get_disks_internal().await;
|
||||
@@ -727,18 +728,18 @@ impl SetDisks {
|
||||
user_defined.insert(key.clone(), value.clone());
|
||||
}
|
||||
}
|
||||
let sc_parity_drives = runtime_sources::storage_class_parity(user_defined.get(AMZ_STORAGE_CLASS).map(String::as_str));
|
||||
|
||||
let mut parity_drives = sc_parity_drives.unwrap_or(self.default_parity_count);
|
||||
if opts.max_parity {
|
||||
parity_drives = disks.len() / 2;
|
||||
}
|
||||
|
||||
let data_drives = disks.len() - parity_drives;
|
||||
let mut write_quorum = data_drives;
|
||||
if data_drives == parity_drives {
|
||||
write_quorum += 1
|
||||
}
|
||||
let WriteLayout {
|
||||
data_drives,
|
||||
parity_drives,
|
||||
write_quorum,
|
||||
} = resolve_write_layout(
|
||||
&storage_class_config,
|
||||
self.pool_index,
|
||||
disks.len(),
|
||||
self.default_parity_count,
|
||||
user_defined.get(AMZ_STORAGE_CLASS).map(String::as_str),
|
||||
opts.max_parity,
|
||||
)?;
|
||||
|
||||
// if filtered_online < write_quorum {
|
||||
// warn!(
|
||||
@@ -776,8 +777,7 @@ impl SetDisks {
|
||||
let erasure = erasure_from_file_info(&fi, false)?;
|
||||
|
||||
let put_object_size = known_put_object_storage_size(data.size());
|
||||
let is_inline_buffer =
|
||||
runtime_sources::storage_class_should_inline(erasure.shard_file_size(put_object_size), opts.versioned);
|
||||
let is_inline_buffer = storage_class_config.should_inline(erasure.shard_file_size(put_object_size), opts.versioned);
|
||||
|
||||
let shard_file_size = erasure.shard_file_size(put_object_size);
|
||||
let shard_size = erasure.shard_size();
|
||||
@@ -1946,7 +1946,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
||||
let inline_data = fi.inline_data();
|
||||
|
||||
for fi in metas.iter_mut() {
|
||||
if fi.is_valid() {
|
||||
if fi.has_valid_erasure_geometry() {
|
||||
fi.metadata = (*src_info.user_defined).clone();
|
||||
if let Some(etag) = &src_info.etag {
|
||||
fi.metadata.insert("etag".to_owned(), etag.clone());
|
||||
@@ -3402,14 +3402,15 @@ pub(in crate::set_disk::ops) mod hermetic_set_disks_support {
|
||||
use tempfile::TempDir;
|
||||
use tokio::sync::RwLock;
|
||||
|
||||
pub(in crate::set_disk::ops) async fn make_formatted_local_disk(
|
||||
async fn make_formatted_local_disk_for_pool(
|
||||
disk_idx: usize,
|
||||
pool_index: usize,
|
||||
format: &FormatV3,
|
||||
) -> (TempDir, Endpoint, DiskStore) {
|
||||
let dir = tempfile::tempdir().expect("tempdir should be created");
|
||||
let mut endpoint =
|
||||
Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse");
|
||||
endpoint.set_pool_index(0);
|
||||
endpoint.set_pool_index(pool_index);
|
||||
endpoint.set_set_index(0);
|
||||
endpoint.set_disk_index(disk_idx);
|
||||
|
||||
@@ -3433,6 +3434,14 @@ pub(in crate::set_disk::ops) mod hermetic_set_disks_support {
|
||||
}
|
||||
|
||||
pub(in crate::set_disk::ops) async fn hermetic_set_disks(disk_count: usize) -> (Vec<TempDir>, Vec<DiskStore>, Arc<SetDisks>) {
|
||||
hermetic_set_disks_for_pool_with_default_parity(disk_count, 0, disk_count / 2).await
|
||||
}
|
||||
|
||||
pub(in crate::set_disk::ops) async fn hermetic_set_disks_for_pool_with_default_parity(
|
||||
disk_count: usize,
|
||||
pool_index: usize,
|
||||
default_parity_count: usize,
|
||||
) -> (Vec<TempDir>, Vec<DiskStore>, Arc<SetDisks>) {
|
||||
let format = FormatV3::new(1, disk_count);
|
||||
|
||||
let mut temp_dirs = Vec::with_capacity(disk_count);
|
||||
@@ -3441,7 +3450,7 @@ pub(in crate::set_disk::ops) mod hermetic_set_disks_support {
|
||||
let mut disks = Vec::with_capacity(disk_count);
|
||||
|
||||
for disk_idx in 0..disk_count {
|
||||
let (temp_dir, endpoint, disk) = make_formatted_local_disk(disk_idx, &format).await;
|
||||
let (temp_dir, endpoint, disk) = make_formatted_local_disk_for_pool(disk_idx, pool_index, &format).await;
|
||||
temp_dirs.push(temp_dir);
|
||||
endpoints.push(endpoint);
|
||||
disk_stores.push(disk.clone());
|
||||
@@ -3452,9 +3461,9 @@ pub(in crate::set_disk::ops) mod hermetic_set_disks_support {
|
||||
"hermetic-ops-test-owner".to_string(),
|
||||
Arc::new(RwLock::new(disks)),
|
||||
disk_count,
|
||||
disk_count / 2,
|
||||
0,
|
||||
default_parity_count,
|
||||
0,
|
||||
pool_index,
|
||||
endpoints,
|
||||
format,
|
||||
Vec::new(),
|
||||
@@ -4201,6 +4210,61 @@ mod transition_source_identity_matrix_tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod heterogeneous_pool_put_tests {
|
||||
use super::hermetic_set_disks_support::hermetic_set_disks_for_pool_with_default_parity;
|
||||
use super::*;
|
||||
use crate::config::storageclass::lookup_config_for_pools_without_env;
|
||||
use crate::disk::{DiskAPI as _, ReadOptions};
|
||||
use rustfs_config::server_config::KVS;
|
||||
use tokio::io::AsyncReadExt;
|
||||
|
||||
#[tokio::test]
|
||||
async fn second_pool_regular_put_uses_its_own_layout_and_round_trips() {
|
||||
// Deliberately inject the first pool's invalid scalar fallback. The
|
||||
// test can pass only if the production PUT uses the held [4, 2]
|
||||
// storage-class snapshot and resolves pool 1 to parity 1.
|
||||
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_for_pool_with_default_parity(2, 1, 2).await;
|
||||
set_disks.set_test_storage_class_config(
|
||||
lookup_config_for_pools_without_env(&KVS::new(), &[4, 2]).expect("heterogeneous pool storage class should resolve"),
|
||||
);
|
||||
|
||||
let bucket = "regular-put-second-pool-bucket";
|
||||
let object = "object";
|
||||
for disk in &disk_stores {
|
||||
disk.make_volume(bucket).await.expect("bucket volume should be created");
|
||||
}
|
||||
|
||||
let payload = vec![0x3c; 4096];
|
||||
let mut reader = PutObjReader::from_vec(payload.clone());
|
||||
set_disks
|
||||
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("second-pool regular PUT should encode without zero data shards");
|
||||
|
||||
for (disk_index, disk) in disk_stores.iter().enumerate() {
|
||||
let file_info = disk
|
||||
.read_version("", bucket, object, "", &ReadOptions::default())
|
||||
.await
|
||||
.unwrap_or_else(|err| panic!("disk {disk_index} should persist valid second-pool metadata: {err}"));
|
||||
assert_eq!(file_info.erasure.data_blocks, 1);
|
||||
assert_eq!(file_info.erasure.parity_blocks, 1);
|
||||
}
|
||||
|
||||
let mut object_reader = set_disks
|
||||
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("second-pool regular PUT should be readable");
|
||||
let mut restored = Vec::new();
|
||||
object_reader
|
||||
.stream
|
||||
.read_to_end(&mut restored)
|
||||
.await
|
||||
.expect("second-pool regular PUT should stream");
|
||||
assert_eq!(restored, payload);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod put_object_tmp_cleanup_tests {
|
||||
//! Regression coverage for backlog#924 (HP-3): the speculative tmp-dir
|
||||
|
||||
@@ -121,7 +121,7 @@ impl SetDisks {
|
||||
online_disks: &[Option<DiskStore>],
|
||||
read_quorum: usize,
|
||||
) {
|
||||
if fi.deleted || !fi.is_valid() {
|
||||
if fi.deleted || !fi.has_valid_erasure_geometry() {
|
||||
return;
|
||||
}
|
||||
let (bucket, object) = identity;
|
||||
@@ -3192,10 +3192,13 @@ mod tests {
|
||||
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 delete_marker = FileInfo {
|
||||
name: "object".to_string(),
|
||||
deleted: true,
|
||||
version_id: Some(Uuid::parse_str("00000000-0000-0000-0000-000000000003").expect("static uuid should parse")),
|
||||
mod_time: Some(OffsetDateTime::from_unix_timestamp(3).expect("static timestamp should parse")),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let diagnostics = MetadataFanoutDiagnostics::new(
|
||||
Duration::from_millis(9),
|
||||
@@ -3296,6 +3299,47 @@ mod tests {
|
||||
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metadata_quorum_accumulator_rejects_semantic_field_splits() {
|
||||
type FileInfoMutation = fn(&mut FileInfo);
|
||||
|
||||
let mutations: &[(&str, FileInfoMutation)] = &[
|
||||
("transition_status", |fi| fi.transition_status = "complete".to_string()),
|
||||
("transitioned_objname", |fi| fi.transitioned_objname = "remote-object".to_string()),
|
||||
("transition_tier", |fi| fi.transition_tier = "WARM".to_string()),
|
||||
("transition_version_id", |fi| fi.transition_version_id = Some(Uuid::from_u128(10))),
|
||||
("expire_restored", |fi| fi.expire_restored = true),
|
||||
("written_by_version", |fi| fi.written_by_version = Some(1)),
|
||||
("replication_state_internal", |fi| {
|
||||
fi.replication_state_internal = Some(Default::default())
|
||||
}),
|
||||
("num_versions", |fi| fi.num_versions = 2),
|
||||
("successor_mod_time", |fi| {
|
||||
fi.successor_mod_time = Some(OffsetDateTime::from_unix_timestamp(10).expect("static timestamp should parse"));
|
||||
}),
|
||||
];
|
||||
|
||||
for (field, mutate) in mutations {
|
||||
let mut accumulator = metadata_early_stop_accumulator();
|
||||
let first = metadata_early_stop_candidate("object", 1);
|
||||
let mut second = metadata_early_stop_candidate("object", 2);
|
||||
let mut third = metadata_early_stop_candidate("object", 3);
|
||||
mutate(&mut second);
|
||||
mutate(&mut third);
|
||||
|
||||
accumulator.observe_file_info(&first);
|
||||
accumulator.observe_file_info(&second);
|
||||
accumulator.observe_file_info(&third);
|
||||
|
||||
assert!(accumulator.conflicting_metadata, "split {field} must block metadata early-stop");
|
||||
assert!(
|
||||
accumulator.early_stop_decision().is_none(),
|
||||
"split {field} must not be mistaken for three matching votes"
|
||||
);
|
||||
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_CONFLICTING_METADATA);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metadata_quorum_accumulator_falls_back_on_split_data_dir() {
|
||||
let mut accumulator = metadata_early_stop_accumulator();
|
||||
@@ -3325,8 +3369,15 @@ mod tests {
|
||||
#[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;
|
||||
let deleted = FileInfo {
|
||||
name: "object".to_string(),
|
||||
deleted: true,
|
||||
version_id: Some(Uuid::new_v4()),
|
||||
mod_time: Some(OffsetDateTime::now_utc()),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
assert!(!deleted.is_valid(), "real delete markers do not carry erasure geometry");
|
||||
|
||||
accumulator.observe_file_info(&deleted);
|
||||
accumulator.observe_file_info(&deleted);
|
||||
@@ -3344,8 +3395,13 @@ mod tests {
|
||||
#[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;
|
||||
let deleted = FileInfo {
|
||||
name: "object".to_string(),
|
||||
deleted: true,
|
||||
version_id: Some(Uuid::parse_str("00000000-0000-0000-0000-000000000004").expect("static uuid should parse")),
|
||||
mod_time: Some(OffsetDateTime::from_unix_timestamp(4).expect("static timestamp should parse")),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
accumulator.observe_file_info(&deleted);
|
||||
|
||||
@@ -3353,6 +3409,49 @@ mod tests {
|
||||
assert_eq!(accumulator.final_miss_reason(), GET_METADATA_EARLY_STOP_REASON_DELETE_MARKER);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metadata_quorum_accumulator_does_not_treat_purge_pending_payload_as_delete_marker() {
|
||||
let mut accumulator = MetadataQuorumAccumulator::new(6, 3, true);
|
||||
let version_id = Uuid::parse_str("00000000-0000-0000-0000-000000000005").expect("static uuid should parse");
|
||||
|
||||
for disk_index in 1..=4 {
|
||||
let mut purge_pending = FileInfo::new("object", 5, 1);
|
||||
purge_pending.name = "object".to_string();
|
||||
purge_pending.version_id = Some(version_id);
|
||||
purge_pending.mod_time = Some(OffsetDateTime::from_unix_timestamp(5).expect("static timestamp should parse"));
|
||||
purge_pending.size = 1;
|
||||
purge_pending.deleted = true;
|
||||
purge_pending.erasure.index = disk_index;
|
||||
purge_pending.add_object_part(1, "part-etag-1".to_string(), 1, None, 1, None, None);
|
||||
|
||||
accumulator.observe_file_info(&purge_pending);
|
||||
}
|
||||
|
||||
assert!(!accumulator.delete_marker_seen);
|
||||
assert_eq!(accumulator.candidate_votes, 4);
|
||||
assert!(
|
||||
accumulator.early_stop_decision().is_none(),
|
||||
"EC:1 purge-pending payload on six disks still requires five matching payload votes"
|
||||
);
|
||||
|
||||
let mut fifth = FileInfo::new("object", 5, 1);
|
||||
fifth.name = "object".to_string();
|
||||
fifth.version_id = Some(version_id);
|
||||
fifth.mod_time = Some(OffsetDateTime::from_unix_timestamp(5).expect("static timestamp should parse"));
|
||||
fifth.size = 1;
|
||||
fifth.deleted = true;
|
||||
fifth.erasure.index = 5;
|
||||
fifth.add_object_part(1, "part-etag-1".to_string(), 1, None, 1, None, None);
|
||||
accumulator.observe_file_info(&fifth);
|
||||
|
||||
assert_eq!(
|
||||
accumulator.early_stop_decision(),
|
||||
Some(MetadataEarlyStopDecision {
|
||||
reason: GET_METADATA_EARLY_STOP_REASON_VALID_QUORUM
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metadata_quorum_accumulator_falls_back_on_object_not_found_quorum() {
|
||||
let mut accumulator = metadata_early_stop_accumulator();
|
||||
|
||||
Reference in New Issue
Block a user