feat(site-replication): support custom TLS peers (#4802)

* feat(madmin): add site replication TLS settings

* feat(site-replication): support custom TLS peers

* test(site-replication): remove redundant clones

* test(site-replication): avoid needless resolver collection
This commit is contained in:
cxymds
2026-07-14 15:33:00 +08:00
committed by GitHub
parent e9a0200a72
commit 25f81f812c
7 changed files with 2749 additions and 159 deletions
+2 -2
View File
@@ -136,7 +136,7 @@ test-group = 'e2e-reliability'
# the nightly profile derives its set as "the replication module MINUS this
# allowlist", so any new replication test lands in nightly by default (never
# silently unrun) until it is explicitly blessed as fast here. Keep the two
# regexes byte-identical. Count invariant: 20 here + 18 nightly = 38 total
# regexes byte-identical. Count invariant: 20 here + 20 nightly = 40 total
# (authority: `cargo nextest list`; docs/testing/e2e-suite-inventory.md).
# HISTORY (2026-07-11): the 20 fast tests were briefly pulled out of this lane
# (#4724) because they set a loopback (127.0.0.1) replication target that the
@@ -162,7 +162,7 @@ fail-fast = false
#
# * 8 bucket-replication data-plane tests — they PUT/delete objects and poll
# until source and target converge; two replicate over HTTPS.
# * 9 `_real_dual_node` site-replication tests — each spawns TWO full rustfs
# * 11 `_real_dual_node` site-replication tests — each spawns TWO full rustfs
# servers and drives the cross-process site-replication control plane.
# * 1 `_real_single_node` service-account round-trip test.
#
@@ -536,6 +536,7 @@ async fn start_https_rustfs_server(env: &mut RustFSTestEnvironment, tls_dir: &Pa
.env("RUST_LOG", "rustfs=info,rustfs_notify=debug")
.env("RUSTFS_TLS_PATH", tls_dir)
.env("RUSTFS_CONSOLE_ENABLE", "false")
.env("RUSTFS_REPLICATION_ALLOW_LOOPBACK_TARGET", "true")
.args([
"--address",
&env.address,
@@ -649,7 +650,7 @@ async fn wait_for_replicated_object_over_https(
}
status if tokio::time::Instant::now() < deadline => {
let body = response.text().await.unwrap_or_default();
if body.contains("NoSuchKey") || body.contains("NotFound") {
if body.contains("NoSuchKey") || body.contains("NoSuchBucket") || body.contains("NotFound") {
sleep(Duration::from_secs(1)).await;
continue;
}
@@ -2932,6 +2933,178 @@ async fn test_replication_recovers_after_runtime_target_cache_is_cleared() -> Re
Ok(())
}
#[tokio::test]
#[serial]
async fn test_site_replication_allows_self_signed_https_with_skip_tls_verify_real_dual_node() -> TestResult {
init_logging();
let mut source_env = new_replication_source_env().await?;
source_env
.start_rustfs_server_with_env(vec![], LOOPBACK_REPLICATION_TARGET_ENV)
.await?;
let mut target_env = new_replication_https_target_env().await?;
let tls_dir = std::path::PathBuf::from(&target_env.temp_dir).join("tls");
let target_host = target_env
.url
.trim_start_matches("https://")
.split(':')
.next()
.ok_or_else(|| std::io::Error::other("target HTTPS URL missing host"))?
.to_string();
generate_self_signed_tls_material(&tls_dir, &target_host).await?;
start_https_rustfs_server(&mut target_env, &tls_dir).await?;
let https_client = insecure_https_client()?;
wait_for_https_server_ready(&https_client, &target_env).await?;
let source_site = PeerSite {
name: "source-site".to_string(),
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
};
let target_site = PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
};
let add_error = site_replication_add(&source_env, &[source_site.clone(), target_site.clone()])
.await
.expect_err("site replication add must reject an untrusted self-signed HTTPS peer");
let add_error = add_error.to_string();
assert!(
add_error.contains("400 Bad Request"),
"unexpected untrusted self-signed peer error: {add_error}"
);
assert!(
add_error.to_ascii_lowercase().contains("tls") || add_error.to_ascii_lowercase().contains("certificate"),
"self-signed peer rejection did not report a TLS certificate error: {add_error}"
);
let disabled = wait_for_site_replication_disabled(&source_env).await?;
assert!(!disabled.enabled && disabled.sites.is_empty());
let add_status = site_replication_add(
&source_env,
&[
source_site,
PeerSite {
skip_tls_verify: true,
..target_site
},
],
)
.await?;
assert!(add_status.success, "unexpected site add result: {add_status:?}");
let _source_info = wait_for_site_replication_enabled(&source_env, 2).await?;
let source_client = source_env.create_s3_client();
let bucket = "site-repl-self-signed-tls";
let key = "self-signed.txt";
let body = "site replication over self-signed https";
source_client.create_bucket().bucket(bucket).send().await?;
source_client
.put_object()
.bucket(bucket)
.key(key)
.body(ByteStream::from(body.as_bytes().to_vec()))
.send()
.await?;
wait_for_replicated_object_over_https(&https_client, &target_env, bucket, key, body).await?;
Ok(())
}
#[tokio::test]
#[serial]
async fn test_site_replication_allows_private_ca_https_with_ca_cert_pem_real_dual_node() -> TestResult {
init_logging();
let mut source_env = new_replication_source_env().await?;
source_env
.start_rustfs_server_with_env(vec![], LOOPBACK_REPLICATION_TARGET_ENV)
.await?;
let mut target_env = new_replication_https_target_env().await?;
let tls_dir = std::path::PathBuf::from(&target_env.temp_dir).join("tls");
let target_host = target_env
.url
.trim_start_matches("https://")
.split(':')
.next()
.ok_or_else(|| std::io::Error::other("target HTTPS URL missing host"))?
.to_string();
let ca_cert_pem = generate_private_ca_tls_material(&tls_dir, &target_host).await?;
start_https_rustfs_server(&mut target_env, &tls_dir).await?;
let https_client = trusted_https_client(&ca_cert_pem)?;
wait_for_https_server_ready(&https_client, &target_env).await?;
let source_site = PeerSite {
name: "source-site".to_string(),
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
};
let target_site = PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
};
let add_error = site_replication_add(&source_env, &[source_site.clone(), target_site.clone()])
.await
.expect_err("site replication add must reject a private CA HTTPS peer without caCertPem");
let add_error = add_error.to_string();
assert!(
add_error.contains("400 Bad Request"),
"unexpected untrusted private CA peer error: {add_error}"
);
assert!(
add_error.to_ascii_lowercase().contains("tls") || add_error.to_ascii_lowercase().contains("certificate"),
"private CA peer rejection did not report a TLS certificate error: {add_error}"
);
let disabled = wait_for_site_replication_disabled(&source_env).await?;
assert!(!disabled.enabled && disabled.sites.is_empty());
let add_status = site_replication_add(
&source_env,
&[
source_site,
PeerSite {
ca_cert_pem,
..target_site
},
],
)
.await?;
assert!(add_status.success, "unexpected site add result: {add_status:?}");
let _source_info = wait_for_site_replication_enabled(&source_env, 2).await?;
let source_client = source_env.create_s3_client();
let bucket = "site-repl-private-ca-tls";
let key = "private-ca.txt";
let body = "site replication over private ca https";
source_client.create_bucket().bucket(bucket).send().await?;
source_client
.put_object()
.bucket(bucket)
.key(key)
.body(ByteStream::from(body.as_bytes().to_vec()))
.send()
.await?;
wait_for_replicated_object_over_https(&https_client, &target_env, bucket, key, body).await?;
Ok(())
}
#[tokio::test]
#[serial]
async fn test_site_replication_resync_start_cancel_restart_real_dual_node() -> Result<(), Box<dyn Error + Send + Sync>> {
@@ -2963,12 +3136,14 @@ async fn test_site_replication_resync_start_cancel_restart_real_dual_node() -> R
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
},
],
)
@@ -3092,18 +3267,21 @@ async fn test_site_replication_edit_and_status_peer_state_real_three_node() -> R
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "relay-site".to_string(),
endpoint: relay_env.url.clone(),
access_key: relay_env.access_key.clone(),
secret_key: relay_env.secret_key.clone(),
..Default::default()
},
],
)
@@ -3326,12 +3504,14 @@ async fn test_site_replication_remove_all_real_dual_node() -> Result<(), Box<dyn
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
},
],
)
@@ -3434,12 +3614,14 @@ async fn test_site_replication_state_edit_fresh_and_stale_real_dual_node() -> Re
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
},
],
)
@@ -3545,12 +3727,14 @@ async fn test_site_replication_replicates_object_with_bucket_versioning_real_dua
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
},
],
)
@@ -3627,12 +3811,14 @@ async fn test_site_replication_replicates_policy_backed_user_access_real_dual_no
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
},
],
)
@@ -3713,12 +3899,14 @@ async fn test_site_replication_replicates_group_policy_backed_access_real_dual_n
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
},
],
)
@@ -3841,12 +4029,14 @@ async fn test_site_replication_replicates_multiple_service_accounts_real_dual_no
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
},
],
)
@@ -3943,12 +4133,14 @@ async fn test_site_replication_replicates_service_accounts_created_from_sts_sess
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_env.secret_key.clone(),
..Default::default()
},
],
)
+255 -52
View File
@@ -1023,20 +1023,49 @@ fn build_insecure_aws_s3_http_client() -> SharedHttpClient {
http_client_fn(move |_settings, _components| connector.clone())
}
fn build_aws_s3_http_client_from_target_ca_pem(ca_cert_pem: &str) -> Result<SharedHttpClient, BucketTargetError> {
let certs = rustls_pki_types::CertificateDer::pem_slice_iter(ca_cert_pem.as_bytes())
fn validate_ca_pem_bundle(ca_cert_pem: &[u8]) -> Result<(), String> {
let certs = rustls_pki_types::CertificateDer::pem_slice_iter(ca_cert_pem)
.collect::<Result<Vec<_>, _>>()
.map_err(|err| BucketTargetError::Io(std::io::Error::other(format!("invalid target CA PEM: {err}"))))?;
.map_err(|err| format!("invalid PEM encoding: {err}"))?;
if certs.is_empty() {
return Err(BucketTargetError::Io(std::io::Error::other(
"invalid target CA PEM: no certificates found",
)));
return Err("no certificates found".to_string());
}
let mut trust_store = smithy_tls::TrustStore::empty();
trust_store.add_pem_certificate(ca_cert_pem.as_bytes());
// Smithy's rustls adapter defers parsing custom certificates and assumes
// they are valid when the HTTPS connector is built. Validate every DER
// certificate first so malformed configuration is reported rather than
// reaching an `expect` in the dependency.
let mut validation_store = rustls::RootCertStore::empty();
for cert in certs {
validation_store
.add(cert)
.map_err(|err| format!("invalid X.509 certificate: {err}"))?;
}
Ok(())
}
fn validate_target_ca_pem(ca_cert_pem: &str) -> Result<(), BucketTargetError> {
validate_ca_pem_bundle(ca_cert_pem.as_bytes())
.map_err(|err| BucketTargetError::Io(std::io::Error::other(format!("invalid target CA PEM: {err}"))))
}
fn compose_replication_trust_store(certificate_bundles: impl IntoIterator<Item = Vec<u8>>) -> (smithy_tls::TrustStore, usize) {
// `TrustStore::default()` keeps the platform-native roots enabled. Target
// and RUSTFS_TLS_PATH certificates extend that baseline instead of
// replacing it with a target-specific trust island.
let mut trust_store = smithy_tls::TrustStore::default();
let mut custom_bundle_count = 0;
for pem in certificate_bundles {
trust_store.add_pem_certificate(pem);
custom_bundle_count += 1;
}
(trust_store, custom_bundle_count)
}
fn build_aws_s3_http_client_with_trust_store(trust_store: smithy_tls::TrustStore) -> Result<SharedHttpClient, BucketTargetError> {
let tls_context = smithy_tls::TlsContext::builder()
.with_trust_store(trust_store)
.build()
@@ -1048,6 +1077,60 @@ fn build_aws_s3_http_client_from_target_ca_pem(ca_cert_pem: &str) -> Result<Shar
.build_https())
}
async fn load_tls_path_ca_bundles(tls_dir: &Path, trust_leaf_cert_as_ca: bool) -> Vec<Vec<u8>> {
let mut certificate_bundles = Vec::new();
let ca_path = tls_dir.join(RUSTFS_CA_CERT);
match tokio::fs::read(&ca_path).await {
Ok(pem) => match validate_ca_pem_bundle(&pem) {
Ok(()) => certificate_bundles.push(pem),
Err(err) => warn!("ignoring invalid custom CA bundle {:?} for replication client: {}", ca_path, err),
},
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => warn!("failed to read custom CA bundle {:?} for replication client: {}", ca_path, e),
}
if trust_leaf_cert_as_ca {
let leaf_cert_path = tls_dir.join(RUSTFS_TLS_CERT);
match tokio::fs::read(&leaf_cert_path).await {
Ok(pem) => match validate_ca_pem_bundle(&pem) {
Ok(()) => certificate_bundles.push(pem),
Err(err) => warn!(
"ignoring invalid leaf certificate {:?} for replication client trust store: {}",
leaf_cert_path, err
),
},
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => warn!("failed to read leaf cert {:?} for replication client trust store: {}", leaf_cert_path, e),
}
}
certificate_bundles
}
async fn load_configured_tls_ca_bundles() -> Vec<Vec<u8>> {
let tls_path = rustfs_utils::get_env_str(rustfs_config::ENV_RUSTFS_TLS_PATH, rustfs_config::DEFAULT_RUSTFS_TLS_PATH);
if tls_path.is_empty() {
return Vec::new();
}
load_tls_path_ca_bundles(
Path::new(&tls_path),
rustfs_utils::get_env_bool(ENV_TRUST_LEAF_CERT_AS_CA, DEFAULT_TRUST_LEAF_CERT_AS_CA),
)
.await
}
async fn build_aws_s3_http_client_from_target_ca_pem(ca_cert_pem: &str) -> Result<SharedHttpClient, BucketTargetError> {
validate_target_ca_pem(ca_cert_pem)?;
let mut certificate_bundles = load_configured_tls_ca_bundles().await;
certificate_bundles.push(ca_cert_pem.as_bytes().to_vec());
let (trust_store, _) = compose_replication_trust_store(certificate_bundles);
build_aws_s3_http_client_with_trust_store(trust_store)
}
async fn build_aws_s3_http_client_for_target(target: &BucketTarget) -> Result<Option<SharedHttpClient>, BucketTargetError> {
if !target.secure {
return Ok(None);
@@ -1058,62 +1141,28 @@ async fn build_aws_s3_http_client_for_target(target: &BucketTarget) -> Result<Op
}
if has_custom_ca_pem(target) {
return build_aws_s3_http_client_from_target_ca_pem(&target.ca_cert_pem).map(Some);
return build_aws_s3_http_client_from_target_ca_pem(&target.ca_cert_pem)
.await
.map(Some);
}
Ok(build_aws_s3_http_client_from_tls_path().await)
}
async fn build_aws_s3_http_client_from_tls_path() -> Option<aws_sdk_s3::config::SharedHttpClient> {
let tls_path = rustfs_utils::get_env_str(rustfs_config::ENV_RUSTFS_TLS_PATH, rustfs_config::DEFAULT_RUSTFS_TLS_PATH);
if tls_path.is_empty() {
let certificate_bundles = load_configured_tls_ca_bundles().await;
if certificate_bundles.is_empty() {
return None;
}
let tls_dir = Path::new(&tls_path);
let mut trust_store = smithy_tls::TrustStore::empty();
let mut has_custom_certs = false;
let ca_path = tls_dir.join(RUSTFS_CA_CERT);
match tokio::fs::read(&ca_path).await {
Ok(pem) => {
trust_store.add_pem_certificate(pem);
has_custom_certs = true;
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => warn!("failed to read custom CA bundle {:?} for replication client: {}", ca_path, e),
}
if rustfs_utils::get_env_bool(ENV_TRUST_LEAF_CERT_AS_CA, DEFAULT_TRUST_LEAF_CERT_AS_CA) {
let leaf_cert_path = tls_dir.join(RUSTFS_TLS_CERT);
match tokio::fs::read(&leaf_cert_path).await {
Ok(pem) => {
trust_store.add_pem_certificate(pem);
has_custom_certs = true;
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => warn!("failed to read leaf cert {:?} for replication client trust store: {}", leaf_cert_path, e),
}
}
if !has_custom_certs {
return None;
}
let tls_context = match smithy_tls::TlsContext::builder().with_trust_store(trust_store).build() {
Ok(ctx) => ctx,
let (trust_store, _) = compose_replication_trust_store(certificate_bundles);
match build_aws_s3_http_client_with_trust_store(trust_store) {
Ok(client) => Some(client),
Err(e) => {
warn!("failed to build AWS SDK TLS context for replication client: {}", e);
return None;
None
}
};
Some(
SmithyHttpClientBuilder::new()
.tls_provider(smithy_tls::Provider::rustls(smithy_tls::rustls_provider::CryptoMode::AwsLc))
.tls_context(tls_context)
.build_https(),
)
}
}
fn should_force_path_style(target: &BucketTarget) -> bool {
@@ -1920,6 +1969,63 @@ mod tests {
use super::*;
use rcgen::generate_simple_self_signed;
fn spawn_single_request_https_server(cert: &rcgen::CertifiedKey<rcgen::KeyPair>) -> (u16, std::thread::JoinHandle<()>) {
use std::io::{Read, Write};
ensure_rustls_crypto_provider();
let listener = std::net::TcpListener::bind(("127.0.0.1", 0)).expect("test TLS listener should bind");
let port = listener
.local_addr()
.expect("test TLS listener should have an address")
.port();
let server_config = rustls::ServerConfig::builder()
.with_no_client_auth()
.with_single_cert(
vec![cert.cert.der().clone()],
rustls_pki_types::PrivateKeyDer::try_from(cert.signing_key.serialize_der())
.expect("test TLS private key should convert"),
)
.expect("test TLS server config should build");
let handle = std::thread::spawn(move || {
let (stream, _) = listener.accept().expect("test TLS client should connect");
stream
.set_read_timeout(Some(Duration::from_secs(10)))
.expect("test TLS read timeout should configure");
stream
.set_write_timeout(Some(Duration::from_secs(10)))
.expect("test TLS write timeout should configure");
let connection = rustls::ServerConnection::new(Arc::new(server_config)).expect("test TLS connection should build");
let mut stream = rustls::StreamOwned::new(connection, stream);
let mut request = [0_u8; 8192];
let _ = stream.read(&mut request).expect("test TLS request should be readable");
stream
.write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 0\r\nConnection: close\r\n\r\n")
.expect("test TLS response should be written");
stream.flush().expect("test TLS response should flush");
});
(port, handle)
}
fn s3_client_with_http_client(port: u16, http_client: SharedHttpClient) -> S3Client {
let credentials = SdkCredentials::builder()
.access_key_id("test-access")
.secret_access_key("test-secret")
.provider_name("bucket_target_tls_test")
.build();
let config = S3Config::builder()
.endpoint_url(format!("https://localhost:{port}"))
.credentials_provider(SharedCredentialsProvider::new(credentials))
.region(SdkRegion::new("us-east-1"))
.force_path_style(true)
.behavior_version(aws_sdk_s3::config::BehaviorVersion::latest())
.http_client(http_client)
.build();
S3Client::from_conf(config)
}
#[test]
fn replication_target_versioning_enabled_requires_enabled_status() {
let enabled = BucketVersioningStatus::Enabled;
@@ -2243,6 +2349,78 @@ mod tests {
assert_eq!(client.endpoint, "https://192.168.1.10:9000");
}
#[tokio::test]
async fn replication_trust_store_composes_system_global_and_target_roots_for_real_tls() {
let tls_dir = tempfile::tempdir().expect("temporary TLS directory should be created");
let global_ca =
generate_simple_self_signed(vec!["localhost".to_string()]).expect("global CA certificate should generate");
let target_ca =
generate_simple_self_signed(vec!["localhost".to_string()]).expect("target CA certificate should generate");
tokio::fs::write(tls_dir.path().join(RUSTFS_CA_CERT), global_ca.cert.pem())
.await
.expect("global CA bundle should be written");
let mut certificate_bundles = load_tls_path_ca_bundles(tls_dir.path(), false).await;
certificate_bundles.push(target_ca.cert.pem().into_bytes());
let (trust_store, _) = compose_replication_trust_store(certificate_bundles);
assert!(
format!("{trust_store:?}").contains("enable_native_roots: true"),
"per-target trust must retain the SDK's platform-native roots"
);
let http_client = build_aws_s3_http_client_with_trust_store(trust_store).expect("composed TLS client should build");
let (global_port, global_server) = spawn_single_request_https_server(&global_ca);
s3_client_with_http_client(global_port, http_client.clone())
.head_bucket()
.bucket("test-bucket")
.send()
.await
.expect("global RUSTFS_TLS_PATH CA should authenticate its TLS server");
global_server.join().expect("global CA TLS server should finish");
let (target_port, target_server) = spawn_single_request_https_server(&target_ca);
s3_client_with_http_client(target_port, http_client)
.head_bucket()
.bucket("test-bucket")
.send()
.await
.expect("per-target CA should authenticate its TLS server alongside global roots");
target_server.join().expect("target CA TLS server should finish");
}
#[tokio::test]
async fn tls_path_leaf_trust_remains_opt_in() {
let tls_dir = tempfile::tempdir().expect("temporary TLS directory should be created");
let global_ca =
generate_simple_self_signed(vec!["global-ca.example".to_string()]).expect("global CA certificate should generate");
let trusted_leaf =
generate_simple_self_signed(vec!["leaf.example".to_string()]).expect("trusted leaf certificate should generate");
tokio::fs::write(tls_dir.path().join(RUSTFS_CA_CERT), global_ca.cert.pem())
.await
.expect("global CA bundle should be written");
tokio::fs::write(tls_dir.path().join(RUSTFS_TLS_CERT), trusted_leaf.cert.pem())
.await
.expect("trusted leaf certificate should be written");
assert_eq!(load_tls_path_ca_bundles(tls_dir.path(), true).await.len(), 2);
assert_eq!(load_tls_path_ca_bundles(tls_dir.path(), false).await.len(), 1);
}
#[tokio::test]
async fn skip_tls_verify_takes_priority_over_invalid_custom_ca_pem() {
let client = build_aws_s3_http_client_for_target(&BucketTarget {
secure: true,
skip_tls_verify: true,
ca_cert_pem: "not a pem".to_string(),
..Default::default()
})
.await
.expect("skip verification should bypass custom CA parsing");
assert!(client.is_some(), "secure targets with skip verification need a custom HTTP client");
}
#[tokio::test]
async fn get_remote_target_client_internal_rejects_invalid_custom_ca_pem() {
let sys = BucketTargetSys::default();
@@ -2267,6 +2445,31 @@ mod tests {
assert!(err.to_string().contains("invalid target CA PEM"));
}
#[test]
fn target_ca_rejects_pem_wrapped_invalid_der_before_smithy_builds() {
let err = validate_target_ca_pem("-----BEGIN CERTIFICATE-----\nAQID\n-----END CERTIFICATE-----\n")
.expect_err("PEM-wrapped invalid DER must be rejected");
assert!(err.to_string().contains("invalid target CA PEM"));
assert!(err.to_string().contains("invalid X.509 certificate"));
}
#[tokio::test]
async fn invalid_global_ca_is_ignored_without_reaching_smithy() {
let tls_dir = tempfile::tempdir().expect("temporary TLS directory should be created");
tokio::fs::write(
tls_dir.path().join(RUSTFS_CA_CERT),
b"-----BEGIN CERTIFICATE-----\nAQID\n-----END CERTIFICATE-----\n",
)
.await
.expect("invalid global CA fixture should be written");
assert!(
load_tls_path_ca_bundles(tls_dir.path(), false).await.is_empty(),
"invalid global CA must fall back to default roots instead of reaching Smithy's panic path"
);
}
// backlog#806-16 regression tests for the rolling one-minute latency window.
#[test]
+138 -2
View File
@@ -16,11 +16,12 @@ use crate::{GroupAddRemove, GroupDesc, SRSvcAccCreate, UserInfo};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::{BTreeMap, HashMap};
use std::fmt;
use time::OffsetDateTime;
pub const SITE_REPL_API_VERSION: &str = "1";
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
#[derive(Clone, Serialize, Deserialize, Default)]
pub struct PeerSite {
#[serde(default)]
pub name: String,
@@ -30,6 +31,21 @@ pub struct PeerSite {
pub access_key: String,
#[serde(rename = "secretKey", default)]
pub secret_key: String,
#[serde(rename = "skipTlsVerify", default)]
pub skip_tls_verify: bool,
#[serde(rename = "caCertPem", default)]
pub ca_cert_pem: String,
}
impl fmt::Debug for PeerSite {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("PeerSite")
.field("name", &self.name)
.field("endpoint", &self.endpoint)
.field("skip_tls_verify", &self.skip_tls_verify)
.field("has_custom_ca", &!self.ca_cert_pem.is_empty())
.finish()
}
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
@@ -106,7 +122,7 @@ pub enum SyncStatus {
Unknown,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
#[derive(Clone, Serialize, Deserialize, Default)]
pub struct PeerInfo {
#[serde(default)]
pub endpoint: String,
@@ -122,10 +138,26 @@ pub struct PeerInfo {
pub replicate_ilm_expiry: bool,
#[serde(rename = "objectNamingMode", default, skip_serializing_if = "String::is_empty")]
pub object_naming_mode: String,
#[serde(rename = "skipTlsVerify", default)]
pub skip_tls_verify: bool,
#[serde(rename = "caCertPem", default)]
pub ca_cert_pem: String,
#[serde(rename = "apiVersion", skip_serializing_if = "Option::is_none")]
pub api_version: Option<String>,
}
impl fmt::Debug for PeerInfo {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("PeerInfo")
.field("endpoint", &self.endpoint)
.field("name", &self.name)
.field("deployment_id", &self.deployment_id)
.field("skip_tls_verify", &self.skip_tls_verify)
.field("has_custom_ca", &!self.ca_cert_pem.is_empty())
.finish()
}
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct SRPolicyMapping {
#[serde(rename = "userOrGroup", default)]
@@ -1167,3 +1199,107 @@ pub struct SiteNetPerfResult {
#[serde(rename = "nodeResults", default, skip_serializing_if = "Vec::is_empty")]
pub node_results: Vec<SiteNetPerfNodeResult>,
}
#[cfg(test)]
mod tests {
use super::{PeerInfo, PeerSite};
use serde_json::{Value, json};
const TEST_CA_CERT: &str = "-----BEGIN CERTIFICATE-----\ntest-ca\n-----END CERTIFICATE-----";
#[test]
fn peer_tls_fields_default_when_missing_from_legacy_json() {
let site: PeerSite = serde_json::from_value(json!({
"name": "site-a",
"endpoints": "https://site-a.example.com"
}))
.expect("legacy PeerSite JSON should deserialize");
let peer: PeerInfo = serde_json::from_value(json!({
"endpoint": "https://site-a.example.com",
"name": "site-a",
"deploymentID": "deployment-a"
}))
.expect("legacy PeerInfo JSON should deserialize");
assert!(!site.skip_tls_verify);
assert_eq!(site.ca_cert_pem, "");
assert!(!peer.skip_tls_verify);
assert_eq!(peer.ca_cert_pem, "");
}
#[test]
fn peer_tls_fields_round_trip_with_exact_json_names() {
let site_json = json!({
"name": "site-a",
"endpoints": "https://site-a.example.com",
"accessKey": "access-key",
"secretKey": "secret-key",
"skipTlsVerify": true,
"caCertPem": TEST_CA_CERT
});
let peer_json = json!({
"endpoint": "https://site-a.example.com",
"name": "site-a",
"deploymentID": "deployment-a",
"sync": "unknown",
"defaultbandwidth": {
"bandwidthLimitPerBucket": 0,
"set": false
},
"replicate-ilm-expiry": false,
"objectNamingMode": "path",
"skipTlsVerify": true,
"caCertPem": TEST_CA_CERT
});
let site: PeerSite = serde_json::from_value(site_json.clone()).expect("PeerSite JSON should deserialize");
let peer: PeerInfo = serde_json::from_value(peer_json.clone()).expect("PeerInfo JSON should deserialize");
assert_eq!(serde_json::to_value(site).expect("PeerSite should serialize"), site_json);
assert_eq!(serde_json::to_value(peer).expect("PeerInfo should serialize"), peer_json);
}
#[test]
fn peer_tls_false_and_empty_ca_are_still_serialized() {
let site = serde_json::to_value(PeerSite::default()).expect("PeerSite should serialize");
let peer = serde_json::to_value(PeerInfo::default()).expect("PeerInfo should serialize");
for value in [site, peer] {
let object = value.as_object().expect("peer JSON should be an object");
assert_eq!(object.get("skipTlsVerify"), Some(&Value::Bool(false)));
assert_eq!(object.get("caCertPem"), Some(&Value::String(String::new())));
}
}
#[test]
fn peer_debug_output_redacts_secrets_and_ca_contents() {
let site = PeerSite {
name: "site-a".to_owned(),
endpoint: "https://site-a.example.com".to_owned(),
access_key: "sensitive-access-key".to_owned(),
secret_key: "sensitive-secret-key".to_owned(),
skip_tls_verify: true,
ca_cert_pem: TEST_CA_CERT.to_owned(),
};
let peer = PeerInfo {
endpoint: "https://site-a.example.com".to_owned(),
name: "site-a".to_owned(),
deployment_id: "deployment-a".to_owned(),
skip_tls_verify: false,
ca_cert_pem: TEST_CA_CERT.to_owned(),
..PeerInfo::default()
};
let site_debug = format!("{site:?}");
let peer_debug = format!("{peer:?}");
assert!(!site_debug.contains("BEGIN CERTIFICATE"));
assert!(!site_debug.contains("sensitive-access-key"));
assert!(!site_debug.contains("sensitive-secret-key"));
assert!(site_debug.contains("skip_tls_verify: true"));
assert!(site_debug.contains("has_custom_ca: true"));
assert!(!peer_debug.contains("BEGIN CERTIFICATE"));
assert!(peer_debug.contains("skip_tls_verify: false"));
assert!(peer_debug.contains("has_custom_ca: true"));
}
}
+10 -5
View File
@@ -3,18 +3,19 @@
> Authoritative per-module test counts for the `e2e_test` crate (backlog#1149
> ci-4), generated from `cargo nextest list -p e2e_test`. Regenerate with:
> ```bash
> cargo nextest list -p e2e_test | grep "^ " | sed "s/^ *//;s/::.*//" | sort | uniq -c
> cargo nextest list -p e2e_test --message-format oneline | awk '{split($2,a,"::"); print a[1]}' | sort | uniq -c
> ```
> Modules marked ✅ are in the PR smoke profile `e2e-smoke`
> (`.config/nextest.toml`); admission criteria: `crates/e2e_test/README.md`.
> 🌙 marks tests in the scheduled `e2e-repl-nightly` profile (backlog#1147
> repl-1): `replication_extension_test` splits 20 fast tests into the PR smoke
> lane and 16 slow / `_real_dual_node` / `_real_single_node` tests into the
> lane and 20 slow / `_real_dual_node` / `_real_single_node` tests into the
> nightly lane (`.github/workflows/e2e-replication-nightly.yml`).
> Note: counts exclude `#[ignore]`d tests (nextest lists them separately).
| module | tests | PR smoke |
|---|---|---|
| admin_auth_test | 3 | |
| admin_timeout_regression_test | 1 | |
| anonymous_access_test | 3 | ✅ |
| archive_download_integrity_test | 13 | |
@@ -29,10 +30,12 @@
| copy_object_version_restore_test | 1 | |
| copy_source_invalid_date_test | 1 | ✅ |
| create_bucket_region_test | 2 | ✅ |
| degraded_read_eof_regression_test | 3 | |
| delete_marker_migration_semantics_test | 2 | ✅ |
| delete_object_no_content_length_test | 1 | |
| delete_objects_versioning_test | 2 | ✅ |
| existing_object_tag_policy_test | 4 | |
| get_codec_streaming_compat_test | 1 | |
| head_object_consistency_test | 1 | ✅ |
| head_object_range_test | 1 | ✅ |
| heal_erasure_disk_rebuild_test | 3 | |
@@ -47,14 +50,16 @@
| mc_mirror_small_bucket_test | 1 | |
| multipart_auth_test | 109 | |
| namespace_lock_quorum_test | 2 | |
| negative_sigv4_test | 6 | |
| object_lambda_test | 16 | |
| object_lock | 33 | |
| overwrite_cleanup_regression_test | 1 | |
| presigned_negative_test | 7 | ✅ |
| protocols | 16 | |
| quota_test | 13 | |
| reliability_disk_fault_test | 3 | |
| reliant | 6 | |
| replication_extension_test | 38 | 20 ✅ +18 🌙 |
| reliant | 9 | 3 ✅ |
| replication_extension_test | 40 | 20 ✅ +20 🌙 |
| security_boundary_test | 4 | |
| server_startup_failfast_test | 1 | |
| snowball_auto_extract_test | 6 | |
@@ -63,4 +68,4 @@
| tls_gen | 3 | |
| version_id_regression_test | 10 | ✅ |
**Total listed: 396 tests across 47 modules · PR smoke subset: 83 tests / 18 modules** (17 full modules + 20 of `replication_extension_test`) **· nightly `e2e-repl-nightly`: 18 tests** · generated 2026-07-12.
**Total listed: 421 tests across 52 modules · PR smoke subset: 93 tests / 20 modules** (18 full modules + 3 `reliant::lifecycle` tests + 20 of `replication_extension_test`) **· nightly `e2e-repl-nightly`: 20 tests** · generated 2026-07-14.
File diff suppressed because it is too large Load Diff
@@ -160,6 +160,8 @@ mod tests {
default_bandwidth: BucketBandwidth::default(),
replicate_ilm_expiry: false,
object_naming_mode: String::new(),
skip_tls_verify: false,
ca_cert_pem: String::new(),
api_version: None,
}
}
@@ -204,6 +206,8 @@ mod tests {
PeerInfo {
api_version: Some("v1".to_string()),
replicate_ilm_expiry: true,
skip_tls_verify: true,
ca_cert_pem: "fallback-ca".to_string(),
..peer("remote-http", "http://node-a.example.com:9000")
},
),
@@ -226,6 +230,8 @@ mod tests {
assert_eq!(peer.deployment_id, "remote-http");
assert_eq!(peer.api_version.as_deref(), Some("v1"));
assert!(peer.replicate_ilm_expiry);
assert!(!peer.skip_tls_verify);
assert_eq!(peer.ca_cert_pem, "");
}
#[test]