diff --git a/cloud/aws-common/src/main/java/org/apache/druid/common/aws/AWSClientConfig.java b/cloud/aws-common/src/main/java/org/apache/druid/common/aws/AWSClientConfig.java index 67f31cfbe981..df8f2db10a7d 100644 --- a/cloud/aws-common/src/main/java/org/apache/druid/common/aws/AWSClientConfig.java +++ b/cloud/aws-common/src/main/java/org/apache/druid/common/aws/AWSClientConfig.java @@ -38,6 +38,7 @@ public class AWSClientConfig // Default values matching AWS SDK v2 defaults private static final boolean DEFAULT_CHUNKED_ENCODING_DISABLED = false; private static final boolean DEFAULT_PATH_STYLE_ACCESS = false; + private static final boolean DEFAULT_LEGACY_MD5_ENABLED = false; private static final int DEFAULT_CONNECTION_TIMEOUT_MILLIS = 10_000; private static final int DEFAULT_SOCKET_TIMEOUT_MILLIS = 50_000; @@ -117,6 +118,9 @@ public static RetryMode fromString(String value) @JsonProperty private boolean enablePathStyleAccess = DEFAULT_PATH_STYLE_ACCESS; + @JsonProperty + private boolean enableLegacyMd5 = DEFAULT_LEGACY_MD5_ENABLED; + /** * @deprecated Use {@link #crossRegionAccessEnabled} instead. */ @@ -178,6 +182,11 @@ public boolean isEnablePathStyleAccess() return enablePathStyleAccess; } + public boolean isEnableLegacyMd5() + { + return enableLegacyMd5; + } + /** * @deprecated Use {@link #isCrossRegionAccessEnabled()} instead. */ @@ -272,6 +281,7 @@ public String toString() "protocol='" + protocol + '\'' + ", disableChunkedEncoding=" + disableChunkedEncoding + ", enablePathStyleAccess=" + enablePathStyleAccess + + ", enableLegacyMd5=" + enableLegacyMd5 + ", crossRegionAccessEnabled=" + isCrossRegionAccessEnabled() + ", connectionTimeout=" + connectionTimeout + ", socketTimeout=" + socketTimeout + diff --git a/cloud/aws-common/src/test/java/org/apache/druid/common/aws/AWSClientConfigTest.java b/cloud/aws-common/src/test/java/org/apache/druid/common/aws/AWSClientConfigTest.java index a86a2900af86..98ead3485673 100644 --- a/cloud/aws-common/src/test/java/org/apache/druid/common/aws/AWSClientConfigTest.java +++ b/cloud/aws-common/src/test/java/org/apache/druid/common/aws/AWSClientConfigTest.java @@ -190,6 +190,18 @@ public void testDeprecatedPropertyStaysUnsetWhenOnlyItsReplacementIsBound() Assertions.assertNull(bind(Map.of("crossRegionAccessEnabled", true)).isForceGlobalBucketAccessEnabled()); } + @Test + public void testLegacyMd5DisabledByDefault() + { + Assertions.assertFalse(bind(Map.of()).isEnableLegacyMd5()); + } + + @Test + public void testLegacyMd5CanBeEnabled() + { + Assertions.assertTrue(bind(Map.of("enableLegacyMd5", true)).isEnableLegacyMd5()); + } + @ParameterizedTest(name = "{0} processors -> {1} connections") @CsvSource({"8, 50", "32, 128"}) public void testDefaultMaxConnectionsTakesTheSdkFloorOrFourPerCore(int processors, int expected) diff --git a/docs/development/extensions-core/s3.md b/docs/development/extensions-core/s3.md index b60e80585e51..9a00e76e4ebe 100644 --- a/docs/development/extensions-core/s3.md +++ b/docs/development/extensions-core/s3.md @@ -130,6 +130,7 @@ For example, to set the region to 'us-east-1' through system properties: |`druid.s3.protocol`|Communication protocol type to use when sending requests to AWS. `http` or `https` can be used. This configuration would be ignored if `druid.s3.endpoint.url` is filled with a URL with a different protocol.|`https`| |`druid.s3.disableChunkedEncoding`|Disables chunked encoding. See [AWS document](https://docs.aws.amazon.com/AWSJavaSDK/latest/javadoc/com/amazonaws/services/s3/AmazonS3Builder.html#disableChunkedEncoding--) for details.|false| |`druid.s3.enablePathStyleAccess`|Enables path style access. See [AWS document](https://docs.aws.amazon.com/AWSJavaSDK/latest/javadoc/com/amazonaws/services/s3/AmazonS3Builder.html#enablePathStyleAccess--) for details.|false| +|`druid.s3.enableLegacyMd5`|Uses `Content-MD5` for operations that require a request checksum and disables the optional `CRC32` checksums that AWS SDK v2 calculates by default since version 2.30.0. Enable this only for S3-compatible storage that rejects the newer checksum headers. See [AWS document](https://docs.aws.amazon.com/sdk-for-java/latest/developer-guide/s3-checksums.html) for details.|false| |`druid.s3.crossRegionAccessEnabled`|Enables cross-region access for S3 requests. When enabled, the S3 client automatically detects the correct region for a bucket on first access and caches it for subsequent requests.|false| |`druid.s3.retryMode`|Retry strategy for AWS clients built from this configuration. One of `standard`, `adaptive` or `legacy`. See [AWS document](https://docs.aws.amazon.com/sdkref/latest/guide/feature-retry-behavior.html) for details about each mode.|`standard`| |`druid.s3.maxAttempts`|Total attempts per HTTP request, including the first, so `1` disables SDK retries. When unset, the SDK default count defined by `druid.s3.retryMode` applies.|null (the retry mode's own default)| diff --git a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3Utils.java b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3Utils.java index 0c2546302a13..0ad2c17cd50e 100644 --- a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3Utils.java +++ b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/S3Utils.java @@ -33,9 +33,12 @@ import org.apache.druid.java.util.common.StringUtils; import org.apache.druid.java.util.common.URIs; import org.apache.druid.java.util.common.logger.Logger; +import software.amazon.awssdk.core.checksums.RequestChecksumCalculation; import software.amazon.awssdk.core.exception.SdkClientException; import software.amazon.awssdk.core.exception.SdkException; import software.amazon.awssdk.http.apache.ProxyConfiguration; +import software.amazon.awssdk.services.s3.LegacyMd5Plugin; +import software.amazon.awssdk.services.s3.S3BaseClientBuilder; import software.amazon.awssdk.services.s3.model.Delete; import software.amazon.awssdk.services.s3.model.DeleteObjectsRequest; import software.amazon.awssdk.services.s3.model.DeleteObjectsResponse; @@ -125,6 +128,28 @@ public boolean apply(Throwable e) } }; + /** + * Restores {@code Content-MD5} for required request checksums and disables optional request checksums on every given + * builder, for S3-compatible stores that reject the CRC32 checksums the SDK sends by default since 2.30.0. + *

+ * Takes all the builders for one client set rather than one builder per call, so the sync and async clients cannot + * end up disagreeing about checksum behavior, and so the switch is logged once per client set. + */ + public static void configureLegacyMd5( + final AWSClientConfig clientConfig, + final S3BaseClientBuilder... s3ClientBuilders + ) + { + if (clientConfig.isEnableLegacyMd5()) { + log.info("Legacy MD5 compatibility mode is enabled for the S3 client."); + for (final S3BaseClientBuilder s3ClientBuilder : s3ClientBuilders) { + s3ClientBuilder + .requestChecksumCalculation(RequestChecksumCalculation.WHEN_REQUIRED) + .addPlugin(LegacyMd5Plugin.create()); + } + } + } + /** * Retries S3 operations that fail intermittently (due to io-related exceptions, during obtaining credentials, etc). * Service-level exceptions (access denied, file not found, etc) are not retried. diff --git a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/ServerSideEncryptingAmazonS3.java b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/ServerSideEncryptingAmazonS3.java index b1e6cb2f0439..81884ea0cf46 100644 --- a/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/ServerSideEncryptingAmazonS3.java +++ b/extensions-core/s3-extensions/src/main/java/org/apache/druid/storage/s3/ServerSideEncryptingAmazonS3.java @@ -389,6 +389,7 @@ public static ServerSideEncryptingAmazonS3.Builder builder( MultipartConfiguration.builder().minimumPartSizeInBytes(transferConfig.getMinimumUploadPartSize()) .thresholdInBytes(transferConfig.getMultipartUploadThreshold()) .build()); + S3Utils.configureLegacyMd5(awsClientConfig, clientBuilder, asyncClientBuilder); } // Configure HTTP client with proxy if needed diff --git a/extensions-core/s3-extensions/src/test/java/org/apache/druid/data/input/s3/S3InputSourceTest.java b/extensions-core/s3-extensions/src/test/java/org/apache/druid/data/input/s3/S3InputSourceTest.java index ba3912c49a32..0e041a7640d0 100644 --- a/extensions-core/s3-extensions/src/test/java/org/apache/druid/data/input/s3/S3InputSourceTest.java +++ b/extensions-core/s3-extensions/src/test/java/org/apache/druid/data/input/s3/S3InputSourceTest.java @@ -516,6 +516,7 @@ public void testS3InputSourceUseEndPointClientProxy() EasyMock.expect(mockAwsClientConfig.isDisableChunkedEncoding()).andStubReturn(false); EasyMock.expect(mockAwsClientConfig.isEnablePathStyleAccess()).andStubReturn(false); EasyMock.expect(mockAwsClientConfig.isCrossRegionAccessEnabled()).andStubReturn(true); + EasyMock.expect(mockAwsClientConfig.isEnableLegacyMd5()).andStubReturn(false); EasyMock.expect(mockAwsClientConfig.getProtocol()).andStubReturn("http"); EasyMock.expect(mockAwsClientConfig.getConnectionTimeoutMillis()).andStubReturn(10_000); EasyMock.expect(mockAwsClientConfig.getSocketTimeoutMillis()).andStubReturn(50_000); diff --git a/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3UtilsTest.java b/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3UtilsTest.java index 16e0257bac21..dbbd272e1ffc 100644 --- a/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3UtilsTest.java +++ b/extensions-core/s3-extensions/src/test/java/org/apache/druid/storage/s3/S3UtilsTest.java @@ -27,16 +27,35 @@ import org.easymock.EasyMock; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import software.amazon.awssdk.auth.credentials.AnonymousCredentialsProvider; +import software.amazon.awssdk.core.checksums.RequestChecksumCalculation; import software.amazon.awssdk.core.exception.SdkClientException; +import software.amazon.awssdk.core.sync.RequestBody; +import software.amazon.awssdk.http.ExecutableHttpRequest; +import software.amazon.awssdk.http.HttpExecuteRequest; +import software.amazon.awssdk.http.HttpExecuteResponse; +import software.amazon.awssdk.http.SdkHttpClient; +import software.amazon.awssdk.http.SdkHttpResponse; +import software.amazon.awssdk.regions.Region; +import software.amazon.awssdk.services.s3.LegacyMd5Plugin; +import software.amazon.awssdk.services.s3.S3AsyncClient; +import software.amazon.awssdk.services.s3.S3AsyncClientBuilder; +import software.amazon.awssdk.services.s3.S3Client; +import software.amazon.awssdk.services.s3.S3ClientBuilder; +import software.amazon.awssdk.services.s3.model.Delete; import software.amazon.awssdk.services.s3.model.DeleteObjectsRequest; import software.amazon.awssdk.services.s3.model.DeleteObjectsResponse; import software.amazon.awssdk.services.s3.model.ObjectIdentifier; +import software.amazon.awssdk.services.s3.model.PutObjectRequest; import software.amazon.awssdk.services.s3.model.S3Error; import software.amazon.awssdk.services.s3.model.S3Exception; +import software.amazon.awssdk.utils.Md5Utils; import javax.crypto.AEADBadTagException; import javax.net.ssl.SSLException; import java.io.IOException; +import java.net.URI; +import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletionException; import java.util.concurrent.atomic.AtomicInteger; @@ -44,6 +63,124 @@ public class S3UtilsTest { + @Test + public void testConfigureLegacyMd5Disabled() + { + final S3ClientBuilder s3ClientBuilder = S3Client.builder(); + + S3Utils.configureLegacyMd5(new AWSClientConfig(), s3ClientBuilder); + + Assertions.assertFalse( + s3ClientBuilder.plugins().stream().anyMatch(LegacyMd5Plugin.class::isInstance) + ); + } + + @Test + public void testConfigureLegacyMd5EnabledForSyncAndAsyncClients() + { + final AWSClientConfig clientConfig = EasyMock.createMock(AWSClientConfig.class); + EasyMock.expect(clientConfig.isEnableLegacyMd5()).andReturn(true).once(); + EasyMock.replay(clientConfig); + final S3ClientBuilder s3ClientBuilder = S3Client.builder(); + final S3AsyncClientBuilder s3AsyncClientBuilder = S3AsyncClient.builder(); + + S3Utils.configureLegacyMd5(clientConfig, s3ClientBuilder, s3AsyncClientBuilder); + + Assertions.assertEquals(1, s3ClientBuilder.plugins().size()); + Assertions.assertEquals(1, s3AsyncClientBuilder.plugins().size()); + Assertions.assertInstanceOf(LegacyMd5Plugin.class, s3ClientBuilder.plugins().get(0)); + Assertions.assertInstanceOf(LegacyMd5Plugin.class, s3AsyncClientBuilder.plugins().get(0)); + try ( + final S3Client s3Client = s3ClientBuilder + .credentialsProvider(AnonymousCredentialsProvider.create()) + .region(Region.US_EAST_1) + .build(); + final S3AsyncClient s3AsyncClient = s3AsyncClientBuilder + .credentialsProvider(AnonymousCredentialsProvider.create()) + .region(Region.US_EAST_1) + .build() + ) { + Assertions.assertSame( + RequestChecksumCalculation.WHEN_REQUIRED, + s3Client.serviceClientConfiguration().requestChecksumCalculation() + ); + Assertions.assertSame( + RequestChecksumCalculation.WHEN_REQUIRED, + s3AsyncClient.serviceClientConfiguration().requestChecksumCalculation() + ); + } + EasyMock.verify(clientConfig); + } + + @Test + public void testConfigureLegacyMd5UsesMd5ForRequiredChecksumsOnly() throws IOException + { + final AWSClientConfig clientConfig = EasyMock.createMock(AWSClientConfig.class); + EasyMock.expect(clientConfig.isEnableLegacyMd5()).andReturn(true).once(); + EasyMock.replay(clientConfig); + final List requests = new ArrayList<>(); + final SdkHttpClient httpClient = new SdkHttpClient() + { + @Override + public ExecutableHttpRequest prepareRequest(final HttpExecuteRequest request) + { + requests.add(request); + return new ExecutableHttpRequest() + { + @Override + public HttpExecuteResponse call() + { + return HttpExecuteResponse.builder() + .response(SdkHttpResponse.builder().statusCode(200).build()) + .build(); + } + + @Override + public void abort() + { + // Nothing to abort in this test client. + } + }; + } + + @Override + public void close() + { + // No resources to close in this test client. + } + }; + final S3ClientBuilder s3ClientBuilder = S3Client.builder() + .credentialsProvider(AnonymousCredentialsProvider.create()) + .region(Region.US_EAST_1) + .endpointOverride(URI.create("http://localhost")) + .forcePathStyle(true) + .httpClient(httpClient); + S3Utils.configureLegacyMd5(clientConfig, s3ClientBuilder); + + try (final S3Client s3Client = s3ClientBuilder.build()) { + s3Client.putObject( + PutObjectRequest.builder().bucket("bucket").key("key").build(), + RequestBody.fromString("payload") + ); + s3Client.deleteObjects( + DeleteObjectsRequest.builder() + .bucket("bucket") + .delete(Delete.builder().objects(ObjectIdentifier.builder().key("key").build()).build()) + .build() + ); + } + + Assertions.assertEquals(2, requests.size()); + Assertions.assertTrue(requests.get(0).httpRequest().firstMatchingHeader("Content-MD5").isEmpty()); + Assertions.assertTrue(requests.get(0).httpRequest().firstMatchingHeader("x-amz-checksum-crc32").isEmpty()); + final HttpExecuteRequest deleteObjectsRequest = requests.get(1); + Assertions.assertEquals( + Md5Utils.md5AsBase64(deleteObjectsRequest.contentStreamProvider().orElseThrow().newStream()), + deleteObjectsRequest.httpRequest().firstMatchingHeader("Content-MD5").orElseThrow() + ); + EasyMock.verify(clientConfig); + } + @Test public void testRetryWithIOExceptions() { diff --git a/website/.spelling b/website/.spelling index a74372f7dc1c..9b8f863885b7 100644 --- a/website/.spelling +++ b/website/.spelling @@ -2659,3 +2659,4 @@ nginx - ../docs/development/extensions-core/s3.md NIO +checksum