mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 693800db8e | |||
| 20d1266496 | |||
| f7003dfddd |
Generated
+1
@@ -12688,6 +12688,7 @@ dependencies = [
|
||||
"js-sys",
|
||||
"rand 0.10.2",
|
||||
"serde_core",
|
||||
"sha1_smol",
|
||||
"wasm-bindgen",
|
||||
]
|
||||
|
||||
|
||||
+627
-112
File diff suppressed because it is too large
Load Diff
@@ -988,14 +988,11 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for Sets {
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl crate::storage_api_contracts::heal::HealOperations for Sets {
|
||||
type Error = Error;
|
||||
type HealResultItem = HealResultItem;
|
||||
type HealOptions = HealOpts;
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
|
||||
impl Sets {
|
||||
pub(crate) async fn heal_format_with_fence<F>(&self, dry_run: bool, fence_lost: F) -> Result<(HealResultItem, Option<Error>)>
|
||||
where
|
||||
F: Fn() -> bool + Send + Sync,
|
||||
{
|
||||
let (disks, init_errs) = init_storage_disks_with_errors(
|
||||
&self.endpoints.endpoints,
|
||||
&DiskOption {
|
||||
@@ -1068,6 +1065,9 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets {
|
||||
// Save new formats `format.json` on unformatted disks.
|
||||
for (index, (fm, disk)) in tmp_new_formats.iter_mut().zip(disks.iter()).enumerate() {
|
||||
if fm.is_some() && disk.is_some() {
|
||||
if fence_lost() {
|
||||
return Ok((res, Some(StorageError::SlowDown)));
|
||||
}
|
||||
if let Err(err) = save_format_file(disk, fm).await {
|
||||
if let Some(disk) = disk.as_ref() {
|
||||
let _ = disk.close().await;
|
||||
@@ -1101,6 +1101,18 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets {
|
||||
}
|
||||
Ok((res, None))
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl crate::storage_api_contracts::heal::HealOperations for Sets {
|
||||
type Error = Error;
|
||||
type HealResultItem = HealResultItem;
|
||||
type HealOptions = HealOpts;
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
|
||||
self.heal_format_with_fence(dry_run, || false).await
|
||||
}
|
||||
#[tracing::instrument(skip(self))]
|
||||
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
|
||||
let mut result = HealResultItem {
|
||||
|
||||
@@ -4736,7 +4736,15 @@ impl SetDisks {
|
||||
achieved: 0,
|
||||
});
|
||||
}
|
||||
let parts_metadata = vec![fi.clone(); disks.len()];
|
||||
// Rebuilt tiered metadata starts with index zero, but shuffling validates
|
||||
// each source slot before assigning the shuffled index below.
|
||||
let parts_metadata: Vec<FileInfo> = (0..disks.len())
|
||||
.map(|disk_index| {
|
||||
let mut part = fi.clone();
|
||||
part.erasure.index = fi.erasure.distribution[disk_index];
|
||||
part
|
||||
})
|
||||
.collect();
|
||||
let (shuffle_disks, parts_metadata) = Self::shuffle_disks_and_parts_metadata(&disks, &parts_metadata, &fi);
|
||||
|
||||
let mut errs = Vec::with_capacity(shuffle_disks.len());
|
||||
|
||||
@@ -13,7 +13,12 @@
|
||||
// limitations under the License.
|
||||
|
||||
use super::*;
|
||||
use crate::core::pools::POOL_META_NAME;
|
||||
use crate::services::rebalance::{REBAL_META_NAME, RebalStatus};
|
||||
use crate::set_disk::get_lock_acquire_timeout;
|
||||
use crate::storage_api_contracts::heal::HealOperations as _;
|
||||
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
|
||||
use rustfs_lock::NamespaceLockGuard;
|
||||
use tracing::trace;
|
||||
|
||||
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
|
||||
@@ -30,7 +35,119 @@ fn invalid_heal_pool_index(pool_idx: usize, pool_count: usize) -> Error {
|
||||
)
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
enum HealFormatPoolSkip {
|
||||
Completed,
|
||||
Retryable,
|
||||
}
|
||||
|
||||
fn classify_heal_format_pool(
|
||||
pool_idx: usize,
|
||||
pool_cmd_line: &str,
|
||||
pool_meta: &PoolMeta,
|
||||
rebalance_meta: Option<&RebalanceMeta>,
|
||||
) -> Option<HealFormatPoolSkip> {
|
||||
let Some(pool) = pool_meta.pools.get(pool_idx) else {
|
||||
return Some(HealFormatPoolSkip::Retryable);
|
||||
};
|
||||
|
||||
if pool.id != pool_idx || pool_cmd_line.is_empty() || pool.cmd_line.is_empty() || pool.cmd_line != pool_cmd_line {
|
||||
return Some(HealFormatPoolSkip::Retryable);
|
||||
}
|
||||
|
||||
if let Some(decommission) = pool.decommission.as_ref() {
|
||||
if decommission.complete {
|
||||
return Some(HealFormatPoolSkip::Completed);
|
||||
}
|
||||
if decommission.failed || decommission.canceled || decommission.queued || pool_meta.is_suspended(pool_idx) {
|
||||
return Some(HealFormatPoolSkip::Retryable);
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(meta) = rebalance_meta {
|
||||
let Some(pool_stats) = meta.pool_stats.get(pool_idx) else {
|
||||
return Some(HealFormatPoolSkip::Retryable);
|
||||
};
|
||||
if pool_stats.info.stopping || (pool_stats.participating && pool_stats.info.status == RebalStatus::Started) {
|
||||
return Some(HealFormatPoolSkip::Retryable);
|
||||
}
|
||||
}
|
||||
|
||||
None
|
||||
}
|
||||
|
||||
fn heal_format_pool_skip_error(skip: HealFormatPoolSkip) -> Error {
|
||||
match skip {
|
||||
HealFormatPoolSkip::Completed => StorageError::NoHealRequired,
|
||||
HealFormatPoolSkip::Retryable => StorageError::SlowDown,
|
||||
}
|
||||
}
|
||||
|
||||
fn heal_format_fence_lost_error() -> Error {
|
||||
StorageError::SlowDown
|
||||
}
|
||||
|
||||
impl ECStore {
|
||||
async fn acquire_heal_format_fence(
|
||||
&self,
|
||||
) -> Result<(NamespaceLockGuard, NamespaceLockGuard, PoolMeta, Option<RebalanceMeta>)> {
|
||||
let metadata_pool = self
|
||||
.pools
|
||||
.first()
|
||||
.cloned()
|
||||
.ok_or_else(|| Error::other("heal format requires at least one storage pool"))?;
|
||||
|
||||
// Metadata fence order is part of the decommission/rebalance protocol:
|
||||
// pool.bin must always be acquired before rebalance.bin.
|
||||
let pool_lock = metadata_pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?;
|
||||
let pool_guard = pool_lock.get_write_lock(get_lock_acquire_timeout()).await?;
|
||||
let rebalance_lock = metadata_pool.new_ns_lock(RUSTFS_META_BUCKET, REBAL_META_NAME).await?;
|
||||
let rebalance_guard = rebalance_lock.get_write_lock(get_lock_acquire_timeout()).await?;
|
||||
|
||||
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
|
||||
return Err(heal_format_fence_lost_error());
|
||||
}
|
||||
|
||||
let mut pool_meta = PoolMeta::default();
|
||||
pool_meta.load_no_lock(metadata_pool.clone()).await?;
|
||||
if pool_meta.pools.len() != self.pools.len()
|
||||
|| pool_meta.pools.iter().enumerate().any(|(pool_idx, pool)| {
|
||||
pool.id != pool_idx || pool.cmd_line.is_empty() || pool.cmd_line != self.pools[pool_idx].endpoints.cmd_line
|
||||
})
|
||||
{
|
||||
return Err(heal_format_fence_lost_error());
|
||||
}
|
||||
|
||||
let mut rebalance_meta = RebalanceMeta::new();
|
||||
let rebalance_meta = match rebalance_meta
|
||||
.load_with_opts(
|
||||
metadata_pool,
|
||||
ObjectOptions {
|
||||
no_lock: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(()) => Some(rebalance_meta),
|
||||
Err(Error::ConfigNotFound) => None,
|
||||
Err(err) => return Err(err),
|
||||
};
|
||||
|
||||
if rebalance_meta
|
||||
.as_ref()
|
||||
.is_some_and(|meta| meta.pool_stats.len() != self.pools.len())
|
||||
{
|
||||
return Err(heal_format_fence_lost_error());
|
||||
}
|
||||
|
||||
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
|
||||
return Err(heal_format_fence_lost_error());
|
||||
}
|
||||
|
||||
Ok((pool_guard, rebalance_guard, pool_meta, rebalance_meta))
|
||||
}
|
||||
|
||||
fn get_pools_for_heal_object(&self, opts: &HealOpts) -> Result<Vec<Arc<Sets>>> {
|
||||
match opts.pool {
|
||||
Some(pool_idx) => Ok(vec![
|
||||
@@ -52,9 +169,26 @@ impl ECStore {
|
||||
};
|
||||
|
||||
let mut count_no_heal = 0;
|
||||
let mut count_completed = 0;
|
||||
let mut first_error = None;
|
||||
for pool in self.pools.iter() {
|
||||
let (mut result, err) = pool.heal_format(dry_run).await?;
|
||||
for (pool_idx, pool) in self.pools.iter().enumerate() {
|
||||
let (pool_guard, rebalance_guard, pool_meta, rebalance_meta) = self.acquire_heal_format_fence().await?;
|
||||
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
|
||||
first_error.get_or_insert(heal_format_fence_lost_error());
|
||||
break;
|
||||
}
|
||||
if let Some(skip) = classify_heal_format_pool(pool_idx, &pool.endpoints.cmd_line, &pool_meta, rebalance_meta.as_ref())
|
||||
{
|
||||
if matches!(skip, HealFormatPoolSkip::Completed) {
|
||||
count_completed += 1;
|
||||
} else {
|
||||
first_error.get_or_insert(heal_format_pool_skip_error(skip));
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
let fence_lost = || pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost();
|
||||
let (mut result, err) = pool.heal_format_with_fence(dry_run, fence_lost).await?;
|
||||
if let Some(err) = err {
|
||||
match err {
|
||||
StorageError::NoHealRequired => {
|
||||
@@ -69,11 +203,18 @@ impl ECStore {
|
||||
r.set_count += result.set_count;
|
||||
r.before.drives.append(&mut result.before.drives);
|
||||
r.after.drives.append(&mut result.after.drives);
|
||||
|
||||
// A lease can be lost after the final write; fail closed before
|
||||
// reporting the pool as successfully healed.
|
||||
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
|
||||
first_error.get_or_insert(heal_format_fence_lost_error());
|
||||
break;
|
||||
}
|
||||
}
|
||||
if let Some(err) = first_error {
|
||||
return Ok((r, Some(err)));
|
||||
}
|
||||
if count_no_heal == self.pools.len() {
|
||||
if count_no_heal + count_completed == self.pools.len() {
|
||||
info!(
|
||||
event = EVENT_HEAL_FORMAT_COMPLETED,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
@@ -302,6 +443,7 @@ mod tests {
|
||||
use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk};
|
||||
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
|
||||
use crate::runtime::instance::InstanceContext;
|
||||
use crate::services::rebalance::{RebalanceInfo, RebalanceStats};
|
||||
use crate::storage_api_contracts::bucket::{BucketOperations, MakeBucketOptions};
|
||||
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations};
|
||||
use crate::store::init_format::{load_format_erasure, save_format_file};
|
||||
@@ -353,6 +495,164 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
fn pool_meta_with_decommission(info: PoolDecommissionInfo) -> PoolMeta {
|
||||
PoolMeta {
|
||||
pools: vec![PoolStatus {
|
||||
id: 0,
|
||||
cmd_line: "pool-0".to_string(),
|
||||
last_update: OffsetDateTime::UNIX_EPOCH,
|
||||
decommission: Some(info),
|
||||
}],
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn heal_format_pool_state_barriers_are_classified() {
|
||||
let active = pool_meta_with_decommission(PoolDecommissionInfo {
|
||||
start_time: Some(OffsetDateTime::UNIX_EPOCH),
|
||||
..Default::default()
|
||||
});
|
||||
assert!(matches!(
|
||||
classify_heal_format_pool(0, "pool-0", &active, None),
|
||||
Some(HealFormatPoolSkip::Retryable)
|
||||
));
|
||||
|
||||
for info in [
|
||||
PoolDecommissionInfo {
|
||||
failed: true,
|
||||
..Default::default()
|
||||
},
|
||||
PoolDecommissionInfo {
|
||||
canceled: true,
|
||||
..Default::default()
|
||||
},
|
||||
] {
|
||||
assert!(matches!(
|
||||
classify_heal_format_pool(0, "pool-0", &pool_meta_with_decommission(info), None),
|
||||
Some(HealFormatPoolSkip::Retryable)
|
||||
));
|
||||
}
|
||||
|
||||
let completed = pool_meta_with_decommission(PoolDecommissionInfo {
|
||||
complete: true,
|
||||
..Default::default()
|
||||
});
|
||||
assert!(matches!(
|
||||
classify_heal_format_pool(0, "pool-0", &completed, None),
|
||||
Some(HealFormatPoolSkip::Completed)
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn heal_format_pool_rebalance_barriers_and_identity_are_fail_closed() {
|
||||
let identity_meta = pool_meta_with_decommission(PoolDecommissionInfo::default());
|
||||
let rebalance = RebalanceMeta {
|
||||
pool_stats: vec![RebalanceStats {
|
||||
participating: true,
|
||||
info: RebalanceInfo {
|
||||
status: RebalStatus::Started,
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
assert!(matches!(
|
||||
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&rebalance)),
|
||||
Some(HealFormatPoolSkip::Retryable)
|
||||
));
|
||||
|
||||
let stopping = RebalanceMeta {
|
||||
pool_stats: vec![RebalanceStats {
|
||||
info: RebalanceInfo {
|
||||
stopping: true,
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
assert!(matches!(
|
||||
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopping)),
|
||||
Some(HealFormatPoolSkip::Retryable)
|
||||
));
|
||||
|
||||
let identity = pool_meta_with_decommission(PoolDecommissionInfo::default());
|
||||
assert!(matches!(
|
||||
classify_heal_format_pool(0, "pool-new", &identity, None),
|
||||
Some(HealFormatPoolSkip::Retryable)
|
||||
));
|
||||
|
||||
let identity_without_decommission = PoolMeta {
|
||||
pools: vec![PoolStatus {
|
||||
id: 0,
|
||||
cmd_line: "pool-0".to_string(),
|
||||
last_update: OffsetDateTime::UNIX_EPOCH,
|
||||
decommission: None,
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
assert!(matches!(
|
||||
classify_heal_format_pool(0, "pool-new", &identity_without_decommission, None),
|
||||
Some(HealFormatPoolSkip::Retryable)
|
||||
));
|
||||
|
||||
assert!(matches!(
|
||||
classify_heal_format_pool(0, "", &identity_meta, None),
|
||||
Some(HealFormatPoolSkip::Retryable)
|
||||
));
|
||||
|
||||
assert!(matches!(
|
||||
classify_heal_format_pool(0, "pool-0", &PoolMeta::default(), None),
|
||||
Some(HealFormatPoolSkip::Retryable)
|
||||
));
|
||||
|
||||
let stopped = RebalanceMeta {
|
||||
stopped_at: Some(OffsetDateTime::UNIX_EPOCH),
|
||||
pool_stats: vec![RebalanceStats {
|
||||
participating: true,
|
||||
info: RebalanceInfo {
|
||||
status: RebalStatus::Stopped,
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
assert!(classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopped)).is_none());
|
||||
|
||||
let stopping_after_stop = RebalanceMeta {
|
||||
stopped_at: Some(OffsetDateTime::UNIX_EPOCH),
|
||||
pool_stats: vec![RebalanceStats {
|
||||
participating: true,
|
||||
info: RebalanceInfo {
|
||||
status: RebalStatus::Started,
|
||||
stopping: true,
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
assert!(matches!(
|
||||
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopping_after_stop)),
|
||||
Some(HealFormatPoolSkip::Retryable)
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn skipped_heal_format_pool_is_never_reported_as_success() {
|
||||
assert!(matches!(
|
||||
heal_format_pool_skip_error(HealFormatPoolSkip::Retryable),
|
||||
StorageError::SlowDown
|
||||
));
|
||||
assert!(matches!(
|
||||
heal_format_pool_skip_error(HealFormatPoolSkip::Completed),
|
||||
StorageError::NoHealRequired
|
||||
));
|
||||
}
|
||||
|
||||
async fn multi_pool_heal_store() -> (tempfile::TempDir, Arc<ECStore>, CancellationToken) {
|
||||
let temp_dir = tempfile::tempdir().expect("multi-pool heal test directory should be created");
|
||||
let mut pool_endpoints = Vec::new();
|
||||
@@ -889,6 +1189,18 @@ mod tests {
|
||||
bucket_fence_registry: std::sync::Arc::default(),
|
||||
};
|
||||
|
||||
let err = store
|
||||
.handle_heal_format(false)
|
||||
.await
|
||||
.expect_err("missing pool metadata must fail closed before format writes");
|
||||
assert!(matches!(err, StorageError::SlowDown));
|
||||
|
||||
let pool_meta = PoolMeta::new(&store.pools, &PoolMeta::default());
|
||||
pool_meta
|
||||
.save(store.pools.clone())
|
||||
.await
|
||||
.expect("pool metadata should be persisted before format heal");
|
||||
|
||||
let (result, err) = store
|
||||
.handle_heal_format(false)
|
||||
.await
|
||||
@@ -902,5 +1214,22 @@ mod tests {
|
||||
.await
|
||||
.expect("the later pool should be healed despite the first pool error");
|
||||
assert_eq!(healed.erasure.this, recoverable_format.erasure.sets[0][2]);
|
||||
|
||||
let mut completed_meta = PoolMeta::new(&store.pools, &PoolMeta::default());
|
||||
for status in &mut completed_meta.pools {
|
||||
status.decommission = Some(PoolDecommissionInfo {
|
||||
complete: true,
|
||||
..Default::default()
|
||||
});
|
||||
}
|
||||
completed_meta
|
||||
.save(store.pools.clone())
|
||||
.await
|
||||
.expect("completed pool metadata should be persisted");
|
||||
let (_, err) = store
|
||||
.handle_heal_format(false)
|
||||
.await
|
||||
.expect("completed pools should be reported as a no-op");
|
||||
assert!(matches!(err, Some(StorageError::NoHealRequired)));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -626,7 +626,7 @@ mod tests {
|
||||
io::Cursor,
|
||||
sync::{
|
||||
Arc,
|
||||
atomic::{AtomicBool, Ordering},
|
||||
atomic::{AtomicBool, AtomicUsize, Ordering},
|
||||
},
|
||||
time::Duration,
|
||||
};
|
||||
@@ -1307,6 +1307,85 @@ mod tests {
|
||||
});
|
||||
}
|
||||
|
||||
const DECOMMISSION_TEST_FAULT_STAGE_DELETE_MARKER: &str = "delete_marker_copy";
|
||||
const DECOMMISSION_TEST_FAULT_STAGE_MIGRATE_OBJECT: &str = "migrate_object";
|
||||
#[cfg(feature = "test-util")]
|
||||
const DECOMMISSION_TEST_FAULT_STAGE_TIERED: &str = "decommission_tiered_object";
|
||||
|
||||
async fn seed_decommission_source(
|
||||
store: &Arc<crate::store::ECStore>,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
body: Vec<u8>,
|
||||
opts: &ObjectOptions,
|
||||
) {
|
||||
let mut reader = PutObjReader::from_vec(body);
|
||||
store.pools[0]
|
||||
.put_object(bucket, object, &mut reader, opts)
|
||||
.await
|
||||
.expect("seed decommission source object");
|
||||
}
|
||||
|
||||
async fn run_decommission_entry_retry_test(
|
||||
store: &Arc<crate::store::ECStore>,
|
||||
rx: CancellationToken,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
expected_bucket_incarnation_id: Option<uuid::Uuid>,
|
||||
source_changed_exhaustions: Arc<AtomicUsize>,
|
||||
) -> crate::error::Result<()> {
|
||||
let source_set = store.pools[0].get_disks_by_key(object);
|
||||
store
|
||||
.decommission_entry_with_retry_state_for_test(
|
||||
rx,
|
||||
0,
|
||||
MetaCacheEntry {
|
||||
name: object.to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
bucket.to_string(),
|
||||
source_set,
|
||||
expected_bucket_incarnation_id,
|
||||
source_changed_exhaustions,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn read_decommission_target_body(
|
||||
store: &Arc<crate::store::ECStore>,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
opts: &ObjectOptions,
|
||||
) -> Vec<u8> {
|
||||
let mut reader = store.pools[1]
|
||||
.get_object_reader(bucket, object, None, HeaderMap::new(), opts)
|
||||
.await
|
||||
.expect("read decommission target object");
|
||||
let mut body = Vec::new();
|
||||
reader
|
||||
.stream
|
||||
.read_to_end(&mut body)
|
||||
.await
|
||||
.expect("drain decommission target body");
|
||||
body
|
||||
}
|
||||
|
||||
async fn assert_decommission_source_absent(
|
||||
store: &Arc<crate::store::ECStore>,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
opts: &ObjectOptions,
|
||||
) {
|
||||
let err = store.pools[0]
|
||||
.get_object_info(bucket, object, opts)
|
||||
.await
|
||||
.expect_err("decommission source must be retained until target commit, then removed");
|
||||
assert!(
|
||||
matches!(err, StorageError::ObjectNotFound(_, _) | StorageError::VersionNotFound(_, _, _)),
|
||||
"unexpected decommission source result: {err:?}"
|
||||
);
|
||||
}
|
||||
|
||||
async fn write_decommission_test_multipart_source(
|
||||
store: &Arc<crate::store::ECStore>,
|
||||
pool_idx: usize,
|
||||
@@ -3102,6 +3181,537 @@ mod tests {
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
fn decommission_entry_retries_source_changed_without_canceling_other_bucket() {
|
||||
let handle = std::thread::Builder::new()
|
||||
.name("decommission_entry_retries_source_changed_without_canceling_other_bucket".to_string())
|
||||
.stack_size(32 * 1024 * 1024)
|
||||
.spawn(|| {
|
||||
let runtime = tokio::runtime::Builder::new_multi_thread()
|
||||
.enable_all()
|
||||
.worker_threads(2)
|
||||
.build()
|
||||
.expect("test runtime should build");
|
||||
runtime.block_on(async {
|
||||
let temp_dir = tempfile::tempdir().expect("create decommission retry store dir");
|
||||
let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store(
|
||||
temp_dir.path(),
|
||||
"decommission-entry-retry",
|
||||
&[4, 4],
|
||||
))
|
||||
.await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
|
||||
let changed_bucket = format!("decom-retry-a-{}", uuid::Uuid::new_v4());
|
||||
let other_bucket = format!("decom-retry-b-{}", uuid::Uuid::new_v4());
|
||||
for bucket in [&changed_bucket, &other_bucket] {
|
||||
store
|
||||
.make_bucket(bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create decommission retry bucket");
|
||||
}
|
||||
|
||||
let changed_object = "changed.bin";
|
||||
let first_version = uuid::Uuid::new_v4();
|
||||
let second_version = uuid::Uuid::new_v4();
|
||||
let base_time = OffsetDateTime::now_utc();
|
||||
seed_decommission_source(
|
||||
&store,
|
||||
&changed_bucket,
|
||||
changed_object,
|
||||
b"first generation".to_vec(),
|
||||
&ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(first_version.to_string()),
|
||||
mod_time: Some(base_time),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await;
|
||||
let other_object = "other.bin";
|
||||
seed_decommission_source(
|
||||
&store,
|
||||
&other_bucket,
|
||||
other_object,
|
||||
b"other bucket generation".to_vec(),
|
||||
&ObjectOptions::default(),
|
||||
)
|
||||
.await;
|
||||
mark_test_pool_decommissioning(&store, 0).await;
|
||||
|
||||
let mutation_calls = Arc::new(AtomicUsize::new(0));
|
||||
let mutation_calls_for_hook = Arc::clone(&mutation_calls);
|
||||
let mutation_store = Arc::clone(&store);
|
||||
let mutation_bucket = changed_bucket.clone();
|
||||
let _mutation_guard = crate::core::pools::DecommissionCleanupMutationGuard::install(Arc::new(
|
||||
move |bucket, object, attempt| {
|
||||
let is_target = bucket == mutation_bucket.as_str() && object == changed_object;
|
||||
let calls = Arc::clone(&mutation_calls_for_hook);
|
||||
let store = Arc::clone(&mutation_store);
|
||||
let bucket = mutation_bucket.clone();
|
||||
Box::pin(async move {
|
||||
if !is_target {
|
||||
return;
|
||||
}
|
||||
calls.fetch_add(1, Ordering::SeqCst);
|
||||
if attempt == 1 {
|
||||
seed_decommission_source(
|
||||
&store,
|
||||
&bucket,
|
||||
changed_object,
|
||||
b"second generation".to_vec(),
|
||||
&ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(second_version.to_string()),
|
||||
mod_time: Some(base_time + time::Duration::seconds(1)),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await;
|
||||
}
|
||||
})
|
||||
},
|
||||
));
|
||||
|
||||
let ordinary_faults = Arc::new(AtomicUsize::new(0));
|
||||
let ordinary_faults_for_hook = Arc::clone(&ordinary_faults);
|
||||
let fault_bucket = other_bucket.clone();
|
||||
let _fault_guard = crate::core::pools::DecommissionTestFaultGuard::install(Arc::new(
|
||||
move |stage, bucket, object, attempt| {
|
||||
let injected = stage == DECOMMISSION_TEST_FAULT_STAGE_MIGRATE_OBJECT
|
||||
&& bucket == fault_bucket.as_str()
|
||||
&& object == other_object
|
||||
&& attempt < crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS;
|
||||
if injected {
|
||||
ordinary_faults_for_hook.fetch_add(1, Ordering::SeqCst);
|
||||
}
|
||||
injected
|
||||
},
|
||||
));
|
||||
|
||||
let rx = CancellationToken::new();
|
||||
let source_changed_exhaustions = Arc::new(AtomicUsize::new(0));
|
||||
let changed_incarnation = Some(
|
||||
store
|
||||
.bucket_incarnation_id(&changed_bucket)
|
||||
.await
|
||||
.expect("changed bucket incarnation"),
|
||||
);
|
||||
let other_incarnation = Some(
|
||||
store
|
||||
.bucket_incarnation_id(&other_bucket)
|
||||
.await
|
||||
.expect("other bucket incarnation"),
|
||||
);
|
||||
let (changed_result, other_result) = tokio::join!(
|
||||
run_decommission_entry_retry_test(
|
||||
&store,
|
||||
rx.clone(),
|
||||
&changed_bucket,
|
||||
changed_object,
|
||||
changed_incarnation,
|
||||
Arc::clone(&source_changed_exhaustions),
|
||||
),
|
||||
run_decommission_entry_retry_test(
|
||||
&store,
|
||||
rx.clone(),
|
||||
&other_bucket,
|
||||
other_object,
|
||||
other_incarnation,
|
||||
Arc::clone(&source_changed_exhaustions),
|
||||
)
|
||||
);
|
||||
changed_result.expect("SourceChanged entry retry must converge");
|
||||
other_result.expect("other bucket entry must continue through ordinary copy retries");
|
||||
|
||||
assert!(!rx.is_cancelled(), "entry-level SourceChanged must not cancel the shared worker token");
|
||||
assert_eq!(mutation_calls.load(Ordering::SeqCst), 2, "entry must be re-listed after SourceChanged");
|
||||
assert_eq!(ordinary_faults.load(Ordering::SeqCst), 2, "ordinary copy must consume the retry budget");
|
||||
assert_eq!(source_changed_exhaustions.load(Ordering::SeqCst), 0);
|
||||
|
||||
for (version_id, expected_body) in [
|
||||
(first_version, b"first generation".as_slice()),
|
||||
(second_version, b"second generation".as_slice()),
|
||||
] {
|
||||
let opts = ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(version_id.to_string()),
|
||||
..Default::default()
|
||||
};
|
||||
assert_decommission_source_absent(&store, &changed_bucket, changed_object, &opts).await;
|
||||
assert_eq!(
|
||||
read_decommission_target_body(&store, &changed_bucket, changed_object, &opts).await,
|
||||
expected_body
|
||||
);
|
||||
}
|
||||
assert_decommission_source_absent(&store, &other_bucket, other_object, &ObjectOptions::default()).await;
|
||||
assert_eq!(
|
||||
read_decommission_target_body(&store, &other_bucket, other_object, &ObjectOptions::default()).await,
|
||||
b"other bucket generation"
|
||||
);
|
||||
|
||||
shutdown.cancel();
|
||||
});
|
||||
})
|
||||
.expect("spawn decommission retry test thread");
|
||||
if let Err(payload) = handle.join() {
|
||||
std::panic::resume_unwind(payload);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
fn decommission_entry_exhausted_source_changed_retains_source_and_records_failure() {
|
||||
let handle = std::thread::Builder::new()
|
||||
.name("decommission_entry_exhausted_source_changed_retains_source_and_records_failure".to_string())
|
||||
.stack_size(32 * 1024 * 1024)
|
||||
.spawn(|| {
|
||||
let runtime = tokio::runtime::Builder::new_multi_thread()
|
||||
.enable_all()
|
||||
.worker_threads(2)
|
||||
.build()
|
||||
.expect("test runtime should build");
|
||||
runtime.block_on(async {
|
||||
let temp_dir = tempfile::tempdir().expect("create decommission exhaustion store dir");
|
||||
let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store(
|
||||
temp_dir.path(),
|
||||
"decommission-entry-exhaustion",
|
||||
&[4, 4],
|
||||
))
|
||||
.await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
|
||||
let bucket = format!("decom-exhausted-{}", uuid::Uuid::new_v4());
|
||||
let object = "exhausted.bin";
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create decommission exhaustion bucket");
|
||||
let original_version = uuid::Uuid::new_v4();
|
||||
let base_time = OffsetDateTime::now_utc();
|
||||
seed_decommission_source(
|
||||
&store,
|
||||
&bucket,
|
||||
object,
|
||||
b"original generation".to_vec(),
|
||||
&ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(original_version.to_string()),
|
||||
mod_time: Some(base_time),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await;
|
||||
mark_test_pool_decommissioning(&store, 0).await;
|
||||
|
||||
let mutation_calls = Arc::new(AtomicUsize::new(0));
|
||||
let mutation_calls_for_hook = Arc::clone(&mutation_calls);
|
||||
let mutation_store = Arc::clone(&store);
|
||||
let mutation_bucket = bucket.clone();
|
||||
let _mutation_guard = crate::core::pools::DecommissionCleanupMutationGuard::install(Arc::new(
|
||||
move |called_bucket, called_object, _attempt| {
|
||||
let is_target = called_bucket == mutation_bucket.as_str() && called_object == object;
|
||||
let calls = Arc::clone(&mutation_calls_for_hook);
|
||||
let store = Arc::clone(&mutation_store);
|
||||
let bucket = mutation_bucket.clone();
|
||||
Box::pin(async move {
|
||||
if !is_target {
|
||||
return;
|
||||
}
|
||||
let call = calls.fetch_add(1, Ordering::SeqCst) + 1;
|
||||
let offset = i64::try_from(call).expect("entry retry count should fit i64");
|
||||
seed_decommission_source(
|
||||
&store,
|
||||
&bucket,
|
||||
object,
|
||||
format!("concurrent generation {call}").into_bytes(),
|
||||
&ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(uuid::Uuid::new_v4().to_string()),
|
||||
mod_time: Some(base_time + time::Duration::seconds(offset)),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await;
|
||||
})
|
||||
},
|
||||
));
|
||||
|
||||
let rx = CancellationToken::new();
|
||||
let source_changed_exhaustions = Arc::new(AtomicUsize::new(0));
|
||||
let incarnation = Some(store.bucket_incarnation_id(&bucket).await.expect("bucket incarnation"));
|
||||
run_decommission_entry_retry_test(
|
||||
&store,
|
||||
rx.clone(),
|
||||
&bucket,
|
||||
object,
|
||||
incarnation,
|
||||
Arc::clone(&source_changed_exhaustions),
|
||||
)
|
||||
.await
|
||||
.expect("entry-level exhaustion must stay local below the pool threshold");
|
||||
|
||||
assert!(!rx.is_cancelled(), "one exhausted entry must not cancel other bucket workers");
|
||||
assert_eq!(mutation_calls.load(Ordering::SeqCst), crate::core::pools::DECOMMISSION_ENTRY_MAX_ATTEMPTS);
|
||||
assert_eq!(source_changed_exhaustions.load(Ordering::SeqCst), 1);
|
||||
store.pools[0]
|
||||
.get_object_info(
|
||||
&bucket,
|
||||
object,
|
||||
&ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(original_version.to_string()),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("retry exhaustion must retain the original source version");
|
||||
let pool_meta = store.pool_meta.read().await;
|
||||
let info = pool_meta.pools[0]
|
||||
.decommission
|
||||
.as_ref()
|
||||
.expect("decommission progress must be initialized");
|
||||
assert_eq!(info.items_decommission_failed, 1, "exhausted entry must be visible as failed");
|
||||
drop(pool_meta);
|
||||
|
||||
shutdown.cancel();
|
||||
});
|
||||
})
|
||||
.expect("spawn decommission retry test thread");
|
||||
if let Err(payload) = handle.join() {
|
||||
std::panic::resume_unwind(payload);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
fn decommission_entry_delete_marker_copy_retries_real_path() {
|
||||
let handle = std::thread::Builder::new()
|
||||
.name("decommission_entry_delete_marker_copy_retries_real_path".to_string())
|
||||
.stack_size(32 * 1024 * 1024)
|
||||
.spawn(|| {
|
||||
let runtime = tokio::runtime::Builder::new_multi_thread()
|
||||
.enable_all()
|
||||
.worker_threads(2)
|
||||
.build()
|
||||
.expect("test runtime should build");
|
||||
runtime.block_on(async {
|
||||
let temp_dir = tempfile::tempdir().expect("create delete marker retry store dir");
|
||||
let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store(
|
||||
temp_dir.path(),
|
||||
"decommission-delete-marker-retry",
|
||||
&[4, 4],
|
||||
))
|
||||
.await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
|
||||
let bucket = format!("decom-marker-{}", uuid::Uuid::new_v4());
|
||||
let object = "marker.bin";
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create delete marker retry bucket");
|
||||
let data_version = uuid::Uuid::new_v4();
|
||||
let marker_version = uuid::Uuid::new_v4();
|
||||
let base_time = OffsetDateTime::now_utc();
|
||||
seed_decommission_source(
|
||||
&store,
|
||||
&bucket,
|
||||
object,
|
||||
b"delete marker data".to_vec(),
|
||||
&ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(data_version.to_string()),
|
||||
mod_time: Some(base_time),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await;
|
||||
store.pools[0]
|
||||
.delete_object(
|
||||
&bucket,
|
||||
object,
|
||||
ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(marker_version.to_string()),
|
||||
delete_marker: true,
|
||||
mod_time: Some(base_time + time::Duration::seconds(1)),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("seed source delete marker");
|
||||
mark_test_pool_decommissioning(&store, 0).await;
|
||||
|
||||
let fault_calls = Arc::new(AtomicUsize::new(0));
|
||||
let fault_calls_for_hook = Arc::clone(&fault_calls);
|
||||
let fault_bucket = bucket.clone();
|
||||
let _fault_guard = crate::core::pools::DecommissionTestFaultGuard::install(Arc::new(
|
||||
move |stage, called_bucket, called_object, attempt| {
|
||||
let injected = stage == DECOMMISSION_TEST_FAULT_STAGE_DELETE_MARKER
|
||||
&& called_bucket == fault_bucket.as_str()
|
||||
&& called_object == object
|
||||
&& attempt < crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS;
|
||||
if injected {
|
||||
fault_calls_for_hook.fetch_add(1, Ordering::SeqCst);
|
||||
}
|
||||
injected
|
||||
},
|
||||
));
|
||||
|
||||
let incarnation = Some(store.bucket_incarnation_id(&bucket).await.expect("bucket incarnation"));
|
||||
run_decommission_entry_retry_test(
|
||||
&store,
|
||||
CancellationToken::new(),
|
||||
&bucket,
|
||||
object,
|
||||
incarnation,
|
||||
Arc::new(AtomicUsize::new(0)),
|
||||
)
|
||||
.await
|
||||
.expect("delete marker copy retries must converge");
|
||||
|
||||
assert_eq!(
|
||||
fault_calls.load(Ordering::SeqCst),
|
||||
crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS - 1
|
||||
);
|
||||
let marker_opts = ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(marker_version.to_string()),
|
||||
..Default::default()
|
||||
};
|
||||
let target_marker = store.pools[1]
|
||||
.get_object_info(&bucket, object, &marker_opts)
|
||||
.await
|
||||
.expect("target delete marker must exist");
|
||||
assert!(target_marker.delete_marker);
|
||||
assert_decommission_source_absent(&store, &bucket, object, &marker_opts).await;
|
||||
|
||||
let data_opts = ObjectOptions {
|
||||
versioned: true,
|
||||
version_id: Some(data_version.to_string()),
|
||||
..Default::default()
|
||||
};
|
||||
assert_eq!(
|
||||
read_decommission_target_body(&store, &bucket, object, &data_opts).await,
|
||||
b"delete marker data"
|
||||
);
|
||||
assert_decommission_source_absent(&store, &bucket, object, &data_opts).await;
|
||||
|
||||
shutdown.cancel();
|
||||
});
|
||||
})
|
||||
.expect("spawn decommission retry test thread");
|
||||
if let Err(payload) = handle.join() {
|
||||
std::panic::resume_unwind(payload);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
fn decommission_entry_tiered_copy_retries_real_path() {
|
||||
let handle = std::thread::Builder::new()
|
||||
.name("decommission_entry_tiered_copy_retries_real_path".to_string())
|
||||
.stack_size(32 * 1024 * 1024)
|
||||
.spawn(|| {
|
||||
let runtime = tokio::runtime::Builder::new_multi_thread()
|
||||
.enable_all()
|
||||
.worker_threads(2)
|
||||
.build()
|
||||
.expect("test runtime should build");
|
||||
runtime.block_on(async {
|
||||
let temp_dir = tempfile::tempdir().expect("create tiered retry store dir");
|
||||
let (ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store(
|
||||
temp_dir.path(),
|
||||
"decommission-tiered-retry",
|
||||
&[4, 4],
|
||||
))
|
||||
.await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
|
||||
let bucket = format!("decom-tiered-{}", uuid::Uuid::new_v4());
|
||||
let object = "tiered.bin";
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create tiered retry bucket");
|
||||
let mut reader = PutObjReader::from_vec(b"tiered generation".to_vec());
|
||||
let original = store.pools[0]
|
||||
.put_object(&bucket, object, &mut reader, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("seed tiered source object");
|
||||
let tier_name = format!("DECOM{}", &uuid::Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
|
||||
register_mock_tier(&ctx.tier_config_mgr(), &tier_name).await;
|
||||
store.pools[0]
|
||||
.transition_object(
|
||||
&bucket,
|
||||
object,
|
||||
&ObjectOptions {
|
||||
transition: TransitionOptions {
|
||||
status: TRANSITION_PENDING.to_string(),
|
||||
tier: tier_name,
|
||||
etag: original.etag.clone().expect("tiered source ETag"),
|
||||
..Default::default()
|
||||
},
|
||||
version_id: original.version_id.map(|version_id| version_id.to_string()),
|
||||
mod_time: original.mod_time,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("transition source object to mock tier");
|
||||
mark_test_pool_decommissioning(&store, 0).await;
|
||||
|
||||
let fault_calls = Arc::new(AtomicUsize::new(0));
|
||||
let fault_calls_for_hook = Arc::clone(&fault_calls);
|
||||
let fault_bucket = bucket.clone();
|
||||
let _fault_guard = crate::core::pools::DecommissionTestFaultGuard::install(Arc::new(
|
||||
move |stage, called_bucket, called_object, attempt| {
|
||||
let injected = stage == DECOMMISSION_TEST_FAULT_STAGE_TIERED
|
||||
&& called_bucket == fault_bucket.as_str()
|
||||
&& called_object == object
|
||||
&& attempt < crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS;
|
||||
if injected {
|
||||
fault_calls_for_hook.fetch_add(1, Ordering::SeqCst);
|
||||
}
|
||||
injected
|
||||
},
|
||||
));
|
||||
|
||||
let incarnation = Some(store.bucket_incarnation_id(&bucket).await.expect("bucket incarnation"));
|
||||
run_decommission_entry_retry_test(
|
||||
&store,
|
||||
CancellationToken::new(),
|
||||
&bucket,
|
||||
object,
|
||||
incarnation,
|
||||
Arc::new(AtomicUsize::new(0)),
|
||||
)
|
||||
.await
|
||||
.expect("tiered copy retries must converge");
|
||||
|
||||
assert_eq!(
|
||||
fault_calls.load(Ordering::SeqCst),
|
||||
crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS - 1
|
||||
);
|
||||
let target = store.pools[1]
|
||||
.get_object_info(&bucket, object, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("tiered target metadata must exist");
|
||||
assert_eq!(target.transitioned_object.status, rustfs_filemeta::TRANSITION_COMPLETE);
|
||||
assert_decommission_source_absent(&store, &bucket, object, &ObjectOptions::default()).await;
|
||||
|
||||
shutdown.cancel();
|
||||
});
|
||||
})
|
||||
.expect("spawn decommission retry test thread");
|
||||
if let Err(payload) = handle.join() {
|
||||
std::panic::resume_unwind(payload);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn decommission_outer_fence_loss_blocks_target_put_commit() {
|
||||
|
||||
@@ -231,6 +231,10 @@ impl HealTask {
|
||||
"Heal erasure set format repair skipped because no format heal was required"
|
||||
);
|
||||
} else {
|
||||
let error = e;
|
||||
if error.is_recoverable_heal() {
|
||||
return Err(error);
|
||||
}
|
||||
error!(
|
||||
target: "rustfs::heal::task",
|
||||
event = EVENT_HEAL_ERASURE_SET_RESULT,
|
||||
@@ -239,7 +243,7 @@ impl HealTask {
|
||||
task_id = %self.id,
|
||||
set_disk_id,
|
||||
result = "format_failed",
|
||||
error = %e,
|
||||
error = %error,
|
||||
"Heal erasure set failed"
|
||||
);
|
||||
{
|
||||
@@ -247,7 +251,7 @@ impl HealTask {
|
||||
progress.update_progress(4, 4, 0, 0);
|
||||
}
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Failed to heal disk format for {set_disk_id}: {e}"),
|
||||
message: format!("Failed to heal disk format for {set_disk_id}: {error}"),
|
||||
});
|
||||
}
|
||||
} else {
|
||||
@@ -284,6 +288,9 @@ impl HealTask {
|
||||
Err(Error::TaskCancelled) => return Err(Error::TaskCancelled),
|
||||
Err(Error::TaskTimeout) => return Err(Error::TaskTimeout),
|
||||
Err(e) => {
|
||||
if e.is_recoverable_heal() {
|
||||
return Err(e);
|
||||
}
|
||||
error!(
|
||||
target: "rustfs::heal::task",
|
||||
event = EVENT_HEAL_ERASURE_SET_RESULT,
|
||||
|
||||
@@ -547,6 +547,7 @@ struct MockStorage {
|
||||
heal_object_outcome: Mutex<Option<MockHealObjectOutcome>>,
|
||||
heal_object_outcomes: Mutex<HashMap<String, VecDeque<MockHealObjectOutcome>>>,
|
||||
format_no_heal_required: Mutex<bool>,
|
||||
format_error: Mutex<Option<Error>>,
|
||||
global_format_calls: Mutex<u32>,
|
||||
replacement_format_calls: Mutex<Vec<(usize, usize, Vec<String>)>>,
|
||||
replacement_targets_ready: Mutex<bool>,
|
||||
@@ -867,6 +868,9 @@ impl HealStorageAPI for MockStorage {
|
||||
|
||||
async fn heal_format(&self, _dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
|
||||
*self.global_format_calls.lock().unwrap() += 1;
|
||||
if let Some(error) = self.format_error.lock().unwrap().take() {
|
||||
return Err(error);
|
||||
}
|
||||
let no_heal_required = *self.format_no_heal_required.lock().unwrap();
|
||||
if no_heal_required {
|
||||
Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::NoHealRequired))))
|
||||
@@ -2052,6 +2056,30 @@ async fn test_erasure_set_heal_continues_after_format_no_heal_required() {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn erasure_set_format_slowdown_is_propagated() {
|
||||
let storage = Arc::new(MockStorage {
|
||||
format_error: Mutex::new(Some(Error::Storage(EcstoreError::SlowDown))),
|
||||
..Default::default()
|
||||
});
|
||||
let request = HealRequest::new(
|
||||
HealType::ErasureSet {
|
||||
buckets: Vec::new(),
|
||||
set_disk_id: "pool_0_set_0".to_string(),
|
||||
},
|
||||
HealOptions::default(),
|
||||
HealPriority::Normal,
|
||||
);
|
||||
let task = HealTask::from_request(request, storage);
|
||||
|
||||
let error = task
|
||||
.execute()
|
||||
.await
|
||||
.expect_err("format SlowDown must remain recoverable for the task manager");
|
||||
|
||||
assert!(matches!(error, Error::Storage(EcstoreError::SlowDown)));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn erasure_set_bucket_prepass_failure_stops_before_object_heal() {
|
||||
let temp = TempDir::new().expect("temporary directory should be created");
|
||||
|
||||
@@ -245,6 +245,18 @@ impl TestECStoreEnvBuilder {
|
||||
.await
|
||||
.expect("build test ECStore");
|
||||
|
||||
// The production bootstrap only persists pool.bin from the elected
|
||||
// first cluster node. Test stores intentionally have no cluster
|
||||
// election, but heal-format still requires that durable fence before
|
||||
// it can write any disk format. Materialize the validated topology
|
||||
// here so the shared fixture models a ready single-node store.
|
||||
let mut pool_meta = ecstore.pool_meta.read().await.clone();
|
||||
pool_meta.dont_save = false;
|
||||
pool_meta
|
||||
.save(ecstore.pools.clone())
|
||||
.await
|
||||
.expect("persist test pool metadata");
|
||||
|
||||
if self.init_bucket_metadata {
|
||||
let buckets_list = ecstore
|
||||
.list_bucket(&BucketOptions {
|
||||
|
||||
+2
-2
@@ -322,7 +322,7 @@ thiserror = { workspace = true }
|
||||
tracing.workspace = true
|
||||
url = { workspace = true }
|
||||
urlencoding = { workspace = true }
|
||||
uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] }
|
||||
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
|
||||
zip = { workspace = true }
|
||||
libc = { workspace = true }
|
||||
rand = { workspace = true, features = ["serde"] }
|
||||
@@ -345,7 +345,7 @@ libsystemd.workspace = true
|
||||
libmimalloc-sys.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] }
|
||||
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
|
||||
serial_test = { workspace = true }
|
||||
tempfile = { workspace = true }
|
||||
aws-config = { workspace = true }
|
||||
|
||||
@@ -41,7 +41,7 @@ use crate::admin::storage_api::config::save_admin_config;
|
||||
use crate::admin::storage_api::contract::bucket::{
|
||||
BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp,
|
||||
};
|
||||
use crate::admin::storage_api::error::Error as StorageError;
|
||||
use crate::admin::storage_api::error::{Error as StorageError, is_err_bucket_not_found};
|
||||
use crate::admin::storage_api::runtime::ECStore;
|
||||
use crate::admin::utils::{encode_compatible_admin_payload, read_compatible_admin_body};
|
||||
use crate::auth::constant_time_eq;
|
||||
@@ -55,6 +55,7 @@ use crate::storage::storage_api::{
|
||||
use base64::Engine;
|
||||
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
|
||||
use base64::engine::general_purpose::URL_SAFE_NO_PAD;
|
||||
use futures::StreamExt;
|
||||
use hmac::{Hmac, Mac};
|
||||
use http::header::{CONTENT_TYPE, HOST};
|
||||
use http::{HeaderMap, HeaderValue, Uri};
|
||||
@@ -2096,6 +2097,18 @@ async fn remote_add_preflight_info(site: &PeerSite) -> S3Result<SiteReplicationA
|
||||
format!("invalid site replication metainfo from `{}`: {e}", site.endpoint),
|
||||
)
|
||||
})?;
|
||||
if info.deployment_id.is_empty() {
|
||||
// The peer will be tracked under a locally derived fallback ID
|
||||
// (deployment_id_for_endpoint) instead of its real deployment ID.
|
||||
warn!(
|
||||
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
|
||||
component = LOG_COMPONENT_ADMIN,
|
||||
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
|
||||
result = "peer_deployment_id_missing",
|
||||
peer_endpoint = %site.endpoint,
|
||||
"admin site replication state"
|
||||
);
|
||||
}
|
||||
|
||||
let idp_body = send_peer_admin_get_request_with_client(
|
||||
&client,
|
||||
@@ -2206,20 +2219,30 @@ fn site_replication_bootstrap_token(uri: &Uri) -> Option<String> {
|
||||
query_pairs(uri).get("bootstrapToken").cloned()
|
||||
}
|
||||
|
||||
fn bootstrap_bucket_make_op_path(bucket: &SRBucketInfo) -> String {
|
||||
/// Query for a peer `make-with-versioning` bucket op. `versioningEnabled`
|
||||
/// always travels so the outbound query matches MinIO's site-replication
|
||||
/// make-bucket wire contract: MinIO's own create-bucket hook sends
|
||||
/// `versioningEnabled=true` on this op. RustFS's inbound handler
|
||||
/// force-enables versioning either way.
|
||||
fn make_with_versioning_bucket_op_path(bucket: &str, created_at: Option<&str>, lock_enabled: bool) -> String {
|
||||
let mut query = form_urlencoded::Serializer::new(String::new());
|
||||
query.append_pair("bucket", &bucket.bucket);
|
||||
query.append_pair("operation", "make-with-versioning");
|
||||
if let Some(created_at) = bucket
|
||||
.created_at
|
||||
.and_then(|value| value.format(&time::format_description::well_known::Rfc3339).ok())
|
||||
{
|
||||
query.append_pair("createdAt", &created_at);
|
||||
query.append_pair("bucket", bucket);
|
||||
query.append_pair("operation", SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING);
|
||||
query.append_pair("versioningEnabled", "true");
|
||||
if let Some(created_at) = created_at {
|
||||
query.append_pair("createdAt", created_at);
|
||||
}
|
||||
if bucket.object_lock_config.is_some() {
|
||||
if lock_enabled {
|
||||
query.append_pair("lockEnabled", "true");
|
||||
}
|
||||
format!("/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", query.finish())
|
||||
format!("{SITE_REPLICATION_PEER_BUCKET_OPS_PATH}?{}", query.finish())
|
||||
}
|
||||
|
||||
fn bootstrap_bucket_make_op_path(bucket: &SRBucketInfo) -> String {
|
||||
let created_at = bucket
|
||||
.created_at
|
||||
.and_then(|value| value.format(&time::format_description::well_known::Rfc3339).ok());
|
||||
make_with_versioning_bucket_op_path(&bucket.bucket, created_at.as_deref(), bucket.object_lock_config.is_some())
|
||||
}
|
||||
|
||||
fn bootstrap_bucket_meta_item(bucket: &SRBucketInfo, item_type: &str, updated_at: Option<OffsetDateTime>) -> SRBucketMeta {
|
||||
@@ -4246,16 +4269,7 @@ async fn broadcast_site_replication_make_bucket(
|
||||
.format(&time::format_description::well_known::Rfc3339)
|
||||
.unwrap_or_default();
|
||||
|
||||
let path = {
|
||||
let mut query = form_urlencoded::Serializer::new(String::new());
|
||||
query.append_pair("bucket", bucket);
|
||||
query.append_pair("operation", "make-with-versioning");
|
||||
query.append_pair("createdAt", &created_at);
|
||||
if lock_enabled {
|
||||
query.append_pair("lockEnabled", "true");
|
||||
}
|
||||
format!("/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", query.finish())
|
||||
};
|
||||
let path = make_with_versioning_bucket_op_path(bucket, Some(&created_at), lock_enabled);
|
||||
let path = if let Some(token) = bootstrap_token {
|
||||
with_site_replication_bootstrap_token(&path, token)
|
||||
} else {
|
||||
@@ -10206,13 +10220,25 @@ impl Operation for SiteReplicationStatusHandler {
|
||||
}
|
||||
}
|
||||
|
||||
/// `POST /v3/site-replication/devnull` — peer link-check upload drain.
|
||||
/// MinIO streams multi-megabyte probe bodies here during site netperf link
|
||||
/// checks and expects an unbounded discard (its handler copies to io.Discard);
|
||||
/// buffering through the 1MB admin body cap turned any larger probe into a
|
||||
/// 400 and a false link failure. Stream and discard instead — no size cap.
|
||||
async fn drain_site_replication_devnull(mut input: Body) -> S3Result<()> {
|
||||
while let Some(chunk) = input.next().await {
|
||||
chunk.map_err(|e| s3_error!(InvalidRequest, "failed to read devnull stream: {}", e))?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub struct SiteReplicationDevNullHandler {}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Operation for SiteReplicationDevNullHandler {
|
||||
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?;
|
||||
let _ = read_plain_admin_body(req.input).await?;
|
||||
drain_site_replication_devnull(req.input).await?;
|
||||
Ok(empty_response(StatusCode::NO_CONTENT))
|
||||
}
|
||||
}
|
||||
@@ -10471,6 +10497,19 @@ impl Operation for SRPeerJoinHandler {
|
||||
}
|
||||
}
|
||||
|
||||
/// Outcome of a peer-driven `purge-deleted-bucket` replay. A bucket that is
|
||||
/// already gone means the purge raced an earlier replay or a local delete —
|
||||
/// that is success — but any other failure must reach the sender like the
|
||||
/// sibling delete branches do: swallowing it answered 200 while the bucket
|
||||
/// survived on this site.
|
||||
fn purge_deleted_bucket_result(result: Result<(), StorageError>) -> S3Result<()> {
|
||||
match result {
|
||||
Ok(()) => Ok(()),
|
||||
Err(err) if is_err_bucket_not_found(&err) => Ok(()),
|
||||
Err(err) => Err(ApiError::from(err).into()),
|
||||
}
|
||||
}
|
||||
|
||||
pub struct SRPeerBucketOpsHandler {}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
@@ -10570,16 +10609,18 @@ impl Operation for SRPeerBucketOpsHandler {
|
||||
.map_err(ApiError::from)?;
|
||||
}
|
||||
"purge-deleted-bucket" => {
|
||||
let _ = store
|
||||
.delete_bucket(
|
||||
&bucket,
|
||||
&DeleteBucketOptions {
|
||||
force: true,
|
||||
srdelete_op: SRBucketDeleteOp::Purge,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await;
|
||||
purge_deleted_bucket_result(
|
||||
store
|
||||
.delete_bucket(
|
||||
&bucket,
|
||||
&DeleteBucketOptions {
|
||||
force: true,
|
||||
srdelete_op: SRBucketDeleteOp::Purge,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await,
|
||||
)?;
|
||||
}
|
||||
_ => return Err(s3_error!(InvalidRequest, "unsupported site replication bucket operation")),
|
||||
}
|
||||
@@ -13925,6 +13966,54 @@ mod tests {
|
||||
assert!(!query_flag(&uri, "missing"));
|
||||
}
|
||||
|
||||
/// A5 red-light: a `purge-deleted-bucket` replay must report success when
|
||||
/// the bucket is already gone, and must propagate every other failure —
|
||||
/// the swallowed error answered 200 while the bucket survived.
|
||||
#[test]
|
||||
fn test_purge_deleted_bucket_result_tolerates_only_missing_bucket() {
|
||||
assert!(purge_deleted_bucket_result(Ok(())).is_ok());
|
||||
assert!(purge_deleted_bucket_result(Err(StorageError::BucketNotFound("photos".to_string()))).is_ok());
|
||||
assert!(purge_deleted_bucket_result(Err(StorageError::VolumeNotFound)).is_ok());
|
||||
let err = purge_deleted_bucket_result(Err(StorageError::StorageFull))
|
||||
.expect_err("non-not-found delete failures must propagate");
|
||||
assert_ne!(*err.code(), S3ErrorCode::NoSuchBucket);
|
||||
}
|
||||
|
||||
/// C5 red-light: the site-replication devnull drain must accept bodies
|
||||
/// beyond the 1MB admin body cap — MinIO's link check streams large
|
||||
/// probe bodies and treats a 400 as a broken link.
|
||||
#[tokio::test]
|
||||
async fn test_site_replication_devnull_drains_body_beyond_admin_cap() {
|
||||
let body = Body::from(vec![0u8; MAX_ADMIN_REQUEST_BODY_SIZE + 1]);
|
||||
drain_site_replication_devnull(body)
|
||||
.await
|
||||
.expect("devnull must drain bodies larger than the admin body cap");
|
||||
}
|
||||
|
||||
/// A3 red-light: `versioningEnabled` must travel on every outbound
|
||||
/// make-with-versioning bucket op so the query matches MinIO's
|
||||
/// site-replication make-bucket wire contract (MinIO's own hook sends
|
||||
/// `versioningEnabled=true` on this op).
|
||||
#[test]
|
||||
fn test_make_with_versioning_op_paths_send_versioning_enabled() {
|
||||
let bucket = SRBucketInfo {
|
||||
bucket: "photos".to_string(),
|
||||
created_at: Some(OffsetDateTime::UNIX_EPOCH),
|
||||
object_lock_config: Some(BASE64_STANDARD.encode("<ObjectLockConfiguration/>")),
|
||||
..Default::default()
|
||||
};
|
||||
let bootstrap = bootstrap_bucket_make_op_path(&bucket);
|
||||
assert!(bootstrap.contains("operation=make-with-versioning"), "{bootstrap}");
|
||||
assert!(bootstrap.contains("versioningEnabled=true"), "{bootstrap}");
|
||||
assert!(bootstrap.contains("createdAt="), "{bootstrap}");
|
||||
assert!(bootstrap.contains("lockEnabled=true"), "{bootstrap}");
|
||||
|
||||
// The broadcast path (create-bucket hook) shares the same builder.
|
||||
let broadcast = make_with_versioning_bucket_op_path("photos", Some("1970-01-01T00:00:00Z"), false);
|
||||
assert!(broadcast.contains("versioningEnabled=true"), "{broadcast}");
|
||||
assert!(!broadcast.contains("lockEnabled"), "{broadcast}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_add_bootstrap_scope_only_allows_expected_bucket_setup_until_guard_drops() {
|
||||
|
||||
@@ -13,9 +13,9 @@
|
||||
// limitations under the License.
|
||||
|
||||
use rustfs_madmin::{PeerInfo, SyncStatus};
|
||||
use std::collections::{BTreeMap, hash_map::DefaultHasher};
|
||||
use std::hash::{Hash, Hasher};
|
||||
use std::collections::BTreeMap;
|
||||
use url::Url;
|
||||
use uuid::Uuid;
|
||||
|
||||
fn has_http_scheme(endpoint: &str) -> bool {
|
||||
endpoint.get(..7).is_some_and(|prefix| prefix.eq_ignore_ascii_case("http://"))
|
||||
@@ -66,10 +66,12 @@ pub fn site_identity_key(endpoint: &str) -> String {
|
||||
.unwrap_or_else(|| trimmed.to_ascii_lowercase())
|
||||
}
|
||||
|
||||
/// Fallback deployment ID for a peer that reported none. UUIDv5 over the
|
||||
/// canonical endpoint: the ID is persisted in site-replication state and
|
||||
/// broadcast to peers, so it must be identical across Rust toolchains
|
||||
/// (`DefaultHasher` is not) and across spellings of the same endpoint.
|
||||
pub fn deployment_id_for_endpoint(endpoint: &str) -> String {
|
||||
let mut hasher = DefaultHasher::new();
|
||||
endpoint.hash(&mut hasher);
|
||||
format!("{:016x}", hasher.finish())
|
||||
Uuid::new_v5(&Uuid::NAMESPACE_URL, canonical_endpoint(endpoint).as_bytes()).to_string()
|
||||
}
|
||||
|
||||
pub fn same_identity_endpoint(left: &str, right: &str) -> bool {
|
||||
@@ -174,6 +176,23 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
/// B8 red-light: the fallback deployment ID must be a toolchain-stable
|
||||
/// UUIDv5 over the canonical endpoint — `DefaultHasher` output is not
|
||||
/// guaranteed stable across Rust releases, yet the ID is persisted in
|
||||
/// site-replication state and broadcast to peers.
|
||||
#[test]
|
||||
fn deployment_id_for_endpoint_is_stable_uuid_v5_over_canonical_endpoint() {
|
||||
let endpoint = "https://node-a.example.com:9000";
|
||||
let id = deployment_id_for_endpoint(endpoint);
|
||||
let parsed = uuid::Uuid::parse_str(&id).expect("fallback deployment ID must be a UUID");
|
||||
assert_eq!(parsed.get_version_num(), 5, "fallback deployment ID must be UUIDv5");
|
||||
// Deterministic for the same endpoint and for spelling variants that
|
||||
// share a canonical form; distinct endpoints stay distinct.
|
||||
assert_eq!(id, deployment_id_for_endpoint(endpoint));
|
||||
assert_eq!(id, deployment_id_for_endpoint(" HTTPS://Node-A.Example.Com:9000/ "));
|
||||
assert_ne!(id, deployment_id_for_endpoint("https://node-b.example.com:9000"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn canonical_endpoint_accepts_case_insensitive_scheme() {
|
||||
assert_eq!(
|
||||
|
||||
@@ -51,7 +51,7 @@ mod ecstore_disk {
|
||||
}
|
||||
|
||||
mod ecstore_error {
|
||||
pub(crate) use crate::storage::storage_api::ecstore_error::StorageError;
|
||||
pub(crate) use crate::storage::storage_api::ecstore_error::{StorageError, is_err_bucket_not_found};
|
||||
}
|
||||
|
||||
#[allow(unused_imports)]
|
||||
@@ -919,6 +919,7 @@ pub(crate) mod contract {
|
||||
}
|
||||
|
||||
pub(crate) mod error {
|
||||
pub(crate) use super::ecstore_error::is_err_bucket_not_found;
|
||||
pub(crate) use super::{Error, StorageError};
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user