Compare commits

..

2 Commits

Author SHA1 Message Date
overtrue 7d6f65734e fix(ecstore): clear clippy warnings 2026-08-23 06:04:29 +08:00
overtrue a560f66cea fix(ecstore): run decommission metadata first 2026-08-23 03:45:38 +08:00
16 changed files with 243 additions and 671 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=af59a0519bfb651bfb6cb57105d7e4eb68bca13cc24887c945908cbf4cb102ff
sha256-linux=5c2d3aeddaaa48e7807039a72e55dad427f971955c8251feec6d0bb66ac25181
sha256-darwin=9f767b37ed8b1c82da62ea441462d75487785c8086e56f08fb6f6cd89c6e2e52
sha256-linux=fbdaf42b220958d4b1e8880e0f8b5a7992d38e21051bb60596dd4538424757d6
Generated
-1
View File
@@ -12688,7 +12688,6 @@ dependencies = [
"js-sys",
"rand 0.10.2",
"serde_core",
"sha1_smol",
"wasm-bindgen",
]
+82 -67
View File
@@ -13,37 +13,55 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::common::{RustFSTestEnvironment, init_logging};
use aws_config::meta::region::RegionProviderChain;
use aws_sdk_s3::Client;
use aws_sdk_s3::error::ProvideErrorMetadata;
use aws_sdk_s3::config::{Credentials, Region};
use aws_sdk_s3::types::{
CsvInput, CsvOutput, ExpressionType, FileHeaderInfo, InputSerialization, JsonInput, JsonOutput, JsonType, OutputSerialization,
};
use bytes::Bytes;
use std::error::Error;
use std::time::Duration;
const ENDPOINT: &str = "http://localhost:9000";
const ACCESS_KEY: &str = "rustfsadmin";
const SECRET_KEY: &str = "rustfsadmin";
const BUCKET: &str = "test-sql-bucket";
const CSV_OBJECT: &str = "test-data.csv";
const JSON_OBJECT: &str = "test-data.json";
const SELECT_RESPONSE_TIMEOUT: Duration = Duration::from_secs(30);
type TestResult<T> = Result<T, Box<dyn Error + Send + Sync>>;
async fn create_aws_s3_client() -> Result<Client, Box<dyn Error>> {
let region_provider = RegionProviderChain::default_provider().or_else(Region::new("us-east-1"));
let shared_config = aws_config::defaults(aws_config::BehaviorVersion::latest())
.region(region_provider)
.credentials_provider(Credentials::new(ACCESS_KEY, SECRET_KEY, None, None, "static"))
.endpoint_url(ENDPOINT)
.load()
.await;
async fn create_test_environment() -> TestResult<(RustFSTestEnvironment, Client)> {
init_logging();
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let client = env.create_s3_client();
Ok((env, client))
let client = Client::from_conf(
aws_sdk_s3::Config::from(&shared_config)
.to_builder()
.force_path_style(true) // Important for S3-compatible services
.build(),
);
Ok(client)
}
async fn setup_test_bucket(client: &Client) -> TestResult<()> {
client.create_bucket().bucket(BUCKET).send().await?;
async fn setup_test_bucket(client: &Client) -> Result<(), Box<dyn Error>> {
match client.create_bucket().bucket(BUCKET).send().await {
Ok(_) => {}
Err(e) => {
let error_str = e.to_string();
if !error_str.contains("BucketAlreadyOwnedByYou") && !error_str.contains("BucketAlreadyExists") {
return Err(e.into());
}
}
}
Ok(())
}
async fn upload_test_csv(client: &Client) -> TestResult<()> {
async fn upload_test_csv(client: &Client) -> Result<(), Box<dyn Error>> {
let csv_data = "name,age,city\nAlice,30,New York\nBob,25,Los Angeles\nCharlie,35,Chicago\nDiana,28,Boston";
client
@@ -57,7 +75,7 @@ async fn upload_test_csv(client: &Client) -> TestResult<()> {
Ok(())
}
async fn upload_test_json(client: &Client) -> TestResult<()> {
async fn upload_test_json(client: &Client) -> Result<(), Box<dyn Error>> {
let json_data = r#"{"name":"Alice","age":30,"city":"New York"}
{"name":"Bob","age":25,"city":"Los Angeles"}
{"name":"Charlie","age":35,"city":"Chicago"}
@@ -75,38 +93,33 @@ async fn upload_test_json(client: &Client) -> TestResult<()> {
async fn process_select_response(
mut event_stream: aws_sdk_s3::operation::select_object_content::SelectObjectContentOutput,
) -> TestResult<String> {
tokio::time::timeout(SELECT_RESPONSE_TIMEOUT, async move {
let mut total_data = Vec::new();
let mut saw_end = false;
) -> Result<String, Box<dyn Error>> {
let mut total_data = Vec::new();
while let Some(event) = event_stream.payload.recv().await? {
match event {
aws_sdk_s3::types::SelectObjectContentEventStream::Records(records_event) => {
if let Some(payload) = records_event.payload {
total_data.extend_from_slice(payload.as_ref());
}
while let Ok(Some(event)) = event_stream.payload.recv().await {
match event {
aws_sdk_s3::types::SelectObjectContentEventStream::Records(records_event) => {
if let Some(payload) = records_event.payload {
let data = payload.into_inner();
total_data.extend_from_slice(&data);
}
aws_sdk_s3::types::SelectObjectContentEventStream::End(_) => {
saw_end = true;
break;
}
_ => {}
}
aws_sdk_s3::types::SelectObjectContentEventStream::End(_) => {
break;
}
_ => {
// Handle other event types (Stats, Progress, Cont, etc.)
}
}
}
if !saw_end {
return Err("Select response ended without an End event".into());
}
Ok(String::from_utf8(total_data)?)
})
.await
.map_err(|_| -> Box<dyn Error + Send + Sync> { "Select response timed out".into() })?
Ok(String::from_utf8(total_data)?)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_select_object_content_csv_basic() -> TestResult<()> {
let (_env, client) = create_test_environment().await?;
#[ignore = "requires running RustFS server at localhost:9000"]
async fn test_select_object_content_csv_basic() -> Result<(), Box<dyn Error>> {
let client = create_aws_s3_client().await?;
setup_test_bucket(&client).await?;
upload_test_csv(&client).await?;
@@ -145,8 +158,9 @@ async fn test_select_object_content_csv_basic() -> TestResult<()> {
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_select_object_content_csv_aggregation() -> TestResult<()> {
let (_env, client) = create_test_environment().await?;
#[ignore = "requires running RustFS server at localhost:9000"]
async fn test_select_object_content_csv_aggregation() -> Result<(), Box<dyn Error>> {
let client = create_aws_s3_client().await?;
setup_test_bucket(&client).await?;
upload_test_csv(&client).await?;
@@ -189,15 +203,16 @@ async fn test_select_object_content_csv_aggregation() -> TestResult<()> {
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_select_object_content_json_basic() -> TestResult<()> {
let (_env, client) = create_test_environment().await?;
#[ignore = "requires running RustFS server at localhost:9000"]
async fn test_select_object_content_json_basic() -> Result<(), Box<dyn Error>> {
let client = create_aws_s3_client().await?;
setup_test_bucket(&client).await?;
upload_test_json(&client).await?;
// Construct JSON query
let sql = "SELECT s.name, s.age FROM S3Object s WHERE s.age > 28";
let json_input = JsonInput::builder().set_type(Some(JsonType::Lines)).build();
let json_input = JsonInput::builder().set_type(Some(JsonType::Document)).build();
let input_serialization = InputSerialization::builder().json(json_input).build();
@@ -229,8 +244,9 @@ async fn test_select_object_content_json_basic() -> TestResult<()> {
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_select_object_content_csv_limit() -> TestResult<()> {
let (_env, client) = create_test_environment().await?;
#[ignore = "requires running RustFS server at localhost:9000"]
async fn test_select_object_content_csv_limit() -> Result<(), Box<dyn Error>> {
let client = create_aws_s3_client().await?;
setup_test_bucket(&client).await?;
upload_test_csv(&client).await?;
@@ -270,8 +286,9 @@ async fn test_select_object_content_csv_limit() -> TestResult<()> {
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_select_object_content_csv_order_by() -> TestResult<()> {
let (_env, client) = create_test_environment().await?;
#[ignore = "requires running RustFS server at localhost:9000"]
async fn test_select_object_content_csv_order_by() -> Result<(), Box<dyn Error>> {
let client = create_aws_s3_client().await?;
setup_test_bucket(&client).await?;
upload_test_csv(&client).await?;
@@ -301,10 +318,9 @@ async fn test_select_object_content_csv_order_by() -> TestResult<()> {
println!("CSV Order By result: {result_str}");
// Verify ordered by age descending
assert_eq!(
result_str.lines().filter(|line| !line.trim().is_empty()).count(),
2,
"Should return exactly 2 records"
assert!(
result_str.lines().filter(|line| !line.trim().is_empty()).count() >= 2,
"Should return at least 2 records"
);
// Check if contains highest age records
@@ -315,8 +331,9 @@ async fn test_select_object_content_csv_order_by() -> TestResult<()> {
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_select_object_content_error_handling() -> TestResult<()> {
let (_env, client) = create_test_environment().await?;
#[ignore = "requires running RustFS server at localhost:9000"]
async fn test_select_object_content_error_handling() -> Result<(), Box<dyn Error>> {
let client = create_aws_s3_client().await?;
setup_test_bucket(&client).await?;
upload_test_csv(&client).await?;
@@ -331,7 +348,7 @@ async fn test_select_object_content_error_handling() -> TestResult<()> {
let output_serialization = OutputSerialization::builder().csv(csv_output).build();
// This query should fail because invalid_column doesn't exist
let error = client
let result = client
.select_object_content()
.bucket(BUCKET)
.key(CSV_OBJECT)
@@ -340,20 +357,18 @@ async fn test_select_object_content_error_handling() -> TestResult<()> {
.input_serialization(input_serialization)
.output_serialization(output_serialization)
.send()
.await
.expect_err("a query referencing an unknown column must fail");
.await;
assert_eq!(
error.as_service_error().and_then(ProvideErrorMetadata::code),
Some("EvaluatorBindingDoesNotExist")
);
// Verify query fails (expected behavior)
assert!(result.is_err(), "Query with invalid column should fail");
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_select_object_content_nonexistent_object() -> TestResult<()> {
let (_env, client) = create_test_environment().await?;
#[ignore = "requires running RustFS server at localhost:9000"]
async fn test_select_object_content_nonexistent_object() -> Result<(), Box<dyn Error>> {
let client = create_aws_s3_client().await?;
setup_test_bucket(&client).await?;
// Test query on nonexistent object
@@ -366,7 +381,7 @@ async fn test_select_object_content_nonexistent_object() -> TestResult<()> {
let csv_output = CsvOutput::builder().build();
let output_serialization = OutputSerialization::builder().csv(csv_output).build();
let error = client
let result = client
.select_object_content()
.bucket(BUCKET)
.key("nonexistent.csv")
@@ -375,10 +390,10 @@ async fn test_select_object_content_nonexistent_object() -> TestResult<()> {
.input_serialization(input_serialization)
.output_serialization(output_serialization)
.send()
.await
.expect_err("selecting a missing object must fail");
.await;
assert_eq!(error.as_service_error().and_then(ProvideErrorMetadata::code), Some("NoSuchKey"));
// Verify query fails (expected behavior)
assert!(result.is_err(), "Query on nonexistent object should fail");
Ok(())
}
+102 -47
View File
@@ -1455,6 +1455,36 @@ where
Ok(())
}
async fn run_decommission_phases<F>(
rx: CancellationToken,
regular_buckets: Vec<DecomBucketInfo>,
meta_buckets: Vec<DecomBucketInfo>,
bucket_concurrency: usize,
mut start_bucket: F,
) -> Result<()>
where
F: FnMut(DecomBucketInfo, CancellationToken) -> BoxFuture<'static, Result<()>>,
{
decommission_cancel_signal_result(rx.is_cancelled())?;
for bucket in meta_buckets {
decommission_cancel_signal_result(rx.is_cancelled())?;
start_bucket(bucket, rx.clone()).await?;
}
decommission_cancel_signal_result(rx.is_cancelled())?;
if bucket_concurrency <= 1 {
for bucket in regular_buckets {
decommission_cancel_signal_result(rx.is_cancelled())?;
start_bucket(bucket, rx.clone()).await?;
}
return Ok(());
}
run_decommission_buckets_bounded(rx, regular_buckets, bucket_concurrency, start_bucket).await
}
#[cfg(test)]
async fn wait_decommission_worker_drain(workers: &Semaphore, limit: usize) -> Result<()> {
let permits = u32::try_from(limit)
@@ -2026,7 +2056,7 @@ impl PoolMeta {
self.load_no_lock(pool).await
}
pub(crate) async fn load_no_lock<S>(&mut self, pool: Arc<S>) -> Result<()>
async fn load_no_lock<S>(&mut self, pool: Arc<S>) -> Result<()>
where
S: EcstoreObjectIO,
{
@@ -4902,25 +4932,6 @@ impl ECStore {
Ok(())
}
async fn decommission_buckets_concurrently(
self: &Arc<Self>,
rx: CancellationToken,
idx: usize,
pool: Arc<Sets>,
buckets: Vec<DecomBucketInfo>,
limit: usize,
entry_budget: Arc<Semaphore>,
) -> Result<()> {
let store = Arc::clone(self);
run_decommission_buckets_bounded(rx, buckets, limit, move |bucket, rx| {
let store = Arc::clone(&store);
let pool = pool.clone();
let entry_budget = entry_budget.clone();
Box::pin(async move { store.decommission_pending_bucket(rx, idx, pool, bucket, entry_budget).await })
})
.await
}
#[tracing::instrument(skip(self, rx))]
async fn decommission_in_background(
self: &Arc<Self>,
@@ -4935,31 +4946,15 @@ impl ECStore {
pool_meta.pending_buckets(idx)
};
let bucket_concurrency = decommission_bucket_concurrency_limit();
if bucket_concurrency <= 1 {
for bucket in pending {
self.decommission_pending_bucket(rx.clone(), idx, pool.clone(), bucket, entry_budget.clone())
.await?;
}
return Ok(());
}
let (regular_buckets, meta_buckets) = split_decommission_buckets(pending);
self.decommission_buckets_concurrently(
rx.clone(),
idx,
pool.clone(),
regular_buckets,
bucket_concurrency,
entry_budget.clone(),
)
.await?;
for bucket in meta_buckets {
self.decommission_pending_bucket(rx.clone(), idx, pool.clone(), bucket, entry_budget.clone())
.await?;
}
Ok(())
let store = Arc::clone(self);
run_decommission_phases(rx, regular_buckets, meta_buckets, bucket_concurrency, move |bucket, rx| {
let store = Arc::clone(&store);
let pool = pool.clone();
let entry_budget = entry_budget.clone();
Box::pin(async move { store.decommission_pending_bucket(rx, idx, pool, bucket, entry_budget).await })
})
.await
}
#[tracing::instrument(skip(self))]
@@ -6377,8 +6372,8 @@ mod pools_tests {
resolve_decommission_terminal_mark_after_error_result, resolve_decommission_terminal_mark_result,
resolve_decommission_update_after_result, resolve_start_decommission_pool_meta_reload_result,
rollback_start_decommission_pool_meta, run_decommission_buckets_bounded, run_decommission_listing_with_retry,
run_decommission_listing_with_retry_and_drain, run_decommission_side_effect, should_cleanup_decommission_source_entry,
should_continue_decommission_queue, should_count_decommission_version_complete,
run_decommission_listing_with_retry_and_drain, run_decommission_phases, run_decommission_side_effect,
should_cleanup_decommission_source_entry, should_continue_decommission_queue, should_count_decommission_version_complete,
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,
spawn_decommission_index_cancelers, split_decommission_buckets, take_and_cancel_decommission_canceler,
@@ -6397,7 +6392,7 @@ mod pools_tests {
use rustfs_filemeta::{MetaCacheEntries, MetadataResolutionParams};
use rustfs_rio::Index;
use std::sync::{
Arc,
Arc, Mutex as StdMutex,
atomic::{AtomicBool, AtomicUsize, Ordering},
};
use std::time::Duration as StdDuration;
@@ -6690,6 +6685,66 @@ mod pools_tests {
);
}
#[tokio::test]
async fn test_decommission_metadata_phase_precedes_regular_failure() {
let events = Arc::new(StdMutex::new(Vec::new()));
let err = run_decommission_phases(
CancellationToken::new(),
vec![
DecomBucketInfo {
name: "regular-fails".to_string(),
..Default::default()
},
DecomBucketInfo {
name: "regular-not-started".to_string(),
..Default::default()
},
],
vec![
DecomBucketInfo {
name: crate::disk::RUSTFS_META_BUCKET.to_string(),
prefix: crate::config::com::CONFIG_PREFIX.to_string(),
},
DecomBucketInfo {
name: crate::disk::RUSTFS_META_BUCKET.to_string(),
prefix: crate::disk::BUCKET_META_PREFIX.to_string(),
},
],
1,
{
let events = Arc::clone(&events);
move |bucket, _rx| {
let events = Arc::clone(&events);
Box::pin(async move {
let event = if bucket.name == crate::disk::RUSTFS_META_BUCKET {
format!("meta:{}", bucket.prefix)
} else {
format!("regular:{}", bucket.name)
};
events.lock().expect("phase event lock should not be poisoned").push(event);
if bucket.name == "regular-fails" {
Err(Error::SlowDown)
} else {
Ok(())
}
})
}
},
)
.await
.expect_err("regular failure should remain fatal after metadata completes");
assert!(matches!(err, Error::SlowDown));
assert_eq!(
*events.lock().expect("phase event lock should not be poisoned"),
vec![
format!("meta:{}", crate::config::com::CONFIG_PREFIX),
format!("meta:{}", crate::disk::BUCKET_META_PREFIX),
"regular:regular-fails".to_string(),
]
);
}
#[tokio::test]
async fn test_run_decommission_buckets_bounded_respects_limit() {
let rx = CancellationToken::new();
+8 -20
View File
@@ -988,11 +988,14 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for Sets {
}
}
impl Sets {
pub(crate) async fn heal_format_with_fence<F>(&self, dry_run: bool, fence_lost: F) -> Result<(HealResultItem, Option<Error>)>
where
F: Fn() -> bool + Send + Sync,
{
#[async_trait::async_trait]
impl crate::storage_api_contracts::heal::HealOperations for Sets {
type Error = Error;
type HealResultItem = HealResultItem;
type HealOptions = HealOpts;
#[tracing::instrument(skip(self))]
async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
let (disks, init_errs) = init_storage_disks_with_errors(
&self.endpoints.endpoints,
&DiskOption {
@@ -1065,9 +1068,6 @@ impl Sets {
// Save new formats `format.json` on unformatted disks.
for (index, (fm, disk)) in tmp_new_formats.iter_mut().zip(disks.iter()).enumerate() {
if fm.is_some() && disk.is_some() {
if fence_lost() {
return Ok((res, Some(StorageError::SlowDown)));
}
if let Err(err) = save_format_file(disk, fm).await {
if let Some(disk) = disk.as_ref() {
let _ = disk.close().await;
@@ -1101,18 +1101,6 @@ impl Sets {
}
Ok((res, None))
}
}
#[async_trait::async_trait]
impl crate::storage_api_contracts::heal::HealOperations for Sets {
type Error = Error;
type HealResultItem = HealResultItem;
type HealOptions = HealOpts;
#[tracing::instrument(skip(self))]
async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
self.heal_format_with_fence(dry_run, || false).await
}
#[tracing::instrument(skip(self))]
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
let mut result = HealResultItem {
+3 -332
View File
@@ -13,12 +13,7 @@
// limitations under the License.
use super::*;
use crate::core::pools::POOL_META_NAME;
use crate::services::rebalance::{REBAL_META_NAME, RebalStatus};
use crate::set_disk::get_lock_acquire_timeout;
use crate::storage_api_contracts::heal::HealOperations as _;
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
use rustfs_lock::NamespaceLockGuard;
use tracing::trace;
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
@@ -35,119 +30,7 @@ fn invalid_heal_pool_index(pool_idx: usize, pool_count: usize) -> Error {
)
}
#[derive(Debug, Clone, Copy)]
enum HealFormatPoolSkip {
Completed,
Retryable,
}
fn classify_heal_format_pool(
pool_idx: usize,
pool_cmd_line: &str,
pool_meta: &PoolMeta,
rebalance_meta: Option<&RebalanceMeta>,
) -> Option<HealFormatPoolSkip> {
let Some(pool) = pool_meta.pools.get(pool_idx) else {
return Some(HealFormatPoolSkip::Retryable);
};
if pool.id != pool_idx || pool_cmd_line.is_empty() || pool.cmd_line.is_empty() || pool.cmd_line != pool_cmd_line {
return Some(HealFormatPoolSkip::Retryable);
}
if let Some(decommission) = pool.decommission.as_ref() {
if decommission.complete {
return Some(HealFormatPoolSkip::Completed);
}
if decommission.failed || decommission.canceled || decommission.queued || pool_meta.is_suspended(pool_idx) {
return Some(HealFormatPoolSkip::Retryable);
}
}
if let Some(meta) = rebalance_meta {
let Some(pool_stats) = meta.pool_stats.get(pool_idx) else {
return Some(HealFormatPoolSkip::Retryable);
};
if pool_stats.info.stopping || (pool_stats.participating && pool_stats.info.status == RebalStatus::Started) {
return Some(HealFormatPoolSkip::Retryable);
}
}
None
}
fn heal_format_pool_skip_error(skip: HealFormatPoolSkip) -> Error {
match skip {
HealFormatPoolSkip::Completed => StorageError::NoHealRequired,
HealFormatPoolSkip::Retryable => StorageError::SlowDown,
}
}
fn heal_format_fence_lost_error() -> Error {
StorageError::SlowDown
}
impl ECStore {
async fn acquire_heal_format_fence(
&self,
) -> Result<(NamespaceLockGuard, NamespaceLockGuard, PoolMeta, Option<RebalanceMeta>)> {
let metadata_pool = self
.pools
.first()
.cloned()
.ok_or_else(|| Error::other("heal format requires at least one storage pool"))?;
// Metadata fence order is part of the decommission/rebalance protocol:
// pool.bin must always be acquired before rebalance.bin.
let pool_lock = metadata_pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?;
let pool_guard = pool_lock.get_write_lock(get_lock_acquire_timeout()).await?;
let rebalance_lock = metadata_pool.new_ns_lock(RUSTFS_META_BUCKET, REBAL_META_NAME).await?;
let rebalance_guard = rebalance_lock.get_write_lock(get_lock_acquire_timeout()).await?;
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
return Err(heal_format_fence_lost_error());
}
let mut pool_meta = PoolMeta::default();
pool_meta.load_no_lock(metadata_pool.clone()).await?;
if pool_meta.pools.len() != self.pools.len()
|| pool_meta.pools.iter().enumerate().any(|(pool_idx, pool)| {
pool.id != pool_idx || pool.cmd_line.is_empty() || pool.cmd_line != self.pools[pool_idx].endpoints.cmd_line
})
{
return Err(heal_format_fence_lost_error());
}
let mut rebalance_meta = RebalanceMeta::new();
let rebalance_meta = match rebalance_meta
.load_with_opts(
metadata_pool,
ObjectOptions {
no_lock: true,
..Default::default()
},
)
.await
{
Ok(()) => Some(rebalance_meta),
Err(Error::ConfigNotFound) => None,
Err(err) => return Err(err),
};
if rebalance_meta
.as_ref()
.is_some_and(|meta| meta.pool_stats.len() != self.pools.len())
{
return Err(heal_format_fence_lost_error());
}
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
return Err(heal_format_fence_lost_error());
}
Ok((pool_guard, rebalance_guard, pool_meta, rebalance_meta))
}
fn get_pools_for_heal_object(&self, opts: &HealOpts) -> Result<Vec<Arc<Sets>>> {
match opts.pool {
Some(pool_idx) => Ok(vec![
@@ -169,26 +52,9 @@ impl ECStore {
};
let mut count_no_heal = 0;
let mut count_completed = 0;
let mut first_error = None;
for (pool_idx, pool) in self.pools.iter().enumerate() {
let (pool_guard, rebalance_guard, pool_meta, rebalance_meta) = self.acquire_heal_format_fence().await?;
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
first_error.get_or_insert(heal_format_fence_lost_error());
break;
}
if let Some(skip) = classify_heal_format_pool(pool_idx, &pool.endpoints.cmd_line, &pool_meta, rebalance_meta.as_ref())
{
if matches!(skip, HealFormatPoolSkip::Completed) {
count_completed += 1;
} else {
first_error.get_or_insert(heal_format_pool_skip_error(skip));
}
continue;
}
let fence_lost = || pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost();
let (mut result, err) = pool.heal_format_with_fence(dry_run, fence_lost).await?;
for pool in self.pools.iter() {
let (mut result, err) = pool.heal_format(dry_run).await?;
if let Some(err) = err {
match err {
StorageError::NoHealRequired => {
@@ -203,18 +69,11 @@ impl ECStore {
r.set_count += result.set_count;
r.before.drives.append(&mut result.before.drives);
r.after.drives.append(&mut result.after.drives);
// A lease can be lost after the final write; fail closed before
// reporting the pool as successfully healed.
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
first_error.get_or_insert(heal_format_fence_lost_error());
break;
}
}
if let Some(err) = first_error {
return Ok((r, Some(err)));
}
if count_no_heal + count_completed == self.pools.len() {
if count_no_heal == self.pools.len() {
info!(
event = EVENT_HEAL_FORMAT_COMPLETED,
component = LOG_COMPONENT_ECSTORE,
@@ -443,7 +302,6 @@ mod tests {
use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk};
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
use crate::runtime::instance::InstanceContext;
use crate::services::rebalance::{RebalanceInfo, RebalanceStats};
use crate::storage_api_contracts::bucket::{BucketOperations, MakeBucketOptions};
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations};
use crate::store::init_format::{load_format_erasure, save_format_file};
@@ -495,164 +353,6 @@ mod tests {
}
}
fn pool_meta_with_decommission(info: PoolDecommissionInfo) -> PoolMeta {
PoolMeta {
pools: vec![PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: OffsetDateTime::UNIX_EPOCH,
decommission: Some(info),
}],
..Default::default()
}
}
#[test]
fn heal_format_pool_state_barriers_are_classified() {
let active = pool_meta_with_decommission(PoolDecommissionInfo {
start_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
});
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &active, None),
Some(HealFormatPoolSkip::Retryable)
));
for info in [
PoolDecommissionInfo {
failed: true,
..Default::default()
},
PoolDecommissionInfo {
canceled: true,
..Default::default()
},
] {
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &pool_meta_with_decommission(info), None),
Some(HealFormatPoolSkip::Retryable)
));
}
let completed = pool_meta_with_decommission(PoolDecommissionInfo {
complete: true,
..Default::default()
});
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &completed, None),
Some(HealFormatPoolSkip::Completed)
));
}
#[test]
fn heal_format_pool_rebalance_barriers_and_identity_are_fail_closed() {
let identity_meta = pool_meta_with_decommission(PoolDecommissionInfo::default());
let rebalance = RebalanceMeta {
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&rebalance)),
Some(HealFormatPoolSkip::Retryable)
));
let stopping = RebalanceMeta {
pool_stats: vec![RebalanceStats {
info: RebalanceInfo {
stopping: true,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopping)),
Some(HealFormatPoolSkip::Retryable)
));
let identity = pool_meta_with_decommission(PoolDecommissionInfo::default());
assert!(matches!(
classify_heal_format_pool(0, "pool-new", &identity, None),
Some(HealFormatPoolSkip::Retryable)
));
let identity_without_decommission = PoolMeta {
pools: vec![PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: OffsetDateTime::UNIX_EPOCH,
decommission: None,
}],
..Default::default()
};
assert!(matches!(
classify_heal_format_pool(0, "pool-new", &identity_without_decommission, None),
Some(HealFormatPoolSkip::Retryable)
));
assert!(matches!(
classify_heal_format_pool(0, "", &identity_meta, None),
Some(HealFormatPoolSkip::Retryable)
));
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &PoolMeta::default(), None),
Some(HealFormatPoolSkip::Retryable)
));
let stopped = RebalanceMeta {
stopped_at: Some(OffsetDateTime::UNIX_EPOCH),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Stopped,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
assert!(classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopped)).is_none());
let stopping_after_stop = RebalanceMeta {
stopped_at: Some(OffsetDateTime::UNIX_EPOCH),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
stopping: true,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopping_after_stop)),
Some(HealFormatPoolSkip::Retryable)
));
}
#[test]
fn skipped_heal_format_pool_is_never_reported_as_success() {
assert!(matches!(
heal_format_pool_skip_error(HealFormatPoolSkip::Retryable),
StorageError::SlowDown
));
assert!(matches!(
heal_format_pool_skip_error(HealFormatPoolSkip::Completed),
StorageError::NoHealRequired
));
}
async fn multi_pool_heal_store() -> (tempfile::TempDir, Arc<ECStore>, CancellationToken) {
let temp_dir = tempfile::tempdir().expect("multi-pool heal test directory should be created");
let mut pool_endpoints = Vec::new();
@@ -1189,18 +889,6 @@ mod tests {
bucket_fence_registry: std::sync::Arc::default(),
};
let err = store
.handle_heal_format(false)
.await
.expect_err("missing pool metadata must fail closed before format writes");
assert!(matches!(err, StorageError::SlowDown));
let pool_meta = PoolMeta::new(&store.pools, &PoolMeta::default());
pool_meta
.save(store.pools.clone())
.await
.expect("pool metadata should be persisted before format heal");
let (result, err) = store
.handle_heal_format(false)
.await
@@ -1214,22 +902,5 @@ mod tests {
.await
.expect("the later pool should be healed despite the first pool error");
assert_eq!(healed.erasure.this, recoverable_format.erasure.sets[0][2]);
let mut completed_meta = PoolMeta::new(&store.pools, &PoolMeta::default());
for status in &mut completed_meta.pools {
status.decommission = Some(PoolDecommissionInfo {
complete: true,
..Default::default()
});
}
completed_meta
.save(store.pools.clone())
.await
.expect("completed pool metadata should be persisted");
let (_, err) = store
.handle_heal_format(false)
.await
.expect("completed pools should be reported as a no-op");
assert!(matches!(err, Some(StorageError::NoHealRequired)));
}
}
+1 -1
View File
@@ -3194,7 +3194,7 @@ impl ECStore {
// Default return value
let mut del_objects = vec![DeletedObject::default(); objects.len()];
let mut accounting = vec![None; objects.len()];
let accounting = vec![None; objects.len()];
let mut del_errs = Vec::with_capacity(objects.len());
for _ in 0..objects.len() {
@@ -271,7 +271,7 @@ pub(super) fn resolve_latest_object_info_candidates(
.filter(|candidate| latest_candidate_mod_time(candidate) == Some(latest_mod_time))
.collect::<Vec<_>>();
latest_candidates.sort_by(|left, right| right.idx.cmp(&left.idx));
latest_candidates.sort_by_key(|candidate| std::cmp::Reverse(candidate.idx));
let Some(winner) = latest_candidates.first() else {
return Err(Error::ErasureReadQuorum);
@@ -231,10 +231,6 @@ impl HealTask {
"Heal erasure set format repair skipped because no format heal was required"
);
} else {
let error = e;
if error.is_recoverable_heal() {
return Err(error);
}
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_RESULT,
@@ -243,7 +239,7 @@ impl HealTask {
task_id = %self.id,
set_disk_id,
result = "format_failed",
error = %error,
error = %e,
"Heal erasure set failed"
);
{
@@ -251,7 +247,7 @@ impl HealTask {
progress.update_progress(4, 4, 0, 0);
}
return Err(Error::TaskExecutionFailed {
message: format!("Failed to heal disk format for {set_disk_id}: {error}"),
message: format!("Failed to heal disk format for {set_disk_id}: {e}"),
});
}
} else {
@@ -288,9 +284,6 @@ impl HealTask {
Err(Error::TaskCancelled) => return Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => return Err(Error::TaskTimeout),
Err(e) => {
if e.is_recoverable_heal() {
return Err(e);
}
error!(
target: "rustfs::heal::task",
event = EVENT_HEAL_ERASURE_SET_RESULT,
-28
View File
@@ -547,7 +547,6 @@ struct MockStorage {
heal_object_outcome: Mutex<Option<MockHealObjectOutcome>>,
heal_object_outcomes: Mutex<HashMap<String, VecDeque<MockHealObjectOutcome>>>,
format_no_heal_required: Mutex<bool>,
format_error: Mutex<Option<Error>>,
global_format_calls: Mutex<u32>,
replacement_format_calls: Mutex<Vec<(usize, usize, Vec<String>)>>,
replacement_targets_ready: Mutex<bool>,
@@ -868,9 +867,6 @@ impl HealStorageAPI for MockStorage {
async fn heal_format(&self, _dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
*self.global_format_calls.lock().unwrap() += 1;
if let Some(error) = self.format_error.lock().unwrap().take() {
return Err(error);
}
let no_heal_required = *self.format_no_heal_required.lock().unwrap();
if no_heal_required {
Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::NoHealRequired))))
@@ -2056,30 +2052,6 @@ async fn test_erasure_set_heal_continues_after_format_no_heal_required() {
);
}
#[tokio::test]
async fn erasure_set_format_slowdown_is_propagated() {
let storage = Arc::new(MockStorage {
format_error: Mutex::new(Some(Error::Storage(EcstoreError::SlowDown))),
..Default::default()
});
let request = HealRequest::new(
HealType::ErasureSet {
buckets: Vec::new(),
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions::default(),
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage);
let error = task
.execute()
.await
.expect_err("format SlowDown must remain recoverable for the task manager");
assert!(matches!(error, Error::Storage(EcstoreError::SlowDown)));
}
#[tokio::test]
async fn erasure_set_bucket_prepass_failure_stops_before_object_heal() {
let temp = TempDir::new().expect("temporary directory should be created");
-12
View File
@@ -245,18 +245,6 @@ impl TestECStoreEnvBuilder {
.await
.expect("build test ECStore");
// The production bootstrap only persists pool.bin from the elected
// first cluster node. Test stores intentionally have no cluster
// election, but heal-format still requires that durable fence before
// it can write any disk format. Materialize the validated topology
// here so the shared fixture models a ready single-node store.
let mut pool_meta = ecstore.pool_meta.read().await.clone();
pool_meta.dont_save = false;
pool_meta
.save(ecstore.pools.clone())
.await
.expect("persist test pool metadata");
if self.init_bucket_metadata {
let buckets_list = ecstore
.list_bucket(&BucketOptions {
+2 -2
View File
@@ -84,7 +84,7 @@
| protocols | 16 | 🌙 |
| quota_test | 14 | |
| reliability_disk_fault_test | 4 | |
| reliant | 32 | 19 ✅ |
| reliant | 25 | 19 ✅ |
| replication_extension_test | 75 | 20 ✅ +55 🌙 |
| security_boundary_test | 4 | |
| server_startup_failfast_test | 1 | |
@@ -99,4 +99,4 @@
| tls_hot_reload_test | 1 | ✅ |
| version_id_regression_test | 10 | ✅ |
**Total listed: 582 tests across 82 modules · PR smoke: 163 tests / 36 modules · merge/main full: 460 tests / 73 modules · nightly replication: 55 tests · nightly cluster faults: 28 tests / 7 modules · nightly protocols: 16 tests** · updated 2026-08-23.
**Total listed: 575 tests across 82 modules · PR smoke: 163 tests / 36 modules · merge/main full: 453 tests / 73 modules · nightly replication: 55 tests · nightly cluster faults: 28 tests / 7 modules · nightly protocols: 16 tests** · updated 2026-08-23.
+2 -2
View File
@@ -322,7 +322,7 @@ thiserror = { workspace = true }
tracing.workspace = true
url = { workspace = true }
urlencoding = { workspace = true }
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] }
zip = { workspace = true }
libc = { workspace = true }
rand = { workspace = true, features = ["serde"] }
@@ -345,7 +345,7 @@ libsystemd.workspace = true
libmimalloc-sys.workspace = true
[dev-dependencies]
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] }
serial_test = { workspace = true }
tempfile = { workspace = true }
aws-config = { workspace = true }
+32 -121
View File
@@ -41,7 +41,7 @@ use crate::admin::storage_api::config::save_admin_config;
use crate::admin::storage_api::contract::bucket::{
BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp,
};
use crate::admin::storage_api::error::{Error as StorageError, is_err_bucket_not_found};
use crate::admin::storage_api::error::Error as StorageError;
use crate::admin::storage_api::runtime::ECStore;
use crate::admin::utils::{encode_compatible_admin_payload, read_compatible_admin_body};
use crate::auth::constant_time_eq;
@@ -55,7 +55,6 @@ use crate::storage::storage_api::{
use base64::Engine;
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
use base64::engine::general_purpose::URL_SAFE_NO_PAD;
use futures::StreamExt;
use hmac::{Hmac, Mac};
use http::header::{CONTENT_TYPE, HOST};
use http::{HeaderMap, HeaderValue, Uri};
@@ -2097,18 +2096,6 @@ async fn remote_add_preflight_info(site: &PeerSite) -> S3Result<SiteReplicationA
format!("invalid site replication metainfo from `{}`: {e}", site.endpoint),
)
})?;
if info.deployment_id.is_empty() {
// The peer will be tracked under a locally derived fallback ID
// (deployment_id_for_endpoint) instead of its real deployment ID.
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
result = "peer_deployment_id_missing",
peer_endpoint = %site.endpoint,
"admin site replication state"
);
}
let idp_body = send_peer_admin_get_request_with_client(
&client,
@@ -2219,30 +2206,20 @@ fn site_replication_bootstrap_token(uri: &Uri) -> Option<String> {
query_pairs(uri).get("bootstrapToken").cloned()
}
/// Query for a peer `make-with-versioning` bucket op. `versioningEnabled`
/// always travels so the outbound query matches MinIO's site-replication
/// make-bucket wire contract: MinIO's own create-bucket hook sends
/// `versioningEnabled=true` on this op. RustFS's inbound handler
/// force-enables versioning either way.
fn make_with_versioning_bucket_op_path(bucket: &str, created_at: Option<&str>, lock_enabled: bool) -> String {
fn bootstrap_bucket_make_op_path(bucket: &SRBucketInfo) -> String {
let mut query = form_urlencoded::Serializer::new(String::new());
query.append_pair("bucket", bucket);
query.append_pair("operation", SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING);
query.append_pair("versioningEnabled", "true");
if let Some(created_at) = created_at {
query.append_pair("createdAt", created_at);
query.append_pair("bucket", &bucket.bucket);
query.append_pair("operation", "make-with-versioning");
if let Some(created_at) = bucket
.created_at
.and_then(|value| value.format(&time::format_description::well_known::Rfc3339).ok())
{
query.append_pair("createdAt", &created_at);
}
if lock_enabled {
if bucket.object_lock_config.is_some() {
query.append_pair("lockEnabled", "true");
}
format!("{SITE_REPLICATION_PEER_BUCKET_OPS_PATH}?{}", query.finish())
}
fn bootstrap_bucket_make_op_path(bucket: &SRBucketInfo) -> String {
let created_at = bucket
.created_at
.and_then(|value| value.format(&time::format_description::well_known::Rfc3339).ok());
make_with_versioning_bucket_op_path(&bucket.bucket, created_at.as_deref(), bucket.object_lock_config.is_some())
format!("/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", query.finish())
}
fn bootstrap_bucket_meta_item(bucket: &SRBucketInfo, item_type: &str, updated_at: Option<OffsetDateTime>) -> SRBucketMeta {
@@ -4269,7 +4246,16 @@ async fn broadcast_site_replication_make_bucket(
.format(&time::format_description::well_known::Rfc3339)
.unwrap_or_default();
let path = make_with_versioning_bucket_op_path(bucket, Some(&created_at), lock_enabled);
let path = {
let mut query = form_urlencoded::Serializer::new(String::new());
query.append_pair("bucket", bucket);
query.append_pair("operation", "make-with-versioning");
query.append_pair("createdAt", &created_at);
if lock_enabled {
query.append_pair("lockEnabled", "true");
}
format!("/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", query.finish())
};
let path = if let Some(token) = bootstrap_token {
with_site_replication_bootstrap_token(&path, token)
} else {
@@ -10220,25 +10206,13 @@ impl Operation for SiteReplicationStatusHandler {
}
}
/// `POST /v3/site-replication/devnull` — peer link-check upload drain.
/// MinIO streams multi-megabyte probe bodies here during site netperf link
/// checks and expects an unbounded discard (its handler copies to io.Discard);
/// buffering through the 1MB admin body cap turned any larger probe into a
/// 400 and a false link failure. Stream and discard instead — no size cap.
async fn drain_site_replication_devnull(mut input: Body) -> S3Result<()> {
while let Some(chunk) = input.next().await {
chunk.map_err(|e| s3_error!(InvalidRequest, "failed to read devnull stream: {}", e))?;
}
Ok(())
}
pub struct SiteReplicationDevNullHandler {}
#[async_trait::async_trait]
impl Operation for SiteReplicationDevNullHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?;
drain_site_replication_devnull(req.input).await?;
let _ = read_plain_admin_body(req.input).await?;
Ok(empty_response(StatusCode::NO_CONTENT))
}
}
@@ -10497,19 +10471,6 @@ impl Operation for SRPeerJoinHandler {
}
}
/// Outcome of a peer-driven `purge-deleted-bucket` replay. A bucket that is
/// already gone means the purge raced an earlier replay or a local delete —
/// that is success — but any other failure must reach the sender like the
/// sibling delete branches do: swallowing it answered 200 while the bucket
/// survived on this site.
fn purge_deleted_bucket_result(result: Result<(), StorageError>) -> S3Result<()> {
match result {
Ok(()) => Ok(()),
Err(err) if is_err_bucket_not_found(&err) => Ok(()),
Err(err) => Err(ApiError::from(err).into()),
}
}
pub struct SRPeerBucketOpsHandler {}
#[async_trait::async_trait]
@@ -10609,18 +10570,16 @@ impl Operation for SRPeerBucketOpsHandler {
.map_err(ApiError::from)?;
}
"purge-deleted-bucket" => {
purge_deleted_bucket_result(
store
.delete_bucket(
&bucket,
&DeleteBucketOptions {
force: true,
srdelete_op: SRBucketDeleteOp::Purge,
..Default::default()
},
)
.await,
)?;
let _ = store
.delete_bucket(
&bucket,
&DeleteBucketOptions {
force: true,
srdelete_op: SRBucketDeleteOp::Purge,
..Default::default()
},
)
.await;
}
_ => return Err(s3_error!(InvalidRequest, "unsupported site replication bucket operation")),
}
@@ -13966,54 +13925,6 @@ mod tests {
assert!(!query_flag(&uri, "missing"));
}
/// A5 red-light: a `purge-deleted-bucket` replay must report success when
/// the bucket is already gone, and must propagate every other failure —
/// the swallowed error answered 200 while the bucket survived.
#[test]
fn test_purge_deleted_bucket_result_tolerates_only_missing_bucket() {
assert!(purge_deleted_bucket_result(Ok(())).is_ok());
assert!(purge_deleted_bucket_result(Err(StorageError::BucketNotFound("photos".to_string()))).is_ok());
assert!(purge_deleted_bucket_result(Err(StorageError::VolumeNotFound)).is_ok());
let err = purge_deleted_bucket_result(Err(StorageError::StorageFull))
.expect_err("non-not-found delete failures must propagate");
assert_ne!(*err.code(), S3ErrorCode::NoSuchBucket);
}
/// C5 red-light: the site-replication devnull drain must accept bodies
/// beyond the 1MB admin body cap — MinIO's link check streams large
/// probe bodies and treats a 400 as a broken link.
#[tokio::test]
async fn test_site_replication_devnull_drains_body_beyond_admin_cap() {
let body = Body::from(vec![0u8; MAX_ADMIN_REQUEST_BODY_SIZE + 1]);
drain_site_replication_devnull(body)
.await
.expect("devnull must drain bodies larger than the admin body cap");
}
/// A3 red-light: `versioningEnabled` must travel on every outbound
/// make-with-versioning bucket op so the query matches MinIO's
/// site-replication make-bucket wire contract (MinIO's own hook sends
/// `versioningEnabled=true` on this op).
#[test]
fn test_make_with_versioning_op_paths_send_versioning_enabled() {
let bucket = SRBucketInfo {
bucket: "photos".to_string(),
created_at: Some(OffsetDateTime::UNIX_EPOCH),
object_lock_config: Some(BASE64_STANDARD.encode("<ObjectLockConfiguration/>")),
..Default::default()
};
let bootstrap = bootstrap_bucket_make_op_path(&bucket);
assert!(bootstrap.contains("operation=make-with-versioning"), "{bootstrap}");
assert!(bootstrap.contains("versioningEnabled=true"), "{bootstrap}");
assert!(bootstrap.contains("createdAt="), "{bootstrap}");
assert!(bootstrap.contains("lockEnabled=true"), "{bootstrap}");
// The broadcast path (create-bucket hook) shares the same builder.
let broadcast = make_with_versioning_bucket_op_path("photos", Some("1970-01-01T00:00:00Z"), false);
assert!(broadcast.contains("versioningEnabled=true"), "{broadcast}");
assert!(!broadcast.contains("lockEnabled"), "{broadcast}");
}
#[tokio::test]
#[serial]
async fn test_add_bootstrap_scope_only_allows_expected_bucket_setup_until_guard_drops() {
+5 -24
View File
@@ -13,9 +13,9 @@
// limitations under the License.
use rustfs_madmin::{PeerInfo, SyncStatus};
use std::collections::BTreeMap;
use std::collections::{BTreeMap, hash_map::DefaultHasher};
use std::hash::{Hash, Hasher};
use url::Url;
use uuid::Uuid;
fn has_http_scheme(endpoint: &str) -> bool {
endpoint.get(..7).is_some_and(|prefix| prefix.eq_ignore_ascii_case("http://"))
@@ -66,12 +66,10 @@ pub fn site_identity_key(endpoint: &str) -> String {
.unwrap_or_else(|| trimmed.to_ascii_lowercase())
}
/// Fallback deployment ID for a peer that reported none. UUIDv5 over the
/// canonical endpoint: the ID is persisted in site-replication state and
/// broadcast to peers, so it must be identical across Rust toolchains
/// (`DefaultHasher` is not) and across spellings of the same endpoint.
pub fn deployment_id_for_endpoint(endpoint: &str) -> String {
Uuid::new_v5(&Uuid::NAMESPACE_URL, canonical_endpoint(endpoint).as_bytes()).to_string()
let mut hasher = DefaultHasher::new();
endpoint.hash(&mut hasher);
format!("{:016x}", hasher.finish())
}
pub fn same_identity_endpoint(left: &str, right: &str) -> bool {
@@ -176,23 +174,6 @@ mod tests {
}
}
/// B8 red-light: the fallback deployment ID must be a toolchain-stable
/// UUIDv5 over the canonical endpoint — `DefaultHasher` output is not
/// guaranteed stable across Rust releases, yet the ID is persisted in
/// site-replication state and broadcast to peers.
#[test]
fn deployment_id_for_endpoint_is_stable_uuid_v5_over_canonical_endpoint() {
let endpoint = "https://node-a.example.com:9000";
let id = deployment_id_for_endpoint(endpoint);
let parsed = uuid::Uuid::parse_str(&id).expect("fallback deployment ID must be a UUID");
assert_eq!(parsed.get_version_num(), 5, "fallback deployment ID must be UUIDv5");
// Deterministic for the same endpoint and for spelling variants that
// share a canonical form; distinct endpoints stay distinct.
assert_eq!(id, deployment_id_for_endpoint(endpoint));
assert_eq!(id, deployment_id_for_endpoint(" HTTPS://Node-A.Example.Com:9000/ "));
assert_ne!(id, deployment_id_for_endpoint("https://node-b.example.com:9000"));
}
#[test]
fn canonical_endpoint_accepts_case_insensitive_scheme() {
assert_eq!(
+1 -2
View File
@@ -51,7 +51,7 @@ mod ecstore_disk {
}
mod ecstore_error {
pub(crate) use crate::storage::storage_api::ecstore_error::{StorageError, is_err_bucket_not_found};
pub(crate) use crate::storage::storage_api::ecstore_error::StorageError;
}
#[allow(unused_imports)]
@@ -919,7 +919,6 @@ pub(crate) mod contract {
}
pub(crate) mod error {
pub(crate) use super::ecstore_error::is_err_bucket_not_found;
pub(crate) use super::{Error, StorageError};
}