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
@@ -0,0 +1,180 @@
/*
* Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
*
* Licensed under the Apache License, Version 2.0 (the "License").
* You may not use this file except in compliance with the License.
* A copy of the License is located at
*
* http://aws.amazon.com/apache2.0
*
* or in the "license" file accompanying this file. This file is distributed
* on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either
* express or implied. See the License for the specific language governing
* permissions and limitations under the License.
*/

package software.amazon.awssdk.services.s3.crt;

import static org.assertj.core.api.Assertions.assertThat;
import static software.amazon.awssdk.testutils.service.S3BucketUtils.temporaryBucketName;

import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.TimeUnit;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import software.amazon.awssdk.core.async.AsyncRequestBody;
import software.amazon.awssdk.core.async.AsyncResponseTransformer;
import software.amazon.awssdk.core.metrics.CoreMetric;
import software.amazon.awssdk.http.HttpMetric;
import software.amazon.awssdk.metrics.MetricCollection;
import software.amazon.awssdk.metrics.MetricPublisher;
import software.amazon.awssdk.services.s3.S3AsyncClient;
import software.amazon.awssdk.services.s3.S3IntegrationTestBase;
import software.amazon.awssdk.services.s3.model.PutObjectRequest;
import software.amazon.awssdk.testutils.RandomTempFile;
import software.amazon.awssdk.testutils.service.AwsTestBase;

/**
* Verifies that the CRT-based S3 client publishes CRT native request telemetry to a client-level {@link MetricPublisher}
* configured via {@code crtBuilder().addMetricPublisher(...)}. Each underlying CRT request attempt is published as its
* own {@code ApiCall -> ApiCallAttempt -> HttpClient} {@link MetricCollection}, so a multipart transfer yields several.
*/
@Timeout(value = 5, unit = TimeUnit.MINUTES)
public class S3CrtClientMetricPublisherIntegrationTest extends S3IntegrationTestBase {
private static final String BUCKET = temporaryBucketName(S3CrtClientMetricPublisherIntegrationTest.class);
private static final String SMALL_KEY = "small-single-part";
private static final String LARGE_KEY = "large-multipart";
private static final long PART_SIZE = 8L * 1024 * 1024;
// Below the 8MB threshold -> single underlying GET.
private static final int SMALL_SIZE = 100 * 1024;
// Comfortably above the threshold -> the meta-request fans out into multiple underlying part requests.
private static final int LARGE_SIZE = 17 * 1024 * 1024;

private static final CapturingMetricPublisher CLIENT_PUBLISHER = new CapturingMetricPublisher();
private static S3AsyncClient crtClient;

@BeforeAll
public static void setup() throws Exception {
S3IntegrationTestBase.setUp();
S3IntegrationTestBase.createBucket(BUCKET);

crtClient = S3AsyncClient.crtBuilder()
.region(S3IntegrationTestBase.DEFAULT_REGION)
.credentialsProvider(AwsTestBase.CREDENTIALS_PROVIDER_CHAIN)
.minimumPartSizeInBytes(PART_SIZE)
.thresholdInBytes(PART_SIZE)
.addMetricPublisher(CLIENT_PUBLISHER)
.build();

S3IntegrationTestBase.s3.putObject(PutObjectRequest.builder().bucket(BUCKET).key(SMALL_KEY).build(),
new RandomTempFile(SMALL_SIZE).toPath());
S3IntegrationTestBase.s3.putObject(PutObjectRequest.builder().bucket(BUCKET).key(LARGE_KEY).build(),
new RandomTempFile(LARGE_SIZE).toPath());
}

@AfterAll
public static void cleanup() {
crtClient.close();
S3IntegrationTestBase.deleteBucketAndAllContents(BUCKET);
}

@BeforeEach
public void resetPublisher() {
CLIENT_PUBLISHER.clear();
}

@Test
void singlePartGetObject_publishesApiCallMetrics() throws InterruptedException {
crtClient.getObject(b -> b.bucket(BUCKET).key(SMALL_KEY), AsyncResponseTransformer.toBytes()).join();

List<MetricCollection> collections = CLIENT_PUBLISHER.awaitAtLeast(1, Duration.ofSeconds(30));
assertThat(collections).isNotEmpty();
collections.forEach(S3CrtClientMetricPublisherIntegrationTest::assertIsCrtApiCallCollection);
}

@Test
void multipartGetObject_publishesOneCollectionPerUnderlyingRequest() throws InterruptedException {
crtClient.getObject(b -> b.bucket(BUCKET).key(LARGE_KEY), AsyncResponseTransformer.toBytes()).join();

// A multipart download fans out into several underlying requests, so we expect more than one collection.
List<MetricCollection> collections = CLIENT_PUBLISHER.awaitAtLeast(2, Duration.ofSeconds(60));
assertThat(collections).hasSizeGreaterThanOrEqualTo(2);
collections.forEach(S3CrtClientMetricPublisherIntegrationTest::assertIsCrtApiCallCollection);
}

@Test
void putObject_publishesApiCallMetrics() throws InterruptedException {
crtClient.putObject(b -> b.bucket(BUCKET).key("put-metrics-key"),
AsyncRequestBody.fromBytes(new byte[SMALL_SIZE])).join();

List<MetricCollection> collections = CLIENT_PUBLISHER.awaitAtLeast(1, Duration.ofSeconds(30));
assertThat(collections).isNotEmpty();
collections.forEach(S3CrtClientMetricPublisherIntegrationTest::assertIsCrtApiCallCollection);
}

@Test
void failedGetObject_publishesUnsuccessfulApiCallMetrics() throws InterruptedException {
try {
crtClient.getObject(b -> b.bucket(BUCKET).key("does-not-exist-" + System.nanoTime()),
AsyncResponseTransformer.toBytes()).join();
} catch (RuntimeException expected) {
// 404 - the metrics for the failed attempt should still be published below.
}

List<MetricCollection> collections = CLIENT_PUBLISHER.awaitAtLeast(1, Duration.ofSeconds(30));
assertThat(collections).isNotEmpty();
collections.forEach(S3CrtClientMetricPublisherIntegrationTest::assertIsCrtApiCallCollection);
assertThat(collections).anySatisfy(
c -> assertThat(c.metricValues(CoreMetric.API_CALL_SUCCESSFUL)).contains(false));
}

private static void assertIsCrtApiCallCollection(MetricCollection apiCall) {
assertThat(apiCall.name()).isEqualTo("ApiCall");
assertThat(apiCall.metricValues(CoreMetric.SERVICE_ID)).containsExactly("S3");
assertThat(apiCall.metricValues(CoreMetric.OPERATION_NAME)).isNotEmpty();
assertThat(apiCall.metricValues(CoreMetric.API_CALL_DURATION)).isNotEmpty();

MetricCollection attempt = apiCall.childrenWithName("ApiCallAttempt").findFirst().orElse(null);
assertThat(attempt).as("ApiCallAttempt child").isNotNull();

MetricCollection httpClient = attempt.childrenWithName("HttpClient").findFirst().orElse(null);
assertThat(httpClient).as("HttpClient child").isNotNull();
assertThat(httpClient.metricValues(HttpMetric.HTTP_CLIENT_NAME)).containsExactly("s3crt");
}

/**
* Thread-safe capture of published collections. {@code onTelemetry} fires on CRT native threads and can lag the
* completion of the operation future, so tests poll {@link #awaitAtLeast(int, Duration)} rather than reading
* immediately.
*/
private static final class CapturingMetricPublisher implements MetricPublisher {
private final List<MetricCollection> collections = new CopyOnWriteArrayList<>();

@Override
public void publish(MetricCollection metricCollection) {
collections.add(metricCollection);
}

@Override
public void close() {
}

void clear() {
collections.clear();
}

List<MetricCollection> awaitAtLeast(int min, Duration timeout) throws InterruptedException {
long deadlineNanos = System.nanoTime() + timeout.toNanos();
while (collections.size() < min && System.nanoTime() < deadlineNanos) {
Thread.sleep(100);
}
return new ArrayList<>(collections);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

import java.net.URI;
import java.nio.file.Path;
import java.util.Collection;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;
Expand All @@ -29,6 +30,7 @@
import software.amazon.awssdk.core.client.config.SdkAdvancedAsyncClientOption;
import software.amazon.awssdk.identity.spi.AwsCredentialsIdentity;
import software.amazon.awssdk.identity.spi.IdentityProvider;
import software.amazon.awssdk.metrics.MetricPublisher;
import software.amazon.awssdk.regions.Region;
import software.amazon.awssdk.services.s3.crt.S3CrtHttpConfiguration;
import software.amazon.awssdk.services.s3.crt.S3CrtRetryConfiguration;
Expand Down Expand Up @@ -379,6 +381,33 @@ default S3CrtAsyncClientBuilder retryConfiguration(Consumer<S3CrtRetryConfigurat
*/
S3CrtAsyncClientBuilder advancedOptions(Map<SdkAdvancedAsyncClientOption<?>, ?> advancedOptions);

/**
* Sets the {@link MetricPublisher}s that this client publishes metrics to. This overrides any previously configured
* publishers.
*
* <p>Each underlying CRT request attempt is published as a separate {@code "ApiCall"} metric collection. A single
* high-level transfer (for example a multipart {@code getObject}) fans out into several underlying requests, so
* several collections are published per transfer.
*
* <p>The SDK does not close these publishers when the client is closed; the caller retains ownership and is
* responsible for closing them.
*
* @param metricPublishers the metric publishers to use
* @return this builder for method chaining.
*/
S3CrtAsyncClientBuilder metricPublishers(Collection<MetricPublisher> metricPublishers);

/**
* Adds a {@link MetricPublisher} that this client publishes metrics to. Can be called multiple times to add several
* publishers.
*
* <p>See {@link #metricPublishers(Collection)} for details on how metrics are published and on publisher ownership.
*
* @param metricPublisher the metric publisher to add
* @return this builder for method chaining.
*/
S3CrtAsyncClientBuilder addMetricPublisher(MetricPublisher metricPublisher);

@Override
S3AsyncClient build();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import java.net.URI;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
Expand Down Expand Up @@ -64,6 +65,7 @@
import software.amazon.awssdk.http.SdkHttpExecutionAttributes;
import software.amazon.awssdk.identity.spi.AwsCredentialsIdentity;
import software.amazon.awssdk.identity.spi.IdentityProvider;
import software.amazon.awssdk.metrics.MetricPublisher;
import software.amazon.awssdk.regions.Region;
import software.amazon.awssdk.services.s3.DelegatingS3AsyncClient;
import software.amazon.awssdk.services.s3.S3AsyncClient;
Expand Down Expand Up @@ -236,7 +238,8 @@ private static S3CrtAsyncHttpClient.Builder initializeS3CrtAsyncHttpClient(Defau
.withMaxRetries(builder.retryConfiguration.numRetries())));
}
return S3CrtAsyncHttpClient.builder()
.s3ClientConfiguration(nativeClientBuilder.build());
.s3ClientConfiguration(nativeClientBuilder.build())
.metricPublishers(builder.metricPublishers);
}

public static final class DefaultS3CrtClientBuilder implements S3CrtAsyncClientBuilder {
Expand All @@ -261,6 +264,7 @@ public static final class DefaultS3CrtClientBuilder implements S3CrtAsyncClientB
private Long thresholdInBytes;
private Executor futureCompletionExecutor;
private Boolean disableS3ExpressSessionAuth;
private List<MetricPublisher> metricPublishers;

private AttributeMap.Builder advancedOptions = AttributeMap.builder();

Expand Down Expand Up @@ -406,6 +410,23 @@ public DefaultS3CrtClientBuilder advancedOptions(Map<SdkAdvancedAsyncClientOptio
return this;
}

@Override
public DefaultS3CrtClientBuilder metricPublishers(Collection<MetricPublisher> metricPublishers) {
Validate.paramNotNull(metricPublishers, "metricPublishers");
this.metricPublishers = new ArrayList<>(metricPublishers);
return this;
}

@Override
public DefaultS3CrtClientBuilder addMetricPublisher(MetricPublisher metricPublisher) {
Validate.paramNotNull(metricPublisher, "metricPublisher");
if (this.metricPublishers == null) {
this.metricPublishers = new ArrayList<>();
}
this.metricPublishers.add(metricPublisher);
return this;
}

@Override
public S3CrtAsyncClient build() {
return new DefaultS3CrtAsyncClient(this);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@

package software.amazon.awssdk.services.s3.internal.crt;

import static software.amazon.awssdk.services.s3.crt.S3CrtSdkHttpExecutionAttribute.CRT_PROGRESS_LISTENER;
import static software.amazon.awssdk.services.s3.crt.S3CrtSdkHttpExecutionAttribute.METAREQUEST_PAUSE_OBSERVABLE;
import static software.amazon.awssdk.services.s3.internal.crt.CrtChecksumUtils.checksumConfig;
import static software.amazon.awssdk.services.s3.internal.crt.S3InternalSdkHttpExecutionAttribute.CRT_PAUSE_RESUME_TOKEN;
Expand All @@ -35,6 +34,7 @@
import java.nio.file.Path;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
Expand All @@ -59,6 +59,7 @@
import software.amazon.awssdk.http.SdkHttpRequest;
import software.amazon.awssdk.http.async.AsyncExecuteRequest;
import software.amazon.awssdk.http.async.SdkAsyncHttpClient;
import software.amazon.awssdk.metrics.MetricPublisher;
import software.amazon.awssdk.regions.Region;
import software.amazon.awssdk.utils.AttributeMap;
import software.amazon.awssdk.utils.NumericUtils;
Expand All @@ -71,14 +72,18 @@
@SdkInternalApi
public final class S3CrtAsyncHttpClient implements SdkAsyncHttpClient {

static final String CLIENT_NAME = "s3crt";

private final S3Client crtS3Client;

private final S3NativeClientConfiguration s3NativeClientConfiguration;
private final S3ClientOptions s3ClientOptions;
private final List<MetricPublisher> metricPublishers;

private S3CrtAsyncHttpClient(Builder builder) {
s3NativeClientConfiguration = builder.clientConfiguration;
this.s3ClientOptions = createS3ClientOption();
this.metricPublishers = builder.metricPublishers == null ? Collections.emptyList() : builder.metricPublishers;

this.crtS3Client = new S3Client(s3ClientOptions);
}
Expand All @@ -88,6 +93,7 @@ private S3CrtAsyncHttpClient(Builder builder) {
Builder builder) {
s3NativeClientConfiguration = builder.clientConfiguration;
s3ClientOptions = createS3ClientOption();
this.metricPublishers = builder.metricPublishers == null ? Collections.emptyList() : builder.metricPublishers;
this.crtS3Client = crtS3Client;
}

Expand Down Expand Up @@ -160,12 +166,22 @@ public CompletableFuture<Void> execute(AsyncExecuteRequest asyncRequest) {
Path responseFilePath = httpExecutionAttributes.getAttribute(RESPONSE_FILE_PATH);
S3MetaRequestOptions.ResponseFileOption responseFileOption = httpExecutionAttributes.getAttribute(RESPONSE_FILE_OPTION);

// The adapter reads its inputs from the execution attributes, so attach the client-level publishers here when
// there are any. toBuilder() preserves everything already in the bag (including CRT_PROGRESS_LISTENER); skip the
// copy entirely when no publishers are configured, to avoid rebuilding the bag on every request.
SdkHttpExecutionAttributes adapterAttributes = httpExecutionAttributes;
if (!metricPublishers.isEmpty()) {
adapterAttributes = httpExecutionAttributes.toBuilder()
.put(S3InternalSdkHttpExecutionAttribute.METRIC_PUBLISHERS,
metricPublishers)
.build();
}

S3CrtResponseHandlerAdapter responseHandler =
new S3CrtResponseHandlerAdapter(
executeFuture,
asyncRequest.responseHandler(),
httpExecutionAttributes.getAttribute(CRT_PROGRESS_LISTENER),
s3MetaRequestFuture);
new S3CrtResponseHandlerAdapter(executeFuture,
asyncRequest.responseHandler(),
adapterAttributes,
s3MetaRequestFuture);

URI endpoint = getEndpoint(uri);

Expand Down Expand Up @@ -244,7 +260,7 @@ private static URI getEndpoint(URI uri) {

@Override
public String clientName() {
return "s3crt";
return CLIENT_NAME;
}

private static S3MetaRequestOptions.MetaRequestType requestType(String operationName) {
Expand Down Expand Up @@ -297,12 +313,18 @@ public static Builder builder() {

public static final class Builder implements SdkAsyncHttpClient.Builder<S3CrtAsyncHttpClient.Builder> {
private S3NativeClientConfiguration clientConfiguration;
private List<MetricPublisher> metricPublishers;

public Builder s3ClientConfiguration(S3NativeClientConfiguration clientConfiguration) {
this.clientConfiguration = clientConfiguration;
return this;
}

public Builder metricPublishers(List<MetricPublisher> metricPublishers) {
this.metricPublishers = metricPublishers;
return this;
}

@Override
public SdkAsyncHttpClient build() {
return new S3CrtAsyncHttpClient(this);
Expand Down
Loading
Loading