Compare commits

..

3 Commits

Author SHA1 Message Date
Zhengchao An 815a256f6f Merge branch 'main' into chore/tier-lint-stage-b 2026-09-07 02:11:48 +08:00
houseme cf9db7a06e Merge branch 'main' into chore/tier-lint-stage-b 2026-09-06 23:50:57 +08:00
cxymds 6c3572873e chore(tier): remove remaining blanket lint allowances 2026-09-06 22:10:04 +08:00
8 changed files with 124 additions and 356 deletions
@@ -3837,157 +3837,6 @@ async fn test_bucket_replication_converges_delete_marker_and_version_purge() ->
Ok(())
}
/// Regression for rustfs/backlog#2340 (not Wasabi specific): a directory
/// marker (`prefix/` with a body) in a versioned bucket is stored as the null
/// version, like MinIO (`putOpts`: "for directory objects skip creating new
/// versions"), and must still replicate to completion instead of staying
/// `PENDING`.
#[tokio::test]
async fn test_bucket_replication_replicates_directory_marker_in_versioned_bucket() -> TestResult {
init_logging();
let mut source_env = RustFSTestEnvironment::new().await?;
let mut source_env_vars = replication_fast_env();
source_env_vars.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
source_env.start_rustfs_server_with_env(vec![], &source_env_vars).await?;
let mut target_env = RustFSTestEnvironment::new().await?;
target_env.start_rustfs_server_without_cleanup(vec![]).await?;
let source_bucket = "replication-dir-marker-src";
let target_bucket = "replication-dir-marker-dst";
let source_client = source_env.create_s3_client();
let target_client = target_env.create_s3_client();
source_client.create_bucket().bucket(source_bucket).send().await?;
target_client.create_bucket().bucket(target_bucket).send().await?;
enable_bucket_versioning(&source_env, source_bucket).await?;
enable_bucket_versioning(&target_env, target_bucket).await?;
let target_arn = set_replication_target(&source_env, source_bucket, &target_env, target_bucket).await?;
put_bucket_replication(&source_env, source_bucket, &target_arn).await?;
let marker_key = "dir/trailing/";
let body = b"directory marker body";
let put = source_client
.put_object()
.bucket(source_bucket)
.key(marker_key)
.body(ByteStream::from_static(body))
.send()
.await?;
assert!(
put.version_id()
.is_none_or(|id| id == "null" || id == uuid::Uuid::nil().to_string()),
"a directory marker is the null version even in a versioned bucket: {:?}",
put.version_id()
);
wait_for_source_replication_status(&source_client, source_bucket, marker_key, "COMPLETED", false).await?;
let replica = target_client
.get_object()
.bucket(target_bucket)
.key(marker_key)
.send()
.await?;
assert_eq!(replica.body.collect().await?.into_bytes().as_ref(), body);
let listed = target_client
.list_object_versions()
.bucket(target_bucket)
.prefix(marker_key)
.send()
.await?;
let marker_versions: Vec<_> = listed.versions().iter().filter(|v| v.key() == Some(marker_key)).collect();
assert_eq!(marker_versions.len(), 1, "the marker must land exactly once: {marker_versions:?}");
assert_eq!(
marker_versions[0].version_id(),
Some("null"),
"the replica keeps the null version identity"
);
Ok(())
}
/// Regression for rustfs/backlog#2340 (not Wasabi specific): permanently
/// deleting a version whose payload lives in a data dir must leave the source
/// clean once the purge replicates. Managed-SSE objects are never inlined and a
/// plain object above the inline threshold takes the same layout. The version
/// retained with a pending purge used to lose its data dir, so the purge state
/// could never be applied (`VersionNotFound` on every retry) and the bucket
/// stayed `BucketNotEmpty` while `ListObjectVersions` was already empty.
#[tokio::test]
async fn test_bucket_replication_version_purge_of_non_inline_object_releases_source_bucket() -> TestResult {
init_logging();
let (source_env, target_env, source_bucket, target_bucket) = build_sse_replication_pair("purge-datadir", true, true).await?;
let target_arn = wait_for_remote_target_arn(&source_env, &source_bucket).await?;
put_bucket_replication_with_delete_statuses(&source_env, &source_bucket, &target_arn, "Enabled", Some("Enabled")).await?;
let source_client = source_env.create_s3_client();
let target_client = target_env.create_s3_client();
let sse_key = "sse-object.bin";
let large_key = "large-object.bin";
let sse_put = source_client
.put_object()
.bucket(&source_bucket)
.key(sse_key)
.body(ByteStream::from_static(b"encrypted source payload"))
.server_side_encryption(ServerSideEncryption::Aes256)
.send()
.await?;
let large_put = source_client
.put_object()
.bucket(&source_bucket)
.key(large_key)
.body(ByteStream::from(vec![0x5a; 2 * 1024 * 1024]))
.send()
.await?;
let purged = [
(sse_key, sse_put.version_id().ok_or("SSE PUT omitted version ID")?.to_string()),
(large_key, large_put.version_id().ok_or("large PUT omitted version ID")?.to_string()),
];
assert_replication_converged(&source_client, &source_bucket, &target_client, &target_bucket).await?;
for (key, version_id) in &purged {
source_client
.delete_object()
.bucket(&source_bucket)
.key(*key)
.version_id(version_id)
.send()
.await?;
}
assert_replication_converged(&source_client, &source_bucket, &target_client, &target_bucket).await?;
let target_state = list_replication_state(&target_client, &target_bucket).await?;
assert!(target_state.is_empty(), "target retained an explicitly purged version: {target_state:?}");
// The purge state is applied on the source asynchronously after the target
// acknowledges the delete; only then does the retained version go away and
// the bucket become deletable. A listing that is empty while DeleteBucket
// keeps answering BucketNotEmpty is exactly the regression.
let deadline = tokio::time::Instant::now() + Duration::from_secs(60);
loop {
let listing = source_client.list_object_versions().bucket(&source_bucket).send().await?;
let listed = listing.versions().len() + listing.delete_markers().len();
match source_client.delete_bucket().bucket(&source_bucket).send().await {
Ok(_) => break,
Err(err) if err.code() == Some("BucketNotEmpty") => {
if tokio::time::Instant::now() >= deadline {
return Err(format!(
"source bucket stayed BucketNotEmpty after the version purge replicated; \
ListObjectVersions shows {listed} entries"
)
.into());
}
sleep(Duration::from_millis(500)).await;
}
Err(err) => return Err(err.into()),
}
}
Ok(())
}
#[tokio::test]
async fn test_bucket_replication_disabled_delete_marker_does_not_propagate() -> TestResult {
init_logging();
+1
View File
@@ -14,6 +14,7 @@
#[cfg(any(test, feature = "test-util"))]
pub mod test_util;
#[allow(clippy::module_inception, reason = "preserve the public services::tier::tier path")]
pub mod tier;
pub mod tier_admin;
pub mod tier_config;
+89 -93
View File
@@ -15,8 +15,6 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use byteorder::{ByteOrder, LittleEndian};
use bytes::Bytes;
@@ -802,7 +800,7 @@ pub enum TierConfigUpdateError {
}
enum TierCandidateMutation {
Add(TierConfig, bool),
Add(Box<TierConfig>, bool),
Edit(String, TierCreds),
Remove(String, bool),
Clear(bool),
@@ -823,7 +821,7 @@ struct PrevalidatedTierCandidateMutation {
impl TierCandidateMutation {
fn add(mut config: TierConfig, force: bool) -> std::result::Result<Self, AdminError> {
normalize_s3_gcs_add_tier_name(&mut config)?;
Ok(Self::Add(config, force))
Ok(Self::Add(Box::new(config), force))
}
fn normalize_add_tier_name(&mut self) -> std::result::Result<(), AdminError> {
@@ -910,7 +908,7 @@ impl TierCandidateMutation {
match self {
Self::Add(config, force) => {
let tier_name = config.name.clone();
candidate.add_with_deadline(config, force, deadline).await?;
candidate.add_with_deadline(*config, force, deadline).await?;
Ok(Some(tier_name))
}
Self::Edit(tier_name, credentials) => {
@@ -2990,7 +2988,7 @@ fn from_external_tier_config(name: String, ext: ExternalTierConfig) -> io::Resul
let tier_type = if wasabi_version {
TierType::Wasabi
} else {
tier_type_from_hint(ext.tier_type_hint.as_deref()).unwrap_or_else(|| match ext.tier_type {
tier_type_from_hint(ext.tier_type_hint.as_deref()).unwrap_or(match ext.tier_type {
EXTERNAL_TIER_TYPE_S3 => TierType::S3,
EXTERNAL_TIER_TYPE_AZURE => TierType::Azure,
EXTERNAL_TIER_TYPE_GCS => TierType::GCS,
@@ -3372,28 +3370,23 @@ impl TierConfigMgr {
pub async fn remove(&mut self, tier_name: &str, force: bool) -> std::result::Result<(), AdminError> {
self.ensure_generation_is_idle(tier_name)?;
let d = self.get_driver(tier_name).await;
if let Err(err) = d {
if err.code == ERR_TIER_NOT_FOUND.code {
return Ok(());
} else {
return Err(err);
}
}
let driver = match self.get_driver(tier_name).await {
Ok(driver) => driver,
Err(err) if err.code == ERR_TIER_NOT_FOUND.code => return Ok(()),
Err(err) => return Err(err),
};
if !force {
if let Ok(driver) = d {
match driver.in_use().await {
Err(err) => {
let mut e = ERR_TIER_PERM_ERR.clone();
e.message.push('.');
e.message.push_str(&err.to_string());
return Err(e);
}
Ok(in_use) if in_use => {
return Err(ERR_TIER_BACKEND_NOT_EMPTY.clone());
}
_ => {}
match driver.in_use().await {
Err(err) => {
let mut e = ERR_TIER_PERM_ERR.clone();
e.message.push('.');
e.message.push_str(&err.to_string());
return Err(e);
}
Ok(in_use) if in_use => {
return Err(ERR_TIER_BACKEND_NOT_EMPTY.clone());
}
_ => {}
}
}
self.tiers.remove(tier_name);
@@ -3402,21 +3395,12 @@ impl TierConfigMgr {
}
pub async fn verify(&mut self, tier_name: &str) -> std::result::Result<(), std::io::Error> {
let d = match self.get_driver(tier_name).await {
Ok(d) => d,
Err(err) => {
return Err(std::io::Error::other(err));
}
};
if let Err(err) = check_warm_backend(Some(d)).await {
return Err(std::io::Error::other(err));
} else {
return Ok(());
}
let driver = self.get_driver(tier_name).await.map_err(std::io::Error::other)?;
check_warm_backend(Some(driver)).await.map_err(std::io::Error::other)
}
pub fn empty(&self) -> bool {
self.list_tiers().len() == 0
self.tiers.is_empty()
}
pub fn tier_type(&self, tier_name: &str) -> String {
@@ -3429,7 +3413,7 @@ impl TierConfigMgr {
pub fn list_tiers(&self) -> Vec<TierConfig> {
let mut tier_cfgs = Vec::<TierConfig>::new();
for (_, tier) in self.tiers.iter() {
for tier in self.tiers.values() {
let tier = tier.redacted();
tier_cfgs.push(tier);
}
@@ -7135,7 +7119,8 @@ mod tests {
let err = expect_decode_err(&encode_fixture(&wrong_hint));
assert!(err.to_string().contains("inconsistent Wasabi type discriminators"), "{err}");
let poison_fields: [(&str, fn(&mut ExternalTierS3)); 6] = [
type WasabiPoisonField = (&'static str, fn(&mut ExternalTierS3));
let poison_fields: [WasabiPoisonField; 6] = [
("storage_class", |s3| s3.storage_class = "GLACIER".to_string()),
("aws_role", |s3| s3.aws_role = true),
("web_identity_token", |s3| s3.aws_role_web_identity_token_file = "/tmp/token".to_string()),
@@ -8256,7 +8241,11 @@ mod tests {
peer_calls.clone(),
Ok(PeerTierMutationState::Committed),
)],
TierConfigMgr::update_candidate_with_config_lock(&manager, store, TierCandidateMutation::Add(tier, true)),
TierConfigMgr::update_candidate_with_config_lock(
&manager,
store,
TierCandidateMutation::Add(Box::new(tier), true),
),
),
)
.await
@@ -8295,7 +8284,7 @@ mod tests {
let add = TIER_DRIVER_TEST_FACTORY.scope(
factory,
apply_tier_candidate_mutation(
TierCandidateMutation::Add(build_rustfs_tier("COLD-DEADLINE"), false),
TierCandidateMutation::Add(Box::new(build_rustfs_tier("COLD-DEADLINE")), false),
&mut candidate,
deadline,
),
@@ -9114,7 +9103,9 @@ mod tests {
fn decode_hex_fixture(hex: &str) -> Vec<u8> {
assert_eq!(hex.len() % 2, 0, "hex fixture must contain complete bytes");
hex.as_bytes()
.chunks_exact(2)
.as_chunks::<2>()
.0
.iter()
.map(|pair| {
let pair = std::str::from_utf8(pair).expect("hex fixture should be ASCII");
u8::from_str_radix(pair, 16).expect("hex fixture should contain only hexadecimal digits")
@@ -11128,7 +11119,7 @@ mod tests {
store.clone(),
candidate,
version,
TierCandidateMutation::Add(build_rustfs_tier("COLD-A"), true),
TierCandidateMutation::Add(Box::new(build_rustfs_tier("COLD-A")), true),
update,
None,
)
@@ -11725,8 +11716,9 @@ mod tests {
assert!(merged[0].has_peer_record && merged[0].has_coordinator_record);
}
let err = TierConfigMgr::merge_mutation_recovery_intents(&[committed.clone()], &[prepared.clone()])
.expect_err("a peer committed record cannot outrun the coordinator commit order");
let err =
TierConfigMgr::merge_mutation_recovery_intents(std::slice::from_ref(&committed), std::slice::from_ref(&prepared))
.expect_err("a peer committed record cannot outrun the coordinator commit order");
assert!(err.to_string().contains("conflicting states"), "{err}");
let mut conflicting_identity = prepared.clone();
@@ -13571,9 +13563,11 @@ mod tests {
let build = tokio::spawn(async move { TierConfigMgr::acquire_operation_lease(&build_manager, cold_tier).await });
barrier.arrived.notified().await;
tokio::time::timeout(Duration::from_millis(100), manager.read())
.await
.expect("cold driver construction must not block manager readers");
drop(
tokio::time::timeout(Duration::from_millis(100), manager.read())
.await
.expect("cold driver construction must not block manager readers"),
);
let tier_b = tokio::time::timeout(Duration::from_millis(100), TierConfigMgr::acquire_operation_lease(&manager, "COLD-B"))
.await
.expect("cold tier A construction must not block tier B")
@@ -13776,9 +13770,11 @@ mod tests {
let verify_manager = manager.clone();
let verify = tokio::spawn(async move { TierConfigMgr::verify_without_manager_lock(&verify_manager, "COLD-A").await });
started.notified().await;
tokio::time::timeout(Duration::from_millis(100), manager.read())
.await
.expect("slow verify must not hold the manager lock");
drop(
tokio::time::timeout(Duration::from_millis(100), manager.read())
.await
.expect("slow verify must not hold the manager lock"),
);
release.add_permits(1);
verify.await.expect("verify task should join").expect("verify should finish");
}
@@ -14186,9 +14182,11 @@ mod tests {
vec!["COLD-A".to_string()]
);
}
tokio::time::timeout(Duration::from_secs(1), manager.read())
.await
.expect("manager reads must not wait for tier A leases");
drop(
tokio::time::timeout(Duration::from_secs(1), manager.read())
.await
.expect("manager reads must not wait for tier A leases"),
);
let next_b = tokio::time::timeout(Duration::from_secs(1), TierConfigMgr::acquire_operation_lease(&manager, "COLD-B"))
.await
.expect("tier B lease acquisition must not wait for tier A")
@@ -14621,7 +14619,7 @@ mod tests {
"https://example-compat.invalid"
);
let runtime = registered_tier_driver_runtime(&manager_guard).expect("runtime sidecar should remain registered");
assert!(lock_unpoisoned(&runtime).generations.get("COLD-A").is_none());
assert!(!lock_unpoisoned(&runtime).generations.contains_key("COLD-A"));
}
#[derive(Debug)]
@@ -15261,15 +15259,12 @@ mod tests {
.filter(|object| object.bucket == bucket && object.name.starts_with(prefix))
.cloned()
.collect();
objects.sort_by(|left, right| tier_test_object_marker(left).cmp(&tier_test_object_marker(right)));
objects.sort_by_key(tier_test_object_marker);
if marker.is_some() || version_marker.is_some() {
let marker = (marker.unwrap_or_default(), version_marker.unwrap_or_default());
objects.retain(|object| tier_test_object_marker(object) > marker);
}
let limit = match usize::try_from(max_keys) {
Ok(limit) => limit,
Err(_) => 0,
};
let limit: usize = usize::try_from(max_keys).unwrap_or_default();
let is_truncated = objects.len() > limit;
if is_truncated {
objects.truncate(limit);
@@ -15299,17 +15294,16 @@ mod tests {
result: Self::WalkResultSender,
opts: Self::WalkOptions,
) -> Result<()> {
if self.fail_reference_walk.load(Ordering::SeqCst) {
if result
if self.fail_reference_walk.load(Ordering::SeqCst)
&& result
.send(StorageObjectInfoOrErr {
item: None,
err: Some(Error::other("injected tier reference walk failure")),
})
.await
.is_err()
{
return Ok(());
}
{
return Ok(());
}
let mut objects = self
.listed_versions
@@ -15320,7 +15314,7 @@ mod tests {
.filter(|object| opts.include_free_versions || !object.transitioned_object.free_version)
.cloned()
.collect::<Vec<_>>();
objects.sort_by(|left, right| tier_test_object_marker(left).cmp(&tier_test_object_marker(right)));
objects.sort_by_key(tier_test_object_marker);
if let Some(marker) = opts.marker.as_deref() {
objects.retain(|object| object.name.as_str() > marker);
}
@@ -15498,17 +15492,18 @@ mod tests {
api_view.rustfs.expect("admin RustFS payload should exist").secret_key,
TIER_CREDENTIAL_REDACTED
);
let observed = lock_unpoisoned(&observed);
assert_eq!(observed.len(), 1);
assert_eq!(
observed[0]
.rustfs
.as_ref()
.expect("backend factory should observe the RustFS payload")
.secret_key,
SECRET_KEY
);
drop(observed);
{
let observed = lock_unpoisoned(&observed);
assert_eq!(observed.len(), 1);
assert_eq!(
observed[0]
.rustfs
.as_ref()
.expect("backend factory should observe the RustFS payload")
.secret_key,
SECRET_KEY
);
}
let operations = backend.op_log().await;
assert_eq!(operations.len(), 5);
@@ -16709,7 +16704,7 @@ mod tests {
candidate.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A"));
candidate.tiers.insert("COLD-B".to_string(), build_rustfs_tier("COLD-B"));
let targets = TierCandidateMutation::Add(build_rustfs_tier("COLD-B"), true)
let targets = TierCandidateMutation::Add(Box::new(build_rustfs_tier("COLD-B")), true)
.affected_targets(&current, &candidate)
.expect("add proof should ignore unchanged durable tiers");
assert_eq!(targets.len(), 1);
@@ -16735,7 +16730,7 @@ mod tests {
TierConfigMgr::update_candidate_with_config_lock(
&manager,
store.clone(),
TierCandidateMutation::Add(build_rustfs_tier("COLD-B"), true),
TierCandidateMutation::Add(Box::new(build_rustfs_tier("COLD-B")), true),
),
)
.await
@@ -16794,14 +16789,15 @@ mod tests {
.await
.expect("legacy nested-name Add must run the full coordinator fanout");
let prepared_intents = lock_unpoisoned(&prepared_intents);
assert_eq!(prepared_intents.len(), 1);
assert_eq!(prepared_intents[0].kind, TierMutationIntentKind::Add);
assert_eq!(prepared_intents[0].affected_targets.len(), 1);
assert_eq!(prepared_intents[0].affected_targets[0].tier_name, "COLD-LEGACY");
assert!(prepared_intents[0].affected_targets[0].old_backend_identity.is_none());
assert!(prepared_intents[0].affected_targets[0].new_backend_identity.is_some());
drop(prepared_intents);
{
let prepared_intents = lock_unpoisoned(&prepared_intents);
assert_eq!(prepared_intents.len(), 1);
assert_eq!(prepared_intents[0].kind, TierMutationIntentKind::Add);
assert_eq!(prepared_intents[0].affected_targets.len(), 1);
assert_eq!(prepared_intents[0].affected_targets[0].tier_name, "COLD-LEGACY");
assert!(prepared_intents[0].affected_targets[0].old_backend_identity.is_none());
assert!(prepared_intents[0].affected_targets[0].new_backend_identity.is_some());
}
let peer_calls = lock_unpoisoned(&peer_calls).clone();
let prepare_index = peer_calls
@@ -16883,7 +16879,7 @@ mod tests {
let err = TierConfigMgr::update_candidate_with_config_lock(
&manager,
store.clone(),
TierCandidateMutation::Add(build_rustfs_tier("COLD-B"), true),
TierCandidateMutation::Add(Box::new(build_rustfs_tier("COLD-B")), true),
)
.await
.expect_err("a new tier config update must wait for pending mutation recovery");
@@ -17159,7 +17155,7 @@ mod tests {
TierConfigMgr::update_candidate_with_config_lock(
&update_manager,
update_store,
TierCandidateMutation::Add(build_rustfs_tier("COLD-A"), true),
TierCandidateMutation::Add(Box::new(build_rustfs_tier("COLD-A")), true),
),
)
.await
@@ -17210,7 +17206,7 @@ mod tests {
TierConfigMgr::update_candidate_with_config_lock(
&update_manager,
update_store,
TierCandidateMutation::Add(build_rustfs_tier("COLD-A"), true),
TierCandidateMutation::Add(Box::new(build_rustfs_tier("COLD-A")), true),
),
),
)
@@ -17273,7 +17269,7 @@ mod tests {
TierConfigMgr::prevalidate_candidate_owned(
empty_mgr(),
None,
TierCandidateMutation::Add(build_rustfs_tier("COLD-A"), true),
TierCandidateMutation::Add(Box::new(build_rustfs_tier("COLD-A")), true),
),
)
.await;
@@ -17937,7 +17933,7 @@ mod tests {
#[tokio::test]
async fn tier_add_succeeds_with_refresh_during_coordinator_commit() {
assert_coordinator_commit_refresh_succeeds(TierCandidateMutation::Add(build_rustfs_tier("COLD-A"), true)).await;
assert_coordinator_commit_refresh_succeeds(TierCandidateMutation::Add(Box::new(build_rustfs_tier("COLD-A")), true)).await;
}
#[tokio::test]
@@ -15,8 +15,6 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use crate::error::is_err_bucket_not_found;
#[cfg(feature = "gcs")]
@@ -719,17 +717,7 @@ async fn check_warm_backend_with_deadlines(
if !matches!(cleanup_result, Ok(Ok(()))) {
return Err(probe_cleanup_incomplete_error());
}
if let Err(err) = read_result {
//if is_err_bucket_not_found(&err) {
// return Err(ERR_TIER_BUCKET_NOT_FOUND);
//}
/*else if is_err_signature_does_not_match(err) {
return Err(ERR_TIER_MISSING_CREDENTIALS);
}*/
//else {
return Err(err);
//}
}
read_result?;
Ok(())
}
@@ -759,7 +747,7 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
message: format!("Unable to setup remote tier, check tier configuration: {err}"),
status_code: StatusCode::BAD_REQUEST,
});
}
@@ -800,7 +788,7 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
message: format!("Unable to setup remote tier, check tier configuration: {err}"),
status_code: StatusCode::BAD_REQUEST,
});
}
@@ -820,7 +808,7 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
message: format!("Unable to setup remote tier, check tier configuration: {err}"),
status_code: StatusCode::BAD_REQUEST,
});
}
@@ -840,7 +828,7 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
message: format!("Unable to setup remote tier, check tier configuration: {err}"),
status_code: StatusCode::BAD_REQUEST,
});
}
@@ -860,7 +848,7 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
message: format!("Unable to setup remote tier, check tier configuration: {err}"),
status_code: StatusCode::BAD_REQUEST,
});
}
@@ -880,7 +868,7 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
message: format!("Unable to setup remote tier, check tier configuration: {err}"),
status_code: StatusCode::BAD_REQUEST,
});
}
@@ -900,7 +888,7 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
message: format!("Unable to setup remote tier, check tier configuration: {err}"),
status_code: StatusCode::BAD_REQUEST,
});
}
@@ -929,7 +917,7 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
message: format!("Unable to setup remote tier, check tier configuration: {err}"),
status_code: StatusCode::BAD_REQUEST,
});
}
@@ -949,7 +937,7 @@ pub async fn new_warm_backend(tier: &TierConfig, probe: bool) -> Result<WarmBack
warn!("{}", err);
return Err(AdminError {
code: "XRustFSAdminTierInvalidConfig".to_string(),
message: format!("Unable to setup remote tier, check tier configuration: {}", err.to_string()),
message: format!("Unable to setup remote tier, check tier configuration: {err}"),
status_code: StatusCode::BAD_REQUEST,
});
}
@@ -15,8 +15,6 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use std::collections::HashMap;
use std::sync::Arc;
@@ -106,39 +104,35 @@ impl WarmBackendS3 {
};
validate_outbound_url(&u).map_err(|err| std::io::Error::other(format!("tier endpoint is not allowed: {err}")))?;
if conf.aws_role_web_identity_token_file == "" && conf.aws_role_arn != ""
|| conf.aws_role_web_identity_token_file != "" && conf.aws_role_arn == ""
{
let has_web_identity_token_file = !conf.aws_role_web_identity_token_file.is_empty();
let has_role_arn = !conf.aws_role_arn.is_empty();
let has_access_key = !conf.access_key.is_empty();
let has_secret_key = !conf.secret_key.is_empty();
if has_web_identity_token_file != has_role_arn {
return Err(std::io::Error::other("both the token file and the role ARN are required"));
} else if conf.access_key == "" && conf.secret_key != "" || conf.access_key != "" && conf.secret_key == "" {
} else if has_access_key != has_secret_key {
return Err(std::io::Error::other("both the access and secret keys are required"));
} else if conf.aws_role
&& (conf.aws_role_web_identity_token_file != ""
|| conf.aws_role_arn != ""
|| conf.access_key != ""
|| conf.secret_key != "")
{
} else if conf.aws_role && (has_web_identity_token_file || has_role_arn || has_access_key || has_secret_key) {
return Err(std::io::Error::other(
"AWS Role cannot be activated with static credentials or the web identity token file",
));
} else if conf.bucket == "" {
} else if conf.bucket.is_empty() {
return Err(std::io::Error::other("no bucket name was provided"));
}
let creds: Credentials<Static>;
if conf.access_key != "" && conf.secret_key != "" {
let creds = if has_access_key && has_secret_key {
//creds = Credentials::new_static_v4(conf.access_key, conf.secret_key, "");
creds = Credentials::new(Static(Value {
Credentials::new(Static(Value {
access_key_id: conf.access_key.clone(),
secret_access_key: conf.secret_key.clone(),
session_token: "".to_string(),
signer_type: SignatureType::SignatureV4,
..Default::default()
}));
}))
} else {
return Err(std::io::Error::other("insufficient parameters for S3 backend authentication"));
}
};
let timeouts = transition_client_timeouts_from_env();
let opts = Options {
creds,
@@ -162,11 +156,11 @@ impl WarmBackendS3 {
}
pub fn get_dest(&self, object: &str) -> String {
let mut dest_obj = object.to_string();
if self.prefix != "" {
dest_obj = format!("{}/{}", &self.prefix, object);
if self.prefix.is_empty() {
object.to_string()
} else {
format!("{}/{}", self.prefix, object)
}
return dest_obj;
}
pub(crate) async fn remove_with_result(&self, object: &str, rv: &str) -> Result<RemoveObjectResult, std::io::Error> {
@@ -413,6 +407,10 @@ impl TransitionCandidateVersions {
}
#[cfg(test)]
#[allow(
clippy::items_after_test_module,
reason = "keep parsing tests adjacent to the helpers they cover"
)]
mod tests {
use super::*;
use rustfs_s3_client::api_s3_datatypes::{ListVersionsResult, Version};
@@ -917,7 +915,7 @@ impl WarmBackend for WarmBackendS3 {
.list_objects_v2(&self.bucket, &self.prefix, "", "", SLASH_SEPARATOR, 1)
.await?;
Ok(result.common_prefixes.len() > 0 || result.contents.len() > 0)
Ok(!result.common_prefixes.is_empty() || !result.contents.is_empty())
}
}
+2 -59
View File
@@ -692,15 +692,10 @@ impl FileMeta {
}
}
// The version stays on disk while the purge replicates
// (status PENDING/FAILED); its data dir must stay with
// it. Returning the dir here made the disk layer delete
// it, which turned every non-inline retained version
// into an unreadable zombie: the purge state could never
// be applied and the bucket could never be deleted.
let old_dir = v.object.as_ref().map(|v| v.data_dir).unwrap_or_default();
self.set_idx(i, v)?;
return Ok(None);
return Ok(old_dir);
}
found_index = Some(i);
}
@@ -2707,58 +2702,6 @@ mod test {
);
}
/// Regression for rustfs/backlog#2340: a version purge that still awaits
/// the replication target keeps the object version on disk with a pending
/// purge status. Its data dir must be retained with it; handing the dir
/// back here made the disk layer delete it, leaving every non-inline
/// retained version unreadable. The dir is released only once the purge
/// completes and the version itself goes away.
#[test]
fn delete_version_pending_version_purge_retains_object_data_dir() {
let version_id = Uuid::new_v4();
let data_dir = Uuid::new_v4();
let mut fm = FileMeta::new();
let mut fi = FileInfo::new("object", 2, 2);
fi.version_id = Some(version_id);
fi.data_dir = Some(data_dir);
fi.mod_time = Some(OffsetDateTime::now_utc());
fm.add_version(fi).unwrap();
let pending_purge = FileInfo {
name: "object".to_string(),
version_id: Some(version_id),
mark_deleted: true,
replication_state_internal: Some(ReplicationState {
version_purge_status_internal: Some("target=PENDING;".to_string()),
purge_targets: version_purge_statuses_map("target=PENDING;"),
..Default::default()
}),
..Default::default()
};
let freed = fm.delete_version(&pending_purge).unwrap();
assert_eq!(freed, None, "a pending purge must not release the retained version's data dir");
assert_eq!(fm.versions.len(), 1, "the version must stay until the purge replicates");
let retained = fm
.into_fileinfo("vol", "object", &version_id.to_string(), false, false, true)
.unwrap();
assert_eq!(retained.data_dir, Some(data_dir));
assert_eq!(retained.version_purge_status(), VersionPurgeStatusType::Pending);
let completed_purge = FileInfo {
name: "object".to_string(),
version_id: Some(version_id),
replication_state_internal: Some(ReplicationState {
version_purge_status_internal: Some("target=COMPLETE;".to_string()),
purge_targets: version_purge_statuses_map("target=COMPLETE;"),
..Default::default()
}),
..Default::default()
};
let freed = fm.delete_version(&completed_purge).unwrap();
assert_eq!(freed, Some(data_dir), "a completed purge removes the version and releases its data dir");
assert!(fm.versions.is_empty());
}
#[test]
fn delete_version_accepts_delete_only_marker_and_free_version_paths() {
let marker_version_id = Uuid::new_v4();
@@ -62,7 +62,6 @@ Object keys are stored as file-system paths under each drive (`{drive}/{bucket}/
| Behavior | RustFS | AWS S3 | Why |
|---|---|---|---|
| Object key with a `.` or `..` path segment, or an empty segment (`//`), such as `a//b/./c/../d` | `400 InvalidArgument` (`check_object_args` in `crates/ecstore/src/bucket/utils.rs`, mirroring MinIO `IsValidObjectPrefix`) | Accepted as an opaque key | A `..` segment would resolve to a parent directory and `.`/`//` segments would alias other keys on disk; encoding them would change the MinIO-compatible on-disk format. |
| Directory marker (key ending in `/`, with or without a body) in a versioned bucket | Stored as the null version: `PutObject`/`HeadObject` report version id `00000000-0000-0000-0000-000000000000`, `ListObjectVersions` reports `null`, and a later PUT of the same key overwrites in place (`put_opts` in `rustfs/src/storage/options.rs`, mirroring MinIO `putOpts`: "for directory objects skip creating new versions") | A real version id per PUT, with a version history | The marker only exists to make an empty prefix listable; keeping a history for it would leave hidden versions behind every prefix delete. Replication still copies the marker as its null version (`test_bucket_replication_replicates_directory_marker_in_versioned_bucket` in `crates/e2e_test/src/replication_extension_test.rs`). |
## Update Rule
-6
View File
@@ -69,12 +69,8 @@ crates/s3-client/src/transition_api.rs|clippy::all
crates/s3-client/src/transition_api.rs|unused_must_use
crates/s3-client/src/transition_api.rs|unused_variables
crates/ecstore/src/services/event_notification.rs|unused_variables
crates/ecstore/src/services/tier/tier.rs|clippy::all
crates/ecstore/src/services/tier/tier.rs|unused_must_use
crates/ecstore/src/services/tier/tier.rs|unused_variables
crates/ecstore/src/services/tier/tier_admin.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_aliyun.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_azure.rs|unused_variables
@@ -83,7 +79,5 @@ crates/ecstore/src/services/tier/warm_backend_huaweicloud.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_minio.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_r2.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_rustfs.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_s3.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend_s3.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend_s3.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_tencent.rs|unused_variables