Compare commits

..

1 Commits

Author SHA1 Message Date
马登山 16c56fe0bd fix(heal): correct progress accounting 2026-08-22 17:58:55 +08:00
23 changed files with 1248 additions and 1222 deletions
+17 -84
View File
@@ -50,10 +50,10 @@ use rustfs_protos::evict_failed_connection;
use rustfs_protos::proto_gen::node_service::RenamePartRequest;
use rustfs_protos::proto_gen::node_service::{
BatchReadVersionRequest, BatchReadVersionResponse, CheckPartsRequest, DeletePathsRequest, DeleteRequest,
DeleteVersionRequest, DeleteVersionsRequest, DeleteVersionsResponse, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest,
ListVolumesRequest, MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest,
ReadMetadataRequest, ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest,
RenameDataRequest, RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest,
DeleteVersionRequest, DeleteVersionsRequest, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest, ListVolumesRequest,
MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest, ReadMetadataRequest,
ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest, RenameDataRequest,
RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest,
SnapshotLeaseRequest, SnapshotLeaseResponse, StatVolumeRequest, UpdateMetadataRequest, VerifyFileRequest, WriteAllRequest,
WriteMetadataRequest, node_service_client::NodeServiceClient,
};
@@ -112,28 +112,6 @@ const EVENT_REMOTE_DISK_RPC: &str = "remote_disk_rpc";
const SNAPSHOT_LEASE_PROTOCOL_VERSION: u32 = 1;
pub const REMOTE_SNAPSHOT_LEASE_TTL: Duration = Duration::from_secs(60);
fn decode_delete_versions_errors(response: DeleteVersionsResponse, expected_len: usize) -> Vec<Option<Error>> {
if !response.item_errors.is_empty() {
if response.item_errors.len() != expected_len {
return vec![Some(Error::other("malformed delete_versions item errors")); expected_len];
}
return response
.item_errors
.into_iter()
.map(|error| (error.code != 0).then(|| error.into()))
.collect();
}
if response.errors.len() != expected_len {
return vec![Some(Error::other("malformed delete_versions errors")); expected_len];
}
response
.errors
.into_iter()
.map(|error| (!error.is_empty()).then(|| Error::other(error)))
.collect()
}
fn snapshot_lease_token_from_response(response: SnapshotLeaseResponse) -> Result<SnapshotLeaseToken> {
if !response.success {
return Err(response.error.unwrap_or_default().into());
@@ -2428,6 +2406,8 @@ impl DiskAPI for RemoteDisk {
return errors;
}
// TODO(backlog): replace string errors with typed `StorageError` variants
let result = self
.execute_with_timeout(
|| async {
@@ -2459,7 +2439,17 @@ impl DiskAPI for RemoteDisk {
}
return errors;
}
decode_delete_versions_errors(response, versions.len())
response
.errors
.iter()
.map(|error| {
if error.is_empty() {
None
} else {
Some(Error::other(error.to_string()))
}
})
.collect()
}
#[tracing::instrument(level = "trace", skip_all)]
@@ -3770,63 +3760,6 @@ mod tests {
static INIT: Once = Once::new();
#[test]
fn delete_versions_response_preserves_typed_item_errors() {
let errors = decode_delete_versions_errors(
DeleteVersionsResponse {
success: true,
errors: vec!["file not found".to_string(), String::new()],
error: None,
item_errors: vec![
rustfs_protos::proto_gen::node_service::Error {
code: DiskError::FileNotFound.to_u32(),
error_info: "file not found".to_string(),
},
rustfs_protos::proto_gen::node_service::Error::default(),
],
},
2,
);
assert!(matches!(errors.as_slice(), [Some(DiskError::FileNotFound), None]));
}
#[test]
fn delete_versions_response_accepts_legacy_string_errors() {
let errors = decode_delete_versions_errors(
DeleteVersionsResponse {
success: true,
errors: vec!["legacy error".to_string(), String::new()],
error: None,
item_errors: Vec::new(),
},
2,
);
assert_eq!(errors.len(), 2);
assert_eq!(errors[0].as_ref().map(ToString::to_string).as_deref(), Some("io error legacy error"));
assert!(errors[1].is_none());
}
#[test]
fn delete_versions_response_rejects_misaligned_item_errors() {
let errors = decode_delete_versions_errors(
DeleteVersionsResponse {
success: true,
errors: vec!["file not found".to_string()],
error: None,
item_errors: vec![rustfs_protos::proto_gen::node_service::Error {
code: DiskError::FileNotFound.to_u32(),
error_info: "file not found".to_string(),
}],
},
2,
);
assert_eq!(errors.len(), 2);
assert!(errors.iter().all(Option::is_some));
}
#[test]
fn disk_mutation_digest_marks_rolling_compatibility() {
let mut request = Request::new(());
+122 -507
View File
@@ -36,8 +36,7 @@ use crate::disk::error::DiskError;
use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET};
use crate::error::{Error, Result};
use crate::error::{
StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_object_not_found, is_err_operation_canceled,
is_err_version_not_found,
StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found,
};
use crate::layout::endpoints::EndpointServerPools;
use crate::object_api::{GetObjectReader, ObjectOptions};
@@ -90,7 +89,6 @@ const DECOMMISSION_STAGE_SOURCE_CLEANUP: &str = "source_cleanup";
const DECOMMISSION_STAGE_ENTRY_FINISHED: &str = "entry_finished";
const DECOMMISSION_PROGRESS_SAVE_INTERVAL: Duration = Duration::seconds(30);
const DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD: usize = 1000;
const DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF: Duration = Duration::seconds(1);
const DECOMMISSION_BUCKET_CONCURRENCY_ENV: &str = "RUSTFS_DECOMMISSION_BUCKET_CONCURRENCY";
const DECOMMISSION_BUCKET_CONCURRENCY_DEFAULT_CAP: usize = 4;
const DECOMMISSION_TARGET_CAPACITY_OVERHEAD_PERCENT: usize = 30;
@@ -640,6 +638,22 @@ fn track_decommission_current_object(meta: &mut PoolMeta, idx: usize, bucket: &s
track_decommission_current_object_stage(meta, idx, bucket, object, "")
}
fn touch_decommission_progress(meta: &mut PoolMeta, idx: usize) -> Result<()> {
let pool_count = meta.pools.len();
ensure_valid_decommission_pool_index(pool_count, idx)?;
let Some(pool) = meta.pools.get_mut(idx) else {
return Err(invalid_decommission_pool_index_error(pool_count, idx));
};
let Some(info) = pool.decommission.as_mut() else {
return Err(decommission_metadata_not_initialized_error("touch decommission progress"));
};
pool.last_update = OffsetDateTime::now_utc();
info.mark_progress_saved();
Ok(())
}
fn resolve_decommission_update_after_result(result: Result<bool>) -> Result<bool> {
result.map_err(|err| Error::other(format!("decommission metadata update failed: {err}")))
}
@@ -759,76 +773,7 @@ async fn load_decommission_entry_exact_versions(
}
fn resolve_decommission_check_after_list_result(list_result: Result<()>, entry_error: Option<Error>) -> Result<()> {
match list_result {
Ok(()) => entry_error.map_or(Ok(()), Err),
Err(list_err) => resolve_decommission_listing_error(Some(list_err), entry_error).map_or(Ok(()), Err),
}
}
fn resolve_decommission_listing_error(listing_error: Option<Error>, entry_error: Option<Error>) -> Option<Error> {
match (listing_error, entry_error) {
(Some(listing_error), Some(entry_error)) if is_err_operation_canceled(&listing_error) => Some(entry_error),
(Some(listing_error), Some(entry_error)) if is_err_operation_canceled(&entry_error) => Some(listing_error),
(Some(listing_error), _) => Some(listing_error),
(None, entry_error) => entry_error,
}
}
fn decommission_unresolved_listing_error(
bucket: &str,
prefix: &str,
candidate: Option<&str>,
candidate_count: usize,
disk_error_count: usize,
pool_index: usize,
set_index: usize,
) -> Error {
let location = candidate.unwrap_or(prefix);
Error::other(format!(
"decommission listing could not resolve metadata for {bucket}/{location} on pool {pool_index} set {set_index} ({candidate_count} candidate(s), {disk_error_count} disk error(s))"
))
}
fn resolve_decommission_partial_listing_entry(
entries: MetaCacheEntries,
resolver: MetadataResolutionParams,
bucket: &str,
prefix: &str,
disk_error_count: usize,
pool_index: usize,
set_index: usize,
) -> Result<MetaCacheEntry> {
let candidate_count = entries.as_ref().iter().flatten().count();
if let Some(entry) = entries.resolve(resolver) {
return Ok(entry);
}
let candidate = entries.as_ref().iter().flatten().map(|entry| entry.name.as_str()).next();
Err(decommission_unresolved_listing_error(
bucket,
prefix,
candidate,
candidate_count,
disk_error_count,
pool_index,
set_index,
))
}
async fn record_decommission_entry_error(
entry_error: &Arc<tokio::sync::Mutex<Option<Error>>>,
rx: &CancellationToken,
err: Error,
) {
if rx.is_cancelled() {
return;
}
let mut first_err = entry_error.lock().await;
if first_err.is_none() && !rx.is_cancelled() {
*first_err = Some(err);
rx.cancel();
}
if let Some(err) = entry_error { Err(err) } else { list_result }
}
fn resolve_decommission_pool_meta_reload_result(result: Result<()>, stage: &str) -> Result<()> {
@@ -1538,7 +1483,6 @@ impl TryFrom<PersistedPoolDecommissionInfo> for PoolDecommissionInfo {
terminal_reload_attempt_at: value.terminal_reload_attempt_at,
terminal_reload_failures: value.terminal_reload_failures,
progress_save_item_baseline: value.items_decommissioned.saturating_add(value.items_decommission_failed),
progress_save_retry_after: None,
})
}
}
@@ -1570,7 +1514,6 @@ impl TryFrom<LegacyPoolDecommissionInfo> for PoolDecommissionInfo {
terminal_reload_attempt_at: None,
terminal_reload_failures: Vec::new(),
progress_save_item_baseline: value.items_decommissioned.saturating_add(value.items_decommission_failed),
progress_save_retry_after: None,
})
}
}
@@ -1684,82 +1627,6 @@ impl PoolMeta {
}
}
fn decommission_progress_checkpoint(
&self,
idx: usize,
duration: Duration,
now: OffsetDateTime,
) -> Result<Option<DecommissionProgressCheckpoint>> {
let pool_count = self.pools.len();
ensure_valid_decommission_pool_index(pool_count, idx)?;
let Some(pool) = self.pools.get(idx) else {
return Err(invalid_decommission_pool_index_error(pool_count, idx));
};
let Some(info) = pool.decommission.as_ref() else {
return Err(decommission_metadata_not_initialized_error("update decommission metadata timestamp"));
};
if info.progress_save_retry_after.is_some_and(|retry_after| now < retry_after) {
return Ok(None);
}
let time_threshold_reached = now.unix_timestamp() - pool.last_update.unix_timestamp() >= duration.whole_seconds();
let item_threshold_reached = info.items_since_last_progress_save() >= DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD;
if !time_threshold_reached && !item_threshold_reached {
return Ok(None);
}
Ok(Some(DecommissionProgressCheckpoint {
start_time: info.start_time,
queued: info.queued,
counted_items: info.counted_items(),
checkpoint_at: now,
}))
}
fn commit_decommission_progress_checkpoint(&mut self, idx: usize, checkpoint: DecommissionProgressCheckpoint) -> bool {
let Some(pool) = self.pools.get_mut(idx) else {
return false;
};
let Some(info) = pool.decommission.as_mut() else {
return false;
};
if info.start_time != checkpoint.start_time
|| info.queued != checkpoint.queued
|| !is_decommission_active(info.complete, info.failed, info.canceled)
{
return false;
}
info.progress_save_item_baseline = info.progress_save_item_baseline.max(checkpoint.counted_items);
info.progress_save_retry_after = None;
pool.last_update = pool.last_update.max(checkpoint.checkpoint_at);
true
}
fn defer_decommission_progress_checkpoint(
&mut self,
idx: usize,
checkpoint: DecommissionProgressCheckpoint,
retry_after: OffsetDateTime,
) {
let Some(pool) = self.pools.get_mut(idx) else {
return;
};
let Some(info) = pool.decommission.as_mut() else {
return;
};
if info.start_time == checkpoint.start_time
&& info.queued == checkpoint.queued
&& is_decommission_active(info.complete, info.failed, info.canceled)
{
info.progress_save_retry_after = Some(retry_after);
}
}
fn load_from_config_data(&mut self, data: Vec<u8>) -> Result<()> {
if data.is_empty() {
return Ok(());
@@ -2120,9 +1987,30 @@ impl PoolMeta {
}
pub fn update_after(&mut self, idx: usize, duration: Duration) -> Result<bool> {
Ok(self
.decommission_progress_checkpoint(idx, duration, OffsetDateTime::now_utc())?
.is_some())
let pool_count = self.pools.len();
ensure_valid_decommission_pool_index(pool_count, idx)?;
let (last_update, item_threshold_reached) = match self.pools.get(idx) {
Some(pool) if let Some(info) = pool.decommission.as_ref() => (
pool.last_update,
info.items_since_last_progress_save() >= DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD,
),
Some(_) => {
return Err(decommission_metadata_not_initialized_error("update decommission metadata timestamp"));
}
None => return Err(invalid_decommission_pool_index_error(pool_count, idx)),
};
let now = OffsetDateTime::now_utc();
if now.unix_timestamp() - last_update.unix_timestamp() >= duration.whole_seconds() || item_threshold_reached {
let Some(pool) = self.pools.get_mut(idx) else {
return Err(invalid_decommission_pool_index_error(pool_count, idx));
};
pool.last_update = now;
return Ok(true);
}
Ok(false)
}
pub fn validate(&self, pools: Vec<Arc<Sets>>) -> Result<bool> {
@@ -2263,16 +2151,6 @@ pub struct PoolDecommissionInfo {
pub terminal_reload_failures: Vec<String>,
#[serde(skip)]
pub progress_save_item_baseline: usize,
#[serde(skip)]
pub progress_save_retry_after: Option<OffsetDateTime>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct DecommissionProgressCheckpoint {
start_time: Option<OffsetDateTime>,
queued: bool,
counted_items: usize,
checkpoint_at: OffsetDateTime,
}
impl PoolDecommissionInfo {
@@ -2307,7 +2185,6 @@ impl PoolDecommissionInfo {
fn mark_progress_saved(&mut self) {
self.progress_save_item_baseline = self.counted_items();
self.progress_save_retry_after = None;
}
pub fn bucket_push(&mut self, bucket: &DecomBucketInfo) {
@@ -2612,40 +2489,6 @@ impl ECStore {
snapshot.save(self.pools.clone()).await
}
async fn save_decommission_progress_checkpoint(&self, idx: usize) -> Result<bool> {
// Lock order: save gate, then the short pool metadata read/write sections. Peer
// reloads are intentionally performed by the caller after both locks are released.
let _save_guard = self.pool_meta_save_gate.lock().await;
let (snapshot, checkpoint) = {
let pool_meta = self.pool_meta.read().await;
let Some(checkpoint) = pool_meta.decommission_progress_checkpoint(
idx,
DECOMMISSION_PROGRESS_SAVE_INTERVAL,
OffsetDateTime::now_utc(),
)?
else {
return Ok(false);
};
let mut snapshot = pool_meta.clone();
let Some(pool) = snapshot.pools.get_mut(idx) else {
return Err(invalid_decommission_pool_index_error(snapshot.pools.len(), idx));
};
pool.last_update = checkpoint.checkpoint_at;
(snapshot, checkpoint)
};
if let Err(err) = snapshot.save(self.pools.clone()).await {
let retry_after = OffsetDateTime::now_utc() + DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF;
let mut pool_meta = self.pool_meta.write().await;
pool_meta.defer_decommission_progress_checkpoint(idx, checkpoint, retry_after);
return Err(err);
}
let mut pool_meta = self.pool_meta.write().await;
Ok(pool_meta.commit_decommission_progress_checkpoint(idx, checkpoint))
}
async fn save_current_pool_meta_for_decommission_start(
&self,
indices: &[usize],
@@ -3028,7 +2871,7 @@ impl ECStore {
Ok(())
}
async fn track_decommission_entry_progress_stage(
async fn save_decommission_entry_progress_stage(
&self,
idx: usize,
bucket: &str,
@@ -3039,6 +2882,22 @@ impl ECStore {
let mut pool_meta = self.pool_meta.write().await;
track_decommission_current_object_stage(&mut pool_meta, idx, bucket, object, stage)
.map_err(|err| with_decommission_entry_context(stage, bucket, object, err))?;
touch_decommission_progress(&mut pool_meta, idx)
.map_err(|err| with_decommission_entry_context(stage, bucket, object, err))?;
}
if let Some(err) = resolve_decommission_progress_save_result(self.save_current_pool_meta().await) {
warn!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket,
object = %object,
stage,
error = ?err,
"Decommission progress stage save failed"
);
}
Ok(())
@@ -3306,7 +3165,7 @@ impl ECStore {
let bucket_name = bucket.clone();
let object_name = rd.object_info.name.clone();
self.track_decommission_entry_progress_stage(
self.save_decommission_entry_progress_stage(
idx,
bucket_name.as_str(),
object_name.as_str(),
@@ -3400,7 +3259,7 @@ impl ECStore {
}
decommission_cancel_signal_result(rx.is_cancelled())?;
self.track_decommission_entry_progress_stage(
self.save_decommission_entry_progress_stage(
idx,
bucket.as_str(),
entry.name.as_str(),
@@ -3408,7 +3267,7 @@ impl ECStore {
)
.await?;
self.track_decommission_entry_progress_stage(
self.save_decommission_entry_progress_stage(
idx,
bucket.as_str(),
entry.name.as_str(),
@@ -3475,42 +3334,34 @@ impl ECStore {
}
};
self.track_decommission_entry_progress_stage(
idx,
bucket.as_str(),
entry.name.as_str(),
DECOMMISSION_STAGE_ENTRY_FINISHED,
)
.await?;
self.save_decommission_entry_progress_stage(idx, bucket.as_str(), entry.name.as_str(), DECOMMISSION_STAGE_ENTRY_FINISHED)
.await?;
if should_save_progress {
match self.save_decommission_progress_checkpoint(idx).await {
Ok(true) => {
if let Some(notification_sys) = runtime_sources::notification_sys()
&& let Err(err) = resolve_decommission_entry_reload_result(
notification_sys.reload_pool_meta().await,
bucket.as_str(),
entry.name.as_str(),
)
{
warn!("{err}");
}
}
Ok(false) => {}
Err(err) => {
if let Some(err) = resolve_decommission_progress_save_result(Err(err)) {
warn!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket,
object = %entry.name,
state = "progress_save_failed",
error = %err,
"Decommission progress save failed; continuing and will retry at the next checkpoint"
);
}
let save_result = self.save_current_pool_meta().await;
if let Some(err) = resolve_decommission_progress_save_result(save_result) {
warn!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket,
object = %entry.name,
state = "progress_save_failed",
error = %err,
"Decommission progress save failed; continuing and will retry at the next checkpoint"
);
} else {
let mut pool_meta = self.pool_meta.write().await;
pool_meta.mark_decommission_progress_saved();
if let Some(notification_sys) = runtime_sources::notification_sys()
&& let Err(err) = resolve_decommission_entry_reload_result(
notification_sys.reload_pool_meta().await,
bucket.as_str(),
entry.name.as_str(),
)
{
warn!("{err}");
}
}
}
@@ -3687,7 +3538,6 @@ impl ECStore {
let rx_clone = rx.clone();
let bi = bi.clone();
let set_id = set_idx;
let listing_entry_error = entry_error.clone();
let worker = tokio::spawn(async move {
let _listing_permit = listing_permit;
run_decommission_listing_with_retry(
@@ -3701,11 +3551,7 @@ impl ECStore {
let set = set.clone();
let rx = rx_clone.clone();
let bucket = bi.clone();
let entry_error = listing_entry_error.clone();
async move {
set.list_objects_to_decommission(rx, bucket, callback, entry_error.clone(), idx, set_id)
.await
}
async move { set.list_objects_to_decommission(rx, bucket, callback).await }
},
)
.await
@@ -3735,7 +3581,11 @@ impl ECStore {
wait_decommission_worker_drain(&workers, worker_limit).await?;
if let Some(err) = resolve_decommission_listing_error(listing_worker_error, entry_error.lock().await.clone()) {
if let Some(err) = listing_worker_error {
return Err(err);
}
if let Some(err) = entry_error.lock().await.clone() {
return Err(err);
}
@@ -4341,7 +4191,7 @@ impl ECStore {
let buckets = self.get_buckets_to_decommission().await?;
let pool = self.pools[idx].clone();
for (set_index, set) in pool.disk_set.iter().enumerate() {
for set in &pool.disk_set {
for bucket_info in &buckets {
let mut lifecycle_config = None;
let mut object_lock_config = None;
@@ -4436,7 +4286,7 @@ impl ECStore {
});
let list_result = set
.list_objects_to_decommission(callback_rx, bucket_info.clone(), callback, entry_error.clone(), idx, set_index)
.list_objects_to_decommission(callback_rx, bucket_info.clone(), callback)
.await;
let entry_error = entry_error.lock().await.clone();
resolve_decommission_check_after_list_result(list_result, entry_error)?;
@@ -5171,15 +5021,12 @@ mod tests {
pub type ListCallback = Arc<dyn Fn(MetaCacheEntry) -> BoxFuture<'static, ()> + Send + Sync + 'static>;
impl SetDisks {
#[tracing::instrument(skip(self, rx, cb_func, entry_error))]
#[tracing::instrument(skip(self, rx, cb_func))]
async fn list_objects_to_decommission(
self: &Arc<Self>,
rx: CancellationToken,
bucket_info: DecomBucketInfo,
cb_func: ListCallback,
entry_error: Arc<tokio::sync::Mutex<Option<Error>>>,
pool_index: usize,
set_index: usize,
) -> Result<()> {
let (disks, _) = self.get_online_disks_with_healing(false).await;
ensure_decommission_listing_disks_available(!disks.is_empty(), &bucket_info.name)?;
@@ -5194,12 +5041,6 @@ impl SetDisks {
};
let cb1 = cb_func.clone();
let unresolved_error = entry_error.clone();
let unresolved_rx = rx.clone();
let unresolved_bucket = bucket_info.name.clone();
let unresolved_prefix = bucket_info.prefix.clone();
let unresolved_pool_index = pool_index;
let unresolved_set_index = set_index;
list_path_raw(
rx,
@@ -5212,51 +5053,20 @@ impl SetDisks {
skip_walkdir_total_timeout: true,
walkdir_stall_timeout: Some(DECOMMISSION_BACKGROUND_WALKDIR_STALL_TIMEOUT),
agreed: Some(Box::new(move |entry: MetaCacheEntry| Box::pin(cb1(entry)))),
partial: Some(Box::new(move |entries: MetaCacheEntries, errs: &[Option<DiskError>]| {
partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option<DiskError>]| {
let resolver = resolver.clone();
let cb_func = cb_func.clone();
let bucket = unresolved_bucket.clone();
let prefix = unresolved_prefix.clone();
let unresolved_error = unresolved_error.clone();
let unresolved_rx = unresolved_rx.clone();
let pool_index = unresolved_pool_index;
let set_index = unresolved_set_index;
let disk_error_count = errs.iter().flatten().count();
if unresolved_rx.is_cancelled() {
return Box::pin(async {});
}
match resolve_decommission_partial_listing_entry(
entries,
resolver,
&bucket,
&prefix,
disk_error_count,
pool_index,
set_index,
) {
Ok(entry) => {
match entries.resolve(resolver) {
Some(entry) => {
warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name);
Box::pin(async move {
cb_func(entry).await;
})
}
Err(err) => Box::pin(async move {
if unresolved_rx.is_cancelled() {
return;
}
warn!(
event = EVENT_DECOMMISSION_BUCKET,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
bucket = %bucket,
prefix = %prefix,
state = "unresolved_entry",
error = %err,
"Decommission listing failed closed on unresolved metadata"
);
record_decommission_entry_error(&unresolved_error, &unresolved_rx, err).await;
}),
None => {
warn!("decommission_pool: list_objects_to_decommission get none");
Box::pin(async {})
}
}
})),
..Default::default()
@@ -5264,10 +5074,6 @@ impl SetDisks {
)
.await?;
if let Some(err) = entry_error.lock().await.clone() {
return Err(err);
}
Ok(())
}
}
@@ -5458,11 +5264,11 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi
#[cfg(test)]
mod pools_tests {
use super::{
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF,
DecomBucketInfo, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta,
PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info, bind_decommission_cancelers,
bind_missing_decommission_cancelers, cancel_decommission_canceler, classify_decommission_terminal_state,
count_decommission_item, decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options,
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DecomBucketInfo,
DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo,
PoolStatus, apply_decommission_status_space_info, bind_decommission_cancelers, bind_missing_decommission_cancelers,
cancel_decommission_canceler, classify_decommission_terminal_state, count_decommission_item,
decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options,
decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency,
ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_listing_disks_available,
ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, ensure_decommission_start_keeps_active_pool,
@@ -5473,12 +5279,11 @@ mod pools_tests {
has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested,
load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done,
merge_pool_status_refresh, missing_decommission_worker_prefix, observe_decommission_terminal_reload_result,
pool_meta_has_active_decommission, record_decommission_entry_error, require_decommission_store,
resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state,
resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result,
resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, resolve_decommission_listing_error,
resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result,
resolve_decommission_partial_listing_entry, resolve_decommission_pool_meta_reload_result,
pool_meta_has_active_decommission, require_decommission_store, resolve_decommission_bucket_done_save_result,
resolve_decommission_bucket_state, resolve_decommission_check_after_list_result,
resolve_decommission_entry_cleanup_delete_result, resolve_decommission_entry_exact_versions,
resolve_decommission_entry_reload_result, resolve_decommission_listing_worker_result,
resolve_decommission_optional_bucket_config_result, resolve_decommission_pool_meta_reload_result,
resolve_decommission_preflight_heal_result, resolve_decommission_progress_save_result,
resolve_decommission_spawn_failure_result, resolve_decommission_terminal_mark_after_error_result,
resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result,
@@ -5488,17 +5293,16 @@ mod pools_tests {
should_preserve_decommission_canceled_state, should_reject_decommission_cancel_as_terminal,
should_retry_decommission_cancel_reload, should_retry_decommission_listing, should_skip_canceled_decommission_routine,
split_decommission_buckets, take_and_cancel_decommission_canceler, take_decommission_canceler,
track_decommission_current_object, track_decommission_current_object_stage, validate_start_decommission_request,
wait_decommission_listing_retry, wait_decommission_worker_drain, with_decommission_entry_context,
touch_decommission_progress, track_decommission_current_object, track_decommission_current_object_stage,
validate_start_decommission_request, wait_decommission_listing_retry, wait_decommission_worker_drain,
with_decommission_entry_context,
};
use crate::data_movement;
use crate::disk::endpoint::Endpoint;
use crate::error::{Error, StorageError};
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats};
use rustfs_filemeta::{
FileInfo, FileInfoVersions, MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ObjectPartInfo,
};
use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo};
use rustfs_rio::Index;
use std::sync::{
Arc,
@@ -6517,65 +6321,6 @@ mod pools_tests {
assert!(matches!(err, Error::SlowDown));
}
#[test]
fn test_resolve_decommission_partial_listing_entry_rejects_unresolved_metadata() {
let err = resolve_decommission_partial_listing_entry(
MetaCacheEntries(vec![None]),
MetadataResolutionParams {
dir_quorum: 2,
obj_quorum: 2,
bucket: "bucket-a".to_string(),
..Default::default()
},
"bucket-a",
"prefix/",
1,
2,
3,
)
.expect_err("unresolved partial listing must fail closed");
let message = err.to_string();
assert!(message.contains("decommission listing could not resolve metadata"));
assert!(message.contains("bucket-a/prefix/"));
assert!(message.contains("pool 2 set 3"));
assert!(message.contains("1 disk error(s)"));
}
#[tokio::test]
async fn test_record_decommission_entry_error_cancels_listing_and_preserves_first_error() {
let entry_error = Arc::new(tokio::sync::Mutex::new(None));
let rx = CancellationToken::new();
record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await;
record_decommission_entry_error(&entry_error, &rx, Error::OperationCanceled).await;
assert!(rx.is_cancelled());
assert!(matches!(*entry_error.lock().await, Some(Error::SlowDown)));
}
#[tokio::test]
async fn test_record_decommission_entry_error_ignores_already_canceled_listing() {
let entry_error = Arc::new(tokio::sync::Mutex::new(None));
let rx = CancellationToken::new();
rx.cancel();
record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await;
assert!(entry_error.lock().await.is_none());
}
#[test]
fn test_resolve_decommission_listing_error_preserves_real_listing_failure() {
let err = resolve_decommission_listing_error(Some(Error::SlowDown), Some(Error::OperationCanceled))
.expect("listing failure should be returned");
assert!(matches!(err, Error::SlowDown));
let err = resolve_decommission_listing_error(Some(Error::OperationCanceled), Some(Error::SlowDown))
.expect("entry failure should be returned");
assert!(matches!(err, Error::SlowDown));
}
#[test]
fn test_resolve_decommission_check_after_list_result_returns_list_result_without_entry_error() {
let err = resolve_decommission_check_after_list_result(Err(Error::OperationCanceled), None)
@@ -6793,7 +6538,7 @@ mod pools_tests {
}
#[test]
fn test_track_decommission_stage_does_not_advance_checkpoint_state() {
fn test_touch_decommission_progress_updates_last_update_and_save_baseline() {
let mut meta = PoolMeta {
pools: vec![PoolStatus {
id: 0,
@@ -6808,13 +6553,11 @@ mod pools_tests {
..Default::default()
};
track_decommission_current_object_stage(&mut meta, 0, "bucket", "object", "migrate_object")
.expect("valid decommission progress should be tracked");
touch_decommission_progress(&mut meta, 0).expect("valid decommission progress should be touched");
assert_eq!(meta.pools[0].last_update, OffsetDateTime::UNIX_EPOCH);
assert!(meta.pools[0].last_update > OffsetDateTime::UNIX_EPOCH);
let info = meta.pools[0].decommission.as_ref().expect("decommission info should exist");
assert_eq!(info.items_since_last_progress_save(), 5);
assert_eq!(info.stage, "migrate_object");
assert_eq!(info.items_since_last_progress_save(), 0);
}
#[test]
@@ -6889,134 +6632,6 @@ mod pools_tests {
assert_eq!(info.items_since_last_progress_save(), 1);
}
#[test]
fn test_pool_meta_update_after_does_not_advance_last_update_before_save() {
let last_update = OffsetDateTime::UNIX_EPOCH;
let mut meta = PoolMeta {
pools: vec![PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update,
decommission: Some(PoolDecommissionInfo {
start_time: Some(last_update),
items_decommissioned: DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD,
..Default::default()
}),
}],
..Default::default()
};
assert!(
meta.update_after(0, DECOMMISSION_PROGRESS_SAVE_INTERVAL)
.expect("item threshold should request a checkpoint")
);
assert_eq!(meta.pools[0].last_update, last_update);
}
#[test]
fn test_decommission_progress_checkpoint_commits_exact_snapshot_watermark() {
let start_time = OffsetDateTime::UNIX_EPOCH;
let checkpoint_at = start_time + Duration::seconds(30);
let mut meta = PoolMeta {
pools: vec![PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: start_time,
decommission: Some(PoolDecommissionInfo {
start_time: Some(start_time),
items_decommissioned: DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD,
..Default::default()
}),
}],
..Default::default()
};
let checkpoint = meta
.decommission_progress_checkpoint(0, DECOMMISSION_PROGRESS_SAVE_INTERVAL, checkpoint_at)
.expect("valid decommission state should produce a checkpoint")
.expect("item threshold should produce a checkpoint");
meta.count_item(0, 1, false);
assert!(meta.commit_decommission_progress_checkpoint(0, checkpoint));
let info = meta.pools[0].decommission.as_ref().expect("decommission info should exist");
assert_eq!(info.progress_save_item_baseline, checkpoint.counted_items);
assert_eq!(info.items_since_last_progress_save(), 1);
assert_eq!(meta.pools[0].last_update, checkpoint_at);
}
#[test]
fn test_decommission_progress_checkpoint_backoff_does_not_advance_baseline() {
let start_time = OffsetDateTime::UNIX_EPOCH;
let checkpoint_at = start_time + Duration::seconds(30);
let retry_after = checkpoint_at + DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF;
let mut meta = PoolMeta {
pools: vec![PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: start_time,
decommission: Some(PoolDecommissionInfo {
start_time: Some(start_time),
items_decommissioned: DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD,
..Default::default()
}),
}],
..Default::default()
};
let checkpoint = meta
.decommission_progress_checkpoint(0, DECOMMISSION_PROGRESS_SAVE_INTERVAL, checkpoint_at)
.expect("valid decommission state should produce a checkpoint")
.expect("item threshold should produce a checkpoint");
meta.defer_decommission_progress_checkpoint(0, checkpoint, retry_after);
assert!(
meta.decommission_progress_checkpoint(0, DECOMMISSION_PROGRESS_SAVE_INTERVAL, checkpoint_at)
.expect("retry backoff check should succeed")
.is_none()
);
assert_eq!(meta.pools[0].last_update, start_time);
assert_eq!(
meta.pools[0]
.decommission
.as_ref()
.expect("decommission info should exist")
.progress_save_item_baseline,
0
);
}
#[test]
fn test_decommission_progress_checkpoint_count_scales_with_threshold() {
let start_time = OffsetDateTime::UNIX_EPOCH;
let checkpoint_at = start_time;
let mut meta = PoolMeta {
pools: vec![PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: start_time,
decommission: Some(PoolDecommissionInfo {
start_time: Some(start_time),
..Default::default()
}),
}],
..Default::default()
};
let mut checkpoint_count = 0;
for _ in 0..(DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD * 10) {
meta.count_item(0, 1, false);
if let Some(checkpoint) = meta
.decommission_progress_checkpoint(0, DECOMMISSION_PROGRESS_SAVE_INTERVAL, checkpoint_at)
.expect("valid decommission state should produce a checkpoint")
{
checkpoint_count += 1;
assert!(meta.commit_decommission_progress_checkpoint(0, checkpoint));
}
}
assert_eq!(checkpoint_count, 10);
}
#[test]
fn test_ensure_decommission_not_rebalancing_rejects_running_rebalance() {
let err = ensure_decommission_not_rebalancing(true).expect_err("rebalance running should be rejected");
-54
View File
@@ -1640,60 +1640,6 @@ mod tests {
assert_eq!(payload["items"].as_array().expect("items should be an array").len(), 0);
}
#[tokio::test]
async fn test_process_query_request_reports_displaced_terminal_detail() {
let heal_manager = Arc::new(HealManager::new(
Arc::new(MockStorage),
Some(HealConfig {
queue_size: 1,
..HealConfig::default()
}),
));
let mut displaced = HealRequest::new(
HealType::Bucket {
bucket: "displaced-channel".to_string(),
},
HealOptions::default(),
HealPriority::Low,
);
displaced.id = "displaced-channel-task".to_string();
let displaced_id = displaced.id.clone();
heal_manager
.submit_heal_request(displaced)
.await
.expect("initial channel task should queue");
heal_manager
.submit_heal_request(HealRequest::new(
HealType::Bucket {
bucket: "successor-channel".to_string(),
},
HealOptions::default(),
HealPriority::High,
))
.await
.expect("successor channel task should displace the initial task");
let processor = HealChannelProcessor::new(heal_manager);
let (tx, rx) = oneshot::channel();
processor
.process_query_request("displaced-channel".to_string(), displaced_id, None, tx)
.await
.expect("displaced query should process");
let response = rx
.await
.expect("query response should be returned")
.expect("displaced query should remain successful");
let payload: serde_json::Value = serde_json::from_slice(response.data.as_deref().expect("status payload should exist"))
.expect("status payload should be json");
assert_eq!(payload["summary"], "stopped");
assert!(
response
.error
.as_deref()
.is_some_and(|detail| detail.contains("reason=displaced"))
);
}
#[tokio::test]
async fn test_process_query_request_reports_running_for_queued_task() {
let heal_manager = create_test_heal_manager();
+246 -25
View File
@@ -13,7 +13,7 @@
// limitations under the License.
use crate::heal::{
progress::HealProgress,
progress::{HealProgress, add_bytes, increment_counter},
resume::{
CheckpointManager, ReplacementTargetIdentity, ResumeManager, ResumeUtils, compose_key,
replacement_target_identities_match,
@@ -410,6 +410,9 @@ 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
@@ -474,6 +477,23 @@ 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;
@@ -483,6 +503,58 @@ 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
@@ -516,13 +588,42 @@ impl ErasureSetHealer {
return bucket_result;
}
// update checkpoint position
checkpoint_manager.update_position(bucket_idx, current_object_index).await?;
// update progress
resume_manager
.update_progress(processed_objects, successful_objects, failed_objects, skipped_objects)
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?;
checkpoint_manager.update_position(bucket_idx, current_object_index).await?;
resume_manager
.update_progress_with_bytes(
processed_objects,
successful_objects,
failed_objects,
skipped_objects,
bytes_processed,
)
.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() {
@@ -781,14 +882,36 @@ impl ErasureSetHealer {
if should_skip_new_version(item.mod_time_unix_nanos, started_at_secs) {
checkpoint_manager.add_processed_object(key).await?;
*processed_objects = processed_objects.saturating_add(1);
let counter_ok = increment_counter(processed_objects);
completed_in_page = completed_in_page.saturating_add(1);
counter!("rustfs_heal_skipped_new_versions_total").increment(1);
{
let (skipped_new, skipped_ilm, 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_progress(*processed_objects, *successful_objects, *failed_objects, bytes_processed);
progress.update_object_progress(
*processed_objects,
*successful_objects,
*failed_objects,
*skipped_objects,
bytes_processed,
);
if !counter_ok {
progress.mark_unknown();
}
(progress.skipped_new_versions, progress.skipped_ilm_expired, progress.counter_unknown)
};
if !counter_ok || counter_unknown {
checkpoint_manager.mark_counter_unknown().await?;
}
checkpoint_manager
.set_skipped_version_counts(skipped_new, skipped_ilm)
.await?;
checkpoint_manager
.update_progress(*successful_objects, *failed_objects, *skipped_objects, bytes_processed)
.await?;
if !counter_ok || counter_unknown {
resume_manager.mark_counter_unknown().await?;
}
debug!(
target: "rustfs::heal::erasure_healer",
@@ -821,14 +944,36 @@ impl ErasureSetHealer {
.await?
{
checkpoint_manager.add_processed_object(key).await?;
*processed_objects = processed_objects.saturating_add(1);
let counter_ok = increment_counter(processed_objects);
completed_in_page = completed_in_page.saturating_add(1);
counter!("rustfs_heal_skipped_ilm_expired_total").increment(1);
{
let (skipped_new, skipped_ilm, 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_progress(*processed_objects, *successful_objects, *failed_objects, bytes_processed);
progress.update_object_progress(
*processed_objects,
*successful_objects,
*failed_objects,
*skipped_objects,
bytes_processed,
);
if !counter_ok {
progress.mark_unknown();
}
(progress.skipped_new_versions, progress.skipped_ilm_expired, progress.counter_unknown)
};
if !counter_ok || counter_unknown {
checkpoint_manager.mark_counter_unknown().await?;
}
checkpoint_manager
.set_skipped_version_counts(skipped_new, skipped_ilm)
.await?;
checkpoint_manager
.update_progress(*successful_objects, *failed_objects, *skipped_objects, bytes_processed)
.await?;
if !counter_ok || counter_unknown {
resume_manager.mark_counter_unknown().await?;
}
debug!(
target: "rustfs::heal::erasure_healer",
@@ -954,10 +1099,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;
match result {
Ok(true) => {
*successful_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
telemetry_unknown |= !increment_counter(successful_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size);
checkpoint_manager.add_processed_object(key).await?;
debug!(
target: "rustfs::heal::erasure_healer",
@@ -974,8 +1120,8 @@ impl ErasureSetHealer {
}
Ok(false) => {
checkpoint_manager.add_processed_object(key).await?;
*successful_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
telemetry_unknown |= !increment_counter(successful_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size);
debug!(
target: "rustfs::heal::erasure_healer",
event = EVENT_HEAL_ERASURE_OBJECT_STATE,
@@ -991,8 +1137,8 @@ impl ErasureSetHealer {
}
Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err),
Err(Error::TransientSkip { message }) => {
*skipped_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
telemetry_unknown |= !increment_counter(skipped_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, 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,
@@ -1008,8 +1154,8 @@ impl ErasureSetHealer {
});
}
Err(err) => {
*failed_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
telemetry_unknown |= !increment_counter(failed_objects);
telemetry_unknown |= !add_bytes(&mut bytes_processed, 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,
@@ -1026,12 +1172,31 @@ impl ErasureSetHealer {
}
}
*processed_objects += 1;
telemetry_unknown |= !increment_counter(processed_objects);
completed_in_page += 1;
{
let progress_unknown = {
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("{bucket}/{object}")));
progress.update_progress(*processed_objects, *successful_objects, *failed_objects, bytes_processed);
progress.update_object_progress(
*processed_objects,
*successful_objects,
*failed_objects,
*skipped_objects,
bytes_processed,
);
if telemetry_unknown {
progress.mark_unknown();
}
progress.counter_unknown
};
if telemetry_unknown || progress_unknown {
checkpoint_manager.mark_counter_unknown().await?;
}
checkpoint_manager
.update_progress(*successful_objects, *failed_objects, *skipped_objects, bytes_processed)
.await?;
if telemetry_unknown || progress_unknown {
resume_manager.mark_counter_unknown().await?;
}
if completed_in_page.is_multiple_of(100) {
@@ -1083,10 +1248,66 @@ impl ErasureSetHealer {
/// initialize progress tracking
async fn initialize_progress(&self, _buckets: &[String], state: &crate::heal::resume::ResumeState) {
let mut progress = self.progress.write().await;
progress.objects_scanned = state.total_objects;
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_healed = state.successful_objects;
progress.objects_failed = state.failed_objects;
progress.bytes_processed = 0; // Resume state tracks object counts, not byte counters.
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.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());
+76 -114
View File
@@ -40,7 +40,6 @@ use tracing::{debug, error, info, warn};
use super::{DiskError, Endpoint, HealDiskExt as _, local_disk_map_read};
const KEEP_HEAL_TASK_STATUS_DURATION: Duration = Duration::from_secs(10 * 60);
const DISPLACED_HEAL_REASON: &str = "reason=displaced; retry_hint=submit_again";
const LOG_COMPONENT_HEAL: &str = "heal";
const LOG_SUBSYSTEM_DISK_SCANNER: &str = "disk_scanner";
const LOG_SUBSYSTEM_MANAGER: &str = "manager";
@@ -121,30 +120,26 @@ struct MrfRepairNoticeTarget {
version_id: Option<[u8; 16]>,
}
#[derive(Debug, Clone)]
#[derive(Debug, Clone, PartialEq, Eq)]
struct HealAdmissionDecision {
result: HealAdmissionResult,
displaced_request: Option<HealRequest>,
displaced_task_id: Option<String>,
}
impl HealAdmissionDecision {
const fn new(result: HealAdmissionResult) -> Self {
Self {
result,
displaced_request: None,
displaced_task_id: None,
}
}
fn accepted_with_displacement(displaced_request: HealRequest) -> Self {
fn accepted_with_displacement(displaced_task_id: String) -> Self {
Self {
result: HealAdmissionResult::Accepted,
displaced_request: Some(displaced_request),
displaced_task_id: Some(displaced_task_id),
}
}
fn displaced_task_id(&self) -> Option<&str> {
self.displaced_request.as_ref().map(|request| request.id.as_str())
}
}
fn lock_mrf_repair_notice_targets(
@@ -156,55 +151,6 @@ fn lock_mrf_repair_notice_targets(
}
}
fn lock_displaced_terminals(
registry: &StdMutex<HashMap<String, Arc<CompletedHealStatus>>>,
) -> StdMutexGuard<'_, HashMap<String, Arc<CompletedHealStatus>>> {
match registry.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
}
}
fn record_displaced_terminal(
registry: &StdMutex<HashMap<String, Arc<CompletedHealStatus>>>,
request: &HealRequest,
) -> Arc<CompletedHealStatus> {
let terminal = Arc::new(CompletedHealStatus {
heal_type: request.heal_type.clone(),
status: HealTaskStatus::Failed {
error: format!("heal task displaced by a higher-priority request ({DISPLACED_HEAL_REASON})"),
},
result_items_truncated: false,
completed_at: SystemTime::now(),
seqed_items: Vec::new(),
next_seq: 0,
min_seq: 0,
});
let mut terminals = lock_displaced_terminals(registry);
prune_completed_heal_statuses(&mut terminals);
terminals.insert(request.id.clone(), Arc::clone(&terminal));
terminal
}
async fn remove_displaced_task_aliases(
aliases: &Arc<Mutex<HashMap<String, HealTaskAlias>>>,
terminals: &StdMutex<HashMap<String, Arc<CompletedHealStatus>>>,
task_id: &str,
terminal: &Arc<CompletedHealStatus>,
) {
let mut aliases = aliases.lock().await;
let alias_ids = aliases
.iter()
.filter_map(|(alias_id, alias)| (alias.task_id == task_id).then_some(alias_id.clone()))
.collect::<Vec<_>>();
let mut displaced_terminals = lock_displaced_terminals(terminals);
prune_completed_heal_statuses(&mut displaced_terminals);
for alias_id in alias_ids {
displaced_terminals.insert(alias_id, Arc::clone(terminal));
}
aliases.retain(|alias_id, alias| alias_id != task_id && alias.task_id != task_id);
}
async fn remove_task_aliases_for_task(registry: &Arc<Mutex<HashMap<String, HealTaskAlias>>>, task_id: &str) {
registry
.lock()
@@ -672,14 +618,6 @@ pub struct HealManager {
/// are shared so the lookup helper can hand a completed entry to a
/// caller without cloning the retained result window.
completed_heals: Arc<Mutex<HashMap<String, Arc<CompletedHealStatus>>>>,
/// Terminals for requests removed by priority displacement. An Accepted
/// task ID remains queryable for the same process lifetime and the normal
/// ten-minute status TTL; clients should treat `reason=displaced` as a
/// terminal result and submit a fresh request. This sidecar is synchronous
/// so admission can publish the terminal while the queue transition is
/// still under its lock, without awaiting another tokio lock. Queue state
/// is process-local, so this guarantee does not extend across restart.
displaced_terminals: Arc<StdMutex<HashMap<String, Arc<CompletedHealStatus>>>>,
/// Client tokens merged into an existing task id.
task_aliases: Arc<Mutex<HashMap<String, HealTaskAlias>>>,
/// Heal tasks waiting for a retry backoff to expire.
@@ -721,7 +659,6 @@ struct HealQueueContext<'a> {
heal_queue: &'a Arc<Mutex<PriorityHealQueue>>,
active_heals: &'a Arc<Mutex<HashMap<String, Arc<HealTask>>>>,
completed_heals: &'a Arc<Mutex<HashMap<String, Arc<CompletedHealStatus>>>>,
displaced_terminals: &'a Arc<StdMutex<HashMap<String, Arc<CompletedHealStatus>>>>,
task_aliases: &'a Arc<Mutex<HashMap<String, HealTaskAlias>>>,
retrying_heals: &'a Arc<Mutex<HashMap<String, RetryingHeal>>>,
mrf_repair_notice_targets: &'a Arc<StdMutex<HashMap<String, Vec<MrfRepairNoticeTarget>>>>,
@@ -937,7 +874,7 @@ impl HealManager {
result = "accepted_by_displacement",
"Heal queue request accepted by displacement"
});
return HealAdmissionDecision::accepted_with_displacement(displaced);
return HealAdmissionDecision::accepted_with_displacement(displaced.id);
}
demote_to_debug_when!(per_object_request, warn, target: "rustfs::heal::manager", {
@@ -1168,7 +1105,6 @@ impl HealManager {
active_heals: Arc::new(Mutex::new(HashMap::new())),
heal_queue: Arc::new(Mutex::new(PriorityHealQueue::new())),
completed_heals: Arc::new(Mutex::new(HashMap::new())),
displaced_terminals: Arc::new(StdMutex::new(HashMap::new())),
task_aliases: Arc::new(Mutex::new(HashMap::new())),
retrying_heals: Arc::new(Mutex::new(HashMap::new())),
mrf_repair_notice_targets: Arc::new(StdMutex::new(HashMap::new())),
@@ -1273,10 +1209,6 @@ impl HealManager {
active_heals.clear();
publish_active_heal_count(&active_heals);
self.completed_heals.lock().await.clear();
// Do not let the synchronous guard live across the following async lock.
{
lock_displaced_terminals(&self.displaced_terminals).clear();
}
self.task_aliases.lock().await.clear();
self.retrying_heals.lock().await.clear();
lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets).clear();
@@ -1527,11 +1459,7 @@ impl HealManager {
task_id = queued_id.to_owned();
}
let should_notify = matches!(admission, HealAdmissionResult::Accepted) && config.event_driven_scheduler_enable;
let displaced_task_id = admission_decision.displaced_task_id().map(ToOwned::to_owned);
let displaced_terminal = admission_decision
.displaced_request
.as_ref()
.map(|request| record_displaced_terminal(&self.displaced_terminals, request));
let displaced_task_id = admission_decision.displaced_task_id;
if matches!(admission, HealAdmissionResult::Accepted | HealAdmissionResult::Merged)
&& let Some(target) = mrf_notice_target
{
@@ -1545,12 +1473,8 @@ impl HealManager {
drop(queue);
drop(active_heals);
if let (Some(displaced_task_id), Some(displaced_terminal)) = (displaced_task_id, displaced_terminal) {
// The queue has already removed the displaced request, so the
// synchronous terminal sidecar was published before aliases and
// MRF ownership are cleaned up.
remove_displaced_task_aliases(&self.task_aliases, &self.displaced_terminals, &displaced_task_id, &displaced_terminal)
.await;
if let Some(displaced_task_id) = displaced_task_id {
self.remove_aliases_for_task(&displaced_task_id).await;
}
if should_notify {
@@ -1625,15 +1549,6 @@ impl HealManager {
}
}
if terminal_completed.is_none() {
let mut displaced_terminals = lock_displaced_terminals(&self.displaced_terminals);
prune_completed_heal_statuses(&mut displaced_terminals);
terminal_completed = displaced_terminals
.get(canonical_task_id)
.filter(|terminal| matches_path(&terminal.heal_type))
.cloned();
}
match terminal_completed {
Some(completed) => TaskStateLookup::Completed(completed),
None => TaskStateLookup::NotFound,
@@ -1754,19 +1669,9 @@ impl HealManager {
let mut completed_heals = self.completed_heals.lock().await;
prune_completed_heal_statuses(&mut completed_heals);
if completed_heals
completed_heals
.values()
.any(|completed| heal_type_matches_path(&completed.heal_type, heal_path))
{
return true;
}
drop(completed_heals);
let mut displaced_terminals = lock_displaced_terminals(&self.displaced_terminals);
prune_completed_heal_statuses(&mut displaced_terminals);
displaced_terminals
.values()
.any(|terminal| heal_type_matches_path(&terminal.heal_type, heal_path))
}
/// Get task progress
@@ -2012,16 +1917,44 @@ impl HealManager {
}
let mut snapshot = HealProgress::default();
let mut has_object_sweep = false;
let mut all_object_baselines_known = true;
let mut counter_overflow = false;
let mut stage_current = 0_u64;
let mut stage_total = 0_u64;
for task in active_tasks {
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);
let object_sweep = matches!(progress.kind, crate::heal::progress::HealProgressKind::ObjectSweep);
has_object_sweep |= object_sweep;
if object_sweep {
all_object_baselines_known &= progress.baseline_known;
}
counter_overflow |=
progress.counter_unknown || matches!(progress.progress_state, crate::heal::progress::HealProgressState::Unknown);
match stage_current.checked_add(progress.stage_current) {
Some(sum) => stage_current = sum,
None => counter_overflow = true,
}
match stage_total.checked_add(progress.stage_total) {
Some(sum) => stage_total = sum,
None => counter_overflow = true,
}
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),
] {
match target.checked_add(value) {
Some(sum) => *target = sum,
None => counter_overflow = true,
}
}
snapshot.start_time = match (snapshot.start_time, progress.start_time) {
(Some(current), Some(next)) => Some(current.min(next)),
(None, next) => next,
@@ -2036,7 +1969,36 @@ impl HealManager {
snapshot.current_object = progress.current_object;
}
}
snapshot.refresh_progress_percentage();
snapshot.kind = if has_object_sweep {
crate::heal::progress::HealProgressKind::ObjectSweep
} else {
crate::heal::progress::HealProgressKind::Stage
};
snapshot.stage_current = stage_current;
snapshot.stage_total = stage_total;
snapshot.baseline_known = has_object_sweep && all_object_baselines_known;
snapshot.progress_state = if counter_overflow {
crate::heal::progress::HealProgressState::Unknown
} else if has_object_sweep && !all_object_baselines_known {
crate::heal::progress::HealProgressState::Indeterminate
} else if has_object_sweep {
crate::heal::progress::HealProgressState::Running
} else if stage_total == 0 {
crate::heal::progress::HealProgressState::Indeterminate
} else {
crate::heal::progress::HealProgressState::Running
};
if counter_overflow {
snapshot.progress_percentage = 0.0;
} else if !has_object_sweep {
snapshot.progress_percentage = if stage_total == 0 {
0.0
} else {
((stage_current as f64 / stage_total as f64) * 100.0).min(99.999)
};
} else {
snapshot.refresh_progress_percentage();
}
snapshot.refresh_estimated_completion_time();
Some(snapshot)
}
+2 -15
View File
@@ -21,7 +21,6 @@ impl HealManager {
let heal_queue = self.heal_queue.clone();
let active_heals = self.active_heals.clone();
let task_aliases = self.task_aliases.clone();
let displaced_terminals = self.displaced_terminals.clone();
let mrf_repair_notice_targets = self.mrf_repair_notice_targets.clone();
let storage = self.storage.clone();
let replacement_recovery_anchors = self.replacement_recovery_anchors.clone();
@@ -482,10 +481,6 @@ impl HealManager {
let admission = admission_decision.result;
let should_notify =
matches!(admission, HealAdmissionResult::Accepted) && config.event_driven_scheduler_enable;
let displaced_terminal = admission_decision
.displaced_request
.as_ref()
.map(|request| record_displaced_terminal(&displaced_terminals, request));
if matches!(admission, HealAdmissionResult::Accepted)
&& let Some(anchor) = recovery_anchor
{
@@ -496,16 +491,8 @@ impl HealManager {
}
drop(queue);
drop(config);
if let (Some(displaced_task_id), Some(displaced_terminal)) =
(admission_decision.displaced_task_id().map(ToOwned::to_owned), displaced_terminal)
{
remove_displaced_task_aliases(
&task_aliases,
&displaced_terminals,
&displaced_task_id,
&displaced_terminal,
)
.await;
if let Some(displaced_task_id) = admission_decision.displaced_task_id {
remove_task_aliases_for_task(&task_aliases, &displaced_task_id).await;
lock_mrf_repair_notice_targets(&mrf_repair_notice_targets).remove(&displaced_task_id);
}
if matches!(admission, HealAdmissionResult::Accepted) {
+3 -25
View File
@@ -21,7 +21,6 @@ impl HealManager {
let heal_queue = self.heal_queue.clone();
let active_heals = self.active_heals.clone();
let completed_heals = self.completed_heals.clone();
let displaced_terminals = self.displaced_terminals.clone();
let task_aliases = self.task_aliases.clone();
let retrying_heals = self.retrying_heals.clone();
let mrf_repair_notice_targets = self.mrf_repair_notice_targets.clone();
@@ -54,7 +53,6 @@ impl HealManager {
heal_queue: &heal_queue,
active_heals: &active_heals,
completed_heals: &completed_heals,
displaced_terminals: &displaced_terminals,
task_aliases: &task_aliases,
retrying_heals: &retrying_heals,
mrf_repair_notice_targets: &mrf_repair_notice_targets,
@@ -73,7 +71,6 @@ impl HealManager {
heal_queue: &heal_queue,
active_heals: &active_heals,
completed_heals: &completed_heals,
displaced_terminals: &displaced_terminals,
task_aliases: &task_aliases,
retrying_heals: &retrying_heals,
mrf_repair_notice_targets: &mrf_repair_notice_targets,
@@ -101,7 +98,6 @@ impl HealManager {
heal_queue,
active_heals,
completed_heals,
displaced_terminals,
task_aliases,
retrying_heals,
mrf_repair_notice_targets,
@@ -187,7 +183,6 @@ impl HealManager {
let active_heals_clone = active_heals.clone();
let heal_queue_clone = heal_queue.clone();
let completed_heals_clone = completed_heals.clone();
let displaced_terminals_clone = displaced_terminals.clone();
let task_aliases_clone = task_aliases.clone();
let retrying_heals_clone = retrying_heals.clone();
let mrf_repair_notice_targets_clone = mrf_repair_notice_targets.clone();
@@ -368,7 +363,6 @@ impl HealManager {
let retry_heal_queue = heal_queue_clone.clone();
let retrying_heals_for_spawn = retrying_heals_clone.clone();
let retry_task_aliases = task_aliases_clone.clone();
let retry_displaced_terminals = displaced_terminals_clone.clone();
let retry_mrf_repair_notice_targets = mrf_repair_notice_targets_clone.clone();
let retry_completed_heals = completed_heals_clone.clone();
let retry_notify = notify_clone.clone();
@@ -436,14 +430,6 @@ impl HealManager {
let admission = admission_decision.result;
let should_notify = matches!(admission, HealAdmissionResult::Accepted)
&& retry_config.event_driven_scheduler_enable;
// Publish the terminal synchronously while the
// queue transition is protected. The subsequent
// queue -> retrying handoff retains the lock order
// used by operations_snapshot.
let displaced_terminal = admission_decision
.displaced_request
.as_ref()
.map(|request| record_displaced_terminal(&retry_displaced_terminals, request));
match admission {
HealAdmissionResult::Accepted => {
// Transfer ownership while holding queue -> retrying,
@@ -451,18 +437,10 @@ impl HealManager {
#[cfg(test)]
pause_retry_ownership_transition(&retry_request_id, true).await;
retrying_heals_for_spawn.lock().await.remove(&retry_request_id);
let displaced_task_id = admission_decision.displaced_task_id().map(ToOwned::to_owned);
let displaced_task_id = admission_decision.displaced_task_id;
drop(queue);
if let (Some(displaced_task_id), Some(displaced_terminal)) =
(displaced_task_id, displaced_terminal)
{
remove_displaced_task_aliases(
&retry_task_aliases,
&retry_displaced_terminals,
&displaced_task_id,
&displaced_terminal,
)
.await;
if let Some(displaced_task_id) = displaced_task_id {
remove_task_aliases_for_task(&retry_task_aliases, &displaced_task_id).await;
remove_mrf_repair_notice_targets(
&retry_mrf_repair_notice_targets,
&displaced_task_id,
+1 -262
View File
@@ -84,7 +84,6 @@ async fn process_manager_queue_once(manager: &HealManager) {
heal_queue: &manager.heal_queue,
active_heals: &manager.active_heals,
completed_heals: &manager.completed_heals,
displaced_terminals: &manager.displaced_terminals,
task_aliases: &manager.task_aliases,
retrying_heals: &manager.retrying_heals,
mrf_repair_notice_targets: &manager.mrf_repair_notice_targets,
@@ -2779,10 +2778,7 @@ async fn test_high_priority_request_displaces_lower_priority_when_queue_full() {
HealAdmissionResult::Accepted
);
assert_eq!(manager.get_queue_length().await, 1);
assert!(matches!(
manager.get_task_status(&low_id).await,
Ok(HealTaskStatus::Failed { error }) if error.contains("reason=displaced")
));
assert!(matches!(manager.get_task_status(&low_id).await, Err(Error::TaskNotFound { .. })));
assert_eq!(
manager
.get_task_status(&high_id)
@@ -2792,263 +2788,6 @@ async fn test_high_priority_request_displaces_lower_priority_when_queue_full() {
);
}
#[tokio::test]
async fn displaced_task_remains_queryable() {
let manager = HealManager::new(
Arc::new(MockStorage),
Some(HealConfig {
queue_size: 1,
..HealConfig::default()
}),
);
let mut displaced = HealRequest::new(
HealType::Bucket {
bucket: "displaced-bucket".to_string(),
},
HealOptions::default(),
HealPriority::Low,
);
displaced.id = "displaced-task".to_string();
let displaced_id = displaced.id.clone();
manager
.submit_heal_request(displaced)
.await
.expect("displaced request should queue");
let successor = HealRequest::new(
HealType::Bucket {
bucket: "successor-bucket".to_string(),
},
HealOptions::default(),
HealPriority::High,
);
manager
.submit_heal_request(successor)
.await
.expect("successor should displace low work");
let report = manager
.get_task_report(&displaced_id)
.await
.expect("displaced report should remain queryable");
assert!(matches!(report.status, HealTaskStatus::Failed { ref error } if error.contains("reason=displaced")));
}
#[tokio::test]
async fn displaced_archive_failure_keeps_queryable_terminal() {
let manager = HealManager::new(Arc::new(MockStorage), None);
let mut request = HealRequest::new(
HealType::Bucket {
bucket: "archive-failure".to_string(),
},
HealOptions::default(),
HealPriority::Low,
);
request.id = "archive-failure-task".to_string();
let request_id = request.id.clone();
// The synchronous sidecar is the authoritative fallback when the normal
// completed-task archive has no entry (the failure window that must not
// turn an Accepted ID into NotFound).
record_displaced_terminal(&manager.displaced_terminals, &request);
assert!(manager.completed_heals.lock().await.is_empty());
assert!(matches!(
manager.get_task_status(&request_id).await,
Ok(HealTaskStatus::Failed { error }) if error.contains("reason=displaced")
));
}
#[tokio::test]
async fn scheduler_retry_displacement_keeps_evicted_task_queryable() {
let manager = Arc::new(HealManager::new(
Arc::new(MockStorage),
Some(HealConfig {
queue_size: 1,
event_driven_scheduler_enable: false,
..HealConfig::default()
}),
));
let mut retry_request = HealRequest::object("retry-transition".to_string(), "object".to_string(), None);
retry_request.priority = HealPriority::High;
let retry_id = retry_request.id.clone();
manager
.submit_heal_request(retry_request)
.await
.expect("retry request should queue");
// Process exactly one queue cycle so the retry task is spawned without a
// background scheduler consuming the filler request before the retry wakes.
process_manager_queue_once(&manager).await;
tokio::time::timeout(Duration::from_secs(1), async {
loop {
if manager.retrying_heals.lock().await.contains_key(&retry_id) {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("retry request should enter backoff");
let filler = HealRequest::new(
HealType::Bucket {
bucket: "retry-displaced-filler".to_string(),
},
HealOptions::default(),
HealPriority::Low,
);
let filler_id = filler.id.clone();
manager
.submit_heal_request(filler)
.await
.expect("filler request should occupy the queue");
tokio::time::timeout(Duration::from_secs(5), async {
loop {
if matches!(
manager.get_task_status(&filler_id).await,
Ok(HealTaskStatus::Failed { ref error }) if error.contains("reason=displaced")
) {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("retry admission should displace the filler request");
assert_eq!(manager.get_queue_length().await, 1);
assert_eq!(
manager.get_task_status(&retry_id).await.expect("retry should be queued"),
HealTaskStatus::Pending
);
}
#[tokio::test]
async fn concurrent_displacers_produce_one_terminal_generation() {
let manager = Arc::new(HealManager::new(
Arc::new(MockStorage),
Some(HealConfig {
queue_size: 1,
..HealConfig::default()
}),
));
let mut displaced = HealRequest::new(
HealType::Bucket {
bucket: "concurrent-displaced".to_string(),
},
HealOptions::default(),
HealPriority::Low,
);
displaced.id = "concurrent-displaced-task".to_string();
let displaced_id = displaced.id.clone();
manager
.submit_heal_request(displaced)
.await
.expect("initial request should queue");
let first = HealRequest::new(
HealType::Bucket {
bucket: "concurrent-successor-a".to_string(),
},
HealOptions::default(),
HealPriority::High,
);
let second = HealRequest::new(
HealType::Bucket {
bucket: "concurrent-successor-b".to_string(),
},
HealOptions::default(),
HealPriority::High,
);
let (first_result, second_result) = tokio::join!(manager.submit_heal_request(first), manager.submit_heal_request(second));
let accepted = [&first_result, &second_result]
.into_iter()
.filter(|result| matches!(result, Ok(HealAdmissionResult::Accepted)))
.count();
assert_eq!(accepted, 1, "exactly one concurrent displacer should win the full queue");
assert!(
first_result.is_ok() && second_result.is_ok(),
"the losing request should receive a typed Full result"
);
let terminals = lock_displaced_terminals(&manager.displaced_terminals);
assert_eq!(terminals.len(), 1);
assert!(terminals.contains_key(&displaced_id));
}
#[tokio::test]
async fn successor_chain_is_bounded_and_authorized() {
let manager = HealManager::new(
Arc::new(MockStorage),
Some(HealConfig {
queue_size: 1,
..HealConfig::default()
}),
);
let mut original = HealRequest::new(
HealType::Bucket {
bucket: "authorized-original".to_string(),
},
HealOptions::default(),
HealPriority::Low,
);
original.id = "authorized-original-task".to_string();
let original_id = original.id.clone();
manager.submit_heal_request(original).await.expect("original should queue");
let mut duplicate = HealRequest::new(
HealType::Bucket {
bucket: "authorized-original".to_string(),
},
HealOptions::default(),
HealPriority::Low,
);
duplicate.id = "authorized-duplicate-task".to_string();
let duplicate_id = duplicate.id.clone();
manager
.submit_heal_request(duplicate)
.await
.expect("same-target duplicate should merge");
let successor = HealRequest::new(
HealType::Bucket {
bucket: "authorized-successor".to_string(),
},
HealOptions::default(),
HealPriority::High,
);
let successor_id = successor.id.clone();
manager.submit_heal_request(successor).await.expect("successor should queue");
assert!(manager.task_aliases.lock().await.is_empty());
assert!(matches!(manager.get_task_status(&original_id).await, Ok(HealTaskStatus::Failed { .. })));
assert!(matches!(manager.get_task_status(&duplicate_id).await, Ok(HealTaskStatus::Failed { .. })));
assert_eq!(
manager
.get_task_status(&successor_id)
.await
.expect("successor should remain queued"),
HealTaskStatus::Pending
);
}
#[tokio::test]
async fn displaced_terminal_expires_after_bounded_ttl() {
let manager = HealManager::new(Arc::new(MockStorage), None);
let mut request = HealRequest::new(
HealType::Bucket {
bucket: "expires".to_string(),
},
HealOptions::default(),
HealPriority::Low,
);
request.id = "expires-task".to_string();
let request_id = request.id.clone();
record_displaced_terminal(&manager.displaced_terminals, &request);
{
let mut terminals = lock_displaced_terminals(&manager.displaced_terminals);
let entry =
Arc::get_mut(terminals.get_mut(&request_id).expect("terminal should be retained")).expect("test owns terminal entry");
entry.completed_at = SystemTime::now() - KEEP_HEAL_TASK_STATUS_DURATION - Duration::from_secs(1);
}
assert!(matches!(manager.get_task_status(&request_id).await, Err(Error::TaskNotFound { .. })));
}
#[tokio::test]
async fn test_displacing_registered_mrf_task_drops_notice_ownership() {
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
+316 -33
View File
@@ -15,15 +15,70 @@
use serde::{Deserialize, Serialize};
use std::time::{Duration, SystemTime};
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)]
#[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, Serialize, Deserialize)]
#[serde(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
@@ -44,11 +99,38 @@ 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()
@@ -56,12 +138,87 @@ 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();
}
@@ -69,50 +226,88 @@ 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) {
self.skipped_new_versions = self.skipped_new_versions.saturating_add(1);
let Some(next) = self.skipped_new_versions.checked_add(1) else {
self.mark_unknown();
return;
};
self.skipped_new_versions = next;
self.last_update_time = Some(SystemTime::now());
self.refresh_progress_percentage();
self.refresh_estimated_completion_time();
}
pub fn record_skipped_ilm_expired(&mut self) {
self.skipped_ilm_expired = self.skipped_ilm_expired.saturating_add(1);
let Some(next) = self.skipped_ilm_expired.checked_add(1) else {
self.mark_unknown();
return;
};
self.skipped_ilm_expired = next;
self.last_update_time = Some(SystemTime::now());
self.refresh_progress_percentage();
self.refresh_estimated_completion_time();
}
fn completed_for_baseline(&self) -> u64 {
fn completed_for_baseline(&self) -> Option<u64> {
self.objects_healed
.saturating_add(self.objects_failed)
.saturating_add(self.skipped_new_versions)
.saturating_add(self.skipped_ilm_expired)
.checked_add(self.objects_failed)?
.checked_add(self.skipped_objects)?
.checked_add(self.skipped_new_versions)?
.checked_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 completed = self.completed_for_baseline();
let Some(completed) = self.completed_for_baseline() else {
self.progress_state = HealProgressState::Unknown;
self.progress_percentage = 0.0;
return;
};
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;
}
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;
if self.baseline_known {
self.progress_state = HealProgressState::Running;
self.progress_percentage = 0.0;
return;
}
self.progress_state = HealProgressState::Indeterminate;
self.progress_percentage = 0.0;
}
pub fn set_current_object(&mut self, object: Option<String>) {
@@ -125,7 +320,11 @@ impl HealProgress {
self.estimated_completion_time = None;
return;
};
if self.is_completed() || !(0.0..100.0).contains(&self.progress_percentage) || self.bytes_processed == 0 {
if self.is_completed()
|| self.progress_percentage <= 0.0
|| self.progress_percentage >= 100.0
|| self.bytes_processed == 0
{
self.estimated_completion_time = None;
return;
}
@@ -142,18 +341,39 @@ impl HealProgress {
}
pub fn is_completed(&self) -> bool {
if self.progress_percentage >= 100.0 {
return true;
}
if self.objects_total_count > 0 || self.objects_total_size > 0 {
return false;
}
self.ledger_complete
}
self.objects_scanned > 0 && self.objects_healed.saturating_add(self.objects_failed) >= self.objects_scanned
/// 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;
}
self.progress_percentage = 100.0;
self.last_update_time = Some(SystemTime::now());
self.estimated_completion_time = None;
}
pub fn get_success_rate(&self) -> f64 {
let total = self.objects_healed + self.objects_failed;
let Some(total) = self.objects_healed.checked_add(self.objects_failed) else {
return 0.0;
};
if total > 0 {
(self.objects_healed as f64 / total as f64) * 100.0
} else {
@@ -230,6 +450,7 @@ 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);
@@ -250,10 +471,8 @@ mod tests {
assert_eq!(progress.objects_healed, 8);
assert_eq!(progress.objects_failed, 2);
assert_eq!(progress.bytes_processed, 1024);
// 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_eq!(progress.progress_state, HealProgressState::Indeterminate);
assert_eq!(progress.progress_percentage, 0.0);
assert!(progress.last_update_time.is_some());
}
@@ -262,7 +481,8 @@ mod tests {
let mut progress = HealProgress::new();
progress.start_time = Some(SystemTime::now() - Duration::from_secs(10));
progress.update_progress(100, 25, 0, 4096);
progress.set_total_baseline(100, 16384);
progress.update_progress(25, 25, 0, 4096);
let eta = progress
.estimated_completion_time
@@ -275,7 +495,7 @@ mod tests {
let mut progress = HealProgress::new();
progress.set_total_baseline(10, 8192);
progress.update_progress(100, 25, 0, 4096);
progress.update_progress(25, 25, 0, 4096);
assert!((progress.progress_percentage - 50.0).abs() < 0.001);
}
@@ -285,7 +505,7 @@ mod tests {
let mut progress = HealProgress::new();
progress.set_total_baseline(10, 0);
progress.update_progress(100, 3, 2, 0);
progress.update_progress(5, 3, 2, 0);
assert!((progress.progress_percentage - 50.0).abs() < 0.001);
}
@@ -295,7 +515,7 @@ mod tests {
let mut progress = HealProgress::new();
progress.set_total_baseline(10, 0);
progress.update_progress(100, 3, 2, 0);
progress.update_progress(5, 3, 2, 0);
progress.record_skipped_new_version();
assert_eq!(progress.skipped_new_versions, 1);
@@ -336,7 +556,8 @@ 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(0, 10, 0, 2048);
progress.update_progress(10, 10, 0, 2048);
progress.mark_completed();
// All healed, should be 100%
assert!((progress.progress_percentage - 100.0).abs() < 0.001);
@@ -394,6 +615,7 @@ 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);
@@ -405,6 +627,7 @@ 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());
}
@@ -415,7 +638,7 @@ mod tests {
progress.objects_scanned = 10;
progress.objects_healed = 8;
progress.objects_failed = 2;
// healed + failed = 8 + 2 = 10 >= scanned = 10
progress.mark_completed();
assert!(progress.is_completed());
}
@@ -455,6 +678,66 @@ 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);
}
#[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();
+127 -2
View File
@@ -340,6 +340,12 @@ 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
@@ -354,6 +360,24 @@ 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 {
@@ -377,6 +401,8 @@ 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(),
@@ -384,6 +410,11 @@ 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,
}
}
@@ -412,6 +443,39 @@ 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;
@@ -454,6 +518,10 @@ 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.
@@ -476,14 +544,28 @@ 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
((self.processed_objects as f64 / self.total_objects as f64) * 100.0).min(99.999)
}
pub fn get_success_rate(&self) -> f64 {
let total = self.successful_objects + self.failed_objects;
let Some(total) = self.successful_objects.checked_add(self.failed_objects) else {
return 0.0;
};
if total == 0 {
return 0.0;
}
@@ -754,6 +836,14 @@ 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;
@@ -838,6 +928,41 @@ 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`.
+113
View File
@@ -57,6 +57,30 @@ 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 {
@@ -70,6 +94,17 @@ 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,
}
}
@@ -91,6 +126,34 @@ 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.
@@ -107,6 +170,17 @@ 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();
}
}
@@ -185,6 +259,17 @@ 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;
@@ -267,6 +352,34 @@ impl CheckpointManager {
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 {
+115
View File
@@ -1296,6 +1296,7 @@ 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);
}
@@ -1639,6 +1640,120 @@ 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;
+26 -2
View File
@@ -19,6 +19,8 @@ use base64::engine::general_purpose::URL_SAFE_NO_PAD;
use rustfs_common::heal_channel::{HealOpts, HealScanMode};
use rustfs_madmin::heal_commands::HealResultItem;
use serde::{Deserialize, Serialize};
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
use std::sync::Arc;
use tracing::{debug, error, warn};
@@ -34,6 +36,9 @@ 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 {
@@ -785,11 +790,30 @@ 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 = baseline.objects_count.saturating_add(usage.objects_count);
baseline.bytes = baseline.bytes.saturating_add(usage.size);
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),
};
}
}
let identity = info.snapshot_identity();
let mut hasher = DefaultHasher::new();
identity.last_update.hash(&mut hasher);
identity.scanner_cycle.hash(&mut hasher);
identity.scanner_epoch.hash(&mut hasher);
let mut scope = buckets.to_vec();
scope.sort_unstable();
scope.hash(&mut hasher);
baseline.generation = Some(hasher.finish());
Ok(Some(baseline))
}
+7 -3
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_progress(0, 1, 0, 0);
progress.update_stage(1, 1);
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_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
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_progress(4, 4, 0, 0);
progress.update_stage(4, 4);
true
}
@@ -831,6 +831,10 @@ 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", {
+37 -14
View File
@@ -13,6 +13,7 @@
// 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};
impl HealTask {
pub(super) async fn heal_bucket(&self, bucket: &str) -> Result<()> {
@@ -32,7 +33,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("bucket: {bucket}")));
progress.update_progress(0, 3, 0, 0);
progress.update_stage(0, 3);
}
// Step 1: Check if bucket exists
@@ -66,7 +67,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(1, 3, 0, 0);
progress.update_stage(1, 3);
}
// Step 2: Perform bucket heal using ecstore
@@ -122,7 +123,7 @@ impl HealTask {
if !self.options.recursive {
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
Ok(())
}
@@ -142,7 +143,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
Err(Error::TaskExecutionFailed {
message: format!("Failed to heal bucket {bucket}: {e}"),
@@ -245,6 +246,7 @@ 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;
@@ -286,14 +288,14 @@ 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);
telemetry_unknown |= !increment_counter(&mut scanned);
}
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("{bucket}/{object}")));
progress.update_progress(scanned, healed, failed, bytes);
}
let error = match self
@@ -304,13 +306,13 @@ impl HealTask {
.await
{
Ok((result, None)) => {
healed = healed.saturating_add(1);
bytes = bytes.saturating_add(u64::try_from(result.object_size).unwrap_or_default());
telemetry_unknown |= !increment_counter(&mut healed);
telemetry_unknown |= !add_bytes(&mut bytes, u64::try_from(result.object_size).unwrap_or(u64::MAX));
self.record_result_item(result).await;
None
}
Ok((_, Some(err))) if is_missing_object_dir_heal_result(object, &err) => {
healed = healed.saturating_add(1);
telemetry_unknown |= !increment_counter(&mut healed);
debug!(
target: "rustfs::heal::task",
event = EVENT_HEAL_BUCKET_RESULT,
@@ -329,6 +331,7 @@ 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,
@@ -357,7 +360,7 @@ impl HealTask {
);
retry.push(item);
} else {
failed = failed.saturating_add(1);
telemetry_unknown |= !increment_counter(&mut failed);
if err.is_recoverable_heal() {
retryable_failed = retryable_failed.saturating_add(1);
} else {
@@ -384,7 +387,10 @@ impl HealTask {
}
let mut progress = self.progress.write().await;
progress.update_progress(scanned, healed, failed, bytes);
progress.update_object_progress(scanned, healed, failed, skipped, bytes);
if telemetry_unknown {
progress.mark_unknown();
}
}
pending = retry;
retry_attempt = retry_attempt.saturating_add(1);
@@ -431,7 +437,7 @@ impl HealTask {
Ok(())
}
pub(super) async fn apply_erasure_set_usage_baseline(&self, buckets: &[String]) -> Result<()> {
pub(super) async fn apply_erasure_set_usage_baseline(&self, buckets: &[String], set_disk_id: &str) -> Result<()> {
let baseline = match self
.await_with_control(self.storage.erasure_set_usage_baseline(buckets))
.await
@@ -442,9 +448,26 @@ impl HealTask {
Err(_) => return Ok(()),
};
let HealBucketUsageBaseline { objects_count, bytes } = baseline;
let HealBucketUsageBaseline {
objects_count,
bytes,
generation,
} = baseline;
let generation = generation.map(|snapshot_generation| {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
snapshot_generation.hash(&mut hasher);
set_disk_id.hash(&mut hasher);
self.options.pool_index.hash(&mut hasher);
self.options.set_index.hash(&mut hasher);
hasher.finish()
});
let mut progress = self.progress.write().await;
progress.set_total_baseline(objects_count, bytes);
if let Some(generation) = generation {
progress.set_total_baseline_with_generation(objects_count, bytes, generation);
} else {
progress.set_total_baseline(objects_count, bytes);
}
Ok(())
}
}
+8 -10
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_progress(0, 4, 0, 0);
progress.update_stage(0, 4);
}
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).await?;
self.apply_erasure_set_usage_baseline(&buckets, &set_disk_id).await?;
let healing_marker = format!("{set_disk_id}:{}", self.id);
if let Some((disk, resume_manager, _)) = replacement_resume.as_ref() {
@@ -244,7 +244,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_progress(4, 4, 0, 0);
progress.update_stage(4, 4);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal disk format for {set_disk_id}: {e}"),
@@ -297,7 +297,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_progress(4, 4, 0, 0);
progress.update_stage(4, 4);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal disk format for {set_disk_id}: {e}"),
@@ -307,7 +307,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(1, 4, 0, 0);
progress.update_stage(1, 4);
}
// The rebuilt disks are formatted now: mark them as healing so
@@ -336,7 +336,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(2, 4, 0, 0);
progress.update_stage(2, 4);
}
// Step 3: Heal bucket structure
@@ -420,7 +420,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 4, 0, 0);
progress.update_stage(3, 4);
}
// Step 4: Execute erasure set heal with resume
@@ -463,9 +463,7 @@ impl HealTask {
};
{
let mut progress = self.progress.write().await;
let bytes_processed = progress.bytes_processed;
progress.update_progress(4, 4, 0, bytes_processed);
self.progress.write().await.update_stage(4, 4);
}
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_progress(0, 3, 0, 0);
progress.update_stage(0, 3);
}
// Step 1: Check if object exists
@@ -74,7 +74,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(1, 3, 0, 0);
progress.update_stage(1, 3);
}
// Step 2: Perform metadata heal using ecstore
@@ -122,7 +122,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
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_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
self.record_result_item(result).await;
Ok(())
@@ -167,7 +167,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
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_progress(0, 3, 0, 0);
progress.update_stage(0, 3);
}
// Step 1: Check if object exists
@@ -236,7 +236,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(1, 3, 0, 0);
progress.update_stage(1, 3);
}
// Step 2: Perform EC decode heal using ecstore
@@ -284,7 +284,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
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_progress(3, 3, 0, object_size);
progress.update_object_progress(1, 1, 0, 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_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
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_progress(0, 4, 0, 0);
progress.update_stage(0, 4);
}
// Step 1: Check if object exists and get metadata
@@ -132,7 +132,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(1, 3, 0, 0);
progress.update_stage(1, 3);
}
// Step 2: directly call ecstore to perform heal
@@ -187,7 +187,7 @@ impl HealTask {
);
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
return Ok(());
}
@@ -207,7 +207,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
if Self::should_return_typed_heal_error(&e) {
@@ -249,7 +249,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, object_size);
progress.update_object_progress(1, 1, 0, 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_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
return Ok(());
}
@@ -295,7 +295,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(3, 3, 0, 0);
progress.update_stage(3, 3);
}
if Self::should_return_typed_heal_error(&e) {
@@ -414,7 +414,7 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(4, 4, 0, object_size);
progress.update_object_progress(1, 1, 0, 0, object_size);
}
self.record_result_item(result).await;
Ok(())
+3
View File
@@ -2096,6 +2096,7 @@ 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()
});
@@ -2119,6 +2120,8 @@ 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);
}
@@ -722,10 +722,6 @@ pub struct DeleteVersionsResponse {
pub errors: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
#[prost(message, optional, tag = "3")]
pub error: ::core::option::Option<Error>,
/// Senders dual-write the legacy strings and typed entries. Receivers prefer typed entries
/// when present and fall back to strings for peers that predate this field. Code zero means success.
#[prost(message, repeated, tag = "4")]
pub item_errors: ::prost::alloc::vec::Vec<Error>,
}
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct ReadMultipleRequest {
-4
View File
@@ -2106,9 +2106,6 @@ pub enum ChannelClass {
Bulk,
}
// Keep multiplexed unary RPCs below h2's per-connection small-frame budget.
const INTERNODE_RPC_CONCURRENCY_LIMIT: usize = 64;
/// Whether control/bulk channel isolation is enabled (env-gated, default off for safe rollout).
fn channel_isolation_enabled() -> bool {
rustfs_utils::get_env_bool(
@@ -2191,7 +2188,6 @@ async fn build_channel(dial_addr: &str, cache_key: &str) -> Result<Channel, Box<
let mut connector = Endpoint::from_shared(dial_addr.to_string())?
// Fast connection timeout for dead peer detection
.connect_timeout(connect_timeout)
.concurrency_limit(INTERNODE_RPC_CONCURRENCY_LIMIT)
// TCP-level keepalive - OS will probe connection
.tcp_keepalive(Some(tcp_keepalive))
// Disable Nagle so latency-sensitive control-plane RPCs (locks/health) are not batched
-3
View File
@@ -493,9 +493,6 @@ message DeleteVersionsResponse {
bool success = 1;
repeated string errors = 2;
optional Error error = 3;
// Senders dual-write the legacy strings and typed entries. Receivers prefer typed entries
// when present and fall back to strings for peers that predate this field. Code zero means success.
repeated Error item_errors = 4;
}
message ReadMultipleRequest {
+11 -43
View File
@@ -146,29 +146,6 @@ fn encode_file_info_msgpack(value: &FileInfo) -> std::result::Result<Vec<u8>, Di
encode_msgpack_with_capacity(value, "FileInfo", FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT)
}
fn encode_delete_versions_errors(disk_errors: Vec<Option<DiskError>>) -> (Vec<String>, Vec<Error>) {
let mut errors = Vec::with_capacity(disk_errors.len());
let mut item_errors = Vec::with_capacity(disk_errors.len());
for error in disk_errors {
match error {
Some(error) => {
let code = match &error {
DiskError::Io(source) if source.kind() == std::io::ErrorKind::NotFound => DiskError::FileNotFound.to_u32(),
_ => error.to_u32(),
};
let error_info = error.to_string();
errors.push(error_info.clone());
item_errors.push(Error { code, error_info });
}
None => {
errors.push(String::new());
item_errors.push(Error::default());
}
}
}
(errors, item_errors)
}
fn encode_msgpack_named<T: serde::Serialize>(value: &T, value_name: &str) -> std::result::Result<Vec<u8>, DiskError> {
let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(MSGPACK_ENCODE_CAPACITY_HINT)).with_struct_map();
value
@@ -575,7 +552,6 @@ impl NodeService {
success: false,
errors: Vec::new(),
error: Some(DiskError::other(format!("decode FileInfoVersions failed: {err}")).into()),
item_errors: Vec::new(),
}));
}
};
@@ -587,26 +563,30 @@ impl NodeService {
success: false,
errors: Vec::new(),
error: Some(DiskError::other(format!("decode DeleteOptions failed: {err}")).into()),
item_errors: Vec::new(),
}));
}
};
let (errors, item_errors) =
encode_delete_versions_errors(disk.delete_versions(&request.volume, versions, opts).await);
let errors = disk
.delete_versions(&request.volume, versions, opts)
.await
.into_iter()
.map(|error| match error {
Some(e) => e.to_string(),
None => "".to_string(),
})
.collect();
Ok(Response::new(DeleteVersionsResponse {
success: true,
errors,
error: None,
item_errors,
}))
} else {
Ok(Response::new(DeleteVersionsResponse {
success: false,
errors: Vec::new(),
error: Some(DiskError::other("cannot find disk".to_string()).into()),
item_errors: Vec::new(),
}))
}
}
@@ -1632,8 +1612,8 @@ impl NodeService {
mod tests {
use super::{
compat_response_json, decode_msgpack_or_json, decode_rename_data_request_file_info,
encode_batch_read_version_response_payloads, encode_delete_versions_errors, encode_file_info_msgpack, encode_msgpack,
encode_msgpack_named, encode_read_multiple_response_payloads, encode_rename_data_response_payloads,
encode_batch_read_version_response_payloads, encode_file_info_msgpack, encode_msgpack, encode_msgpack_named,
encode_read_multiple_response_payloads, encode_rename_data_response_payloads,
};
use crate::storage::rpc::node_service::make_server;
use crate::storage::storage_api::ReadMultipleResp;
@@ -1652,18 +1632,6 @@ mod tests {
count: u32,
}
#[test]
fn delete_versions_response_dual_writes_typed_item_errors() {
let raw_not_found = super::DiskError::Io(std::io::Error::from(std::io::ErrorKind::NotFound));
let (errors, item_errors) = encode_delete_versions_errors(vec![Some(raw_not_found), None]);
assert!(errors[0].starts_with("io error "));
assert!(errors[1].is_empty());
assert_eq!(item_errors[0].code, super::DiskError::FileNotFound.to_u32());
assert_eq!(item_errors[0].error_info, errors[0]);
assert_eq!(item_errors[1].code, 0);
}
#[tokio::test]
#[serial]
async fn handle_read_version_records_attribution_for_missing_disk() {