mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 12:49:04 +00:00
chore: merge main into e2e prerequisites
This commit is contained in:
@@ -39,9 +39,6 @@ mod kms_edge_cases_test;
|
||||
#[cfg(test)]
|
||||
mod kms_fault_recovery_test;
|
||||
|
||||
#[cfg(test)]
|
||||
mod test_runner;
|
||||
|
||||
#[cfg(test)]
|
||||
mod bucket_default_encryption_test;
|
||||
|
||||
|
||||
@@ -1,499 +0,0 @@
|
||||
// Copyright 2024 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
#![allow(dead_code)]
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
//! Unified KMS test suite runner
|
||||
//!
|
||||
//! This module provides a unified interface for running KMS tests with categorization,
|
||||
//! filtering, and comprehensive reporting capabilities.
|
||||
|
||||
use crate::common::init_logging;
|
||||
use std::time::Instant;
|
||||
use tokio::time::{Duration, sleep};
|
||||
use tracing::{debug, error, info, warn};
|
||||
|
||||
/// Test category for organization and filtering
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
||||
pub enum TestCategory {
|
||||
CoreFunctionality,
|
||||
MultipartEncryption,
|
||||
EdgeCases,
|
||||
FaultRecovery,
|
||||
Comprehensive,
|
||||
Performance,
|
||||
}
|
||||
|
||||
impl TestCategory {
|
||||
pub fn as_str(&self) -> &'static str {
|
||||
match self {
|
||||
TestCategory::CoreFunctionality => "core-functionality",
|
||||
TestCategory::MultipartEncryption => "multipart-encryption",
|
||||
TestCategory::EdgeCases => "edge-cases",
|
||||
TestCategory::FaultRecovery => "fault-recovery",
|
||||
TestCategory::Comprehensive => "comprehensive",
|
||||
TestCategory::Performance => "performance",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Test definition with metadata
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct TestDefinition {
|
||||
pub name: String,
|
||||
pub description: String,
|
||||
pub category: TestCategory,
|
||||
pub estimated_duration: Duration,
|
||||
pub is_critical: bool,
|
||||
}
|
||||
|
||||
impl TestDefinition {
|
||||
pub fn new(
|
||||
name: impl Into<String>,
|
||||
description: impl Into<String>,
|
||||
category: TestCategory,
|
||||
estimated_duration: Duration,
|
||||
is_critical: bool,
|
||||
) -> Self {
|
||||
Self {
|
||||
name: name.into(),
|
||||
description: description.into(),
|
||||
category,
|
||||
estimated_duration,
|
||||
is_critical,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Test execution result
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct TestResult {
|
||||
pub test_name: String,
|
||||
pub category: TestCategory,
|
||||
pub success: bool,
|
||||
pub duration: Duration,
|
||||
pub error_message: Option<String>,
|
||||
}
|
||||
|
||||
impl TestResult {
|
||||
pub fn success(test_name: String, category: TestCategory, duration: Duration) -> Self {
|
||||
Self {
|
||||
test_name,
|
||||
category,
|
||||
success: true,
|
||||
duration,
|
||||
error_message: None,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn failure(test_name: String, category: TestCategory, duration: Duration, error: String) -> Self {
|
||||
Self {
|
||||
test_name,
|
||||
category,
|
||||
success: false,
|
||||
duration,
|
||||
error_message: Some(error),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Comprehensive test suite configuration
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct TestSuiteConfig {
|
||||
pub categories: Vec<TestCategory>,
|
||||
pub include_critical_only: bool,
|
||||
pub max_duration: Option<Duration>,
|
||||
pub parallel_execution: bool,
|
||||
}
|
||||
|
||||
impl Default for TestSuiteConfig {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
categories: vec![
|
||||
TestCategory::CoreFunctionality,
|
||||
TestCategory::MultipartEncryption,
|
||||
TestCategory::EdgeCases,
|
||||
TestCategory::FaultRecovery,
|
||||
TestCategory::Comprehensive,
|
||||
],
|
||||
include_critical_only: false,
|
||||
max_duration: None,
|
||||
parallel_execution: false,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Unified KMS test suite runner
|
||||
pub struct KMSTestSuite {
|
||||
tests: Vec<TestDefinition>,
|
||||
config: TestSuiteConfig,
|
||||
}
|
||||
|
||||
impl KMSTestSuite {
|
||||
/// Create a new test suite with default configuration
|
||||
pub fn new() -> Self {
|
||||
let tests = vec![
|
||||
// Core Functionality Tests
|
||||
TestDefinition::new(
|
||||
"test_local_kms_end_to_end",
|
||||
"End-to-end KMS test with all encryption types",
|
||||
TestCategory::CoreFunctionality,
|
||||
Duration::from_secs(60),
|
||||
true,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_local_kms_key_isolation",
|
||||
"Test KMS key isolation and security",
|
||||
TestCategory::CoreFunctionality,
|
||||
Duration::from_secs(45),
|
||||
true,
|
||||
),
|
||||
// Multipart Encryption Tests
|
||||
TestDefinition::new(
|
||||
"test_local_kms_multipart_upload",
|
||||
"Test large file multipart upload with encryption",
|
||||
TestCategory::MultipartEncryption,
|
||||
Duration::from_secs(120),
|
||||
true,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_step1_basic_single_file_encryption",
|
||||
"Basic single file encryption test",
|
||||
TestCategory::MultipartEncryption,
|
||||
Duration::from_secs(30),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_step2_basic_multipart_upload_without_encryption",
|
||||
"Basic multipart upload without encryption",
|
||||
TestCategory::MultipartEncryption,
|
||||
Duration::from_secs(45),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_step3_multipart_upload_with_sse_s3",
|
||||
"Multipart upload with SSE-S3 encryption",
|
||||
TestCategory::MultipartEncryption,
|
||||
Duration::from_secs(60),
|
||||
true,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_step4_large_multipart_upload_with_encryption",
|
||||
"Large file multipart upload with encryption",
|
||||
TestCategory::MultipartEncryption,
|
||||
Duration::from_secs(90),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_step5_all_encryption_types_multipart",
|
||||
"All encryption types multipart test",
|
||||
TestCategory::MultipartEncryption,
|
||||
Duration::from_secs(120),
|
||||
true,
|
||||
),
|
||||
// Edge Cases Tests
|
||||
TestDefinition::new(
|
||||
"test_kms_zero_byte_file_encryption",
|
||||
"Test encryption of zero-byte files",
|
||||
TestCategory::EdgeCases,
|
||||
Duration::from_secs(20),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_kms_single_byte_file_encryption",
|
||||
"Test encryption of single-byte files",
|
||||
TestCategory::EdgeCases,
|
||||
Duration::from_secs(20),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_kms_multipart_boundary_conditions",
|
||||
"Test multipart upload boundary conditions",
|
||||
TestCategory::EdgeCases,
|
||||
Duration::from_secs(45),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_kms_invalid_key_scenarios",
|
||||
"Test invalid key scenarios",
|
||||
TestCategory::EdgeCases,
|
||||
Duration::from_secs(30),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_kms_concurrent_encryption",
|
||||
"Test concurrent encryption operations",
|
||||
TestCategory::EdgeCases,
|
||||
Duration::from_secs(60),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_kms_key_validation_security",
|
||||
"Test key validation security",
|
||||
TestCategory::EdgeCases,
|
||||
Duration::from_secs(30),
|
||||
false,
|
||||
),
|
||||
// Fault Recovery Tests
|
||||
TestDefinition::new(
|
||||
"test_kms_key_directory_unavailable",
|
||||
"Test KMS when key directory is unavailable",
|
||||
TestCategory::FaultRecovery,
|
||||
Duration::from_secs(45),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_kms_corrupted_key_files",
|
||||
"Test KMS with corrupted key files",
|
||||
TestCategory::FaultRecovery,
|
||||
Duration::from_secs(30),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_kms_multipart_upload_interruption",
|
||||
"Test multipart upload interruption recovery",
|
||||
TestCategory::FaultRecovery,
|
||||
Duration::from_secs(60),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_kms_resource_constraints",
|
||||
"Test KMS under resource constraints",
|
||||
TestCategory::FaultRecovery,
|
||||
Duration::from_secs(90),
|
||||
false,
|
||||
),
|
||||
// Comprehensive Tests
|
||||
TestDefinition::new(
|
||||
"test_comprehensive_kms_full_workflow",
|
||||
"Full KMS workflow comprehensive test",
|
||||
TestCategory::Comprehensive,
|
||||
Duration::from_secs(300),
|
||||
true,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_comprehensive_stress_test",
|
||||
"KMS stress test with large datasets",
|
||||
TestCategory::Comprehensive,
|
||||
Duration::from_secs(400),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_comprehensive_key_isolation",
|
||||
"Comprehensive key isolation test",
|
||||
TestCategory::Comprehensive,
|
||||
Duration::from_secs(180),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_comprehensive_concurrent_operations",
|
||||
"Comprehensive concurrent operations test",
|
||||
TestCategory::Comprehensive,
|
||||
Duration::from_secs(240),
|
||||
false,
|
||||
),
|
||||
TestDefinition::new(
|
||||
"test_comprehensive_performance_benchmark",
|
||||
"KMS performance benchmark test",
|
||||
TestCategory::Comprehensive,
|
||||
Duration::from_secs(360),
|
||||
false,
|
||||
),
|
||||
];
|
||||
|
||||
Self {
|
||||
tests,
|
||||
config: TestSuiteConfig::default(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Configure the test suite
|
||||
pub fn with_config(mut self, config: TestSuiteConfig) -> Self {
|
||||
self.config = config;
|
||||
self
|
||||
}
|
||||
|
||||
/// Filter tests based on category
|
||||
pub fn filter_by_category(&self, category: &TestCategory) -> Vec<&TestDefinition> {
|
||||
self.tests.iter().filter(|test| &test.category == category).collect()
|
||||
}
|
||||
|
||||
/// Filter tests based on criticality
|
||||
pub fn filter_critical_tests(&self) -> Vec<&TestDefinition> {
|
||||
self.tests.iter().filter(|test| test.is_critical).collect()
|
||||
}
|
||||
|
||||
/// Get test summary by category
|
||||
pub fn get_category_summary(&self) -> std::collections::HashMap<TestCategory, Vec<&TestDefinition>> {
|
||||
let mut summary = std::collections::HashMap::new();
|
||||
for test in &self.tests {
|
||||
summary.entry(test.category.clone()).or_insert_with(Vec::new).push(test);
|
||||
}
|
||||
summary
|
||||
}
|
||||
|
||||
/// Run the complete test suite
|
||||
pub async fn run_test_suite(&self) -> Vec<TestResult> {
|
||||
init_logging();
|
||||
info!("🚀 Starting unified KMS test suite");
|
||||
|
||||
let start_time = Instant::now();
|
||||
let mut results = Vec::new();
|
||||
|
||||
// Filter tests based on configuration
|
||||
let tests_to_run: Vec<&TestDefinition> = self
|
||||
.tests
|
||||
.iter()
|
||||
.filter(|test| self.config.categories.contains(&test.category))
|
||||
.filter(|test| !self.config.include_critical_only || test.is_critical)
|
||||
.collect();
|
||||
|
||||
info!("📊 Test plan: {} test(s) scheduled", tests_to_run.len());
|
||||
for (i, test) in tests_to_run.iter().enumerate() {
|
||||
info!(" {}. {} ({})", i + 1, test.name, test.category.as_str());
|
||||
}
|
||||
|
||||
// Execute tests
|
||||
for (i, test_def) in tests_to_run.iter().enumerate() {
|
||||
info!("🧪 Running test {}/{}: {}", i + 1, tests_to_run.len(), test_def.name);
|
||||
info!(" 📝 Description: {}", test_def.description);
|
||||
info!(" 🏷️ Category: {}", test_def.category.as_str());
|
||||
info!(" ⏱️ Estimated duration: {:?}", test_def.estimated_duration);
|
||||
|
||||
let test_start = Instant::now();
|
||||
let result = self.run_single_test(test_def).await;
|
||||
let test_duration = test_start.elapsed();
|
||||
|
||||
match result {
|
||||
Ok(_) => {
|
||||
info!("✅ Test passed: {} ({:.2}s)", test_def.name, test_duration.as_secs_f64());
|
||||
results.push(TestResult::success(test_def.name.clone(), test_def.category.clone(), test_duration));
|
||||
}
|
||||
Err(e) => {
|
||||
error!("❌ Test failed: {} ({:.2}s): {}", test_def.name, test_duration.as_secs_f64(), e);
|
||||
results.push(TestResult::failure(
|
||||
test_def.name.clone(),
|
||||
test_def.category.clone(),
|
||||
test_duration,
|
||||
e.to_string(),
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
// Add delay between tests to avoid resource conflicts
|
||||
if i < tests_to_run.len() - 1 {
|
||||
debug!("⏸️ Waiting two seconds before the next test...");
|
||||
sleep(Duration::from_secs(2)).await;
|
||||
}
|
||||
}
|
||||
|
||||
let total_duration = start_time.elapsed();
|
||||
self.print_test_summary(&results, total_duration);
|
||||
|
||||
results
|
||||
}
|
||||
|
||||
/// Run a single test by dispatching to the appropriate test function
|
||||
async fn run_single_test(&self, test_def: &TestDefinition) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
// This is a placeholder for test dispatch logic
|
||||
// In a real implementation, this would dispatch to actual test functions
|
||||
warn!("⚠️ Test '{}' is not implemented in the unified runner; skipping", test_def.name);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Print comprehensive test summary
|
||||
fn print_test_summary(&self, results: &[TestResult], total_duration: Duration) {
|
||||
info!("📊 KMS test suite summary");
|
||||
info!("⏱️ Total duration: {:.2} seconds", total_duration.as_secs_f64());
|
||||
info!("📈 Total tests: {}", results.len());
|
||||
|
||||
let passed = results.iter().filter(|r| r.success).count();
|
||||
let failed = results.iter().filter(|r| !r.success).count();
|
||||
|
||||
info!("✅ Passed: {}", passed);
|
||||
info!("❌ Failed: {}", failed);
|
||||
info!("📊 Success rate: {:.1}%", (passed as f64 / results.len() as f64) * 100.0);
|
||||
|
||||
// Summary by category
|
||||
let mut category_summary: std::collections::HashMap<TestCategory, (usize, usize)> = std::collections::HashMap::new();
|
||||
for result in results {
|
||||
let (total, passed_count) = category_summary.entry(result.category.clone()).or_insert((0, 0));
|
||||
*total += 1;
|
||||
if result.success {
|
||||
*passed_count += 1;
|
||||
}
|
||||
}
|
||||
|
||||
info!("📊 Category summary:");
|
||||
for (category, (total, passed_count)) in category_summary {
|
||||
info!(
|
||||
" 🏷️ {}: {}/{} ({:.1}%)",
|
||||
category.as_str(),
|
||||
passed_count,
|
||||
total,
|
||||
(passed_count as f64 / total as f64) * 100.0
|
||||
);
|
||||
}
|
||||
|
||||
// List failed tests
|
||||
if failed > 0 {
|
||||
warn!("❌ Failing tests:");
|
||||
for result in results.iter().filter(|r| !r.success) {
|
||||
warn!(" - {}: {}", result.test_name, result.error_message.as_deref().unwrap_or("Unknown error"));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Quick test suite for critical tests only
|
||||
#[tokio::test]
|
||||
async fn test_kms_critical_suite() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
let config = TestSuiteConfig {
|
||||
categories: vec![TestCategory::CoreFunctionality, TestCategory::MultipartEncryption],
|
||||
include_critical_only: true,
|
||||
max_duration: Some(Duration::from_secs(600)), // 10 minutes max
|
||||
parallel_execution: false,
|
||||
};
|
||||
|
||||
let suite = KMSTestSuite::new().with_config(config);
|
||||
let results = suite.run_test_suite().await;
|
||||
|
||||
let failed_count = results.iter().filter(|r| !r.success).count();
|
||||
if failed_count > 0 {
|
||||
return Err(format!("Critical test suite failed: {failed_count} tests failed").into());
|
||||
}
|
||||
|
||||
info!("✅ All critical tests passed");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Full comprehensive test suite
|
||||
#[tokio::test]
|
||||
async fn test_kms_full_suite() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
let suite = KMSTestSuite::new();
|
||||
let results = suite.run_test_suite().await;
|
||||
|
||||
let total_tests = results.len();
|
||||
let failed_count = results.iter().filter(|r| !r.success).count();
|
||||
let success_rate = ((total_tests - failed_count) as f64 / total_tests as f64) * 100.0;
|
||||
|
||||
info!("📊 Full suite success rate: {:.1}%", success_rate);
|
||||
|
||||
// Allow up to 10% failure rate for non-critical tests
|
||||
if success_rate < 90.0 {
|
||||
return Err(format!("Test suite success rate too low: {success_rate:.1}%").into());
|
||||
}
|
||||
|
||||
info!("✅ Full test suite succeeded");
|
||||
Ok(())
|
||||
}
|
||||
@@ -16,6 +16,7 @@ use crate::common::{RustFSTestClusterEnvironment, RustFSTestEnvironment, init_lo
|
||||
use aws_sdk_s3::primitives::ByteStream;
|
||||
use http::header::{CONTENT_TYPE, HOST};
|
||||
use reqwest::StatusCode;
|
||||
use rustfs_config::{ENV_DRIVE_ACTIVE_CHECK_INTERVAL_SECS, ENV_NOTIFY_ENABLE};
|
||||
use rustfs_signer::pre_sign_v4;
|
||||
use rustfs_utils::egress::ENV_OUTBOUND_ALLOW_ORIGINS;
|
||||
use s3s::Body;
|
||||
@@ -976,7 +977,8 @@ async fn test_get_object_lambda_rejects_disabled_target() -> Result<(), Box<dyn
|
||||
init_logging();
|
||||
|
||||
let mut env = RustFSTestEnvironment::new().await?;
|
||||
env.start_rustfs_server(vec![]).await?;
|
||||
env.start_rustfs_server_with_env(vec![], &[(ENV_NOTIFY_ENABLE, "true")])
|
||||
.await?;
|
||||
|
||||
let bucket = "object-lambda-e2e-disabled-target";
|
||||
let key = "input.txt";
|
||||
@@ -992,17 +994,24 @@ async fn test_get_object_lambda_rejects_disabled_target() -> Result<(), Box<dyn
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
configure_webhook_target_with_key_values(
|
||||
&env,
|
||||
"transformer",
|
||||
vec![
|
||||
("endpoint", "http://127.0.0.1:9/transform".to_string()),
|
||||
("auth_token", "secret-token".to_string()),
|
||||
("enable", "off".to_string()),
|
||||
],
|
||||
let queue_dir = format!("{}/disabled-target-queue", env.temp_dir);
|
||||
tokio::fs::create_dir_all(&queue_dir).await?;
|
||||
let config_url = format!("{}/rustfs/admin/v3/set-config-kv", env.url);
|
||||
let directive = format!(
|
||||
"notify_webhook:transformer enable=off endpoint=\"http://127.0.0.1:9/transform\" auth_token=\"secret-token\" queue_dir=\"{queue_dir}\""
|
||||
);
|
||||
let disable_response = signed_request(
|
||||
http::Method::PUT,
|
||||
&config_url,
|
||||
&env.access_key,
|
||||
&env.secret_key,
|
||||
Some(directive.into_bytes()),
|
||||
Some("text/plain"),
|
||||
)
|
||||
.await?;
|
||||
wait_for_target_visibility(&env, "transformer").await?;
|
||||
let disable_status = disable_response.status();
|
||||
let disable_body = disable_response.text().await?;
|
||||
assert_eq!(disable_status, StatusCode::OK, "failed to disable target: {disable_body}");
|
||||
|
||||
let lambda_url = format!("{}/{}/{}?lambdaArn={}", env.url, bucket, key, urlencoding::encode(lambda_arn));
|
||||
let response = signed_request(http::Method::GET, &lambda_url, &env.access_key, &env.secret_key, None, None).await?;
|
||||
@@ -1021,7 +1030,8 @@ async fn test_configure_object_lambda_target_rejects_invalid_endpoint() -> Resul
|
||||
init_logging();
|
||||
|
||||
let mut env = RustFSTestEnvironment::new().await?;
|
||||
env.start_rustfs_server(vec![]).await?;
|
||||
env.start_rustfs_server_with_env(vec![], &[(ENV_NOTIFY_ENABLE, "true")])
|
||||
.await?;
|
||||
|
||||
let bucket = "object-lambda-e2e-invalid-endpoint";
|
||||
|
||||
@@ -1064,7 +1074,8 @@ async fn test_configure_object_lambda_notify_webhook_rejects_response_header_tim
|
||||
init_logging();
|
||||
|
||||
let mut env = RustFSTestEnvironment::new().await?;
|
||||
env.start_rustfs_server(vec![]).await?;
|
||||
env.start_rustfs_server_with_env(vec![], &[(ENV_NOTIFY_ENABLE, "true")])
|
||||
.await?;
|
||||
|
||||
let response = send_configure_webhook_target_request(
|
||||
&env,
|
||||
@@ -1173,6 +1184,8 @@ async fn test_listen_notification_fans_in_remote_node_events() -> Result<(), Box
|
||||
init_logging();
|
||||
|
||||
let mut cluster = RustFSTestClusterEnvironment::new(2).await?;
|
||||
cluster.set_env(ENV_NOTIFY_ENABLE, "true");
|
||||
cluster.set_env(ENV_DRIVE_ACTIVE_CHECK_INTERVAL_SECS, "1");
|
||||
cluster.start().await?;
|
||||
|
||||
let bucket = "listen-notification-cluster";
|
||||
|
||||
@@ -15,7 +15,6 @@
|
||||
use crate::common::{RustFSTestClusterEnvironment, init_logging};
|
||||
use aws_sdk_s3::error::SdkError;
|
||||
use aws_sdk_s3::primitives::ByteStream;
|
||||
use aws_sdk_s3::types::CompletedMultipartUpload;
|
||||
use tokio::time::{Duration, sleep};
|
||||
use tracing::info;
|
||||
use uuid::Uuid;
|
||||
@@ -43,32 +42,18 @@ async fn list_parts_reports_missing_upload(
|
||||
}
|
||||
}
|
||||
|
||||
async fn complete_reports_missing_upload(
|
||||
async fn multipart_listing_reports_missing_upload(
|
||||
client: &aws_sdk_s3::Client,
|
||||
bucket: &str,
|
||||
key: &str,
|
||||
upload_id: &str,
|
||||
) -> Result<bool, Box<dyn std::error::Error + Send + Sync>> {
|
||||
let result = client
|
||||
.complete_multipart_upload()
|
||||
.bucket(bucket)
|
||||
.key(key)
|
||||
.upload_id(upload_id)
|
||||
.multipart_upload(CompletedMultipartUpload::builder().build())
|
||||
.send()
|
||||
.await;
|
||||
match result {
|
||||
Ok(_) => Ok(false),
|
||||
Err(SdkError::ServiceError(err)) => {
|
||||
let code = err.err().meta().code().unwrap_or("");
|
||||
if code == "NoSuchUpload" {
|
||||
Ok(true)
|
||||
} else {
|
||||
Err(format!("unexpected complete_multipart_upload service error: code={code}, err={err:?}").into())
|
||||
}
|
||||
}
|
||||
Err(err) => Err(format!("unexpected complete_multipart_upload error: {err:?}").into()),
|
||||
}
|
||||
let result = client.list_multipart_uploads().bucket(bucket).prefix(key).send().await?;
|
||||
|
||||
Ok(!result
|
||||
.uploads()
|
||||
.iter()
|
||||
.any(|upload| upload.key() == Some(key) && upload.upload_id() == Some(upload_id)))
|
||||
}
|
||||
|
||||
async fn wait_for_cleanup_on_all_nodes(
|
||||
@@ -81,8 +66,8 @@ async fn wait_for_cleanup_on_all_nodes(
|
||||
let mut all_cleaned = true;
|
||||
for (idx, client) in clients.iter().enumerate() {
|
||||
let list_parts_missing = list_parts_reports_missing_upload(client, bucket, key, upload_id).await?;
|
||||
let complete_missing = complete_reports_missing_upload(client, bucket, key, upload_id).await?;
|
||||
if !(list_parts_missing && complete_missing) {
|
||||
let listing_missing = multipart_listing_reports_missing_upload(client, bucket, key, upload_id).await?;
|
||||
if !(list_parts_missing && listing_missing) {
|
||||
info!("stale multipart still visible on node {} at attempt {}", idx, attempt + 1);
|
||||
all_cleaned = false;
|
||||
break;
|
||||
@@ -146,6 +131,10 @@ async fn test_stale_multipart_cleanup_removes_incomplete_upload_across_cluster()
|
||||
1,
|
||||
"multipart upload should be visible before background cleanup"
|
||||
);
|
||||
assert!(
|
||||
!multipart_listing_reports_missing_upload(&clients[2], CLEANUP_BUCKET, &key, &upload_id).await?,
|
||||
"multipart upload listing should contain the upload before background cleanup"
|
||||
);
|
||||
|
||||
wait_for_cleanup_on_all_nodes(&clients, CLEANUP_BUCKET, &key, &upload_id).await?;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user