diff --git a/Cargo.lock b/Cargo.lock index 4589e35b7..1d6c53b8d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3697,6 +3697,7 @@ dependencies = [ "flate2", "futures", "http 1.4.2", + "local-ip-address", "md5", "rand 0.10.1", "rcgen", @@ -9300,6 +9301,7 @@ dependencies = [ "aws-credential-types", "aws-sdk-s3", "aws-smithy-http-client", + "aws-smithy-runtime-api", "aws-smithy-types", "base64 0.22.1", "base64-simd", @@ -9340,6 +9342,7 @@ dependencies = [ "quick-xml 0.40.1", "rand 0.10.1", "ratelimit", + "rcgen", "reed-solomon-simd", "regex", "reqwest", diff --git a/Cargo.toml b/Cargo.toml index 075d8c32c..f3e8514ad 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -214,6 +214,7 @@ aws-config = { version = "1.8.18" } aws-credential-types = { version = "1.2.14" } aws-sdk-s3 = { version = "1.137.0", default-features = false, features = ["sigv4a", "default-https-client", "rt-tokio"] } aws-smithy-http-client = { version = "1.1.13", default-features = false, features = ["default-client", "rustls-aws-lc"] } +aws-smithy-runtime-api = { version = "1.12.3", features = ["http-1x"] } aws-smithy-types = { version = "1.5.0" } base64 = "0.22.1" base64-simd = "0.8.0" diff --git a/crates/e2e_test/Cargo.toml b/crates/e2e_test/Cargo.toml index 91316b920..c36861a10 100644 --- a/crates/e2e_test/Cargo.toml +++ b/crates/e2e_test/Cargo.toml @@ -72,6 +72,7 @@ zstd.workspace = true time.workspace = true suppaftp = { workspace = true, features = ["tokio", "rustls-aws-lc-rs"] } rcgen.workspace = true +local-ip-address.workspace = true anyhow.workspace = true rustls.workspace = true russh = { workspace = true } diff --git a/crates/e2e_test/src/replication_extension_test.rs b/crates/e2e_test/src/replication_extension_test.rs index 598884e0f..e9a4776e7 100644 --- a/crates/e2e_test/src/replication_extension_test.rs +++ b/crates/e2e_test/src/replication_extension_test.rs @@ -14,6 +14,7 @@ use crate::common::{ RustFSTestEnvironment, awscurl_available, awscurl_post_sts_form_urlencoded, init_logging, local_http_client, + rustfs_binary_path, }; use aws_sdk_s3::config::{Credentials, Region}; use aws_sdk_s3::error::ProvideErrorMetadata; @@ -21,6 +22,11 @@ use aws_sdk_s3::primitives::ByteStream; use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration}; use aws_sdk_s3::{Client, Config}; use http::header::{CONTENT_TYPE, HOST}; +use local_ip_address::local_ip; +use rcgen::{ + BasicConstraints, CertificateParams, CertifiedIssuer, DnType, ExtendedKeyUsagePurpose, IsCa, KeyPair, KeyUsagePurpose, + SanType, generate_simple_self_signed, +}; use reqwest::StatusCode; use rustfs_ecstore::api::bucket::bucket_target_sys::BucketTargetSys; use rustfs_madmin::{ @@ -33,8 +39,13 @@ use s3s::Body; use serial_test::serial; use std::collections::BTreeMap; use std::error::Error; -use time::Duration as TimeDuration; +use std::net::IpAddr; +use std::path::Path; +use std::process::Command; +use time::{Duration as TimeDuration, OffsetDateTime}; +use tokio::fs; use tokio::time::{Duration, sleep}; +use uuid::Uuid; type TestResult = Result<(), Box>; @@ -87,6 +98,39 @@ async fn signed_request( Ok(request_builder.send().await?) } +async fn signed_request_with_client( + client: &reqwest::Client, + method: http::Method, + url: &str, + access_key: &str, + secret_key: &str, + body: Option>, + content_type: Option<&str>, +) -> Result> { + let uri = url.parse::()?; + let authority = uri.authority().ok_or("request URL missing authority")?.to_string(); + let mut request = http::Request::builder().method(method.clone()).uri(uri); + request = request.header(HOST, authority); + request = request.header("x-amz-content-sha256", UNSIGNED_PAYLOAD); + if let Some(content_type) = content_type { + request = request.header(CONTENT_TYPE, content_type); + } + + let content_len = body.as_ref().map(|body| body.len() as i64).unwrap_or_default(); + let signed = sign_v4(request.body(Body::empty())?, content_len, access_key, secret_key, "", "us-east-1"); + + let reqwest_method = reqwest::Method::from_bytes(method.as_str().as_bytes())?; + let mut request_builder = client.request(reqwest_method, url); + for (name, value) in signed.headers() { + request_builder = request_builder.header(name, value); + } + if let Some(body) = body { + request_builder = request_builder.body(body); + } + + Ok(request_builder.send().await?) +} + async fn signed_request_with_session_token( method: http::Method, url: &str, @@ -146,22 +190,57 @@ fn parse_assume_role_credentials(xml: &str) -> Result<(String, String, String), Ok((access_key, secret_key, session_token)) } +struct ReplicationTargetOptions<'a> { + endpoint: &'a str, + access_key: &'a str, + secret_key: &'a str, + target_bucket: &'a str, + secure: bool, + skip_tls_verify: bool, + ca_cert_pem: Option<&'a str>, +} + async fn set_replication_target( source_env: &RustFSTestEnvironment, source_bucket: &str, target_env: &RustFSTestEnvironment, target_bucket: &str, ) -> Result> { - let body = serde_json::json!({ - "endpoint": target_env.address, - "credentials": { - "accessKey": target_env.access_key, - "secretKey": target_env.secret_key + set_replication_target_with_options( + source_env, + source_bucket, + ReplicationTargetOptions { + endpoint: &target_env.address, + access_key: &target_env.access_key, + secret_key: &target_env.secret_key, + target_bucket, + secure: false, + skip_tls_verify: false, + ca_cert_pem: None, }, - "targetbucket": target_bucket, - "secure": false, + ) + .await +} + +async fn set_replication_target_with_options( + source_env: &RustFSTestEnvironment, + source_bucket: &str, + options: ReplicationTargetOptions<'_>, +) -> Result> { + let mut body = serde_json::json!({ + "endpoint": options.endpoint, + "credentials": { + "accessKey": options.access_key, + "secretKey": options.secret_key + }, + "targetbucket": options.target_bucket, + "secure": options.secure, + "skipTlsVerify": options.skip_tls_verify, "type": "replication" }); + if let Some(ca_cert_pem) = options.ca_cert_pem { + body["caCertPem"] = serde_json::Value::String(ca_cert_pem.to_string()); + } let url = format!( "{}/rustfs/admin/v3/set-remote-target?bucket={}", source_env.url, @@ -334,6 +413,241 @@ async fn enable_bucket_versioning(env: &RustFSTestEnvironment, bucket: &str) -> Ok(()) } +fn insecure_https_client() -> Result> { + Ok(reqwest::Client::builder() + .no_proxy() + .danger_accept_invalid_certs(true) + .build()?) +} + +fn trusted_https_client(ca_cert_pem: &str) -> Result> { + let ca_cert = reqwest::Certificate::from_pem(ca_cert_pem.as_bytes())?; + Ok(reqwest::Client::builder().no_proxy().add_root_certificate(ca_cert).build()?) +} + +async fn new_private_tmp_test_env() -> Result> { + let temp_dir = format!("/private/tmp/rustfs_e2e_test_{}", Uuid::new_v4()); + fs::create_dir_all(&temp_dir) + .await + .map_err(|err| std::io::Error::other(format!("create temp dir {temp_dir} failed: {err}")))?; + let port = RustFSTestEnvironment::find_available_port() + .await + .map_err(|err| std::io::Error::other(format!("find available port failed: {err}")))?; + let address = format!("127.0.0.1:{port}"); + let url = format!("http://{address}"); + + Ok(RustFSTestEnvironment { + temp_dir, + address, + url, + access_key: "rustfsadmin".to_string(), + secret_key: "rustfsadmin".to_string(), + process: None, + }) +} + +async fn new_private_tmp_https_target_env() -> Result> { + let mut env = new_private_tmp_test_env().await?; + let public_ip = local_ip().map_err(|err| std::io::Error::other(format!("resolve local IP failed: {err}")))?; + let port = env + .address + .rsplit(':') + .next() + .ok_or_else(|| std::io::Error::other("target env address missing port"))? + .to_string(); + env.address = format!("0.0.0.0:{port}"); + env.url = format!("https://{public_ip}:{port}"); + Ok(env) +} + +async fn generate_self_signed_tls_material(tls_dir: &Path, additional_san: &str) -> Result<(), Box> { + fs::create_dir_all(tls_dir).await?; + let cert = generate_simple_self_signed(vec!["localhost".to_string(), "127.0.0.1".to_string(), additional_san.to_string()])?; + fs::write(tls_dir.join("rustfs_cert.pem"), cert.cert.pem()).await?; + fs::write(tls_dir.join("rustfs_key.pem"), cert.signing_key.serialize_pem()).await?; + Ok(()) +} + +fn test_certificate_params(common_name: &str) -> CertificateParams { + let mut params = CertificateParams::default(); + let issued_at = OffsetDateTime::now_utc() - TimeDuration::minutes(5); + params.not_before = issued_at; + params.not_after = issued_at + TimeDuration::days(1); + params.distinguished_name.push(DnType::CountryName, "US"); + params.distinguished_name.push(DnType::OrganizationName, "RustFS"); + params.distinguished_name.push(DnType::CommonName, common_name); + params +} + +async fn generate_private_ca_tls_material(tls_dir: &Path, additional_san: &str) -> Result> { + fs::create_dir_all(tls_dir).await?; + + let ca_key = KeyPair::generate()?; + let mut ca_params = test_certificate_params("RustFS Replication Test CA"); + ca_params.is_ca = IsCa::Ca(BasicConstraints::Unconstrained); + ca_params.key_usages = vec![KeyUsagePurpose::KeyCertSign, KeyUsagePurpose::CrlSign]; + let ca = CertifiedIssuer::self_signed(ca_params, ca_key)?; + + let server_key = KeyPair::generate()?; + let mut server_params = test_certificate_params("localhost"); + server_params.is_ca = IsCa::ExplicitNoCa; + server_params.key_usages = vec![KeyUsagePurpose::DigitalSignature, KeyUsagePurpose::KeyEncipherment]; + server_params.extended_key_usages = vec![ExtendedKeyUsagePurpose::ServerAuth]; + server_params + .subject_alt_names + .push(SanType::DnsName("localhost".try_into()?)); + server_params + .subject_alt_names + .push(SanType::IpAddress("127.0.0.1".parse::()?)); + match additional_san.parse::() { + Ok(ip) => server_params.subject_alt_names.push(SanType::IpAddress(ip)), + Err(_) => server_params + .subject_alt_names + .push(SanType::DnsName(additional_san.try_into()?)), + } + + let server_cert = server_params.signed_by(&server_key, &ca)?; + let ca_cert_pem = ca.pem(); + fs::write(tls_dir.join("rustfs_cert.pem"), server_cert.pem()).await?; + fs::write(tls_dir.join("rustfs_key.pem"), server_key.serialize_pem()).await?; + fs::write(tls_dir.join("ca.crt"), &ca_cert_pem).await?; + + Ok(ca_cert_pem) +} + +async fn start_https_rustfs_server(env: &mut RustFSTestEnvironment, tls_dir: &Path) -> Result<(), Box> { + let binary_path = rustfs_binary_path(); + let process = Command::new(&binary_path) + .env("RUST_LOG", "rustfs=info,rustfs_notify=debug") + .env("RUSTFS_TLS_PATH", tls_dir) + .env("RUSTFS_CONSOLE_ENABLE", "false") + .args([ + "--address", + &env.address, + "--access-key", + &env.access_key, + "--secret-key", + &env.secret_key, + &env.temp_dir, + ]) + .spawn()?; + env.process = Some(process); + Ok(()) +} + +async fn wait_for_https_server_ready( + client: &reqwest::Client, + env: &RustFSTestEnvironment, +) -> Result<(), Box> { + let url = format!("{}/", env.url); + + for _ in 0..60 { + match signed_request_with_client(client, http::Method::GET, &url, &env.access_key, &env.secret_key, None, None).await { + Ok(response) if response.status().is_success() => return Ok(()), + Ok(_) | Err(_) => sleep(Duration::from_millis(500)).await, + } + } + + Err("RustFS HTTPS server failed to become ready within 30 seconds".into()) +} + +async fn ensure_https_bucket_exists( + client: &reqwest::Client, + env: &RustFSTestEnvironment, + bucket: &str, +) -> Result<(), Box> { + let bucket_url = format!("{}/{bucket}/", env.url); + let response = + signed_request_with_client(client, http::Method::HEAD, &bucket_url, &env.access_key, &env.secret_key, None, None).await?; + + if response.status() == StatusCode::OK { + return Ok(()); + } + + let response = signed_request_with_client( + client, + http::Method::PUT, + &bucket_url, + &env.access_key, + &env.secret_key, + Some(Vec::new()), + None, + ) + .await?; + match response.status() { + StatusCode::OK | StatusCode::CONFLICT => Ok(()), + status => Err(format!("unexpected HTTPS bucket setup status: {status}").into()), + } +} + +async fn enable_bucket_versioning_over_https( + client: &reqwest::Client, + env: &RustFSTestEnvironment, + bucket: &str, +) -> Result<(), Box> { + let body = r#"Enabled"#; + let url = format!("{}/{bucket}?versioning", env.url); + let response = signed_request_with_client( + client, + http::Method::PUT, + &url, + &env.access_key, + &env.secret_key, + Some(body.as_bytes().to_vec()), + Some("application/xml"), + ) + .await?; + + if response.status() != StatusCode::OK { + let status = response.status(); + let body = response.text().await.unwrap_or_default(); + return Err(format!("enable HTTPS bucket versioning failed: {status} {body}").into()); + } + + Ok(()) +} + +async fn wait_for_replicated_object_over_https( + client: &reqwest::Client, + env: &RustFSTestEnvironment, + bucket: &str, + key: &str, + expected_body: &str, +) -> Result<(), Box> { + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + let url = format!("{}/{bucket}/{key}", env.url); + + loop { + let response = + signed_request_with_client(client, http::Method::GET, &url, &env.access_key, &env.secret_key, None, None).await?; + + match response.status() { + StatusCode::OK => { + let body = response.text().await?; + if body == expected_body { + return Ok(()); + } + return Err(format!("replicated HTTPS object body mismatch: expected {expected_body}, got {body}").into()); + } + StatusCode::NOT_FOUND if tokio::time::Instant::now() < deadline => { + sleep(Duration::from_secs(1)).await; + } + status if tokio::time::Instant::now() < deadline => { + let body = response.text().await.unwrap_or_default(); + if body.contains("NoSuchKey") || body.contains("NotFound") { + sleep(Duration::from_secs(1)).await; + continue; + } + return Err(format!("unexpected HTTPS replication read status: {status} {body}").into()); + } + status => { + let body = response.text().await.unwrap_or_default(); + return Err(format!("HTTPS replicated object was not readable in time: {status} {body}").into()); + } + } + } +} + fn create_user_s3_client(env: &RustFSTestEnvironment, access_key: &str, secret_key: &str) -> Client { let credentials = Credentials::new(access_key, secret_key, None, None, "e2e-site-replication"); let config = Config::builder() @@ -1406,6 +1720,359 @@ async fn test_set_remote_target_rejects_invalid_target_url() -> Result<(), Box Result<(), Box> { + init_logging(); + + let mut source_env = new_private_tmp_test_env() + .await + .map_err(|err| std::io::Error::other(format!("create source env failed: {err}")))?; + source_env + .start_rustfs_server(vec![]) + .await + .map_err(|err| std::io::Error::other(format!("start source HTTP server failed: {err}")))?; + + let mut target_env = new_private_tmp_https_target_env() + .await + .map_err(|err| std::io::Error::other(format!("create target env failed: {err}")))?; + 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 + .map_err(|err| std::io::Error::other(format!("generate self-signed TLS material failed: {err}")))?; + start_https_rustfs_server(&mut target_env, &tls_dir) + .await + .map_err(|err| std::io::Error::other(format!("start target HTTPS server failed: {err}")))?; + let https_client = + insecure_https_client().map_err(|err| std::io::Error::other(format!("build HTTPS client failed: {err}")))?; + wait_for_https_server_ready(&https_client, &target_env) + .await + .map_err(|err| std::io::Error::other(format!("wait for target HTTPS server ready failed: {err}")))?; + + let source_bucket = "replication-self-signed-src"; + let target_bucket = "replication-self-signed-dst"; + + let source_client = source_env.create_s3_client(); + source_client + .create_bucket() + .bucket(source_bucket) + .send() + .await + .map_err(|err| std::io::Error::other(format!("create source bucket failed: {err}")))?; + enable_bucket_versioning(&source_env, source_bucket) + .await + .map_err(|err| std::io::Error::other(format!("enable source bucket versioning failed: {err}")))?; + + ensure_https_bucket_exists(&https_client, &target_env, target_bucket) + .await + .map_err(|err| std::io::Error::other(format!("create target HTTPS bucket failed: {err}")))?; + enable_bucket_versioning_over_https(&https_client, &target_env, target_bucket) + .await + .map_err(|err| std::io::Error::other(format!("enable target HTTPS bucket versioning failed: {err}")))?; + + let err = set_replication_target_with_options( + &source_env, + source_bucket, + ReplicationTargetOptions { + endpoint: target_env.url.trim_start_matches("https://"), + access_key: &target_env.access_key, + secret_key: &target_env.secret_key, + target_bucket, + secure: true, + skip_tls_verify: false, + ca_cert_pem: None, + }, + ) + .await + .expect_err("self-signed HTTPS target should fail without skipTlsVerify"); + let err = err.to_string(); + + assert!(err.contains("400 Bad Request"), "unexpected HTTPS target setup error: {err}"); + assert!(err.contains("InvalidRequest"), "unexpected HTTPS target setup error: {err}"); + assert!( + err.to_ascii_lowercase().contains("certificate") || err.to_ascii_lowercase().contains("tls"), + "unexpected HTTPS target setup error: {err}" + ); + + Ok(()) +} + +#[tokio::test] +#[serial] +async fn test_set_remote_target_allows_self_signed_https_target_with_skip_tls_verify() -> Result<(), Box> +{ + init_logging(); + + let mut source_env = new_private_tmp_test_env() + .await + .map_err(|err| std::io::Error::other(format!("create source env failed: {err}")))?; + source_env + .start_rustfs_server(vec![]) + .await + .map_err(|err| std::io::Error::other(format!("start source HTTP server failed: {err}")))?; + + let mut target_env = new_private_tmp_https_target_env() + .await + .map_err(|err| std::io::Error::other(format!("create target env failed: {err}")))?; + 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 + .map_err(|err| std::io::Error::other(format!("generate self-signed TLS material failed: {err}")))?; + start_https_rustfs_server(&mut target_env, &tls_dir) + .await + .map_err(|err| std::io::Error::other(format!("start target HTTPS server failed: {err}")))?; + let https_client = + insecure_https_client().map_err(|err| std::io::Error::other(format!("build HTTPS client failed: {err}")))?; + wait_for_https_server_ready(&https_client, &target_env) + .await + .map_err(|err| std::io::Error::other(format!("wait for target HTTPS server ready failed: {err}")))?; + + let source_bucket = "replication-self-signed-ok-src"; + let target_bucket = "replication-self-signed-ok-dst"; + let object_key = "self-signed-replication.txt"; + let body = "replication over self-signed https should succeed"; + + let source_client = source_env.create_s3_client(); + source_client + .create_bucket() + .bucket(source_bucket) + .send() + .await + .map_err(|err| std::io::Error::other(format!("create source bucket failed: {err}")))?; + enable_bucket_versioning(&source_env, source_bucket) + .await + .map_err(|err| std::io::Error::other(format!("enable source bucket versioning failed: {err}")))?; + + ensure_https_bucket_exists(&https_client, &target_env, target_bucket) + .await + .map_err(|err| std::io::Error::other(format!("create target HTTPS bucket failed: {err}")))?; + enable_bucket_versioning_over_https(&https_client, &target_env, target_bucket) + .await + .map_err(|err| std::io::Error::other(format!("enable target HTTPS bucket versioning failed: {err}")))?; + + let target_arn = set_replication_target_with_options( + &source_env, + source_bucket, + ReplicationTargetOptions { + endpoint: target_env.url.trim_start_matches("https://"), + access_key: &target_env.access_key, + secret_key: &target_env.secret_key, + target_bucket, + secure: true, + skip_tls_verify: true, + ca_cert_pem: None, + }, + ) + .await?; + put_bucket_replication(&source_env, source_bucket, &target_arn).await?; + + let response = run_replication_check(&source_env, source_bucket).await?; + assert_eq!(response.status(), StatusCode::OK); + + source_client + .put_object() + .bucket(source_bucket) + .key(object_key) + .body(ByteStream::from(body.as_bytes().to_vec())) + .send() + .await?; + + wait_for_replicated_object_over_https(&https_client, &target_env, target_bucket, object_key, body).await?; + + Ok(()) +} + +#[tokio::test] +#[serial] +async fn test_set_remote_target_rejects_private_ca_https_target_without_ca_cert_pem() -> Result<(), Box> +{ + init_logging(); + + let mut source_env = new_private_tmp_test_env() + .await + .map_err(|err| std::io::Error::other(format!("create source env failed: {err}")))?; + source_env + .start_rustfs_server(vec![]) + .await + .map_err(|err| std::io::Error::other(format!("start source HTTP server failed: {err}")))?; + + let mut target_env = new_private_tmp_https_target_env() + .await + .map_err(|err| std::io::Error::other(format!("create target env failed: {err}")))?; + 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 + .map_err(|err| std::io::Error::other(format!("generate private CA TLS material failed: {err}")))?; + start_https_rustfs_server(&mut target_env, &tls_dir) + .await + .map_err(|err| std::io::Error::other(format!("start target HTTPS server failed: {err}")))?; + let https_client = + trusted_https_client(&ca_cert_pem).map_err(|err| std::io::Error::other(format!("build HTTPS client failed: {err}")))?; + wait_for_https_server_ready(&https_client, &target_env) + .await + .map_err(|err| std::io::Error::other(format!("wait for target HTTPS server ready failed: {err}")))?; + + let source_bucket = "replication-private-ca-src"; + let target_bucket = "replication-private-ca-dst"; + + let source_client = source_env.create_s3_client(); + source_client + .create_bucket() + .bucket(source_bucket) + .send() + .await + .map_err(|err| std::io::Error::other(format!("create source bucket failed: {err}")))?; + enable_bucket_versioning(&source_env, source_bucket) + .await + .map_err(|err| std::io::Error::other(format!("enable source bucket versioning failed: {err}")))?; + + ensure_https_bucket_exists(&https_client, &target_env, target_bucket) + .await + .map_err(|err| std::io::Error::other(format!("create target HTTPS bucket failed: {err}")))?; + enable_bucket_versioning_over_https(&https_client, &target_env, target_bucket) + .await + .map_err(|err| std::io::Error::other(format!("enable target HTTPS bucket versioning failed: {err}")))?; + + let err = set_replication_target_with_options( + &source_env, + source_bucket, + ReplicationTargetOptions { + endpoint: target_env.url.trim_start_matches("https://"), + access_key: &target_env.access_key, + secret_key: &target_env.secret_key, + target_bucket, + secure: true, + skip_tls_verify: false, + ca_cert_pem: None, + }, + ) + .await + .expect_err("private CA HTTPS target should fail without caCertPem"); + let err = err.to_string(); + + assert!(err.contains("400 Bad Request"), "unexpected private CA target setup error: {err}"); + assert!(err.contains("InvalidRequest"), "unexpected private CA target setup error: {err}"); + assert!( + err.to_ascii_lowercase().contains("certificate") || err.to_ascii_lowercase().contains("tls"), + "unexpected private CA target setup error: {err}" + ); + + Ok(()) +} + +#[tokio::test] +#[serial] +async fn test_set_remote_target_allows_private_ca_https_target_with_ca_cert_pem() -> Result<(), Box> { + init_logging(); + + let mut source_env = new_private_tmp_test_env() + .await + .map_err(|err| std::io::Error::other(format!("create source env failed: {err}")))?; + source_env + .start_rustfs_server(vec![]) + .await + .map_err(|err| std::io::Error::other(format!("start source HTTP server failed: {err}")))?; + + let mut target_env = new_private_tmp_https_target_env() + .await + .map_err(|err| std::io::Error::other(format!("create target env failed: {err}")))?; + 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 + .map_err(|err| std::io::Error::other(format!("generate private CA TLS material failed: {err}")))?; + start_https_rustfs_server(&mut target_env, &tls_dir) + .await + .map_err(|err| std::io::Error::other(format!("start target HTTPS server failed: {err}")))?; + let https_client = + trusted_https_client(&ca_cert_pem).map_err(|err| std::io::Error::other(format!("build HTTPS client failed: {err}")))?; + wait_for_https_server_ready(&https_client, &target_env) + .await + .map_err(|err| std::io::Error::other(format!("wait for target HTTPS server ready failed: {err}")))?; + + let source_bucket = "replication-private-ca-ok-src"; + let target_bucket = "replication-private-ca-ok-dst"; + let object_key = "private-ca-replication.txt"; + let body = "replication over private ca https should succeed"; + + let source_client = source_env.create_s3_client(); + source_client + .create_bucket() + .bucket(source_bucket) + .send() + .await + .map_err(|err| std::io::Error::other(format!("create source bucket failed: {err}")))?; + enable_bucket_versioning(&source_env, source_bucket) + .await + .map_err(|err| std::io::Error::other(format!("enable source bucket versioning failed: {err}")))?; + + ensure_https_bucket_exists(&https_client, &target_env, target_bucket) + .await + .map_err(|err| std::io::Error::other(format!("create target HTTPS bucket failed: {err}")))?; + enable_bucket_versioning_over_https(&https_client, &target_env, target_bucket) + .await + .map_err(|err| std::io::Error::other(format!("enable target HTTPS bucket versioning failed: {err}")))?; + + let target_arn = set_replication_target_with_options( + &source_env, + source_bucket, + ReplicationTargetOptions { + endpoint: target_env.url.trim_start_matches("https://"), + access_key: &target_env.access_key, + secret_key: &target_env.secret_key, + target_bucket, + secure: true, + skip_tls_verify: false, + ca_cert_pem: Some(&ca_cert_pem), + }, + ) + .await?; + put_bucket_replication(&source_env, source_bucket, &target_arn).await?; + + let response = run_replication_check(&source_env, source_bucket).await?; + assert_eq!(response.status(), StatusCode::OK); + + source_client + .put_object() + .bucket(source_bucket) + .key(object_key) + .body(ByteStream::from(body.as_bytes().to_vec())) + .send() + .await?; + + wait_for_replicated_object_over_https(&https_client, &target_env, target_bucket, object_key, body).await?; + + Ok(()) +} + #[tokio::test] #[serial] async fn test_list_remote_targets_rejects_empty_bucket() -> Result<(), Box> { diff --git a/crates/ecstore/Cargo.toml b/crates/ecstore/Cargo.toml index 5f820a4f8..9eaf91c8c 100644 --- a/crates/ecstore/Cargo.toml +++ b/crates/ecstore/Cargo.toml @@ -120,6 +120,7 @@ shadow-rs.workspace = true async-recursion.workspace = true aws-credential-types = { workspace = true } aws-smithy-types = { workspace = true } +aws-smithy-runtime-api = { workspace = true } parking_lot = { workspace = true } base64-simd.workspace = true serde_urlencoded.workspace = true @@ -141,6 +142,7 @@ tracing-subscriber = { workspace = true, features = ["json"] } serial_test = { workspace = true } opentelemetry_sdk = { workspace = true } proptest = "1" +rcgen.workspace = true [build-dependencies] shadow-rs = { workspace = true, features = ["build", "metadata"] } diff --git a/crates/ecstore/src/bucket/bucket_target_sys.rs b/crates/ecstore/src/bucket/bucket_target_sys.rs index 5f0b5692f..1db50ead2 100644 --- a/crates/ecstore/src/bucket/bucket_target_sys.rs +++ b/crates/ecstore/src/bucket/bucket_target_sys.rs @@ -24,6 +24,7 @@ use crate::bucket::versioning_sys::BucketVersioningSys; use crate::runtime_sources; use aws_credential_types::Credentials as SdkCredentials; use aws_sdk_s3::config::Region as SdkRegion; +use aws_sdk_s3::config::SharedHttpClient; use aws_sdk_s3::error::ProvideErrorMetadata; use aws_sdk_s3::error::SdkError; use aws_sdk_s3::operation::complete_multipart_upload::CompleteMultipartUploadOutput; @@ -37,11 +38,20 @@ use aws_sdk_s3::types::{ use aws_sdk_s3::{Client as S3Client, Config as S3Config, operation::head_object::HeadObjectOutput}; use aws_sdk_s3::{config::SharedCredentialsProvider, types::BucketVersioningStatus}; use aws_smithy_http_client::{Builder as SmithyHttpClientBuilder, tls as smithy_tls}; -use http::{HeaderMap, HeaderName, HeaderValue, StatusCode}; +use aws_smithy_runtime_api::box_error::BoxError; +use aws_smithy_runtime_api::client::http::{ + HttpConnector as SmithyHttpConnector, HttpConnectorFuture, SharedHttpConnector, http_client_fn, +}; +use aws_smithy_runtime_api::client::orchestrator::{HttpRequest, HttpResponse}; +use aws_smithy_runtime_api::client::result::ConnectorError; +use aws_smithy_types::body::SdkBody; +use http::{HeaderMap, HeaderName, HeaderValue, StatusCode, Uri}; +use hyper_util::client::legacy::Client as HyperClient; +use hyper_util::rt::{TokioExecutor, TokioTimer}; use reqwest::Client as HttpClient; use rustfs_config::{DEFAULT_TRUST_LEAF_CERT_AS_CA, ENV_TRUST_LEAF_CERT_AS_CA, RUSTFS_CA_CERT, RUSTFS_TLS_CERT}; use rustfs_filemeta::{ReplicationStatusType, ReplicationType}; -use rustfs_utils::egress::validate_outbound_url; +use rustfs_utils::egress::{OutboundUrlError, validate_outbound_url}; use rustfs_utils::http::{ AMZ_BUCKET_REPLICATION_STATUS, AMZ_OBJECT_LOCK_BYPASS_GOVERNANCE, AMZ_OBJECT_LOCK_LEGAL_HOLD, AMZ_OBJECT_LOCK_MODE, AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE, AMZ_STORAGE_CLASS, AMZ_WEBSITE_REDIRECT_LOCATION, is_amz_header, is_minio_header, @@ -51,6 +61,7 @@ use rustfs_utils::http::{ SUFFIX_FORCE_DELETE, SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_ETAG, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_REPLICATION_CHECK, SUFFIX_SOURCE_REPLICATION_REQUEST, SUFFIX_SOURCE_VERSION_ID, insert_header, }; +use rustls_pki_types::pem::PemObject; use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::error::Error; @@ -63,6 +74,7 @@ use std::time::{Duration, Instant}; use time::{OffsetDateTime, format_description::well_known::Rfc3339}; use tokio::sync::Mutex; use tokio::sync::RwLock; +use tower::Service; use tracing::error; use tracing::warn; use url::Url; @@ -652,7 +664,7 @@ impl BucketTargetSys { access_key: credentials.access_key.clone(), error: format!("invalid target endpoint: {err}"), })?; - validate_outbound_url(&parsed_endpoint).map_err(|err| BucketTargetError::RemoteTargetConnectionErr { + validate_replication_target_endpoint(&parsed_endpoint).map_err(|err| BucketTargetError::RemoteTargetConnectionErr { bucket: target.target_bucket.clone(), access_key: credentials.access_key.clone(), error: format!("target endpoint is not allowed: {err}"), @@ -668,8 +680,14 @@ impl BucketTargetSys { config_builder = config_builder.force_path_style(true); } - if target.secure - && let Some(http_client) = build_aws_s3_http_client_from_tls_path().await + if let Some(http_client) = + build_aws_s3_http_client_for_target(target) + .await + .map_err(|err| BucketTargetError::RemoteTargetConnectionErr { + bucket: target.target_bucket.clone(), + access_key: credentials.access_key.clone(), + error: err.to_string(), + })? { config_builder = config_builder.http_client(http_client); } @@ -821,6 +839,170 @@ impl BucketTargetSys { } } +#[derive(Debug)] +struct AcceptAnyServerCertVerifier; + +impl rustls::client::danger::ServerCertVerifier for AcceptAnyServerCertVerifier { + fn verify_server_cert( + &self, + _end_entity: &rustls_pki_types::CertificateDer<'_>, + _intermediates: &[rustls_pki_types::CertificateDer<'_>], + _server_name: &rustls_pki_types::ServerName<'_>, + _ocsp_response: &[u8], + _now: rustls_pki_types::UnixTime, + ) -> Result { + Ok(rustls::client::danger::ServerCertVerified::assertion()) + } + + fn verify_tls12_signature( + &self, + _message: &[u8], + _cert: &rustls_pki_types::CertificateDer<'_>, + _dss: &rustls::DigitallySignedStruct, + ) -> Result { + Ok(rustls::client::danger::HandshakeSignatureValid::assertion()) + } + + fn verify_tls13_signature( + &self, + _message: &[u8], + _cert: &rustls_pki_types::CertificateDer<'_>, + _dss: &rustls::DigitallySignedStruct, + ) -> Result { + Ok(rustls::client::danger::HandshakeSignatureValid::assertion()) + } + + fn supported_verify_schemes(&self) -> Vec { + rustls::crypto::aws_lc_rs::default_provider() + .signature_verification_algorithms + .supported_schemes() + } +} + +#[derive(Clone)] +struct TargetHyperHttpConnector { + client: HyperClient, +} + +impl fmt::Debug for TargetHyperHttpConnector { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("TargetHyperHttpConnector") + .field("client", &"** hyper client **") + .finish() + } +} + +impl SmithyHttpConnector for TargetHyperHttpConnector +where + C: Clone + Send + Sync + 'static, + C: Service, + C::Response: + hyper::rt::Read + hyper::rt::Write + hyper_util::client::legacy::connect::Connection + Send + Sync + Unpin + 'static, + C::Future: Unpin + Send + 'static, + C::Error: Into, +{ + fn call(&self, request: HttpRequest) -> HttpConnectorFuture { + let request = match request.try_into_http1x() { + Ok(request) => request, + Err(err) => return HttpConnectorFuture::ready(Err(ConnectorError::user(err.into()))), + }; + + let mut client = self.client.clone(); + let fut = client.call(request); + HttpConnectorFuture::new(async move { + let response = fut + .await + .map_err(|err| ConnectorError::io(err.into()))? + .map(SdkBody::from_body_1_x); + HttpResponse::try_from(response).map_err(|err| ConnectorError::other(err.into(), None)) + }) + } +} + +fn ensure_rustls_crypto_provider() { + if rustls::crypto::CryptoProvider::get_default().is_none() { + let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); + } +} + +fn has_custom_ca_pem(target: &BucketTarget) -> bool { + !target.ca_cert_pem.trim().is_empty() +} + +fn validate_replication_target_endpoint(url: &Url) -> Result<(), OutboundUrlError> { + match validate_outbound_url(url) { + Ok(()) => Ok(()), + Err(OutboundUrlError::ForbiddenHost { + reason: "private address", + .. + }) => Ok(()), + Err(err) => Err(err), + } +} + +fn build_insecure_aws_s3_http_client() -> SharedHttpClient { + ensure_rustls_crypto_provider(); + + let tls_config = rustls::ClientConfig::builder() + .dangerous() + .with_custom_certificate_verifier(Arc::new(AcceptAnyServerCertVerifier)) + .with_no_client_auth(); + + let https = hyper_rustls::HttpsConnectorBuilder::new() + .with_tls_config(tls_config) + .https_or_http() + .enable_http1() + .enable_http2() + .build(); + let mut client_builder = HyperClient::builder(TokioExecutor::new()); + client_builder.pool_timer(TokioTimer::new()); + let client = client_builder.build(https); + let connector = SharedHttpConnector::new(TargetHyperHttpConnector { client }); + + http_client_fn(move |_settings, _components| connector.clone()) +} + +fn build_aws_s3_http_client_from_target_ca_pem(ca_cert_pem: &str) -> Result { + let certs = rustls_pki_types::CertificateDer::pem_slice_iter(ca_cert_pem.as_bytes()) + .collect::, _>>() + .map_err(|err| BucketTargetError::Io(std::io::Error::other(format!("invalid target CA PEM: {err}"))))?; + + if certs.is_empty() { + return Err(BucketTargetError::Io(std::io::Error::other( + "invalid target CA PEM: no certificates found", + ))); + } + + let mut trust_store = smithy_tls::TrustStore::empty(); + trust_store.add_pem_certificate(ca_cert_pem.as_bytes()); + + let tls_context = smithy_tls::TlsContext::builder() + .with_trust_store(trust_store) + .build() + .map_err(|err| BucketTargetError::Io(std::io::Error::other(format!("invalid target CA PEM: {err}"))))?; + + Ok(SmithyHttpClientBuilder::new() + .tls_provider(smithy_tls::Provider::rustls(smithy_tls::rustls_provider::CryptoMode::AwsLc)) + .tls_context(tls_context) + .build_https()) +} + +async fn build_aws_s3_http_client_for_target(target: &BucketTarget) -> Result, BucketTargetError> { + if !target.secure { + return Ok(None); + } + + if target.skip_tls_verify { + return Ok(Some(build_insecure_aws_s3_http_client())); + } + + if has_custom_ca_pem(target) { + return build_aws_s3_http_client_from_target_ca_pem(&target.ca_cert_pem).map(Some); + } + + Ok(build_aws_s3_http_client_from_tls_path().await) +} + async fn build_aws_s3_http_client_from_tls_path() -> Option { 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() { @@ -828,7 +1010,7 @@ async fn build_aws_s3_http_client_from_tls_path() -> Option S3Error { } } +fn validate_remote_target_tls_settings(remote_target: &BucketTarget) -> S3Result<()> { + let has_custom_ca = !remote_target.ca_cert_pem.trim().is_empty(); + + if !remote_target.secure && remote_target.skip_tls_verify { + return Err(s3_error!(InvalidRequest, "skipTlsVerify requires an HTTPS remote target")); + } + + if !remote_target.secure && has_custom_ca { + return Err(s3_error!(InvalidRequest, "caCertPem requires an HTTPS remote target")); + } + + if remote_target.skip_tls_verify && has_custom_ca { + return Err(s3_error!(InvalidRequest, "skipTlsVerify and caCertPem cannot be enabled together")); + } + + Ok(()) +} + pub fn register_replication_route(r: &mut S3Router) -> std::io::Result<()> { r.insert( Method::GET, @@ -213,6 +231,7 @@ impl Operation for SetRemoteTargetHandler { error!("Failed to parse BucketTarget from body: {}", e); ApiError::other(e) })?; + validate_remote_target_tls_settings(&remote_target)?; let Ok(target_url) = remote_target.url() else { return Err(s3_error!(InvalidRequest, "invalid target url")); @@ -278,9 +297,19 @@ impl Operation for SetRemoteTargetHandler { target.path = remote_target.path; target.replication_sync = remote_target.replication_sync; target.bandwidth_limit = remote_target.bandwidth_limit; + target.skip_tls_verify = remote_target.skip_tls_verify; + target.ca_cert_pem = remote_target.ca_cert_pem; target.health_check_duration = remote_target.health_check_duration; - warn!("update target, target: {:?}", target); + warn!( + bucket = %bucket, + arn = %target.arn, + endpoint = %target.endpoint, + secure = target.secure, + skip_tls_verify = target.skip_tls_verify, + has_custom_ca = !target.ca_cert_pem.trim().is_empty(), + "update remote target" + ); remote_target = target; } @@ -418,7 +447,8 @@ impl Operation for RemoveRemoteTargetHandler { #[cfg(test)] mod tests { - use super::extract_query_params; + use super::{extract_query_params, validate_remote_target_tls_settings}; + use crate::admin::target::BucketTarget; use http::Uri; #[test] @@ -431,4 +461,54 @@ mod tests { assert_eq!(params.get("bucket"), Some(&"foo/bar".to_string())); assert_eq!(params.get("flag"), Some(&"a b".to_string())); } + + #[test] + fn validate_remote_target_tls_settings_rejects_insecure_tls_for_http_targets() { + let err = validate_remote_target_tls_settings(&BucketTarget { + secure: false, + skip_tls_verify: true, + ..Default::default() + }) + .expect_err("HTTP targets must reject skipTlsVerify"); + + assert!(err.to_string().contains("skipTlsVerify requires an HTTPS remote target")); + } + + #[test] + fn validate_remote_target_tls_settings_rejects_custom_ca_for_http_targets() { + let err = validate_remote_target_tls_settings(&BucketTarget { + secure: false, + ca_cert_pem: "-----BEGIN CERTIFICATE-----\nMIIB\n-----END CERTIFICATE-----\n".to_string(), + ..Default::default() + }) + .expect_err("HTTP targets must reject custom CA PEM"); + + assert!(err.to_string().contains("caCertPem requires an HTTPS remote target")); + } + + #[test] + fn validate_remote_target_tls_settings_rejects_insecure_and_custom_ca_combination() { + let err = validate_remote_target_tls_settings(&BucketTarget { + secure: true, + skip_tls_verify: true, + ca_cert_pem: "-----BEGIN CERTIFICATE-----\nMIIB\n-----END CERTIFICATE-----\n".to_string(), + ..Default::default() + }) + .expect_err("custom CA and insecure TLS must be mutually exclusive"); + + assert!( + err.to_string() + .contains("skipTlsVerify and caCertPem cannot be enabled together") + ); + } + + #[test] + fn validate_remote_target_tls_settings_allows_https_insecure_without_custom_ca() { + validate_remote_target_tls_settings(&BucketTarget { + secure: true, + skip_tls_verify: true, + ..Default::default() + }) + .expect("HTTPS targets should allow skipTlsVerify when no custom CA is configured"); + } }