diff --git a/src/main/java/org/springframework/integration/aws/outbound/S3MessageHandler.java b/src/main/java/org/springframework/integration/aws/outbound/S3MessageHandler.java index c3c4873..4b357b9 100644 --- a/src/main/java/org/springframework/integration/aws/outbound/S3MessageHandler.java +++ b/src/main/java/org/springframework/integration/aws/outbound/S3MessageHandler.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2017 the original author or authors. + * Copyright 2016-2019 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. @@ -69,11 +69,13 @@ import com.amazonaws.util.Md5Utils; * invocation. Consider to use an async upstream hand off if this blocking behavior isn't appropriate. *
* The "request-reply" behavior is async and the {@link Transfer} result from the {@link TransferManager} - * operation is sent to the {@link #outputChannel}, assuming the transfer progress observation in the + * operation is sent to the {@link #getOutputChannel()}, assuming the transfer progress observation in the * downstream flow. *
* The {@link S3ProgressListener} can be supplied to track the transfer progress. * Also the listener can be populated into the returned {@link Transfer} afterwards in the downstream flow. + * If the context of the {@code requestMessage} is important in the {@code progressChanged} event, it is + * recommended to use a {@link MessageS3ProgressListener} implementation instead. * *
* For the upload operation the {@link UploadMetadataProvider} callback can be supplied to populate required
* {@link ObjectMetadata} options, as for a single entry, as well as for each file in directory to upload.
@@ -131,12 +133,8 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
}
public S3MessageHandler(AmazonS3 amazonS3, String bucket, boolean produceReply) {
- this(TransferManagerBuilder.standard()
- .withS3Client(amazonS3)
- .build(),
- bucket,
- produceReply);
- Assert.notNull(amazonS3, "'amazonS3' must not be null");
+ this(amazonS3, new LiteralExpression(bucket), produceReply);
+ Assert.notNull(bucket, "'bucket' must not be null");
}
public S3MessageHandler(AmazonS3 amazonS3, Expression bucketExpression, boolean produceReply) {
@@ -228,6 +226,7 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
/**
* Specify a {@link S3ProgressListener} for upload and download operations.
* @param s3ProgressListener the {@link S3ProgressListener} to use.
+ * @see MessageS3ProgressListener
*/
public void setProgressListener(S3ProgressListener s3ProgressListener) {
this.s3ProgressListener = s3ProgressListener;
@@ -260,9 +259,8 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
@Override
protected Object handleRequestMessage(Message> requestMessage) {
Command command = this.commandExpression.getValue(this.evaluationContext, requestMessage, Command.class);
- Assert.state(command != null, "'commandExpression' ["
- + this.commandExpression.getExpressionString()
- + "] cannot evaluate to null.");
+ Assert.state(command != null, () ->
+ "'commandExpression' [" + this.commandExpression.getExpressionString() + "] cannot evaluate to null.");
Transfer transfer = null;
@@ -327,8 +325,8 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
InputStream inputStream = (InputStream) payload;
if (metadata.getContentMD5() == null) {
Assert.state(inputStream.markSupported(),
- "For an upload InputStream with no MD5 digest metadata, the " +
- "markSupported() method must evaluate to true. ");
+ "For an upload InputStream with no MD5 digest metadata, " +
+ "the markSupported() method must evaluate to true.");
String contentMd5 = Md5Utils.md5AsBase64(inputStream);
metadata.setContentMD5(contentMd5);
inputStream.reset();
@@ -350,7 +348,9 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
if (metadata.getContentType() == null) {
metadata.setContentType(Mimetypes.getInstance().getMimetype(fileToUpload));
}
- putObjectRequest = new PutObjectRequest(bucketName, key, fileToUpload).withMetadata(metadata);
+ putObjectRequest =
+ new PutObjectRequest(bucketName, key, fileToUpload)
+ .withMetadata(metadata);
}
else if (payload instanceof byte[]) {
byte[] payloadBytes = (byte[]) payload;
@@ -387,12 +387,30 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
}
}
- S3ProgressListener progressListener = this.s3ProgressListener;
+ S3ProgressListener configuredProgressListener = this.s3ProgressListener;
+ if (this.s3ProgressListener instanceof MessageS3ProgressListener) {
+ configuredProgressListener = new S3ProgressListener() {
+
+ @Override
+ public void onPersistableTransfer(PersistableTransfer persistableTransfer) {
+ S3MessageHandler.this.s3ProgressListener.onPersistableTransfer(persistableTransfer);
+ }
+
+ @Override
+ public void progressChanged(ProgressEvent progressEvent) {
+ ((MessageS3ProgressListener) S3MessageHandler.this.s3ProgressListener)
+ .progressChanged(progressEvent, requestMessage);
+ }
+
+ };
+ }
+
+ S3ProgressListener progressListener = configuredProgressListener;
if (this.objectAclExpression != null) {
Object acl = this.objectAclExpression.getValue(this.evaluationContext, requestMessage);
- Assert.state(acl instanceof AccessControlList || acl instanceof CannedAccessControlList,
- "The 'objectAclExpression' ["
+ Assert.state(acl == null || acl instanceof AccessControlList || acl instanceof CannedAccessControlList,
+ () -> "The 'objectAclExpression' ["
+ this.objectAclExpression.getExpressionString()
+ "] must evaluate to com.amazonaws.services.s3.model.AccessControlList " +
"or must evaluate to com.amazonaws.services.s3.model.CannedAccessControlList. " +
@@ -424,8 +442,8 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
};
- if (this.s3ProgressListener != null) {
- progressListener = new S3ProgressListenerChain(this.s3ProgressListener, progressListener);
+ if (configuredProgressListener != null) {
+ progressListener = new S3ProgressListenerChain(configuredProgressListener, progressListener);
}
}
@@ -441,8 +459,9 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
private Transfer download(Message> requestMessage) {
Object payload = requestMessage.getPayload();
- Assert.state(payload instanceof File, "For the 'DOWNLOAD' operation the 'payload' must be of " +
- "'java.io.File' type, but gotten: [" + payload.getClass() + ']');
+ Assert.state(payload instanceof File,
+ () -> "For the 'DOWNLOAD' operation the 'payload' must be of " +
+ "'java.io.File' type, but gotten: [" + payload.getClass() + ']');
File targetFile = (File) payload;
@@ -457,7 +476,7 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
}
Assert.state(key != null,
- "The 'keyExpression' must not be null for non-File payloads and can't evaluate to null. " +
+ () -> "The 'keyExpression' must not be null for non-File payloads and can't evaluate to null. " +
"Root object is: " + requestMessage);
if (targetFile.isDirectory()) {
@@ -483,7 +502,7 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
}
Assert.state(sourceKey != null,
- "The 'keyExpression' must not be null for 'copy' operation " +
+ () -> "The 'keyExpression' must not be null for 'copy' operation " +
"and 'keyExpression' can't evaluate to null. " +
"Root object is: " + requestMessage);
@@ -498,8 +517,8 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
}
Assert.state(destinationBucketName != null,
- "The 'destinationBucketExpression' must not be null for 'copy' operation and can't evaluate to null. " +
- "Root object is: " + requestMessage);
+ () -> "The 'destinationBucketExpression' must not be null for 'copy' operation " +
+ "and can't evaluate to null. Root object is: " + requestMessage);
String destinationKey = null;
if (this.destinationKeyExpression != null) {
@@ -508,8 +527,8 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
}
Assert.state(destinationKey != null,
- "The 'destinationKeyExpression' must not be null for 'copy' operation and can't evaluate to null. " +
- "Root object is: " + requestMessage);
+ () -> "The 'destinationKeyExpression' must not be null for 'copy' operation " +
+ "and can't evaluate to null. Root object is: " + requestMessage);
CopyObjectRequest copyObjectRequest =
@@ -525,9 +544,10 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
else {
bucketName = this.bucketExpression.getValue(this.evaluationContext, requestMessage, String.class);
}
- Assert.state(bucketName != null, "The 'bucketExpression' ["
- + this.bucketExpression.getExpressionString()
- + "] must not evaluate to null. Root object is: " + requestMessage);
+ Assert.state(bucketName != null,
+ () -> "The 'bucketExpression' ["
+ + this.bucketExpression.getExpressionString()
+ + "] must not evaluate to null. Root object is: " + requestMessage);
if (this.resourceIdResolver != null) {
bucketName = this.resourceIdResolver.resolveToPhysicalResourceId(bucketName);
@@ -560,6 +580,23 @@ public class S3MessageHandler extends AbstractReplyProducingMessageHandler {
}
+ /**
+ * An {@link S3ProgressListener} extension to provide a {@code requestMessage}
+ * context for the {@code progressChanged} event.
+ *
+ * @since 2.1
+ */
+ public interface MessageS3ProgressListener extends S3ProgressListener {
+
+ @Override
+ default void progressChanged(ProgressEvent progressEvent) {
+ throw new UnsupportedOperationException("Use progressChanged(ProgressEvent, Message>) instead.");
+ }
+
+ void progressChanged(ProgressEvent progressEvent, Message> message);
+
+ }
+
/**
* The callback to populate an {@link ObjectMetadata} for upload operation.
* The message can be used as a metadata source.
diff --git a/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws-2.0.xsd b/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws-2.0.xsd
index ec0629f..8442d38 100644
--- a/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws-2.0.xsd
+++ b/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws-2.0.xsd
@@ -179,6 +179,9 @@