From 197da9fdb2ca13018fd53f4663e2da27b9154b67 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 29 Aug 2019 18:40:09 -0400 Subject: [PATCH] Add S3 remote file info into headers Related to https://github.com/spring-projects/spring-integration/issues/3043 --- .../inbound/S3InboundFileSynchronizer.java | 5 +++ .../integration/aws/support/S3Session.java | 39 ++++++++++--------- .../inbound/S3InboundChannelAdapterTests.java | 8 ++++ .../S3StreamingChannelAdapterTests.java | 6 ++- 4 files changed, 39 insertions(+), 19 deletions(-) diff --git a/src/main/java/org/springframework/integration/aws/inbound/S3InboundFileSynchronizer.java b/src/main/java/org/springframework/integration/aws/inbound/S3InboundFileSynchronizer.java index 2391721..dbce24a 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/S3InboundFileSynchronizer.java +++ b/src/main/java/org/springframework/integration/aws/inbound/S3InboundFileSynchronizer.java @@ -80,4 +80,9 @@ public class S3InboundFileSynchronizer extends AbstractInboundFileSynchronizer { } @Override - public S3ObjectSummary[] list(String path) throws IOException { + public S3ObjectSummary[] list(String path) { String[] bucketPrefix = splitPathToBucketAndKey(path, false); ListObjectsRequest listObjectsRequest = new ListObjectsRequest().withBucketName(bucketPrefix[0]); @@ -86,7 +87,7 @@ public class S3Session implements Session { } while (objectListing.isTruncated()); - return objectSummaries.toArray(new S3ObjectSummary[objectSummaries.size()]); + return objectSummaries.toArray(new S3ObjectSummary[0]); } private String resolveBucket(String bucket) { @@ -99,7 +100,7 @@ public class S3Session implements Session { } @Override - public String[] listNames(String path) throws IOException { + public String[] listNames(String path) { String[] bucketPrefix = splitPathToBucketAndKey(path, false); ListObjectsRequest listObjectsRequest = new ListObjectsRequest().withBucketName(bucketPrefix[0]); @@ -123,18 +124,18 @@ public class S3Session implements Session { } while (objectListing.isTruncated()); - return names.toArray(new String[names.size()]); + return names.toArray(new String[0]); } @Override - public boolean remove(String path) throws IOException { + public boolean remove(String path) { String[] bucketKey = splitPathToBucketAndKey(path, true); this.amazonS3.deleteObject(bucketKey[0], bucketKey[1]); return true; } @Override - public void rename(String pathFrom, String pathTo) throws IOException { + public void rename(String pathFrom, String pathTo) { String[] bucketKeyFrom = splitPathToBucketAndKey(pathFrom, true); String[] bucketKeyTo = splitPathToBucketAndKey(pathTo, true); CopyObjectRequest copyRequest = new CopyObjectRequest(bucketKeyFrom[0], bucketKeyFrom[1], bucketKeyTo[0], @@ -149,41 +150,37 @@ public class S3Session implements Session { public void read(String source, OutputStream outputStream) throws IOException { String[] bucketKey = splitPathToBucketAndKey(source, true); S3Object s3Object = this.amazonS3.getObject(bucketKey[0], bucketKey[1]); - S3ObjectInputStream objectContent = s3Object.getObjectContent(); - try { + try (S3ObjectInputStream objectContent = s3Object.getObjectContent()) { StreamUtils.copy(objectContent, outputStream); } - finally { - objectContent.close(); - } } @Override - public void write(InputStream inputStream, String destination) throws IOException { + public void write(InputStream inputStream, String destination) { Assert.notNull(inputStream, "'inputStream' must not be null."); String[] bucketKey = splitPathToBucketAndKey(destination, true); this.amazonS3.putObject(bucketKey[0], bucketKey[1], inputStream, new ObjectMetadata()); } @Override - public void append(InputStream inputStream, String destination) throws IOException { + public void append(InputStream inputStream, String destination) { throw new UnsupportedOperationException("The 'append' operation isn't supported by the Amazon S3 protocol."); } @Override - public boolean mkdir(String directory) throws IOException { + public boolean mkdir(String directory) { this.amazonS3.createBucket(directory); return true; } @Override - public boolean rmdir(String directory) throws IOException { + public boolean rmdir(String directory) { this.amazonS3.deleteBucket(resolveBucket(directory)); return true; } @Override - public boolean exists(String path) throws IOException { + public boolean exists(String path) { String[] bucketKey = splitPathToBucketAndKey(path, true); try { this.amazonS3.getObjectMetadata(bucketKey[0], bucketKey[1]); @@ -200,7 +197,7 @@ public class S3Session implements Session { } @Override - public InputStream readRaw(String source) throws IOException { + public InputStream readRaw(String source) { String[] bucketKey = splitPathToBucketAndKey(source, true); S3Object s3Object = this.amazonS3.getObject(bucketKey[0], bucketKey[1]); return s3Object.getObjectContent(); @@ -217,7 +214,7 @@ public class S3Session implements Session { } @Override - public boolean finalizeRaw() throws IOException { + public boolean finalizeRaw() { return true; } @@ -226,6 +223,12 @@ public class S3Session implements Session { return this.amazonS3; } + @Override + public String getHostPort() { + Region region = this.amazonS3.getRegion().toAWSRegion(); + return String.format("%s.%s.%s:%d", AmazonS3.ENDPOINT_PREFIX, region.getName(), region.getDomain(), 443); + } + public String normalizeBucketName(String path) { return splitPathToBucketAndKey(path, false)[0]; } diff --git a/src/test/java/org/springframework/integration/aws/inbound/S3InboundChannelAdapterTests.java b/src/test/java/org/springframework/integration/aws/inbound/S3InboundChannelAdapterTests.java index 2409de9..a953e87 100644 --- a/src/test/java/org/springframework/integration/aws/inbound/S3InboundChannelAdapterTests.java +++ b/src/test/java/org/springframework/integration/aws/inbound/S3InboundChannelAdapterTests.java @@ -19,6 +19,7 @@ package org.springframework.integration.aws.inbound; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; import static org.mockito.BDDMockito.willAnswer; +import static org.mockito.BDDMockito.willReturn; import java.io.File; import java.io.FileInputStream; @@ -46,6 +47,7 @@ import org.springframework.integration.annotation.Poller; import org.springframework.integration.aws.support.filters.S3RegexPatternFileListFilter; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.config.EnableIntegration; +import org.springframework.integration.file.FileHeaders; import org.springframework.integration.file.filters.AcceptOnceFileListFilter; import org.springframework.messaging.Message; import org.springframework.messaging.PollableChannel; @@ -57,6 +59,7 @@ import org.springframework.util.FileCopyUtils; import com.amazonaws.services.s3.AmazonS3; import com.amazonaws.services.s3.model.ListObjectsRequest; import com.amazonaws.services.s3.model.ObjectListing; +import com.amazonaws.services.s3.model.Region; import com.amazonaws.services.s3.model.S3Object; import com.amazonaws.services.s3.model.S3ObjectSummary; @@ -126,6 +129,9 @@ public class S3InboundChannelAdapterTests { assertThat(localFile.lastModified()).isGreaterThan(System.currentTimeMillis()); + assertThat(message.getHeaders()) + .containsKeys(FileHeaders.REMOTE_DIRECTORY, FileHeaders.REMOTE_HOST_PORT, FileHeaders.REMOTE_FILE); + assertThat(this.s3FilesChannel.receive(10)).isNull(); File file = new File(LOCAL_FOLDER, "A.TEST.a"); @@ -168,6 +174,8 @@ public class S3InboundChannelAdapterTests { willAnswer(invocation -> s3Object).given(amazonS3).getObject(S3_BUCKET, s3Object.getKey()); } + willReturn(Region.US_West).given(amazonS3).getRegion(); + return amazonS3; } diff --git a/src/test/java/org/springframework/integration/aws/inbound/S3StreamingChannelAdapterTests.java b/src/test/java/org/springframework/integration/aws/inbound/S3StreamingChannelAdapterTests.java index 5a575fc..ab7817c 100644 --- a/src/test/java/org/springframework/integration/aws/inbound/S3StreamingChannelAdapterTests.java +++ b/src/test/java/org/springframework/integration/aws/inbound/S3StreamingChannelAdapterTests.java @@ -19,6 +19,7 @@ package org.springframework.integration.aws.inbound; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; import static org.mockito.BDDMockito.willAnswer; +import static org.mockito.BDDMockito.willReturn; import java.io.File; import java.io.FileInputStream; @@ -60,6 +61,7 @@ import org.springframework.util.FileCopyUtils; import com.amazonaws.services.s3.AmazonS3; import com.amazonaws.services.s3.model.ListObjectsRequest; import com.amazonaws.services.s3.model.ObjectListing; +import com.amazonaws.services.s3.model.Region; import com.amazonaws.services.s3.model.S3Object; import com.amazonaws.services.s3.model.S3ObjectSummary; @@ -119,6 +121,8 @@ public class S3StreamingChannelAdapterTests { assertThat(message).isNotNull(); assertThat(message.getPayload()).isInstanceOf(InputStream.class); assertThat(message.getHeaders().get(FileHeaders.REMOTE_FILE)).isEqualTo("subdir/b.test"); + assertThat(message.getHeaders()) + .containsKeys(FileHeaders.REMOTE_DIRECTORY, FileHeaders.REMOTE_HOST_PORT, FileHeaders.REMOTE_FILE); InputStream inputStreamB = (InputStream) message.getPayload(); assertThat(IOUtils.toString(inputStreamB, Charset.defaultCharset())).isEqualTo("Bye"); @@ -149,7 +153,7 @@ public class S3StreamingChannelAdapterTests { for (final S3Object s3Object : S3_OBJECTS) { willAnswer(invocation -> s3Object).given(amazonS3).getObject(S3_BUCKET, s3Object.getKey()); } - + willReturn(Region.US_West).given(amazonS3).getRegion(); return amazonS3; }