diff --git a/.config/nextest.toml b/.config/nextest.toml index 804888924..b4a1fbe54 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -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 diff --git a/crates/ecstore/src/services/tier/warm_backend.rs b/crates/ecstore/src/services/tier/warm_backend.rs index 8e7a4c8f4..6ae923b44 100644 --- a/crates/ecstore/src/services/tier/warm_backend.rs +++ b/crates/ecstore/src/services/tier/warm_backend.rs @@ -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::()) 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>>, } + struct AuthRejectingProbePutBackend { + status_code: StatusCode, + code: S3ErrorCode, + probes: Arc, + removes: Arc, + } + #[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 { + 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, + ) -> Result { + self.put(object, r, length).await + } + + async fn get(&self, _object: &str, _rv: &str, _opts: WarmBackendGetOpts) -> Result { + 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 { + 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 { + 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(); diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index 385db4dda..0b3683655 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -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)]), }; diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index caa7e36b6..ccd2135fd 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -2056,6 +2056,9 @@ fn data_movement_delete_marker_metadata_identity(metadata: &HashMap Option { 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);