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 @@ -37,13 +37,15 @@
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.services.s3.model.Tag;
import software.amazon.awssdk.services.s3.model.TaggingDirective;
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.
* Verifies that the CRT-based S3 client publishes CRT native request telemetry to client-level and request-level
* {@link MetricPublisher}s. 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 {
Expand Down Expand Up @@ -134,6 +136,87 @@ void failedGetObject_publishesUnsuccessfulApiCallMetrics() throws InterruptedExc
c -> assertThat(c.metricValues(CoreMetric.API_CALL_SUCCESSFUL)).contains(false));
}

@Test
void requestLevelPublisher_overridesClientLevelPublisher() throws InterruptedException {
CapturingMetricPublisher clientPublisher = new CapturingMetricPublisher();
CapturingMetricPublisher requestPublisher = new CapturingMetricPublisher();

try (S3AsyncClient client = crtClientWith(clientPublisher)) {
client.getObject(b -> b.bucket(BUCKET).key(SMALL_KEY)
.overrideConfiguration(o -> o.addMetricPublisher(requestPublisher)),
AsyncResponseTransformer.toBytes()).join();

// The request-level publisher receives the CRT telemetry ...
List<MetricCollection> requestCollections = requestPublisher.awaitAtLeast(1, Duration.ofSeconds(30));
assertThat(requestCollections).isNotEmpty();
requestCollections.forEach(S3CrtClientMetricPublisherIntegrationTest::assertIsCrtApiCallCollection);

// ... and this client's own client-level publisher receives nothing, since request-level takes precedence.
assertThat(clientPublisher.awaitAtLeast(1, Duration.ofSeconds(1))).isEmpty();
}
}

@Test
void copyObjectWithRequestLevelPublisher_publishesSubRequestTelemetryToIt() throws InterruptedException {
CapturingMetricPublisher clientPublisher = new CapturingMetricPublisher();
CapturingMetricPublisher requestPublisher = new CapturingMetricPublisher();
String destinationKey = "copy-dest-" + System.nanoTime();

try (S3AsyncClient client = crtClientWith(clientPublisher)) {
client.copyObject(b -> b.sourceBucket(BUCKET).sourceKey(LARGE_KEY)
.destinationBucket(BUCKET).destinationKey(destinationKey)
.overrideConfiguration(o -> o.addMetricPublisher(requestPublisher)))
.join();

List<MetricCollection> collections = requestPublisher.awaitAtLeast(2, Duration.ofSeconds(60));
assertThat(collections).hasSizeGreaterThanOrEqualTo(2);
collections.forEach(S3CrtClientMetricPublisherIntegrationTest::assertIsCrtApiCallCollection);

// request-level takes precedence, so this client's own client-level publisher sees none of the sub-requests.
assertThat(clientPublisher.awaitAtLeast(1, Duration.ofSeconds(1))).isEmpty();
}
}

@Test

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we add assertThat(clientPublisher.awaitAtLeast(1, Duration.ofSeconds(1))).isEmpty();?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch, added.

void copyObjectWithTaggingDirective_requestLevelPublisher_receivesTaggingSubRequestMetrics() throws InterruptedException {
// Tag the source so a taggingDirective(COPY) copy actually issues GetObjectTagging + PutObjectTagging sub-requests.
S3IntegrationTestBase.s3.putObjectTagging(
r -> r.bucket(BUCKET).key(LARGE_KEY)
.tagging(t -> t.tagSet(Tag.builder().key("env").value("t").build())));

CapturingMetricPublisher clientPublisher = new CapturingMetricPublisher();
CapturingMetricPublisher requestPublisher = new CapturingMetricPublisher();
String destinationKey = "copy-tagged-dest-" + System.nanoTime();

try (S3AsyncClient client = crtClientWith(clientPublisher)) {
client.copyObject(b -> b.sourceBucket(BUCKET).sourceKey(LARGE_KEY)
.destinationBucket(BUCKET).destinationKey(destinationKey)
.taggingDirective(TaggingDirective.COPY)
.overrideConfiguration(o -> o.addMetricPublisher(requestPublisher)))
.join();

List<MetricCollection> collections = requestPublisher.awaitAtLeast(2, Duration.ofSeconds(60));
collections.forEach(S3CrtClientMetricPublisherIntegrationTest::assertIsCrtApiCallCollection);
assertThat(collections).anySatisfy(
c -> assertThat(c.metricValues(CoreMetric.OPERATION_NAME)).contains("GetObjectTagging"));
assertThat(collections).anySatisfy(
c -> assertThat(c.metricValues(CoreMetric.OPERATION_NAME)).contains("PutObjectTagging"));

// request-level takes precedence, so this client's own client-level publisher sees none of the sub-requests.
assertThat(clientPublisher.awaitAtLeast(1, Duration.ofSeconds(1))).isEmpty();
}
}

private static S3AsyncClient crtClientWith(MetricPublisher clientLevelPublisher) {
return S3AsyncClient.crtBuilder()
.region(S3IntegrationTestBase.DEFAULT_REGION)
.credentialsProvider(AwsTestBase.CREDENTIALS_PROVIDER_CHAIN)
.minimumPartSizeInBytes(PART_SIZE)
.thresholdInBytes(PART_SIZE)
.addMetricPublisher(clientLevelPublisher)
.build();
}

private static void assertIsCrtApiCallCollection(MetricCollection apiCall) {
assertThat(apiCall.name()).isEqualTo("ApiCall");
assertThat(apiCall.metricValues(CoreMetric.SERVICE_ID)).containsExactly("S3");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,10 +31,12 @@
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;
import java.util.function.Function;
import software.amazon.awssdk.annotations.SdkInternalApi;
import software.amazon.awssdk.annotations.SdkTestInternalApi;
import software.amazon.awssdk.auth.credentials.AwsCredentialsProvider;
Expand Down Expand Up @@ -82,6 +84,7 @@
import software.amazon.awssdk.services.s3.model.GetObjectResponse;
import software.amazon.awssdk.services.s3.model.PutObjectRequest;
import software.amazon.awssdk.services.s3.model.PutObjectResponse;
import software.amazon.awssdk.services.s3.model.S3Request;
import software.amazon.awssdk.services.s3.presignedurl.AsyncPresignedUrlExtension;
import software.amazon.awssdk.utils.AttributeMap;
import software.amazon.awssdk.utils.CollectionUtils;
Expand All @@ -93,6 +96,8 @@ public final class DefaultS3CrtAsyncClient extends DelegatingS3AsyncClient imple
public static final ExecutionAttribute<Path> RESPONSE_FILE_PATH = new ExecutionAttribute<>("responseFilePath");
public static final ExecutionAttribute<S3MetaRequestOptions.ResponseFileOption> RESPONSE_FILE_OPTION =
new ExecutionAttribute<>("responseFileOption");
public static final ExecutionAttribute<List<MetricPublisher>> REQUEST_METRIC_PUBLISHERS =
new ExecutionAttribute<>("requestMetricPublishers");
private static final String CRT_CLIENT_CLASSPATH = "software.amazon.awssdk.crt.s3.S3Client";
private final CopyObjectHelper copyObjectHelper;

Expand Down Expand Up @@ -137,7 +142,35 @@ public CompletableFuture<GetObjectResponse> getObject(GetObjectRequest getObject

@Override
public CompletableFuture<CopyObjectResponse> copyObject(CopyObjectRequest copyObjectRequest) {
return copyObjectHelper.copyObject(copyObjectRequest);
// copyObject's sub-requests bypass invokeOperation, so stash once here; CopyObjectHelper propagates the copy's
// override (with the stashed attribute) to each sub-request.
return copyObjectHelper.copyObject(stashRequestMetricPublishers(copyObjectRequest));
}

/**
* All operations funnel through here. If the request carries request-level metric publishers, move them off the
* request override (so the inner standard client's resolveMetricPublishers stays empty -> NoOp, no hollow ApiCall)
* and stash them in an execution attribute that the CRT transport reads to publish telemetry to them instead of
* (overriding) the client-level publishers.
*/
@Override
protected <T extends S3Request, ReturnT> CompletableFuture<ReturnT> invokeOperation(
T request, Function<T, CompletableFuture<ReturnT>> operation) {
return operation.apply(stashRequestMetricPublishers(request));
}

@SuppressWarnings("unchecked")
static <T extends S3Request> T stashRequestMetricPublishers(T request) {
AwsRequestOverrideConfiguration override = request.overrideConfiguration().orElse(null);
if (override == null || CollectionUtils.isNullOrEmpty(override.metricPublishers())) {
return request;
}
AwsRequestOverrideConfiguration newOverride =
override.toBuilder()
.metricPublishers(Collections.emptyList())
.putExecutionAttribute(REQUEST_METRIC_PUBLISHERS, override.metricPublishers())
.build();
return (T) request.toBuilder().overrideConfiguration(newOverride).build();
}

private static S3AsyncClient initializeS3AsyncClient(DefaultS3CrtClientBuilder builder) {
Expand Down Expand Up @@ -460,7 +493,9 @@ public void afterMarshalling(Context.AfterMarshalling context,
.put(S3InternalSdkHttpExecutionAttribute.RESPONSE_FILE_PATH,
executionAttributes.getAttribute(RESPONSE_FILE_PATH))
.put(S3InternalSdkHttpExecutionAttribute.RESPONSE_FILE_OPTION,
executionAttributes.getAttribute(RESPONSE_FILE_OPTION));
executionAttributes.getAttribute(RESPONSE_FILE_OPTION))
.put(S3InternalSdkHttpExecutionAttribute.METRIC_PUBLISHERS,
executionAttributes.getAttribute(REQUEST_METRIC_PUBLISHERS));

SdkRequest request = context.request();
if (request instanceof AwsRequest) {
Expand Down Expand Up @@ -517,10 +552,6 @@ private static void validateOverrideConfiguration(SdkRequest request) {
throw new UnsupportedOperationException("Request-level signer override is not supported");
}

if (!CollectionUtils.isNullOrEmpty(overrideConfiguration.metricPublishers())) {
throw new UnsupportedOperationException("Request-level Metric Publishers override is not supported");
}

if (overrideConfiguration.apiCallAttemptTimeout().isPresent()) {
throw new UnsupportedOperationException("Request-level apiCallAttemptTimeout override is not supported");
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -166,14 +166,15 @@ 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.
List<MetricPublisher> requestMetricPublishers =
httpExecutionAttributes.getAttribute(S3InternalSdkHttpExecutionAttribute.METRIC_PUBLISHERS);
List<MetricPublisher> effectivePublishers = resolveEffectiveMetricPublishers(requestMetricPublishers,
metricPublishers);
SdkHttpExecutionAttributes adapterAttributes = httpExecutionAttributes;
if (!metricPublishers.isEmpty()) {
if (!effectivePublishers.isEmpty()) {
adapterAttributes = httpExecutionAttributes.toBuilder()
.put(S3InternalSdkHttpExecutionAttribute.METRIC_PUBLISHERS,
metricPublishers)
effectivePublishers)
.build();
}

Expand Down Expand Up @@ -232,6 +233,18 @@ public CompletableFuture<Void> execute(AsyncExecuteRequest asyncRequest) {
return executeFuture;
}

/**
* Request-level publishers (if any were set on the request override) take precedence over the client-level
* publishers; otherwise the client-level publishers are used.
*/
static List<MetricPublisher> resolveEffectiveMetricPublishers(List<MetricPublisher> requestLevel,
List<MetricPublisher> clientLevel) {
if (requestLevel != null && !requestLevel.isEmpty()) {
return requestLevel;
}
return clientLevel;
}

private AwsSigningConfig awsSigningConfig(Region signingRegion, SdkHttpExecutionAttributes httpExecutionAttributes) {
CrtCredentialsProviderAdapter requestAdapter =
httpExecutionAttributes.getAttribute(S3InternalSdkHttpExecutionAttribute.CRT_CREDENTIALS_PROVIDER_ADAPTER);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,11 +70,12 @@ public final class S3InternalSdkHttpExecutionAttribute<T> extends SdkHttpExecuti
new S3InternalSdkHttpExecutionAttribute<>(CrtCredentialsProviderAdapter.class);

/**
* Metric publishers that this request's CRT telemetry is published to.
* The metric publishers this request's CRT telemetry is published to: the request-level publishers if the request
* set any, otherwise the client-level publishers. The CRT transport resolves the two and folds the effective set in.
*/
@SuppressWarnings("unchecked")
public static final S3InternalSdkHttpExecutionAttribute<List<MetricPublisher>> METRIC_PUBLISHERS =
new S3InternalSdkHttpExecutionAttribute<>((Class<List<MetricPublisher>>) (Class<?>) List.class);
new S3InternalSdkHttpExecutionAttribute<>((Class<List<MetricPublisher>>) (Class<?>) List.class);

private S3InternalSdkHttpExecutionAttribute(Class<T> valueClass) {
super(valueClass);
Expand Down
Loading
Loading