fix(storage): harden rebalance decommission state (#3515)

This commit is contained in:
cxymds
2026-06-22 12:03:13 +08:00
committed by GitHub
parent 2dcc32db5a
commit 2f25cf606e
43 changed files with 10151 additions and 835 deletions
@@ -0,0 +1,165 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
#[cfg(test)]
mod tests {
use crate::common::{RustFSTestEnvironment, init_logging};
use aws_sdk_s3::Client;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration};
use serial_test::serial;
async fn create_versioned_bucket(client: &Client, bucket: &str) {
client
.create_bucket()
.bucket(bucket)
.send()
.await
.expect("create versioning proof bucket");
client
.put_bucket_versioning()
.bucket(bucket)
.versioning_configuration(
VersioningConfiguration::builder()
.status(BucketVersioningStatus::Enabled)
.build(),
)
.send()
.await
.expect("enable versioning for proof bucket");
}
async fn assert_current_get_is_delete_marker_not_found(client: &Client, bucket: &str, key: &str) {
let err = client
.get_object()
.bucket(bucket)
.key(key)
.send()
.await
.expect_err("current get should fail when latest version is a delete marker");
let service_err = err.into_service_error();
let code = service_err.meta().code();
assert!(
matches!(code, Some("NoSuchKey") | Some("NotFound")),
"current get should expose delete-marker not-found semantics, got {service_err:?}"
);
}
#[tokio::test]
#[serial]
async fn test_versioning_only_delete_marker_has_minio_compatible_visibility_for_migration_proof() {
init_logging();
let mut env = RustFSTestEnvironment::new().await.expect("create test environment");
env.start_rustfs_server(vec![]).await.expect("start RustFS");
let client = env.create_s3_client();
let bucket = "delete-marker-migration-proof-a";
let key = "only-delete-marker.txt";
create_versioned_bucket(&client, bucket).await;
let delete_marker = client
.delete_object()
.bucket(bucket)
.key(key)
.send()
.await
.expect("create delete marker for absent object");
let delete_marker_version_id = delete_marker
.version_id()
.expect("MinIO-compatible delete marker should have a version id");
let listed = client
.list_object_versions()
.bucket(bucket)
.prefix(key)
.send()
.await
.expect("list delete marker only object");
let versions = listed.versions();
let markers = listed.delete_markers();
assert!(versions.is_empty(), "only-delete-marker case must not report data versions");
assert_eq!(markers.len(), 1, "only-delete-marker case must report the delete marker");
assert_eq!(markers[0].version_id(), Some(delete_marker_version_id));
assert_eq!(markers[0].is_latest(), Some(true));
assert_current_get_is_delete_marker_not_found(&client, bucket, key).await;
}
#[tokio::test]
#[serial]
async fn test_versioning_delete_marker_plus_history_remains_visible_for_migration_proof() {
init_logging();
let mut env = RustFSTestEnvironment::new().await.expect("create test environment");
env.start_rustfs_server(vec![]).await.expect("start RustFS");
let client = env.create_s3_client();
let bucket = "delete-marker-migration-proof-b";
let key = "delete-marker-with-history.txt";
let body = b"historical version body";
create_versioned_bucket(&client, bucket).await;
let put = client
.put_object()
.bucket(bucket)
.key(key)
.body(ByteStream::from_static(body))
.send()
.await
.expect("put historical version");
let data_version_id = put.version_id().expect("put should return data version id");
let delete_marker = client
.delete_object()
.bucket(bucket)
.key(key)
.send()
.await
.expect("create delete marker over historical version");
let delete_marker_version_id = delete_marker.version_id().expect("delete marker should have a version id");
let listed = client
.list_object_versions()
.bucket(bucket)
.prefix(key)
.send()
.await
.expect("list delete marker plus historical version");
let versions = listed.versions();
let markers = listed.delete_markers();
assert_eq!(versions.len(), 1, "history case must report the historical data version");
assert_eq!(markers.len(), 1, "history case must report the latest delete marker");
assert_eq!(versions[0].version_id(), Some(data_version_id));
assert_eq!(versions[0].is_latest(), Some(false));
assert_eq!(markers[0].version_id(), Some(delete_marker_version_id));
assert_eq!(markers[0].is_latest(), Some(true));
assert_current_get_is_delete_marker_not_found(&client, bucket, key).await;
let historical = client
.get_object()
.bucket(bucket)
.key(key)
.version_id(data_version_id)
.send()
.await
.expect("historical version get should succeed");
let bytes = historical
.body
.collect()
.await
.expect("collect historical version body")
.into_bytes();
assert_eq!(bytes.as_ref(), body);
}
}
+4
View File
@@ -83,6 +83,10 @@ mod delete_objects_versioning_test;
#[cfg(test)]
mod delete_object_no_content_length_test;
// Delete-marker visibility baseline for data-movement migration proof.
#[cfg(test)]
mod delete_marker_migration_semantics_test;
// Regression test for Issue #2252: ListObjectVersions misses newest version after put -> delete -> put
#[cfg(test)]
mod list_object_versions_regression_test;
+3 -1
View File
@@ -140,7 +140,9 @@ pub mod object {
pub mod rebalance {
pub use crate::rebalance::{
DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarnings, RebalanceInfo, RebalanceMeta, RebalanceStats,
DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo,
RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, decode_rebalance_stop_propagation_record,
encode_rebalance_stop_propagation_record,
};
}
File diff suppressed because it is too large Load Diff
+200 -32
View File
@@ -32,10 +32,13 @@ use std::hash::{Hash, Hasher};
use std::sync::{Mutex, OnceLock};
use std::time::{Duration, SystemTime};
use tokio::time::timeout;
use tracing::{error, warn};
use tracing::{debug, error, info, warn};
/// After this many consecutive admin-call failures, mark the peer as offline.
const CONSECUTIVE_FAILURE_THRESHOLD: u32 = 3;
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
const LOG_SUBSYSTEM_NOTIFICATION: &str = "notification";
const EVENT_NOTIFICATION_PEER_PROPAGATION: &str = "notification_peer_propagation";
/// Cached result from the last successful admin call to a peer.
struct PeerAdminCache {
@@ -472,6 +475,15 @@ impl NotificationSys {
let host = client.grid_host.clone();
futures.push(async move { client.reload_pool_meta().await.map_err(|err| (host, err)) });
} else {
warn!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "reload_pool_meta",
result = "peer_unreachable",
peer_index = idx,
"notification peer propagation"
);
failures.push(format!("peer[{idx}] reload_pool_meta failed: peer is not reachable"));
}
}
@@ -479,7 +491,16 @@ impl NotificationSys {
for result in join_all(futures).await {
if let Err((host, err)) = result {
let failure = format!("peer {host} reload_pool_meta failed: {err}");
error!("notification reload_pool_meta err {}", failure);
error!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "reload_pool_meta",
result = "peer_failed",
peer = %host,
error = %err,
"notification peer propagation"
);
failures.push(failure);
}
}
@@ -489,87 +510,187 @@ impl NotificationSys {
#[tracing::instrument(skip(self))]
pub async fn load_rebalance_meta(&self, start: bool) -> Result<()> {
let failures = self.load_rebalance_meta_failures(start).await?;
aggregate_notification_failures("load_rebalance_meta", failures)
}
#[tracing::instrument(skip(self))]
pub async fn load_rebalance_meta_failures(&self, start: bool) -> Result<Vec<String>> {
let operation = format!("load_rebalance_meta(start={start})");
let mut failures = Vec::new();
let mut futures = Vec::with_capacity(self.peer_clients.len());
for (idx, client) in self.peer_clients.iter().enumerate() {
if let Some(client) = client {
warn!(
"notification load_rebalance_meta start: {}, index: {}, client: {:?}",
start, idx, client.host
);
let host = client.grid_host.clone();
futures.push(async move { client.load_rebalance_meta(start).await.map_err(|err| (host, err)) });
futures.push(async move {
let result = client.load_rebalance_meta(start).await;
(host, result)
});
} else {
warn!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "load_rebalance_meta",
result = "peer_unreachable",
peer_index = idx,
start_rebalance = start,
"notification peer propagation"
);
failures.push(format!("peer[{idx}] {operation} failed: peer is not reachable"));
}
}
for result in join_all(futures).await {
if let Err((host, err)) = result {
for (host, result) in join_all(futures).await {
if let Err(err) = result {
let failure = format!("peer {host} {operation} failed: {err}");
error!("notification load_rebalance_meta err {}", failure);
error!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "load_rebalance_meta",
result = "peer_failed",
peer = %host,
start_rebalance = start,
error = %err,
"notification peer propagation"
);
failures.push(failure);
} else {
warn!("notification load_rebalance_meta success");
debug!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "load_rebalance_meta",
result = "peer_success",
peer = %host,
start_rebalance = start,
"notification peer propagation"
);
}
}
aggregate_notification_failures("load_rebalance_meta", failures)
Ok(failures)
}
pub async fn stop_rebalance(&self) -> Result<()> {
warn!("notification stop_rebalance start");
pub async fn stop_rebalance(&self, expected_rebalance_id: Option<&str>) -> Result<()> {
let failures = self.stop_rebalance_failures(expected_rebalance_id).await?;
aggregate_notification_failures("stop_rebalance", failures)
}
pub async fn stop_rebalance_failures(&self, expected_rebalance_id: Option<&str>) -> Result<Vec<String>> {
info!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "stop_rebalance",
state = "started",
"notification peer propagation"
);
let Some(store) = resolve_object_store_handle() else {
error!("stop_rebalance: not init");
error!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "stop_rebalance",
result = "failed",
reason = "object_layer_not_initialized",
"notification peer propagation"
);
return Err(Error::other("stop_rebalance: object layer not initialized"));
};
// warn!("notification stop_rebalance load_rebalance_meta");
// self.load_rebalance_meta(false).await;
// warn!("notification stop_rebalance load_rebalance_meta done");
let mut failures = Vec::new();
let mut futures = Vec::with_capacity(self.peer_clients.len());
for (idx, client) in self.peer_clients.iter().enumerate() {
if let Some(client) = client {
let host = client.grid_host.clone();
futures.push(async move { client.stop_rebalance().await.map_err(|err| (host, err)) });
futures.push(async move {
let result = client.stop_rebalance(expected_rebalance_id).await;
(host, result)
});
} else {
warn!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "stop_rebalance",
result = "peer_unreachable",
peer_index = idx,
"notification peer propagation"
);
failures.push(format!("peer[{idx}] stop_rebalance failed: peer is not reachable"));
}
}
for result in join_all(futures).await {
if let Err((host, err)) = result {
for (host, result) in join_all(futures).await {
if let Err(err) = result {
let failure = format!("peer {host} stop_rebalance failed: {err}");
error!("notification stop_rebalance err {}", failure);
error!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "stop_rebalance",
result = "peer_failed",
peer = %host,
error = %err,
"notification peer propagation"
);
failures.push(failure);
} else {
debug!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "stop_rebalance",
result = "peer_success",
peer = %host,
"notification peer propagation"
);
}
}
warn!("notification stop_rebalance stop_rebalance start");
match store.stop_rebalance().await {
match store.stop_rebalance_for_id(expected_rebalance_id).await {
Ok(_) => {
if let Err(err) = store.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt).await {
error!("notification stop_rebalance local save err {:?}", err);
error!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "stop_rebalance",
result = "local_save_failed",
error = %err,
"notification peer propagation"
);
return Err(Error::other(format!(
"local stop_rebalance save_rebalance_stats(stopped_at) failed: {err}"
)));
}
}
Err(err) => {
error!("notification stop_rebalance local stop err {:?}", err);
error!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "stop_rebalance",
result = "local_stop_failed",
error = %err,
"notification peer propagation"
);
return Err(Error::other(format!("local stop_rebalance stop failed: {err}")));
}
}
if let Err(err) = aggregate_notification_failures("stop_rebalance", failures) {
warn!("{err}");
}
warn!("notification stop_rebalance stop_rebalance done");
Ok(())
info!(
event = EVENT_NOTIFICATION_PEER_PROPAGATION,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_NOTIFICATION,
action = "stop_rebalance",
result = if failures.is_empty() { "success" } else { "partial_failure" },
"notification peer propagation"
);
Ok(failures)
}
pub async fn load_bucket_metadata(&self, bucket: &str) -> Result<()> {
@@ -1094,6 +1215,53 @@ mod tests {
assert!(msg.contains("local save failed"));
}
#[test]
fn load_rebalance_meta_aggregate_failures_return_error() {
let err = aggregate_notification_failures(
"load_rebalance_meta(start=true)",
vec!["peer[0] load_rebalance_meta failed: peer is not reachable".to_string()],
)
.expect_err("load_rebalance_meta peer failures must be returned");
let msg = err.to_string();
assert!(msg.contains("load_rebalance_meta(start=true)"));
assert!(msg.contains("1 failure(s)"));
assert!(msg.contains("peer[0]"));
}
#[test]
fn stop_rebalance_aggregate_failures_return_error() {
let err = aggregate_notification_failures(
"stop_rebalance",
vec!["peer[0] stop_rebalance failed: peer is not reachable".to_string()],
)
.expect_err("stop_rebalance peer failures must be returned");
let msg = err.to_string();
assert!(msg.contains("stop_rebalance"));
assert!(msg.contains("1 failure(s)"));
assert!(msg.contains("peer[0]"));
}
#[tokio::test]
async fn reload_pool_meta_reports_unreachable_peers() {
let sys = NotificationSys {
peer_clients: vec![None],
all_peer_clients: Vec::new(),
peer_admin_caches: vec![Mutex::new(PeerAdminCache::new())],
};
let err = sys
.reload_pool_meta()
.await
.expect_err("unreachable peers should fail pool metadata reload");
let msg = err.to_string();
assert!(msg.contains("reload_pool_meta"));
assert!(msg.contains("1 failure(s)"));
assert!(msg.contains("peer[0]"));
}
#[tokio::test]
async fn load_bucket_metadata_reports_unreachable_peers() {
let sys = NotificationSys {
+113
View File
@@ -314,6 +314,24 @@ impl ReadPlan {
rs = http_range_spec_from_object_info(oi, part_number);
}
if opts.raw_data_movement_read {
let (visible_offset, visible_length) = if let Some(rs) = rs {
rs.get_offset_length(oi.size)?
} else {
(0, oi.size)
};
return Ok(Self {
storage_offset: visible_offset,
storage_length: visible_length,
object_size: oi.size,
transform: ReadTransform::Plain {
visible_offset,
visible_length,
},
});
}
let mut is_encrypted = oi.is_encrypted();
let (algo, compression_backend, mut is_compressed) = oi.compression_read_plan()?;
@@ -2162,6 +2180,101 @@ mod tests {
));
}
#[tokio::test]
async fn test_raw_data_movement_read_plan_bypasses_compression_transform() {
let object_info = ObjectInfo {
size: 3_000_000,
user_defined: Arc::new(HashMap::from([
("x-minio-internal-compression".to_string(), "klauspost/compress/s2".to_string()),
("x-minio-internal-actual-size".to_string(), "4194304".to_string()),
])),
..Default::default()
};
let opts = ObjectOptions {
raw_data_movement_read: true,
..Default::default()
};
let plan = ReadPlan::build(None, &object_info, &opts, &HeaderMap::new())
.await
.expect("raw data movement read should bypass compression planning");
assert_eq!(plan.storage_offset, 0);
assert_eq!(plan.storage_length, object_info.size);
assert_eq!(plan.object_size, object_info.size);
assert!(matches!(
plan.transform,
ReadTransform::Plain {
visible_offset: 0,
visible_length: 3_000_000
}
));
}
#[tokio::test]
async fn test_raw_data_movement_read_plan_bypasses_encryption_transform() {
let object_info = ObjectInfo {
size: 128,
user_defined: Arc::new(HashMap::from([
("X-Amz-Server-Side-Encryption".to_string(), "aws:kms".to_string()),
("X-Amz-Server-Side-Encryption-Iv".to_string(), "AAAAAAAAAAAAAAAA".to_string()),
("X-Amz-Server-Side-Encryption-Key".to_string(), BASE64_STANDARD.encode([7_u8; 32])),
("x-rustfs-encryption-original-size".to_string(), "64".to_string()),
])),
..Default::default()
};
let opts = ObjectOptions {
raw_data_movement_read: true,
..Default::default()
};
let plan = ReadPlan::build(None, &object_info, &opts, &HeaderMap::new())
.await
.expect("raw data movement read should not require decryption material");
assert_eq!(plan.storage_offset, 0);
assert_eq!(plan.storage_length, object_info.size);
assert_eq!(plan.object_size, object_info.size);
assert!(matches!(
plan.transform,
ReadTransform::Plain {
visible_offset: 0,
visible_length: 128
}
));
}
#[tokio::test]
async fn test_raw_data_movement_read_plan_bypasses_ssec_header_resolution() {
let object_info = ObjectInfo {
size: 256,
user_defined: Arc::new(HashMap::from([
(SSEC_ALGORITHM_HEADER.to_string(), "AES256".to_string()),
(SSEC_KEY_MD5_HEADER.to_string(), "stored-key-md5".to_string()),
])),
..Default::default()
};
let opts = ObjectOptions {
raw_data_movement_read: true,
..Default::default()
};
let plan = ReadPlan::build(None, &object_info, &opts, &HeaderMap::new())
.await
.expect("raw data movement read should not require SSE-C request headers");
assert_eq!(plan.storage_offset, 0);
assert_eq!(plan.storage_length, object_info.size);
assert_eq!(plan.object_size, object_info.size);
assert!(matches!(
plan.transform,
ReadTransform::Plain {
visible_offset: 0,
visible_length: 256
}
));
}
#[tokio::test]
async fn test_get_object_reader_allows_encrypted_full_object_passthrough() {
async_with_vars([("__RUSTFS_SSE_SIMPLE_CMK", Some(BASE64_STANDARD.encode([0u8; 32])))], async {
+1
View File
@@ -25,6 +25,7 @@ pub struct ObjectOptions {
pub skip_free_version: bool,
pub data_movement: bool,
pub raw_data_movement_read: bool,
pub src_pool_idx: usize,
pub user_defined: HashMap<String, String>,
pub preserve_etag: Option<String>,
+2199 -319
View File
File diff suppressed because it is too large Load Diff
+9 -2
View File
@@ -25,13 +25,15 @@ const EVENT_REBALANCE_LISTING: &str = "rebalance_listing";
const REBAL_META_FMT: u16 = 1; // Replace with actual format value
const REBAL_META_VER: u16 = 1; // Replace with actual version value
const REBAL_META_NAME: &str = "rebalance.bin";
pub(crate) const REBAL_META_NAME: &str = "rebalance.bin";
const DEFAULT_REBALANCE_MAX_ATTEMPTS: usize = 3;
const REBALANCE_MAX_ATTEMPTS_ENV: &str = "RUSTFS_REBALANCE_MAX_ATTEMPTS";
const REBALANCE_STOP_PROPAGATION_ERROR_PREFIX: &str = "rebalance stop propagation incomplete: ";
const REBALANCE_LISTING_RETRY_BASE_DELAY: Duration = Duration::from_millis(250);
const REBALANCE_MIGRATION_RETRY_BASE_DELAY: Duration = Duration::from_millis(250);
const REBALANCE_MIGRATION_LOCK_RETRY_CAP: Duration = Duration::from_secs(10);
const REBALANCE_DEFERRED_ENTRY_ERROR_PREFIX: &str = "deferred transient rebalance entry failure:";
const REBALANCE_CLEANUP_WARNING_ENTRY_LIMIT: usize = 10;
mod control;
mod entry;
@@ -41,7 +43,12 @@ mod runtime;
mod types;
mod worker;
pub use types::{DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarnings, RebalanceInfo, RebalanceMeta, RebalanceStats};
pub(crate) use meta::is_rebalance_conflicting_with_decommission;
pub use meta::{decode_rebalance_stop_propagation_record, encode_rebalance_stop_propagation_record};
pub use types::{
DiskStat, RebalSaveOpt, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceInfo, RebalanceMeta,
RebalanceStats, RebalanceStopPropagationRecord,
};
use types::{RebalanceBucketConfigs, RebalanceBucketOutcome, RebalanceEntryOutcome};
#[cfg(test)]
+98 -4
View File
@@ -2,7 +2,9 @@ use super::meta::{
clone_first_arc, clone_rebalance_pool_stats, defer_bucket_in_rebalance_queue, ensure_valid_rebalance_pool_index,
invalid_rebalance_pool_index_error, is_rebalance_conflicting_with_decommission, mark_rebalance_bucket_done,
merge_rebalance_meta, percent_free_ratio, rebalance_metadata_not_initialized_error, record_rebalance_cleanup_warning_in_meta,
resolve_next_rebalance_bucket, should_accept_rebalance_stats_update, should_pool_participate, stop_rebalance_meta_snapshot,
record_rebalance_stop_propagation_snapshot, resolve_next_rebalance_bucket, rollback_rebalance_start_meta_snapshot_for_id,
should_accept_rebalance_stats_update, should_pool_participate, stop_rebalance_meta_snapshot_for_id,
validate_init_rebalance_state,
};
use super::worker::{
rebalance_meta_lock_error, resolve_load_rebalance_stats_update_result, resolve_rebalance_meta_load_result,
@@ -10,7 +12,8 @@ use super::worker::{
};
use super::{
DiskStat, EVENT_REBALANCE_BUCKET, EVENT_REBALANCE_STATE, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_REBALANCE, REBAL_META_NAME,
RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats,
RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord,
encode_rebalance_stop_propagation_record,
};
use crate::error::{Error, Result};
use crate::object_api::ObjectOptions;
@@ -247,6 +250,35 @@ impl ECStore {
Ok(id)
}
#[tracing::instrument(skip(self, bucktes))]
pub async fn init_and_start_rebalance(self: &Arc<Self>, bucktes: Vec<String>) -> Result<String> {
let _start_guard = self.start_gate.lock().await;
let decommission_running = self.is_decommission_running().await;
{
let rebalance_meta = self.rebalance_meta.read().await;
validate_init_rebalance_state(decommission_running, rebalance_meta.as_ref())?;
}
let id = self.init_rebalance_meta(bucktes).await?;
if let Err(start_err) = self.start_rebalance().await {
if let Err(rollback_err) = self
.rollback_rebalance_start_without_worker_for_id(Some(&id), start_err.to_string())
.await
{
return Err(Error::other(format!(
"failed to start rebalance after metadata initialized for {id}: {start_err}; rollback failed: {rollback_err}"
)));
}
return Err(Error::other(format!(
"failed to start rebalance after metadata initialized for {id}; local metadata was finalized as failed: {start_err}"
)));
}
Ok(id)
}
#[tracing::instrument(skip(self, fi))]
pub async fn update_pool_stats(&self, pool_index: usize, bucket: String, fi: &FileInfo) -> Result<()> {
self.update_pool_stats_batch(pool_index, bucket, &[fi]).await
@@ -389,12 +421,24 @@ impl ECStore {
false
}
pub async fn current_rebalance_id(&self) -> Option<String> {
let rebalance_meta = self.rebalance_meta.read().await;
rebalance_meta
.as_ref()
.and_then(|meta| (!meta.id.is_empty()).then(|| meta.id.clone()))
}
#[tracing::instrument(skip(self))]
pub async fn stop_rebalance(self: &Arc<Self>) -> Result<()> {
self.stop_rebalance_for_id(None).await
}
#[tracing::instrument(skip(self))]
pub async fn stop_rebalance_for_id(self: &Arc<Self>, expected_id: Option<&str>) -> Result<()> {
let meta_to_save = {
let mut rebalance_meta = self.rebalance_meta.write().await;
stop_rebalance_meta_snapshot(rebalance_meta.as_mut(), OffsetDateTime::now_utc())
};
stop_rebalance_meta_snapshot_for_id(rebalance_meta.as_mut(), OffsetDateTime::now_utc(), expected_id)
}?;
if let Some(meta_to_save) = meta_to_save {
let pool = clone_first_arc(self.pools.as_slice(), "stop_rebalance: no pools available")?;
@@ -407,4 +451,54 @@ impl ECStore {
Ok(())
}
async fn rollback_rebalance_start_without_worker_for_id(
self: &Arc<Self>,
expected_id: Option<&str>,
start_error: String,
) -> Result<()> {
let meta_to_save = {
let mut rebalance_meta = self.rebalance_meta.write().await;
rollback_rebalance_start_meta_snapshot_for_id(
rebalance_meta.as_mut(),
OffsetDateTime::now_utc(),
expected_id,
start_error,
)
};
if let Some(meta_to_save) = meta_to_save {
let pool = clone_first_arc(self.pools.as_slice(), "rollback_rebalance_start: no pools available")?;
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, &meta_to_save, "rollback_rebalance_start")
.await,
"rollback_rebalance_start",
)?;
}
Ok(())
}
pub async fn record_rebalance_stop_propagation(self: &Arc<Self>, record: RebalanceStopPropagationRecord) -> Result<()> {
if !record.has_failures() {
return Ok(());
}
let encoded_error = encode_rebalance_stop_propagation_record(&record);
let meta_to_save = {
let mut rebalance_meta = self.rebalance_meta.write().await;
record_rebalance_stop_propagation_snapshot(rebalance_meta.as_mut(), encoded_error, OffsetDateTime::now_utc())
};
if let Some(meta_to_save) = meta_to_save {
let pool = clone_first_arc(self.pools.as_slice(), "record_rebalance_stop_propagation: no pools available")?;
resolve_rebalance_meta_save_result(
self.save_rebalance_meta_with_merge(pool, &meta_to_save, "record_rebalance_stop_propagation")
.await,
"record_rebalance_stop_propagation",
)?;
}
Ok(())
}
}
+1 -1
View File
@@ -122,7 +122,7 @@ impl ECStore {
true,
&crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc::Rebal,
)
.await
.await?
{
expired += 1;
debug!(
+192 -25
View File
@@ -1,8 +1,9 @@
use super::{
EVENT_REBALANCE_BUCKET, EVENT_REBALANCE_STATE, Error, GetObjectReader, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_REBALANCE,
ObjectInfo, ObjectOptions, PutObjReader, REBAL_META_FMT, REBAL_META_NAME, REBAL_META_VER,
REBALANCE_DEFERRED_ENTRY_ERROR_PREFIX, RebalSaveOpt, RebalStatus, RebalanceCleanupWarnings, RebalanceMeta, RebalanceStats,
Result,
REBALANCE_CLEANUP_WARNING_ENTRY_LIMIT, REBALANCE_DEFERRED_ENTRY_ERROR_PREFIX, REBALANCE_STOP_PROPAGATION_ERROR_PREFIX,
RebalSaveOpt, RebalStatus, RebalanceCleanupWarningEntry, RebalanceCleanupWarnings, RebalanceMeta, RebalanceStats,
RebalanceStopPropagationRecord, Result,
};
use crate::config::com::{read_config_with_metadata, save_config_with_opts};
use crate::error::is_err_operation_canceled;
@@ -197,15 +198,15 @@ pub(super) fn is_rebalance_pool_started(pool_stat: &RebalanceStats) -> bool {
pool_stat.participating && pool_stat.info.status == RebalStatus::Started
}
pub(super) fn is_rebalance_in_progress(meta: &RebalanceMeta) -> bool {
if meta.stopped_at.is_some() {
return false;
}
meta.pool_stats.iter().any(is_rebalance_pool_started)
pub(super) fn is_rebalance_pool_active(pool_stat: &RebalanceStats) -> bool {
is_rebalance_pool_started(pool_stat) || pool_stat.info.stopping
}
pub(super) fn is_rebalance_conflicting_with_decommission(meta: &RebalanceMeta) -> bool {
pub(super) fn is_rebalance_in_progress(meta: &RebalanceMeta) -> bool {
meta.pool_stats.iter().any(is_rebalance_pool_active)
}
pub(crate) fn is_rebalance_conflicting_with_decommission(meta: &RebalanceMeta) -> bool {
is_rebalance_in_progress(meta)
}
@@ -224,6 +225,30 @@ pub(super) fn rebalance_meta_load_unknown_format_error(fmt: u16) -> Error {
pub(super) fn rebalance_meta_load_unknown_version_error(ver: u16) -> Error {
Error::other(format!("rebalance metadata load failed: unknown version {ver}"))
}
pub fn encode_rebalance_stop_propagation_record(record: &RebalanceStopPropagationRecord) -> String {
match serde_json::to_string(record) {
Ok(payload) => format!("{REBALANCE_STOP_PROPAGATION_ERROR_PREFIX}{payload}"),
Err(err) => {
let payload = serde_json::json!({
"encodeError": err.to_string(),
"stopFailures": [],
"terminalReloadFailures": [],
});
format!("{REBALANCE_STOP_PROPAGATION_ERROR_PREFIX}{payload}")
}
}
}
pub fn decode_rebalance_stop_propagation_record(message: &str) -> Option<RebalanceStopPropagationRecord> {
let payload = message.strip_prefix(REBALANCE_STOP_PROPAGATION_ERROR_PREFIX)?;
serde_json::from_str(payload).ok()
}
fn is_rebalance_stop_propagation_error(message: Option<&str>) -> bool {
message.is_some_and(|message| message.starts_with(REBALANCE_STOP_PROPAGATION_ERROR_PREFIX))
}
pub(super) fn rebalance_goal_reached(init_free_space: u64, init_capacity: u64, bytes: u64, percent_free_goal: f64) -> bool {
if init_capacity == 0 {
return false;
@@ -384,10 +409,17 @@ pub(super) fn record_rebalance_cleanup_warning_in_meta(
};
pool_stat.cleanup_warnings.count = pool_stat.cleanup_warnings.count.saturating_add(1);
pool_stat.cleanup_warnings.last_message = Some(message);
pool_stat.cleanup_warnings.last_message = Some(message.clone());
pool_stat.cleanup_warnings.last_bucket = Some(bucket.to_string());
pool_stat.cleanup_warnings.last_object = Some(object.to_string());
pool_stat.cleanup_warnings.last_at = Some(now);
pool_stat.cleanup_warnings.entries.push(RebalanceCleanupWarningEntry {
bucket: bucket.to_string(),
object: object.to_string(),
message,
timestamp: Some(now),
});
truncate_rebalance_cleanup_warning_entries(&mut pool_stat.cleanup_warnings.entries);
meta.last_refreshed_at = Some(now);
Ok(())
}
@@ -572,6 +604,17 @@ pub(super) fn validate_start_rebalance_state(decommission_running: bool, meta_lo
Ok(())
}
pub(super) fn validate_init_rebalance_state(decommission_running: bool, current_meta: Option<&RebalanceMeta>) -> Result<()> {
if !ensure_rebalance_not_decommissioning(decommission_running) {
return Err(Error::DecommissionAlreadyRunning);
}
if current_meta.is_some_and(is_rebalance_in_progress) {
return Err(Error::RebalanceAlreadyRunning);
}
Ok(())
}
pub(super) fn should_skip_start_rebalance(cancel_attached: bool, in_progress: bool) -> bool {
cancel_attached && in_progress
}
@@ -654,6 +697,31 @@ pub(super) fn merge_rebalance_cleanup_warnings(remote: &mut RebalanceCleanupWarn
remote.last_object = local.last_object.clone();
remote.last_at = local.last_at;
}
merge_rebalance_cleanup_warning_entries(&mut remote.entries, &local.entries);
let retained_entries = u64::try_from(remote.entries.len()).unwrap_or(u64::MAX);
remote.count = remote.count.max(retained_entries);
}
pub(super) fn merge_rebalance_cleanup_warning_entries(
remote: &mut Vec<RebalanceCleanupWarningEntry>,
local: &[RebalanceCleanupWarningEntry],
) {
for entry in local {
if !remote.iter().any(|existing| existing == entry) {
remote.push(entry.clone());
}
}
remote.sort_by_key(|entry| entry.timestamp);
truncate_rebalance_cleanup_warning_entries(remote);
}
fn truncate_rebalance_cleanup_warning_entries(entries: &mut Vec<RebalanceCleanupWarningEntry>) {
if entries.len() > REBALANCE_CLEANUP_WARNING_ENTRY_LIMIT {
let remove_count = entries.len() - REBALANCE_CLEANUP_WARNING_ENTRY_LIMIT;
entries.drain(0..remove_count);
}
}
pub(super) fn should_replace_rebalance_cleanup_warning(
@@ -667,6 +735,16 @@ pub(super) fn should_replace_rebalance_cleanup_warning(
}
}
pub(super) fn merge_rebalance_stop_propagation_error(remote: Option<String>, local: Option<String>) -> Option<String> {
if is_rebalance_stop_propagation_error(local.as_deref()) {
local
} else if is_rebalance_stop_propagation_error(remote.as_deref()) {
remote
} else {
None
}
}
pub(super) fn merge_rebalance_pool_stats(remote: &mut RebalanceStats, local: &RebalanceStats) {
remote.init_free_space = remote.init_free_space.max(local.init_free_space);
remote.init_capacity = remote.init_capacity.max(local.init_capacity);
@@ -702,19 +780,23 @@ pub(super) fn merge_rebalance_pool_stats(remote: &mut RebalanceStats, local: &Re
match local.info.status {
RebalStatus::Failed => {
remote.info.status = RebalStatus::Failed;
remote.info.stopping = false;
remote.info.end_time = local.info.end_time.or(remote.info.end_time);
remote.info.last_error = local.info.last_error.clone().or_else(|| remote.info.last_error.clone());
}
RebalStatus::Stopped => {
if remote.info.status != RebalStatus::Failed {
remote.info.status = RebalStatus::Stopped;
remote.info.stopping = false;
remote.info.end_time = local.info.end_time.or(remote.info.end_time);
remote.info.last_error = None;
remote.info.last_error =
merge_rebalance_stop_propagation_error(remote.info.last_error.clone(), local.info.last_error.clone());
}
}
RebalStatus::Completed => {
if !matches!(remote.info.status, RebalStatus::Failed | RebalStatus::Stopped) {
remote.info.status = RebalStatus::Completed;
remote.info.stopping = false;
remote.info.end_time = local.info.end_time.or(remote.info.end_time);
remote.info.last_error = None;
}
@@ -722,6 +804,7 @@ pub(super) fn merge_rebalance_pool_stats(remote: &mut RebalanceStats, local: &Re
RebalStatus::Started => {
if !is_rebalance_terminal_status(remote.info.status) {
remote.info.status = RebalStatus::Started;
remote.info.stopping |= local.info.stopping;
remote.info.last_error = local.info.last_error.clone();
}
}
@@ -736,7 +819,6 @@ pub(super) fn merge_rebalance_meta(remote: &mut RebalanceMeta, local: &Rebalance
}
if !local.id.is_empty() && remote.id != local.id {
*remote = local.clone();
return;
}
@@ -761,35 +843,120 @@ pub(super) fn mark_started_rebalance_pools_stopped(meta: &mut RebalanceMeta, sto
for pool_stat in meta.pool_stats.iter_mut() {
if pool_stat.info.status == RebalStatus::Started {
pool_stat.info.status = RebalStatus::Stopped;
pool_stat.info.stopping = false;
pool_stat.info.end_time.get_or_insert(stop_time);
if !is_rebalance_stop_propagation_error(pool_stat.info.last_error.as_deref()) {
pool_stat.info.last_error = None;
}
}
}
}
pub(super) fn mark_started_rebalance_pools_stopping(meta: &mut RebalanceMeta) {
for pool_stat in meta.pool_stats.iter_mut() {
if pool_stat.info.status == RebalStatus::Started {
pool_stat.info.stopping = true;
pool_stat.info.last_error = None;
}
}
}
pub(super) fn apply_stopped_at(meta: &mut RebalanceMeta, now: OffsetDateTime) {
meta.stopped_at = Some(now);
mark_started_rebalance_pools_stopped(meta, now);
meta.stopped_at.get_or_insert(now);
mark_started_rebalance_pools_stopping(meta);
}
pub(super) fn clear_rebalance_cancel_token(meta: Option<&mut RebalanceMeta>) -> bool {
let Some(meta) = meta else {
return false;
};
if let Some(cancel_tx) = meta.cancel.take() {
cancel_tx.cancel();
return true;
}
false
}
pub(super) fn stop_rebalance_state(meta: &mut RebalanceMeta, now: OffsetDateTime) {
if let Some(cancel_tx) = meta.cancel.take() {
cancel_tx.cancel();
}
let stop_time = meta.stopped_at.unwrap_or(now);
clear_rebalance_cancel_token(Some(meta));
if meta.stopped_at.is_none() && is_rebalance_in_progress(meta) {
meta.stopped_at = Some(stop_time);
}
if meta.stopped_at.is_some() {
mark_started_rebalance_pools_stopped(meta, stop_time);
apply_stopped_at(meta, now);
} else if meta.stopped_at.is_some() {
mark_started_rebalance_pools_stopping(meta);
}
}
pub(super) fn stop_rebalance_meta_snapshot_for_id(
meta: Option<&mut RebalanceMeta>,
now: OffsetDateTime,
expected_id: Option<&str>,
) -> Result<Option<RebalanceMeta>> {
let Some(meta) = meta else {
return Ok(None);
};
if let Some(expected_id) = expected_id
&& !expected_id.is_empty()
&& meta.id != expected_id
{
return Err(Error::other(format!(
"rebalance stop id mismatch: expected {expected_id}, found {}",
meta.id
)));
}
stop_rebalance_state(meta, now);
meta.last_refreshed_at = Some(now);
Ok(Some(meta.clone()))
}
pub(super) fn rollback_rebalance_start_meta_snapshot_for_id(
meta: Option<&mut RebalanceMeta>,
now: OffsetDateTime,
expected_id: Option<&str>,
start_error: String,
) -> Option<RebalanceMeta> {
meta.and_then(|meta| {
if let Some(expected_id) = expected_id
&& !expected_id.is_empty()
&& meta.id != expected_id
{
return None;
}
clear_rebalance_cancel_token(Some(meta));
meta.stopped_at.get_or_insert(now);
meta.last_refreshed_at = Some(now);
for pool_stat in meta.pool_stats.iter_mut() {
if pool_stat.info.status == RebalStatus::Started {
pool_stat.info.status = RebalStatus::Failed;
pool_stat.info.stopping = false;
pool_stat.info.end_time.get_or_insert(now);
pool_stat.info.last_error = Some(start_error.clone());
}
}
Some(meta.clone())
})
}
pub(super) fn stop_rebalance_meta_snapshot(meta: Option<&mut RebalanceMeta>, now: OffsetDateTime) -> Option<RebalanceMeta> {
let meta = meta?;
stop_rebalance_state(meta, now);
meta.last_refreshed_at = Some(now);
Some(meta.clone())
}
pub(super) fn record_rebalance_stop_propagation_snapshot(
meta: Option<&mut RebalanceMeta>,
encoded_error: String,
now: OffsetDateTime,
) -> Option<RebalanceMeta> {
meta.map(|meta| {
stop_rebalance_state(meta, now);
for pool_stat in meta.pool_stats.iter_mut() {
if pool_stat.participating || pool_stat.info.stopping || pool_stat.info.status == RebalStatus::Stopped {
pool_stat.info.last_error = Some(encoded_error.clone());
}
}
meta.last_refreshed_at = Some(now);
meta.clone()
})
@@ -25,7 +25,7 @@ use super::meta::{
record_rebalance_cleanup_warning_in_meta, remove_rebalanced_buckets_from_queue, resolve_next_rebalance_bucket,
resolve_rebalance_participants, should_accept_rebalance_stats_update, should_ignore_rebalance_data_usage_cache,
should_pool_participate, should_preserve_rebalance_stopped_state, should_skip_start_rebalance, stop_rebalance_meta_snapshot,
stop_rebalance_state, take_bucket_from_rebalance_queue, validate_start_rebalance_state,
stop_rebalance_state, take_bucket_from_rebalance_queue, validate_init_rebalance_state, validate_start_rebalance_state,
};
use super::migration::{
MigrationBackend, MigrationVersionResult, migrate_entry_version, migrate_entry_version_with_retry_wait,
@@ -1199,6 +1199,7 @@ fn test_merge_rebalance_meta_preserves_updates_from_multiple_pools() {
last_bucket: Some("bucket-a".to_string()),
last_object: Some("local-object".to_string()),
last_at: Some(warning_at),
entries: Vec::new(),
},
..Default::default()
},
@@ -1242,7 +1243,7 @@ fn test_merge_rebalance_meta_does_not_overwrite_failed_with_started_stats() {
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
status: RebalStatus::Stopped,
..Default::default()
},
num_versions: 8,
@@ -1307,7 +1308,7 @@ fn test_merge_rebalance_meta_preserves_failed_status_over_stopped() {
stopped_at: Some(stopped_at),
pool_stats: vec![RebalanceStats {
info: RebalanceInfo {
status: RebalStatus::Stopped,
status: RebalStatus::Started,
end_time: Some(stopped_at),
..Default::default()
},
@@ -1812,7 +1813,7 @@ fn test_should_accept_rebalance_stats_update_rejects_stopped_meta() {
stopped_at: Some(OffsetDateTime::UNIX_EPOCH),
pool_stats: vec![RebalanceStats {
info: RebalanceInfo {
status: RebalStatus::Started,
status: RebalStatus::Stopped,
..Default::default()
},
..Default::default()
@@ -2173,6 +2174,93 @@ fn test_validate_start_rebalance_state_allows_loaded_meta() {
validate_start_rebalance_state(false, true).expect("loaded rebalance meta should allow start");
}
#[test]
fn test_validate_init_rebalance_state_rejects_running_decommission() {
let err = validate_init_rebalance_state(true, None).expect_err("running decommission should block rebalance init");
assert!(matches!(err, Error::DecommissionAlreadyRunning));
}
#[test]
fn test_validate_init_rebalance_state_rejects_active_rebalance() {
let now = OffsetDateTime::now_utc();
let meta = RebalanceMeta {
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
start_time: Some(now),
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let err = validate_init_rebalance_state(false, Some(&meta)).expect_err("active rebalance should block rebalance init");
assert!(matches!(err, Error::RebalanceAlreadyRunning));
}
#[test]
fn test_validate_init_rebalance_state_allows_terminal_or_missing_rebalance() {
let completed = RebalanceMeta {
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Completed,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
validate_init_rebalance_state(false, None).expect("missing rebalance meta should allow init");
validate_init_rebalance_state(false, Some(&completed)).expect("terminal rebalance meta should allow init");
}
#[tokio::test]
async fn test_init_and_start_rebalance_rejects_second_start_after_gate() {
let now = OffsetDateTime::now_utc();
let active_meta = RebalanceMeta {
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
start_time: Some(now),
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
let endpoint_pools: crate::endpoints::EndpointServerPools = Vec::new().into();
let store = Arc::new(crate::store::ECStore {
id: uuid::Uuid::new_v4(),
disk_map: std::collections::HashMap::new(),
pools: Vec::new(),
peer_sys: crate::rpc::S3PeerSys::new(&endpoint_pools),
pool_meta: tokio::sync::RwLock::new(crate::pools::PoolMeta::default()),
rebalance_meta: tokio::sync::RwLock::new(Some(active_meta)),
decommission_cancelers: tokio::sync::RwLock::new(Vec::new()),
start_gate: tokio::sync::Mutex::new(()),
pool_meta_save_gate: tokio::sync::Mutex::new(()),
local_disk_map: crate::global::GLOBAL_LOCAL_DISK_MAP.clone(),
local_disk_id_map: crate::global::GLOBAL_LOCAL_DISK_ID_MAP.clone(),
local_disk_set_drives: crate::global::GLOBAL_LOCAL_DISK_SET_DRIVES.clone(),
tier_config_mgr: crate::tier::tier::TierConfigMgr::new(),
event_notifier: crate::event_notification::EventNotifier::new(),
bucket_monitor: std::sync::OnceLock::new(),
});
let err = store
.init_and_start_rebalance(vec!["bucket".to_string()])
.await
.expect_err("second rebalance start should be rejected before metadata init");
assert!(matches!(err, Error::RebalanceAlreadyRunning));
}
#[test]
fn test_percent_free_ratio_zero_capacity_is_zero() {
assert_eq!(percent_free_ratio(100, 0), 0.0);
@@ -2412,7 +2500,7 @@ fn test_resolve_rebalance_participants_respects_runtime_pool_count() {
RebalanceStats {
participating: false,
info: RebalanceInfo {
status: RebalStatus::Started,
status: RebalStatus::Stopped,
start_time: Some(now),
..Default::default()
},
@@ -2421,7 +2509,7 @@ fn test_resolve_rebalance_participants_respects_runtime_pool_count() {
RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
status: RebalStatus::Stopped,
start_time: Some(now),
..Default::default()
},
@@ -2542,7 +2630,7 @@ fn test_is_rebalance_conflicting_with_decommission_false_when_stopped() {
participating: true,
info: RebalanceInfo {
start_time: Some(now),
status: RebalStatus::Started,
status: RebalStatus::Stopped,
..Default::default()
},
..Default::default()
@@ -2562,7 +2650,7 @@ fn test_is_rebalance_in_progress_stopped_takes_precedence() {
participating: true,
info: RebalanceInfo {
start_time: Some(now),
status: RebalStatus::Started,
status: RebalStatus::Stopped,
..Default::default()
},
..Default::default()
@@ -2871,6 +2959,7 @@ fn test_complete_rebalance_pools_with_empty_queue_preserves_cleanup_warnings() {
last_bucket: Some("bucket-a".to_string()),
last_object: Some("obj.txt".to_string()),
last_at: Some(warning_at),
entries: Vec::new(),
},
..Default::default()
}],
@@ -2916,8 +3005,9 @@ fn test_apply_stopped_at_transitions_started_pools_only() {
apply_stopped_at(&mut meta, now);
assert_eq!(meta.stopped_at, Some(now));
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Stopped);
assert_eq!(meta.pool_stats[0].info.end_time, Some(now));
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started);
assert!(meta.pool_stats[0].info.stopping);
assert_eq!(meta.pool_stats[0].info.end_time, None);
assert_eq!(meta.pool_stats[0].info.last_error, None);
assert_eq!(meta.pool_stats[1].info.status, RebalStatus::Failed);
@@ -2948,8 +3038,9 @@ fn test_stop_rebalance_state_cancels_token_and_marks_stopped_when_in_progress()
assert!(cancel_clone.is_cancelled());
assert!(meta.cancel.is_none());
assert_eq!(meta.stopped_at, Some(now));
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Stopped);
assert_eq!(meta.pool_stats[0].info.end_time, Some(now));
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started);
assert!(meta.pool_stats[0].info.stopping);
assert_eq!(meta.pool_stats[0].info.end_time, None);
}
#[test]
@@ -3005,8 +3096,9 @@ fn test_stop_rebalance_state_normalizes_started_pool_when_stopped_at_already_set
assert!(cancel_clone.is_cancelled());
assert!(meta.cancel.is_none());
assert_eq!(meta.stopped_at, Some(stopped_at));
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Stopped);
assert_eq!(meta.pool_stats[0].info.end_time, Some(stopped_at));
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started);
assert!(meta.pool_stats[0].info.stopping);
assert_eq!(meta.pool_stats[0].info.end_time, None);
assert_eq!(meta.pool_stats[0].info.last_error, None);
}
@@ -3040,13 +3132,15 @@ fn test_stop_rebalance_meta_snapshot_stops_meta_and_returns_snapshot() {
assert!(meta.cancel.is_none());
assert_eq!(meta.stopped_at, Some(now));
assert_eq!(meta.last_refreshed_at, Some(now));
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Stopped);
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started);
assert!(meta.pool_stats[0].info.stopping);
assert!(snapshot.cancel.is_none());
assert_eq!(snapshot.stopped_at, Some(now));
assert_eq!(snapshot.last_refreshed_at, Some(now));
assert_eq!(snapshot.pool_stats[0].info.status, RebalStatus::Stopped);
assert_eq!(snapshot.pool_stats[0].info.end_time, Some(now));
assert_eq!(snapshot.pool_stats[0].info.status, RebalStatus::Started);
assert!(snapshot.pool_stats[0].info.stopping);
assert_eq!(snapshot.pool_stats[0].info.end_time, None);
}
#[test]
@@ -3108,8 +3202,9 @@ fn test_apply_rebalance_save_option_stopped_at_updates_refresh_and_statuses() {
assert_eq!(meta.stopped_at, Some(now));
assert_eq!(meta.last_refreshed_at, Some(now));
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Stopped);
assert_eq!(meta.pool_stats[0].info.end_time, Some(now));
assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started);
assert!(meta.pool_stats[0].info.stopping);
assert_eq!(meta.pool_stats[0].info.end_time, None);
assert!(meta.pool_stats[0].info.last_error.is_none());
assert_eq!(meta.pool_stats[1].info.status, RebalStatus::Failed);
assert_eq!(meta.pool_stats[1].info.last_error.as_deref(), Some("previous failure"));
+1
View File
@@ -225,6 +225,7 @@ impl ECStore {
"Preserved stopped rebalance status"
);
} else {
pool_stat.info.stopping = false;
apply_rebalance_terminal_event(
&mut pool_stat.info.status,
&mut pool_stat.info.end_time,
+40
View File
@@ -4,6 +4,7 @@ use time::OffsetDateTime;
use tokio_util::sync::CancellationToken;
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RebalanceStats {
#[serde(rename = "ifs")]
pub init_free_space: u64, // Pool free space at the start of rebalance
@@ -70,6 +71,7 @@ pub enum RebalSaveOpt {
}
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RebalanceInfo {
#[serde(rename = "startTs")]
pub start_time: Option<OffsetDateTime>, // Time at which rebalance-start was issued
@@ -79,9 +81,25 @@ pub struct RebalanceInfo {
pub last_error: Option<String>, // Last rebalance error message
#[serde(rename = "status")]
pub status: RebalStatus, // Current state of rebalance operation
#[serde(rename = "stopping", default)]
pub stopping: bool, // True after stop is requested and before worker terminal acknowledgement
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct RebalanceCleanupWarningEntry {
#[serde(rename = "bucket", default)]
pub bucket: String,
#[serde(rename = "object", default)]
pub object: String,
#[serde(rename = "message", default)]
pub message: String,
#[serde(rename = "timestamp", default)]
pub timestamp: Option<OffsetDateTime>,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct RebalanceCleanupWarnings {
#[serde(rename = "count", default)]
pub count: u64,
@@ -93,6 +111,27 @@ pub struct RebalanceCleanupWarnings {
pub last_object: Option<String>,
#[serde(rename = "lastAt", default)]
pub last_at: Option<OffsetDateTime>,
#[serde(rename = "entries", default)]
pub entries: Vec<RebalanceCleanupWarningEntry>,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct RebalanceStopPropagationRecord {
#[serde(rename = "stopAttemptAt", default)]
pub stop_attempt_at: Option<OffsetDateTime>,
#[serde(rename = "stopFailures", default)]
pub stop_failures: Vec<String>,
#[serde(rename = "terminalReloadAttemptAt", default)]
pub terminal_reload_attempt_at: Option<OffsetDateTime>,
#[serde(rename = "terminalReloadFailures", default)]
pub terminal_reload_failures: Vec<String>,
}
impl RebalanceStopPropagationRecord {
pub fn has_failures(&self) -> bool {
!self.stop_failures.is_empty() || !self.terminal_reload_failures.is_empty()
}
}
#[allow(dead_code)]
@@ -103,6 +142,7 @@ pub struct DiskStat {
}
#[derive(Debug, Default, Serialize, Deserialize, Clone)]
#[serde(deny_unknown_fields)]
pub struct RebalanceMeta {
#[serde(skip)]
pub cancel: Option<CancellationToken>, // To be invoked on rebalance-stop
+16 -4
View File
@@ -52,7 +52,7 @@ use tokio::{net::TcpStream, time::Duration};
use tonic::Request;
use tonic::service::interceptor::InterceptedService;
use tonic::transport::Channel;
use tracing::{Instrument, warn};
use tracing::{Instrument, debug, warn};
pub const PEER_RESTSIGNAL: &str = "signal";
pub const PEER_RESTSUB_SYS: &str = "sub-sys";
@@ -916,11 +916,13 @@ impl PeerRestClient {
.await
}
pub async fn stop_rebalance(&self) -> Result<()> {
pub async fn stop_rebalance(&self, expected_rebalance_id: Option<&str>) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let request = Request::new(StopRebalanceRequest {});
let request = Request::new(StopRebalanceRequest {
expected_rebalance_id: expected_rebalance_id.unwrap_or_default().to_string(),
});
let response = client.stop_rebalance(request).await?.into_inner();
if !response.success {
@@ -945,7 +947,17 @@ impl PeerRestClient {
let response = client.load_rebalance_meta(request).await?.into_inner();
warn!("load_rebalance_meta response {:?}, grid_host: {:?}", response, &self.grid_host);
debug!(
event = "peer_rebalance_meta",
component = "ecstore",
subsystem = "peer_rest_client",
action = "load_rebalance_meta",
result = "response_received",
peer = %self.grid_host,
success = response.success,
start_rebalance = start_rebalance,
"peer rebalance metadata response"
);
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
+24 -7
View File
@@ -346,7 +346,8 @@ impl S3PeerSys {
#[derive(Debug)]
pub struct LocalPeerS3Client {
// pub local_disks: Vec<DiskStore>,
#[cfg(test)]
local_disks: Option<Vec<DiskStore>>,
// pub node: Node,
pub pools: Option<Vec<usize>>,
}
@@ -354,13 +355,29 @@ pub struct LocalPeerS3Client {
impl LocalPeerS3Client {
pub fn new(_node: Option<Node>, pools: Option<Vec<usize>>) -> Self {
Self {
// local_disks,
#[cfg(test)]
local_disks: None,
// node,
pools,
}
}
#[cfg(test)]
fn new_with_local_disks(_node: Option<Node>, pools: Option<Vec<usize>>, local_disks: Vec<DiskStore>) -> Self {
Self {
local_disks: Some(local_disks),
pools,
}
}
async fn local_disks_for_pools(&self) -> Vec<DiskStore> {
#[cfg(test)]
let local_disks = if let Some(local_disks) = self.local_disks.as_ref() {
local_disks.clone()
} else {
all_local_disk().await
};
#[cfg(not(test))]
let local_disks = all_local_disk().await;
let Some(pools) = self.pools.as_ref() else {
return local_disks;
@@ -1323,7 +1340,7 @@ mod tests {
assert_eq!(walk_err, DiskError::Timeout);
let info = LocalPeerS3Client::new(None, Some(vec![0]))
let info = LocalPeerS3Client::new_with_local_disks(None, Some(vec![0]), disks.clone())
.get_bucket_info(bucket, &BucketOptions::default())
.await
.expect("bucket info should still succeed after prior walk timeout");
@@ -1347,7 +1364,7 @@ mod tests {
.await
.expect("bucket should be created on one disk");
let err = LocalPeerS3Client::new(None, Some(vec![0]))
let err = LocalPeerS3Client::new_with_local_disks(None, Some(vec![0]), disks.clone())
.get_bucket_info("partial-bucket", &BucketOptions::default())
.await
.expect_err("partial bucket should not satisfy local write quorum");
@@ -1375,19 +1392,19 @@ mod tests {
.await
.expect("bucket should be created on pool 0 disk 1");
let pool0_info = LocalPeerS3Client::new(None, Some(vec![0]))
let pool0_info = LocalPeerS3Client::new_with_local_disks(None, Some(vec![0]), disks.clone())
.get_bucket_info(bucket, &BucketOptions::default())
.await
.expect("pool 0 peer should see bucket on pool 0 disks");
assert_eq!(pool0_info.name, bucket);
let pool1_err = LocalPeerS3Client::new(None, Some(vec![1]))
let pool1_err = LocalPeerS3Client::new_with_local_disks(None, Some(vec![1]), disks.clone())
.get_bucket_info(bucket, &BucketOptions::default())
.await
.expect_err("pool 1 peer should not count pool 0 disks");
assert_eq!(pool1_err, Error::VolumeNotFound);
let pool1_buckets = LocalPeerS3Client::new(None, Some(vec![1]))
let pool1_buckets = LocalPeerS3Client::new_with_local_disks(None, Some(vec![1]), disks.clone())
.list_bucket(&BucketOptions::default())
.await
.expect("pool 1 local listing should succeed against its own disks");
+154 -14
View File
@@ -1234,7 +1234,10 @@ impl rustfs_storage_api::ObjectIO for SetDisks {
//TODO: userDefined
let etag = data.stream.try_resolve_etag().unwrap_or_default();
let mut etag = data.stream.try_resolve_etag().unwrap_or_default();
if let Some(ref tag) = opts.preserve_etag {
etag = tag.clone();
}
user_defined.insert("etag".to_owned(), etag.clone());
@@ -2275,7 +2278,7 @@ impl rustfs_storage_api::ObjectOperations for SetDisks {
if let Some(err) = &gerr
&& goi.name.is_empty()
{
if opts.delete_marker {
if should_force_delete_marker_for_missing_version(&opts) {
version_found = false;
} else {
return Err(err.clone());
@@ -2336,7 +2339,7 @@ impl rustfs_storage_api::ObjectOperations for SetDisks {
fi.set_skip_tier_free_version();
}
fi.version_id = if let Some(vid) = opts.version_id {
fi.version_id = if let Some(vid) = opts.version_id.as_ref() {
Some(Uuid::parse_str(vid.as_str())?)
} else if opts.versioned {
Some(Uuid::new_v4())
@@ -2344,7 +2347,7 @@ impl rustfs_storage_api::ObjectOperations for SetDisks {
None
};
self.delete_object_version(bucket, object, &fi, opts.delete_marker)
self.delete_object_version(bucket, object, &fi, should_force_delete_marker_for_missing_version(&opts))
.await
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
@@ -2905,8 +2908,12 @@ fn should_preserve_delete_replication_state(opts: &ObjectOptions) -> bool {
}) || opts.version_purge_status() == VersionPurgeStatusType::Complete
}
fn should_force_delete_marker_for_missing_version(opts: &ObjectOptions) -> bool {
opts.delete_marker || (opts.versioned && opts.version_id.is_none() && !opts.data_movement)
}
fn resolve_delete_version_state(opts: &ObjectOptions, goi: &ObjectInfo, version_found: bool) -> (bool, bool) {
let mut mark_delete = goi.version_id.is_some();
let mut mark_delete = goi.version_id.is_some() || (opts.versioned && opts.version_id.is_none());
let mut delete_marker = opts.versioned;
if opts.version_id.is_some() {
@@ -3956,15 +3963,7 @@ impl rustfs_storage_api::MultipartOperations for SetDisks {
object_size += ext_part.size;
object_actual_size += ext_part.actual_size;
fi.parts.push(ObjectPartInfo {
etag: ext_part.etag.clone(),
number: p.part_num,
size: ext_part.size,
mod_time: ext_part.mod_time,
actual_size: ext_part.actual_size,
index: ext_part.index.clone(),
..Default::default()
});
fi.parts.push(completed_multipart_object_part(p.part_num, ext_part));
}
if let Some(wtcs) = opts.want_checksum.as_ref() {
@@ -4913,6 +4912,19 @@ fn get_complete_multipart_md5(parts: &[CompletePart]) -> String {
format!("{}-{}", etag_hex, parts.len())
}
fn completed_multipart_object_part(part_num: usize, ext_part: &ObjectPartInfo) -> ObjectPartInfo {
ObjectPartInfo {
etag: ext_part.etag.clone(),
number: part_num,
size: ext_part.size,
mod_time: ext_part.mod_time,
actual_size: ext_part.actual_size,
index: ext_part.index.clone(),
checksums: ext_part.checksums.clone(),
..Default::default()
}
}
fn complete_part_checksum(part: &CompletePart, checksum_type: rustfs_rio::ChecksumType) -> Option<Option<String>> {
match checksum_type.base() {
rustfs_rio::ChecksumType::SHA256 => Some(part.checksum_sha256.clone()),
@@ -5300,6 +5312,42 @@ mod tests {
assert!(delete_marker);
}
#[test]
fn resolve_delete_version_state_creates_marker_for_missing_latest_versioned_delete() {
let opts = ObjectOptions {
versioned: true,
..Default::default()
};
let (mark_delete, delete_marker) = resolve_delete_version_state(&opts, &ObjectInfo::default(), false);
assert!(mark_delete);
assert!(delete_marker);
}
#[test]
fn should_force_delete_marker_for_missing_version_rejects_data_movement_latest_delete() {
let opts = ObjectOptions {
versioned: true,
data_movement: true,
..Default::default()
};
assert!(!should_force_delete_marker_for_missing_version(&opts));
}
#[test]
fn should_force_delete_marker_for_missing_version_allows_explicit_marker_creation() {
let opts = ObjectOptions {
versioned: true,
data_movement: true,
delete_marker: true,
..Default::default()
};
assert!(should_force_delete_marker_for_missing_version(&opts));
}
#[test]
fn resolve_delete_version_state_skips_marker_creation_for_replica_purge_when_version_missing() {
let opts = ObjectOptions {
@@ -5433,6 +5481,33 @@ mod tests {
assert!(single_result.ends_with("-1"));
}
#[test]
fn test_completed_multipart_object_part_preserves_checksums() {
let checksums = HashMap::from([
(rustfs_rio::ChecksumType::CRC32.to_string(), "crc32-value".to_string()),
(rustfs_rio::ChecksumType::CRC32C.to_string(), "crc32c-value".to_string()),
]);
let ext_part = ObjectPartInfo {
number: 7,
etag: "etag-7".to_string(),
size: 123,
actual_size: 456,
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
index: Some(Bytes::from_static(&[1, 2, 3])),
checksums: Some(checksums.clone()),
..Default::default()
};
let completed = completed_multipart_object_part(7, &ext_part);
assert_eq!(completed.number, 7);
assert_eq!(completed.etag, ext_part.etag);
assert_eq!(completed.size, ext_part.size);
assert_eq!(completed.actual_size, ext_part.actual_size);
assert_eq!(completed.index, ext_part.index);
assert_eq!(completed.checksums, Some(checksums));
}
#[test]
fn test_get_upload_id_dir() {
// Test upload ID directory path generation
@@ -6305,6 +6380,71 @@ mod tests {
drop(temp_dirs);
}
#[tokio::test]
async fn load_file_info_versions_exact_returns_none_for_explicit_not_found() {
let format = FormatV3::new(1, 1);
let (temp_dir, endpoint, disk) = make_formatted_local_disk_for_info_test(0, &format).await;
let bucket = "bucket";
disk.make_volume(bucket).await.expect("bucket should be created");
let set_disks = SetDisks::new(
"test-owner".to_string(),
Arc::new(RwLock::new(vec![Some(disk)])),
1,
0,
0,
0,
vec![endpoint],
format,
Vec::new(),
)
.await;
let versions = set_disks
.load_file_info_versions_exact(bucket, "missing-object")
.await
.expect("explicit object not found should be accepted");
assert!(versions.is_none());
drop(temp_dir);
}
#[tokio::test]
async fn load_file_info_versions_exact_rejects_corrupt_metadata() {
let format = FormatV3::new(1, 1);
let (temp_dir, endpoint, disk) = make_formatted_local_disk_for_info_test(0, &format).await;
let bucket = "bucket";
let object = "object.txt";
disk.make_volume(bucket).await.expect("bucket should be created");
let metadata_path = format!("{object}/{STORAGE_FORMAT_FILE}");
disk.write_all(bucket, &metadata_path, bytes::Bytes::from_static(b"not-xl-meta"))
.await
.expect("corrupt metadata file should be written");
let set_disks = SetDisks::new(
"test-owner".to_string(),
Arc::new(RwLock::new(vec![Some(disk)])),
1,
0,
0,
0,
vec![endpoint],
format,
Vec::new(),
)
.await;
let err = set_disks
.load_file_info_versions_exact(bucket, object)
.await
.expect_err("corrupt exact metadata must fail closed");
assert!(!is_err_object_not_found(&err), "corrupt metadata must not be treated as not found: {err}");
drop(temp_dir);
}
#[tokio::test]
async fn list_path_still_uses_disk_after_prior_walk_timeout() {
use std::pin::Pin;
+48
View File
@@ -342,6 +342,54 @@ impl SetDisks {
Self::pick_latest_quorum_files_info(fileinfos, errs, bucket, object, read_data, incl_free_vers).await
}
pub(crate) async fn load_file_info_versions_exact(
&self,
bucket: &str,
object: &str,
) -> Result<Option<rustfs_filemeta::FileInfoVersions>> {
let disks = self.get_disks_internal().await;
if disks.is_empty() {
return Err(to_object_err(StorageError::ErasureReadQuorum, vec![bucket, object]));
}
let read_quorum = disks.len().div_ceil(2).max(1);
let (raw_fileinfos, errs) = Self::read_all_raw_file_info(&disks, bucket, object, false).await;
if let Some(err) = reduce_read_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, read_quorum) {
let object_err = to_object_err(err.into(), vec![bucket, object]);
if is_err_object_not_found(&object_err) || is_err_version_not_found(&object_err) {
return Ok(None);
}
return Err(object_err);
}
let mut shallow_versions = Vec::with_capacity(raw_fileinfos.len());
for raw_fileinfo in raw_fileinfos.into_iter().flatten() {
let meta = FileMeta::load(&raw_fileinfo.buf)
.map_err(|err| Error::other(format!("exact object metadata decode failed for {bucket}/{object}: {err}")))?;
shallow_versions.push(meta.versions);
}
if shallow_versions.len() < read_quorum {
return Err(to_object_err(StorageError::ErasureReadQuorum, vec![bucket, object]));
}
let versions = merge_file_meta_versions(read_quorum, true, 0, &shallow_versions);
if versions.is_empty() {
return Err(Error::other(format!(
"exact object metadata read returned no quorum versions for {bucket}/{object}"
)));
}
FileMeta {
versions,
..Default::default()
}
.get_all_file_info_versions(bucket, object, true)
.map(Some)
.map_err(|err| Error::other(format!("exact object versions decode failed for {bucket}/{object}: {err}")))
}
pub(super) async fn read_all_raw_file_info(
disks: &[Option<DiskStore>],
bucket: &str,
+13 -1
View File
@@ -93,7 +93,7 @@ use std::{
};
use time::OffsetDateTime;
use tokio::select;
use tokio::sync::RwLock;
use tokio::sync::{Mutex, RwLock};
use tokio::time::sleep;
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, instrument, warn};
@@ -188,6 +188,18 @@ pub struct ECStore {
pub pool_meta: RwLock<PoolMeta>,
pub rebalance_meta: RwLock<Option<RebalanceMeta>>,
pub decommission_cancelers: RwLock<Vec<Option<CancellationToken>>>,
/// Serializes rebalance/decommission start transitions.
///
/// Lock order: acquire `start_gate` before `pool_meta`, `rebalance_meta`,
/// or `decommission_cancelers`. The guarded sections may perform bounded
/// async metadata work so check/init/start cannot race across operations.
pub(crate) start_gate: Mutex<()>,
/// Serializes full-document pool metadata saves.
///
/// Lock order: acquire `pool_meta_save_gate` without holding `pool_meta`.
/// The saver then clones the latest `pool_meta` under a short read lock and
/// releases it before awaiting disk writes.
pub(crate) pool_meta_save_gate: Mutex<()>,
// Phase 2 migration pending - do not use directly.
/// Local disk maps (migrated from GLOBAL_LOCAL_DISK_MAP/ID_MAP/SET_DRIVES)
+108 -26
View File
@@ -18,6 +18,7 @@ use crate::global::{
GLOBAL_EventNotifier, GLOBAL_LOCAL_DISK_ID_MAP, GLOBAL_LOCAL_DISK_MAP, GLOBAL_LOCAL_DISK_SET_DRIVES, GLOBAL_TierConfigMgr,
get_global_bucket_monitor, is_dist_erasure, is_first_cluster_node_local,
};
use crate::pools::local_decommission_queue_prefix;
use tracing::{debug, error, info, warn};
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
@@ -54,6 +55,22 @@ fn should_retry_local_decommission_resume(err: &Error, attempt: usize) -> bool {
matches!(err, Error::ConfigNotFound) && attempt < LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES
}
fn should_auto_start_rebalance_after_init(decommission_running: bool, rebalance_meta_loaded: bool) -> bool {
rebalance_meta_loaded && !decommission_running
}
fn pool_meta_has_active_decommission(meta: &PoolMeta) -> bool {
meta.pools.iter().any(|pool| {
pool.decommission
.as_ref()
.is_some_and(|info| !info.complete && !info.failed && !info.canceled)
})
}
fn should_auto_start_rebalance_after_recovered_meta(pool_meta: &PoolMeta, rebalance_meta_loaded: bool) -> bool {
should_auto_start_rebalance_after_init(pool_meta_has_active_decommission(pool_meta), rebalance_meta_loaded)
}
async fn wait_for_local_decommission_resume_delay(rx: &CancellationToken, delay: Duration) -> bool {
tokio::select! {
_ = rx.cancelled() => false,
@@ -305,6 +322,8 @@ impl ECStore {
pool_meta: RwLock::new(pool_meta),
rebalance_meta: RwLock::new(None),
decommission_cancelers,
start_gate: tokio::sync::Mutex::new(()),
pool_meta_save_gate: tokio::sync::Mutex::new(()),
local_disk_map: GLOBAL_LOCAL_DISK_MAP.clone(),
local_disk_id_map: GLOBAL_LOCAL_DISK_ID_MAP.clone(),
@@ -354,11 +373,6 @@ impl ECStore {
pub async fn init(self: &Arc<Self>, rx: CancellationToken) -> Result<()> {
GLOBAL_BOOT_TIME.get_or_init(|| async { SystemTime::now() }).await;
resolve_store_init_stage_result(self.load_rebalance_meta().await, "load_rebalance_meta")?;
if self.rebalance_meta.read().await.is_some() {
resolve_store_init_stage_result(self.start_rebalance().await, "start_rebalance")?;
}
let mut meta = PoolMeta::default();
resolve_store_init_stage_result(
meta.load(
@@ -375,11 +389,8 @@ impl ECStore {
let endpoints = get_global_endpoints();
let should_persist_pool_meta = is_first_cluster_node_local().await;
if !update {
{
let mut pool_meta = self.pool_meta.write().await;
*pool_meta = meta.clone();
}
let installed_pool_meta = if !update {
meta.clone()
} else {
let new_meta = PoolMeta::new(&self.pools, &meta);
// Only one local node should persist validated pool metadata here; otherwise
@@ -387,13 +398,32 @@ impl ECStore {
if should_persist_pool_meta {
resolve_store_init_stage_result(new_meta.save(self.pools.clone()).await, "save_validated_pool_meta")?;
}
{
let mut pool_meta = self.pool_meta.write().await;
*pool_meta = new_meta;
}
new_meta
};
{
let mut pool_meta = self.pool_meta.write().await;
*pool_meta = installed_pool_meta.clone();
}
let pools = meta.return_resumable_pools();
resolve_store_init_stage_result(self.load_rebalance_meta().await, "load_rebalance_meta")?;
let rebalance_meta_loaded = self.rebalance_meta.read().await.is_some();
let decommission_running =
pool_meta_has_active_decommission(&installed_pool_meta) || self.is_decommission_running().await;
if should_auto_start_rebalance_after_init(decommission_running, rebalance_meta_loaded) {
resolve_store_init_stage_result(self.start_rebalance().await, "start_rebalance")?;
} else if decommission_running && rebalance_meta_loaded {
warn!(
event = EVENT_ECSTORE_INIT_STATUS,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_STORE_INIT,
stage = "start_rebalance",
reason = "active_decommission",
"Deferred rebalance auto-start during store init because decommission is active"
);
}
let pools = installed_pool_meta.return_resumable_pools();
let mut pool_indices = Vec::with_capacity(pools.len());
for p in pools.iter() {
@@ -407,18 +437,16 @@ impl ECStore {
}
}
if !pool_indices.is_empty() {
let idx = pool_indices[0];
if should_resume_local_decommission(&endpoints, idx)? {
let store = self.clone();
let local_pool_indices = local_decommission_queue_prefix(&endpoints, &pool_indices)?;
if !local_pool_indices.is_empty() {
let store = self.clone();
tokio::spawn(async move {
if !wait_for_local_decommission_resume_delay(&rx, LOCAL_DECOMMISSION_INITIAL_RESUME_DELAY).await {
return;
}
resume_local_decommission_after_init(store, rx, pool_indices).await;
});
}
tokio::spawn(async move {
if !wait_for_local_decommission_resume_delay(&rx, LOCAL_DECOMMISSION_INITIAL_RESUME_DELAY).await {
return;
}
resume_local_decommission_after_init(store, rx, local_pool_indices).await;
});
}
let num_nodes = get_global_endpoints().get_nodes().len() as u64;
@@ -448,16 +476,33 @@ impl ECStore {
mod tests {
use super::{
LOCAL_DECOMMISSION_RESUME_MAX_CONFIG_RETRIES, pool_first_endpoint_is_local, resolve_store_init_stage_result,
should_auto_start_rebalance_after_init, should_auto_start_rebalance_after_recovered_meta,
should_resume_local_decommission, should_retry_local_decommission_resume, wait_for_local_decommission_resume_delay,
};
use crate::{
disk::endpoint::Endpoint,
endpoints::{EndpointServerPools, Endpoints, PoolEndpoints},
error::StorageError,
pools::{POOL_META_VERSION, PoolDecommissionInfo, PoolMeta, PoolStatus},
rebalance::RebalanceMeta,
};
use std::time::Duration;
use time::OffsetDateTime;
use tokio_util::sync::CancellationToken;
fn init_test_pool_meta(decommission: Option<PoolDecommissionInfo>) -> PoolMeta {
PoolMeta {
version: POOL_META_VERSION,
pools: vec![PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: OffsetDateTime::UNIX_EPOCH,
decommission,
}],
dont_save: false,
}
}
#[test]
fn test_should_resume_local_decommission_respects_local_flag() {
let mut local_endpoint = Endpoint::try_from("http://127.0.0.1:9000/data").expect("endpoint should parse");
@@ -519,6 +564,43 @@ mod tests {
assert!(!should_retry_local_decommission_resume(&StorageError::SlowDown, 0));
}
#[test]
fn test_should_auto_start_rebalance_after_init_allows_loaded_rebalance_without_decommission() {
assert!(should_auto_start_rebalance_after_init(false, true));
}
#[test]
fn test_should_auto_start_rebalance_after_init_rejects_active_decommission() {
assert!(!should_auto_start_rebalance_after_init(true, true));
}
#[test]
fn test_should_auto_start_rebalance_after_init_rejects_missing_rebalance_meta() {
assert!(!should_auto_start_rebalance_after_init(false, false));
}
#[test]
fn test_store_init_recovery_skips_rebalance_when_decommission_metadata_is_active() {
let pool_meta = init_test_pool_meta(Some(PoolDecommissionInfo {
start_time: Some(OffsetDateTime::UNIX_EPOCH),
complete: false,
failed: false,
canceled: false,
..Default::default()
}));
let rebalance_meta = Some(RebalanceMeta::default());
assert!(!should_auto_start_rebalance_after_recovered_meta(&pool_meta, rebalance_meta.is_some()));
}
#[test]
fn test_store_init_recovery_allows_rebalance_when_only_rebalance_metadata_exists() {
let pool_meta = init_test_pool_meta(None);
let rebalance_meta = Some(RebalanceMeta::default());
assert!(should_auto_start_rebalance_after_recovered_meta(&pool_meta, rebalance_meta.is_some()));
}
#[test]
fn test_resolve_store_init_stage_result_passthrough_ok() {
resolve_store_init_stage_result(Ok(()), "load_rebalance_meta").expect("successful stage should pass through");
+212 -116
View File
@@ -216,6 +216,10 @@ fn resolve_latest_object_access(
Ok((info, idx))
}
fn should_create_delete_marker_for_missing_object(opts: &ObjectOptions) -> bool {
opts.versioned && opts.version_id.is_none() && !opts.delete_marker && !opts.data_movement
}
fn version_aware_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions {
let mut lookup_opts = opts.clone();
lookup_opts.no_lock = no_lock;
@@ -234,6 +238,12 @@ fn data_movement_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> Object
lookup_opts
}
fn transition_restore_pool_opts(opts: &ObjectOptions) -> ObjectOptions {
let mut lookup_opts = opts.clone();
lookup_opts.skip_decommissioned = true;
lookup_opts
}
fn effective_object_actual_size(info: &ObjectInfo) -> Option<i64> {
info.get_actual_size().ok()
}
@@ -243,27 +253,44 @@ fn is_equivalent_data_movement_delete_marker(source: &ObjectInfo, target: &Objec
&& is_data_movement_delete_marker(target)
&& source.version_id == target.version_id
&& source.mod_time == target.mod_time
&& source.user_defined == target.user_defined
&& source.user_tags == target.user_tags
&& source.replication_status_internal == target.replication_status_internal
&& source.replication_status == target.replication_status
&& source.version_purge_status_internal == target.version_purge_status_internal
&& source.version_purge_status == target.version_purge_status
}
fn is_data_movement_delete_marker(info: &ObjectInfo) -> bool {
info.delete_marker
}
fn expected_data_movement_tiered_object(source: &rustfs_filemeta::FileInfo) -> ObjectInfo {
ObjectInfo::from_file_info(source, "", &source.name, source.version_id.is_some())
}
fn is_equivalent_data_movement_tiered_object(source: &rustfs_filemeta::FileInfo, target: &ObjectInfo) -> bool {
let expected = expected_data_movement_tiered_object(source);
source.version_id == target.version_id
&& !target.delete_marker
&& source.size == target.size
&& source.get_etag() == target.etag
&& source.checksum == target.checksum
&& source.mod_time == target.mod_time
&& source.transition_status == target.transitioned_object.status
&& source.transitioned_objname == target.transitioned_object.name
&& source.transition_tier == target.transitioned_object.tier
&& source
.transition_version_id
.map(|version_id| version_id.to_string())
.unwrap_or_default()
== target.transitioned_object.version_id
&& expected.user_defined == target.user_defined
&& expected.user_tags == target.user_tags
&& expected.expires == target.expires
&& expected.storage_class == target.storage_class
&& expected.replication_status_internal == target.replication_status_internal
&& expected.replication_status == target.replication_status
&& expected.version_purge_status_internal == target.version_purge_status_internal
&& expected.version_purge_status == target.version_purge_status
&& expected.transitioned_object.status == target.transitioned_object.status
&& expected.transitioned_object.name == target.transitioned_object.name
&& expected.transitioned_object.tier == target.transitioned_object.tier
&& expected.transitioned_object.version_id == target.transitioned_object.version_id
&& expected.transitioned_object.free_version == target.transitioned_object.free_version
&& effective_object_actual_size(target) == Some(source.size)
}
@@ -836,7 +863,7 @@ impl ECStore {
let target_pool_idx =
resolve_data_movement_resume_target_pool(selected_target_pool_idx, resume_target_pool_idx, opts.src_pool_idx);
if opts.src_pool_idx == selected_target_pool_idx {
if !should_check_data_movement_resume_target(opts.src_pool_idx, target_pool_idx) {
if let Ok((source_pool_info, _)) = existing_pool_info
&& opts.delete_marker
&& is_data_movement_delete_marker(&source_pool_info.object_info)
@@ -862,24 +889,23 @@ impl ECStore {
));
}
let mut obj = self.pools[selected_target_pool_idx]
.delete_object(bucket, object, opts)
.await?;
let mut obj = self.pools[target_pool_idx].delete_object(bucket, object, opts).await?;
obj.name = decode_dir_object(obj.name.as_str());
return Ok(obj);
}
// Determine which pool contains it
let (mut pinfo, errs) = self
.get_pool_info_existing_with_opts(bucket, object, &gopts)
.await
.map_err(|e| {
if is_err_read_quorum(&e) {
StorageError::ErasureWriteQuorum
} else {
e
}
})?;
let (mut pinfo, errs) = match self.get_pool_info_existing_with_opts(bucket, object, &gopts).await {
Ok(res) => res,
Err(err) if is_err_read_quorum(&err) => return Err(StorageError::ErasureWriteQuorum),
Err(err) if is_err_object_not_found(&err) && should_create_delete_marker_for_missing_object(&opts) => {
let target_pool_idx = self.get_pool_idx_no_lock(bucket, object, 0).await?;
let mut obj = self.pools[target_pool_idx].delete_object(bucket, object, opts).await?;
obj.name = decode_dir_object(object);
return Ok(obj);
}
Err(err) => return Err(err),
};
if pinfo.object_info.delete_marker && opts.version_id.is_none() {
pinfo.object_info.name = decode_dir_object(object);
@@ -1133,11 +1159,12 @@ impl ECStore {
return self.pools[0].transition_object(bucket, &object, opts).await;
}
//opts.skip_decommissioned = true;
//opts.no_lock = true;
let (_, idx) = self.get_latest_accessible_object_info_with_idx(bucket, &object, opts).await?;
let opts = transition_restore_pool_opts(opts);
let (_, idx) = self
.get_latest_accessible_object_info_with_idx(bucket, &object, &opts)
.await?;
self.pools[idx].transition_object(bucket, &object, opts).await
self.pools[idx].transition_object(bucket, &object, &opts).await
}
#[instrument(skip(self))]
@@ -1152,15 +1179,14 @@ impl ECStore {
return self.pools[0].clone().restore_transitioned_object(bucket, &object, opts).await;
}
//opts.skip_decommissioned = true;
//opts.nolock = true;
let opts = transition_restore_pool_opts(opts);
let (_, idx) = self
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), opts)
.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), &opts)
.await?;
self.pools[idx]
.clone()
.restore_transitioned_object(bucket, &object, opts)
.restore_transitioned_object(bucket, &object, &opts)
.await
}
@@ -1269,8 +1295,8 @@ mod tests {
use super::*;
use crate::bucket::lifecycle::core::TRANSITION_COMPLETE;
use bytes::Bytes;
use rustfs_storage_api::TransitionedObject;
use std::io::Cursor;
use std::sync::Arc;
use tokio::io::AsyncReadExt;
#[test]
@@ -1313,6 +1339,33 @@ mod tests {
assert!(!is_equivalent_data_movement_delete_marker(&source, &mismatched));
}
#[test]
fn equivalent_data_movement_delete_marker_rejects_metadata_and_replication_mismatch() {
let version_id = Uuid::nil();
let mod_time = OffsetDateTime::UNIX_EPOCH;
let source = ObjectInfo {
version_id: Some(version_id),
delete_marker: true,
mod_time: Some(mod_time),
user_defined: Arc::new(HashMap::from([("x-amz-meta-source".to_string(), "true".to_string())])),
replication_status_internal: Some("arn:minio:replication:target=COMPLETED;".to_string()),
version_purge_status_internal: Some("arn:minio:replication:target=PENDING;".to_string()),
..Default::default()
};
let mut target = source.clone();
target.user_defined = Arc::new(HashMap::from([("x-amz-meta-source".to_string(), "false".to_string())]));
assert!(!is_equivalent_data_movement_delete_marker(&source, &target));
let mut target = source.clone();
target.replication_status_internal = Some("arn:minio:replication:target=FAILED;".to_string());
assert!(!is_equivalent_data_movement_delete_marker(&source, &target));
let mut target = source.clone();
target.version_purge_status_internal = Some("arn:minio:replication:target=COMPLETE;".to_string());
assert!(!is_equivalent_data_movement_delete_marker(&source, &target));
}
#[test]
fn equivalent_data_movement_delete_marker_rejects_live_object() {
let source = ObjectInfo {
@@ -1367,6 +1420,7 @@ mod tests {
fn data_movement_resume_target_uses_resolved_non_source_pool_when_selected_is_source() {
let target_pool_idx = resolve_data_movement_resume_target_pool(1, Some(3), 1);
assert_eq!(target_pool_idx, 3);
assert!(should_check_data_movement_resume_target(1, target_pool_idx));
}
#[test]
@@ -1386,12 +1440,12 @@ mod tests {
assert!(matches!(result, Err(Error::SlowDown)));
}
#[test]
fn equivalent_data_movement_tiered_object_accepts_matching_transition_metadata() {
fn tiered_equivalence_source() -> FileInfo {
let version_id = Uuid::nil();
let transition_version_id = Uuid::new_v4();
let mod_time = OffsetDateTime::UNIX_EPOCH;
let source = FileInfo {
FileInfo {
version_id: Some(version_id),
size: 1024,
mod_time: Some(mod_time),
@@ -1400,80 +1454,86 @@ mod tests {
transitioned_objname: "remote/object".to_string(),
transition_tier: "WARM".to_string(),
transition_version_id: Some(transition_version_id),
metadata: HashMap::from([("etag".to_string(), "etag-value".to_string())]),
..Default::default()
};
let target = ObjectInfo {
version_id: Some(version_id),
size: 1024,
mod_time: Some(mod_time),
checksum: Some(Bytes::from_static(b"checksum")),
etag: Some("etag-value".to_string()),
transitioned_object: TransitionedObject {
name: "remote/object".to_string(),
version_id: transition_version_id.to_string(),
tier: "WARM".to_string(),
status: TRANSITION_COMPLETE.to_string(),
replication_state_internal: Some(rustfs_filemeta::ReplicationState {
replication_status_internal: Some("arn:minio:replication:target=COMPLETED;".to_string()),
targets: rustfs_filemeta::replication_statuses_map("arn:minio:replication:target=COMPLETED;"),
version_purge_status_internal: Some("arn:minio:replication:target=PENDING;".to_string()),
purge_targets: rustfs_filemeta::version_purge_statuses_map("arn:minio:replication:target=PENDING;"),
..Default::default()
},
}),
metadata: HashMap::from([
("etag".to_string(), "etag-value".to_string()),
("x-amz-meta-key".to_string(), "metadata-value".to_string()),
(rustfs_utils::http::AMZ_OBJECT_TAGGING.to_string(), "tag=value".to_string()),
("expires".to_string(), "1970-01-01T00:33:20Z".to_string()),
]),
..Default::default()
};
}
}
fn tiered_equivalence_target(source: &FileInfo) -> ObjectInfo {
ObjectInfo::from_file_info(source, "bucket", "object", source.version_id.is_some())
}
#[test]
fn equivalent_data_movement_tiered_object_accepts_matching_persisted_metadata() {
let source = tiered_equivalence_source();
let target = tiered_equivalence_target(&source);
assert!(is_equivalent_data_movement_tiered_object(&source, &target));
}
#[test]
fn equivalent_data_movement_tiered_object_rejects_transition_mismatch() {
let source = FileInfo {
version_id: Some(Uuid::nil()),
size: 1024,
transition_status: TRANSITION_COMPLETE.to_string(),
transitioned_objname: "remote/source".to_string(),
transition_tier: "WARM".to_string(),
..Default::default()
};
let target = ObjectInfo {
version_id: source.version_id,
size: 1024,
transitioned_object: TransitionedObject {
name: "remote/target".to_string(),
tier: "WARM".to_string(),
status: TRANSITION_COMPLETE.to_string(),
..Default::default()
},
..Default::default()
};
let source = tiered_equivalence_source();
let mut target = tiered_equivalence_target(&source);
target.transitioned_object.name = "remote/target".to_string();
assert!(!is_equivalent_data_movement_tiered_object(&source, &target));
}
#[test]
fn equivalent_data_movement_tiered_object_rejects_user_metadata_mismatch() {
let source = tiered_equivalence_source();
let mut target = tiered_equivalence_target(&source);
target.user_defined = Arc::new(HashMap::from([("x-amz-meta-key".to_string(), "target-value".to_string())]));
assert!(!is_equivalent_data_movement_tiered_object(&source, &target));
}
#[test]
fn equivalent_data_movement_tiered_object_rejects_tag_mismatch() {
let source = tiered_equivalence_source();
let mut target = tiered_equivalence_target(&source);
target.user_tags = Arc::new("tag=target".to_string());
assert!(!is_equivalent_data_movement_tiered_object(&source, &target));
}
#[test]
fn equivalent_data_movement_tiered_object_rejects_replication_mismatch() {
let source = tiered_equivalence_source();
let mut target = tiered_equivalence_target(&source);
target.replication_status_internal = Some("arn:minio:replication:target=FAILED;".to_string());
target.replication_status = rustfs_filemeta::ReplicationStatusType::Failed;
assert!(!is_equivalent_data_movement_tiered_object(&source, &target));
}
#[test]
fn equivalent_data_movement_tiered_object_rejects_version_purge_mismatch() {
let source = tiered_equivalence_source();
let mut target = tiered_equivalence_target(&source);
target.version_purge_status_internal = Some("arn:minio:replication:target=COMPLETE;".to_string());
target.version_purge_status = rustfs_filemeta::VersionPurgeStatusType::Complete;
assert!(!is_equivalent_data_movement_tiered_object(&source, &target));
}
#[test]
fn data_movement_tiered_resume_accepts_equivalent_target() {
let version_id = Uuid::nil();
let transition_version_id = Uuid::new_v4();
let source = FileInfo {
version_id: Some(version_id),
size: 1024,
transition_status: TRANSITION_COMPLETE.to_string(),
transitioned_objname: "remote/object".to_string(),
transition_tier: "WARM".to_string(),
transition_version_id: Some(transition_version_id),
metadata: HashMap::from([("etag".to_string(), "etag-value".to_string())]),
..Default::default()
};
let target = ObjectInfo {
version_id: Some(version_id),
size: 1024,
etag: Some("etag-value".to_string()),
transitioned_object: TransitionedObject {
name: "remote/object".to_string(),
version_id: transition_version_id.to_string(),
tier: "WARM".to_string(),
status: TRANSITION_COMPLETE.to_string(),
..Default::default()
},
..Default::default()
};
let source = tiered_equivalence_source();
let target = tiered_equivalence_target(&source);
let should_resume = resolve_data_movement_tiered_resume_result(Ok(Some(target)), &source, 0, 1)
.expect("equivalent tiered target should be evaluated");
@@ -1483,28 +1543,8 @@ mod tests {
#[test]
fn data_movement_tiered_resume_rejects_source_pool_target() {
let version_id = Uuid::nil();
let source = FileInfo {
version_id: Some(version_id),
size: 1024,
transition_status: TRANSITION_COMPLETE.to_string(),
transitioned_objname: "remote/object".to_string(),
transition_tier: "WARM".to_string(),
metadata: HashMap::from([("etag".to_string(), "etag-value".to_string())]),
..Default::default()
};
let target = ObjectInfo {
version_id: Some(version_id),
size: 1024,
etag: Some("etag-value".to_string()),
transitioned_object: TransitionedObject {
name: "remote/object".to_string(),
tier: "WARM".to_string(),
status: TRANSITION_COMPLETE.to_string(),
..Default::default()
},
..Default::default()
};
let source = tiered_equivalence_source();
let target = tiered_equivalence_target(&source);
let should_resume = resolve_data_movement_tiered_resume_result(Ok(Some(target)), &source, 0, 0)
.expect("source-pool target should be rejected before target lookup");
@@ -1623,6 +1663,39 @@ mod tests {
assert!(matches!(err, Error::MethodNotAllowed));
}
#[test]
fn should_create_delete_marker_for_missing_object_allows_latest_versioned_delete() {
let opts = ObjectOptions {
versioned: true,
..Default::default()
};
assert!(should_create_delete_marker_for_missing_object(&opts));
}
#[test]
fn should_create_delete_marker_for_missing_object_rejects_specialized_deletes() {
let version_delete = ObjectOptions {
versioned: true,
version_id: Some("vid-1".to_string()),
..Default::default()
};
let delete_marker_replication = ObjectOptions {
versioned: true,
delete_marker: true,
..Default::default()
};
let data_movement = ObjectOptions {
versioned: true,
data_movement: true,
..Default::default()
};
assert!(!should_create_delete_marker_for_missing_object(&version_delete));
assert!(!should_create_delete_marker_for_missing_object(&delete_marker_replication));
assert!(!should_create_delete_marker_for_missing_object(&data_movement));
}
#[test]
fn resolve_decommission_target_pool_idx_result_passthrough_ok() {
let idx = ECStore::resolve_decommission_target_pool_idx_result(Ok(3), "bucket", "object").unwrap();
@@ -1715,6 +1788,29 @@ mod tests {
assert!(lookup_opts.skip_rebalancing);
}
#[test]
fn transition_restore_pool_opts_skips_decommissioned_and_preserves_locking() {
let lookup_opts = transition_restore_pool_opts(&ObjectOptions {
no_lock: false,
skip_decommissioned: false,
..Default::default()
});
assert!(lookup_opts.skip_decommissioned);
assert!(!lookup_opts.no_lock);
}
#[test]
fn transition_restore_pool_opts_preserves_existing_no_lock() {
let lookup_opts = transition_restore_pool_opts(&ObjectOptions {
no_lock: true,
..Default::default()
});
assert!(lookup_opts.skip_decommissioned);
assert!(lookup_opts.no_lock);
}
#[tokio::test]
#[serial_test::serial]
async fn reader_lock_is_held_when_optimization_is_disabled() {
@@ -1080,8 +1080,11 @@ pub struct ReloadPoolMetaResponse {
#[prost(string, optional, tag = "2")]
pub error_info: ::core::option::Option<::prost::alloc::string::String>,
}
#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)]
pub struct StopRebalanceRequest {}
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct StopRebalanceRequest {
#[prost(string, tag = "1")]
pub expected_rebalance_id: ::prost::alloc::string::String,
}
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct StopRebalanceResponse {
#[prost(bool, tag = "1")]
+3 -1
View File
@@ -764,7 +764,9 @@ message ReloadPoolMetaResponse {
optional string error_info = 2;
}
message StopRebalanceRequest {}
message StopRebalanceRequest {
string expected_rebalance_id = 1;
}
message StopRebalanceResponse {
bool success = 1;