Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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.
*/
Expand Down Expand Up @@ -178,6 +182,11 @@ public boolean isEnablePathStyleAccess()
return enablePathStyleAccess;
}

public boolean isEnableLegacyMd5()
{
return enableLegacyMd5;
}

/**
* @deprecated Use {@link #isCrossRegionAccessEnabled()} instead.
*/
Expand Down Expand Up @@ -272,6 +281,7 @@ public String toString()
"protocol='" + protocol + '\'' +
", disableChunkedEncoding=" + disableChunkedEncoding +
", enablePathStyleAccess=" + enablePathStyleAccess +
", enableLegacyMd5=" + enableLegacyMd5 +
", crossRegionAccessEnabled=" + isCrossRegionAccessEnabled() +
", connectionTimeout=" + connectionTimeout +
", socketTimeout=" + socketTimeout +
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
1 change: 1 addition & 0 deletions docs/development/extensions-core/s3.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)|
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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.
* <p>
* 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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,23 +27,160 @@
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;
import java.util.stream.Collectors;

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<HttpExecuteRequest> 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()
{
Expand Down
1 change: 1 addition & 0 deletions website/.spelling
Original file line number Diff line number Diff line change
Expand Up @@ -2659,3 +2659,4 @@ nginx

- ../docs/development/extensions-core/s3.md
NIO
checksum
Loading