Compare commits

..

31 Commits

Author SHA1 Message Date
houseme 894e342b48 Merge branch 'main' into overtrue/fix-1905-activation-fence 2026-08-23 12:36:14 +08:00
Zhengchao An 23a2c7d776 test(kms): stabilize Vault failover validation (#6385)
* test(kms): bound Vault failover progress wait

* ci(nightly): honor manual dispatch ref

* test(kms): preserve Vault worker failures

* test(kms): validate Vault circuit recovery
2026-08-23 12:32:11 +08:00
overtrue 6f115ea2ec fix(ecstore): resolve CI clippy failures 2026-08-23 05:45:10 +08:00
overtrue 4d0ffa680a Merge remote-tracking branch 'origin/main' into overtrue/fix-1905-activation-fence
# Conflicts:
#	crates/ecstore/src/core/pools.rs
2026-08-23 02:36:20 +08:00
overtrue d5648b52a7 fix(ecstore): remove duplicate activation test import 2026-08-23 01:09:40 +08:00
overtrue e05fe7c014 Merge origin/main into overtrue/fix-1905-activation-fence 2026-08-23 00:52:29 +08:00
overtrue d926642713 Merge origin/main into overtrue/fix-1905-activation-fence 2026-08-23 00:23:44 +08:00
overtrue ac57641e9b fix(ecstore): align activation fence test imports 2026-08-23 00:22:53 +08:00
overtrue f5f5212abd Merge origin/main into overtrue/fix-1905-activation-fence 2026-08-22 23:58:53 +08:00
overtrue 07bac4978e test(ecstore): observe decommission lock attempt 2026-08-22 12:37:47 +08:00
overtrue d908f00f01 test(ecstore): scope rebalance disk trait import 2026-08-22 12:34:25 +08:00
overtrue b8e9170cb0 test(ecstore): fix activation fence synchronization 2026-08-22 12:34:25 +08:00
overtrue 4fdc1c3a85 fix(ecstore): repair rebalance entry runtime failures 2026-08-22 12:34:25 +08:00
overtrue 89e910df3a fix(ecstore): repair rebalance test imports 2026-08-22 12:34:25 +08:00
overtrue 7b12d621a7 fix(rebalance): make prepared stop terminal-safe 2026-08-22 12:34:25 +08:00
overtrue 04bafeb662 fix(rebalance): preserve committed activation recovery 2026-08-22 12:34:25 +08:00
overtrue afbb842dbf fix(ecstore): adopt activations after durable commit 2026-08-22 12:34:25 +08:00
overtrue 8d97f3570d fix(ecstore): fence multipart staging on rebalance lock loss 2026-08-22 12:34:25 +08:00
overtrue 4cce1b2e3b fix(rebalance): cancel admin stop before activation wait 2026-08-22 12:34:25 +08:00
overtrue 8b0e96c314 test(ecstore): exercise real rebalance fences 2026-08-22 12:34:25 +08:00
overtrue 3677482f9f fix(ecstore): fence rebalance commits and unblock stop 2026-08-22 12:34:24 +08:00
overtrue ce3cd2d890 fix(ecstore): commit rebalance activation after persistence 2026-08-22 12:34:24 +08:00
overtrue c20fe73d6f fix(ecstore): fence stale rebalance workers 2026-08-22 12:34:24 +08:00
overtrue dc0c64689b fix(ecstore): satisfy rebalance activation clippy checks 2026-08-22 12:34:24 +08:00
overtrue 1a59f8ef6b test(ecstore): reuse rebalance metadata fixture 2026-08-22 12:34:24 +08:00
overtrue a69fb0882d fix(ecstore): repair rebalance fence test wiring 2026-08-22 12:34:24 +08:00
overtrue 209d3481da test(ecstore): exercise lost rebalance commit fence 2026-08-22 12:34:24 +08:00
overtrue d4b297186f fix(ecstore): close rebalance activation races 2026-08-22 12:34:24 +08:00
overtrue 814e17b02a fix(ecstore): bind rebalance workers to activation id 2026-08-22 12:34:24 +08:00
overtrue a72deafc9f fix(ecstore): fence lost activation locks 2026-08-22 12:34:24 +08:00
overtrue bb2ac2758f fix(ecstore): fence rebalance and decommission activation 2026-08-22 12:34:24 +08:00
47 changed files with 4299 additions and 2207 deletions
+3 -6
View File
@@ -39,11 +39,10 @@ jobs:
env:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
steps:
- name: Checkout main branch
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
@@ -89,11 +88,10 @@ jobs:
# either casing.
NO_PROXY: 127.0.0.1,localhost
steps:
- name: Checkout main branch
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
@@ -178,11 +176,10 @@ jobs:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
NO_PROXY: 127.0.0.1,localhost
steps:
- name: Checkout main branch
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
+6
View File
@@ -440,6 +440,12 @@ pub mod rebalance {
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, decode_rebalance_stop_propagation_record,
encode_rebalance_stop_propagation_record,
};
#[cfg(feature = "test-util")]
pub mod test_util {
pub use crate::services::rebalance::PausedRebalanceEntryTestFixture;
pub use crate::services::rebalance::test_store_with_persisted_rebalance_meta;
}
}
pub mod rio {
@@ -4324,6 +4324,16 @@ pub async fn expire_transitioned_object(
lc_event: &lifecycle::Event,
_src: &LcEventSrc,
bucket_incarnation_id: Uuid,
) -> Result<ObjectInfo, std::io::Error> {
expire_transitioned_object_with_lock_lost_signal(api, oi, lc_event, bucket_incarnation_id, None).await
}
async fn expire_transitioned_object_with_lock_lost_signal(
api: Arc<ECStore>,
oi: &ObjectInfo,
lc_event: &lifecycle::Event,
bucket_incarnation_id: Uuid,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> Result<ObjectInfo, std::io::Error> {
let publication_guard = lifecycle_expiry_publication_guard(&api, oi, bucket_incarnation_id)
.await
@@ -4335,6 +4345,9 @@ pub async fn expire_transitioned_object(
let mut opts = transitioned_object_delete_opts(oi, lc_event.action, versioned, version_suspended, bucket_incarnation_id)
.map_err(std::io::Error::other)?;
opts.add_namespace_lock_guard(&publication_guard);
if let Some(signal) = lock_lost_signal {
opts.add_namespace_lock_lost_signal(signal);
}
opts.delete_replication_config_snapshot = Some(Arc::new(snapshot));
//let tags = LcAuditEvent::new(src, lcEvent).Tags();
if lc_event.action.delete_restored() {
@@ -4993,12 +5006,32 @@ pub async fn apply_expiry_on_transitioned_object(
lc_event: &lifecycle::Event,
src: &LcEventSrc,
bucket_incarnation_id: Uuid,
) -> bool {
apply_expiry_on_transitioned_object_with_lock_lost_signal(api, oi, lc_event, src, bucket_incarnation_id, None).await
}
async fn apply_expiry_on_transitioned_object_with_lock_lost_signal(
api: Arc<ECStore>,
oi: &ObjectInfo,
lc_event: &lifecycle::Event,
_src: &LcEventSrc,
bucket_incarnation_id: Uuid,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> bool {
if lc_event.action.delete_all() {
return apply_expiry_on_non_transitioned_objects(api, oi, lc_event, src, bucket_incarnation_id).await;
return apply_expiry_on_non_transitioned_objects_with_lock_lost_signal(
api,
oi,
lc_event,
bucket_incarnation_id,
lock_lost_signal,
)
.await;
}
let time_ilm = Metrics::time_ilm(lc_event.action);
if let Err(_err) = expire_transitioned_object(api, oi, lc_event, src, bucket_incarnation_id).await {
if let Err(_err) =
expire_transitioned_object_with_lock_lost_signal(api, oi, lc_event, bucket_incarnation_id, lock_lost_signal).await
{
return false;
}
time_ilm(1)();
@@ -5012,6 +5045,16 @@ pub async fn apply_expiry_on_non_transitioned_objects(
lc_event: &lifecycle::Event,
_src: &LcEventSrc,
bucket_incarnation_id: Uuid,
) -> bool {
apply_expiry_on_non_transitioned_objects_with_lock_lost_signal(api, oi, lc_event, bucket_incarnation_id, None).await
}
async fn apply_expiry_on_non_transitioned_objects_with_lock_lost_signal(
api: Arc<ECStore>,
oi: &ObjectInfo,
lc_event: &lifecycle::Event,
bucket_incarnation_id: Uuid,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> bool {
let Some(publication_guard) = lifecycle_expiry_publication_guard(&api, oi, bucket_incarnation_id).await else {
return false;
@@ -5042,6 +5085,9 @@ pub async fn apply_expiry_on_non_transitioned_objects(
..Default::default()
};
opts.add_namespace_lock_guard(&publication_guard);
if let Some(signal) = lock_lost_signal {
opts.add_namespace_lock_lost_signal(signal);
}
if lc_event.action.delete_versioned() {
opts.version_id = oi.version_id.map(|v| v.to_string());
@@ -5123,6 +5169,61 @@ async fn enqueue_expiry_rule_with_incarnation(
expiry_state.enqueue_by_days(oi, event, src, bucket_incarnation_id)
}
fn lifecycle_expiry_object_matches(current: &ObjectInfo, expected: &ObjectInfo) -> bool {
current.version_id == expected.version_id
&& current.data_dir == expected.data_dir
&& current.mod_time == expected.mod_time
&& current.etag == expected.etag
&& current.delete_marker == expected.delete_marker
&& current.transitioned_object.name == expected.transitioned_object.name
&& current.transitioned_object.version_id == expected.transitioned_object.version_id
&& current.transitioned_object.tier == expected.transitioned_object.tier
&& current.transitioned_object.status == expected.transitioned_object.status
&& current.restore_expires == expected.restore_expires
}
pub(crate) async fn apply_expiry_rule_for_data_movement(
api: Arc<ECStore>,
event: &lifecycle::Event,
src: &LcEventSrc,
oi: &ObjectInfo,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> bool {
let Ok(_lifecycle_guard) = api.acquire_bucket_lifecycle_read_lock(&oi.bucket).await else {
return false;
};
let Ok(bucket_incarnation_id) = api.bucket_incarnation_id_from_disk(&oi.bucket).await else {
return false;
};
let current = match api
.get_object_info(
&oi.bucket,
&oi.name,
&ObjectOptions {
version_id: oi.version_id.map(|version_id| version_id.to_string()),
versioned: oi.version_id.is_some(),
expected_bucket_incarnation_id: Some(bucket_incarnation_id),
..Default::default()
},
)
.await
{
Ok(current) => current,
Err(_) => return false,
};
if !lifecycle_expiry_object_matches(&current, oi) {
return false;
}
if oi.transitioned_object.status.is_empty() {
apply_expiry_on_non_transitioned_objects_with_lock_lost_signal(api, oi, event, bucket_incarnation_id, lock_lost_signal)
.await
} else {
apply_expiry_on_transitioned_object_with_lock_lost_signal(api, oi, event, src, bucket_incarnation_id, lock_lost_signal)
.await
}
}
pub(crate) async fn apply_expiry_rule_in(api: Arc<ECStore>, event: &lifecycle::Event, src: &LcEventSrc, oi: &ObjectInfo) -> bool {
let Ok(_lifecycle_guard) = api.acquire_bucket_lifecycle_read_lock(&oi.bucket).await else {
return false;
@@ -5146,17 +5247,7 @@ pub(crate) async fn apply_expiry_rule_in(api: Arc<ECStore>, event: &lifecycle::E
Ok(current) => current,
Err(_) => return false,
};
if current.version_id != oi.version_id
|| current.data_dir != oi.data_dir
|| current.mod_time != oi.mod_time
|| current.etag != oi.etag
|| current.delete_marker != oi.delete_marker
|| current.transitioned_object.name != oi.transitioned_object.name
|| current.transitioned_object.version_id != oi.transitioned_object.version_id
|| current.transitioned_object.tier != oi.transitioned_object.tier
|| current.transitioned_object.status != oi.transitioned_object.status
|| current.restore_expires != oi.restore_expires
{
if !lifecycle_expiry_object_matches(&current, oi) {
return false;
}
enqueue_expiry_rule_with_incarnation(event, src, oi, bucket_incarnation_id).await
@@ -1089,9 +1089,7 @@ impl PeerRestClient {
.await?
.max_decoding_message_size(BACKGROUND_HEAL_STATUS_MAX_MESSAGE_SIZE);
let response = match client
.background_heal_status(Request::new(BackgroundHealStatusRequest {
protocol_version: rustfs_protos::BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION,
}))
.background_heal_status(Request::new(BackgroundHealStatusRequest::default()))
.await
{
Ok(response) => response.into_inner(),
+607 -30
View File
@@ -19,8 +19,8 @@ use crate::bucket::{
LifecycleExpiryConfigs,
bucket_lifecycle_audit::LcEventSrc,
bucket_lifecycle_ops::{
LifecycleOps, apply_expiry_on_transitioned_object, apply_expiry_rule_in, eval_action_from_lifecycle,
lifecycle_delete_all_versions_blocked_by_replication,
LifecycleOps, apply_expiry_on_transitioned_object, apply_expiry_rule_for_data_movement, apply_expiry_rule_in,
eval_action_from_lifecycle, lifecycle_delete_all_versions_blocked_by_replication,
},
get_expiry_configs,
lifecycle::IlmAction,
@@ -28,7 +28,9 @@ use crate::bucket::{
metadata_sys,
};
use crate::cache_value::metacache_set::{ListPathRawOptions, list_path_raw};
use crate::config::com::{CONFIG_PREFIX, read_config, read_config_no_lock, save_config, save_config_with_opts};
use crate::config::com::{
CONFIG_PREFIX, read_config, read_config_no_lock, save_config, save_config_with_opts, save_config_with_opts_quiet,
};
use crate::data_movement;
use crate::data_movement::backpressure::{self, DataMovementOperation};
use crate::data_usage::DATA_USAGE_CACHE_NAME;
@@ -48,7 +50,6 @@ use crate::storage_api_contracts::{
admin::StorageAdminApi,
bucket::{BucketOperations, BucketOptions, MakeBucketOptions},
heal::HealOperations as _,
namespace::NamespaceLocking as _,
object::{EcstoreObjectIO, ObjectIO as _, ObjectOperations as _},
};
use crate::{core::sets::Sets, store::ECStore};
@@ -1086,7 +1087,7 @@ fn resolve_start_decommission_pool_meta_reload_result(result: Result<()>) -> Res
resolve_decommission_pool_meta_reload_result(result, "start_decommission")
}
fn decommission_rebalance_meta_lock_error(err: rustfs_lock::LockError) -> Error {
fn activation_rebalance_meta_lock_error(err: rustfs_lock::LockError) -> Error {
match err {
rustfs_lock::LockError::QuorumNotReached { required, achieved } => Error::NamespaceLockQuorumUnavailable {
mode: "write",
@@ -1096,12 +1097,12 @@ fn decommission_rebalance_meta_lock_error(err: rustfs_lock::LockError) -> Error
achieved,
},
other => Error::other(format!(
"failed to acquire rebalance metadata write lock before decommission start on {RUSTFS_META_BUCKET}/{REBAL_META_NAME}: {other}"
"failed to acquire rebalance activation lock on {RUSTFS_META_BUCKET}/{REBAL_META_NAME}: {other}"
)),
}
}
fn decommission_pool_meta_lock_error(err: rustfs_lock::LockError) -> Error {
fn activation_pool_meta_lock_error(err: rustfs_lock::LockError) -> Error {
match err {
rustfs_lock::LockError::QuorumNotReached { required, achieved } => Error::NamespaceLockQuorumUnavailable {
mode: "write",
@@ -1111,11 +1112,227 @@ fn decommission_pool_meta_lock_error(err: rustfs_lock::LockError) -> Error {
achieved,
},
other => Error::other(format!(
"failed to acquire pool metadata write lock before decommission start on {RUSTFS_META_BUCKET}/{POOL_META_NAME}: {other}"
"failed to acquire pool activation lock on {RUSTFS_META_BUCKET}/{POOL_META_NAME}: {other}"
)),
}
}
pub(crate) struct PoolRebalanceActivationFence {
pool_meta_guard: rustfs_lock::NamespaceLockGuard,
rebalance_meta_guard: rustfs_lock::NamespaceLockGuard,
#[cfg(test)]
forced_lost: Arc<AtomicBool>,
}
impl PoolRebalanceActivationFence {
pub(crate) fn ensure_held(&self) -> Result<()> {
#[cfg(test)]
let forced_lost = self.forced_lost.load(Ordering::Acquire);
#[cfg(not(test))]
let forced_lost = false;
if forced_lost || self.pool_meta_guard.is_lock_lost() || self.rebalance_meta_guard.is_lock_lost() {
return Err(Error::other("activation lock lost before metadata commit or worker admission"));
}
Ok(())
}
pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
opts.add_namespace_lock_guard(&self.pool_meta_guard);
opts.add_namespace_lock_guard(&self.rebalance_meta_guard);
}
#[cfg(test)]
fn force_lost_for_test(&self) {
self.forced_lost.store(true, Ordering::Release);
}
}
pub(crate) async fn acquire_pool_rebalance_activation_locks<S>(pool: Arc<S>) -> Result<PoolRebalanceActivationFence>
where
S: crate::storage_api_contracts::namespace::NamespaceLocking<
Error = Error,
NamespaceLock = rustfs_lock::NamespaceLockWrapper,
>,
{
// Activation lock order is always pool.bin -> rebalance.bin.
let pool_meta_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?;
let pool_meta_guard = pool_meta_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(activation_pool_meta_lock_error)?;
let rebalance_meta_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, REBAL_META_NAME).await?;
let rebalance_meta_guard = rebalance_meta_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(activation_rebalance_meta_lock_error)?;
Ok(PoolRebalanceActivationFence {
pool_meta_guard,
rebalance_meta_guard,
#[cfg(test)]
forced_lost: Arc::new(AtomicBool::new(false)),
})
}
#[cfg(test)]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum PoolActivationStartKind {
Rebalance,
Decommission,
}
#[cfg(test)]
struct PoolActivationDurableSaveBarrierState {
pool_key: usize,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
static POOL_ACTIVATION_DURABLE_SAVE_BARRIER: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<PoolActivationDurableSaveBarrierState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
pub(crate) struct PoolActivationDurableSaveBarrier {
state: Arc<PoolActivationDurableSaveBarrierState>,
}
#[cfg(test)]
fn pool_activation_test_pool_key<S>(pool: &Arc<S>) -> usize {
Arc::as_ptr(pool).cast::<()>() as usize
}
#[cfg(test)]
impl PoolActivationDurableSaveBarrier {
pub(crate) fn install<S>(pool: &Arc<S>) -> Self {
let state = Arc::new(PoolActivationDurableSaveBarrierState {
pool_key: pool_activation_test_pool_key(pool),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let mut barrier = POOL_ACTIVATION_DURABLE_SAVE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("pool activation durable save barrier should not be poisoned");
assert!(barrier.is_none(), "pool activation durable save barrier must be unique");
*barrier = Some(Arc::clone(&state));
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(std::time::Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("activation should reach the post-durable-save barrier");
}
pub(crate) fn release_after_fence_loss(&self) {
self.state.release.notify_one();
}
}
#[cfg(test)]
impl Drop for PoolActivationDurableSaveBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut barrier = POOL_ACTIVATION_DURABLE_SAVE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("pool activation durable save barrier should not be poisoned");
if barrier.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*barrier = None;
}
}
}
#[cfg(test)]
pub(crate) async fn pause_pool_activation_after_durable_save<S>(pool: &Arc<S>, fence: &PoolRebalanceActivationFence) {
let pool_key = pool_activation_test_pool_key(pool);
let barrier = {
let mut barrier = POOL_ACTIVATION_DURABLE_SAVE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("pool activation durable save barrier should not be poisoned");
if barrier.as_ref().is_some_and(|state| state.pool_key == pool_key) {
barrier.take()
} else {
None
}
};
if let Some(barrier) = barrier {
barrier.arrived.notify_one();
barrier.release.notified().await;
fence.force_lost_for_test();
}
}
#[cfg(test)]
struct PoolActivationStartProbeState {
kind: PoolActivationStartKind,
attempted: std::sync::atomic::AtomicBool,
notify: tokio::sync::Notify,
}
#[cfg(test)]
static POOL_ACTIVATION_START_PROBES: std::sync::OnceLock<std::sync::Mutex<Vec<Arc<PoolActivationStartProbeState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
pub(crate) struct PoolActivationStartProbe {
state: Arc<PoolActivationStartProbeState>,
}
#[cfg(test)]
impl PoolActivationStartProbe {
pub(crate) fn install(kind: PoolActivationStartKind) -> Self {
let state = Arc::new(PoolActivationStartProbeState {
kind,
attempted: std::sync::atomic::AtomicBool::new(false),
notify: tokio::sync::Notify::new(),
});
POOL_ACTIVATION_START_PROBES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("pool activation start probe should not be poisoned")
.push(Arc::clone(&state));
Self { state }
}
pub(crate) async fn wait_until_attempted(&self) {
while !self.state.attempted.load(Ordering::Acquire) {
self.state.notify.notified().await;
}
}
}
#[cfg(test)]
impl Drop for PoolActivationStartProbe {
fn drop(&mut self) {
let mut probes = POOL_ACTIVATION_START_PROBES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("pool activation start probe should not be poisoned");
probes.retain(|state| !Arc::ptr_eq(state, &self.state));
}
}
#[cfg(test)]
pub(crate) fn observe_pool_activation_start_attempt(kind: PoolActivationStartKind) {
let probes = POOL_ACTIVATION_START_PROBES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("pool activation start probe should not be poisoned")
.iter()
.filter(|state| state.kind == kind)
.cloned()
.collect::<Vec<_>>();
for state in probes {
state.attempted.store(true, Ordering::Release);
state.notify.notify_one();
}
}
fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: PoolMeta) {
*pool_meta = previous_pool_meta;
}
@@ -2100,6 +2317,67 @@ impl PoolMeta {
Ok(())
}
async fn save_no_lock_with_activation_fence<S>(
&self,
pools: Vec<Arc<S>>,
activation_fence: PoolRebalanceActivationFence,
) -> Result<()>
where
S: EcstoreObjectIO,
{
let data = self.encode_config_data()?;
if data.is_empty() {
return Ok(());
}
let mut pools = pools.into_iter();
let Some(canonical_pool) = pools.next() else {
return Ok(());
};
let mut opts = ObjectOptions {
max_parity: true,
no_lock: true,
..Default::default()
};
activation_fence.add_namespace_lock_fence(&mut opts);
activation_fence.ensure_held()?;
#[cfg(test)]
let barrier_pool = canonical_pool.clone();
save_config_with_opts(canonical_pool, POOL_META_NAME, data.clone(), &opts).await?;
#[cfg(test)]
pause_pool_activation_after_durable_save(&barrier_pool, &activation_fence).await;
// Pool zero is canonical. Once its save succeeds, later writes only
// replicate committed state and must not reuse the admission fence.
// Failed replicas remain repairable by the next full pool metadata save.
drop(activation_fence);
for (pool_index, pool) in pools.enumerate() {
if let Err(err) = save_config_with_opts_quiet(
pool,
POOL_META_NAME,
data.clone(),
&ObjectOptions {
max_parity: true,
no_lock: true,
..Default::default()
},
)
.await
{
warn!(
event = EVENT_DECOMMISSION_STATE,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = pool_index + 1,
state = "activation_replica_repair_pending",
error = %err,
"Decommission activation replica repair pending"
);
}
}
Ok(())
}
pub fn decommission_cancel(&mut self, idx: usize) -> bool {
if let Some(stats) = self.pools.get_mut(idx) {
if let Some(d) = &stats.decommission {
@@ -2697,6 +2975,85 @@ fn lifecycle_action_skips_heal_version(action: IlmAction) -> bool {
action.delete()
}
#[cfg(test)]
struct LifecycleDataMovementMutationBarrierState {
bucket: String,
object: String,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
pub(crate) struct LifecycleDataMovementMutationBarrier {
state: Arc<LifecycleDataMovementMutationBarrierState>,
}
#[cfg(test)]
static LIFECYCLE_DATA_MOVEMENT_MUTATION_BARRIER: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<LifecycleDataMovementMutationBarrierState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl LifecycleDataMovementMutationBarrier {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(LifecycleDataMovementMutationBarrierState {
bucket: bucket.to_string(),
object: object.to_string(),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let mut slot = LIFECYCLE_DATA_MOVEMENT_MUTATION_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("lifecycle data movement mutation barrier should not be poisoned");
assert!(slot.is_none(), "lifecycle data movement mutation barrier must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(std::time::Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("lifecycle data movement should reach its mutation boundary");
}
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
#[cfg(test)]
impl Drop for LifecycleDataMovementMutationBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut slot = LIFECYCLE_DATA_MOVEMENT_MUTATION_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("lifecycle data movement mutation barrier should not be poisoned");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
async fn pause_lifecycle_data_movement_mutation(bucket: &str, object: &str, has_run_fence_signal: bool) {
if !has_run_fence_signal {
return;
}
let barrier = LIFECYCLE_DATA_MOVEMENT_MUTATION_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("lifecycle data movement mutation barrier should not be poisoned")
.as_ref()
.filter(|barrier| barrier.bucket == bucket && barrier.object == object)
.cloned();
if let Some(barrier) = barrier {
barrier.arrived.notify_one();
barrier.release.notified().await;
}
}
fn resolve_data_movement_lifecycle_expiry_result(action: IlmAction, apply_actions: bool, applied: bool) -> Result<bool> {
if !apply_actions || applied {
return Ok(true);
@@ -2716,6 +3073,7 @@ pub(crate) async fn should_skip_lifecycle_for_data_movement(
object_lock_config: Option<&ObjectLockConfiguration>,
apply_actions: bool,
event_source: &LcEventSrc,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> Result<bool> {
let Some(lifecycle_config) = lifecycle_config else {
return Ok(false);
@@ -2731,16 +3089,31 @@ pub(crate) async fn should_skip_lifecycle_for_data_movement(
let Ok(bucket_incarnation_id) = store.bucket_incarnation_id_from_disk(bucket).await else {
return Ok(false);
};
let _ =
apply_expiry_on_transitioned_object(store, &object_info, &event, event_source, bucket_incarnation_id).await;
let _ = match lock_lost_signal {
Some(signal) => {
apply_expiry_rule_for_data_movement(store, &event, event_source, &object_info, Some(signal)).await
}
None => {
apply_expiry_on_transitioned_object(store, &object_info, &event, event_source, bucket_incarnation_id)
.await
}
};
}
Ok(false)
}
action if lifecycle_action_removes_data_movement_version(action) => {
#[cfg(test)]
pause_lifecycle_data_movement_mutation(bucket, &version.name, lock_lost_signal.is_some()).await;
if lifecycle_delete_all_versions_blocked_by_replication(store.clone(), bucket, &object_info.name, action).await? {
return Ok(false);
}
let applied = !apply_actions || apply_expiry_rule_in(store, &event, event_source, &object_info).await;
let applied = !apply_actions
|| match lock_lost_signal {
Some(signal) => {
apply_expiry_rule_for_data_movement(store, &event, event_source, &object_info, Some(signal)).await
}
None => apply_expiry_rule_in(store, &event, event_source, &object_info).await,
};
resolve_data_movement_lifecycle_expiry_result(action, apply_actions, applied)
}
_ => Ok(false),
@@ -2877,16 +3250,9 @@ impl ECStore {
.first()
.cloned()
.ok_or_else(|| Error::other("decommission start rebalance metadata load failed: no storage pools available"))?;
let pool_meta_lock = rebalance_pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?;
let _pool_meta_guard = pool_meta_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(decommission_pool_meta_lock_error)?;
let ns_lock = rebalance_pool.new_ns_lock(RUSTFS_META_BUCKET, REBAL_META_NAME).await?;
let _guard = ns_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(decommission_rebalance_meta_lock_error)?;
#[cfg(test)]
observe_pool_activation_start_attempt(PoolActivationStartKind::Decommission);
let activation_fence = acquire_pool_rebalance_activation_locks(rebalance_pool.clone()).await?;
let mut rebalance_meta = RebalanceMeta::new();
match rebalance_meta
@@ -2931,7 +3297,10 @@ impl ECStore {
latest_pool_meta.queue_buckets(idx, decom_buckets.clone());
}
latest_pool_meta.save_no_lock(self.pools.clone()).await?;
activation_fence.ensure_held()?;
latest_pool_meta
.save_no_lock_with_activation_fence(self.pools.clone(), activation_fence)
.await?;
{
let mut pool_meta = self.pool_meta.write().await;
*pool_meta = latest_pool_meta;
@@ -2945,6 +3314,11 @@ impl ECStore {
ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await)
}
async fn ensure_decommission_rebalance_idle_after_refresh_under_start_gate(&self) -> Result<()> {
self.load_rebalance_meta_under_start_gate().await?;
ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await)
}
pub async fn status(&self, idx: usize) -> Result<PoolStatus> {
let space_info = self.get_decommission_pool_space_info(idx).await?;
@@ -3839,6 +4213,7 @@ impl ECStore {
object_lock_config.as_ref(),
true,
&LcEventSrc::Decom,
None,
)
.await
})
@@ -4169,6 +4544,7 @@ impl ECStore {
lifecycle_guard: bucket_incarnation_fence
.as_ref()
.and_then(|guard| guard.namespace_lock_guard()),
namespace_lock_lost_signal: None,
object_mutation_fence: Some(&source_cleanup_mutation_fence),
},
"decommission",
@@ -4985,7 +5361,11 @@ impl ECStore {
validate_start_decommission_request(&indices, self.single_pool())?;
self.ensure_decommission_rebalance_idle_after_refresh().await?;
ensure_decommission_start_local_leader(&self.endpoints(), &indices)?;
#[cfg(test)]
let endpoints = self.instance_endpoints().unwrap_or_else(|| self.endpoints());
#[cfg(not(test))]
let endpoints = self.endpoints();
ensure_decommission_start_local_leader(&endpoints, &indices)?;
for idx in indices.iter().copied() {
ensure_valid_decommission_pool_index(self.pools.len(), idx)?;
@@ -5020,7 +5400,8 @@ impl ECStore {
}
let _start_guard = self.start_gate.lock().await;
self.ensure_decommission_rebalance_idle_after_refresh().await?;
self.ensure_decommission_rebalance_idle_after_refresh_under_start_gate()
.await?;
let all_space_infos = self.get_decommission_all_pool_space_infos().await?;
self.cancel_decommission_routines_and_wait(&indices).await;
@@ -5214,6 +5595,7 @@ impl ECStore {
object_lock_config.as_ref(),
false,
&LcEventSrc::Decom,
None,
)
.await
{
@@ -5294,6 +5676,128 @@ mod tests {
use crate::bucket::replication::{ReplicationState, ReplicationStatusType};
use serde::Serialize;
#[tokio::test]
#[serial_test::serial]
async fn decommission_activation_replicates_commit_after_post_save_fence_loss() {
let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(None).await;
let barrier = PoolActivationDurableSaveBarrier::install(&store.pools[0]);
let start_store = Arc::clone(&store);
let start_task = tokio::spawn(async move {
start_store
.save_current_pool_meta_for_decommission_start(
&[0],
vec![(
0,
PoolSpaceInfo {
free: 50,
total: 100,
used: 50,
},
)],
Vec::new(),
)
.await
});
barrier.wait_until_paused().await;
barrier.release_after_fence_loss();
tokio::time::timeout(std::time::Duration::from_secs(30), start_task)
.await
.expect("decommission activation should finish after its canonical commit")
.expect("decommission activation task should not panic")
.expect("post-commit fence loss must not report the committed activation as failed");
let local = store.pool_meta.read().await;
assert!(pool_meta_has_active_decommission(&local));
drop(local);
for pool in &store.pools {
let mut persisted = PoolMeta::default();
persisted
.load_no_lock(pool.clone())
.await
.expect("every pool should retain readable committed decommission metadata");
assert!(
pool_meta_has_active_decommission(&persisted),
"every pool must adopt the canonical committed activation"
);
}
}
#[tokio::test]
#[serial_test::serial]
async fn decommission_activation_adopts_canonical_commit_after_replica_failure() {
let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(None).await;
let barrier = PoolActivationDurableSaveBarrier::install(&store.pools[0]);
let start_store = Arc::clone(&store);
let start_task = tokio::spawn(async move {
start_store
.save_current_pool_meta_for_decommission_start(
&[0],
vec![(
0,
PoolSpaceInfo {
free: 50,
total: 100,
used: 50,
},
)],
Vec::new(),
)
.await
});
barrier.wait_until_paused().await;
let mut replica_disks = Vec::new();
for set in &store.pools[1].disk_set {
let mut disks = set.disks.write().await;
let saved = std::mem::take(&mut *disks);
*disks = vec![None; saved.len()];
replica_disks.push(saved);
}
barrier.release_after_fence_loss();
tokio::time::timeout(std::time::Duration::from_secs(30), start_task)
.await
.expect("decommission activation should finish after its canonical commit")
.expect("decommission activation task should not panic")
.expect("a replica save failure must not report the committed activation as failed");
for (set, disks) in store.pools[1].disk_set.iter().zip(replica_disks) {
*set.disks.write().await = disks;
}
let local = store.pool_meta.read().await;
assert!(pool_meta_has_active_decommission(&local));
drop(local);
let mut canonical = PoolMeta::default();
canonical
.load_no_lock(store.pools[0].clone())
.await
.expect("the canonical committed decommission metadata should remain readable");
assert!(pool_meta_has_active_decommission(&canonical));
let mut replica = PoolMeta::default();
replica
.load_no_lock(store.pools[1].clone())
.await
.expect("the stale replica metadata should remain readable after disks recover");
assert!(!pool_meta_has_active_decommission(&replica));
let worker_cancel = CancellationToken::new();
store
.spawn_decommission_routines(Arc::clone(&store), worker_cancel.clone(), vec![0])
.await
.expect("the committed activation should admit its decommission worker");
let admitted_cancel = store.decommission_cancelers.read().await[0]
.clone()
.expect("the admitted decommission worker should have a cancellation token");
assert!(!admitted_cancel.is_cancelled());
worker_cancel.cancel();
assert!(admitted_cancel.is_cancelled());
}
#[test]
fn ensure_pool_not_left_in_cmdline_after_decommission_allows_active_pool() {
assert!(ensure_pool_not_left_in_cmdline_after_decommission(0, "http://node{1...4}/disk{1...4}", false).is_ok());
@@ -6351,11 +6855,12 @@ mod pools_tests {
use super::{
DECOMMISSION_ENTRY_CONCURRENCY_DEFAULT_CAP, DECOMMISSION_ENTRY_CONCURRENCY_HARD_CAP, DECOMMISSION_ENTRY_QUEUE_HARD_CAP,
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo, DecommissionCanceler,
DecommissionEntryEnqueueResult, DecommissionStartPoolState, DecommissionTerminalState, ListCallback,
PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, QueuedDecommissionEntry, apply_decommission_status_space_info,
await_decommission_worker, bind_decommission_cancelers, bind_missing_decommission_cancelers,
cancel_decommission_canceler, clamp_decommission_entry_concurrency, classify_decommission_terminal_state,
count_decommission_item, decommission_cancel_signal_result, decommission_entry_queue_capacity, decommission_item_size,
DecommissionEntryEnqueueResult, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, POOL_META_NAME,
PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, QueuedDecommissionEntry, REBAL_META_NAME,
acquire_pool_rebalance_activation_locks, apply_decommission_status_space_info, await_decommission_worker,
bind_decommission_cancelers, bind_missing_decommission_cancelers, cancel_decommission_canceler,
clamp_decommission_entry_concurrency, classify_decommission_terminal_state, count_decommission_item,
decommission_cancel_signal_result, decommission_entry_queue_capacity, decommission_item_size,
decommission_meta_bucket_options, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency,
default_decommission_entry_concurrency, drain_decommission_entry_queue, enqueue_decommission_entry,
ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_generation,
@@ -6396,15 +6901,45 @@ mod pools_tests {
use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo};
use rustfs_filemeta::{MetaCacheEntries, MetadataResolutionParams};
use rustfs_rio::Index;
use std::future::Future;
use std::sync::{
Arc,
Arc, Mutex,
atomic::{AtomicBool, AtomicUsize, Ordering},
};
use std::task::{Context, Poll};
use std::time::Duration as StdDuration;
use time::{Duration, OffsetDateTime};
use tokio::sync::Semaphore;
use tokio_util::sync::CancellationToken;
#[derive(Debug)]
struct ActivationLockRecorder {
lock_manager: Arc<rustfs_lock::GlobalLockManager>,
owner: &'static str,
resources: Mutex<Vec<String>>,
}
#[async_trait::async_trait]
impl crate::storage_api_contracts::namespace::NamespaceLocking for ActivationLockRecorder {
type Error = Error;
type NamespaceLock = rustfs_lock::NamespaceLockWrapper;
async fn new_ns_lock(&self, bucket: &str, object: &str) -> crate::error::Result<Self::NamespaceLock> {
self.resources
.lock()
.expect("activation lock recorder should not be poisoned")
.push(object.to_string());
Ok(rustfs_lock::NamespaceLockWrapper::new(
rustfs_lock::NamespaceLock::with_local_manager(
"activation-lock-test".to_string(),
Arc::clone(&self.lock_manager),
),
rustfs_lock::ObjectKey::new(bucket, object),
self.owner.to_string(),
))
}
}
fn noop_decommission_list_callback() -> ListCallback {
Arc::new(|_| Box::pin(async {}))
}
@@ -6453,6 +6988,48 @@ mod pools_tests {
}
}
#[tokio::test]
async fn test_activation_fence_uses_one_lock_order_and_serializes_callers() {
let manager = Arc::new(rustfs_lock::GlobalLockManager::new());
let first = Arc::new(ActivationLockRecorder {
lock_manager: Arc::clone(&manager),
owner: "first",
resources: Mutex::new(Vec::new()),
});
let second = Arc::new(ActivationLockRecorder {
lock_manager: manager,
owner: "second",
resources: Mutex::new(Vec::new()),
});
let first_guards = acquire_pool_rebalance_activation_locks(first.clone())
.await
.expect("first activation should acquire both locks");
assert_eq!(
*first
.resources
.lock()
.expect("activation lock recorder should not be poisoned"),
vec![POOL_META_NAME.to_string(), REBAL_META_NAME.to_string()]
);
let mut second_acquire = Box::pin(acquire_pool_rebalance_activation_locks(second.clone()));
let mut context = Context::from_waker(futures::task::noop_waker_ref());
assert!(matches!(second_acquire.as_mut().poll(&mut context), Poll::Pending));
drop(first_guards);
second_acquire
.await
.expect("second activation should acquire both locks after the first releases them");
assert_eq!(
*second
.resources
.lock()
.expect("activation lock recorder should not be poisoned"),
vec![POOL_META_NAME.to_string(), REBAL_META_NAME.to_string()]
);
}
#[test]
fn test_apply_decommission_status_space_info_adds_idle_pool_usage() {
let status = apply_decommission_status_space_info(
+12 -4
View File
@@ -1241,8 +1241,16 @@ pub(crate) async fn make_local_two_set_sets() -> (Vec<tempfile::TempDir>, Arc<Se
make_local_two_set_sets_with_ctx(bootstrap_ctx()).await
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>) -> (Vec<tempfile::TempDir>, Arc<Sets>) {
make_local_two_set_sets_for_pool_with_ctx(ctx, 0).await
}
#[cfg(any(test, feature = "test-util"))]
pub(crate) async fn make_local_two_set_sets_for_pool_with_ctx(
ctx: Arc<InstanceContext>,
pool_idx: usize,
) -> (Vec<tempfile::TempDir>, Arc<Sets>) {
use crate::layout::endpoint::Endpoint;
use rustfs_lock::client::local::LocalClient;
@@ -1258,7 +1266,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>)
let temp_dir = tempfile::tempdir().expect("tempdir should be created");
let mut endpoint = Endpoint::try_from(temp_dir.path().to_str().expect("tempdir path should be utf8"))
.expect("endpoint should parse");
endpoint.set_pool_index(0);
endpoint.set_pool_index(pool_idx);
endpoint.set_set_index(set_index);
endpoint.set_disk_index(disk_index);
let disk = new_disk(
@@ -1294,7 +1302,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>)
2,
1,
set_index,
0,
pool_idx,
endpoints,
format.clone(),
lockers,
@@ -1307,7 +1315,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc<InstanceContext>)
let sets = Arc::new(Sets {
id: format.id,
disk_set: disk_sets,
pool_idx: 0,
pool_idx,
endpoints: PoolEndpoints {
legacy: false,
set_count: 2,
+69 -26
View File
@@ -1023,10 +1023,11 @@ pub(crate) enum SourceCleanupError {
Storage(#[from] Error),
}
#[derive(Clone, Copy, Default)]
#[derive(Clone, Default)]
pub(crate) struct SourceCleanupBucketFence<'a> {
pub(crate) expected_incarnation_id: Option<uuid::Uuid>,
pub(crate) lifecycle_guard: Option<&'a rustfs_lock::NamespaceLockGuard>,
pub(crate) namespace_lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
pub(crate) object_mutation_fence: Option<&'a SourceCleanupMutationFence>,
}
@@ -1061,7 +1062,7 @@ pub(crate) async fn ensure_source_cleanup_versions_unchanged(
ensure_source_cleanup_versions_match(expected, &current, allowed_missing)
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
struct SourceCleanupDeleteBarrierState {
bucket: String,
object: String,
@@ -1071,7 +1072,7 @@ struct SourceCleanupDeleteBarrierState {
release: tokio::sync::Notify,
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
#[allow(
dead_code,
reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)"
@@ -1080,11 +1081,11 @@ pub(crate) struct SourceCleanupDeleteBarrier {
state: Arc<SourceCleanupDeleteBarrierState>,
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
static SOURCE_CLEANUP_DELETE_BARRIERS: std::sync::OnceLock<std::sync::Mutex<Vec<Arc<SourceCleanupDeleteBarrierState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
#[allow(
dead_code,
reason = "installed by set_disk object tests behind `--features test-util` (backlog#1823)"
@@ -1148,7 +1149,7 @@ pub(crate) fn notify_source_cleanup_mutation_fence_pending(bucket: &str, object:
}
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
impl Drop for SourceCleanupDeleteBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
@@ -1160,7 +1161,7 @@ impl Drop for SourceCleanupDeleteBarrier {
}
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
async fn pause_source_cleanup_before_delete(bucket: &str, object: &str) {
let barrier = SOURCE_CLEANUP_DELETE_BARRIERS
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
@@ -1220,7 +1221,7 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
ensure_source_cleanup_versions_unchanged(set.clone(), bucket, object, expected, allowed_missing, op_label).await?;
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
pause_source_cleanup_before_delete(bucket, object).await;
let mut opts = ObjectOptions {
@@ -1240,6 +1241,9 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
if let Some(bucket_lifecycle_guard) = bucket_fence.lifecycle_guard {
opts.add_bucket_lifecycle_lock_guard(bucket_lifecycle_guard);
}
if let Some(signal) = bucket_fence.namespace_lock_lost_signal {
opts.add_namespace_lock_lost_signal(signal);
}
let result = set.delete_object(bucket, cleanup_key.as_str(), opts).await;
if result.is_ok() {
crate::store::list_objects::observe_scanner_namespace_mutations(bucket, 1);
@@ -1410,11 +1414,13 @@ pub(crate) async fn migrate_decommission_object(
rd,
source_bucket_incarnation_id,
op_label,
None,
Some(&_mutation_fence),
)
.await
}
#[cfg(test)]
pub(crate) async fn migrate_object(
store: Arc<ECStore>,
pool_idx: usize,
@@ -1423,9 +1429,33 @@ pub(crate) async fn migrate_object(
source_bucket_incarnation_id: Option<uuid::Uuid>,
op_label: &str,
) -> Result<()> {
migrate_object_inner(store, pool_idx, bucket, rd, source_bucket_incarnation_id, op_label, None).await
migrate_object_with_lock_lost_signal(store, pool_idx, bucket, rd, source_bucket_incarnation_id, op_label, None).await
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn migrate_object_with_lock_lost_signal(
store: Arc<ECStore>,
pool_idx: usize,
bucket: String,
rd: GetObjectReader,
source_bucket_incarnation_id: Option<uuid::Uuid>,
op_label: &str,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> Result<()> {
migrate_object_inner(
store,
pool_idx,
bucket,
rd,
source_bucket_incarnation_id,
op_label,
lock_lost_signal,
None,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn migrate_object_inner(
store: Arc<ECStore>,
pool_idx: usize,
@@ -1433,6 +1463,7 @@ async fn migrate_object_inner(
rd: GetObjectReader,
source_bucket_incarnation_id: Option<uuid::Uuid>,
op_label: &str,
lock_lost_signal: Option<Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
mutation_fence: Option<&ObjectLockDiagGuard>,
) -> Result<()> {
let object_info = rd.object_info.clone();
@@ -1446,6 +1477,9 @@ async fn migrate_object_inner(
if should_use_multipart_data_movement(&object_info, has_part_checksums) {
let mut new_multipart_opts = data_movement_new_multipart_opts(&object_info, pool_idx);
new_multipart_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
if let Some(signal) = lock_lost_signal.as_ref() {
new_multipart_opts.add_namespace_lock_lost_signal(Arc::clone(signal));
}
let (res, target_pool_idx, expected_bucket_incarnation_id) = match store
.handle_new_multipart_upload_with_pool_idx(&bucket, &object_info.name, &new_multipart_opts, mutation_fence)
.await
@@ -1490,7 +1524,7 @@ async fn migrate_object_inner(
err,
)
})?;
let part_opts = ObjectOptions {
let mut part_opts = ObjectOptions {
part_number: Some(part.number),
preserve_etag: Some(part.etag.clone()),
data_movement: true,
@@ -1498,6 +1532,9 @@ async fn migrate_object_inner(
expected_bucket_incarnation_id,
..Default::default()
};
if let Some(signal) = lock_lost_signal.as_ref() {
part_opts.add_namespace_lock_lost_signal(Arc::clone(signal));
}
let pi = match store
.put_object_part_for_data_movement(
target_pool_idx,
@@ -1542,6 +1579,9 @@ async fn migrate_object_inner(
)
})?;
complete_multipart_opts.expected_bucket_incarnation_id = expected_bucket_incarnation_id;
if let Some(signal) = lock_lost_signal.as_ref() {
complete_multipart_opts.add_namespace_lock_lost_signal(Arc::clone(signal));
}
if let Err(err) = store
.clone()
.complete_multipart_upload_for_data_movement(
@@ -1590,18 +1630,18 @@ async fn migrate_object_inner(
if multipart_result.is_ok() && should_abort_multipart_upload(&abort_multipart_flag) {
let abort_result = store
.abort_multipart_upload_for_data_movement(
target_pool_idx,
&bucket,
&object_info.name,
&res.upload_id,
&ObjectOptions {
.abort_multipart_upload_for_data_movement(target_pool_idx, &bucket, &object_info.name, &res.upload_id, &{
let mut opts = ObjectOptions {
data_movement: true,
src_pool_idx: pool_idx,
expected_bucket_incarnation_id,
..Default::default()
},
)
};
if let Some(signal) = lock_lost_signal.as_ref() {
opts.add_namespace_lock_lost_signal(Arc::clone(signal));
}
opts
})
.await;
match abort_result {
Ok(()) => return Ok(()),
@@ -1659,18 +1699,18 @@ async fn migrate_object_inner(
if let Err(primary_err) = multipart_result {
if should_abort_multipart_upload(&abort_multipart_flag) {
return match store
.abort_multipart_upload_for_data_movement(
target_pool_idx,
&bucket,
&object_info.name,
&res.upload_id,
&ObjectOptions {
.abort_multipart_upload_for_data_movement(target_pool_idx, &bucket, &object_info.name, &res.upload_id, &{
let mut opts = ObjectOptions {
data_movement: true,
src_pool_idx: pool_idx,
expected_bucket_incarnation_id,
..Default::default()
},
)
};
if let Some(signal) = lock_lost_signal.as_ref() {
opts.add_namespace_lock_lost_signal(Arc::clone(signal));
}
opts
})
.await
{
Ok(()) => Err(primary_err),
@@ -1705,6 +1745,9 @@ async fn migrate_object_inner(
let mut put_opts = data_movement_put_object_opts(&object_info, pool_idx);
put_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
if let Some(signal) = lock_lost_signal {
put_opts.add_namespace_lock_lost_signal(signal);
}
let (target_pool_idx, put_result) = store
.put_object_for_data_movement(&bucket, &object_info.name, &mut data, &put_opts, mutation_fence)
.await
+69
View File
@@ -84,6 +84,61 @@ impl NamespaceLockFence {
}
}
#[cfg(test)]
static NAMESPACE_LOCK_SIGNAL_TEST_FENCES: std::sync::OnceLock<std::sync::Mutex<Vec<(usize, NamespaceLockFence)>>> =
std::sync::OnceLock::new();
#[cfg(test)]
pub(crate) struct NamespaceLockSignalTestFence {
signal_key: usize,
}
#[cfg(test)]
impl NamespaceLockSignalTestFence {
pub(crate) fn install_with_loss_handle(
signal: &Arc<rustfs_lock::distributed_lock::LockLostSignal>,
loss_handle: Arc<std::sync::atomic::AtomicBool>,
) -> Self {
let fence = NamespaceLockFence {
signals: Arc::default(),
forced_lost: Arc::new(vec![loss_handle]),
};
let signal_key = Arc::as_ptr(signal) as usize;
let mut fences = NAMESPACE_LOCK_SIGNAL_TEST_FENCES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("namespace lock signal test fence should not be poisoned");
assert!(
!fences.iter().any(|(key, _)| *key == signal_key),
"namespace lock signal test fence must be unique"
);
fences.push((signal_key, fence));
Self { signal_key }
}
}
#[cfg(test)]
impl Drop for NamespaceLockSignalTestFence {
fn drop(&mut self) {
let mut fences = NAMESPACE_LOCK_SIGNAL_TEST_FENCES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("namespace lock signal test fence should not be poisoned");
fences.retain(|(key, _)| *key != self.signal_key);
}
}
#[cfg(test)]
pub(crate) fn namespace_lock_signal_test_fence_is_lost(signal: &Arc<rustfs_lock::distributed_lock::LockLostSignal>) -> bool {
NAMESPACE_LOCK_SIGNAL_TEST_FENCES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("namespace lock signal test fence should not be poisoned")
.iter()
.find(|(key, _)| *key == Arc::as_ptr(signal) as usize)
.is_some_and(|(_, fence)| fence.is_lock_lost())
}
#[derive(Debug)]
pub struct ObjectLockConfigSnapshot {
store_id: Option<Uuid>,
@@ -405,9 +460,23 @@ impl ObjectOptions {
}
pub(crate) fn add_namespace_lock_lost_signal(&mut self, signal: Arc<rustfs_lock::distributed_lock::LockLostSignal>) {
#[cfg(test)]
let test_fence = NAMESPACE_LOCK_SIGNAL_TEST_FENCES
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("namespace lock signal test fence should not be poisoned")
.iter()
.find(|(key, _)| *key == Arc::as_ptr(&signal) as usize)
.map(|(_, fence)| fence.clone());
self.namespace_lock_fence
.get_or_insert_with(NamespaceLockFence::new)
.add_signal(signal);
#[cfg(test)]
if let Some(test_fence) = test_fence {
self.namespace_lock_fence
.get_or_insert_with(NamespaceLockFence::new)
.extend(&test_fence);
}
}
pub(crate) fn ensure_namespace_lock_fence(&mut self) {
+2 -2
View File
@@ -119,8 +119,8 @@ pub(crate) fn endpoint_erasure_set_count() -> Option<usize> {
endpoint_pools().map(|endpoints| endpoints.es_count())
}
pub(crate) fn endpoint_pool_is_local(pool_index: usize) -> bool {
get_global_endpoints()
pub(crate) fn endpoint_pool_is_local(endpoints: &EndpointServerPools, pool_index: usize) -> bool {
endpoints
.as_ref()
.get(pool_index)
.is_some_and(|pool| pool.endpoints.as_ref().first().is_some_and(|endpoint| endpoint.is_local))
@@ -1121,9 +1121,21 @@ impl NotificationSys {
}
}
match store.stop_rebalance_for_id(expected_rebalance_id).await {
let local_rebalance_id = match expected_rebalance_id {
Some(expected_id) => Some(expected_id.to_owned()),
None => store.current_rebalance_id().await,
};
match store.stop_rebalance_for_id(local_rebalance_id.as_deref()).await {
Ok(_) => {
if let Err(err) = store.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt).await {
let save_result = match local_rebalance_id.as_deref() {
Some(expected_id) => {
store
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_id)
.await
}
None => Ok(()),
};
if let Err(err) = save_result {
error!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
@@ -102,11 +102,20 @@ pub(crate) trait MigrationBackend: Send + Sync {
pub(crate) struct RebalanceMigrationBackend<'a> {
source: &'a SetDisks,
store: &'a ECStore,
lock_lost_signal: Option<std::sync::Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
}
impl<'a> RebalanceMigrationBackend<'a> {
pub(crate) fn new(source: &'a SetDisks, store: &'a ECStore) -> Self {
Self { source, store }
pub(crate) fn new(
source: &'a SetDisks,
store: &'a ECStore,
lock_lost_signal: Option<std::sync::Arc<rustfs_lock::distributed_lock::LockLostSignal>>,
) -> Self {
Self {
source,
store,
lock_lost_signal,
}
}
}
@@ -130,7 +139,11 @@ impl MigrationBackend for RebalanceMigrationBackend<'_> {
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<()> {
self.store.decommission_tiered_object(bucket, object, fi, opts).await
let mut opts = opts.clone();
if let Some(signal) = self.lock_lost_signal.as_ref() {
opts.add_namespace_lock_lost_signal(std::sync::Arc::clone(signal));
}
self.store.decommission_tiered_object(bucket, object, fi, &opts).await
}
}
@@ -12,6 +12,8 @@
// See the License for the specific language governing permissions and
// limitations under the License.
#[cfg(test)]
use crate::disk::DiskAPI;
use crate::error::{Error, Result};
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader};
use tokio::time::Duration;
@@ -45,6 +47,8 @@ mod runtime;
mod types;
mod worker;
#[cfg(feature = "test-util")]
pub use entry::test_util::PausedRebalanceEntryTestFixture;
pub(crate) use meta::is_rebalance_conflicting_with_decommission;
pub use meta::{decode_rebalance_stop_propagation_record, encode_rebalance_stop_propagation_record};
pub use types::{
@@ -53,5 +57,91 @@ pub use types::{
};
use types::{RebalanceBucketConfigs, RebalanceBucketOutcome, RebalanceEntryOutcome};
#[cfg(any(test, feature = "test-util"))]
pub async fn test_store_with_persisted_rebalance_meta(
meta: RebalanceMeta,
) -> (Vec<tempfile::TempDir>, std::sync::Arc<crate::store::ECStore>) {
let ctx = std::sync::Arc::new(crate::runtime::instance::InstanceContext::new());
let (temp_dirs, pool) = crate::core::sets::make_local_two_set_sets_with_ctx(ctx.clone()).await;
meta.save(pool.clone())
.await
.expect("rebalance test metadata should be persisted");
let endpoint_pools: crate::layout::endpoints::EndpointServerPools = vec![pool.endpoints.clone()].into();
let store = std::sync::Arc::new(crate::store::ECStore {
id: uuid::Uuid::new_v4(),
disk_map: std::collections::HashMap::new(),
pools: vec![pool],
peer_sys: crate::cluster::rpc::S3PeerSys::new_with_instance_ctx(&endpoint_pools, ctx.clone()),
pool_meta: tokio::sync::RwLock::new(crate::core::pools::PoolMeta::default()),
rebalance_meta: tokio::sync::RwLock::new(Some(meta)),
decommission_cancelers: tokio::sync::RwLock::new(vec![None]),
start_gate: tokio::sync::Mutex::new(()),
pool_meta_save_gate: tokio::sync::Mutex::new(()),
ctx,
bucket_fence_registry: std::sync::Arc::default(),
});
(temp_dirs, store)
}
#[cfg(test)]
pub(crate) async fn test_two_pool_stores(
rebalance_meta: Option<RebalanceMeta>,
) -> (
Vec<tempfile::TempDir>,
std::sync::Arc<crate::store::ECStore>,
std::sync::Arc<crate::store::ECStore>,
) {
use crate::core::pools::PoolMeta;
use crate::layout::endpoints::{EndpointServerPools, SetupType};
let ctx = std::sync::Arc::new(crate::runtime::instance::InstanceContext::new());
ctx.update_erasure_type(SetupType::DistErasure).await;
let (mut temp_dirs, first_pool) =
crate::core::sets::make_local_two_set_sets_for_pool_with_ctx(std::sync::Arc::clone(&ctx), 0).await;
let (second_temp_dirs, second_pool) =
crate::core::sets::make_local_two_set_sets_for_pool_with_ctx(std::sync::Arc::clone(&ctx), 1).await;
temp_dirs.extend(second_temp_dirs);
let pools = vec![first_pool, second_pool];
{
let local_disk_map = ctx.local_disk_map();
let mut local_disk_map = local_disk_map.write().await;
for pool in &pools {
for set in &pool.disk_set {
for disk in set.disks.read().await.iter().flatten() {
local_disk_map.insert(disk.endpoint().to_string(), Some(disk.clone()));
}
}
}
}
let pool_meta = PoolMeta::new(&pools, &PoolMeta::default());
pool_meta
.save(pools.clone())
.await
.expect("baseline pool metadata should be persisted");
if let Some(meta) = rebalance_meta.as_ref() {
meta.save(pools[0].clone())
.await
.expect("active rebalance metadata should be persisted");
}
let endpoint_pools: EndpointServerPools = pools.iter().map(|pool| pool.endpoints.clone()).collect::<Vec<_>>().into();
ctx.set_endpoints(endpoint_pools.clone());
let make_store = || {
std::sync::Arc::new(crate::store::ECStore {
id: uuid::Uuid::new_v4(),
disk_map: std::collections::HashMap::new(),
pools: pools.clone(),
peer_sys: crate::cluster::rpc::S3PeerSys::new_with_instance_ctx(&endpoint_pools, std::sync::Arc::clone(&ctx)),
pool_meta: tokio::sync::RwLock::new(pool_meta.clone()),
rebalance_meta: tokio::sync::RwLock::new(rebalance_meta.clone()),
decommission_cancelers: tokio::sync::RwLock::new(vec![None, None]),
start_gate: tokio::sync::Mutex::new(()),
pool_meta_save_gate: tokio::sync::Mutex::new(()),
ctx: std::sync::Arc::clone(&ctx),
bucket_fence_registry: std::sync::Arc::default(),
})
};
(temp_dirs, make_store(), make_store())
}
#[cfg(test)]
mod rebalance_unit_tests;
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use super::control::validate_rebalance_disk_stats_coverage;
use super::control::{fail_next_rebalance_activation_save_for_test, validate_rebalance_disk_stats_coverage};
use super::meta::{
RebalanceMetaMergeOutcome, RebalanceTerminalEvent, apply_rebalance_save_option, apply_rebalance_terminal_event,
apply_stopped_at, classify_rebalance_terminal_event, clone_arc_by_index, clone_first_arc, clone_rebalance_pool_stats,
@@ -32,7 +32,11 @@ use super::migration::{
MigrationBackend, MigrationVersionResult, migrate_entry_version, migrate_entry_version_with_retry_wait,
rebalance_delete_marker_opts,
};
use super::runtime::{should_fail_repeated_rebalance_bucket_defer, source_cleanup_defer_attempt};
use super::runtime::{
RebalanceLocalActivationOutcome, commit_local_rebalance_worker_activation,
commit_local_rebalance_worker_activation_candidate, should_fail_repeated_rebalance_bucket_defer,
source_cleanup_defer_attempt, stage_local_rebalance_worker_activation,
};
use super::worker::{
RebalanceEntryCleanupResult, ensure_rebalance_listing_disks_available, is_transient_rebalance_error,
parse_rebalance_max_attempts, rebalance_listing_retry_delay, rebalance_migration_retry_delay,
@@ -48,6 +52,7 @@ use super::worker::{
use super::{
DiskStat, GetObjectReader, ObjectInfo, ObjectOptions, RebalSaveOpt, RebalStatus, RebalanceBucketConfigs,
RebalanceBucketOutcome, RebalanceCleanupWarnings, RebalanceEntryOutcome, RebalanceInfo, RebalanceMeta, RebalanceStats,
RebalanceStopPropagationRecord,
};
use super::{REBALANCE_DEFERRED_ENTRY_ERROR_PREFIX, REBALANCE_SOURCE_CLEANUP_DEFERRED_ERROR_PREFIX};
use crate::bucket::replication::{ReplicationState, ReplicationStatusType, replication_state_to_filemeta};
@@ -2708,6 +2713,281 @@ async fn test_start_rebalance_for_id_rejects_stopped_metadata() {
assert!(err.to_string().contains("was stopped before start"));
}
#[test]
fn test_stopped_activation_state_prevents_worker_token_commit() {
let mut meta = RebalanceMeta {
id: "rebalance-a".to_string(),
stopped_at: Some(OffsetDateTime::now_utc()),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let outcome = commit_local_rebalance_worker_activation(&mut meta, "rebalance-a", tokio_util::sync::CancellationToken::new())
.expect("stopped metadata should produce a non-start outcome");
assert_eq!(outcome, RebalanceLocalActivationOutcome::NotStartedTerminal);
assert!(meta.cancel.is_none(), "stopped rebalance must not receive a worker token");
}
#[test]
fn test_rebalance_activation_candidate_does_not_clobber_replacement_token() {
let mut local = RebalanceMeta {
id: "rebalance-a".to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
buckets: vec!["bucket-a".to_string()],
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let (candidate, outcome, must_persist) =
stage_local_rebalance_worker_activation(&local, "rebalance-a", CancellationToken::new(), OffsetDateTime::UNIX_EPOCH)
.expect("active activation candidate should be staged");
assert_eq!(outcome, RebalanceLocalActivationOutcome::Started);
assert!(!must_persist);
let replacement = CancellationToken::new();
local.cancel = Some(replacement.clone());
let err = commit_local_rebalance_worker_activation_candidate(&mut local, "rebalance-a", None, candidate)
.expect_err("a replacement token must reject the stale activation candidate");
assert!(err.to_string().contains("worker token changed"));
assert_eq!(local.cancel.as_ref(), Some(&replacement));
}
#[tokio::test]
#[serial_test::serial]
async fn test_rebalance_start_save_failure_retries_persisted_completed_state() {
let active = RebalanceMeta {
id: "rebalance-real-save-completed".to_string(),
percent_free_goal: 0.5,
pool_stats: vec![RebalanceStats {
participating: true,
init_free_space: 400,
init_capacity: 1_000,
bytes: 100,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let (_temp_dirs, store) = super::test_store_with_persisted_rebalance_meta(active).await;
fail_next_rebalance_activation_save_for_test("rebalance-real-save-completed");
let err = store
.start_rebalance_under_gate()
.await
.expect_err("the injected first activation save must fail through the real start path");
assert!(err.to_string().contains("injected rebalance activation save failure"));
{
let local = store.rebalance_meta.read().await;
let local = local.as_ref().expect("local rebalance metadata should remain present");
assert_eq!(local.pool_stats[0].info.status, RebalStatus::Started);
assert!(local.cancel.is_none(), "failed persistence must not publish a worker token");
}
let mut after_failure = RebalanceMeta::new();
after_failure
.load(store.pools[0].clone())
.await
.expect("active metadata should remain readable after the failed save");
assert_eq!(after_failure.pool_stats[0].info.status, RebalStatus::Started);
store
.start_rebalance_under_gate()
.await
.expect("the real start path must retry and persist the terminal candidate");
let mut persisted = RebalanceMeta::new();
persisted
.load(store.pools[0].clone())
.await
.expect("retry-persisted completed metadata should be readable");
assert_eq!(persisted.pool_stats[0].info.status, RebalStatus::Completed);
let local = store.rebalance_meta.read().await;
let local = local.as_ref().expect("local rebalance metadata should remain present");
assert_eq!(local.pool_stats[0].info.status, RebalStatus::Completed);
assert!(local.cancel.is_none());
}
#[tokio::test]
#[serial_test::serial]
async fn test_rebalance_start_save_failure_retries_persisted_stopped_state() {
let active = RebalanceMeta {
id: "rebalance-real-save-stopped".to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let (_temp_dirs, store) = super::test_store_with_persisted_rebalance_meta(active).await;
let stopped_at = OffsetDateTime::from_unix_timestamp(1_000).expect("test timestamp should be valid");
{
let mut local = store.rebalance_meta.write().await;
let local = local.as_mut().expect("local rebalance metadata should remain present");
local.stopped_at = Some(stopped_at);
local.pool_stats[0].info.status = RebalStatus::Stopped;
local.pool_stats[0].info.end_time = Some(stopped_at);
}
fail_next_rebalance_activation_save_for_test("rebalance-real-save-stopped");
let err = store
.start_rebalance_under_gate()
.await
.expect_err("the injected first stopped-state save must fail through the real start path");
assert!(err.to_string().contains("injected rebalance activation save failure"));
let mut after_failure = RebalanceMeta::new();
after_failure
.load(store.pools[0].clone())
.await
.expect("active metadata should remain readable after the failed save");
assert_eq!(after_failure.pool_stats[0].info.status, RebalStatus::Started);
{
let local = store.rebalance_meta.read().await;
let local = local.as_ref().expect("local rebalance metadata should remain present");
assert_eq!(local.pool_stats[0].info.status, RebalStatus::Stopped);
assert!(local.cancel.is_none(), "failed persistence must not publish a worker token");
}
store
.start_rebalance_under_gate()
.await
.expect("the real start path must retry and persist the stopped candidate");
let mut persisted = RebalanceMeta::new();
persisted
.load(store.pools[0].clone())
.await
.expect("retry-persisted stopped metadata should be readable");
assert_eq!(persisted.stopped_at, Some(stopped_at));
assert_eq!(persisted.pool_stats[0].info.status, RebalStatus::Stopped);
let local = store.rebalance_meta.read().await;
let local = local.as_ref().expect("local rebalance metadata should remain present");
assert_eq!(local.pool_stats[0].info.status, RebalStatus::Stopped);
assert!(local.cancel.is_none());
}
#[tokio::test]
async fn test_old_worker_cannot_mutate_replacement_rebalance_state() {
let meta = RebalanceMeta {
id: "rebalance-b".to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
buckets: vec!["bucket-a".to_string()],
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let store = test_store_with_rebalance_meta(meta);
let fi = FileInfo {
size: 128,
..Default::default()
};
for err in [
store
.next_rebal_bucket(0, "rebalance-a")
.await
.expect_err("old worker must not read replacement work"),
store
.bucket_rebalance_done(0, "bucket-a".to_string(), "rebalance-a")
.await
.expect_err("old worker must not complete replacement bucket"),
store
.update_pool_stats_batch_for_rebalance(0, "bucket-a".to_string(), &[&fi], "rebalance-a")
.await
.expect_err("old worker must not update replacement stats"),
store
.check_if_rebalance_done(0, "rebalance-a")
.await
.expect_err("old worker must not complete replacement pool"),
store
.save_rebalance_stats_for_id(0, RebalSaveOpt::Stats, "rebalance-a")
.await
.expect_err("old save task must not persist replacement metadata"),
store
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, "rebalance-a")
.await
.expect_err("old stop path must not persist replacement metadata"),
] {
assert!(err.to_string().contains("stale rebalance worker rejected"));
}
let meta = store.rebalance_meta.read().await;
let meta = meta.as_ref().expect("replacement metadata should remain present");
assert_eq!(meta.id, "rebalance-b");
assert!(meta.pool_stats[0].rebalanced_buckets.is_empty());
assert_eq!(meta.pool_stats[0].bytes, 0);
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started);
assert!(meta.stopped_at.is_none());
}
#[tokio::test]
async fn test_rebalance_metadata_reload_under_start_gate_does_not_reacquire_gate() {
let store = test_store_with_rebalance_meta(RebalanceMeta::default());
let _start_guard = store.start_gate.lock().await;
let err = store
.load_rebalance_meta_under_start_gate()
.await
.expect_err("empty test store should reach the metadata load without waiting on start_gate again");
assert!(err.to_string().contains("no pools available"));
}
#[tokio::test]
async fn test_stale_stop_propagation_cannot_mutate_replacement_rebalance() {
let meta = RebalanceMeta {
id: "rebalance-b".to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let store = test_store_with_rebalance_meta(meta);
let record = RebalanceStopPropagationRecord {
stop_failures: vec!["old rebalance stop failed".to_string()],
..Default::default()
};
let err = store
.record_rebalance_stop_propagation("rebalance-a", record)
.await
.expect_err("old propagation failure must not mutate replacement metadata");
assert!(err.to_string().contains("stale rebalance worker rejected"));
let meta = store.rebalance_meta.read().await;
let meta = meta.as_ref().expect("replacement metadata should remain present");
assert_eq!(meta.id, "rebalance-b");
assert!(meta.last_refreshed_at.is_none());
assert!(meta.pool_stats[0].info.last_error.is_none());
}
fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc<crate::store::ECStore> {
let endpoint_pools: crate::layout::endpoints::EndpointServerPools = Vec::new().into();
Arc::new(crate::store::ECStore {
+218 -33
View File
@@ -1,3 +1,4 @@
use super::control::RebalanceWorkerActivationFence;
use super::meta::{
apply_rebalance_save_option, apply_rebalance_terminal_event, classify_rebalance_terminal_event, clone_first_arc,
complete_rebalance_pools_at_goal, complete_rebalance_pools_with_empty_queue, ensure_valid_rebalance_pool_index,
@@ -39,9 +40,97 @@ pub(super) fn source_cleanup_defer_attempt(deferred_attempts: &mut HashMap<Strin
*attempts
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum RebalanceLocalActivationOutcome {
Started,
NotStartedTerminal,
}
pub(super) fn commit_local_rebalance_worker_activation(
meta: &mut super::RebalanceMeta,
expected_id: &str,
cancel: CancellationToken,
) -> Result<RebalanceLocalActivationOutcome> {
if meta.id != expected_id {
return Err(Error::other(format!(
"rebalance metadata changed before local worker activation: expected {expected_id}, found {}",
meta.id
)));
}
if meta.stopped_at.is_some() || !is_rebalance_in_progress(meta) {
return Ok(RebalanceLocalActivationOutcome::NotStartedTerminal);
}
meta.cancel = Some(cancel);
Ok(RebalanceLocalActivationOutcome::Started)
}
pub(super) fn stage_local_rebalance_worker_activation(
meta: &super::RebalanceMeta,
expected_id: &str,
cancel: CancellationToken,
now: OffsetDateTime,
) -> Result<(super::RebalanceMeta, RebalanceLocalActivationOutcome, bool)> {
let mut candidate = meta.clone();
let completed_at_goal = complete_rebalance_pools_at_goal(&mut candidate, now);
let completed_empty_queue = complete_rebalance_pools_with_empty_queue(&mut candidate, now);
let outcome = commit_local_rebalance_worker_activation(&mut candidate, expected_id, cancel)?;
let must_persist =
completed_at_goal || completed_empty_queue || outcome == RebalanceLocalActivationOutcome::NotStartedTerminal;
Ok((candidate, outcome, must_persist))
}
pub(super) fn commit_local_rebalance_worker_activation_candidate(
current: &mut super::RebalanceMeta,
expected_id: &str,
expected_cancel: Option<&CancellationToken>,
candidate: super::RebalanceMeta,
) -> Result<()> {
if current.id != expected_id || candidate.id != expected_id {
return Err(Error::other(format!(
"rebalance metadata changed before local worker activation commit: expected {expected_id}, found {}",
current.id
)));
}
if !Arc::ptr_eq(&current.activation_gate, &candidate.activation_gate) {
return Err(Error::other(format!(
"rebalance activation gate changed before local worker activation commit: {expected_id}"
)));
}
if current.cancel.as_ref() != expected_cancel {
return Err(Error::other(format!(
"rebalance worker token changed before local worker activation commit: {expected_id}"
)));
}
*current = candidate;
Ok(())
}
pub(super) fn rollback_local_rebalance_worker_activation(
meta: Option<&mut super::RebalanceMeta>,
expected_id: &str,
activation_token: &CancellationToken,
) -> bool {
let Some(meta) = meta else {
return false;
};
if meta.id != expected_id || meta.cancel.as_ref() != Some(activation_token) {
return false;
}
if let Some(cancel) = meta.cancel.take() {
cancel.cancel();
return true;
}
false
}
impl ECStore {
#[tracing::instrument(skip_all)]
pub async fn start_rebalance(self: &Arc<Self>) -> Result<()> {
let _start_guard = self.start_gate.lock().await;
self.start_rebalance_under_gate().await
}
pub(super) async fn start_rebalance_under_gate(self: &Arc<Self>) -> Result<()> {
info!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -49,12 +138,27 @@ impl ECStore {
state = "starting",
"Starting rebalance"
);
let expected_id: Arc<str> = {
let rebalance_meta = self.rebalance_meta.read().await;
Arc::from(rebalance_meta.as_ref().ok_or(Error::ConfigNotFound)?.id.as_str())
};
let pool = clone_first_arc(self.pools.as_slice(), "start_rebalance: no pools available")?;
let activation_fence = match self
.fence_rebalance_worker_activation(pool.clone(), expected_id.as_ref())
.await?
{
RebalanceWorkerActivationFence::Ready(fence) => fence,
RebalanceWorkerActivationFence::NotStartedTerminal => return Ok(()),
};
let decommission_running = self.is_decommission_running().await;
// let rebalance_meta = self.rebalance_meta.read().await;
let cancel_tx = CancellationToken::new();
let rx = cancel_tx.clone();
let mut meta_to_save = None;
let activation_outcome;
let candidate;
let expected_cancel;
let must_persist;
{
let mut rebalance_meta = self.rebalance_meta.write().await;
@@ -74,25 +178,70 @@ impl ECStore {
);
return Ok(());
}
let now = OffsetDateTime::now_utc();
if complete_rebalance_pools_at_goal(meta, now) {
meta_to_save = Some(meta.clone());
expected_cancel = meta.cancel.clone();
(candidate, activation_outcome, must_persist) = stage_local_rebalance_worker_activation(
meta,
expected_id.as_ref(),
cancel_tx.clone(),
OffsetDateTime::now_utc(),
)?;
if let Err(err) = activation_fence.ensure_held() {
cancel_tx.cancel();
return Err(err);
}
if complete_rebalance_pools_with_empty_queue(meta, now) {
meta_to_save = Some(meta.clone());
if !must_persist
&& let Err(err) = commit_local_rebalance_worker_activation_candidate(
meta,
expected_id.as_ref(),
expected_cancel.as_ref(),
candidate.clone(),
)
{
cancel_tx.cancel();
return Err(err);
}
meta.cancel = Some(cancel_tx);
drop(rebalance_meta);
}
if let Some(meta) = meta_to_save {
let pool = clone_first_arc(self.pools.as_slice(), "start_rebalance: no pools available")?;
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, &meta, "start_rebalance complete pools at goal")
.await,
"start_rebalance complete pools at goal",
)?;
if must_persist {
let save_result = resolve_rebalance_meta_save_result(
self.save_rebalance_meta_under_activation_fence(
pool,
&candidate,
"start_rebalance persist activation candidate",
activation_fence.as_ref(),
expected_id.as_ref(),
)
.await,
"start_rebalance persist activation candidate",
);
if let Err(err) = save_result {
cancel_tx.cancel();
return Err(err);
}
let mut rebalance_meta = self.rebalance_meta.write().await;
let Some(meta) = rebalance_meta.as_mut() else {
cancel_tx.cancel();
return Err(Error::ConfigNotFound);
};
if let Err(err) = commit_local_rebalance_worker_activation_candidate(
meta,
expected_id.as_ref(),
expected_cancel.as_ref(),
candidate,
) {
cancel_tx.cancel();
return Err(err);
}
}
if !must_persist && let Err(err) = activation_fence.ensure_held() {
let mut rebalance_meta = self.rebalance_meta.write().await;
rollback_local_rebalance_worker_activation(rebalance_meta.as_mut(), expected_id.as_ref(), &rx);
return Err(err);
}
drop(activation_fence);
if activation_outcome != RebalanceLocalActivationOutcome::Started {
return Ok(());
}
let participants = if let Some(ref meta) = *self.rebalance_meta.read().await {
@@ -110,6 +259,8 @@ impl ECStore {
};
if !participants.iter().any(|participating| *participating) {
let mut rebalance_meta = self.rebalance_meta.write().await;
rollback_local_rebalance_worker_activation(rebalance_meta.as_mut(), expected_id.as_ref(), &rx);
debug!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -121,6 +272,11 @@ impl ECStore {
return Ok(());
}
#[cfg(test)]
let endpoints = self.instance_endpoints().unwrap_or_else(|| self.endpoints());
#[cfg(not(test))]
let endpoints = self.endpoints();
let mut workers_started = 0usize;
for (idx, participating) in participants.iter().enumerate() {
if !*participating {
@@ -136,7 +292,7 @@ impl ECStore {
continue;
}
if !runtime_sources::endpoint_pool_is_local(idx) {
if !runtime_sources::endpoint_pool_is_local(&endpoints, idx) {
debug!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -152,9 +308,10 @@ impl ECStore {
let pool_idx = idx;
let store = self.clone();
let rx_clone = rx.clone();
let worker_id = Arc::clone(&expected_id);
workers_started += 1;
tokio::spawn(async move {
if let Err(err) = store.rebalance_buckets(rx_clone, pool_idx).await {
if let Err(err) = store.rebalance_buckets(rx_clone, pool_idx, worker_id).await {
error!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -178,6 +335,8 @@ impl ECStore {
}
if workers_started == 0 {
let mut rebalance_meta = self.rebalance_meta.write().await;
rollback_local_rebalance_worker_activation(rebalance_meta.as_mut(), expected_id.as_ref(), &rx);
debug!(
event = EVENT_REBALANCE_STATE,
component = LOG_COMPONENT_ECSTORE,
@@ -201,13 +360,14 @@ impl ECStore {
}
#[tracing::instrument(skip(self, rx))]
async fn rebalance_buckets(self: &Arc<Self>, rx: CancellationToken, pool_index: usize) -> Result<()> {
async fn rebalance_buckets(self: &Arc<Self>, rx: CancellationToken, pool_index: usize, rebalance_id: Arc<str>) -> Result<()> {
ensure_valid_rebalance_pool_index(self.pools.len(), pool_index)?;
let (done_tx, mut done_rx) = tokio::sync::mpsc::channel::<Result<()>>(1);
// Save rebalance metadata periodically
let store = self.clone();
let save_rebalance_id = Arc::clone(&rebalance_id);
let save_task = tokio::spawn(async move {
let mut timer = tokio::time::interval_at(Instant::now() + Duration::from_secs(30), Duration::from_secs(10));
let mut msg: String;
@@ -221,6 +381,11 @@ impl ECStore {
let terminal_event = classify_rebalance_terminal_event(result, now);
msg = terminal_event.message().to_string();
let mut rebalance_meta = store.rebalance_meta.write().await;
super::control::ensure_rebalance_run_id(
rebalance_meta.as_ref(),
save_rebalance_id.as_ref(),
"apply rebalance terminal event",
)?;
if let Some(meta) = rebalance_meta.as_mut() {
let meta_stopped = meta.stopped_at.is_some();
if let Some(pool_stat) = meta.pool_stats.get_mut(pool_index) {
@@ -269,7 +434,10 @@ impl ECStore {
}
}
if let Err(err) = store.save_rebalance_stats(pool_index, RebalSaveOpt::Stats).await {
if let Err(err) = store
.save_rebalance_stats_for_id(pool_index, RebalSaveOpt::Stats, save_rebalance_id.as_ref())
.await
{
let wrapped = Error::other(format!("rebalance save_task stats save failed for pool {pool_index}: {err}"));
error!("{} err: {:?}", msg, wrapped);
if quit {
@@ -335,7 +503,7 @@ impl ECStore {
break;
}
let next_bucket = match self.next_rebal_bucket(pool_index).await {
let next_bucket = match self.next_rebal_bucket(pool_index, rebalance_id.as_ref()).await {
Ok(bucket) => bucket,
Err(err) => {
error!(
@@ -367,7 +535,8 @@ impl ECStore {
);
let outcome = match resolve_rebalance_bucket_result(
self.rebalance_bucket(rx.clone(), bucket.clone(), pool_index).await,
self.rebalance_bucket(rx.clone(), bucket.clone(), pool_index, Arc::clone(&rebalance_id))
.await,
pool_index,
&bucket,
) {
@@ -430,7 +599,7 @@ impl ECStore {
"Deferred rebalance bucket after transient object failures"
);
if let Err(err) = self
.defer_rebalance_bucket(pool_index, bucket.clone(), last_error.clone())
.defer_rebalance_bucket(pool_index, bucket.clone(), last_error.clone(), rebalance_id.as_ref())
.await
{
error!(
@@ -494,7 +663,7 @@ impl ECStore {
"Completed rebalance bucket"
);
source_cleanup_deferred_attempts.remove(&bucket);
if let Err(err) = self.bucket_rebalance_done(pool_index, bucket).await {
if let Err(err) = self.bucket_rebalance_done(pool_index, bucket, rebalance_id.as_ref()).await {
error!(
event = EVENT_REBALANCE_BUCKET,
component = LOG_COMPONENT_ECSTORE,
@@ -555,8 +724,9 @@ impl ECStore {
final_result
}
pub(super) async fn check_if_rebalance_done(&self, pool_index: usize) -> bool {
pub(super) async fn check_if_rebalance_done(&self, pool_index: usize, expected_id: &str) -> Result<bool> {
let mut rebalance_meta = self.rebalance_meta.write().await;
super::control::ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "check rebalance completion")?;
if let Some(meta) = rebalance_meta.as_mut()
&& let Some(pool_stat) = meta.pool_stats.get_mut(pool_index)
@@ -571,7 +741,7 @@ impl ECStore {
state = "already_completed",
"Rebalance pool is already completed"
);
return true;
return Ok(true);
}
// Mark pool rebalance as done only after it reaches the PercentFreeGoal.
@@ -601,19 +771,30 @@ impl ECStore {
percent_free = pfi,
"Marked rebalance pool completed"
);
return true;
return Ok(true);
}
}
false
Ok(false)
}
}
impl ECStore {
#[tracing::instrument(skip(self))]
pub async fn save_rebalance_stats(&self, pool_idx: usize, opt: RebalSaveOpt) -> Result<()> {
self.save_rebalance_stats_inner(pool_idx, opt, None).await
}
pub async fn save_rebalance_stats_for_id(&self, pool_idx: usize, opt: RebalSaveOpt, expected_id: &str) -> Result<()> {
self.save_rebalance_stats_inner(pool_idx, opt, Some(expected_id)).await
}
async fn save_rebalance_stats_inner(&self, pool_idx: usize, opt: RebalSaveOpt, expected_id: Option<&str>) -> Result<()> {
let meta_to_save = {
let mut rebalance_meta = self.rebalance_meta.write().await;
if let Some(expected_id) = expected_id {
super::control::ensure_rebalance_run_id(rebalance_meta.as_ref(), expected_id, "save rebalance stats")?;
}
let Some(meta) = rebalance_meta.as_mut() else {
return Ok(());
};
@@ -635,10 +816,14 @@ impl ECStore {
"Rebalance metadata save requested"
);
let stage = format!("save_rebalance_stats for pool {pool_idx} opt {opt:?}");
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, &meta_to_save, stage.as_str()).await,
stage.as_str(),
)?;
let save_result = match expected_id {
Some(expected_id) => {
self.save_rebalance_meta_for_id_with_merge(pool, &meta_to_save, stage.as_str(), expected_id)
.await
}
None => self.save_rebalance_meta_with_merge(pool, &meta_to_save, stage.as_str()).await,
};
resolve_rebalance_meta_save_result(save_result, stage.as_str())?;
Ok(())
}
@@ -144,6 +144,8 @@ pub struct RebalanceMeta {
#[serde(skip)]
pub cancel: Option<CancellationToken>, // To be invoked on rebalance-stop
#[serde(skip)]
pub activation_gate: std::sync::Arc<tokio::sync::RwLock<()>>,
#[serde(skip)]
pub last_refreshed_at: Option<OffsetDateTime>,
#[serde(rename = "stopTs")]
pub stopped_at: Option<OffsetDateTime>, // Time when rebalance-stop was issued
@@ -100,17 +100,17 @@ pub(super) fn resolve_rebalance_meta_save_result(result: Result<()>, stage: &str
result.map_err(|err| Error::other(format!("rebalance meta save failed during {stage}: {err}")))
}
pub(super) fn rebalance_meta_lock_error(err: rustfs_lock::LockError) -> Error {
pub(super) fn rebalance_meta_lock_error(err: rustfs_lock::LockError, mode: &'static str) -> Error {
match err {
rustfs_lock::LockError::QuorumNotReached { required, achieved } => Error::NamespaceLockQuorumUnavailable {
mode: "write",
mode,
bucket: crate::disk::RUSTFS_META_BUCKET.to_string(),
object: REBAL_META_NAME.to_string(),
required,
achieved,
},
other => Error::other(format!(
"failed to acquire rebalance metadata write lock on {}/{}: {other}",
"failed to acquire rebalance metadata {mode} lock on {}/{}: {other}",
crate::disk::RUSTFS_META_BUCKET,
REBAL_META_NAME
)),
+80
View File
@@ -735,6 +735,9 @@ pub(crate) use core::io_primitives::disk_call_counters;
mod ctx;
mod metadata;
mod ops;
#[cfg(test)]
pub(crate) use ops::hermetic_set_disks_isolated;
#[cfg(test)]
pub(crate) use ops::multipart::NewMultipartUploadCommitObservation;
#[cfg(any(test, feature = "test-util"))]
@@ -1449,6 +1452,21 @@ mod prepared_get_object_metadata_tests {
}
impl SetDisks {
#[cfg(test)]
async fn pause_tiered_metadata_commit(bucket: &str, object: &str) {
let barrier = TIERED_METADATA_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("tiered metadata commit barrier should not be poisoned")
.as_ref()
.filter(|barrier| barrier.bucket == bucket && barrier.object == object)
.cloned();
if let Some(barrier) = barrier {
barrier.arrived.notify_one();
barrier.release.notified().await;
}
}
pub(crate) async fn prepare_get_object_metadata(
&self,
bucket: &str,
@@ -4717,6 +4735,8 @@ impl SetDisks {
)?;
let fi = build_tiered_decommission_file_info(bucket, object, fi, layout);
let write_quorum = layout.write_quorum;
#[cfg(test)]
Self::pause_tiered_metadata_commit(bucket, object).await;
if _lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost())
|| opts
.namespace_lock_fence
@@ -4764,6 +4784,66 @@ impl SetDisks {
}
}
#[cfg(test)]
struct TieredMetadataCommitBarrierState {
bucket: String,
object: String,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
pub(crate) struct TieredMetadataCommitBarrier {
state: Arc<TieredMetadataCommitBarrierState>,
}
#[cfg(test)]
static TIERED_METADATA_COMMIT_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<TieredMetadataCommitBarrierState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
impl TieredMetadataCommitBarrier {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(TieredMetadataCommitBarrierState {
bucket: bucket.to_string(),
object: object.to_string(),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let mut slot = TIERED_METADATA_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("tiered metadata commit barrier should not be poisoned");
assert!(slot.is_none(), "tiered metadata commit barrier must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("tiered metadata write should reach its deterministic commit barrier");
}
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
#[cfg(test)]
impl Drop for TieredMetadataCommitBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut slot = TIERED_METADATA_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("tiered metadata commit barrier should not be poisoned");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[derive(Debug, PartialEq, Eq)]
struct ObjProps {
successor_mod_time: Option<OffsetDateTime>,
+3
View File
@@ -25,3 +25,6 @@ pub(crate) mod list;
pub(crate) mod locking;
pub(crate) mod multipart;
pub(crate) mod object;
#[cfg(test)]
pub(crate) use object::hermetic_set_disks_support::hermetic_set_disks_isolated;
+42 -1
View File
@@ -86,6 +86,8 @@ struct MultipartCommitBarrierState {
pause: MultipartCommitPause,
expected_arrivals: usize,
arrivals: AtomicUsize,
#[cfg(test)]
committed: AtomicBool,
arrived: tokio::sync::Notify,
release: tokio::sync::Semaphore,
}
@@ -113,6 +115,8 @@ impl MultipartCommitBarrier {
pause,
expected_arrivals,
arrivals: AtomicUsize::new(0),
#[cfg(test)]
committed: AtomicBool::new(false),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Semaphore::new(0),
});
@@ -143,6 +147,11 @@ impl MultipartCommitBarrier {
pub fn release(&self) {
self.state.release.add_permits(self.state.expected_arrivals);
}
#[cfg(test)]
pub(crate) fn commit_observed(&self) -> bool {
self.state.committed.load(Ordering::Acquire)
}
}
#[cfg(any(test, feature = "test-util"))]
@@ -263,6 +272,20 @@ async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartComm
}
}
#[cfg(test)]
fn observe_multipart_commit(bucket: &str, object: &str, pause: MultipartCommitPause) {
let slot = MULTIPART_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("multipart commit barrier mutex should not poison");
if let Some(barrier) = slot
.as_ref()
.filter(|barrier| barrier.bucket == bucket && barrier.object == object && barrier.pause == pause)
{
barrier.committed.store(true, Ordering::Release);
}
}
fn map_upload_id_metadata_error(bucket: &str, object: &str, upload_id: &str, err: DiskError) -> Error {
if err == DiskError::FileNotFound {
return StorageError::InvalidUploadID(bucket.to_owned(), object.to_owned(), upload_id.to_owned());
@@ -1320,6 +1343,19 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
pause_multipart_commit(bucket, object, MultipartCommitPause::PutPartBeforeLockLost).await;
fence_commit_on_lock_loss(_upload_commit_guard.as_ref(), "put_object_part_commit", &upload_id_path)?;
fence_commit_on_lock_loss(_part_commit_guard.as_ref(), "put_object_part_commit", &part_lock_path)?;
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "put_object_part_outer_lock",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?;
let _ = self
@@ -1340,6 +1376,8 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
}),
)
.await?;
#[cfg(test)]
observe_multipart_commit(bucket, object, MultipartCommitPause::PutPartBeforeLockLost);
#[cfg(test)]
pause_multipart_commit(bucket, object, MultipartCommitPause::PutPartAfterRename).await;
@@ -1720,7 +1758,10 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
.await
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
#[cfg(test)]
observe_new_multipart_upload_commit(bucket, object);
{
observe_multipart_commit(bucket, object, MultipartCommitPause::NewUploadBeforeLockLost);
observe_new_multipart_upload_commit(bucket, object);
}
// evalDisks
+4 -5
View File
@@ -6519,6 +6519,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
};
let find_vid = Uuid::new_v4();
#[cfg(test)]
pause_delete_object_commit(bucket, object).await;
if mark_delete && (opts.versioned || opts.version_suspended) {
if !delete_marker {
@@ -6587,8 +6589,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
dfi.set_skip_tier_free_version();
}
#[cfg(test)]
pause_delete_object_commit(bucket, object).await;
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
self.delete_object_version(bucket, object, &dfi, opts.delete_marker)
.await
@@ -7750,9 +7750,7 @@ pub(in crate::set_disk::ops) mod hermetic_set_disks_support {
/// for tests that never touch context-resolved services registered on the
/// ambient context (tier config manager, expiry state, ...), because the
/// isolated context starts every one of those cells fresh.
pub(in crate::set_disk::ops) async fn hermetic_set_disks_isolated(
disk_count: usize,
) -> (Vec<TempDir>, Vec<DiskStore>, Arc<SetDisks>) {
pub(crate) async fn hermetic_set_disks_isolated(disk_count: usize) -> (Vec<TempDir>, Vec<DiskStore>, Arc<SetDisks>) {
hermetic_set_disks_for_pool_with_default_parity_isolated(disk_count, 0, disk_count / 2).await
}
@@ -11815,6 +11813,7 @@ mod transition_upload_integrity_tests {
crate::data_movement::SourceCleanupBucketFence {
expected_incarnation_id: None,
lifecycle_guard: Some(&bucket_guard),
namespace_lock_lost_signal: None,
..Default::default()
},
"test_data_movement",
+54 -513
View File
@@ -13,10 +13,10 @@
// limitations under the License.
use crate::heal::{
progress::{HealProgress, add_bytes, increment_counter},
progress::HealProgress,
resume::{
CheckpointManager, CheckpointObjectOutcome, CheckpointObjectOutcomeRecord, ReplacementTargetIdentity, ResumeManager,
ResumeUtils, compose_key, replacement_target_identities_match,
CheckpointManager, ReplacementTargetIdentity, ResumeManager, ResumeUtils, compose_key,
replacement_target_identities_match,
},
storage::{HealStorageAPI, next_heal_listing_token},
task::{demote_to_debug_when, is_missing_object_dir_heal_result, take_failure_log_sample},
@@ -410,9 +410,6 @@ impl ErasureSetHealer {
&& state.successful_objects == 0
&& state.failed_objects == 0
&& state.skipped_objects == 0
&& state.skipped_new_versions == 0
&& state.skipped_ilm_expired == 0
&& state.processed_bytes == 0
{
// schedule_retry persists the authoritative resume reset before
// resetting the checkpoint. Reapply the checkpoint reset after
@@ -477,23 +474,6 @@ impl ErasureSetHealer {
// 2. initialize progress
self.initialize_progress(buckets, &state).await;
let (baseline_known, baseline_count, baseline_size, baseline_generation) = {
let baseline = self.progress.read().await;
(
baseline.baseline_known,
baseline.objects_total_count,
baseline.objects_total_size,
baseline.baseline_generation,
)
};
if baseline_known {
resume_manager
.set_progress_baseline(baseline_count, baseline_size, baseline_generation)
.await?;
checkpoint_manager
.set_progress_baseline(baseline_count, baseline_size, baseline_generation)
.await?;
}
// 3. continue from checkpoint
let current_bucket_index = checkpoint.current_bucket_index;
@@ -503,66 +483,12 @@ impl ErasureSetHealer {
let mut successful_objects = state.successful_objects;
let mut failed_objects = state.failed_objects;
let mut skipped_objects = state.skipped_objects;
let checkpoint_has_progress = checkpoint.baseline_known
|| checkpoint.successful_objects > 0
|| checkpoint.failed_object_count > 0
|| checkpoint.skipped_object_count > 0
|| checkpoint.skipped_new_versions > 0
|| checkpoint.skipped_ilm_expired > 0
|| checkpoint.processed_bytes > 0
|| checkpoint.total_objects > 0
|| checkpoint.total_bytes > 0
|| checkpoint.baseline_generation.is_some()
|| checkpoint.counter_unknown;
let checkpoint_generation_mismatch = checkpoint.baseline_known && checkpoint.baseline_generation != baseline_generation;
let mut restored_counter_unknown = state.counter_unknown || checkpoint.counter_unknown;
if checkpoint_has_progress {
successful_objects = checkpoint.successful_objects;
failed_objects = checkpoint.failed_object_count;
skipped_objects = checkpoint.skipped_object_count;
let restored_processed_objects = successful_objects
.checked_add(failed_objects)
.and_then(|value| value.checked_add(skipped_objects))
.and_then(|value| value.checked_add(checkpoint.skipped_new_versions))
.and_then(|value| value.checked_add(checkpoint.skipped_ilm_expired));
let checkpoint_counter_overflow = restored_processed_objects.is_none();
restored_counter_unknown |= checkpoint_counter_overflow;
processed_objects = restored_processed_objects.unwrap_or(u64::MAX);
let mut progress = self.progress.write().await;
progress.objects_scanned = processed_objects;
progress.objects_healed = successful_objects;
progress.objects_failed = failed_objects;
progress.skipped_objects = skipped_objects;
progress.skipped_new_versions = checkpoint.skipped_new_versions;
progress.skipped_ilm_expired = checkpoint.skipped_ilm_expired;
if checkpoint.baseline_known && !checkpoint_generation_mismatch {
progress.objects_total_count = checkpoint.total_objects;
progress.objects_total_size = checkpoint.total_bytes;
progress.baseline_generation = checkpoint.baseline_generation;
progress.baseline_known = true;
}
progress.bytes_processed = checkpoint.processed_bytes;
progress.counter_unknown = state.counter_unknown || checkpoint.counter_unknown;
progress.refresh_progress_percentage();
if checkpoint_generation_mismatch || checkpoint_counter_overflow || progress.counter_unknown {
progress.mark_unknown();
}
}
if checkpoint_generation_mismatch {
restored_counter_unknown = true;
}
if restored_counter_unknown {
checkpoint_manager.mark_counter_unknown().await?;
resume_manager.mark_counter_unknown().await?;
}
let mut failed_buckets = 0u64;
// 4. process remaining buckets
for (bucket_idx, bucket) in buckets.iter().enumerate().skip(current_bucket_index) {
// check if completed
if state.completed_buckets.contains(bucket) {
checkpoint_manager.complete_bucket(bucket_idx.saturating_add(1)).await?;
current_object_index = 0;
continue;
}
@@ -590,42 +516,13 @@ impl ErasureSetHealer {
return bucket_result;
}
// update progress
let progress_snapshot = self.progress.read().await;
let bytes_processed = progress_snapshot.bytes_processed;
let skipped_new_versions = progress_snapshot.skipped_new_versions;
let skipped_ilm_expired = progress_snapshot.skipped_ilm_expired;
let counter_unknown = progress_snapshot.counter_unknown;
drop(progress_snapshot);
// The checkpoint is the recovery authority for object progress.
// Publish its counters and fence before the resume summary so a
// crash between the two stores cannot make recovery select newer
// summary bytes with an older checkpoint ledger.
if counter_unknown {
checkpoint_manager.mark_counter_unknown().await?;
}
checkpoint_manager
.update_progress(successful_objects, failed_objects, skipped_objects, bytes_processed)
.await?;
checkpoint_manager
.set_skipped_version_counts(skipped_new_versions, skipped_ilm_expired)
.await?;
// update checkpoint position
checkpoint_manager.update_position(bucket_idx, current_object_index).await?;
// update progress
resume_manager
.update_progress_with_bytes(
processed_objects,
successful_objects,
failed_objects,
skipped_objects,
bytes_processed,
)
.update_progress(processed_objects, successful_objects, failed_objects, skipped_objects)
.await?;
resume_manager
.set_skipped_version_counts(skipped_new_versions, skipped_ilm_expired)
.await?;
if counter_unknown {
resume_manager.mark_counter_unknown().await?;
}
// check cancel status
if self.cancel_token.is_cancelled() {
@@ -645,7 +542,6 @@ impl ErasureSetHealer {
match bucket_result {
Ok(_) => {
resume_manager.complete_bucket(bucket).await?;
checkpoint_manager.complete_bucket(bucket_idx.saturating_add(1)).await?;
debug!(
target: "rustfs::heal::erasure_healer",
event = EVENT_HEAL_ERASURE_BUCKET_STATE,
@@ -671,9 +567,7 @@ impl ErasureSetHealer {
error = %e,
"Erasure set bucket heal failed"
);
// A single durable cursor and ledger cannot safely preserve
// this bucket while processing a later one.
break;
// continue to next bucket, do not interrupt the whole process
}
}
@@ -881,49 +775,20 @@ impl ErasureSetHealer {
// Per-version dedup identity — the single canonical key.
let key = compose_key(&item.name, item.version_id.as_deref());
if checkpoint.processed_objects.contains(&key)
|| checkpoint.failed_objects.contains(&key)
|| checkpoint.skipped_objects.contains(&key)
{
if checkpoint.processed_objects.contains(&key) || checkpoint.skipped_objects.contains(&key) {
continue;
}
if should_skip_new_version(item.mod_time_unix_nanos, started_at_secs) {
let counter_ok = increment_counter(processed_objects);
checkpoint_manager.add_processed_object(key).await?;
*processed_objects = processed_objects.saturating_add(1);
completed_in_page = completed_in_page.saturating_add(1);
counter!("rustfs_heal_skipped_new_versions_total").increment(1);
let (outcome_record, counter_unknown) = {
{
let mut progress = self.progress.write().await;
progress.record_skipped_new_version();
progress.set_current_object(Some(format!("skipped_new: {bucket}/{}", item.name)));
progress.update_object_progress(
*processed_objects,
*successful_objects,
*failed_objects,
*skipped_objects,
bytes_processed,
);
if !counter_ok {
progress.mark_unknown();
}
(
CheckpointObjectOutcomeRecord {
object: key,
outcome: CheckpointObjectOutcome::Processed,
successful: progress.objects_healed,
failed: progress.objects_failed,
skipped: progress.skipped_objects,
bytes: progress.bytes_processed,
skipped_new_versions: progress.skipped_new_versions,
skipped_ilm_expired: progress.skipped_ilm_expired,
counter_unknown: progress.counter_unknown,
},
progress.counter_unknown,
)
};
checkpoint_manager.record_object_outcome(outcome_record).await?;
if counter_unknown {
resume_manager.mark_counter_unknown().await?;
progress.update_progress(*processed_objects, *successful_objects, *failed_objects, bytes_processed);
}
debug!(
target: "rustfs::heal::erasure_healer",
@@ -955,41 +820,15 @@ impl ErasureSetHealer {
)
.await?
{
let counter_ok = increment_counter(processed_objects);
checkpoint_manager.add_processed_object(key).await?;
*processed_objects = processed_objects.saturating_add(1);
completed_in_page = completed_in_page.saturating_add(1);
counter!("rustfs_heal_skipped_ilm_expired_total").increment(1);
let (outcome_record, counter_unknown) = {
{
let mut progress = self.progress.write().await;
progress.record_skipped_ilm_expired();
progress.set_current_object(Some(format!("skipped_ilm: {bucket}/{}", item.name)));
progress.update_object_progress(
*processed_objects,
*successful_objects,
*failed_objects,
*skipped_objects,
bytes_processed,
);
if !counter_ok {
progress.mark_unknown();
}
(
CheckpointObjectOutcomeRecord {
object: key,
outcome: CheckpointObjectOutcome::Processed,
successful: progress.objects_healed,
failed: progress.objects_failed,
skipped: progress.skipped_objects,
bytes: progress.bytes_processed,
skipped_new_versions: progress.skipped_new_versions,
skipped_ilm_expired: progress.skipped_ilm_expired,
counter_unknown: progress.counter_unknown,
},
progress.counter_unknown,
)
};
checkpoint_manager.record_object_outcome(outcome_record).await?;
if counter_unknown {
resume_manager.mark_counter_unknown().await?;
progress.update_progress(*processed_objects, *successful_objects, *failed_objects, bytes_processed);
}
debug!(
target: "rustfs::heal::erasure_healer",
@@ -1115,11 +954,11 @@ impl ErasureSetHealer {
while let Some((key, object, version_id, result)) = page_tasks.next().await {
let (object_size, result) = result;
let mut telemetry_unknown = false;
let checkpoint_outcome = match result {
match result {
Ok(true) => {
telemetry_unknown |= !increment_counter(successful_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size);
*successful_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
checkpoint_manager.add_processed_object(key).await?;
debug!(
target: "rustfs::heal::erasure_healer",
event = EVENT_HEAL_ERASURE_OBJECT_STATE,
@@ -1132,11 +971,11 @@ impl ErasureSetHealer {
state = "healed",
"Erasure set object healed"
);
CheckpointObjectOutcome::Processed
}
Ok(false) => {
telemetry_unknown |= !increment_counter(successful_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size);
checkpoint_manager.add_processed_object(key).await?;
*successful_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
debug!(
target: "rustfs::heal::erasure_healer",
event = EVENT_HEAL_ERASURE_OBJECT_STATE,
@@ -1149,12 +988,12 @@ impl ErasureSetHealer {
state = "missing_treated_as_ok",
"Erasure set missing object treated as ok"
);
CheckpointObjectOutcome::Processed
}
Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err),
Err(Error::TransientSkip { message }) => {
telemetry_unknown |= !increment_counter(skipped_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size);
*skipped_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
checkpoint_manager.add_skipped_object(key).await?;
demote_to_debug_when!(!take_failure_log_sample(&mut transient_skip_samples_logged), warn, target: "rustfs::heal::erasure_healer", {
event = EVENT_HEAL_ERASURE_OBJECT_STATE,
component = LOG_COMPONENT_HEAL,
@@ -1167,11 +1006,11 @@ impl ErasureSetHealer {
error = %message,
"Erasure set object heal skipped due to transient error"
});
CheckpointObjectOutcome::Skipped
}
Err(err) => {
telemetry_unknown |= !increment_counter(failed_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size);
*failed_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
checkpoint_manager.add_failed_object(key).await?;
demote_to_debug_when!(!take_failure_log_sample(&mut failure_samples_logged), warn, target: "rustfs::heal::erasure_healer", {
event = EVENT_HEAL_ERASURE_OBJECT_STATE,
component = LOG_COMPONENT_HEAL,
@@ -1184,43 +1023,15 @@ impl ErasureSetHealer {
error = %err,
"Erasure set object heal failed"
});
CheckpointObjectOutcome::Failed
}
};
}
telemetry_unknown |= !increment_counter(processed_objects);
*processed_objects += 1;
completed_in_page += 1;
let (outcome_record, counter_unknown) = {
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("{bucket}/{object}")));
progress.update_object_progress(
*processed_objects,
*successful_objects,
*failed_objects,
*skipped_objects,
bytes_processed,
);
if telemetry_unknown {
progress.mark_unknown();
}
(
CheckpointObjectOutcomeRecord {
object: key,
outcome: checkpoint_outcome,
successful: progress.objects_healed,
failed: progress.objects_failed,
skipped: progress.skipped_objects,
bytes: progress.bytes_processed,
skipped_new_versions: progress.skipped_new_versions,
skipped_ilm_expired: progress.skipped_ilm_expired,
counter_unknown: progress.counter_unknown,
},
progress.counter_unknown,
)
};
checkpoint_manager.record_object_outcome(outcome_record).await?;
if counter_unknown {
resume_manager.mark_counter_unknown().await?;
progress.update_progress(*processed_objects, *successful_objects, *failed_objects, bytes_processed);
}
if completed_in_page.is_multiple_of(100) {
@@ -1230,22 +1041,16 @@ impl ErasureSetHealer {
*current_object_index = global_obj_idx;
// Persist the checkpoint ledger and page position before exposing
// the next resume cursor. A crash before cursor publication keeps
// the page identities available for exact-once replay.
checkpoint_manager.advance_page(bucket_index, *current_object_index).await?;
// Persist the authoritative cursor FIRST (points at the next page
// boundary), then prune the per-version dedup sets. Both are
// idempotent under crash: heal_object re-heals safely.
let next_cursor = if is_truncated { next_token.clone() } else { None };
resume_manager.set_resume_cursor(next_cursor.clone()).await?;
checkpoint_manager.complete_page(bucket_index, *current_object_index).await?;
// Check if there are more pages
if !is_truncated {
break;
}
continuation_token = next_heal_listing_token(bucket, "", next_token, is_truncated)?;
if continuation_token.is_none() {
// A truncated page without a continuation token is terminal.
// Retain its ledger until bucket completion is durable.
break;
}
resume_manager.set_resume_cursor(continuation_token.clone()).await?;
checkpoint_manager.prune_completed_page().await?;
// Anti-loop guard: an empty page reported as truncated cannot advance
// the cursor (there is no last identity to move past), so treat it as a
@@ -1264,6 +1069,12 @@ impl ErasureSetHealer {
)));
}
previous_page_last = page_last;
continuation_token = next_heal_listing_token(bucket, "", next_token, is_truncated)?;
if continuation_token.is_none() {
// Truncated but no continuation token: treat as end of listing.
break;
}
}
Ok(())
@@ -1272,66 +1083,10 @@ impl ErasureSetHealer {
/// initialize progress tracking
async fn initialize_progress(&self, _buckets: &[String], state: &crate::heal::resume::ResumeState) {
let mut progress = self.progress.write().await;
let existing_baseline = (
progress.objects_total_count,
progress.objects_total_size,
progress.baseline_generation,
progress.progress_state,
progress.baseline_known,
);
let baseline_generation_mismatch =
state.baseline_known && existing_baseline.4 && state.baseline_generation != existing_baseline.2;
let use_persisted_baseline = state.baseline_known && !baseline_generation_mismatch;
progress.objects_scanned = state.processed_objects;
progress.objects_scanned = state.total_objects;
progress.objects_healed = state.successful_objects;
progress.objects_failed = state.failed_objects;
progress.skipped_objects = state.skipped_objects;
progress.skipped_new_versions = state.skipped_new_versions;
progress.skipped_ilm_expired = state.skipped_ilm_expired;
progress.bytes_processed = state.processed_bytes;
progress.counter_unknown = state.counter_unknown;
if use_persisted_baseline
|| existing_baseline.0 > 0
|| existing_baseline.1 > 0
|| existing_baseline.2.is_some()
|| existing_baseline.4
{
progress.objects_total_count = if use_persisted_baseline {
state.total_objects
} else {
existing_baseline.0
};
progress.objects_total_size = if use_persisted_baseline {
state.total_bytes
} else {
existing_baseline.1
};
progress.baseline_generation = if use_persisted_baseline {
state.baseline_generation
} else {
existing_baseline.2
};
progress.baseline_known = use_persisted_baseline
|| existing_baseline.0 > 0
|| existing_baseline.1 > 0
|| existing_baseline.2.is_some()
|| existing_baseline.4;
}
progress.progress_state = if use_persisted_baseline
|| existing_baseline.0 > 0
|| existing_baseline.1 > 0
|| existing_baseline.2.is_some()
|| existing_baseline.4
{
crate::heal::progress::HealProgressState::Running
} else {
crate::heal::progress::HealProgressState::Indeterminate
};
if baseline_generation_mismatch || state.counter_unknown {
progress.mark_unknown();
}
progress.ledger_complete = false;
progress.refresh_progress_percentage();
progress.bytes_processed = 0; // Resume state tracks object counts, not byte counters.
progress.start_time = UNIX_EPOCH.checked_add(Duration::from_secs(state.start_time));
progress.last_update_time = UNIX_EPOCH.checked_add(Duration::from_secs(state.last_update));
progress.set_current_object(state.current_object.clone());
@@ -1509,8 +1264,8 @@ mod resume_loop_tests {
};
use crate::heal::progress::HealProgress;
use crate::heal::resume::{
CheckpointManager, CheckpointObjectOutcome, CheckpointObjectOutcomeRecord, RESUME_CHECKPOINT_FILE,
ReplacementTargetIdentity, ResumeDeleteFailure, ResumeManager, ResumeUtils, compose_key,
CheckpointManager, RESUME_CHECKPOINT_FILE, ReplacementTargetIdentity, ResumeDeleteFailure, ResumeManager, ResumeUtils,
compose_key,
};
use crate::heal::storage::{HealLifecycleExpiryContext, HealListItem, HealObjectInfo, HealStorageAPI};
use crate::heal::storage_api::status::BucketInfo;
@@ -1650,7 +1405,6 @@ mod resume_loop_tests {
list_include_lifecycle_object_info: Mutex<Vec<bool>>,
replacement_target_identity_sequences: Mutex<VecDeque<Vec<ReplacementTargetIdentity>>>,
fail_listing: AtomicBool,
fail_listing_buckets: Mutex<HashSet<String>>,
}
impl FakeStorage {
@@ -1687,9 +1441,6 @@ mod resume_loop_tests {
fn fail_listing(&self) {
self.fail_listing.store(true, Ordering::SeqCst);
}
fn fail_bucket_listing(&self, bucket: &str) {
self.fail_listing_buckets.lock().unwrap().insert(bucket.to_string());
}
}
#[async_trait::async_trait]
@@ -1780,7 +1531,7 @@ mod resume_loop_tests {
}
async fn list_objects_for_heal_page(
&self,
bucket: &str,
_bucket: &str,
_prefix: &str,
continuation_token: Option<&str>,
include_lifecycle_object_info: bool,
@@ -1789,7 +1540,7 @@ mod resume_loop_tests {
.lock()
.unwrap()
.push(include_lifecycle_object_info);
if self.fail_listing.load(Ordering::SeqCst) || self.fail_listing_buckets.lock().unwrap().contains(bucket) {
if self.fail_listing.load(Ordering::SeqCst) {
return Err(Error::other("injected listing failure"));
}
let key = continuation_token.map(str::to_string);
@@ -2167,49 +1918,6 @@ mod resume_loop_tests {
assert!(state.completed_buckets.is_empty(), "the failed bucket must remain resumable");
}
#[tokio::test]
async fn bucket_failure_stops_before_a_later_bucket_checkpoint() {
let env = make_env().await;
let task_id = ResumeUtils::generate_task_id();
let buckets = vec!["a".to_string(), "b".to_string()];
let resume = ResumeManager::new(
env.healer.disk.clone(),
task_id.clone(),
"erasure_set".to_string(),
"pool_0_set_0".to_string(),
buckets.clone(),
)
.await
.unwrap();
let checkpoint = CheckpointManager::new(env.healer.disk.clone(), task_id.clone())
.await
.unwrap();
env.storage.fail_bucket_listing("a");
for _ in 0..3 {
assert!(resume.schedule_retry().await.unwrap());
}
env.healer
.execute_heal_with_resume(&buckets, "pool_0_set_0", &resume, &checkpoint)
.await
.expect_err("the first bucket failure must keep the pass incomplete");
let persisted = checkpoint.get_checkpoint().await;
assert_eq!(persisted.current_bucket_index, 0);
assert!(resume.get_state().await.completed_buckets.is_empty());
let resumed = ResumeManager::load_from_disk(env.healer.disk.clone(), &task_id)
.await
.unwrap();
let checkpoint = CheckpointManager::load_from_disk(env.healer.disk.clone(), &task_id)
.await
.unwrap();
env.healer
.execute_heal_with_resume(&buckets, "pool_0_set_0", &resumed, &checkpoint)
.await
.expect_err("recovery must retry the earlier failed bucket");
assert!(!resumed.get_state().await.completed);
}
#[tokio::test]
async fn completed_resume_state_is_not_selected_for_a_new_heal() {
let env = make_env().await;
@@ -2399,175 +2107,8 @@ mod resume_loop_tests {
let mut names: Vec<String> = env.storage.calls().into_iter().map(|(n, _)| n).collect();
names.sort();
assert_eq!(names, vec!["a", "b", "c", "d"], "every object exactly once, none dropped/doubled");
// Keep the final page cursor until the outer loop durably completes the
// bucket, so a crash can replay only this page against its identities.
assert_eq!(env.resume.resume_cursor().await, Some("t1".to_string()));
}
#[tokio::test]
async fn persisted_failure_waits_for_the_bounded_retry_after_page_replay() {
let env = make_env().await;
env.storage.set_page(
None,
Page {
items: vec![item("object", Some("v1"), false)],
next: None,
truncated: false,
},
);
env.checkpoint
.record_object_outcome(CheckpointObjectOutcomeRecord {
object: compose_key("object", Some("v1")),
outcome: CheckpointObjectOutcome::Failed,
successful: 0,
failed: 1,
skipped: 0,
bytes: 0,
skipped_new_versions: 0,
skipped_ilm_expired: 0,
counter_unknown: false,
})
.await
.unwrap();
env.checkpoint.advance_page(0, 1).await.unwrap();
let resumed = ResumeManager::load_from_disk(env.healer.disk.clone(), &env.task_id)
.await
.unwrap();
let checkpoint = CheckpointManager::load_from_disk(env.healer.disk.clone(), &env.task_id)
.await
.unwrap();
env.healer
.execute_heal_with_resume(&["b".to_string()], "pool_0_set_0", &resumed, &checkpoint)
.await
.expect_err("the persisted failure must schedule a bounded retry");
assert!(
env.storage.calls().is_empty(),
"the failed identity must not be repeated in the same pass"
);
env.healer
.execute_heal_with_resume(&["b".to_string()], "pool_0_set_0", &resumed, &checkpoint)
.await
.expect("the bounded retry must heal the object");
assert_eq!(env.storage.calls(), vec![("object".to_string(), Some("v1".to_string()))]);
let state = resumed.get_state().await;
assert_eq!(state.successful_objects, 1);
assert_eq!(state.failed_objects, 0);
}
#[tokio::test]
async fn final_page_crash_replays_only_the_retained_page_identities() {
let env = make_env().await;
env.storage.set_page(
None,
Page {
items: vec![item("first", Some("v1"), false)],
next: Some("final-page".to_string()),
truncated: true,
},
);
env.storage.set_page(
Some("final-page"),
Page {
items: vec![item("last", Some("v1"), false)],
next: None,
truncated: false,
},
);
let (processed, successful, failed, skipped, result) = run(&env).await;
result.expect("the bucket pass must finish before the simulated crash");
assert_eq!((processed, successful, failed, skipped), (2, 2, 0, 0));
let resumed = ResumeManager::load_from_disk(env.healer.disk.clone(), &env.task_id)
.await
.unwrap();
let checkpoint = CheckpointManager::load_from_disk(env.healer.disk.clone(), &env.task_id)
.await
.unwrap();
env.healer
.execute_heal_with_resume(&["b".to_string()], "pool_0_set_0", &resumed, &checkpoint)
.await
.expect("the retained final-page ledger must make recovery exact");
assert_eq!(
env.storage.calls(),
vec![
("first".to_string(), Some("v1".to_string())),
("last".to_string(), Some("v1".to_string()))
]
);
let state = resumed.get_state().await;
assert_eq!(state.successful_objects, 2);
assert_eq!(state.processed_objects, 2);
}
#[tokio::test]
async fn truncated_page_without_token_retains_its_replay_ledger() {
let env = make_env().await;
env.storage.set_page(
None,
Page {
items: vec![item("object", Some("v1"), false)],
next: None,
truncated: true,
},
);
let (processed, successful, failed, skipped, result) = run(&env).await;
result.expect("the tokenless truncated page is a terminal page");
assert_eq!((processed, successful, failed, skipped), (1, 1, 0, 0));
let resumed = ResumeManager::load_from_disk(env.healer.disk.clone(), &env.task_id)
.await
.unwrap();
let checkpoint = CheckpointManager::load_from_disk(env.healer.disk.clone(), &env.task_id)
.await
.unwrap();
env.healer
.execute_heal_with_resume(&["b".to_string()], "pool_0_set_0", &resumed, &checkpoint)
.await
.expect("terminal-page recovery must not replay a durable identity");
assert_eq!(env.storage.calls(), vec![("object".to_string(), Some("v1".to_string()))]);
let state = resumed.get_state().await;
assert_eq!(state.successful_objects, 1);
assert_eq!(state.processed_objects, 1);
}
#[tokio::test]
async fn completed_bucket_reconciles_its_final_page_checkpoint_after_crash() {
let env = make_env().await;
env.storage.set_page(
None,
Page {
items: vec![item("object", Some("v1"), false)],
next: None,
truncated: false,
},
);
let (_, _, _, _, result) = run(&env).await;
result.expect("the bucket pass must finish before the simulated crash");
env.resume.complete_bucket("b").await.unwrap();
let resumed = ResumeManager::load_from_disk(env.healer.disk.clone(), &env.task_id)
.await
.unwrap();
let checkpoint = CheckpointManager::load_from_disk(env.healer.disk.clone(), &env.task_id)
.await
.unwrap();
env.healer
.execute_heal_with_resume(&["b".to_string()], "pool_0_set_0", &resumed, &checkpoint)
.await
.expect("recovery must finish the checkpoint transition without replaying the bucket");
assert_eq!(env.storage.calls(), vec![("object".to_string(), Some("v1".to_string()))]);
let checkpoint = checkpoint.get_checkpoint().await;
assert_eq!(checkpoint.current_bucket_index, 1);
assert!(checkpoint.processed_objects.is_empty());
// Final page not truncated => cursor cleared.
assert_eq!(env.resume.resume_cursor().await, None);
}
#[tokio::test]
+26 -3
View File
@@ -2011,11 +2011,34 @@ impl HealManager {
return None;
}
let mut progresses = Vec::with_capacity(active_tasks.len());
let mut snapshot = HealProgress::default();
for task in active_tasks {
progresses.push(task.get_progress().await);
let progress = task.get_progress().await;
snapshot.objects_scanned = snapshot.objects_scanned.saturating_add(progress.objects_scanned);
snapshot.objects_healed = snapshot.objects_healed.saturating_add(progress.objects_healed);
snapshot.objects_failed = snapshot.objects_failed.saturating_add(progress.objects_failed);
snapshot.skipped_new_versions = snapshot.skipped_new_versions.saturating_add(progress.skipped_new_versions);
snapshot.skipped_ilm_expired = snapshot.skipped_ilm_expired.saturating_add(progress.skipped_ilm_expired);
snapshot.objects_total_count = snapshot.objects_total_count.saturating_add(progress.objects_total_count);
snapshot.objects_total_size = snapshot.objects_total_size.saturating_add(progress.objects_total_size);
snapshot.bytes_processed = snapshot.bytes_processed.saturating_add(progress.bytes_processed);
snapshot.start_time = match (snapshot.start_time, progress.start_time) {
(Some(current), Some(next)) => Some(current.min(next)),
(None, next) => next,
(current, None) => current,
};
snapshot.last_update_time = match (snapshot.last_update_time, progress.last_update_time) {
(Some(current), Some(next)) => Some(current.max(next)),
(None, next) => next,
(current, None) => current,
};
if progress.current_object.is_some() {
snapshot.current_object = progress.current_object;
}
}
crate::heal::progress::aggregate_heal_progress(progresses)
snapshot.refresh_progress_percentage();
snapshot.refresh_estimated_completion_time();
Some(snapshot)
}
}
+34 -475
View File
@@ -15,91 +15,15 @@
use serde::{Deserialize, Serialize};
use std::time::{Duration, SystemTime};
pub(crate) fn stable_generation(parts: &[&[u8]]) -> u64 {
let mut hash = 0xcbf29ce484222325u64;
for part in parts {
for byte in (part.len() as u64).to_be_bytes().into_iter().chain(part.iter().copied()) {
hash ^= u64::from(byte);
hash = hash.wrapping_mul(0x100000001b3);
}
}
hash
}
#[cfg(test)]
mod stable_generation_tests {
use super::stable_generation;
#[test]
fn stable_generation_has_a_fixed_vector() {
assert_eq!(stable_generation(&[b"rustfs", b"heal", b"42"]), 11_007_672_338_488_385_056);
}
}
pub(crate) fn increment_counter(counter: &mut u64) -> bool {
match counter.checked_add(1) {
Some(next) => {
*counter = next;
true
}
None => {
*counter = u64::MAX;
false
}
}
}
pub(crate) fn add_bytes(total: &mut u64, amount: u64) -> bool {
match total.checked_add(amount) {
Some(next) => {
*total = next;
true
}
None => {
*total = u64::MAX;
false
}
}
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub enum HealProgressKind {
#[default]
Unknown,
Stage,
ObjectSweep,
}
/// Whether the object ledger can produce a meaningful percentage.
///
/// A zero-valued baseline is not a completed scan: it means that no complete
/// usage snapshot was available. Keep this state explicit so callers do not
/// mistake the legacy `0.0` wire value for a measured zero-percent result.
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub enum HealProgressState {
#[default]
Unknown,
Indeterminate,
Running,
Completed,
}
#[derive(Debug, Default, Clone, PartialEq, Serialize, Deserialize)]
#[serde(default, rename_all = "camelCase")]
pub struct HealProgress {
#[serde(default)]
pub kind: HealProgressKind,
/// Objects scanned
pub objects_scanned: u64,
/// Objects healed
pub objects_healed: u64,
/// Objects failed
pub objects_failed: u64,
/// Versions deferred for a later retry pass.
#[serde(default)]
pub skipped_objects: u64,
/// Versions skipped because they were written after this heal started
pub skipped_new_versions: u64,
/// Versions skipped because lifecycle already selected them for expiry
@@ -120,38 +44,11 @@ pub struct HealProgress {
pub last_update_time: Option<SystemTime>,
/// Estimated completion time
pub estimated_completion_time: Option<SystemTime>,
/// Current stage number. Stage updates are intentionally independent from
/// the object ledger below.
#[serde(default)]
pub stage_current: u64,
/// Number of stages in the current task.
#[serde(default)]
pub stage_total: u64,
/// Explicitly distinguishes a missing usage baseline from measured 0%.
#[serde(default)]
pub progress_state: HealProgressState,
/// True only after the task's durable completion ledger was committed.
#[serde(default)]
pub ledger_complete: bool,
/// Generation of the usage snapshot used for the baseline, if available.
#[serde(default)]
pub baseline_generation: Option<u64>,
/// Whether the baseline was explicitly observed. This is separate from
/// the counters so a known empty scope (0 objects, 0 bytes) is not
/// confused with a legacy snapshot that omitted the baseline fields.
#[serde(default)]
pub baseline_known: bool,
/// Internal telemetry fence set when an aggregate counter overflows or
/// becomes inconsistent. It prevents a later refresh from fabricating a
/// percentage from the poisoned values.
#[serde(default)]
pub counter_unknown: bool,
}
impl HealProgress {
pub fn new() -> Self {
Self {
kind: HealProgressKind::Unknown,
start_time: Some(SystemTime::now()),
last_update_time: Some(SystemTime::now()),
..Default::default()
@@ -159,87 +56,12 @@ impl HealProgress {
}
pub fn update_progress(&mut self, scanned: u64, healed: u64, failed: u64, bytes: u64) {
self.update_object_sweep_progress(scanned, healed, failed, bytes);
}
pub fn update_object_sweep_progress(&mut self, scanned: u64, healed: u64, failed: u64, bytes: u64) {
self.kind = HealProgressKind::ObjectSweep;
self.objects_scanned = scanned;
self.objects_healed = healed;
self.objects_failed = failed;
self.bytes_processed = bytes;
self.last_update_time = Some(SystemTime::now());
let explicit_skipped = match self.skipped_new_versions.checked_add(self.skipped_ilm_expired) {
Some(value) => value,
None => {
self.mark_unknown();
0
}
};
let skipped = healed
.checked_add(failed)
.and_then(|value| value.checked_add(explicit_skipped))
.and_then(|value| scanned.checked_sub(value))
.unwrap_or(0);
self.update_object_progress(scanned, healed, failed, skipped, bytes);
}
/// Update task stage progress without modifying object counters.
pub fn update_stage(&mut self, current: u64, total: u64) {
let object_sweep_active = matches!(self.kind, HealProgressKind::ObjectSweep);
if !object_sweep_active {
self.kind = HealProgressKind::Stage;
}
self.ledger_complete = false;
self.stage_current = current.min(total);
self.stage_total = total;
if object_sweep_active {
self.last_update_time = Some(SystemTime::now());
self.refresh_progress_percentage();
return;
}
self.progress_state = if total == 0 {
HealProgressState::Indeterminate
} else {
HealProgressState::Running
};
self.progress_percentage = if total == 0 {
0.0
} else {
(current as f64 / total as f64 * 100.0).min(100.0)
};
self.last_update_time = Some(SystemTime::now());
}
/// Update the disjoint object ledger. `scanned` is the number of terminal
/// object outcomes and must equal healed + failed + deferred skipped plus
/// the two terminal skip classes. Overflow is a corrupt/unknown counter
/// state, not a reason to abort a completed heal.
pub fn update_object_progress(&mut self, scanned: u64, healed: u64, failed: u64, skipped: u64, bytes: u64) {
self.kind = HealProgressKind::ObjectSweep;
// `skipped` is the transient/deferred class. The two explicit skip
// counters are terminal classifications too, so include them in the
// same ledger without making callers maintain a second aggregate.
let outcomes = healed
.checked_add(failed)
.and_then(|value| value.checked_add(skipped))
.and_then(|value| value.checked_add(self.skipped_new_versions))
.and_then(|value| value.checked_add(self.skipped_ilm_expired));
self.objects_scanned = scanned;
self.objects_healed = healed;
self.objects_failed = failed;
self.skipped_objects = skipped;
self.bytes_processed = bytes;
self.last_update_time = Some(SystemTime::now());
self.ledger_complete = false;
if outcomes != Some(scanned) {
// Telemetry corruption must not abort a heal. Preserve the
// counters for diagnostics, but do not derive a percentage from a
// double-counted or overflowing ledger.
self.mark_unknown();
return;
}
self.refresh_progress_percentage();
self.refresh_estimated_completion_time();
}
@@ -247,88 +69,50 @@ impl HealProgress {
pub fn set_total_baseline(&mut self, objects_total_count: u64, objects_total_size: u64) {
self.objects_total_count = objects_total_count;
self.objects_total_size = objects_total_size;
self.baseline_known = true;
self.last_update_time = Some(SystemTime::now());
self.refresh_progress_percentage();
self.refresh_estimated_completion_time();
}
pub fn set_total_baseline_with_generation(&mut self, objects_total_count: u64, objects_total_size: u64, generation: u64) {
self.baseline_generation = Some(generation);
self.set_total_baseline(objects_total_count, objects_total_size);
}
pub fn record_skipped_new_version(&mut self) {
let Some(next) = self.skipped_new_versions.checked_add(1) else {
self.mark_unknown();
return;
};
self.skipped_new_versions = next;
self.skipped_new_versions = self.skipped_new_versions.saturating_add(1);
self.last_update_time = Some(SystemTime::now());
self.refresh_progress_percentage();
self.refresh_estimated_completion_time();
}
pub fn record_skipped_ilm_expired(&mut self) {
let Some(next) = self.skipped_ilm_expired.checked_add(1) else {
self.mark_unknown();
return;
};
self.skipped_ilm_expired = next;
self.skipped_ilm_expired = self.skipped_ilm_expired.saturating_add(1);
self.last_update_time = Some(SystemTime::now());
self.refresh_progress_percentage();
self.refresh_estimated_completion_time();
}
fn completed_for_baseline(&self) -> Option<u64> {
fn completed_for_baseline(&self) -> u64 {
self.objects_healed
.checked_add(self.objects_failed)?
.checked_add(self.skipped_objects)?
.checked_add(self.skipped_new_versions)?
.checked_add(self.skipped_ilm_expired)
.saturating_add(self.objects_failed)
.saturating_add(self.skipped_new_versions)
.saturating_add(self.skipped_ilm_expired)
}
pub(crate) fn refresh_progress_percentage(&mut self) {
if self.ledger_complete {
self.progress_state = HealProgressState::Completed;
self.progress_percentage = 100.0;
return;
}
if self.counter_unknown {
self.progress_state = HealProgressState::Unknown;
self.progress_percentage = 0.0;
return;
}
if !self.baseline_known {
self.progress_state = HealProgressState::Indeterminate;
self.progress_percentage = 0.0;
self.estimated_completion_time = None;
return;
}
if self.objects_total_size > 0 {
self.progress_percentage = ((self.bytes_processed as f64 / self.objects_total_size as f64) * 100.0).min(100.0);
self.progress_percentage = self.progress_percentage.min(99.999);
self.progress_state = HealProgressState::Running;
return;
}
if self.objects_total_count > 0 {
let Some(completed) = self.completed_for_baseline() else {
self.progress_state = HealProgressState::Unknown;
self.progress_percentage = 0.0;
return;
};
let completed = self.completed_for_baseline();
self.progress_percentage = ((completed as f64 / self.objects_total_count as f64) * 100.0).min(100.0);
self.progress_percentage = self.progress_percentage.min(99.999);
self.progress_state = HealProgressState::Running;
return;
}
if self.baseline_known {
self.progress_state = HealProgressState::Running;
self.progress_percentage = 0.0;
return;
let total = self
.objects_scanned
.saturating_add(self.objects_healed)
.saturating_add(self.objects_failed);
if total > 0 {
self.progress_percentage = (self.objects_healed as f64 / total as f64) * 100.0;
}
self.progress_state = HealProgressState::Indeterminate;
self.progress_percentage = 0.0;
}
pub fn set_current_object(&mut self, object: Option<String>) {
@@ -341,11 +125,7 @@ impl HealProgress {
self.estimated_completion_time = None;
return;
};
if self.is_completed()
|| self.progress_percentage <= 0.0
|| self.progress_percentage >= 100.0
|| self.bytes_processed == 0
{
if self.is_completed() || !(0.0..100.0).contains(&self.progress_percentage) || self.bytes_processed == 0 {
self.estimated_completion_time = None;
return;
}
@@ -362,39 +142,18 @@ impl HealProgress {
}
pub fn is_completed(&self) -> bool {
self.ledger_complete
}
/// Mark telemetry unknown while allowing the underlying heal operation to
/// continue. This is used for corrupt/overflowing counters at the
/// observability boundary; it must never turn a successful heal into an
/// execution error.
pub fn mark_unknown(&mut self) {
self.counter_unknown = true;
self.progress_state = HealProgressState::Unknown;
self.ledger_complete = false;
self.progress_percentage = 0.0;
self.estimated_completion_time = None;
self.last_update_time = Some(SystemTime::now());
}
/// Mark the object ledger terminal only after the enclosing task has
/// committed all durable resume state and cleanup fences.
pub fn mark_completed(&mut self) {
let telemetry_unknown = self.counter_unknown || self.progress_state == HealProgressState::Unknown;
self.ledger_complete = true;
if !telemetry_unknown {
self.progress_state = HealProgressState::Completed;
if self.progress_percentage >= 100.0 {
return true;
}
self.progress_percentage = 100.0;
self.last_update_time = Some(SystemTime::now());
self.estimated_completion_time = None;
if self.objects_total_count > 0 || self.objects_total_size > 0 {
return false;
}
self.objects_scanned > 0 && self.objects_healed.saturating_add(self.objects_failed) >= self.objects_scanned
}
pub fn get_success_rate(&self) -> f64 {
let Some(total) = self.objects_healed.checked_add(self.objects_failed) else {
return 0.0;
};
let total = self.objects_healed + self.objects_failed;
if total > 0 {
(self.objects_healed as f64 / total as f64) * 100.0
} else {
@@ -403,101 +162,6 @@ impl HealProgress {
}
}
pub fn aggregate_heal_progress(progresses: impl IntoIterator<Item = HealProgress>) -> Option<HealProgress> {
let mut snapshot = HealProgress::default();
let mut found = false;
let mut has_object_sweep = false;
let mut all_object_baselines_known = true;
let mut baseline_generation = None;
let mut baseline_generation_consistent = true;
let mut all_ledgers_complete = true;
let mut counter_overflow = false;
for progress in progresses {
found = true;
let object_sweep = matches!(progress.kind, HealProgressKind::ObjectSweep);
has_object_sweep |= object_sweep;
all_ledgers_complete &= progress.ledger_complete;
if object_sweep {
all_object_baselines_known &= progress.baseline_known;
match baseline_generation {
None => baseline_generation = Some(progress.baseline_generation),
Some(generation) => baseline_generation_consistent &= generation == progress.baseline_generation,
}
}
counter_overflow |= progress.counter_unknown || matches!(progress.progress_state, HealProgressState::Unknown);
for (target, value) in [
(&mut snapshot.objects_scanned, progress.objects_scanned),
(&mut snapshot.objects_healed, progress.objects_healed),
(&mut snapshot.objects_failed, progress.objects_failed),
(&mut snapshot.skipped_objects, progress.skipped_objects),
(&mut snapshot.skipped_new_versions, progress.skipped_new_versions),
(&mut snapshot.skipped_ilm_expired, progress.skipped_ilm_expired),
(&mut snapshot.objects_total_count, progress.objects_total_count),
(&mut snapshot.objects_total_size, progress.objects_total_size),
(&mut snapshot.bytes_processed, progress.bytes_processed),
(&mut snapshot.stage_current, progress.stage_current),
(&mut snapshot.stage_total, progress.stage_total),
] {
match target.checked_add(value) {
Some(sum) => *target = sum,
None => {
*target = u64::MAX;
counter_overflow = true;
}
}
}
snapshot.start_time = match (snapshot.start_time, progress.start_time) {
(Some(current), Some(next)) => Some(current.min(next)),
(None, next) => next,
(current, None) => current,
};
snapshot.last_update_time = match (snapshot.last_update_time, progress.last_update_time) {
(Some(current), Some(next)) => Some(current.max(next)),
(None, next) => next,
(current, None) => current,
};
if progress.current_object.is_some() {
snapshot.current_object = progress.current_object;
}
}
if !found {
return None;
}
snapshot.kind = if has_object_sweep {
HealProgressKind::ObjectSweep
} else {
HealProgressKind::Stage
};
snapshot.baseline_known = has_object_sweep && all_object_baselines_known && baseline_generation_consistent;
snapshot.baseline_generation = if snapshot.baseline_known && baseline_generation_consistent {
baseline_generation.flatten()
} else {
None
};
snapshot.ledger_complete = all_ledgers_complete;
snapshot.counter_unknown = counter_overflow;
if counter_overflow {
snapshot.progress_state = HealProgressState::Unknown;
snapshot.progress_percentage = if snapshot.ledger_complete { 100.0 } else { 0.0 };
} else if snapshot.ledger_complete {
snapshot.progress_state = HealProgressState::Completed;
snapshot.progress_percentage = 100.0;
} else if has_object_sweep {
snapshot.refresh_progress_percentage();
} else if snapshot.stage_total == 0 {
snapshot.progress_state = HealProgressState::Indeterminate;
snapshot.progress_percentage = 0.0;
} else {
snapshot.progress_state = HealProgressState::Running;
snapshot.progress_percentage = ((snapshot.stage_current as f64 / snapshot.stage_total as f64) * 100.0).min(99.999);
}
snapshot.refresh_estimated_completion_time();
Some(snapshot)
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HealStatistics {
/// Total heal tasks
@@ -566,7 +230,6 @@ mod tests {
assert_eq!(progress.objects_scanned, 0);
assert_eq!(progress.objects_healed, 0);
assert_eq!(progress.objects_failed, 0);
assert_eq!(progress.skipped_objects, 0);
assert_eq!(progress.skipped_new_versions, 0);
assert_eq!(progress.skipped_ilm_expired, 0);
assert_eq!(progress.objects_total_count, 0);
@@ -587,8 +250,10 @@ mod tests {
assert_eq!(progress.objects_healed, 8);
assert_eq!(progress.objects_failed, 2);
assert_eq!(progress.bytes_processed, 1024);
assert_eq!(progress.progress_state, HealProgressState::Indeterminate);
assert_eq!(progress.progress_percentage, 0.0);
// Progress percentage should be calculated based on healed/total
// total = scanned + healed + failed = 10 + 8 + 2 = 20
// healed/total = 8/20 = 0.4 = 40%
assert!((progress.progress_percentage - 40.0).abs() < 0.001);
assert!(progress.last_update_time.is_some());
}
@@ -597,8 +262,7 @@ mod tests {
let mut progress = HealProgress::new();
progress.start_time = Some(SystemTime::now() - Duration::from_secs(10));
progress.set_total_baseline(100, 16384);
progress.update_progress(25, 25, 0, 4096);
progress.update_progress(100, 25, 0, 4096);
let eta = progress
.estimated_completion_time
@@ -611,7 +275,7 @@ mod tests {
let mut progress = HealProgress::new();
progress.set_total_baseline(10, 8192);
progress.update_progress(25, 25, 0, 4096);
progress.update_progress(100, 25, 0, 4096);
assert!((progress.progress_percentage - 50.0).abs() < 0.001);
}
@@ -621,7 +285,7 @@ mod tests {
let mut progress = HealProgress::new();
progress.set_total_baseline(10, 0);
progress.update_progress(5, 3, 2, 0);
progress.update_progress(100, 3, 2, 0);
assert!((progress.progress_percentage - 50.0).abs() < 0.001);
}
@@ -631,7 +295,7 @@ mod tests {
let mut progress = HealProgress::new();
progress.set_total_baseline(10, 0);
progress.update_progress(5, 3, 2, 0);
progress.update_progress(100, 3, 2, 0);
progress.record_skipped_new_version();
assert_eq!(progress.skipped_new_versions, 1);
@@ -672,8 +336,7 @@ mod tests {
fn test_heal_progress_update_progress_all_healed() {
let mut progress = HealProgress::new();
// When scanned=0, healed=10, failed=0: total=10, progress = 10/10 = 100%
progress.update_progress(10, 10, 0, 2048);
progress.mark_completed();
progress.update_progress(0, 10, 0, 2048);
// All healed, should be 100%
assert!((progress.progress_percentage - 100.0).abs() < 0.001);
@@ -731,7 +394,6 @@ mod tests {
assert_eq!(json["objectsScanned"], 10);
assert_eq!(json["objectsHealed"], 8);
assert_eq!(json["objectsFailed"], 2);
assert_eq!(json["skippedObjects"], 0);
assert_eq!(json["skippedNewVersions"], 0);
assert_eq!(json["skippedIlmExpired"], 0);
assert_eq!(json["bytesProcessed"], 1024);
@@ -743,7 +405,6 @@ mod tests {
fn test_heal_progress_is_completed_by_percentage() {
let mut progress = HealProgress::new();
progress.update_progress(10, 10, 0, 1024);
progress.mark_completed();
assert!(progress.is_completed());
}
@@ -754,7 +415,7 @@ mod tests {
progress.objects_scanned = 10;
progress.objects_healed = 8;
progress.objects_failed = 2;
progress.mark_completed();
// healed + failed = 8 + 2 = 10 >= scanned = 10
assert!(progress.is_completed());
}
@@ -794,108 +455,6 @@ mod tests {
assert!((progress.get_success_rate() - 100.0).abs() < 0.001);
}
#[test]
fn single_object_progress_reaches_terminal_100() {
let mut progress = HealProgress::new();
progress.update_object_progress(1, 1, 0, 0, 128);
assert!(!progress.is_completed());
progress.mark_completed();
assert!(progress.is_completed());
assert_eq!(progress.progress_percentage, 100.0);
}
#[test]
fn progress_without_baseline_is_indeterminate() {
let mut progress = HealProgress::new();
progress.update_object_progress(1, 1, 0, 0, 128);
assert_eq!(progress.progress_state, HealProgressState::Indeterminate);
assert_eq!(progress.progress_percentage, 0.0);
assert!(progress.estimated_completion_time.is_none());
}
#[test]
fn progress_retry_is_exactly_once() {
let mut progress = HealProgress::new();
progress.set_total_baseline(1, 128);
progress.update_object_progress(1, 1, 0, 0, 128);
progress.update_object_progress(1, 1, 0, 0, 128);
assert_eq!(progress.objects_scanned, 1);
assert_eq!(progress.objects_healed, 1);
assert_eq!(progress.bytes_processed, 128);
}
#[test]
fn progress_never_triggers_cleanup_before_terminal_ledger_empty() {
let mut progress = HealProgress::new();
progress.progress_percentage = 100.0;
assert!(!progress.is_completed());
progress.mark_completed();
assert!(progress.is_completed());
}
#[test]
fn progress_counter_overflow_is_marked_unknown_without_aborting_completed_heal() {
let mut progress = HealProgress::new();
progress.update_object_progress(u64::MAX, u64::MAX, 1, 0, 0);
assert_eq!(progress.progress_state, HealProgressState::Unknown);
progress.mark_completed();
assert!(progress.is_completed());
assert_eq!(progress.progress_state, HealProgressState::Unknown);
let aggregate = aggregate_heal_progress([progress]).expect("progress should aggregate");
assert!(aggregate.ledger_complete);
assert!(aggregate.counter_unknown);
assert_eq!(aggregate.progress_state, HealProgressState::Unknown);
assert_eq!(aggregate.progress_percentage, 100.0);
}
#[test]
fn aggregate_rejects_mixed_baseline_generations() {
let progress = |generation| HealProgress {
kind: HealProgressKind::ObjectSweep,
objects_scanned: 5,
objects_total_count: 10,
progress_state: HealProgressState::Running,
baseline_generation: Some(generation),
baseline_known: true,
..Default::default()
};
let aggregate = aggregate_heal_progress([progress(1), progress(2)]).expect("progress should aggregate");
assert!(!aggregate.baseline_known);
assert_eq!(aggregate.baseline_generation, None);
assert_eq!(aggregate.progress_state, HealProgressState::Indeterminate);
assert_eq!(aggregate.progress_percentage, 0.0);
}
#[test]
fn aggregate_accepts_multiple_sets_from_one_snapshot_generation() {
let progress = |objects_scanned| HealProgress {
kind: HealProgressKind::ObjectSweep,
objects_scanned,
objects_total_count: 10,
progress_state: HealProgressState::Running,
baseline_generation: Some(7),
baseline_known: true,
..Default::default()
};
let aggregate = aggregate_heal_progress([progress(5), progress(3)]).expect("progress should aggregate");
assert!(aggregate.baseline_known);
assert_eq!(aggregate.baseline_generation, Some(7));
}
#[test]
fn stage_updates_do_not_double_count_object_outcomes() {
let mut progress = HealProgress::new();
progress.update_object_progress(2, 1, 0, 1, 256);
progress.update_stage(3, 4);
assert_eq!(progress.kind, HealProgressKind::ObjectSweep);
assert_eq!(progress.objects_scanned, 2);
assert_eq!(progress.objects_healed, 1);
assert_eq!(progress.skipped_objects, 1);
}
#[test]
fn test_heal_statistics_new() {
let stats = HealStatistics::new();
+4 -130
View File
@@ -31,7 +31,7 @@ mod checkpoint;
mod replacement;
mod utils;
pub use checkpoint::{CheckpointManager, CheckpointObjectOutcome, CheckpointObjectOutcomeRecord, ResumeCheckpoint};
pub use checkpoint::{CheckpointManager, ResumeCheckpoint};
pub(crate) use replacement::replacement_target_identities_match;
use replacement::replacement_targets_match_identities;
pub use replacement::{
@@ -340,12 +340,6 @@ pub struct ResumeState {
pub failed_objects: u64,
/// skipped objects
pub skipped_objects: u64,
/// Terminal versions skipped because they were newer than the heal start.
#[serde(default)]
pub skipped_new_versions: u64,
/// Terminal versions handed to lifecycle expiry.
#[serde(default)]
pub skipped_ilm_expired: u64,
/// current bucket
pub current_bucket: Option<String>,
/// current object
@@ -360,24 +354,6 @@ pub struct ResumeState {
pub retry_count: u32,
/// max retries
pub max_retries: u32,
/// Bytes accounted by the object ledger; additive for old snapshots.
#[serde(default)]
pub processed_bytes: u64,
/// Total bytes from a complete usage snapshot, when available.
#[serde(default)]
pub total_bytes: u64,
/// Generation of the usage snapshot used for the baseline.
#[serde(default)]
pub baseline_generation: Option<u64>,
/// Whether the usage baseline is known. Missing in old snapshots means
/// indeterminate rather than a measured zero baseline.
#[serde(default)]
pub baseline_known: bool,
/// Persistent telemetry fence for counter/byte overflow or corruption.
/// It must survive a restart so a saturated snapshot is never presented as
/// a measured percentage on the next resume.
#[serde(default)]
pub counter_unknown: bool,
}
impl ResumeState {
@@ -401,8 +377,6 @@ impl ResumeState {
successful_objects: 0,
failed_objects: 0,
skipped_objects: 0,
skipped_new_versions: 0,
skipped_ilm_expired: 0,
current_bucket: None,
current_object: None,
completed_buckets: Vec::new(),
@@ -410,11 +384,6 @@ impl ResumeState {
error_message: None,
retry_count: 0,
max_retries: 3,
processed_bytes: 0,
total_bytes: 0,
baseline_generation: None,
baseline_known: false,
counter_unknown: false,
}
}
@@ -443,39 +412,6 @@ impl ResumeState {
self.last_update = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
}
pub fn update_progress_with_bytes(
&mut self,
processed: u64,
successful: u64,
failed: u64,
skipped: u64,
processed_bytes: u64,
) {
self.update_progress(processed, successful, failed, skipped);
self.processed_bytes = processed_bytes;
}
pub fn set_skipped_version_counts(&mut self, new_versions: u64, ilm_expired: u64) {
self.skipped_new_versions = new_versions;
self.skipped_ilm_expired = ilm_expired;
self.last_update = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
}
pub fn set_progress_baseline(&mut self, total_objects: u64, total_bytes: u64, generation: Option<u64>) {
self.total_objects = total_objects;
self.total_bytes = total_bytes;
self.baseline_generation = generation;
// This method is called only after a complete usage snapshot has been
// validated. A complete but empty snapshot is still a known baseline.
self.baseline_known = true;
self.last_update = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
}
pub fn mark_counter_unknown(&mut self) {
self.counter_unknown = true;
self.last_update = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
}
pub fn set_current_item(&mut self, bucket: Option<String>, object: Option<String>) {
self.current_bucket = bucket;
self.current_object = object;
@@ -501,7 +437,6 @@ impl ResumeState {
if let Some(pos) = self.pending_buckets.iter().position(|b| b == bucket) {
self.pending_buckets.remove(pos);
}
self.resume_cursor = None;
self.last_update = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
}
@@ -519,10 +454,6 @@ impl ResumeState {
self.successful_objects = 0;
self.failed_objects = 0;
self.skipped_objects = 0;
self.skipped_new_versions = 0;
self.skipped_ilm_expired = 0;
self.processed_bytes = 0;
self.counter_unknown = false;
self.completed = false;
// A retry re-scans every bucket from the beginning, so the version
// cursor must be cleared too — otherwise the retry would resume mid-scan.
@@ -545,28 +476,14 @@ impl ResumeState {
}
pub fn get_progress_percentage(&self) -> f64 {
if self.completed {
return 100.0;
}
if self.counter_unknown {
return 0.0;
}
if !self.baseline_known {
return 0.0;
}
if self.total_bytes > 0 {
return ((self.processed_bytes as f64 / self.total_bytes as f64) * 100.0).min(99.999);
}
if self.total_objects == 0 {
return 0.0;
}
((self.processed_objects as f64 / self.total_objects as f64) * 100.0).min(99.999)
(self.processed_objects as f64 / self.total_objects as f64) * 100.0
}
pub fn get_success_rate(&self) -> f64 {
let Some(total) = self.successful_objects.checked_add(self.failed_objects) else {
return 0.0;
};
let total = self.successful_objects + self.failed_objects;
if total == 0 {
return 0.0;
}
@@ -837,14 +754,6 @@ impl ResumeManager {
state.successful_objects = 0;
state.failed_objects = 0;
state.skipped_objects = 0;
state.skipped_new_versions = 0;
state.skipped_ilm_expired = 0;
state.processed_bytes = 0;
state.total_objects = 0;
state.total_bytes = 0;
state.baseline_generation = None;
state.baseline_known = false;
state.counter_unknown = false;
state.completed = false;
state.completed_buckets.clear();
state.schema_version = CURRENT_RESUME_SCHEMA;
@@ -929,41 +838,6 @@ impl ResumeManager {
self.save_state_throttled().await
}
pub async fn update_progress_with_bytes(
&self,
processed: u64,
successful: u64,
failed: u64,
skipped: u64,
processed_bytes: u64,
) -> Result<()> {
let mut state = self.state.write().await;
state.update_progress_with_bytes(processed, successful, failed, skipped, processed_bytes);
drop(state);
self.save_state_throttled().await
}
pub async fn set_progress_baseline(&self, total_objects: u64, total_bytes: u64, generation: Option<u64>) -> Result<()> {
let mut state = self.state.write().await;
state.set_progress_baseline(total_objects, total_bytes, generation);
drop(state);
self.save_state_throttled().await
}
pub async fn mark_counter_unknown(&self) -> Result<()> {
let mut state = self.state.write().await;
state.mark_counter_unknown();
drop(state);
self.save_state().await
}
pub async fn set_skipped_version_counts(&self, new_versions: u64, ilm_expired: u64) -> Result<()> {
let mut state = self.state.write().await;
state.set_skipped_version_counts(new_versions, ilm_expired);
drop(state);
self.save_state_throttled().await
}
/// Set current item. Called once per healed object, so persistence is
/// throttled: the in-memory state always updates, but the snapshot is only
/// written every `PERSIST_EVERY_MUTATIONS` calls or `PERSIST_INTERVAL`.
@@ -1008,7 +882,7 @@ impl ResumeManager {
let mut state = self.state.write().await;
state.complete_bucket(bucket);
drop(state);
self.save_state().await
self.save_state_throttled().await
}
/// mark task completed
+9 -198
View File
@@ -29,30 +29,10 @@ use super::{
const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state";
/// Current on-disk schema version for `ResumeCheckpoint`. Schema 5 could
/// persist dedup identities without the aggregate counters needed to restore
/// them safely, so stale checkpoints are discarded and replayed.
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 6;
#[derive(Debug, Clone, Copy)]
pub enum CheckpointObjectOutcome {
Processed,
Failed,
Skipped,
}
#[derive(Debug)]
pub struct CheckpointObjectOutcomeRecord {
pub object: String,
pub outcome: CheckpointObjectOutcome,
pub successful: u64,
pub failed: u64,
pub skipped: u64,
pub bytes: u64,
pub skipped_new_versions: u64,
pub skipped_ilm_expired: u64,
pub counter_unknown: bool,
}
/// Current on-disk schema version for `ResumeCheckpoint`. Same rationale as
/// `CURRENT_RESUME_SCHEMA`: pre-per-version dedup identities are not comparable
/// to the new `compose_key` identities, so a stale checkpoint is discarded.
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5;
/// resume checkpoint
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -77,30 +57,6 @@ pub struct ResumeCheckpoint {
pub failed_objects: HashSet<String>,
/// skipped objects
pub skipped_objects: HashSet<String>,
/// Aggregate object ledger counters restored alongside the dedup sets.
#[serde(default)]
pub successful_objects: u64,
#[serde(default)]
pub failed_object_count: u64,
#[serde(default)]
pub skipped_object_count: u64,
#[serde(default)]
pub skipped_new_versions: u64,
#[serde(default)]
pub skipped_ilm_expired: u64,
#[serde(default)]
pub processed_bytes: u64,
#[serde(default)]
pub total_objects: u64,
#[serde(default)]
pub total_bytes: u64,
#[serde(default)]
pub baseline_generation: Option<u64>,
#[serde(default)]
pub baseline_known: bool,
/// Persistent telemetry fence for counter/byte overflow or corruption.
#[serde(default)]
pub counter_unknown: bool,
}
impl ResumeCheckpoint {
@@ -114,17 +70,6 @@ impl ResumeCheckpoint {
processed_objects: HashSet::new(),
failed_objects: HashSet::new(),
skipped_objects: HashSet::new(),
successful_objects: 0,
failed_object_count: 0,
skipped_object_count: 0,
skipped_new_versions: 0,
skipped_ilm_expired: 0,
processed_bytes: 0,
total_objects: 0,
total_bytes: 0,
baseline_generation: None,
baseline_known: false,
counter_unknown: false,
}
}
@@ -146,34 +91,6 @@ impl ResumeCheckpoint {
self.skipped_objects.insert(object);
}
pub fn update_progress(&mut self, successful: u64, failed: u64, skipped: u64, bytes: u64) {
self.successful_objects = successful;
self.failed_object_count = failed;
self.skipped_object_count = skipped;
self.processed_bytes = bytes;
}
pub fn set_progress_baseline(&mut self, total_objects: u64, total_bytes: u64, generation: Option<u64>) {
self.total_objects = total_objects;
self.total_bytes = total_bytes;
self.baseline_generation = generation;
// The caller has already validated that this is a complete snapshot;
// preserve the distinction between a known empty scope and an old
// checkpoint that omitted all baseline fields.
self.baseline_known = true;
}
pub fn mark_counter_unknown(&mut self) {
self.counter_unknown = true;
self.checkpoint_time = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
}
pub fn set_skipped_version_counts(&mut self, new_versions: u64, ilm_expired: u64) {
self.skipped_new_versions = new_versions;
self.skipped_ilm_expired = ilm_expired;
self.checkpoint_time = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
}
/// Advance past a fully-processed page: objects below `object_index` are
/// skipped by position on resume, so the per-object sets no longer need
/// their entries and would otherwise grow with the whole bucket.
@@ -190,17 +107,6 @@ impl ResumeCheckpoint {
self.update_position(0, 0);
self.processed_objects.clear();
self.skipped_objects.clear();
self.successful_objects = 0;
self.failed_object_count = 0;
self.skipped_object_count = 0;
self.skipped_new_versions = 0;
self.skipped_ilm_expired = 0;
self.processed_bytes = 0;
self.total_objects = 0;
self.total_bytes = 0;
self.baseline_generation = None;
self.baseline_known = false;
self.counter_unknown = false;
self.failed_objects.clear();
}
}
@@ -252,9 +158,10 @@ impl CheckpointManager {
});
}
// Older checkpoints can contain identities that are not comparable to
// the current keys or lack their corresponding aggregate counters.
// Discard the stale sets and position so the scan restarts cleanly.
// A checkpoint from an older schema stored latest-only dedup identities
// that are not comparable to the new per-version `compose_key`
// identities. Discard the stale sets and position, then stamp the
// current schema so the scan restarts cleanly.
if checkpoint.schema_version > CURRENT_CHECKPOINT_SCHEMA {
return Err(Error::TaskExecutionFailed {
message: format!(
@@ -278,17 +185,6 @@ impl CheckpointManager {
checkpoint.processed_objects.clear();
checkpoint.failed_objects.clear();
checkpoint.skipped_objects.clear();
checkpoint.successful_objects = 0;
checkpoint.failed_object_count = 0;
checkpoint.skipped_object_count = 0;
checkpoint.skipped_new_versions = 0;
checkpoint.skipped_ilm_expired = 0;
checkpoint.processed_bytes = 0;
checkpoint.total_objects = 0;
checkpoint.total_bytes = 0;
checkpoint.baseline_generation = None;
checkpoint.baseline_known = false;
checkpoint.counter_unknown = false;
checkpoint.current_bucket_index = 0;
checkpoint.current_object_index = 0;
checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA;
@@ -329,7 +225,7 @@ impl CheckpointManager {
self.save_checkpoint_throttled().await
}
/// Persist a completed page position while retaining its identities.
/// Advance past a completed page and prune the per-object sets, then persist.
pub async fn complete_page(&self, bucket_index: usize, object_index: usize) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await;
checkpoint.complete_page(bucket_index, object_index);
@@ -337,35 +233,6 @@ impl CheckpointManager {
self.save_checkpoint_throttled().await
}
/// Persist the page position while retaining identities until the resume
/// cursor is durable.
pub async fn advance_page(&self, bucket_index: usize, object_index: usize) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await;
checkpoint.update_position(bucket_index, object_index);
drop(checkpoint);
self.save_checkpoint().await
}
/// Remove the previous page's dedup identities only after its resume cursor
/// has been durably exposed.
pub async fn prune_completed_page(&self) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await;
checkpoint.processed_objects.clear();
checkpoint.skipped_objects.clear();
checkpoint.failed_objects.clear();
drop(checkpoint);
self.save_checkpoint().await
}
/// Advance to the next bucket and clear the final page identities after the
/// resume state has durably recorded the completed bucket.
pub async fn complete_bucket(&self, next_bucket_index: usize) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await;
checkpoint.complete_page(next_bucket_index, 0);
drop(checkpoint);
self.save_checkpoint().await
}
/// Reset the checkpoint to the start of the scan for a retry, then persist.
pub async fn reset_for_retry(&self) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await;
@@ -400,62 +267,6 @@ impl CheckpointManager {
self.save_checkpoint_if_due().await
}
/// Atomically persist an object's dedup identity with its aggregate result.
pub async fn record_object_outcome(&self, record: CheckpointObjectOutcomeRecord) -> Result<()> {
let CheckpointObjectOutcomeRecord {
object,
outcome,
successful,
failed,
skipped,
bytes,
skipped_new_versions,
skipped_ilm_expired,
counter_unknown,
} = record;
let mut checkpoint = self.checkpoint.write().await;
match outcome {
CheckpointObjectOutcome::Processed => checkpoint.add_processed_object(object),
CheckpointObjectOutcome::Failed => checkpoint.add_failed_object(object),
CheckpointObjectOutcome::Skipped => checkpoint.add_skipped_object(object),
}
checkpoint.update_progress(successful, failed, skipped, bytes);
checkpoint.set_skipped_version_counts(skipped_new_versions, skipped_ilm_expired);
if counter_unknown {
checkpoint.mark_counter_unknown();
}
drop(checkpoint);
self.save_checkpoint_if_due().await
}
pub async fn update_progress(&self, successful: u64, failed: u64, skipped: u64, bytes: u64) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await;
checkpoint.update_progress(successful, failed, skipped, bytes);
drop(checkpoint);
self.save_checkpoint_if_due().await
}
pub async fn set_progress_baseline(&self, total_objects: u64, total_bytes: u64, generation: Option<u64>) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await;
checkpoint.set_progress_baseline(total_objects, total_bytes, generation);
drop(checkpoint);
self.save_checkpoint_throttled().await
}
pub async fn mark_counter_unknown(&self) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await;
checkpoint.mark_counter_unknown();
drop(checkpoint);
self.save_checkpoint().await
}
pub async fn set_skipped_version_counts(&self, new_versions: u64, ilm_expired: u64) -> Result<()> {
let mut checkpoint = self.checkpoint.write().await;
checkpoint.set_skipped_version_counts(new_versions, ilm_expired);
drop(checkpoint);
self.save_checkpoint_throttled().await
}
async fn save_checkpoint_if_due(&self) -> Result<()> {
let should_save = self.throttle.lock().map(|mut throttle| throttle.record()).unwrap_or(true);
if !should_save {
+4 -153
View File
@@ -1296,7 +1296,6 @@ async fn test_resume_state_progress() {
assert_eq!(progress, 0.0); // total_objects is 0
state.total_objects = 100;
state.baseline_known = true;
let progress = state.get_progress_percentage();
assert_eq!(progress, 10.0);
}
@@ -1476,40 +1475,6 @@ fn test_checkpoint_object_sets_dedupe_and_prune() {
assert!(checkpoint.failed_objects.is_empty());
}
#[tokio::test]
async fn checkpoint_page_commit_keeps_ledger_until_cursor_is_durable() {
let (_temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let checkpoint = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
checkpoint
.record_object_outcome(CheckpointObjectOutcomeRecord {
object: "bucket/object:v1".to_string(),
outcome: CheckpointObjectOutcome::Processed,
successful: 1,
failed: 0,
skipped: 0,
bytes: 128,
skipped_new_versions: 0,
skipped_ilm_expired: 0,
counter_unknown: false,
})
.await
.unwrap();
checkpoint.advance_page(0, 1).await.unwrap();
let reloaded = CheckpointManager::load_from_disk(disk.clone(), &task_id).await.unwrap();
let snapshot = reloaded.get_checkpoint().await;
assert_eq!(snapshot.current_object_index, 1);
assert_eq!(snapshot.successful_objects, 1);
assert_eq!(snapshot.processed_bytes, 128);
assert!(snapshot.processed_objects.contains("bucket/object:v1"));
checkpoint.prune_completed_page().await.unwrap();
let reloaded = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap();
assert!(reloaded.get_checkpoint().await.processed_objects.is_empty());
}
#[test]
fn test_checkpoint_loads_legacy_vec_format() {
// Checkpoints written before the HashSet migration stored the object
@@ -1603,14 +1568,14 @@ async fn test_resumestate_schema_v0_discarded_on_load() {
}
#[tokio::test]
async fn test_checkpoint_schema_v5_discarded_on_load() {
async fn test_checkpoint_schema_v4_discarded_on_load() {
let (temp_dir, disk) = schema_test_disk().await;
// Schema v5 can persist failed identities without the aggregate counters
// that make those identities safe to deduplicate after an upgrade.
// The previous checkpoint schema is unsafe once its paired resume
// state is discarded: retaining either position would skip work.
let task_id = "00000000-0000-4000-8000-000000000002";
let legacy = r#"{
"schema_version": 5,
"schema_version": 4,
"task_id": "00000000-0000-4000-8000-000000000002",
"checkpoint_time": 1700000000,
"current_bucket_index": 2,
@@ -1674,120 +1639,6 @@ async fn current_normal_resume_schema_preserves_progress() {
temp_dir.close().expect("remove schema test directory");
}
#[test]
fn progress_checkpoint_restores_bytes_and_generation() {
let mut checkpoint = ResumeCheckpoint::new("progress-checkpoint".to_string());
checkpoint.set_progress_baseline(9, 4096, Some(77));
checkpoint.update_progress(4, 1, 2, 2048);
checkpoint.set_skipped_version_counts(3, 1);
checkpoint.mark_counter_unknown();
let restored: ResumeCheckpoint =
serde_json::from_slice(&serde_json::to_vec(&checkpoint).expect("serialize checkpoint")).expect("deserialize checkpoint");
assert_eq!(restored.processed_bytes, 2048);
assert_eq!(restored.total_objects, 9);
assert_eq!(restored.total_bytes, 4096);
assert_eq!(restored.baseline_generation, Some(77));
assert!(restored.baseline_known);
assert_eq!(restored.skipped_new_versions, 3);
assert_eq!(restored.skipped_ilm_expired, 1);
assert!(restored.counter_unknown);
}
#[test]
fn old_progress_schema_migrates_missing_fields_to_unknown() {
let state = ResumeState::new(
"legacy-progress".to_string(),
"erasure_set".to_string(),
"pool_0_set_0".to_string(),
Vec::new(),
);
let mut value = serde_json::to_value(state).expect("serialize legacy-compatible state");
let object = value.as_object_mut().expect("state must be an object");
for field in [
"processed_bytes",
"total_bytes",
"baseline_generation",
"baseline_known",
"skipped_new_versions",
"skipped_ilm_expired",
] {
object.remove(field);
}
object.insert("total_objects".to_string(), serde_json::json!(10));
object.insert("processed_objects".to_string(), serde_json::json!(5));
let restored: ResumeState = serde_json::from_value(value).expect("deserialize old progress state");
assert_eq!(restored.processed_bytes, 0);
assert_eq!(restored.total_bytes, 0);
assert_eq!(restored.baseline_generation, None);
assert!(!restored.baseline_known, "missing baseline must remain unknown");
assert_eq!(restored.get_progress_percentage(), 0.0);
assert_eq!(restored.skipped_new_versions, 0);
assert_eq!(restored.skipped_ilm_expired, 0);
}
#[test]
fn progress_counter_unknown_survives_resume_round_trip() {
let mut state = ResumeState::new(
"overflow-progress".to_string(),
"erasure_set".to_string(),
"pool_0_set_0".to_string(),
Vec::new(),
);
state.mark_counter_unknown();
let restored: ResumeState =
serde_json::from_slice(&serde_json::to_vec(&state).expect("serialize resume state")).expect("deserialize resume state");
assert!(restored.counter_unknown);
}
#[tokio::test]
async fn checkpoint_progress_survives_a_torn_resume_summary_write() {
let (_temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let _resume = ResumeManager::new(
disk.clone(),
task_id.clone(),
"erasure_set".to_string(),
"pool_0_set_0".to_string(),
vec!["bucket".to_string()],
)
.await
.expect("resume state should persist");
let checkpoint = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("checkpoint should persist");
// This is the ordering used by the erasure-set loop: the checkpoint is
// durable before the summary write. Stop here to model a crash in the
// inter-store window and verify that the recovery authority retains the
// telemetry fence and bytes.
checkpoint
.update_progress(3, 0, 0, 1024)
.await
.expect("checkpoint progress should persist");
checkpoint.mark_counter_unknown().await.expect("unknown fence should persist");
checkpoint
.update_position(0, 3)
.await
.expect("checkpoint position should persist");
let restored_checkpoint = CheckpointManager::load_from_disk(disk.clone(), &task_id)
.await
.expect("checkpoint should reload")
.get_checkpoint()
.await;
let restored_resume = ResumeManager::load_from_disk(disk, &task_id)
.await
.expect("resume summary should reload")
.get_state()
.await;
assert!(restored_checkpoint.counter_unknown);
assert_eq!(restored_checkpoint.processed_bytes, 1024);
assert_eq!(restored_checkpoint.current_object_index, 3);
assert!(!restored_resume.counter_unknown, "summary is intentionally the torn/older store");
}
#[tokio::test]
async fn future_resume_and_checkpoint_schemas_are_rejected() {
let (temp_dir, disk) = schema_test_disk().await;
+2 -47
View File
@@ -22,7 +22,6 @@ use serde::{Deserialize, Serialize};
use std::sync::Arc;
use tracing::{debug, error, warn};
use super::progress::stable_generation;
use super::storage_api::owner::{EcstoreHealLifecycleExpiryContext, ecstore_load_admin_data_usage_from_backend_cached};
use super::storage_api::storage::{
BucketInfo, BucketOperations, DiskSetSelector, HealOperations as _, ListOperations as _, ObjectIO as _,
@@ -35,9 +34,6 @@ pub use super::{HealObjectInfo, HealObjectOptions, HealPutObjReader};
pub struct HealBucketUsageBaseline {
pub objects_count: u64,
pub bytes: u64,
/// Stable identity of the validated usage snapshot and selected scope.
/// `None` is retained for test/legacy providers that cannot expose one.
pub generation: Option<u64>,
}
pub struct HealLifecycleExpiryContext {
@@ -789,52 +785,11 @@ impl HealStorageAPI for ECStoreHealStorage {
let mut baseline = HealBucketUsageBaseline::default();
for bucket in buckets {
if let Some(usage) = info.buckets_usage.get(bucket) {
baseline.objects_count = match baseline.objects_count.checked_add(usage.objects_count) {
Some(total) => total,
// A corrupt/overflowing usage snapshot is not a usable
// denominator. Leave progress indeterminate instead of
// turning saturation into a plausible percentage.
None => return Ok(None),
};
baseline.bytes = match baseline.bytes.checked_add(usage.size) {
Some(total) => total,
None => return Ok(None),
};
baseline.objects_count = baseline.objects_count.saturating_add(usage.objects_count);
baseline.bytes = baseline.bytes.saturating_add(usage.size);
}
}
let identity = info.snapshot_identity();
let mut canonical = Vec::new();
match identity.last_update {
Some(last_update) => {
canonical.push(1);
canonical.extend_from_slice(
&last_update
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos()
.to_be_bytes(),
);
}
None => canonical.push(0),
}
for value in [identity.scanner_cycle, identity.scanner_epoch] {
match value {
Some(value) => {
canonical.push(1);
canonical.extend_from_slice(&value.to_be_bytes());
}
None => canonical.push(0),
}
}
let mut scope = buckets.to_vec();
scope.sort_unstable();
for bucket in scope {
canonical.extend_from_slice(&(bucket.len() as u64).to_be_bytes());
canonical.extend_from_slice(bucket.as_bytes());
}
baseline.generation = Some(stable_generation(&[&canonical]));
Ok(Some(baseline))
}
+3 -7
View File
@@ -649,7 +649,7 @@ impl HealTask {
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("skipped: {bucket}/{object}")));
progress.update_stage(1, 1);
progress.update_progress(0, 1, 0, 0);
Ok(())
}
@@ -733,7 +733,7 @@ impl HealTask {
"Heal object skipped for data usage cache after transient error"
);
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
true
}
@@ -757,7 +757,7 @@ impl HealTask {
);
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("skipped: {bucket}/{object}")));
progress.update_stage(4, 4);
progress.update_progress(4, 4, 0, 0);
true
}
@@ -831,10 +831,6 @@ impl HealTask {
match &result {
Ok(_) => {
// A stage can reach its final step before the durable resume
// ledger and cleanup fences commit. Publish terminal 100 only
// after the enclosing operation has returned success.
self.progress.write().await.mark_completed();
let mut status = self.status.write().await;
*status = HealTaskStatus::Completed;
demote_to_debug_when!(self.heal_type.is_per_object(), info, target: "rustfs::heal::task", {
+16 -41
View File
@@ -13,7 +13,6 @@
// limitations under the License.
/// bucket/cluster/prefix heal: the recursive bucket-objects sweep and the erasure-set usage baseline
use super::*;
use crate::heal::progress::{add_bytes, increment_counter, stable_generation};
impl HealTask {
pub(super) async fn heal_bucket(&self, bucket: &str) -> Result<()> {
@@ -33,7 +32,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("bucket: {bucket}")));
progress.update_stage(0, 3);
progress.update_progress(0, 3, 0, 0);
}
// Step 1: Check if bucket exists
@@ -67,7 +66,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_stage(1, 3);
progress.update_progress(1, 3, 0, 0);
}
// Step 2: Perform bucket heal using ecstore
@@ -123,7 +122,7 @@ impl HealTask {
if !self.options.recursive {
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
Ok(())
}
@@ -143,7 +142,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
Err(Error::TaskExecutionFailed {
message: format!("Failed to heal bucket {bucket}: {e}"),
@@ -246,7 +245,6 @@ impl HealTask {
let mut scanned = 0u64;
let mut healed = 0u64;
let mut failed = 0u64;
let mut skipped = 0u64;
let mut retryable_failed = 0u64;
let mut permanent_failed = 0u64;
let mut bytes = 0u64;
@@ -288,14 +286,16 @@ impl HealTask {
let mut retry = Vec::with_capacity(pending.len());
for item in pending {
self.check_control_flags().await?;
let mut telemetry_unknown = false;
let object = item.name.as_str();
if retry_attempt == 0 {
scanned = scanned.saturating_add(1);
}
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("{bucket}/{object}")));
progress.update_progress(scanned, healed, failed, bytes);
}
let mut terminal_outcome = true;
let error = match self
.await_with_control(
self.storage
@@ -304,13 +304,13 @@ impl HealTask {
.await
{
Ok((result, None)) => {
telemetry_unknown |= !increment_counter(&mut healed);
telemetry_unknown |= !add_bytes(&mut bytes, u64::try_from(result.object_size).unwrap_or(u64::MAX));
healed = healed.saturating_add(1);
bytes = bytes.saturating_add(u64::try_from(result.object_size).unwrap_or_default());
self.record_result_item(result).await;
None
}
Ok((_, Some(err))) if is_missing_object_dir_heal_result(object, &err) => {
telemetry_unknown |= !increment_counter(&mut healed);
healed = healed.saturating_add(1);
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
@@ -329,7 +329,6 @@ impl HealTask {
if let Some(err) = error {
if Self::should_skip_data_usage_cache_heal_error(bucket, object, &err) {
telemetry_unknown |= !increment_counter(&mut skipped);
warn!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
@@ -343,7 +342,6 @@ impl HealTask {
"Heal bucket object repair skipped due to transient metadata error"
);
} else if err.is_recoverable_heal() && retry_attempt < MAX_BUCKET_OBJECT_HEAL_RETRIES {
terminal_outcome = false;
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
@@ -359,7 +357,7 @@ impl HealTask {
);
retry.push(item);
} else {
telemetry_unknown |= !increment_counter(&mut failed);
failed = failed.saturating_add(1);
if err.is_recoverable_heal() {
retryable_failed = retryable_failed.saturating_add(1);
} else {
@@ -385,19 +383,8 @@ impl HealTask {
}
}
if terminal_outcome {
telemetry_unknown |= !increment_counter(&mut scanned);
}
if !terminal_outcome {
continue;
}
let mut progress = self.progress.write().await;
progress.update_object_progress(scanned, healed, failed, skipped, bytes);
if telemetry_unknown {
progress.mark_unknown();
}
progress.update_progress(scanned, healed, failed, bytes);
}
pending = retry;
retry_attempt = retry_attempt.saturating_add(1);
@@ -444,10 +431,7 @@ impl HealTask {
Ok(())
}
pub(super) async fn apply_erasure_set_usage_baseline(&self, buckets: &[String], _set_disk_id: &str) -> Result<()> {
if matches!(self.options.scan_mode, HealScanMode::Deep) || matches!(self.source, HealRequestSource::AutoHeal) {
return Ok(());
}
pub(super) async fn apply_erasure_set_usage_baseline(&self, buckets: &[String]) -> Result<()> {
let baseline = match self
.await_with_control(self.storage.erasure_set_usage_baseline(buckets))
.await
@@ -458,18 +442,9 @@ impl HealTask {
Err(_) => return Ok(()),
};
let HealBucketUsageBaseline {
objects_count,
bytes,
generation,
} = baseline;
let generation = generation.map(|snapshot_generation| stable_generation(&[&snapshot_generation.to_be_bytes()]));
let HealBucketUsageBaseline { objects_count, bytes } = baseline;
let mut progress = self.progress.write().await;
if let Some(generation) = generation {
progress.set_total_baseline_with_generation(objects_count, bytes, generation);
} else {
progress.set_total_baseline(objects_count, bytes);
}
progress.set_total_baseline(objects_count, bytes);
Ok(())
}
}
+10 -8
View File
@@ -32,7 +32,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("erasure_set: {} ({} buckets)", set_disk_id, buckets.len())));
progress.update_stage(0, 4);
progress.update_progress(0, 4, 0, 0);
}
let is_auto_replacement = matches!(self.source, HealRequestSource::AutoHeal) && !self.heal_endpoints.is_empty();
@@ -158,7 +158,7 @@ impl HealTask {
None
};
self.apply_erasure_set_usage_baseline(&buckets, &set_disk_id).await?;
self.apply_erasure_set_usage_baseline(&buckets).await?;
let healing_marker = format!("{set_disk_id}:{}", self.id);
if let Some((disk, resume_manager, _)) = replacement_resume.as_ref() {
@@ -248,7 +248,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_stage(4, 4);
progress.update_progress(4, 4, 0, 0);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal disk format for {set_disk_id}: {error}"),
@@ -304,7 +304,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_stage(4, 4);
progress.update_progress(4, 4, 0, 0);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal disk format for {set_disk_id}: {e}"),
@@ -314,7 +314,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_stage(1, 4);
progress.update_progress(1, 4, 0, 0);
}
// The rebuilt disks are formatted now: mark them as healing so
@@ -343,7 +343,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_stage(2, 4);
progress.update_progress(2, 4, 0, 0);
}
// Step 3: Heal bucket structure
@@ -427,7 +427,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 4);
progress.update_progress(3, 4, 0, 0);
}
// Step 4: Execute erasure set heal with resume
@@ -470,7 +470,9 @@ impl HealTask {
};
{
self.progress.write().await.update_stage(4, 4);
let mut progress = self.progress.write().await;
let bytes_processed = progress.bytes_processed;
progress.update_progress(4, 4, 0, bytes_processed);
}
match result {
+10 -10
View File
@@ -32,7 +32,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("metadata: {bucket}/{object}")));
progress.update_stage(0, 3);
progress.update_progress(0, 3, 0, 0);
}
// Step 1: Check if object exists
@@ -74,7 +74,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_stage(1, 3);
progress.update_progress(1, 3, 0, 0);
}
// Step 2: Perform metadata heal using ecstore
@@ -122,7 +122,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal metadata {bucket}/{object}: {e}"),
@@ -145,7 +145,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
self.record_result_item(result).await;
Ok(())
@@ -167,7 +167,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
Err(Error::TaskExecutionFailed {
message: format!("Failed to heal metadata {bucket}/{object}: {e}"),
@@ -194,7 +194,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("ec_decode: {bucket}/{object}")));
progress.update_stage(0, 3);
progress.update_progress(0, 3, 0, 0);
}
// Step 1: Check if object exists
@@ -236,7 +236,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_stage(1, 3);
progress.update_progress(1, 3, 0, 0);
}
// Step 2: Perform EC decode heal using ecstore
@@ -284,7 +284,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal EC decode {bucket}/{object}: {e}"),
@@ -309,7 +309,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_object_progress(1, 1, 0, 0, object_size);
progress.update_progress(3, 3, 0, object_size);
}
self.record_result_item(result).await;
Ok(())
@@ -331,7 +331,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
Err(Error::TaskExecutionFailed {
message: format!("Failed to heal EC decode {bucket}/{object}: {e}"),
+8 -8
View File
@@ -36,7 +36,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("{bucket}/{object}")));
progress.update_stage(0, 4);
progress.update_progress(0, 4, 0, 0);
}
// Step 1: Check if object exists and get metadata
@@ -132,7 +132,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_stage(1, 3);
progress.update_progress(1, 3, 0, 0);
}
// Step 2: directly call ecstore to perform heal
@@ -187,7 +187,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
return Ok(());
}
@@ -207,7 +207,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
if Self::should_return_typed_heal_error(&e) {
@@ -249,7 +249,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_object_progress(1, 1, 0, 0, object_size);
progress.update_progress(3, 3, 0, object_size);
}
self.record_result_item(result).await;
Ok(())
@@ -275,7 +275,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
return Ok(());
}
@@ -295,7 +295,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_stage(3, 3);
progress.update_progress(3, 3, 0, 0);
}
if Self::should_return_typed_heal_error(&e) {
@@ -414,7 +414,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_object_progress(1, 1, 0, 0, object_size);
progress.update_progress(4, 4, 0, object_size);
}
self.record_result_item(result).await;
Ok(())
-47
View File
@@ -22,7 +22,6 @@ use std::sync::Mutex;
use tempfile::TempDir;
use super::super::storage_api::status::BucketInfo;
use crate::heal::progress::HealProgressState;
#[tokio::test]
async fn retry_request_carries_remaining_timeout_budget() {
@@ -2125,7 +2124,6 @@ async fn erasure_set_heal_applies_usage_baseline_to_progress() {
usage_baseline: Mutex::new(Some(HealBucketUsageBaseline {
objects_count: 10,
bytes: 8,
generation: Some(1),
})),
..Default::default()
});
@@ -2149,55 +2147,10 @@ async fn erasure_set_heal_applies_usage_baseline_to_progress() {
let progress = task.get_progress().await;
assert_eq!(progress.objects_total_count, 10);
assert_eq!(progress.objects_total_size, 8);
assert!(progress.baseline_generation.is_some());
assert!(progress.baseline_known);
assert_eq!(progress.bytes_processed, 2);
assert!((progress.progress_percentage - 25.0).abs() < 0.001);
}
#[tokio::test]
async fn erasure_set_disk_walk_keeps_cluster_usage_baseline_indeterminate() {
for (scan_mode, source) in [
(HealScanMode::Deep, HealRequestSource::Admin),
(HealScanMode::Normal, HealRequestSource::AutoHeal),
] {
let temp = TempDir::new().expect("temporary directory should be created");
let disk = make_resume_disk(&temp).await;
let storage = Arc::new(MockStorage {
resume_disk: Mutex::new(Some(disk)),
usage_baseline: Mutex::new(Some(HealBucketUsageBaseline {
objects_count: 10,
bytes: 8,
generation: Some(1),
})),
..Default::default()
});
let mut request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
scan_mode,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = source;
let task = HealTask::from_request(request, storage);
task.heal_erasure_set(vec!["bucket-a".to_string()], "pool_0_set_0".to_string())
.await
.expect("erasure set heal should complete");
let progress = task.get_progress().await;
assert!(!progress.baseline_known);
assert_eq!(progress.baseline_generation, None);
assert_eq!(progress.progress_state, HealProgressState::Indeterminate);
}
}
#[tokio::test]
async fn erasure_set_heal_ignores_usage_baseline_errors() {
let temp = TempDir::new().expect("temporary directory should be created");
+1 -1
View File
@@ -19,7 +19,7 @@ pub use error::{Error, Result};
pub use heal::{
HealManager, HealOperationsSnapshot, HealOptions, HealPriority, HealPriorityCounts, HealRequest, HealSourceCounts, HealType,
channel::HealChannelProcessor,
progress::{HealProgress, aggregate_heal_progress},
progress::HealProgress,
resume::{ReplacementRecoveryRecord, ReplacementRecoveryState, ResumeUtils},
};
use rustfs_concurrency::WorkloadAdmissionSnapshotProvider;
+86 -25
View File
@@ -16,14 +16,14 @@
//!
//! `scripts/test/vault_ha_kms_live.sh` owns the official Vault containers and
//! kills the active node while this test continuously decrypts through a
//! surviving standby. KV2 and Transit requests must remain successful, use a
//! bounded number of attempts, and leave the circuit and in-flight gauges at
//! zero after a new leader is elected.
//! surviving standby. KV2 and Transit must recover after the bounded circuit
//! interval, use a bounded number of attempts, and leave the circuit and
//! in-flight gauges at zero after a new leader is elected.
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use metrics_util::MetricKind;
@@ -43,6 +43,11 @@ const OPERATION_ATTEMPTS: &str = "rustfs_kms_backend_operation_attempts";
const IN_FLIGHT: &str = "rustfs_kms_backend_in_flight";
const CIRCUIT_OPEN: &str = "rustfs_kms_backend_circuit_open";
const MAX_ATTEMPTS: u32 = 10;
const ATTEMPT_TIMEOUT: Duration = Duration::from_secs(2);
const HEALTHY_PROGRESS_TIMEOUT: Duration = Duration::from_secs(20);
// The circuit remains open for 30s after five failed attempts.
const POST_FAILOVER_PROGRESS_TIMEOUT: Duration = Duration::from_secs(35);
const FAILOVER_ERROR_POLL_INTERVAL: Duration = Duration::from_millis(100);
type MetricEntry = (
metrics_util::CompositeKey,
@@ -64,7 +69,7 @@ fn config(backend: KmsBackend, backend_config: BackendConfig) -> KmsConfig {
backend,
backend_config,
allow_insecure_dev_defaults: true,
timeout: Duration::from_secs(2),
timeout: ATTEMPT_TIMEOUT,
retry_attempts: MAX_ATTEMPTS,
enable_cache: false,
..KmsConfig::default()
@@ -164,14 +169,31 @@ fn retryable_failures(snapshot: &[MetricEntry], operation: &str) -> u64 {
.sum()
}
async fn wait_for_count(counter: &AtomicU64, minimum: u64, description: &str) {
tokio::time::timeout(Duration::from_secs(20), async {
async fn wait_for_count(
counter: &AtomicU64,
failure: &Mutex<Option<String>>,
minimum: u64,
description: &str,
timeout: Duration,
) {
tokio::time::timeout(timeout, async {
while counter.load(Ordering::SeqCst) < minimum {
if let Some(error) = failure.lock().expect("decrypt failure lock poisoned").as_ref() {
panic!(
"{description} worker failed after {} successful decrypts: {error}",
counter.load(Ordering::SeqCst)
);
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await
.unwrap_or_else(|_| panic!("timed out waiting for {description}"));
.unwrap_or_else(|_| {
panic!(
"timed out after {timeout:?} waiting for {description}: completed {}, expected {minimum}",
counter.load(Ordering::SeqCst)
)
});
}
async fn wait_for_file(path: &Path, description: &str) {
@@ -189,7 +211,8 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
request: DecryptRequest,
expected: Vec<u8>,
completed: Arc<AtomicU64>,
failed: Arc<AtomicBool>,
allow_failover_errors: Arc<AtomicBool>,
failure: Arc<Mutex<Option<String>>>,
stop: CancellationToken,
) {
while !stop.is_cancelled() {
@@ -197,8 +220,18 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
Ok(response) if response.plaintext == expected => {
completed.fetch_add(1, Ordering::SeqCst);
}
Ok(_) | Err(_) => {
failed.store(true, Ordering::SeqCst);
Ok(_) => {
*failure.lock().expect("decrypt failure lock poisoned") =
Some("decrypt returned unexpected plaintext".to_string());
return;
}
Err(rustfs_kms::KmsError::BackendError { .. } | rustfs_kms::KmsError::OperationTimedOut { .. })
if allow_failover_errors.load(Ordering::SeqCst) =>
{
tokio::time::sleep(FAILOVER_ERROR_POLL_INTERVAL).await;
}
Err(error) => {
*failure.lock().expect("decrypt failure lock poisoned") = Some(error.to_string());
return;
}
}
@@ -296,7 +329,9 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
);
let stop = CancellationToken::new();
let failed = Arc::new(AtomicBool::new(false));
let allow_failover_errors = Arc::new(AtomicBool::new(false));
let kv2_failure = Arc::new(Mutex::new(None));
let transit_failure = Arc::new(Mutex::new(None));
let kv2_completed = Arc::new(AtomicU64::new(0));
let transit_completed = Arc::new(AtomicU64::new(0));
let kv2_worker = tokio::spawn(decrypt_loop(
@@ -304,7 +339,8 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
kv2_request,
kv2_data_key.plaintext_key,
Arc::clone(&kv2_completed),
Arc::clone(&failed),
Arc::clone(&allow_failover_errors),
Arc::clone(&kv2_failure),
stop.clone(),
));
let transit_worker = tokio::spawn(decrypt_loop(
@@ -312,12 +348,21 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
transit_request,
transit_data_key.plaintext_key,
Arc::clone(&transit_completed),
Arc::clone(&failed),
Arc::clone(&allow_failover_errors),
Arc::clone(&transit_failure),
stop.clone(),
));
wait_for_count(&kv2_completed, 2, "two healthy KV2 decrypts").await;
wait_for_count(&transit_completed, 2, "two healthy Transit decrypts").await;
wait_for_count(&kv2_completed, &kv2_failure, 2, "two healthy KV2 decrypts", HEALTHY_PROGRESS_TIMEOUT).await;
wait_for_count(
&transit_completed,
&transit_failure,
2,
"two healthy Transit decrypts",
HEALTHY_PROGRESS_TIMEOUT,
)
.await;
allow_failover_errors.store(true, Ordering::SeqCst);
std::fs::write(&marker, b"ready").expect("publish failover readiness marker");
wait_for_file(&elected, "the replacement Vault leader").await;
@@ -326,18 +371,39 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
let kv2_after_election = kv2_completed.load(Ordering::SeqCst) + 2;
let transit_after_election = transit_completed.load(Ordering::SeqCst) + 2;
wait_for_count(&kv2_completed, kv2_after_election, "post-failover KV2 decrypts").await;
wait_for_count(&transit_completed, transit_after_election, "post-failover Transit decrypts").await;
wait_for_count(
&kv2_completed,
&kv2_failure,
kv2_after_election,
"post-failover KV2 decrypts",
POST_FAILOVER_PROGRESS_TIMEOUT,
)
.await;
wait_for_count(
&transit_completed,
&transit_failure,
transit_after_election,
"post-failover Transit decrypts",
POST_FAILOVER_PROGRESS_TIMEOUT,
)
.await;
stop.cancel();
kv2_worker.await.expect("KV2 decrypt worker must join");
transit_worker.await.expect("Transit decrypt worker must join");
assert!(!failed.load(Ordering::SeqCst), "no decrypt may fail or return different plaintext");
assert!(
kv2_failure.lock().expect("KV2 failure lock poisoned").is_none(),
"no KV2 decrypt may fail or return different plaintext"
);
assert!(
transit_failure.lock().expect("Transit failure lock poisoned").is_none(),
"no Transit decrypt may fail or return different plaintext"
);
}
#[test]
#[ignore = "requires a real three-node Vault Raft cluster; run scripts/test/vault_ha_kms_live.sh"]
fn vault_raft_leader_failure_preserves_kv2_and_transit_decrypts() {
fn vault_raft_leader_failure_recovers_kv2_and_transit_decrypts() {
let recorder = DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
metrics::with_local_recorder(&recorder, || {
@@ -349,11 +415,6 @@ fn vault_raft_leader_failure_preserves_kv2_and_transit_decrypts() {
});
let snapshot = snapshotter.snapshot().into_vec();
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "circuit_open")]),
0,
"a bounded leader election must not open the circuit"
);
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "budget_exhausted")]),
0,
@@ -1215,10 +1215,7 @@ pub struct ScannerActivityResponse {
pub dirty_usage_pending: bool,
}
#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
pub struct BackgroundHealStatusRequest {
#[prost(uint32, tag = "1")]
pub protocol_version: u32,
}
pub struct BackgroundHealStatusRequest {}
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct BackgroundHealStatusResponse {
#[prost(bool, tag = "1")]
-19
View File
@@ -170,7 +170,6 @@ pub fn internode_rpc_max_message_size() -> usize {
pub const HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE: usize = heal_control::RESULT_MAX_SIZE + 1024;
pub const HEAL_CONTROL_PROTOCOL_VERSION: u32 = 3;
pub const DYNAMIC_CONFIG_PROTOCOL_VERSION: u32 = 1;
pub const BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION: u32 = 2;
pub const HEAL_CONTROL_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-heal-control-capability-v3\0";
pub const REMOTE_VERSION_STATE_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-tier-remote-version-state-capability-v1\0";
pub const CROSS_POOL_FENCE_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-cross-pool-fence-capability-v1\0";
@@ -2347,28 +2346,10 @@ pub async fn evict_failed_connection_with_log_level(addr: &str, log_level: Conne
#[cfg(test)]
mod tests {
use super::*;
use prost::Message as _;
use std::sync::Mutex;
static INTERNODE_RPC_MSGPACK_ONLY_ENV_LOCK: Mutex<()> = Mutex::new(());
#[derive(Clone, PartialEq, prost::Message)]
struct BackgroundHealStatusRequestV1 {}
#[test]
fn background_heal_status_request_remains_rolling_upgrade_compatible() {
let current = proto_gen::node_service::BackgroundHealStatusRequest {
protocol_version: BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION,
};
let encoded = current.encode_to_vec();
BackgroundHealStatusRequestV1::decode(encoded.as_slice()).expect("v1 server should ignore the version field");
let encoded = BackgroundHealStatusRequestV1 {}.encode_to_vec();
let decoded = proto_gen::node_service::BackgroundHealStatusRequest::decode(encoded.as_slice())
.expect("v2 server should accept a v1 request");
assert_eq!(decoded.protocol_version, 0);
}
#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
struct CompatPayloadField {
message: &'static str,
+1 -3
View File
@@ -846,9 +846,7 @@ message ScannerActivityResponse {
bool dirty_usage_pending = 9;
}
message BackgroundHealStatusRequest {
uint32 protocol_version = 1;
}
message BackgroundHealStatusRequest {}
message BackgroundHealStatusResponse {
bool success = 1;
+26 -46
View File
@@ -20,7 +20,7 @@ use crate::admin::storage_api::bucket::utils::is_valid_object_prefix;
use crate::server::ADMIN_PREFIX;
use crate::server::RemoteAddr;
use crate::storage::rpc::node_service::heal::{
HealControlCoordinator, NodeHealStatusSnapshot, capture_node_heal_status, decode_node_heal_status,
HealControlCoordinator, NodeHealProgress, NodeHealStatusSnapshot, capture_node_heal_status, decode_node_heal_status,
decode_node_replacement_recovery_status, heal_control_coordinator, heal_topology_fingerprint,
};
use bytes::Bytes;
@@ -298,7 +298,14 @@ fn background_heal_runtime_state(
}
}
type BackgroundHealProgress = rustfs_heal::HealProgress;
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct BackgroundHealProgress {
objects_scanned: u64,
objects_healed: u64,
objects_failed: u64,
bytes_processed: u64,
}
#[derive(Debug)]
struct ClusterHealStatusSnapshot {
@@ -337,10 +344,17 @@ fn add_operations(total: &mut rustfs_heal::HealOperationsSnapshot, next: rustfs_
add_source_counts(&mut total.retrying_by_source, next.retrying_by_source);
}
fn add_progress(total: &mut BackgroundHealProgress, next: NodeHealProgress) {
total.objects_scanned = total.objects_scanned.saturating_add(next.objects_scanned);
total.objects_healed = total.objects_healed.saturating_add(next.objects_healed);
total.objects_failed = total.objects_failed.saturating_add(next.objects_failed);
total.bytes_processed = total.bytes_processed.saturating_add(next.bytes_processed);
}
fn aggregate_cluster_heal_status(snapshots: Vec<NodeHealStatusSnapshot>) -> ClusterHealStatusSnapshot {
let mut info = BackgroundHealInfo::default();
let mut operations = rustfs_heal::HealOperationsSnapshot::default();
let mut progress = Vec::new();
let mut progress = None;
let mut any_services_enabled = false;
let mut any_initialized = false;
@@ -357,12 +371,18 @@ fn aggregate_cluster_heal_status(snapshots: Vec<NodeHealStatusSnapshot>) -> Clus
}
add_operations(&mut operations, snapshot.operations);
if let Some(next) = snapshot.progress {
progress.push(next);
add_progress(
progress.get_or_insert(BackgroundHealProgress {
objects_scanned: 0,
objects_healed: 0,
objects_failed: 0,
bytes_processed: 0,
}),
next,
);
}
}
let progress = rustfs_heal::aggregate_heal_progress(progress);
let state = if operations.queue_length > 0 || operations.active_tasks > 0 || operations.retrying_tasks > 0 {
HealRuntimeState::Active
} else if any_initialized {
@@ -2181,19 +2201,10 @@ mod tests {
};
let progress = BackgroundHealProgress {
kind: rustfs_heal::heal::progress::HealProgressKind::ObjectSweep,
objects_scanned: 7,
objects_healed: 3,
objects_failed: 1,
skipped_objects: 3,
objects_total_count: 10,
objects_total_size: 8192,
bytes_processed: 4096,
progress_percentage: 50.0,
progress_state: rustfs_heal::heal::progress::HealProgressState::Running,
baseline_generation: Some(42),
baseline_known: true,
..Default::default()
};
let encoded = encode_background_heal_status(
@@ -2209,14 +2220,7 @@ mod tests {
assert_eq!(json["progress"]["objectsScanned"], 7);
assert_eq!(json["progress"]["objectsHealed"], 3);
assert_eq!(json["progress"]["objectsFailed"], 1);
assert_eq!(json["progress"]["skippedObjects"], 3);
assert_eq!(json["progress"]["objectsTotalCount"], 10);
assert_eq!(json["progress"]["objectsTotalSize"], 8192);
assert_eq!(json["progress"]["bytesProcessed"], 4096);
assert_eq!(json["progress"]["progressState"], "running");
assert_eq!(json["progress"]["baselineGeneration"], 42);
assert_eq!(json["progress"]["baselineKnown"], true);
assert_eq!(json["progress"]["counterUnknown"], false);
}
#[test]
@@ -2238,18 +2242,10 @@ mod tests {
..Default::default()
},
Some(NodeHealProgress {
kind: rustfs_heal::heal::progress::HealProgressKind::ObjectSweep,
objects_scanned: 3,
objects_healed: 1,
objects_failed: 0,
skipped_objects: 2,
objects_total_count: 6,
objects_total_size: 400,
bytes_processed: 100,
progress_state: rustfs_heal::heal::progress::HealProgressState::Running,
baseline_generation: Some(9),
baseline_known: true,
..Default::default()
}),
);
let peer = NodeHealStatusSnapshot::for_test(
@@ -2269,17 +2265,10 @@ mod tests {
..Default::default()
},
Some(NodeHealProgress {
kind: rustfs_heal::heal::progress::HealProgressKind::ObjectSweep,
objects_scanned: 5,
objects_healed: 4,
objects_failed: 1,
objects_total_count: 4,
objects_total_size: 1600,
bytes_processed: 900,
progress_state: rustfs_heal::heal::progress::HealProgressState::Running,
baseline_generation: Some(9),
baseline_known: true,
..Default::default()
}),
);
@@ -2298,15 +2287,7 @@ mod tests {
assert_eq!(progress.objects_scanned, 8);
assert_eq!(progress.objects_healed, 5);
assert_eq!(progress.objects_failed, 1);
assert_eq!(progress.skipped_objects, 2);
assert_eq!(progress.objects_total_count, 10);
assert_eq!(progress.objects_total_size, 2000);
assert_eq!(progress.bytes_processed, 1000);
assert_eq!(progress.progress_percentage, 50.0);
assert_eq!(progress.progress_state, rustfs_heal::heal::progress::HealProgressState::Running);
assert_eq!(progress.baseline_generation, Some(9));
assert!(progress.baseline_known);
assert!(!progress.counter_unknown);
assert_eq!(peer_first.state, HealRuntimeState::Active);
assert_eq!(peer_first.operations, local_first.operations);
@@ -2346,7 +2327,6 @@ mod tests {
objects_healed: value,
objects_failed: value,
bytes_processed: value,
..Default::default()
};
let saturated = NodeHealStatusSnapshot::for_test(
true,
+255 -38
View File
@@ -177,12 +177,15 @@ async fn rollback_cluster_rebalance_start(
terminal_reload_attempt_at: Some(terminal_reload_attempt_at),
terminal_reload_failures: terminal_reload_failures.clone(),
};
store.record_rebalance_stop_propagation(record).await.map_err(|err| {
format!(
"cluster rebalance rollback for {rebalance_id} partial; failed to persist stop propagation: {err}; {}",
rebalance_rollback_failure_message(rebalance_id, &stop_failures, &terminal_reload_failures)
)
})?;
store
.record_rebalance_stop_propagation(rebalance_id, record)
.await
.map_err(|err| {
format!(
"cluster rebalance rollback for {rebalance_id} partial; failed to persist stop propagation: {err}; {}",
rebalance_rollback_failure_message(rebalance_id, &stop_failures, &terminal_reload_failures)
)
})?;
return Err(rebalance_rollback_failure_message(
rebalance_id,
&stop_failures,
@@ -197,7 +200,7 @@ async fn rollback_cluster_rebalance_start(
.await
.map_err(|err| format!("local stop_rebalance rollback for {rebalance_id} failed: {err}"))?;
store
.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt)
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, rebalance_id)
.await
.map_err(|err| format!("local rollback stop metadata save for {rebalance_id} failed: {err}"))?;
Ok(())
@@ -679,7 +682,7 @@ impl Operation for RebalanceStart {
terminal_reload_attempt_at: Some(terminal_reload_attempt_at),
terminal_reload_failures: terminal_reload_failures.clone(),
};
store.record_rebalance_stop_propagation(record).await.map_err(|err| {
store.record_rebalance_stop_propagation(&id, record).await.map_err(|err| {
rebalance_internal_error(format!(
"failed to persist rebalance local-start rollback propagation metadata: {err}"
))
@@ -869,6 +872,38 @@ impl Operation for RebalanceStatus {
}
}
async fn rebalance_stop_target_id(store: &Arc<ECStore>) -> S3Result<Option<String>> {
store
.prepare_rebalance_stop()
.await
.map_err(|e| s3_error!(InternalError, "failed to prepare rebalance metadata for stop: {}", e))
}
async fn stop_rebalance_admission_first(
store: &Arc<ECStore>,
notification_sys: Option<&NotificationSys>,
expected_rebalance_id: &str,
) -> S3Result<Vec<String>> {
// prepare_rebalance_stop already closed admission for this exact run.
if let Some(notification_sys) = notification_sys {
return notification_sys
.stop_rebalance_failures(Some(expected_rebalance_id))
.await
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance via notification system: {}", e));
}
store
.stop_rebalance_for_id(Some(expected_rebalance_id))
.await
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance: {}", e))?;
store
.save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_rebalance_id)
.await
.map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop metadata: {}", e))?;
Ok(Vec::new())
}
// RebalanceStop
pub struct RebalanceStop {}
@@ -916,36 +951,15 @@ impl Operation for RebalanceStop {
return Err(s3_error!(InternalError, "object layer is not initialized"));
};
store
.load_rebalance_meta()
.await
.map_err(|e| s3_error!(InternalError, "failed to load rebalance metadata before stop: {}", e))?;
let expected_rebalance_id = store.current_rebalance_id().await;
if !store.is_rebalance_conflicting_with_decommission().await {
let Some(expected_rebalance_id) = rebalance_stop_target_id(&store).await? else {
log_rebalance_request_rejected("stop", "rebalance_not_started", &request_id, &actor, &remote_addr);
return Err(s3_error!(NoSuchResource, "pool rebalance is not started"));
}
};
let notification_sys = current_notification_system();
let stop_attempt_at = OffsetDateTime::now_utc();
let mut stop_failures = Vec::new();
if let Some(notification_sys) = notification_sys.as_ref() {
stop_failures = notification_sys
.stop_rebalance_failures(expected_rebalance_id.as_deref())
.await
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance via notification system: {}", e))?;
} else {
store
.stop_rebalance_for_id(expected_rebalance_id.as_deref())
.await
.map_err(|e| s3_error!(InternalError, "failed to stop rebalance: {}", e))?;
store
.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt)
.await
.map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop metadata: {}", e))?;
}
let stop_failures =
stop_rebalance_admission_first(&store, notification_sys.as_deref(), expected_rebalance_id.as_str()).await?;
info!(
event = EVENT_ADMIN_REBALANCE_STATE,
@@ -1007,7 +1021,7 @@ impl Operation for RebalanceStop {
terminal_reload_failures: terminal_reload_failures.clone(),
};
store
.record_rebalance_stop_propagation(record)
.record_rebalance_stop_propagation(expected_rebalance_id.as_str(), record)
.await
.map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop propagation metadata: {}", e))?;
@@ -1081,15 +1095,218 @@ mod rebalance_handler_tests {
RebalPoolProgress, RebalanceAdminStatus, RebalancePoolStatus, RebalanceStartStep, RebalanceStopPropagationStatus,
build_rebalance_admin_status, build_rebalance_pool_statuses, build_rebalance_stop_propagation_status,
rebalance_pool_used, rebalance_query_present, rebalance_remaining_buckets, rebalance_rollback_failure_message,
rebalance_rollback_stop_failure_message, rebalance_start_rollback_error, rebalance_start_steps, rebalance_used_pct,
rollback_result_label,
rebalance_rollback_stop_failure_message, rebalance_start_rollback_error, rebalance_start_steps, rebalance_stop_target_id,
rebalance_used_pct, rollback_result_label, stop_rebalance_admission_first,
};
use crate::admin::storage_api::rebalance::{
DiskStat, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo, RebalanceMeta,
RebalanceStats, RebalanceStopPropagationRecord, encode_rebalance_stop_propagation_record,
DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo,
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, encode_rebalance_stop_propagation_record,
};
use time::OffsetDateTime;
fn started_rebalance_meta(id: &str) -> RebalanceMeta {
RebalanceMeta {
id: id.to_string(),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
start_time: Some(OffsetDateTime::now_utc()),
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
}
}
#[tokio::test]
#[serial_test::serial]
async fn real_admin_stop_cancels_paused_entry_before_waiting_for_activation_gate() {
const REBALANCE_ID: &str = "admin-stop-paused-entry";
let mut fixture =
crate::admin::storage_api::ecstore_rebalance::test_util::PausedRebalanceEntryTestFixture::new(REBALANCE_ID).await;
fixture.wait_until_entry_paused().await;
let stop_store = fixture.store();
let mut stop_task = tokio::spawn(async move {
let expected_rebalance_id = rebalance_stop_target_id(&stop_store)
.await
.expect("admin stop target resolution should succeed")
.expect("the active rebalance should remain stoppable");
stop_rebalance_admission_first(&stop_store, None, expected_rebalance_id.as_str()).await
});
fixture.wait_until_admission_cancelled().await;
fixture.wait_until_stop_waiting_for_entry().await;
assert!(!stop_task.is_finished(), "admin stop must wait for the paused entry guard to drain");
fixture.release_entry();
fixture.assert_entry_cancelled().await;
let stop_failures = tokio::time::timeout(std::time::Duration::from_secs(5), &mut stop_task)
.await
.expect("admin stop should finish after the entry guard drains")
.expect("admin stop task should not panic")
.expect("admin stop should persist the terminal state");
assert!(stop_failures.is_empty());
assert!(!fixture.store().is_rebalance_conflicting_with_decommission().await);
}
#[tokio::test]
#[serial_test::serial]
async fn real_admin_stop_accepts_same_run_terminalization_after_prepare() {
const REBALANCE_ID: &str = "admin-stop-terminal-after-prepare";
const REPLACEMENT_ID: &str = "admin-stop-replacement";
let (_temp_dirs, store) =
crate::admin::storage_api::ecstore_rebalance::test_util::test_store_with_persisted_rebalance_meta(
started_rebalance_meta(REBALANCE_ID),
)
.await;
let terminal_barrier = std::sync::Arc::new(tokio::sync::Barrier::new(2));
let worker_barrier = std::sync::Arc::clone(&terminal_barrier);
let worker_store = std::sync::Arc::clone(&store);
let terminal_task = tokio::spawn(async move {
worker_barrier.wait().await;
{
let mut rebalance_meta = worker_store.rebalance_meta.write().await;
let meta = rebalance_meta
.as_mut()
.expect("the prepared rebalance metadata should remain installed");
assert_eq!(meta.id, REBALANCE_ID);
let pool = meta
.pool_stats
.first_mut()
.expect("the prepared rebalance should have a pool");
pool.info.status = RebalStatus::Stopped;
pool.info.end_time = Some(OffsetDateTime::now_utc());
}
worker_store
.save_rebalance_stats_for_id(0, RebalSaveOpt::Stats, REBALANCE_ID)
.await
});
let expected_rebalance_id = rebalance_stop_target_id(&store)
.await
.expect("admin stop target resolution should succeed")
.expect("the active rebalance should remain stoppable");
assert_eq!(expected_rebalance_id, REBALANCE_ID);
let cancel = store
.rebalance_meta
.read()
.await
.as_ref()
.and_then(|meta| meta.cancel.clone())
.expect("prepare should install the admission cancellation token");
assert!(cancel.is_cancelled());
terminal_barrier.wait().await;
tokio::time::timeout(std::time::Duration::from_secs(30), terminal_task)
.await
.expect("worker terminalization should finish after the barrier opens")
.expect("worker terminalization task should not panic")
.expect("worker terminalization should persist");
assert!(!store.is_rebalance_conflicting_with_decommission().await);
*store.rebalance_meta.write().await = None;
store
.load_rebalance_meta()
.await
.expect("the worker terminal state should reload before the final stop");
{
let terminal = store.rebalance_meta.read().await;
let terminal = terminal.as_ref().expect("the worker terminal state should remain persisted");
assert_eq!(terminal.id, REBALANCE_ID);
assert_eq!(terminal.pool_stats[0].info.status, RebalStatus::Stopped);
}
let stop_failures = stop_rebalance_admission_first(&store, None, expected_rebalance_id.as_str())
.await
.expect("same-run terminalization after prepare should be a successful stop");
assert!(stop_failures.is_empty());
*store.rebalance_meta.write().await = Some(started_rebalance_meta(REPLACEMENT_ID));
let error = stop_rebalance_admission_first(&store, None, expected_rebalance_id.as_str())
.await
.expect_err("the prepared stop must not mutate a replacement run");
assert!(error.to_string().contains(REBALANCE_ID));
let replacement = store.rebalance_meta.read().await;
let replacement = replacement.as_ref().expect("the replacement run should remain installed");
assert_eq!(replacement.id, REPLACEMENT_ID);
assert_eq!(replacement.pool_stats[0].info.status, RebalStatus::Started);
assert!(replacement.cancel.is_none());
assert!(replacement.stopped_at.is_none());
}
#[tokio::test]
#[serial_test::serial]
async fn real_admin_stop_loads_persisted_active_rebalance_from_cold_memory() {
const REBALANCE_ID: &str = "admin-stop-cold-memory";
let (_temp_dirs, store) =
crate::admin::storage_api::ecstore_rebalance::test_util::test_store_with_persisted_rebalance_meta(
started_rebalance_meta(REBALANCE_ID),
)
.await;
*store.rebalance_meta.write().await = None;
assert!(store.current_rebalance_id().await.is_none());
let expected_rebalance_id = rebalance_stop_target_id(&store)
.await
.expect("admin stop should load persisted rebalance metadata")
.expect("persisted active rebalance should be stoppable");
assert_eq!(expected_rebalance_id, REBALANCE_ID);
let stop_failures = stop_rebalance_admission_first(&store, None, expected_rebalance_id.as_str())
.await
.expect("admin stop should persist the terminal state after a cold load");
assert!(stop_failures.is_empty());
assert!(!store.is_rebalance_conflicting_with_decommission().await);
*store.rebalance_meta.write().await = None;
store
.load_rebalance_meta()
.await
.expect("the persisted terminal rebalance metadata should remain readable");
assert_eq!(store.current_rebalance_id().await.as_deref(), Some(REBALANCE_ID));
assert!(!store.is_rebalance_conflicting_with_decommission().await);
}
#[tokio::test]
#[serial_test::serial]
async fn real_admin_stop_refreshes_persisted_active_over_stale_inactive_memory() {
const PERSISTED_REBALANCE_ID: &str = "admin-stop-persisted-active";
const STALE_REBALANCE_ID: &str = "admin-stop-stale-terminal";
let (_temp_dirs, store) =
crate::admin::storage_api::ecstore_rebalance::test_util::test_store_with_persisted_rebalance_meta(
started_rebalance_meta(PERSISTED_REBALANCE_ID),
)
.await;
*store.rebalance_meta.write().await = Some(RebalanceMeta {
id: STALE_REBALANCE_ID.to_string(),
stopped_at: Some(OffsetDateTime::now_utc()),
..Default::default()
});
assert_eq!(store.current_rebalance_id().await.as_deref(), Some(STALE_REBALANCE_ID));
assert!(!store.is_rebalance_conflicting_with_decommission().await);
let expected_rebalance_id = rebalance_stop_target_id(&store)
.await
.expect("admin stop should refresh stale inactive local metadata")
.expect("persisted active rebalance should replace the stale local terminal state");
assert_eq!(expected_rebalance_id, PERSISTED_REBALANCE_ID);
let stop_failures = stop_rebalance_admission_first(&store, None, expected_rebalance_id.as_str())
.await
.expect("admin stop should persist the refreshed run's terminal state");
assert!(stop_failures.is_empty());
*store.rebalance_meta.write().await = None;
store
.load_rebalance_meta()
.await
.expect("the refreshed run's persisted terminal metadata should remain readable");
assert_eq!(store.current_rebalance_id().await.as_deref(), Some(PERSISTED_REBALANCE_ID));
assert!(!store.is_rebalance_conflicting_with_decommission().await);
}
#[test]
fn test_calculate_rebalance_progress_running() {
let start = OffsetDateTime::from_unix_timestamp(1_000).unwrap();
+3 -1
View File
@@ -70,7 +70,9 @@ mod ecstore_notification {
}
#[allow(unused_imports)]
mod ecstore_rebalance {
pub(crate) mod ecstore_rebalance {
#[cfg(test)]
pub(crate) use crate::storage::storage_api::ecstore_rebalance::test_util;
pub(crate) use crate::storage::storage_api::ecstore_rebalance::{
DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo,
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, decode_rebalance_stop_propagation_record,
+2 -2
View File
@@ -1912,7 +1912,7 @@ impl Node for NodeService {
async fn background_heal_status(
&self,
request: Request<BackgroundHealStatusRequest>,
_request: Request<BackgroundHealStatusRequest>,
) -> Result<Response<BackgroundHealStatusResponse>, Status> {
if self.resolve_object_store().is_none() {
return Ok(Response::new(BackgroundHealStatusResponse {
@@ -1922,7 +1922,7 @@ impl Node for NodeService {
}));
}
let snapshot = heal::capture_node_heal_status(rustfs_scanner::scanner::BackgroundHealInfo::default()).await;
match heal::encode_node_heal_status(&snapshot, request.into_inner().protocol_version) {
match heal::encode_node_heal_status(&snapshot) {
Ok(bg_heal_state) => Ok(Response::new(BackgroundHealStatusResponse {
success: true,
bg_heal_state: bg_heal_state.into(),
+29 -174
View File
@@ -25,8 +25,7 @@ use std::io::Cursor;
use super::super::encode_msgpack_map;
const NODE_HEAL_STATUS_PREVIOUS_VERSION: u8 = 1;
const NODE_HEAL_STATUS_VERSION: u8 = 2;
const NODE_HEAL_STATUS_VERSION: u8 = 1;
const NODE_HEAL_STATUS_MAX_SIZE: usize = 64 * 1024;
const NODE_REPLACEMENT_RECOVERY_STATUS_VERSION: u8 = 1;
const NODE_REPLACEMENT_RECOVERY_STATUS_MAX_SIZE: usize = 64 * 1024;
@@ -173,38 +172,13 @@ pub(crate) fn heal_topology_fingerprint(endpoint_pools: &EndpointServerPools) ->
Ok(hex_simd::encode_to_string(hasher.finalize(), hex_simd::AsciiCase::Lower))
}
pub(crate) type NodeHealProgress = rustfs_heal::HealProgress;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct NodeHealProgressV1 {
objects_scanned: u64,
objects_healed: u64,
objects_failed: u64,
bytes_processed: u64,
}
impl From<&NodeHealProgress> for NodeHealProgressV1 {
fn from(progress: &NodeHealProgress) -> Self {
Self {
objects_scanned: progress.objects_scanned,
objects_healed: progress.objects_healed,
objects_failed: progress.objects_failed,
bytes_processed: progress.bytes_processed,
}
}
}
impl From<NodeHealProgressV1> for NodeHealProgress {
fn from(progress: NodeHealProgressV1) -> Self {
Self {
objects_scanned: progress.objects_scanned,
objects_healed: progress.objects_healed,
objects_failed: progress.objects_failed,
bytes_processed: progress.bytes_processed,
..Default::default()
}
}
pub(crate) struct NodeHealProgress {
pub objects_scanned: u64,
pub objects_healed: u64,
pub objects_failed: u64,
pub bytes_processed: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -246,48 +220,6 @@ pub(crate) struct NodeHealStatusSnapshot {
pub progress: Option<NodeHealProgress>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct NodeHealStatusSnapshotV1 {
version: u8,
services_enabled: bool,
initialized: bool,
info: NodeHealInfo,
operations: HealOperationsSnapshot,
progress: Option<NodeHealProgressV1>,
}
impl From<&NodeHealStatusSnapshot> for NodeHealStatusSnapshotV1 {
fn from(snapshot: &NodeHealStatusSnapshot) -> Self {
Self {
version: NODE_HEAL_STATUS_PREVIOUS_VERSION,
services_enabled: snapshot.services_enabled,
initialized: snapshot.initialized,
info: snapshot.info.clone(),
operations: snapshot.operations,
progress: snapshot.progress.as_ref().map(NodeHealProgressV1::from),
}
}
}
impl From<NodeHealStatusSnapshotV1> for NodeHealStatusSnapshot {
fn from(snapshot: NodeHealStatusSnapshotV1) -> Self {
Self {
version: snapshot.version,
services_enabled: snapshot.services_enabled,
initialized: snapshot.initialized,
info: snapshot.info,
operations: snapshot.operations,
progress: snapshot.progress.map(NodeHealProgress::from),
}
}
}
#[derive(Deserialize)]
struct NodeHealStatusVersion {
version: u8,
}
impl NodeHealStatusSnapshot {
#[cfg(test)]
pub(crate) fn for_test(
@@ -317,7 +249,14 @@ impl NodeHealStatusSnapshot {
}
pub(crate) async fn capture_node_heal_status(info: BackgroundHealInfo) -> NodeHealStatusSnapshot {
let progress = rustfs_heal::current_heal_progress_snapshot().await;
let progress = rustfs_heal::current_heal_progress_snapshot()
.await
.map(|progress| NodeHealProgress {
objects_scanned: progress.objects_scanned,
objects_healed: progress.objects_healed,
objects_failed: progress.objects_failed,
bytes_processed: progress.bytes_processed,
});
NodeHealStatusSnapshot {
version: NODE_HEAL_STATUS_VERSION,
@@ -329,44 +268,23 @@ pub(crate) async fn capture_node_heal_status(info: BackgroundHealInfo) -> NodeHe
}
}
pub(crate) fn encode_node_heal_status(snapshot: &NodeHealStatusSnapshot, protocol_version: u32) -> Result<Vec<u8>, String> {
let encoded = if protocol_version < rustfs_protos::BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION {
encode_msgpack_map(&NodeHealStatusSnapshotV1::from(snapshot))
} else {
let mut snapshot = snapshot.clone();
snapshot.version = NODE_HEAL_STATUS_VERSION;
encode_msgpack_map(&snapshot)
};
encoded.map_err(|err| format!("failed to encode node heal status: {err}"))
pub(crate) fn encode_node_heal_status(snapshot: &NodeHealStatusSnapshot) -> Result<Vec<u8>, String> {
encode_msgpack_map(snapshot).map_err(|err| format!("failed to encode node heal status: {err}"))
}
pub(crate) fn decode_node_heal_status(data: &[u8]) -> Result<NodeHealStatusSnapshot, String> {
if data.len() > NODE_HEAL_STATUS_MAX_SIZE {
return Err("node heal status exceeds size limit".to_string());
}
let decode_version = || {
let mut deserializer = Deserializer::new(Cursor::new(data));
NodeHealStatusVersion::deserialize(&mut deserializer)
.map(|version| (version, deserializer))
.map_err(|err| format!("failed to decode node heal status: {err}"))
};
let (version, deserializer) = decode_version()?;
if usize::try_from(deserializer.get_ref().position()).ok() != Some(data.len()) {
return Err("node heal status contains trailing data".to_string());
}
let mut deserializer = Deserializer::new(Cursor::new(data));
let snapshot = match version.version {
NODE_HEAL_STATUS_PREVIOUS_VERSION => {
NodeHealStatusSnapshotV1::deserialize(&mut deserializer).map(NodeHealStatusSnapshot::from)
}
NODE_HEAL_STATUS_VERSION => NodeHealStatusSnapshot::deserialize(&mut deserializer),
version => return Err(format!("unsupported node heal status version: {version}")),
}
.map_err(|err| format!("failed to decode node heal status: {err}"))?;
let snapshot = NodeHealStatusSnapshot::deserialize(&mut deserializer)
.map_err(|err| format!("failed to decode node heal status: {err}"))?;
if usize::try_from(deserializer.get_ref().position()).ok() != Some(data.len()) {
return Err("node heal status contains trailing data".to_string());
}
if snapshot.version != NODE_HEAL_STATUS_VERSION {
return Err(format!("unsupported node heal status version: {}", snapshot.version));
}
Ok(snapshot)
}
@@ -469,10 +387,9 @@ pub(crate) fn decode_node_replacement_recovery_status(data: &[u8]) -> Result<Nod
#[cfg(test)]
mod tests {
use super::{
NODE_HEAL_STATUS_MAX_SIZE, NODE_HEAL_STATUS_PREVIOUS_VERSION, NODE_HEAL_STATUS_VERSION, NodeHealProgress,
NodeHealStatusSnapshot, NodeReplacementRecoveryStatusSnapshot, decode_node_heal_status,
decode_node_replacement_recovery_status, encode_node_heal_status, encode_node_replacement_recovery_status,
heal_control_coordinator, heal_topology_fingerprint,
NODE_HEAL_STATUS_MAX_SIZE, NODE_HEAL_STATUS_VERSION, NodeHealProgress, NodeHealStatusSnapshot,
NodeReplacementRecoveryStatusSnapshot, decode_node_heal_status, decode_node_replacement_recovery_status,
encode_node_heal_status, encode_node_replacement_recovery_status, heal_control_coordinator, heal_topology_fingerprint,
};
use crate::storage::storage_api::{
Endpoint,
@@ -616,72 +533,20 @@ mod tests {
..Default::default()
},
Some(NodeHealProgress {
kind: rustfs_heal::heal::progress::HealProgressKind::ObjectSweep,
objects_scanned: 7,
objects_healed: 5,
objects_failed: 1,
skipped_objects: 1,
objects_total_count: 10,
objects_total_size: 2048,
bytes_processed: 1024,
progress_percentage: 50.0,
progress_state: rustfs_heal::heal::progress::HealProgressState::Running,
baseline_generation: Some(42),
baseline_known: true,
..Default::default()
}),
);
let encoded = encode_node_heal_status(&snapshot, rustfs_protos::BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION)
.expect("snapshot should encode");
let encoded = encode_node_heal_status(&snapshot).expect("snapshot should encode");
let decoded = decode_node_heal_status(&encoded).expect("snapshot should decode");
assert_eq!(decoded.version, NODE_HEAL_STATUS_VERSION);
assert_eq!(decoded.operations.queue_length, 2);
assert_eq!(decoded.progress, snapshot.progress);
}
#[test]
fn node_heal_status_v1_encoding_preserves_rolling_compatibility() {
let snapshot = NodeHealStatusSnapshot::for_test(
true,
true,
BackgroundHealInfo::default(),
HealOperationsSnapshot::default(),
Some(NodeHealProgress {
objects_scanned: 7,
objects_healed: 5,
objects_failed: 1,
skipped_objects: 3,
bytes_processed: 1024,
progress_state: rustfs_heal::heal::progress::HealProgressState::Running,
baseline_known: true,
..Default::default()
}),
);
let encoded =
encode_node_heal_status(&snapshot, u32::from(NODE_HEAL_STATUS_PREVIOUS_VERSION)).expect("v1 snapshot should encode");
let wire: serde_json::Value = rmp_serde::from_slice(&encoded).expect("v1 snapshot should decode as JSON");
let progress = wire["progress"].as_object().expect("v1 progress should be a map");
assert_eq!(wire["version"], NODE_HEAL_STATUS_PREVIOUS_VERSION);
assert_eq!(progress.len(), 4);
assert_eq!(progress["objectsScanned"], 7);
assert!(!progress.contains_key("skippedObjects"));
let decoded = decode_node_heal_status(&encoded).expect("v1 snapshot should decode");
let progress = decoded.progress.as_ref().expect("v1 progress should be present");
assert_eq!(decoded.version, NODE_HEAL_STATUS_PREVIOUS_VERSION);
assert_eq!(progress.objects_scanned, 7);
assert_eq!(progress.progress_state, rustfs_heal::heal::progress::HealProgressState::Unknown);
assert!(!progress.baseline_known);
let upgraded = encode_node_heal_status(&decoded, rustfs_protos::BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION)
.expect("decoded v1 snapshot should upgrade to v2");
let upgraded = decode_node_heal_status(&upgraded).expect("upgraded snapshot should decode");
assert_eq!(upgraded.version, NODE_HEAL_STATUS_VERSION);
assert_eq!(upgraded.progress.as_ref(), Some(progress));
}
#[test]
fn node_heal_status_rejects_unknown_version() {
let mut snapshot = NodeHealStatusSnapshot::for_test(
@@ -693,7 +558,7 @@ mod tests {
);
snapshot.version += 1;
let encoded = rmp_serde::to_vec_named(&snapshot).expect("snapshot should encode");
let encoded = encode_node_heal_status(&snapshot).expect("snapshot should encode");
let err = decode_node_heal_status(&encoded).expect_err("unknown version should fail closed");
assert!(err.contains("unsupported node heal status version"));
}
@@ -714,22 +579,13 @@ mod tests {
"activeBySource": {"scanner": 0, "admin": 1, "autoHeal": 0, "internal": 0, "readRepair": 0},
"retryingBySource": {"scanner": 0, "admin": 0, "autoHeal": 0, "internal": 0, "readRepair": 0}
},
"progress": {
"objectsScanned": 7,
"objectsHealed": 5,
"objectsFailed": 1,
"bytesProcessed": 1024
}
"progress": null
});
let encoded = rmp_serde::to_vec_named(&fixture).expect("fixture should encode");
let decoded = decode_node_heal_status(&encoded).expect("fixed v1 fixture should decode");
assert_eq!(decoded.info().bitrot_start_cycle, 9);
assert_eq!(decoded.operations.queue_length, 2);
assert_eq!(decoded.operations.queued_by_source.mrf, 0);
let progress = decoded.progress.expect("legacy progress should decode");
assert_eq!(progress.objects_scanned, 7);
assert!(!progress.baseline_known);
assert_eq!(progress.progress_state, rustfs_heal::heal::progress::HealProgressState::Unknown);
}
#[test]
@@ -775,8 +631,7 @@ mod tests {
HealOperationsSnapshot::default(),
None,
);
let encoded = encode_node_heal_status(&snapshot, rustfs_protos::BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION)
.expect("snapshot should encode");
let encoded = encode_node_heal_status(&snapshot).expect("snapshot should encode");
let encoded_json: serde_json::Value = rmp_serde::from_slice(&encoded).expect("encoded snapshot should decode as JSON");
assert_eq!(encoded_json["info"]["bitrotStartTime"], serde_json::json!("2023-11-14T22:13:20.123456Z"));
}
+2
View File
@@ -501,6 +501,8 @@ pub(crate) mod ecstore_notification {
#[allow(unused_imports)]
pub(crate) mod ecstore_rebalance {
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::rebalance::test_util;
pub(crate) use rustfs_ecstore::api::rebalance::{
DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo,
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, decode_rebalance_stop_propagation_record,
+1 -1
View File
@@ -241,7 +241,7 @@ env \
RUSTFS_TEST_VAULT_FAILOVER_MARKER="$MARKER" \
RUSTFS_TEST_VAULT_OLD_LEADER="$OLD_LEADER" \
cargo test -p rustfs-kms --test vault_ha_failover_live \
vault_raft_leader_failure_preserves_kv2_and_transit_decrypts -- \
vault_raft_leader_failure_recovers_kv2_and_transit_decrypts -- \
--ignored --nocapture --test-threads=1 &
TEST_PID=$!