fix(tier): classify invalid remote credentials (#7753)

* fix(tier): classify invalid remote credentials

Classify deterministic remote-tier auth rejections before the probe cleanup fallback, and preserve the admin 4xx status for invalid credential responses.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* fix(ecstore): stabilize delete-marker movement regressions

Ignore target-local bucket incarnation fence metadata when comparing data-movement delete-marker identities, keep the release null-version listing fixture valid, and fix the nextest decommission-family filter so fault-hook tests actually join the serial group.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

---------

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
Hauser
2026-09-14 01:04:40 +08:00
committed by GitHub
parent 8f6d06db21
commit 0900d78877
5 changed files with 161 additions and 6 deletions
+2 -2
View File
@@ -152,7 +152,7 @@ test-group = 'ecstore-serial-flaky'
# deterministic commit barriers. Keep the whole init decommission family in one
# nextest group; serial_test alone cannot isolate separate test processes.
[[profile.default.overrides]]
filter = 'package(rustfs-ecstore) & test(/^store::init::tests::(decommission_|suspended_.*decommission)$/)'
filter = 'package(rustfs-ecstore) & test(/^store::init::tests::(decommission_.*|suspended_.*decommission.*)$/)'
test-group = 'ecstore-serial-flaky'
# Serialize the bucket-incarnation / lifecycle-fence tests. They drive
@@ -323,7 +323,7 @@ filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_reg
test-group = 'ecstore-serial-flaky'
[[profile.ci.overrides]]
filter = 'package(rustfs-ecstore) & test(/^store::init::tests::(decommission_|suspended_.*decommission)$/)'
filter = 'package(rustfs-ecstore) & test(/^store::init::tests::(decommission_.*|suspended_.*decommission.*)$/)'
test-group = 'ecstore-serial-flaky'
# Serialize the bucket-incarnation / lifecycle-fence tests under the ci profile
@@ -22,7 +22,7 @@ use crate::services::tier::warm_backend_gcs::WarmBackendGCS;
use crate::services::tier::{
tier::{ERR_TIER_BACKEND_IN_USE, ERR_TIER_INVALID_CONFIG, ERR_TIER_TYPE_UNSUPPORTED},
tier_config::{TierConfig, TierType},
tier_handlers::{ERR_TIER_BUCKET_NOT_FOUND, ERR_TIER_NOT_FOUND, ERR_TIER_PERM_ERR},
tier_handlers::{ERR_TIER_BUCKET_NOT_FOUND, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_NOT_FOUND, ERR_TIER_PERM_ERR},
warm_backend_aliyun::WarmBackendAliyun,
warm_backend_azure::WarmBackendAzure,
warm_backend_huaweicloud::WarmBackendHuaweicloud,
@@ -39,7 +39,7 @@ use rustfs_s3_client::credentials::{Credentials, SignatureType, Static, Value};
use rustfs_s3_client::transition_api::{BucketLookupType, Options, TransitionClient, TransitionClientTimeouts, TransitionCore};
use rustfs_s3_client::{
admin_handler_utils::AdminError,
api_error_response::to_error_response,
api_error_response::{ErrorResponse, to_error_response},
api_put_object::{AdvancedPutOptions, PutObjectOptions},
transition_api::{ReadCloser, ReaderImpl},
};
@@ -669,6 +669,17 @@ fn probe_cleanup_incomplete_error() -> AdminError {
err
}
fn is_remote_tier_auth_rejection(err: &std::io::Error) -> bool {
let Some(response) = err.get_ref().and_then(|err| err.downcast_ref::<ErrorResponse>()) else {
return false;
};
matches!(response.status_code, StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN)
|| matches!(
response.code,
S3ErrorCode::AccessDenied | S3ErrorCode::InvalidAccessKeyId | S3ErrorCode::SignatureDoesNotMatch
)
}
async fn check_warm_backend_with_deadlines(
w: Option<&WarmBackendImpl>,
deadline: tokio::time::Instant,
@@ -689,6 +700,9 @@ async fn check_warm_backend_with_deadlines(
tokio::time::timeout_at(deadline, w.put(&probe_object, ReaderImpl::Body(Bytes::from_static(b"RustFS")), 6)).await;
let remote_version_id = match put_result {
Ok(Ok(remote_version_id)) => remote_version_id,
Ok(Err(err)) if is_remote_tier_auth_rejection(&err) => {
return Err(ERR_TIER_INVALID_CREDENTIALS.clone());
}
Ok(Err(_)) => {
return Err(match compensate_uncertain_probe_put(w, &probe_object, cleanup_deadline).await {
Ok(()) => ERR_TIER_PERM_ERR.clone(),
@@ -1285,6 +1299,13 @@ mod tests {
removed_versions: Arc<tokio::sync::Mutex<Vec<String>>>,
}
struct AuthRejectingProbePutBackend {
status_code: StatusCode,
code: S3ErrorCode,
probes: Arc<AtomicUsize>,
removes: Arc<AtomicUsize>,
}
#[derive(Clone, Copy)]
enum ProbeBody {
Exact,
@@ -1520,6 +1541,46 @@ mod tests {
}
}
#[async_trait::async_trait]
impl WarmBackend for AuthRejectingProbePutBackend {
async fn put(&self, _object: &str, _r: ReaderImpl, _length: i64) -> Result<String, std::io::Error> {
Err(std::io::Error::other(rustfs_s3_client::api_error_response::ErrorResponse {
status_code: self.status_code,
code: self.code.clone(),
message: "remote credentials rejected".to_string(),
..Default::default()
}))
}
async fn put_with_meta(
&self,
object: &str,
r: ReaderImpl,
length: i64,
_meta: HashMap<String, String>,
) -> Result<String, std::io::Error> {
self.put(object, r, length).await
}
async fn get(&self, _object: &str, _rv: &str, _opts: WarmBackendGetOpts) -> Result<ReadCloser, std::io::Error> {
Err(std::io::Error::other("GET must not run after an auth-rejected probe PUT"))
}
async fn remove(&self, _object: &str, _rv: &str) -> Result<(), std::io::Error> {
self.removes.fetch_add(1, Ordering::SeqCst);
Err(std::io::Error::other("remove must not run after a deterministic auth rejection"))
}
async fn probe_transition_candidate(&self, _object: &str) -> Result<TransitionCandidateProbe, std::io::Error> {
self.probes.fetch_add(1, Ordering::SeqCst);
Err(std::io::Error::other("probe must not run after a deterministic auth rejection"))
}
async fn in_use(&self) -> Result<bool, std::io::Error> {
Ok(false)
}
}
#[tokio::test]
async fn check_warm_backend_validates_before_probe_io() {
let validations = Arc::new(AtomicUsize::new(0));
@@ -1700,6 +1761,34 @@ mod tests {
assert_eq!(removed_versions.lock().await.as_slice(), [PROBE_VERSION]);
}
#[tokio::test]
async fn check_warm_backend_reports_auth_rejected_probe_put_as_invalid_credentials() {
for (status_code, code) in [
(StatusCode::FORBIDDEN, S3ErrorCode::AccessDenied),
(StatusCode::UNAUTHORIZED, S3ErrorCode::Custom("".into())),
(StatusCode::OK, S3ErrorCode::InvalidAccessKeyId),
(StatusCode::OK, S3ErrorCode::SignatureDoesNotMatch),
] {
let probes = Arc::new(AtomicUsize::new(0));
let removes = Arc::new(AtomicUsize::new(0));
let backend: WarmBackendImpl = Box::new(AuthRejectingProbePutBackend {
status_code,
code,
probes: probes.clone(),
removes: removes.clone(),
});
let err = check_warm_backend(Some(&backend))
.await
.expect_err("a deterministic remote-auth rejection should not become an uncertain probe error");
assert_eq!(err.code, ERR_TIER_INVALID_CREDENTIALS.code);
assert_eq!(err.status_code, StatusCode::BAD_REQUEST);
assert_eq!(probes.load(Ordering::SeqCst), 0);
assert_eq!(removes.load(Ordering::SeqCst), 0);
}
}
#[tokio::test(start_paused = true)]
async fn check_warm_backend_reconciles_a_lost_put_response() {
let backend = MockWarmBackend::new();
+31 -1
View File
@@ -7330,6 +7330,36 @@ mod test {
}
}
fn test_null_version_meta_entry(name: &str, mod_time: time::OffsetDateTime) -> MetaCacheEntry {
let mut fi = FileInfo::new(name, 2, 2);
fi.erasure.index = 1;
fi.data_dir = Some(Uuid::from_u128(0x5678));
fi.volume = "bucket".to_owned();
fi.name = name.to_owned();
fi.version_id = None;
fi.versioned = false;
fi.size = 1;
fi.parts = vec![ObjectPartInfo {
number: 1,
size: 1,
actual_size: 1,
..Default::default()
}];
fi.mod_time = Some(mod_time);
fi.metadata.insert("etag".to_string(), "null-etag".to_string());
let mut meta = FileMeta::new();
meta.add_version(fi).expect("test metadata should accept null object version");
let metadata = meta.marshal_msg().expect("test metadata should marshal");
MetaCacheEntry {
name: name.to_owned(),
metadata,
cached: Some(meta),
reusable: false,
}
}
#[test]
fn fallback_entries_for_object_filters_claimed_physical_disks() {
let mut entries = FallbackListingEntries::new();
@@ -7738,7 +7768,7 @@ mod test {
};
let entry = match kind {
"deletes" => test_delete_marker_meta_entry(&name, mod_time),
"null" => test_object_meta_entry(&name),
"null" => test_null_version_meta_entry(&name, mod_time),
"mixed" => test_object_with_delete_marker_meta_entry(&name, mod_time, mod_time + time::Duration::SECOND),
_ => test_object_meta_entry_with_erasure_versions(&name, &[(mod_time, "etag", 2, 2)]),
};
+21
View File
@@ -2056,6 +2056,9 @@ fn data_movement_delete_marker_metadata_identity(metadata: &HashMap<String, Stri
local_tier_free_version_id = Some(version_id);
continue;
}
if suffix.eq_ignore_ascii_case(rustfs_utils::http::SUFFIX_BUCKET_INCARNATION_ID) {
continue;
}
let canonical_suffix = [
rustfs_utils::http::SUFFIX_REPLICA_TIMESTAMP,
@@ -7110,6 +7113,24 @@ mod tests {
assert!(!is_equivalent_data_movement_delete_marker(&source, &target));
}
#[test]
fn equivalent_data_movement_delete_marker_ignores_target_bucket_incarnation_fence() {
let source = ObjectInfo {
version_id: Some(Uuid::from_u128(1)),
delete_marker: true,
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
};
let mut target = source.clone();
rustfs_utils::http::insert_str(
Arc::make_mut(&mut target.user_defined),
rustfs_utils::http::SUFFIX_BUCKET_INCARNATION_ID,
Uuid::from_u128(2).to_string(),
);
assert!(is_equivalent_data_movement_delete_marker(&source, &target));
}
#[test]
fn data_movement_delete_marker_source_requires_persisted_mod_time() {
let source = ObjectInfo {
+16 -1
View File
@@ -229,6 +229,12 @@ fn tier_backend_error_response(err: &AdminError) -> Option<S3Error> {
Some(response)
}
fn tier_admin_custom_error_response(err: &AdminError) -> S3Error {
let mut response = S3Error::with_message(S3ErrorCode::Custom(err.code.clone().into()), err.message.clone());
response.set_status_code(err.status_code);
response
}
fn clear_tier_error_response(err: &AdminError) -> S3Error {
S3Error::with_message(S3ErrorCode::Custom("TierClearFailed".into()), format!("tier clear failed. {err}"))
}
@@ -406,7 +412,7 @@ impl Operation for AddTier {
} else if err.code == ERR_TIER_INVALID_CONFIG.code {
Err(S3Error::with_message(S3ErrorCode::InvalidArgument, err.message))
} else if err.code == ERR_TIER_INVALID_CREDENTIALS.code {
Err(S3Error::with_message(S3ErrorCode::Custom(err.code.clone().into()), err.message))
Err(tier_admin_custom_error_response(&err))
} else {
warn!(
event = EVENT_ADMIN_TIER_STATE,
@@ -1399,6 +1405,15 @@ mod tests {
assert!(tier_backend_error_response(&unknown).is_none());
}
#[test]
fn tier_admin_custom_error_response_preserves_admin_status() {
let response = tier_admin_custom_error_response(&ERR_TIER_INVALID_CREDENTIALS);
assert_eq!(response.code(), &S3ErrorCode::Custom("XRustFSAdminTierInvalidCredentials".into()));
assert_eq!(response.message(), Some("Invalid remote tier credentials"));
assert_eq!(response.status_code(), Some(StatusCode::BAD_REQUEST));
}
#[test]
fn clear_preserves_legacy_backend_not_empty_error_code() {
let response = clear_tier_error_response(&ERR_TIER_BACKEND_NOT_EMPTY);