mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| f4da36b364 | |||
| 20d1266496 | |||
| f7003dfddd |
@@ -1,2 +1,2 @@
|
|||||||
sha256-darwin=9f767b37ed8b1c82da62ea441462d75487785c8086e56f08fb6f6cd89c6e2e52
|
sha256-darwin=03bdfb6a9d6e25d744c385f1461e651f05e4da78e5b0ead2adb3e8b2463e3834
|
||||||
sha256-linux=fbdaf42b220958d4b1e8880e0f8b5a7992d38e21051bb60596dd4538424757d6
|
sha256-linux=78c46adad135231017fb8679fb91877b5ae3ec6ce1f75d3d045d2c1076c12a49
|
||||||
|
|||||||
Generated
+1
@@ -12688,6 +12688,7 @@ dependencies = [
|
|||||||
"js-sys",
|
"js-sys",
|
||||||
"rand 0.10.2",
|
"rand 0.10.2",
|
||||||
"serde_core",
|
"serde_core",
|
||||||
|
"sha1_smol",
|
||||||
"wasm-bindgen",
|
"wasm-bindgen",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
|||||||
@@ -1,52 +1,26 @@
|
|||||||
#![cfg(test)]
|
#![cfg(test)]
|
||||||
|
|
||||||
use aws_config::meta::region::RegionProviderChain;
|
use crate::common::{RustFSTestEnvironment, TEST_BUCKET, init_logging};
|
||||||
use aws_sdk_s3::Client;
|
use aws_sdk_s3::Client;
|
||||||
use aws_sdk_s3::config::{Credentials, Region};
|
use aws_sdk_s3::error::{ProvideErrorMetadata, SdkError};
|
||||||
use aws_sdk_s3::error::SdkError;
|
|
||||||
use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart};
|
use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart};
|
||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use std::error::Error;
|
use std::error::Error;
|
||||||
|
use std::fmt::Debug;
|
||||||
|
|
||||||
const ENDPOINT: &str = "http://localhost:9000";
|
type TestResult = Result<(), Box<dyn Error + Send + Sync>>;
|
||||||
const ACCESS_KEY: &str = "rustfsadmin";
|
|
||||||
const SECRET_KEY: &str = "rustfsadmin";
|
|
||||||
const BUCKET: &str = "api-test";
|
|
||||||
|
|
||||||
async fn create_aws_s3_client() -> Result<Client, Box<dyn Error>> {
|
fn assert_s3_error_code<T, E>(result: Result<T, SdkError<E>>, expected: &str)
|
||||||
let region_provider = RegionProviderChain::default_provider().or_else(Region::new("us-east-1"));
|
where
|
||||||
let shared_config = aws_config::defaults(aws_config::BehaviorVersion::latest())
|
T: Debug,
|
||||||
.region(region_provider)
|
E: ProvideErrorMetadata + Debug,
|
||||||
.credentials_provider(Credentials::new(ACCESS_KEY, SECRET_KEY, None, None, "static"))
|
{
|
||||||
.endpoint_url(ENDPOINT)
|
let error = result.expect_err("conditional request must fail");
|
||||||
.load()
|
assert_eq!(
|
||||||
.await;
|
error.as_service_error().and_then(ProvideErrorMetadata::code),
|
||||||
|
Some(expected),
|
||||||
let client = Client::from_conf(
|
"unexpected conditional request error: {error:?}"
|
||||||
aws_sdk_s3::Config::from(&shared_config)
|
|
||||||
.to_builder()
|
|
||||||
.force_path_style(true)
|
|
||||||
.build(),
|
|
||||||
);
|
);
|
||||||
Ok(client)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Setup test bucket, creating it if it doesn't exist
|
|
||||||
async fn setup_test_bucket(client: &Client) -> Result<(), Box<dyn Error>> {
|
|
||||||
match client.create_bucket().bucket(BUCKET).send().await {
|
|
||||||
Ok(_) => {}
|
|
||||||
Err(SdkError::ServiceError(e)) => {
|
|
||||||
let e = e.into_err();
|
|
||||||
let error_code = e.meta().code().unwrap_or("");
|
|
||||||
if !error_code.eq("BucketAlreadyExists") {
|
|
||||||
return Err(e.into());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Err(e) => {
|
|
||||||
return Err(e.into());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Ok(())
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Generate test data of specified size
|
/// Generate test data of specified size
|
||||||
@@ -60,7 +34,12 @@ fn generate_test_data(size: usize) -> Vec<u8> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Upload an object and return its ETag
|
/// Upload an object and return its ETag
|
||||||
async fn upload_object_with_metadata(client: &Client, bucket: &str, key: &str, data: &[u8]) -> Result<String, Box<dyn Error>> {
|
async fn upload_object_with_metadata(
|
||||||
|
client: &Client,
|
||||||
|
bucket: &str,
|
||||||
|
key: &str,
|
||||||
|
data: &[u8],
|
||||||
|
) -> Result<String, Box<dyn Error + Send + Sync>> {
|
||||||
let response = client
|
let response = client
|
||||||
.put_object()
|
.put_object()
|
||||||
.bucket(bucket)
|
.bucket(bucket)
|
||||||
@@ -69,188 +48,164 @@ async fn upload_object_with_metadata(client: &Client, bucket: &str, key: &str, d
|
|||||||
.send()
|
.send()
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
let etag = response.e_tag().unwrap_or("").to_string();
|
response
|
||||||
Ok(etag)
|
.e_tag()
|
||||||
|
.map(str::to_owned)
|
||||||
|
.ok_or_else(|| std::io::Error::other("put object response did not include an ETag").into())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Cleanup test objects from bucket
|
async fn object_body(client: &Client, key: &str) -> Result<Bytes, Box<dyn Error + Send + Sync>> {
|
||||||
async fn cleanup_objects(client: &Client, bucket: &str, keys: &[&str]) {
|
let response = client.get_object().bucket(TEST_BUCKET).key(key).send().await?;
|
||||||
for key in keys {
|
Ok(response.body.collect().await?.into_bytes())
|
||||||
let _ = client.delete_object().bucket(bucket).key(*key).send().await;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Generate unique test object key
|
|
||||||
fn generate_test_key(prefix: &str) -> String {
|
|
||||||
use std::time::{SystemTime, UNIX_EPOCH};
|
|
||||||
let timestamp = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos();
|
|
||||||
format!("{prefix}-{timestamp}")
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[ignore = "requires running RustFS server at localhost:9000"]
|
async fn test_conditional_put_okay() -> TestResult {
|
||||||
async fn test_conditional_put_okay() -> Result<(), Box<dyn std::error::Error>> {
|
init_logging();
|
||||||
let client = create_aws_s3_client().await?;
|
let mut env = RustFSTestEnvironment::new().await?;
|
||||||
setup_test_bucket(&client).await?;
|
env.start_rustfs_server(vec![]).await?;
|
||||||
|
env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
|
let client = env.create_s3_client();
|
||||||
|
|
||||||
let test_key = generate_test_key("conditional-put-ok");
|
let test_key = "conditional-put-ok";
|
||||||
let initial_data = generate_test_data(1024); // 1KB test data
|
let initial_data = generate_test_data(1024); // 1KB test data
|
||||||
let updated_data = generate_test_data(2048); // 2KB updated data
|
let matching_data = generate_test_data(2048); // 2KB updated data
|
||||||
|
let non_matching_data = generate_test_data(3072); // 3KB updated data
|
||||||
|
|
||||||
// Upload initial object and get its ETag
|
// Upload initial object and get its ETag
|
||||||
let initial_etag = upload_object_with_metadata(&client, BUCKET, &test_key, &initial_data).await?;
|
let initial_etag = upload_object_with_metadata(&client, TEST_BUCKET, test_key, &initial_data).await?;
|
||||||
|
|
||||||
// Test 1: PUT with matching If-Match condition (should succeed)
|
// Test 1: PUT with matching If-Match condition (should succeed)
|
||||||
let response1 = client
|
client
|
||||||
.put_object()
|
.put_object()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(&test_key)
|
.key(test_key)
|
||||||
.body(Bytes::from(updated_data.clone()).into())
|
.body(Bytes::from(matching_data.clone()).into())
|
||||||
.if_match(&initial_etag)
|
.if_match(&initial_etag)
|
||||||
.send()
|
.send()
|
||||||
.await;
|
.await?;
|
||||||
assert!(response1.is_ok(), "PUT with matching If-Match should succeed");
|
assert_eq!(object_body(&client, test_key).await?.as_ref(), matching_data);
|
||||||
|
|
||||||
// Test 2: PUT with non-matching If-None-Match condition (should succeed)
|
// Test 2: PUT with non-matching If-None-Match condition (should succeed)
|
||||||
let fake_etag = "\"fake-etag-12345\"";
|
let fake_etag = "\"fake-etag-12345\"";
|
||||||
let response2 = client
|
client
|
||||||
.put_object()
|
.put_object()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(&test_key)
|
.key(test_key)
|
||||||
.body(Bytes::from(updated_data.clone()).into())
|
.body(Bytes::from(non_matching_data.clone()).into())
|
||||||
.if_none_match(fake_etag)
|
.if_none_match(fake_etag)
|
||||||
.send()
|
.send()
|
||||||
.await;
|
.await?;
|
||||||
assert!(response2.is_ok(), "PUT with non-matching If-None-Match should succeed");
|
assert_eq!(object_body(&client, test_key).await?.as_ref(), non_matching_data);
|
||||||
|
|
||||||
// Cleanup
|
|
||||||
cleanup_objects(&client, BUCKET, &[&test_key]).await;
|
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[ignore = "requires running RustFS server at localhost:9000"]
|
async fn test_conditional_put_failed() -> TestResult {
|
||||||
async fn test_conditional_put_failed() -> Result<(), Box<dyn std::error::Error>> {
|
init_logging();
|
||||||
let client = create_aws_s3_client().await?;
|
let mut env = RustFSTestEnvironment::new().await?;
|
||||||
setup_test_bucket(&client).await?;
|
env.start_rustfs_server(vec![]).await?;
|
||||||
|
env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
|
let client = env.create_s3_client();
|
||||||
|
|
||||||
let test_key = generate_test_key("conditional-put-failed");
|
let test_key = "conditional-put-failed";
|
||||||
let initial_data = generate_test_data(1024);
|
let initial_data = generate_test_data(1024);
|
||||||
let updated_data = generate_test_data(2048);
|
let updated_data = generate_test_data(2048);
|
||||||
|
|
||||||
// Upload initial object and get its ETag
|
// Upload initial object and get its ETag
|
||||||
let initial_etag = upload_object_with_metadata(&client, BUCKET, &test_key, &initial_data).await?;
|
let initial_etag = upload_object_with_metadata(&client, TEST_BUCKET, test_key, &initial_data).await?;
|
||||||
|
|
||||||
// Test 1: PUT with non-matching If-Match condition (should fail with 412)
|
// Test 1: PUT with non-matching If-Match condition (should fail with 412)
|
||||||
let fake_etag = "\"fake-etag-should-not-match\"";
|
let fake_etag = "\"fake-etag-should-not-match\"";
|
||||||
let response1 = client
|
let response1 = client
|
||||||
.put_object()
|
.put_object()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(&test_key)
|
.key(test_key)
|
||||||
.body(Bytes::from(updated_data.clone()).into())
|
.body(Bytes::from(updated_data.clone()).into())
|
||||||
.if_match(fake_etag)
|
.if_match(fake_etag)
|
||||||
.send()
|
.send()
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
assert!(response1.is_err(), "PUT with non-matching If-Match should fail");
|
assert_s3_error_code(response1, "PreconditionFailed");
|
||||||
if let Err(e) = response1 {
|
assert_eq!(object_body(&client, test_key).await?.as_ref(), initial_data);
|
||||||
if let SdkError::ServiceError(e) = e {
|
|
||||||
let e = e.into_err();
|
|
||||||
let error_code = e.meta().code().unwrap_or("");
|
|
||||||
assert_eq!("PreconditionFailed", error_code);
|
|
||||||
} else {
|
|
||||||
panic!("Unexpected error: {e:?}");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Test 2: PUT with matching If-None-Match condition (should fail with 412)
|
// Test 2: PUT with matching If-None-Match condition (should fail with 412)
|
||||||
let response2 = client
|
let response2 = client
|
||||||
.put_object()
|
.put_object()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(&test_key)
|
.key(test_key)
|
||||||
.body(Bytes::from(updated_data.clone()).into())
|
.body(Bytes::from(updated_data.clone()).into())
|
||||||
.if_none_match(&initial_etag)
|
.if_none_match(&initial_etag)
|
||||||
.send()
|
.send()
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
assert!(response2.is_err(), "PUT with matching If-None-Match should fail");
|
assert_s3_error_code(response2, "PreconditionFailed");
|
||||||
if let Err(e) = response2 {
|
assert_eq!(object_body(&client, test_key).await?.as_ref(), initial_data);
|
||||||
if let SdkError::ServiceError(e) = e {
|
|
||||||
let e = e.into_err();
|
|
||||||
let error_code = e.meta().code().unwrap_or("");
|
|
||||||
assert_eq!("PreconditionFailed", error_code);
|
|
||||||
} else {
|
|
||||||
panic!("Unexpected error: {e:?}");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Cleanup - only need to clean up the initial object since failed PUTs shouldn't create objects
|
|
||||||
cleanup_objects(&client, BUCKET, &[&test_key]).await;
|
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[ignore = "requires running RustFS server at localhost:9000"]
|
async fn test_conditional_put_when_object_does_not_exist() -> TestResult {
|
||||||
async fn test_conditional_put_when_object_does_not_exist() -> Result<(), Box<dyn std::error::Error>> {
|
init_logging();
|
||||||
let client = create_aws_s3_client().await?;
|
let mut env = RustFSTestEnvironment::new().await?;
|
||||||
setup_test_bucket(&client).await?;
|
env.start_rustfs_server(vec![]).await?;
|
||||||
|
env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
|
let client = env.create_s3_client();
|
||||||
|
|
||||||
let key = "some_key";
|
let key = "conditional-put-missing";
|
||||||
cleanup_objects(&client, BUCKET, &[key]).await;
|
|
||||||
|
|
||||||
// When the object does not exist, the If-Match condition should always fail
|
// When the object does not exist, the If-Match condition should always fail
|
||||||
let response1 = client
|
let response1 = client
|
||||||
.put_object()
|
.put_object()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(key)
|
.key(key)
|
||||||
.body(Bytes::from(generate_test_data(1024)).into())
|
.body(Bytes::from(generate_test_data(1024)).into())
|
||||||
.if_match("*")
|
.if_match("*")
|
||||||
.send()
|
.send()
|
||||||
.await;
|
.await;
|
||||||
assert!(response1.is_err());
|
assert_s3_error_code(response1, "NoSuchKey");
|
||||||
if let Err(e) = response1 {
|
|
||||||
if let SdkError::ServiceError(e) = e {
|
|
||||||
let e = e.into_err();
|
|
||||||
let error_code = e.meta().code().unwrap_or("");
|
|
||||||
assert_eq!("NoSuchKey", error_code);
|
|
||||||
} else {
|
|
||||||
panic!("Unexpected error: {e:?}");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// When the object does not exist, the If-None-Match condition should be able to succeed
|
// When the object does not exist, the If-None-Match condition should be able to succeed
|
||||||
let response2 = client
|
let created_data = generate_test_data(1024);
|
||||||
|
client
|
||||||
.put_object()
|
.put_object()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(key)
|
.key(key)
|
||||||
.body(Bytes::from(generate_test_data(1024)).into())
|
.body(Bytes::from(created_data.clone()).into())
|
||||||
.if_none_match("*")
|
.if_none_match("*")
|
||||||
.send()
|
.send()
|
||||||
.await;
|
.await?;
|
||||||
assert!(response2.is_ok());
|
assert_eq!(object_body(&client, key).await?.as_ref(), created_data);
|
||||||
|
|
||||||
cleanup_objects(&client, BUCKET, &[key]).await;
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[ignore = "requires running RustFS server at localhost:9000"]
|
async fn test_conditional_multi_part_upload() -> TestResult {
|
||||||
async fn test_conditional_multi_part_upload() -> Result<(), Box<dyn std::error::Error>> {
|
init_logging();
|
||||||
let client = create_aws_s3_client().await?;
|
let mut env = RustFSTestEnvironment::new().await?;
|
||||||
setup_test_bucket(&client).await?;
|
env.start_rustfs_server(vec![]).await?;
|
||||||
|
env.create_test_bucket(TEST_BUCKET).await?;
|
||||||
|
let client = env.create_s3_client();
|
||||||
|
|
||||||
let test_key = generate_test_key("multipart-upload-ok");
|
let test_key = "conditional-multipart-upload";
|
||||||
let test_data = generate_test_data(1024);
|
let test_data = generate_test_data(1024);
|
||||||
let initial_etag = upload_object_with_metadata(&client, BUCKET, &test_key, &test_data).await?;
|
let initial_etag = upload_object_with_metadata(&client, TEST_BUCKET, test_key, &test_data).await?;
|
||||||
|
|
||||||
let part_size = 5 * 1024 * 1024; // 5MB per part (minimum for multipart)
|
let part_size = 5 * 1024 * 1024; // 5MB per part (minimum for multipart)
|
||||||
let num_parts = 3;
|
let num_parts = 3;
|
||||||
let mut parts = Vec::new();
|
let mut parts = Vec::new();
|
||||||
|
let mut expected_data = Vec::with_capacity(part_size * usize::try_from(num_parts)?);
|
||||||
|
|
||||||
// Initiate multipart upload
|
// Initiate multipart upload
|
||||||
let initiate_response = client.create_multipart_upload().bucket(BUCKET).key(&test_key).send().await?;
|
let initiate_response = client
|
||||||
|
.create_multipart_upload()
|
||||||
|
.bucket(TEST_BUCKET)
|
||||||
|
.key(test_key)
|
||||||
|
.send()
|
||||||
|
.await?;
|
||||||
|
|
||||||
let upload_id = initiate_response
|
let upload_id = initiate_response
|
||||||
.upload_id()
|
.upload_id()
|
||||||
@@ -258,12 +213,13 @@ async fn test_conditional_multi_part_upload() -> Result<(), Box<dyn std::error::
|
|||||||
|
|
||||||
// Upload parts
|
// Upload parts
|
||||||
for part_number in 1..=num_parts {
|
for part_number in 1..=num_parts {
|
||||||
let part_data = generate_test_data(part_size);
|
let part_data = vec![u8::try_from(part_number)?; part_size];
|
||||||
|
expected_data.extend_from_slice(&part_data);
|
||||||
|
|
||||||
let upload_part_response = client
|
let upload_part_response = client
|
||||||
.upload_part()
|
.upload_part()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(&test_key)
|
.key(test_key)
|
||||||
.upload_id(upload_id)
|
.upload_id(upload_id)
|
||||||
.part_number(part_number)
|
.part_number(part_number)
|
||||||
.body(Bytes::from(part_data).into())
|
.body(Bytes::from(part_data).into())
|
||||||
@@ -286,57 +242,62 @@ async fn test_conditional_multi_part_upload() -> Result<(), Box<dyn std::error::
|
|||||||
// Test 1: Multipart upload with wildcard If-None-Match, should fail
|
// Test 1: Multipart upload with wildcard If-None-Match, should fail
|
||||||
let complete_response = client
|
let complete_response = client
|
||||||
.complete_multipart_upload()
|
.complete_multipart_upload()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(&test_key)
|
.key(test_key)
|
||||||
.upload_id(upload_id)
|
.upload_id(upload_id)
|
||||||
.multipart_upload(completed_upload.clone())
|
.multipart_upload(completed_upload.clone())
|
||||||
.if_none_match("*")
|
.if_none_match("*")
|
||||||
.send()
|
.send()
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
assert!(complete_response.is_err());
|
assert_s3_error_code(complete_response, "PreconditionFailed");
|
||||||
|
|
||||||
// Test 2: Multipart upload with matching If-None-Match, should fail
|
// Test 2: Multipart upload with matching If-None-Match, should fail
|
||||||
let complete_response = client
|
let complete_response = client
|
||||||
.complete_multipart_upload()
|
.complete_multipart_upload()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(&test_key)
|
.key(test_key)
|
||||||
.upload_id(upload_id)
|
.upload_id(upload_id)
|
||||||
.multipart_upload(completed_upload.clone())
|
.multipart_upload(completed_upload.clone())
|
||||||
.if_none_match(initial_etag.clone())
|
.if_none_match(initial_etag.clone())
|
||||||
.send()
|
.send()
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
assert!(complete_response.is_err());
|
assert_s3_error_code(complete_response, "PreconditionFailed");
|
||||||
|
|
||||||
// Test 3: Multipart upload with unmatching If-Match, should fail
|
// Test 3: Multipart upload with unmatching If-Match, should fail
|
||||||
let complete_response = client
|
let complete_response = client
|
||||||
.complete_multipart_upload()
|
.complete_multipart_upload()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(&test_key)
|
.key(test_key)
|
||||||
.upload_id(upload_id)
|
.upload_id(upload_id)
|
||||||
.multipart_upload(completed_upload.clone())
|
.multipart_upload(completed_upload.clone())
|
||||||
.if_match("\"abcdef\"")
|
.if_match("\"abcdef\"")
|
||||||
.send()
|
.send()
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
assert!(complete_response.is_err());
|
assert_s3_error_code(complete_response, "PreconditionFailed");
|
||||||
|
|
||||||
|
let staged_parts = client
|
||||||
|
.list_parts()
|
||||||
|
.bucket(TEST_BUCKET)
|
||||||
|
.key(test_key)
|
||||||
|
.upload_id(upload_id)
|
||||||
|
.send()
|
||||||
|
.await?;
|
||||||
|
assert_eq!(staged_parts.parts().len(), usize::try_from(num_parts)?);
|
||||||
|
|
||||||
// Test 4: Multipart upload with matching If-Match, should succeed
|
// Test 4: Multipart upload with matching If-Match, should succeed
|
||||||
let complete_response = client
|
client
|
||||||
.complete_multipart_upload()
|
.complete_multipart_upload()
|
||||||
.bucket(BUCKET)
|
.bucket(TEST_BUCKET)
|
||||||
.key(&test_key)
|
.key(test_key)
|
||||||
.upload_id(upload_id)
|
.upload_id(upload_id)
|
||||||
.multipart_upload(completed_upload.clone())
|
.multipart_upload(completed_upload)
|
||||||
.if_match(initial_etag)
|
.if_match(initial_etag)
|
||||||
.send()
|
.send()
|
||||||
.await;
|
.await?;
|
||||||
|
assert_eq!(object_body(&client, test_key).await?.as_ref(), expected_data);
|
||||||
assert!(complete_response.is_ok());
|
|
||||||
|
|
||||||
// Cleanup
|
|
||||||
cleanup_objects(&client, BUCKET, &[&test_key]).await;
|
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2026,7 +2026,7 @@ impl PoolMeta {
|
|||||||
self.load_no_lock(pool).await
|
self.load_no_lock(pool).await
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn load_no_lock<S>(&mut self, pool: Arc<S>) -> Result<()>
|
pub(crate) async fn load_no_lock<S>(&mut self, pool: Arc<S>) -> Result<()>
|
||||||
where
|
where
|
||||||
S: EcstoreObjectIO,
|
S: EcstoreObjectIO,
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -988,14 +988,11 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for Sets {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[async_trait::async_trait]
|
impl Sets {
|
||||||
impl crate::storage_api_contracts::heal::HealOperations for Sets {
|
pub(crate) async fn heal_format_with_fence<F>(&self, dry_run: bool, fence_lost: F) -> Result<(HealResultItem, Option<Error>)>
|
||||||
type Error = Error;
|
where
|
||||||
type HealResultItem = HealResultItem;
|
F: Fn() -> bool + Send + Sync,
|
||||||
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(
|
let (disks, init_errs) = init_storage_disks_with_errors(
|
||||||
&self.endpoints.endpoints,
|
&self.endpoints.endpoints,
|
||||||
&DiskOption {
|
&DiskOption {
|
||||||
@@ -1068,6 +1065,9 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets {
|
|||||||
// Save new formats `format.json` on unformatted disks.
|
// Save new formats `format.json` on unformatted disks.
|
||||||
for (index, (fm, disk)) in tmp_new_formats.iter_mut().zip(disks.iter()).enumerate() {
|
for (index, (fm, disk)) in tmp_new_formats.iter_mut().zip(disks.iter()).enumerate() {
|
||||||
if fm.is_some() && disk.is_some() {
|
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 Err(err) = save_format_file(disk, fm).await {
|
||||||
if let Some(disk) = disk.as_ref() {
|
if let Some(disk) = disk.as_ref() {
|
||||||
let _ = disk.close().await;
|
let _ = disk.close().await;
|
||||||
@@ -1101,6 +1101,18 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets {
|
|||||||
}
|
}
|
||||||
Ok((res, None))
|
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))]
|
#[tracing::instrument(skip(self))]
|
||||||
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
|
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
|
||||||
let mut result = HealResultItem {
|
let mut result = HealResultItem {
|
||||||
|
|||||||
@@ -13,7 +13,12 @@
|
|||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
use super::*;
|
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::heal::HealOperations as _;
|
||||||
|
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
|
||||||
|
use rustfs_lock::NamespaceLockGuard;
|
||||||
use tracing::trace;
|
use tracing::trace;
|
||||||
|
|
||||||
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
|
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
|
||||||
@@ -30,7 +35,119 @@ fn invalid_heal_pool_index(pool_idx: usize, pool_count: usize) -> Error {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy)]
|
||||||
|
enum HealFormatPoolSkip {
|
||||||
|
Completed,
|
||||||
|
Retryable,
|
||||||
|
}
|
||||||
|
|
||||||
|
fn classify_heal_format_pool(
|
||||||
|
pool_idx: usize,
|
||||||
|
pool_cmd_line: &str,
|
||||||
|
pool_meta: &PoolMeta,
|
||||||
|
rebalance_meta: Option<&RebalanceMeta>,
|
||||||
|
) -> Option<HealFormatPoolSkip> {
|
||||||
|
let Some(pool) = pool_meta.pools.get(pool_idx) else {
|
||||||
|
return Some(HealFormatPoolSkip::Retryable);
|
||||||
|
};
|
||||||
|
|
||||||
|
if pool.id != pool_idx || pool_cmd_line.is_empty() || pool.cmd_line.is_empty() || pool.cmd_line != pool_cmd_line {
|
||||||
|
return Some(HealFormatPoolSkip::Retryable);
|
||||||
|
}
|
||||||
|
|
||||||
|
if let Some(decommission) = pool.decommission.as_ref() {
|
||||||
|
if decommission.complete {
|
||||||
|
return Some(HealFormatPoolSkip::Completed);
|
||||||
|
}
|
||||||
|
if decommission.failed || decommission.canceled || decommission.queued || pool_meta.is_suspended(pool_idx) {
|
||||||
|
return Some(HealFormatPoolSkip::Retryable);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if let Some(meta) = rebalance_meta {
|
||||||
|
let Some(pool_stats) = meta.pool_stats.get(pool_idx) else {
|
||||||
|
return Some(HealFormatPoolSkip::Retryable);
|
||||||
|
};
|
||||||
|
if pool_stats.info.stopping || (pool_stats.participating && pool_stats.info.status == RebalStatus::Started) {
|
||||||
|
return Some(HealFormatPoolSkip::Retryable);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
None
|
||||||
|
}
|
||||||
|
|
||||||
|
fn heal_format_pool_skip_error(skip: HealFormatPoolSkip) -> Error {
|
||||||
|
match skip {
|
||||||
|
HealFormatPoolSkip::Completed => StorageError::NoHealRequired,
|
||||||
|
HealFormatPoolSkip::Retryable => StorageError::SlowDown,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn heal_format_fence_lost_error() -> Error {
|
||||||
|
StorageError::SlowDown
|
||||||
|
}
|
||||||
|
|
||||||
impl ECStore {
|
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>>> {
|
fn get_pools_for_heal_object(&self, opts: &HealOpts) -> Result<Vec<Arc<Sets>>> {
|
||||||
match opts.pool {
|
match opts.pool {
|
||||||
Some(pool_idx) => Ok(vec![
|
Some(pool_idx) => Ok(vec![
|
||||||
@@ -52,9 +169,26 @@ impl ECStore {
|
|||||||
};
|
};
|
||||||
|
|
||||||
let mut count_no_heal = 0;
|
let mut count_no_heal = 0;
|
||||||
|
let mut count_completed = 0;
|
||||||
let mut first_error = None;
|
let mut first_error = None;
|
||||||
for pool in self.pools.iter() {
|
for (pool_idx, pool) in self.pools.iter().enumerate() {
|
||||||
let (mut result, err) = pool.heal_format(dry_run).await?;
|
let (pool_guard, rebalance_guard, pool_meta, rebalance_meta) = self.acquire_heal_format_fence().await?;
|
||||||
|
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
|
||||||
|
first_error.get_or_insert(heal_format_fence_lost_error());
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
if let Some(skip) = classify_heal_format_pool(pool_idx, &pool.endpoints.cmd_line, &pool_meta, rebalance_meta.as_ref())
|
||||||
|
{
|
||||||
|
if matches!(skip, HealFormatPoolSkip::Completed) {
|
||||||
|
count_completed += 1;
|
||||||
|
} else {
|
||||||
|
first_error.get_or_insert(heal_format_pool_skip_error(skip));
|
||||||
|
}
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
let fence_lost = || pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost();
|
||||||
|
let (mut result, err) = pool.heal_format_with_fence(dry_run, fence_lost).await?;
|
||||||
if let Some(err) = err {
|
if let Some(err) = err {
|
||||||
match err {
|
match err {
|
||||||
StorageError::NoHealRequired => {
|
StorageError::NoHealRequired => {
|
||||||
@@ -69,11 +203,18 @@ impl ECStore {
|
|||||||
r.set_count += result.set_count;
|
r.set_count += result.set_count;
|
||||||
r.before.drives.append(&mut result.before.drives);
|
r.before.drives.append(&mut result.before.drives);
|
||||||
r.after.drives.append(&mut result.after.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 {
|
if let Some(err) = first_error {
|
||||||
return Ok((r, Some(err)));
|
return Ok((r, Some(err)));
|
||||||
}
|
}
|
||||||
if count_no_heal == self.pools.len() {
|
if count_no_heal + count_completed == self.pools.len() {
|
||||||
info!(
|
info!(
|
||||||
event = EVENT_HEAL_FORMAT_COMPLETED,
|
event = EVENT_HEAL_FORMAT_COMPLETED,
|
||||||
component = LOG_COMPONENT_ECSTORE,
|
component = LOG_COMPONENT_ECSTORE,
|
||||||
@@ -302,6 +443,7 @@ mod tests {
|
|||||||
use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk};
|
use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk};
|
||||||
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
|
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
|
||||||
use crate::runtime::instance::InstanceContext;
|
use crate::runtime::instance::InstanceContext;
|
||||||
|
use crate::services::rebalance::{RebalanceInfo, RebalanceStats};
|
||||||
use crate::storage_api_contracts::bucket::{BucketOperations, MakeBucketOptions};
|
use crate::storage_api_contracts::bucket::{BucketOperations, MakeBucketOptions};
|
||||||
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations};
|
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations};
|
||||||
use crate::store::init_format::{load_format_erasure, save_format_file};
|
use crate::store::init_format::{load_format_erasure, save_format_file};
|
||||||
@@ -353,6 +495,164 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn pool_meta_with_decommission(info: PoolDecommissionInfo) -> PoolMeta {
|
||||||
|
PoolMeta {
|
||||||
|
pools: vec![PoolStatus {
|
||||||
|
id: 0,
|
||||||
|
cmd_line: "pool-0".to_string(),
|
||||||
|
last_update: OffsetDateTime::UNIX_EPOCH,
|
||||||
|
decommission: Some(info),
|
||||||
|
}],
|
||||||
|
..Default::default()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn heal_format_pool_state_barriers_are_classified() {
|
||||||
|
let active = pool_meta_with_decommission(PoolDecommissionInfo {
|
||||||
|
start_time: Some(OffsetDateTime::UNIX_EPOCH),
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
assert!(matches!(
|
||||||
|
classify_heal_format_pool(0, "pool-0", &active, None),
|
||||||
|
Some(HealFormatPoolSkip::Retryable)
|
||||||
|
));
|
||||||
|
|
||||||
|
for info in [
|
||||||
|
PoolDecommissionInfo {
|
||||||
|
failed: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
PoolDecommissionInfo {
|
||||||
|
canceled: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
] {
|
||||||
|
assert!(matches!(
|
||||||
|
classify_heal_format_pool(0, "pool-0", &pool_meta_with_decommission(info), None),
|
||||||
|
Some(HealFormatPoolSkip::Retryable)
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
let completed = pool_meta_with_decommission(PoolDecommissionInfo {
|
||||||
|
complete: true,
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
assert!(matches!(
|
||||||
|
classify_heal_format_pool(0, "pool-0", &completed, None),
|
||||||
|
Some(HealFormatPoolSkip::Completed)
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn heal_format_pool_rebalance_barriers_and_identity_are_fail_closed() {
|
||||||
|
let identity_meta = pool_meta_with_decommission(PoolDecommissionInfo::default());
|
||||||
|
let rebalance = RebalanceMeta {
|
||||||
|
pool_stats: vec![RebalanceStats {
|
||||||
|
participating: true,
|
||||||
|
info: RebalanceInfo {
|
||||||
|
status: RebalStatus::Started,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
..Default::default()
|
||||||
|
}],
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
assert!(matches!(
|
||||||
|
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&rebalance)),
|
||||||
|
Some(HealFormatPoolSkip::Retryable)
|
||||||
|
));
|
||||||
|
|
||||||
|
let stopping = RebalanceMeta {
|
||||||
|
pool_stats: vec![RebalanceStats {
|
||||||
|
info: RebalanceInfo {
|
||||||
|
stopping: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
..Default::default()
|
||||||
|
}],
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
assert!(matches!(
|
||||||
|
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopping)),
|
||||||
|
Some(HealFormatPoolSkip::Retryable)
|
||||||
|
));
|
||||||
|
|
||||||
|
let identity = pool_meta_with_decommission(PoolDecommissionInfo::default());
|
||||||
|
assert!(matches!(
|
||||||
|
classify_heal_format_pool(0, "pool-new", &identity, None),
|
||||||
|
Some(HealFormatPoolSkip::Retryable)
|
||||||
|
));
|
||||||
|
|
||||||
|
let identity_without_decommission = PoolMeta {
|
||||||
|
pools: vec![PoolStatus {
|
||||||
|
id: 0,
|
||||||
|
cmd_line: "pool-0".to_string(),
|
||||||
|
last_update: OffsetDateTime::UNIX_EPOCH,
|
||||||
|
decommission: None,
|
||||||
|
}],
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
assert!(matches!(
|
||||||
|
classify_heal_format_pool(0, "pool-new", &identity_without_decommission, None),
|
||||||
|
Some(HealFormatPoolSkip::Retryable)
|
||||||
|
));
|
||||||
|
|
||||||
|
assert!(matches!(
|
||||||
|
classify_heal_format_pool(0, "", &identity_meta, None),
|
||||||
|
Some(HealFormatPoolSkip::Retryable)
|
||||||
|
));
|
||||||
|
|
||||||
|
assert!(matches!(
|
||||||
|
classify_heal_format_pool(0, "pool-0", &PoolMeta::default(), None),
|
||||||
|
Some(HealFormatPoolSkip::Retryable)
|
||||||
|
));
|
||||||
|
|
||||||
|
let stopped = RebalanceMeta {
|
||||||
|
stopped_at: Some(OffsetDateTime::UNIX_EPOCH),
|
||||||
|
pool_stats: vec![RebalanceStats {
|
||||||
|
participating: true,
|
||||||
|
info: RebalanceInfo {
|
||||||
|
status: RebalStatus::Stopped,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
..Default::default()
|
||||||
|
}],
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
assert!(classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopped)).is_none());
|
||||||
|
|
||||||
|
let stopping_after_stop = RebalanceMeta {
|
||||||
|
stopped_at: Some(OffsetDateTime::UNIX_EPOCH),
|
||||||
|
pool_stats: vec![RebalanceStats {
|
||||||
|
participating: true,
|
||||||
|
info: RebalanceInfo {
|
||||||
|
status: RebalStatus::Started,
|
||||||
|
stopping: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
..Default::default()
|
||||||
|
}],
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
assert!(matches!(
|
||||||
|
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopping_after_stop)),
|
||||||
|
Some(HealFormatPoolSkip::Retryable)
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn skipped_heal_format_pool_is_never_reported_as_success() {
|
||||||
|
assert!(matches!(
|
||||||
|
heal_format_pool_skip_error(HealFormatPoolSkip::Retryable),
|
||||||
|
StorageError::SlowDown
|
||||||
|
));
|
||||||
|
assert!(matches!(
|
||||||
|
heal_format_pool_skip_error(HealFormatPoolSkip::Completed),
|
||||||
|
StorageError::NoHealRequired
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
async fn multi_pool_heal_store() -> (tempfile::TempDir, Arc<ECStore>, CancellationToken) {
|
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 temp_dir = tempfile::tempdir().expect("multi-pool heal test directory should be created");
|
||||||
let mut pool_endpoints = Vec::new();
|
let mut pool_endpoints = Vec::new();
|
||||||
@@ -889,6 +1189,18 @@ mod tests {
|
|||||||
bucket_fence_registry: std::sync::Arc::default(),
|
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
|
let (result, err) = store
|
||||||
.handle_heal_format(false)
|
.handle_heal_format(false)
|
||||||
.await
|
.await
|
||||||
@@ -902,5 +1214,22 @@ mod tests {
|
|||||||
.await
|
.await
|
||||||
.expect("the later pool should be healed despite the first pool error");
|
.expect("the later pool should be healed despite the first pool error");
|
||||||
assert_eq!(healed.erasure.this, recoverable_format.erasure.sets[0][2]);
|
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)));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -231,6 +231,10 @@ impl HealTask {
|
|||||||
"Heal erasure set format repair skipped because no format heal was required"
|
"Heal erasure set format repair skipped because no format heal was required"
|
||||||
);
|
);
|
||||||
} else {
|
} else {
|
||||||
|
let error = e;
|
||||||
|
if error.is_recoverable_heal() {
|
||||||
|
return Err(error);
|
||||||
|
}
|
||||||
error!(
|
error!(
|
||||||
target: "rustfs::heal::task",
|
target: "rustfs::heal::task",
|
||||||
event = EVENT_HEAL_ERASURE_SET_RESULT,
|
event = EVENT_HEAL_ERASURE_SET_RESULT,
|
||||||
@@ -239,7 +243,7 @@ impl HealTask {
|
|||||||
task_id = %self.id,
|
task_id = %self.id,
|
||||||
set_disk_id,
|
set_disk_id,
|
||||||
result = "format_failed",
|
result = "format_failed",
|
||||||
error = %e,
|
error = %error,
|
||||||
"Heal erasure set failed"
|
"Heal erasure set failed"
|
||||||
);
|
);
|
||||||
{
|
{
|
||||||
@@ -247,7 +251,7 @@ impl HealTask {
|
|||||||
progress.update_progress(4, 4, 0, 0);
|
progress.update_progress(4, 4, 0, 0);
|
||||||
}
|
}
|
||||||
return Err(Error::TaskExecutionFailed {
|
return Err(Error::TaskExecutionFailed {
|
||||||
message: format!("Failed to heal disk format for {set_disk_id}: {e}"),
|
message: format!("Failed to heal disk format for {set_disk_id}: {error}"),
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
@@ -284,6 +288,9 @@ impl HealTask {
|
|||||||
Err(Error::TaskCancelled) => return Err(Error::TaskCancelled),
|
Err(Error::TaskCancelled) => return Err(Error::TaskCancelled),
|
||||||
Err(Error::TaskTimeout) => return Err(Error::TaskTimeout),
|
Err(Error::TaskTimeout) => return Err(Error::TaskTimeout),
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
|
if e.is_recoverable_heal() {
|
||||||
|
return Err(e);
|
||||||
|
}
|
||||||
error!(
|
error!(
|
||||||
target: "rustfs::heal::task",
|
target: "rustfs::heal::task",
|
||||||
event = EVENT_HEAL_ERASURE_SET_RESULT,
|
event = EVENT_HEAL_ERASURE_SET_RESULT,
|
||||||
|
|||||||
@@ -547,6 +547,7 @@ struct MockStorage {
|
|||||||
heal_object_outcome: Mutex<Option<MockHealObjectOutcome>>,
|
heal_object_outcome: Mutex<Option<MockHealObjectOutcome>>,
|
||||||
heal_object_outcomes: Mutex<HashMap<String, VecDeque<MockHealObjectOutcome>>>,
|
heal_object_outcomes: Mutex<HashMap<String, VecDeque<MockHealObjectOutcome>>>,
|
||||||
format_no_heal_required: Mutex<bool>,
|
format_no_heal_required: Mutex<bool>,
|
||||||
|
format_error: Mutex<Option<Error>>,
|
||||||
global_format_calls: Mutex<u32>,
|
global_format_calls: Mutex<u32>,
|
||||||
replacement_format_calls: Mutex<Vec<(usize, usize, Vec<String>)>>,
|
replacement_format_calls: Mutex<Vec<(usize, usize, Vec<String>)>>,
|
||||||
replacement_targets_ready: Mutex<bool>,
|
replacement_targets_ready: Mutex<bool>,
|
||||||
@@ -867,6 +868,9 @@ impl HealStorageAPI for MockStorage {
|
|||||||
|
|
||||||
async fn heal_format(&self, _dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
|
async fn heal_format(&self, _dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
|
||||||
*self.global_format_calls.lock().unwrap() += 1;
|
*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();
|
let no_heal_required = *self.format_no_heal_required.lock().unwrap();
|
||||||
if no_heal_required {
|
if no_heal_required {
|
||||||
Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::NoHealRequired))))
|
Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::NoHealRequired))))
|
||||||
@@ -2052,6 +2056,30 @@ async fn test_erasure_set_heal_continues_after_format_no_heal_required() {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn erasure_set_format_slowdown_is_propagated() {
|
||||||
|
let storage = Arc::new(MockStorage {
|
||||||
|
format_error: Mutex::new(Some(Error::Storage(EcstoreError::SlowDown))),
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
let request = HealRequest::new(
|
||||||
|
HealType::ErasureSet {
|
||||||
|
buckets: Vec::new(),
|
||||||
|
set_disk_id: "pool_0_set_0".to_string(),
|
||||||
|
},
|
||||||
|
HealOptions::default(),
|
||||||
|
HealPriority::Normal,
|
||||||
|
);
|
||||||
|
let task = HealTask::from_request(request, storage);
|
||||||
|
|
||||||
|
let error = task
|
||||||
|
.execute()
|
||||||
|
.await
|
||||||
|
.expect_err("format SlowDown must remain recoverable for the task manager");
|
||||||
|
|
||||||
|
assert!(matches!(error, Error::Storage(EcstoreError::SlowDown)));
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn erasure_set_bucket_prepass_failure_stops_before_object_heal() {
|
async fn erasure_set_bucket_prepass_failure_stops_before_object_heal() {
|
||||||
let temp = TempDir::new().expect("temporary directory should be created");
|
let temp = TempDir::new().expect("temporary directory should be created");
|
||||||
|
|||||||
@@ -245,6 +245,18 @@ impl TestECStoreEnvBuilder {
|
|||||||
.await
|
.await
|
||||||
.expect("build test ECStore");
|
.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 {
|
if self.init_bucket_metadata {
|
||||||
let buckets_list = ecstore
|
let buckets_list = ecstore
|
||||||
.list_bucket(&BucketOptions {
|
.list_bucket(&BucketOptions {
|
||||||
|
|||||||
@@ -84,7 +84,7 @@
|
|||||||
| protocols | 16 | 🌙 |
|
| protocols | 16 | 🌙 |
|
||||||
| quota_test | 14 | |
|
| quota_test | 14 | |
|
||||||
| reliability_disk_fault_test | 4 | |
|
| reliability_disk_fault_test | 4 | |
|
||||||
| reliant | 25 | 19 ✅ |
|
| reliant | 29 | 19 ✅ |
|
||||||
| replication_extension_test | 75 | 20 ✅ +55 🌙 |
|
| replication_extension_test | 75 | 20 ✅ +55 🌙 |
|
||||||
| security_boundary_test | 4 | |
|
| security_boundary_test | 4 | |
|
||||||
| server_startup_failfast_test | 1 | |
|
| server_startup_failfast_test | 1 | |
|
||||||
|
|||||||
+2
-2
@@ -322,7 +322,7 @@ thiserror = { workspace = true }
|
|||||||
tracing.workspace = true
|
tracing.workspace = true
|
||||||
url = { workspace = true }
|
url = { workspace = true }
|
||||||
urlencoding = { workspace = true }
|
urlencoding = { workspace = true }
|
||||||
uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] }
|
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
|
||||||
zip = { workspace = true }
|
zip = { workspace = true }
|
||||||
libc = { workspace = true }
|
libc = { workspace = true }
|
||||||
rand = { workspace = true, features = ["serde"] }
|
rand = { workspace = true, features = ["serde"] }
|
||||||
@@ -345,7 +345,7 @@ libsystemd.workspace = true
|
|||||||
libmimalloc-sys.workspace = true
|
libmimalloc-sys.workspace = true
|
||||||
|
|
||||||
[dev-dependencies]
|
[dev-dependencies]
|
||||||
uuid = { workspace = true, features = ["v4", "fast-rng", "macro-diagnostics"] }
|
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
|
||||||
serial_test = { workspace = true }
|
serial_test = { workspace = true }
|
||||||
tempfile = { workspace = true }
|
tempfile = { workspace = true }
|
||||||
aws-config = { workspace = true }
|
aws-config = { workspace = true }
|
||||||
|
|||||||
@@ -41,7 +41,7 @@ use crate::admin::storage_api::config::save_admin_config;
|
|||||||
use crate::admin::storage_api::contract::bucket::{
|
use crate::admin::storage_api::contract::bucket::{
|
||||||
BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp,
|
BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp,
|
||||||
};
|
};
|
||||||
use crate::admin::storage_api::error::Error as StorageError;
|
use crate::admin::storage_api::error::{Error as StorageError, is_err_bucket_not_found};
|
||||||
use crate::admin::storage_api::runtime::ECStore;
|
use crate::admin::storage_api::runtime::ECStore;
|
||||||
use crate::admin::utils::{encode_compatible_admin_payload, read_compatible_admin_body};
|
use crate::admin::utils::{encode_compatible_admin_payload, read_compatible_admin_body};
|
||||||
use crate::auth::constant_time_eq;
|
use crate::auth::constant_time_eq;
|
||||||
@@ -55,6 +55,7 @@ use crate::storage::storage_api::{
|
|||||||
use base64::Engine;
|
use base64::Engine;
|
||||||
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
|
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
|
||||||
use base64::engine::general_purpose::URL_SAFE_NO_PAD;
|
use base64::engine::general_purpose::URL_SAFE_NO_PAD;
|
||||||
|
use futures::StreamExt;
|
||||||
use hmac::{Hmac, Mac};
|
use hmac::{Hmac, Mac};
|
||||||
use http::header::{CONTENT_TYPE, HOST};
|
use http::header::{CONTENT_TYPE, HOST};
|
||||||
use http::{HeaderMap, HeaderValue, Uri};
|
use http::{HeaderMap, HeaderValue, Uri};
|
||||||
@@ -2096,6 +2097,18 @@ async fn remote_add_preflight_info(site: &PeerSite) -> S3Result<SiteReplicationA
|
|||||||
format!("invalid site replication metainfo from `{}`: {e}", site.endpoint),
|
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(
|
let idp_body = send_peer_admin_get_request_with_client(
|
||||||
&client,
|
&client,
|
||||||
@@ -2206,20 +2219,30 @@ fn site_replication_bootstrap_token(uri: &Uri) -> Option<String> {
|
|||||||
query_pairs(uri).get("bootstrapToken").cloned()
|
query_pairs(uri).get("bootstrapToken").cloned()
|
||||||
}
|
}
|
||||||
|
|
||||||
fn bootstrap_bucket_make_op_path(bucket: &SRBucketInfo) -> String {
|
/// Query for a peer `make-with-versioning` bucket op. `versioningEnabled`
|
||||||
|
/// always travels so the outbound query matches MinIO's site-replication
|
||||||
|
/// make-bucket wire contract: MinIO's own create-bucket hook sends
|
||||||
|
/// `versioningEnabled=true` on this op. RustFS's inbound handler
|
||||||
|
/// force-enables versioning either way.
|
||||||
|
fn make_with_versioning_bucket_op_path(bucket: &str, created_at: Option<&str>, lock_enabled: bool) -> String {
|
||||||
let mut query = form_urlencoded::Serializer::new(String::new());
|
let mut query = form_urlencoded::Serializer::new(String::new());
|
||||||
query.append_pair("bucket", &bucket.bucket);
|
query.append_pair("bucket", bucket);
|
||||||
query.append_pair("operation", "make-with-versioning");
|
query.append_pair("operation", SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING);
|
||||||
if let Some(created_at) = bucket
|
query.append_pair("versioningEnabled", "true");
|
||||||
.created_at
|
if let Some(created_at) = created_at {
|
||||||
.and_then(|value| value.format(&time::format_description::well_known::Rfc3339).ok())
|
query.append_pair("createdAt", created_at);
|
||||||
{
|
|
||||||
query.append_pair("createdAt", &created_at);
|
|
||||||
}
|
}
|
||||||
if bucket.object_lock_config.is_some() {
|
if lock_enabled {
|
||||||
query.append_pair("lockEnabled", "true");
|
query.append_pair("lockEnabled", "true");
|
||||||
}
|
}
|
||||||
format!("/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", query.finish())
|
format!("{SITE_REPLICATION_PEER_BUCKET_OPS_PATH}?{}", query.finish())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn bootstrap_bucket_make_op_path(bucket: &SRBucketInfo) -> String {
|
||||||
|
let created_at = bucket
|
||||||
|
.created_at
|
||||||
|
.and_then(|value| value.format(&time::format_description::well_known::Rfc3339).ok());
|
||||||
|
make_with_versioning_bucket_op_path(&bucket.bucket, created_at.as_deref(), bucket.object_lock_config.is_some())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn bootstrap_bucket_meta_item(bucket: &SRBucketInfo, item_type: &str, updated_at: Option<OffsetDateTime>) -> SRBucketMeta {
|
fn bootstrap_bucket_meta_item(bucket: &SRBucketInfo, item_type: &str, updated_at: Option<OffsetDateTime>) -> SRBucketMeta {
|
||||||
@@ -4246,16 +4269,7 @@ async fn broadcast_site_replication_make_bucket(
|
|||||||
.format(&time::format_description::well_known::Rfc3339)
|
.format(&time::format_description::well_known::Rfc3339)
|
||||||
.unwrap_or_default();
|
.unwrap_or_default();
|
||||||
|
|
||||||
let path = {
|
let path = make_with_versioning_bucket_op_path(bucket, Some(&created_at), lock_enabled);
|
||||||
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 {
|
let path = if let Some(token) = bootstrap_token {
|
||||||
with_site_replication_bootstrap_token(&path, token)
|
with_site_replication_bootstrap_token(&path, token)
|
||||||
} else {
|
} else {
|
||||||
@@ -10206,13 +10220,25 @@ impl Operation for SiteReplicationStatusHandler {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// `POST /v3/site-replication/devnull` — peer link-check upload drain.
|
||||||
|
/// MinIO streams multi-megabyte probe bodies here during site netperf link
|
||||||
|
/// checks and expects an unbounded discard (its handler copies to io.Discard);
|
||||||
|
/// buffering through the 1MB admin body cap turned any larger probe into a
|
||||||
|
/// 400 and a false link failure. Stream and discard instead — no size cap.
|
||||||
|
async fn drain_site_replication_devnull(mut input: Body) -> S3Result<()> {
|
||||||
|
while let Some(chunk) = input.next().await {
|
||||||
|
chunk.map_err(|e| s3_error!(InvalidRequest, "failed to read devnull stream: {}", e))?;
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
pub struct SiteReplicationDevNullHandler {}
|
pub struct SiteReplicationDevNullHandler {}
|
||||||
|
|
||||||
#[async_trait::async_trait]
|
#[async_trait::async_trait]
|
||||||
impl Operation for SiteReplicationDevNullHandler {
|
impl Operation for SiteReplicationDevNullHandler {
|
||||||
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||||
validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?;
|
validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?;
|
||||||
let _ = read_plain_admin_body(req.input).await?;
|
drain_site_replication_devnull(req.input).await?;
|
||||||
Ok(empty_response(StatusCode::NO_CONTENT))
|
Ok(empty_response(StatusCode::NO_CONTENT))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -10471,6 +10497,19 @@ impl Operation for SRPeerJoinHandler {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Outcome of a peer-driven `purge-deleted-bucket` replay. A bucket that is
|
||||||
|
/// already gone means the purge raced an earlier replay or a local delete —
|
||||||
|
/// that is success — but any other failure must reach the sender like the
|
||||||
|
/// sibling delete branches do: swallowing it answered 200 while the bucket
|
||||||
|
/// survived on this site.
|
||||||
|
fn purge_deleted_bucket_result(result: Result<(), StorageError>) -> S3Result<()> {
|
||||||
|
match result {
|
||||||
|
Ok(()) => Ok(()),
|
||||||
|
Err(err) if is_err_bucket_not_found(&err) => Ok(()),
|
||||||
|
Err(err) => Err(ApiError::from(err).into()),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub struct SRPeerBucketOpsHandler {}
|
pub struct SRPeerBucketOpsHandler {}
|
||||||
|
|
||||||
#[async_trait::async_trait]
|
#[async_trait::async_trait]
|
||||||
@@ -10570,16 +10609,18 @@ impl Operation for SRPeerBucketOpsHandler {
|
|||||||
.map_err(ApiError::from)?;
|
.map_err(ApiError::from)?;
|
||||||
}
|
}
|
||||||
"purge-deleted-bucket" => {
|
"purge-deleted-bucket" => {
|
||||||
let _ = store
|
purge_deleted_bucket_result(
|
||||||
.delete_bucket(
|
store
|
||||||
&bucket,
|
.delete_bucket(
|
||||||
&DeleteBucketOptions {
|
&bucket,
|
||||||
force: true,
|
&DeleteBucketOptions {
|
||||||
srdelete_op: SRBucketDeleteOp::Purge,
|
force: true,
|
||||||
..Default::default()
|
srdelete_op: SRBucketDeleteOp::Purge,
|
||||||
},
|
..Default::default()
|
||||||
)
|
},
|
||||||
.await;
|
)
|
||||||
|
.await,
|
||||||
|
)?;
|
||||||
}
|
}
|
||||||
_ => return Err(s3_error!(InvalidRequest, "unsupported site replication bucket operation")),
|
_ => return Err(s3_error!(InvalidRequest, "unsupported site replication bucket operation")),
|
||||||
}
|
}
|
||||||
@@ -13925,6 +13966,54 @@ mod tests {
|
|||||||
assert!(!query_flag(&uri, "missing"));
|
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]
|
#[tokio::test]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn test_add_bootstrap_scope_only_allows_expected_bucket_setup_until_guard_drops() {
|
async fn test_add_bootstrap_scope_only_allows_expected_bucket_setup_until_guard_drops() {
|
||||||
|
|||||||
@@ -13,9 +13,9 @@
|
|||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
use rustfs_madmin::{PeerInfo, SyncStatus};
|
use rustfs_madmin::{PeerInfo, SyncStatus};
|
||||||
use std::collections::{BTreeMap, hash_map::DefaultHasher};
|
use std::collections::BTreeMap;
|
||||||
use std::hash::{Hash, Hasher};
|
|
||||||
use url::Url;
|
use url::Url;
|
||||||
|
use uuid::Uuid;
|
||||||
|
|
||||||
fn has_http_scheme(endpoint: &str) -> bool {
|
fn has_http_scheme(endpoint: &str) -> bool {
|
||||||
endpoint.get(..7).is_some_and(|prefix| prefix.eq_ignore_ascii_case("http://"))
|
endpoint.get(..7).is_some_and(|prefix| prefix.eq_ignore_ascii_case("http://"))
|
||||||
@@ -66,10 +66,12 @@ pub fn site_identity_key(endpoint: &str) -> String {
|
|||||||
.unwrap_or_else(|| trimmed.to_ascii_lowercase())
|
.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 {
|
pub fn deployment_id_for_endpoint(endpoint: &str) -> String {
|
||||||
let mut hasher = DefaultHasher::new();
|
Uuid::new_v5(&Uuid::NAMESPACE_URL, canonical_endpoint(endpoint).as_bytes()).to_string()
|
||||||
endpoint.hash(&mut hasher);
|
|
||||||
format!("{:016x}", hasher.finish())
|
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn same_identity_endpoint(left: &str, right: &str) -> bool {
|
pub fn same_identity_endpoint(left: &str, right: &str) -> bool {
|
||||||
@@ -174,6 +176,23 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// B8 red-light: the fallback deployment ID must be a toolchain-stable
|
||||||
|
/// UUIDv5 over the canonical endpoint — `DefaultHasher` output is not
|
||||||
|
/// guaranteed stable across Rust releases, yet the ID is persisted in
|
||||||
|
/// site-replication state and broadcast to peers.
|
||||||
|
#[test]
|
||||||
|
fn deployment_id_for_endpoint_is_stable_uuid_v5_over_canonical_endpoint() {
|
||||||
|
let endpoint = "https://node-a.example.com:9000";
|
||||||
|
let id = deployment_id_for_endpoint(endpoint);
|
||||||
|
let parsed = uuid::Uuid::parse_str(&id).expect("fallback deployment ID must be a UUID");
|
||||||
|
assert_eq!(parsed.get_version_num(), 5, "fallback deployment ID must be UUIDv5");
|
||||||
|
// Deterministic for the same endpoint and for spelling variants that
|
||||||
|
// share a canonical form; distinct endpoints stay distinct.
|
||||||
|
assert_eq!(id, deployment_id_for_endpoint(endpoint));
|
||||||
|
assert_eq!(id, deployment_id_for_endpoint(" HTTPS://Node-A.Example.Com:9000/ "));
|
||||||
|
assert_ne!(id, deployment_id_for_endpoint("https://node-b.example.com:9000"));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn canonical_endpoint_accepts_case_insensitive_scheme() {
|
fn canonical_endpoint_accepts_case_insensitive_scheme() {
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
|
|||||||
@@ -51,7 +51,7 @@ mod ecstore_disk {
|
|||||||
}
|
}
|
||||||
|
|
||||||
mod ecstore_error {
|
mod ecstore_error {
|
||||||
pub(crate) use crate::storage::storage_api::ecstore_error::StorageError;
|
pub(crate) use crate::storage::storage_api::ecstore_error::{StorageError, is_err_bucket_not_found};
|
||||||
}
|
}
|
||||||
|
|
||||||
#[allow(unused_imports)]
|
#[allow(unused_imports)]
|
||||||
@@ -919,6 +919,7 @@ pub(crate) mod contract {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) mod error {
|
pub(crate) mod error {
|
||||||
|
pub(crate) use super::ecstore_error::is_err_bucket_not_found;
|
||||||
pub(crate) use super::{Error, StorageError};
|
pub(crate) use super::{Error, StorageError};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user