From 3a2c14be82ad5adcb8a2350551bbc63a094d1566 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 1 Aug 2023 13:02:29 -0400 Subject: [PATCH] Update `S3SourceTests` to use AWS SDK v2 * Update `stream-applications-integration-tests/pom.xml` to rely on versions from the parent * Use `s3-supplier` as dependency to rely on its transitives for AWS SDK --- .../pom.xml | 13 +-- .../test/source/s3/S3SourceTests.java | 105 +++++++++--------- 2 files changed, 57 insertions(+), 61 deletions(-) diff --git a/applications/stream-applications-integration-tests/pom.xml b/applications/stream-applications-integration-tests/pom.xml index 1243a8fd..75b8d7b5 100644 --- a/applications/stream-applications-integration-tests/pom.xml +++ b/applications/stream-applications-integration-tests/pom.xml @@ -18,7 +18,6 @@ 2.6.2 2.27.1 - 1.12.261 8.0.28 @@ -31,36 +30,30 @@ org.testcontainers testcontainers - ${testcontainers.version} test org.testcontainers junit-jupiter - ${testcontainers.version} test org.testcontainers kafka - ${testcontainers.version} test org.testcontainers rabbitmq - ${testcontainers.version} org.testcontainers mongodb - ${testcontainers.version} test org.testcontainers mysql - ${testcontainers.version} test @@ -81,10 +74,8 @@ test - com.amazonaws - aws-java-sdk-s3 - ${aws.version} - test + org.springframework.cloud.fn + s3-supplier com.squareup.okhttp3 diff --git a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/source/s3/S3SourceTests.java b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/source/s3/S3SourceTests.java index 243e7888..3a43b840 100644 --- a/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/source/s3/S3SourceTests.java +++ b/applications/stream-applications-integration-tests/src/test/java/org/springframework/cloud/stream/app/integration/test/source/s3/S3SourceTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2020-2020 the original author or authors. + * Copyright 2020-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -16,19 +16,11 @@ package org.springframework.cloud.stream.app.integration.test.source.s3; +import java.net.URI; import java.time.Duration; import java.util.Map; import java.util.function.Consumer; -import com.amazonaws.ClientConfiguration; -import com.amazonaws.auth.AWSCredentials; -import com.amazonaws.auth.AWSStaticCredentialsProvider; -import com.amazonaws.auth.BasicAWSCredentials; -import com.amazonaws.client.builder.AwsClientBuilder; -import com.amazonaws.regions.Regions; -import com.amazonaws.services.s3.AmazonS3; -import com.amazonaws.services.s3.AmazonS3ClientBuilder; -import com.amazonaws.services.s3.model.PutObjectRequest; import com.github.dockerjava.api.command.CreateContainerCmd; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeAll; @@ -42,6 +34,12 @@ import org.testcontainers.containers.GenericContainer; import org.testcontainers.containers.wait.strategy.Wait; import org.testcontainers.junit.jupiter.Container; import org.testcontainers.utility.DockerImageName; +import software.amazon.awssdk.auth.credentials.AwsBasicCredentials; +import software.amazon.awssdk.auth.credentials.AwsCredentials; +import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider; +import software.amazon.awssdk.regions.Region; +import software.amazon.awssdk.services.s3.S3Client; +import software.amazon.awssdk.services.s3.model.NoSuchBucketException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.stream.app.test.integration.OutputMatcher; @@ -53,11 +51,15 @@ import static org.springframework.cloud.stream.app.integration.test.common.Confi import static org.springframework.cloud.stream.app.test.integration.AppLog.appLog; import static org.springframework.cloud.stream.app.test.integration.FluentMap.fluentMap; import static org.springframework.cloud.stream.app.test.integration.StreamAppContainerTestUtils.resourceAsFile; + @Tag("integration") @ExtendWith(BaseContainerExtension.class) abstract class S3SourceTests { private static final Logger logger = LoggerFactory.getLogger(S3SourceTests.class); - private static AmazonS3 s3Client; + + private static S3Client s3Client; + + private static String minioAddress; private StreamAppContainer source; @@ -67,31 +69,29 @@ abstract class S3SourceTests { @Container private static final GenericContainer minio = new GenericContainer( DockerImageName.parse("minio/minio:latest")) - .withExposedPorts(9000) - .withNetworkAliases("minio-host") - .withEnv("MINIO_ACCESS_KEY", "minio") - .withEnv("MINIO_SECRET_KEY", "minio123") - .waitingFor(Wait.forHttp("/minio/health/live")) - .withCreateContainerCmdModifier( - (Consumer) createContainerCmd -> createContainerCmd - .withHostName("minio")) - .withLogConsumer(appLog("minio")) - .withCommand("minio", "server", "/data") - .withStartupTimeout(Duration.ofSeconds(120)) - .withStartupAttempts(3); + .withExposedPorts(9000) + .withNetworkAliases("minio-host") + .withEnv("MINIO_ACCESS_KEY", "minio") + .withEnv("MINIO_SECRET_KEY", "minio123") + .waitingFor(Wait.forHttp("/minio/health/live")) + .withCreateContainerCmdModifier( + (Consumer) createContainerCmd -> createContainerCmd + .withHostName("minio")) + .withLogConsumer(appLog("minio")) + .withCommand("minio", "server", "/data") + .withStartupTimeout(Duration.ofSeconds(120)) + .withStartupAttempts(3); @BeforeAll static void initS3() { - AWSCredentials credentials = new BasicAWSCredentials("minio", "minio123"); - ClientConfiguration clientConfiguration = new ClientConfiguration(); - String minioAddress = "http://" + minio.getHost() + ":" + minio.getMappedPort(9000); + minioAddress = "http://" + minio.getHost() + ":" + minio.getMappedPort(9000); + AwsCredentials credentials = AwsBasicCredentials.create("minio", "minio123"); logger.info("minio:address={}", minioAddress); - s3Client = AmazonS3ClientBuilder - .standard() - .withEndpointConfiguration(new AwsClientBuilder.EndpointConfiguration(minioAddress, Regions.US_EAST_1.name())) - .withPathStyleAccessEnabled(true) - .withClientConfiguration(clientConfiguration) - .withCredentials(new AWSStaticCredentialsProvider(credentials)) + s3Client = S3Client.builder() + .endpointOverride(URI.create(minioAddress)) + .region(Region.US_EAST_1) + .forcePathStyle(true) + .credentialsProvider(StaticCredentialsProvider.create(credentials)) .build(); } @@ -99,22 +99,21 @@ abstract class S3SourceTests { @BeforeEach void configureSource() { source = BaseContainerExtension.containerInstance() - .withEnv("S3_COMMON_ENDPOINT_URL", "http://minio-host:" + minio.getMappedPort(9000)) - .withEnv("S3_COMMON_PATH_STYLE_ACCESS", "true") - .withEnv("CLOUD_AWS_STACK_AUTO", "false") - .withEnv("CLOUD_AWS_CREDENTIALS_ACCESS_KEY", "minio") - .withEnv("CLOUD_AWS_CREDENTIALS_SECRET_KEY", "minio123") + .withEnv("SPRING_CLOUD_AWS_S3_ENDPOINT", minioAddress) + .withEnv("SPRING_CLOUD_AWS_S3_PATH_STYLE_ACCESS_ENABLED", "true") + .withEnv("SPRING_CLOUD_AWS_CREDENTIALS_ACCESS_KEY", "minio") + .withEnv("SPRING_CLOUD_AWS_CREDENTIALS_SECRET_KEY", "minio123") .withEnv("LOGGING_LEVEL_ORG_SPRINGFRAMEWORK_INTEGRATION", "DEBUG") - .withEnv("CLOUD_AWS_REGION_STATIC", "us-east-1").log(); - s3Client.createBucket("bucket"); + .withEnv("SPRING_CLOUD_AWS_REGION_STATIC", Region.US_EAST_1.id()) + .log(); + + s3Client.createBucket(r -> r.bucket("bucket")); } @Test void testLines() { - startContainer( - fluentMap().withEntry("FILE_CONSUMER_MODE", "lines")); - s3Client.putObject(new PutObjectRequest("bucket", "test", - resourceAsFile("minio/data"))); + startContainer(fluentMap().withEntry("FILE_CONSUMER_MODE", "lines")); + s3Client.putObject(r -> r.bucket("bucket").key("test"), resourceAsFile("minio/data").toPath()); await().atMost(DEFAULT_DURATION).until(outputMatcher.payloadMatches((String s) -> s.contains("Bart Simpson"))); @@ -129,8 +128,8 @@ abstract class S3SourceTests { .withEntry("TASK_LAUNCH_REQUEST_TASK_NAME", "myTask") .withEntry("FILE_CONSUMER_MODE", "ref")); - s3Client.putObject(new PutObjectRequest("bucket", "test", - resourceAsFile("minio/data"))); + s3Client.putObject(r -> r.bucket("bucket").key("test"), resourceAsFile("minio/data").toPath()); + await().atMost(DEFAULT_DURATION).until(outputMatcher.payloadMatches(s -> s.equals( "{\"args\":[\"filename=/tmp/s3-supplier/test\"],\"deploymentProps\":{},\"name\":\"myTask\"}"))); } @@ -142,8 +141,8 @@ abstract class S3SourceTests { .withEntry("FILE_CONSUMER_MODE", "ref") .withEntry("S3_SUPPLIER_LIST_ONLY", "true")); - s3Client.putObject(new PutObjectRequest("bucket", "test", - resourceAsFile("minio/data"))); + s3Client.putObject(r -> r.bucket("bucket").key("test"), resourceAsFile("minio/data").toPath()); + await().atMost(DEFAULT_DURATION) .until(outputMatcher .payloadMatches((String s) -> s.contains("\"bucketName\":\"bucket\",\"key\":\"test\""))); @@ -157,10 +156,16 @@ abstract class S3SourceTests { @AfterEach void stop() { - if (s3Client.doesBucketExistV2("bucket")) { - s3Client.deleteObject("bucket", "test"); - s3Client.deleteBucket("bucket"); + try { + s3Client.headBucket(r -> r.bucket("bucket")); + // NoSuchBucketException - no deletion attempts + s3Client.deleteObject(r -> r.bucket("bucket").key("test")); + s3Client.deleteBucket(r -> r.bucket("bucket")); } + catch (NoSuchBucketException exception) { + logger.warn("No bucket 'bucket' to remove"); + } + source.stop(); outputMatcher.clearMessageMatchers(); }