From ae97e3e31d5ff5d90ff5fabe9b9b4971e0b1ec89 Mon Sep 17 00:00:00 2001 From: Corneil du Plessis Date: Fri, 20 Oct 2023 17:05:30 +0200 Subject: [PATCH] Fix S3 Source Integration tests. --- .../test/source/s3/S3SourceTests.java | 17 ++++++++++------- .../supplier/s3/AwsS3SupplierConfiguration.java | 11 ++++++----- stream-applications-build/pom.xml | 4 ++-- 3 files changed, 18 insertions(+), 14 deletions(-) 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 86740022..21cbe05a 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 @@ -19,14 +19,15 @@ package org.springframework.cloud.stream.app.integration.test.source.s3; import java.util.Map; import java.util.function.Predicate; +import com.fasterxml.jackson.databind.ObjectMapper; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.testcontainers.containers.Network; import org.testcontainers.containers.wait.strategy.Wait; import software.amazon.awssdk.services.s3.S3Client; import software.amazon.awssdk.services.s3.model.NoSuchBucketException; @@ -57,12 +58,11 @@ abstract class S3SourceTests implements LocalstackContainerTest { @BeforeEach void configureSource() { s3Client.createBucket(r -> r.bucket("bucket")); - // Use LocalStack container network String endpoint = String.format("http://localstack:%d", LOCAL_STACK_CONTAINER.getExposedPorts().get(0)); String region = LOCAL_STACK_CONTAINER.getRegion(); logger.info("creating S3 source with region={}, endpoint={}, container={}", region, endpoint, LOCAL_STACK_CONTAINER.getEndpoint()); source = BaseContainerExtension.containerInstance() - .withNetwork(LOCAL_STACK_CONTAINER.getNetwork()) + .withNetwork(Network.SHARED) .withEnv("SPRING_CLOUD_CONFIG_ENABLED", "false") .withEnv("SPRING_CLOUD_AWS_S3_ENDPOINT", endpoint) .withEnv("SPRING_CLOUD_AWS_S3_PATH_STYLE_ACCESS_ENABLED", "true") @@ -75,7 +75,6 @@ abstract class S3SourceTests implements LocalstackContainerTest { @SuppressWarnings("unchecked") @Test - @Disabled void testLines() { startContainer(fluentStringMap().withEntry("FILE_CONSUMER_MODE", "lines")); s3Client.putObject(r -> r.bucket("bucket").key("test"), resourceAsFile("s3/data").toPath()); @@ -83,8 +82,8 @@ abstract class S3SourceTests implements LocalstackContainerTest { .until(outputMatcher.payloadMatches((String s) -> s.contains("Bart Simpson"))); } + @SuppressWarnings("unchecked") @Test - @Disabled void testTaskLaunchRequest() { startContainer(fluentStringMap().withEntry("SPRING_CLOUD_FUNCTION_DEFINITION", "s3Supplier|taskLaunchRequestFunction") .withEntry("TASK_LAUNCH_REQUEST_ARG_EXPRESSIONS", "filename=payload") @@ -99,13 +98,18 @@ abstract class S3SourceTests implements LocalstackContainerTest { .until(outputMatcher.payloadMatches(predicate)); } + @SuppressWarnings("unchecked") @Test void testListOnly() { + ObjectMapper mapper = new ObjectMapper(); startContainer(fluentStringMap() .withEntry("FILE_CONSUMER_MODE", "ref") .withEntry("S3_SUPPLIER_LIST_ONLY", "true")); s3Client.putObject(r -> r.bucket("bucket").key("test"), resourceAsFile("s3/data").toPath()); - Predicate predicate = (String s) -> s.contains("\"bucketName\":\"bucket\",\"key\":\"test\""); + Predicate predicate = (String s) -> { + logger.info("payload:{}", s); + return s.contains("\"key\":\"test\""); + }; await().atMost(DEFAULT_DURATION) .until(outputMatcher.payloadMatches(predicate)); } @@ -127,7 +131,6 @@ abstract class S3SourceTests implements LocalstackContainerTest { catch (NoSuchBucketException exception) { logger.warn("No bucket 'bucket' to remove"); } - source.stop(); outputMatcher.clearMessageMatchers(); } diff --git a/functions/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java b/functions/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java index 5ad44a6b..be44f4fd 100644 --- a/functions/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java +++ b/functions/supplier/s3-supplier/src/main/java/org/springframework/cloud/fn/supplier/s3/AwsS3SupplierConfiguration.java @@ -208,9 +208,12 @@ public class AwsS3SupplierConfiguration { } @Bean - ReactiveMessageSourceProducer s3ListingMessageProducer(S3Client amazonS3, ObjectMapper objectMapper, - AwsS3SupplierProperties awsS3SupplierProperties, Predicate filter) { - + ReactiveMessageSourceProducer s3ListingMessageProducer( + S3Client amazonS3, + ObjectMapper objectMapper, + AwsS3SupplierProperties awsS3SupplierProperties, + Predicate filter + ) { return new ReactiveMessageSourceProducer( (MessageSource>) () -> { List summaryList = @@ -232,7 +235,5 @@ public class AwsS3SupplierConfiguration { return summaryList.isEmpty() ? null : new GenericMessage<>(summaryList); }); } - } - } diff --git a/stream-applications-build/pom.xml b/stream-applications-build/pom.xml index a4d13d56..091e4719 100644 --- a/stream-applications-build/pom.xml +++ b/stream-applications-build/pom.xml @@ -36,8 +36,8 @@ true 0.0.35 - 3.1.3 - 6.0.11 + 3.1.5 + 6.0.13 3.0.7 3.0.4