test(e2e): make security boundary oracles fail closed (#6542)

This commit is contained in:
Zhengchao An
2026-08-25 21:19:28 +08:00
committed by GitHub
parent 9a89434644
commit bcfed065c1
3 changed files with 64 additions and 34 deletions
+55 -30
View File
@@ -21,11 +21,13 @@
//! - SSRF prevention (internal/private endpoints rejected for tiering) //! - SSRF prevention (internal/private endpoints rejected for tiering)
//! - Race condition handling (concurrent writes converge without corruption) //! - Race condition handling (concurrent writes converge without corruption)
use crate::common::{RustFSTestEnvironment, awscurl_put, init_logging, require_awscurl}; use crate::common::{RustFSTestEnvironment, init_logging, signed_s3_request};
use aws_sdk_s3::error::ProvideErrorMetadata; use aws_sdk_s3::error::ProvideErrorMetadata;
use aws_sdk_s3::primitives::ByteStream; use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart, Tag, Tagging}; use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart, Tag, Tagging};
use std::error::Error; use std::error::Error;
use std::time::Duration;
use tokio::net::TcpListener;
/// Oversized tagging payloads must be rejected by the per-object tag limit. /// Oversized tagging payloads must be rejected by the per-object tag limit.
/// ///
@@ -74,8 +76,12 @@ async fn test_large_xml_body_rejection() -> Result<(), Box<dyn Error + Send + Sy
let _ = client.delete_object().bucket(&bucket_name).key(object_key).send().await; let _ = client.delete_object().bucket(&bucket_name).key(object_key).send().await;
let _ = client.delete_bucket().bucket(&bucket_name).send().await; let _ = client.delete_bucket().bucket(&bucket_name).send().await;
assert!(result.is_err(), "Server must reject an oversized tagging payload, but it was accepted"); let err = result.expect_err("Server must reject an oversized tagging payload");
let err = result.expect_err("checked is_err above"); assert_eq!(
err.raw_response().map(|response| response.status().as_u16()),
Some(400),
"Oversized tagging should return HTTP 400, got: {err:?}"
);
let code = err.as_service_error().and_then(|e| e.code()); let code = err.as_service_error().and_then(|e| e.code());
assert_eq!( assert_eq!(
code, code,
@@ -146,6 +152,11 @@ async fn test_excessive_multipart_parts() -> Result<(), Box<dyn Error + Send + S
let _ = client.delete_bucket().bucket(&bucket_name).send().await; let _ = client.delete_bucket().bucket(&bucket_name).send().await;
let error = result.expect_err("server should reject a completion part number above 10000"); let error = result.expect_err("server should reject a completion part number above 10000");
assert_eq!(
error.raw_response().map(|response| response.status().as_u16()),
Some(400),
"part-limit rejection must return HTTP 400, got: {error:?}"
);
let service_error = error let service_error = error
.as_service_error() .as_service_error()
.expect("part-limit rejection must be an S3 service error"); .expect("part-limit rejection must be an S3 service error");
@@ -219,15 +230,14 @@ async fn test_concurrent_object_operations() -> Result<(), Box<dyn Error + Send
// After deletion, the object must be absent. // After deletion, the object must be absent.
client.delete_object().bucket(&bucket_name).key(key).send().await?; client.delete_object().bucket(&bucket_name).key(key).send().await?;
let get_after_delete = client.get_object().bucket(&bucket_name).key(key).send().await; let get_after_delete = client.get_object().bucket(&bucket_name).key(key).send().await;
assert!(get_after_delete.is_err(), "Object must be absent after delete"); let error = get_after_delete.expect_err("Object must be absent after delete");
let code = get_after_delete
.err()
.and_then(|e| e.as_service_error().and_then(|se| se.code()).map(str::to_string));
assert_eq!( assert_eq!(
code.as_deref(), error.raw_response().map(|response| response.status().as_u16()),
Some("NoSuchKey"), Some(404),
"GET after delete should return NoSuchKey, got {code:?}" "GET after delete should return HTTP 404, got: {error:?}"
); );
let code = error.as_service_error().and_then(|service_error| service_error.code());
assert_eq!(code, Some("NoSuchKey"), "GET after delete should return NoSuchKey, got: {error:?}");
let _ = client.delete_bucket().bucket(&bucket_name).send().await; let _ = client.delete_bucket().bucket(&bucket_name).send().await;
@@ -235,35 +245,31 @@ async fn test_concurrent_object_operations() -> Result<(), Box<dyn Error + Send
Ok(()) Ok(())
} }
/// Internal/private endpoints must be rejected as remote tier backends (SSRF). /// Internal/private endpoints must be rejected as remote tier backends (SSRF)
/// before RustFS opens a connection to them.
/// ///
/// This issues a real admin AddTier call (`PUT /rustfs/admin/v3/tier`) for each /// This issues a real admin AddTier call (`PUT /rustfs/admin/v3/tier`) for each
/// internal/private endpoint and asserts the server rejects it (non-2xx, so the /// internal/private endpoint and requires the exact validation response. A
/// signed request helper returns an error). An internal endpoint must never be /// loopback listener additionally proves the literal loopback case is rejected
/// accepted as a tier backend. The rejection may originate from explicit /// without an outbound connection; a generic connectivity failure is not
/// SSRF/internal-address filtering or from the backend connectivity/credential /// sufficient evidence of SSRF protection.
/// validation performed during AddTier; either way the security-relevant
/// outcome — the internal endpoint is not accepted — is asserted here.
///
/// The admin API is exercised via signed `awscurl` requests, matching the
/// pattern used by the other admin-API E2E tests in this crate. The full E2E
/// lane installs and verifies the pinned `awscurl` prerequisite.
#[tokio::test] #[tokio::test]
async fn test_tiering_url_validation() -> Result<(), Box<dyn Error + Send + Sync>> { async fn test_tiering_url_validation() -> Result<(), Box<dyn Error + Send + Sync>> {
init_logging(); init_logging();
require_awscurl()?;
let mut env = RustFSTestEnvironment::new().await?; let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?; env.start_rustfs_server(vec![]).await?;
let tier_url = format!("{}/rustfs/admin/v3/tier", env.url); let tier_url = format!("{}/rustfs/admin/v3/tier", env.url);
let sentinel = TcpListener::bind("127.0.0.1:0").await?;
let sentinel_endpoint = format!("http://{}", sentinel.local_addr()?);
let internal_endpoints = [ let internal_endpoints = [
"http://127.0.0.1:8080", (sentinel_endpoint.as_str(), true),
"http://localhost:8080", ("http://localhost:8080", false),
"http://169.254.169.254", // cloud instance metadata endpoint ("http://169.254.169.254", false), // cloud instance metadata endpoint
"http://[::1]:8080", ("http://[::1]:8080", false),
]; ];
for endpoint in internal_endpoints { for (endpoint, checks_connection_attempt) in internal_endpoints {
// AddTier expects an uppercase tier name and a backend configuration. // AddTier expects an uppercase tier name and a backend configuration.
let body = serde_json::json!({ let body = serde_json::json!({
"type": "s3", "type": "s3",
@@ -280,11 +286,30 @@ async fn test_tiering_url_validation() -> Result<(), Box<dyn Error + Send + Sync
}) })
.to_string(); .to_string();
let result = awscurl_put(&tier_url, &body, &env.access_key, &env.secret_key).await; let response = signed_s3_request(
http::Method::PUT,
&tier_url,
Some(body),
Some("application/json"),
&env.access_key,
&env.secret_key,
)
.await?;
let status = response.status();
let response_body = response.text().await?;
assert_eq!(status, 400, "AddTier must reject internal endpoint {endpoint} with HTTP 400");
assert!( assert!(
result.is_err(), response_body.contains("<Code>InvalidArgument</Code>") && response_body.contains("tier endpoint is not allowed"),
"AddTier must reject internal endpoint {endpoint}, but it was accepted: {result:?}" "AddTier must reject internal endpoint {endpoint} during URL validation, got: {response_body}"
); );
if checks_connection_attempt {
match tokio::time::timeout(Duration::from_secs(1), sentinel.accept()).await {
Err(_) => {}
Ok(Ok((_, peer))) => panic!("SSRF guard connected to loopback endpoint {endpoint} from {peer}"),
Ok(Err(error)) => return Err(error.into()),
}
}
} }
env.stop_server(); env.stop_server();
+5 -2
View File
@@ -17,8 +17,9 @@ use crate::admin::runtime_sources::object_store_from_extensions;
use crate::admin::storage_api::runtime_sources::TierConfigMgr; use crate::admin::storage_api::runtime_sources::TierConfigMgr;
use crate::admin::storage_api::tier::{ use crate::admin::storage_api::tier::{
AdminError, DailyAllTierStats, ERR_TIER_ALREADY_EXISTS, ERR_TIER_BACKEND_IN_USE, ERR_TIER_BACKEND_NOT_EMPTY, AdminError, DailyAllTierStats, ERR_TIER_ALREADY_EXISTS, ERR_TIER_BACKEND_IN_USE, ERR_TIER_BACKEND_NOT_EMPTY,
ERR_TIER_CONNECT_ERR, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_MISSING_CREDENTIALS, ERR_TIER_NAME_NOT_UPPERCASE, ERR_TIER_CONNECT_ERR, ERR_TIER_INVALID_CONFIG, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_MISSING_CREDENTIALS,
ERR_TIER_NOT_FOUND, ERR_TIER_RESERVED_NAME, TierConfig, TierConfigUpdateError, TierCreds, TierType, ERR_TIER_NAME_NOT_UPPERCASE, ERR_TIER_NOT_FOUND, ERR_TIER_RESERVED_NAME, TierConfig, TierConfigUpdateError, TierCreds,
TierType,
}; };
use crate::{ use crate::{
admin::runtime_sources::{current_daily_tier_stats, current_notification_system, current_tier_config_handle}, admin::runtime_sources::{current_daily_tier_stats, current_notification_system, current_tier_config_handle},
@@ -391,6 +392,8 @@ impl Operation for AddTier {
S3ErrorCode::Custom("TierConnectError".into()), S3ErrorCode::Custom("TierConnectError".into()),
"tier connectivity check failed", "tier connectivity check failed",
)) ))
} 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 { } else if err.code == ERR_TIER_INVALID_CREDENTIALS.code {
Err(S3Error::with_message(S3ErrorCode::Custom(err.code.clone().into()), err.message)) Err(S3Error::with_message(S3ErrorCode::Custom(err.code.clone().into()), err.message))
} else { } else {
+4 -2
View File
@@ -817,6 +817,7 @@ pub(crate) static ERR_TIER_MISSING_CREDENTIALS: AdminErrorRef =
pub(crate) static ERR_TIER_ALREADY_EXISTS: AdminErrorRef = pub(crate) static ERR_TIER_ALREADY_EXISTS: AdminErrorRef =
AdminErrorRef(|| &ecstore_tier::tier_handlers::ERR_TIER_ALREADY_EXISTS); AdminErrorRef(|| &ecstore_tier::tier_handlers::ERR_TIER_ALREADY_EXISTS);
pub(crate) static ERR_TIER_CONNECT_ERR: AdminErrorRef = AdminErrorRef(|| &ecstore_tier::tier_handlers::ERR_TIER_CONNECT_ERR); pub(crate) static ERR_TIER_CONNECT_ERR: AdminErrorRef = AdminErrorRef(|| &ecstore_tier::tier_handlers::ERR_TIER_CONNECT_ERR);
pub(crate) static ERR_TIER_INVALID_CONFIG: AdminErrorRef = AdminErrorRef(|| &ecstore_tier::tier::ERR_TIER_INVALID_CONFIG);
pub(crate) static ERR_TIER_INVALID_CREDENTIALS: AdminErrorRef = pub(crate) static ERR_TIER_INVALID_CREDENTIALS: AdminErrorRef =
AdminErrorRef(|| &ecstore_tier::tier_handlers::ERR_TIER_INVALID_CREDENTIALS); AdminErrorRef(|| &ecstore_tier::tier_handlers::ERR_TIER_INVALID_CREDENTIALS);
pub(crate) static ERR_TIER_NAME_NOT_UPPERCASE: AdminErrorRef = pub(crate) static ERR_TIER_NAME_NOT_UPPERCASE: AdminErrorRef =
@@ -963,7 +964,8 @@ pub(crate) mod s3 {
pub(crate) mod tier { pub(crate) mod tier {
pub(crate) use super::{ pub(crate) use super::{
AdminError, DailyAllTierStats, ERR_TIER_ALREADY_EXISTS, ERR_TIER_BACKEND_IN_USE, ERR_TIER_BACKEND_NOT_EMPTY, AdminError, DailyAllTierStats, ERR_TIER_ALREADY_EXISTS, ERR_TIER_BACKEND_IN_USE, ERR_TIER_BACKEND_NOT_EMPTY,
ERR_TIER_CONNECT_ERR, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_MISSING_CREDENTIALS, ERR_TIER_NAME_NOT_UPPERCASE, ERR_TIER_CONNECT_ERR, ERR_TIER_INVALID_CONFIG, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_MISSING_CREDENTIALS,
ERR_TIER_NOT_FOUND, ERR_TIER_RESERVED_NAME, TierConfig, TierConfigUpdateError, TierCreds, TierType, ERR_TIER_NAME_NOT_UPPERCASE, ERR_TIER_NOT_FOUND, ERR_TIER_RESERVED_NAME, TierConfig, TierConfigUpdateError, TierCreds,
TierType,
}; };
} }