From 247232bdde24b81814a82100743f77d881aaf06b Mon Sep 17 00:00:00 2001 From: Gunnar Hillert Date: Wed, 17 Jun 2015 17:17:08 -0400 Subject: [PATCH] INT-3659: InputStream support in FileWritingMH https://jira.spring.io/browse/INT-3659 Add `InputStream` `payload` handling to the `FileWritingMessageHandler` --- .../file/FileWritingMessageHandler.java | 109 +++++++++++------- .../file/FileWritingMessageHandlerTests.java | 65 ++++++++++- src/reference/asciidoc/file.adoc | 7 +- src/reference/asciidoc/whats-new.adoc | 2 + 4 files changed, 139 insertions(+), 44 deletions(-) diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java b/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java index afe877b51b..f11f0deba6 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java @@ -23,6 +23,7 @@ import java.io.File; import java.io.FileInputStream; import java.io.FileOutputStream; import java.io.IOException; +import java.io.InputStream; import java.io.OutputStreamWriter; import java.nio.charset.Charset; @@ -81,11 +82,11 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand private static final String LINE_SEPARATOR = System.getProperty("line.separator"); - private volatile String temporaryFileSuffix =".writing"; + private volatile String temporaryFileSuffix = ".writing"; private volatile boolean temporaryFileSuffixSet = false; - private volatile FileExistsMode fileExistsMode = FileExistsMode.REPLACE; + private volatile FileExistsMode fileExistsMode = FileExistsMode.REPLACE; private final Log logger = LogFactory.getLog(this.getClass()); @@ -164,11 +165,11 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand * case the destination exists. For example {@link FileExistsMode#APPEND} * instructs this handler to append data to the existing file rather then * creating a new file for each {@link Message}. - * + *

* If set to {@link FileExistsMode#APPEND}, the adapter will also * create a real instance of the {@link LockRegistry} to ensure that there * is no collisions when multiple threads are writing to the same file. - * + *

* Otherwise the LockRegistry is set to {@link PassThruLockRegistry} which * has no effect. * @@ -179,7 +180,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand Assert.notNull(fileExistsMode, "'fileExistsMode' must not be null."); this.fileExistsMode = fileExistsMode; - if (FileExistsMode.APPEND.equals(fileExistsMode)){ + if (FileExistsMode.APPEND.equals(fileExistsMode)) { this.lockRegistry = this.lockRegistry instanceof PassThruLockRegistry ? new DefaultLockRegistry() : this.lockRegistry; @@ -198,6 +199,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand /** * If 'true' will append a new-line after each write. It is 'false' by default. + * * @param appendNewLine true if a new-line should be written to the file after payload is written * @since 4.0.7 */ @@ -266,7 +268,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand if (!destinationDirectory.exists() && autoCreateDirectory) { Assert.isTrue(destinationDirectory.mkdirs(), - "Destination directory [" + destinationDirectory + "] could not be created."); + "Destination directory [" + destinationDirectory + "] could not be created."); } Assert.isTrue(destinationDirectory.exists(), @@ -286,11 +288,11 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand Object payload = requestMessage.getPayload(); Assert.notNull(payload, "message payload must not be null"); String generatedFileName = this.fileNameGenerator.generateFileName(requestMessage); - File originalFileFromHeader = this.retrieveOriginalFileFromHeader(requestMessage); + File originalFileFromHeader = retrieveOriginalFileFromHeader(requestMessage); final File destinationDirectoryToUse = evaluateDestinationDirectoryExpression(requestMessage); - File tempFile = new File(destinationDirectoryToUse, generatedFileName + temporaryFileSuffix); + File tempFile = new File(destinationDirectoryToUse, generatedFileName + this.temporaryFileSuffix); File resultFile = new File(destinationDirectoryToUse, generatedFileName); if (FileExistsMode.FAIL.equals(this.fileExistsMode) && resultFile.exists()) { @@ -306,7 +308,11 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand try { if (payload instanceof File) { - resultFile = this.handleFileMessage((File) payload, tempFile, resultFile); + resultFile = handleFileMessage((File) payload, tempFile, resultFile); + } + else if (payload instanceof InputStream) { + resultFile = handleInputStreamMessage((InputStream) payload, originalFileFromHeader, tempFile, + resultFile); } else if (payload instanceof byte[]) { resultFile = this.handleByteArrayMessage( @@ -357,17 +363,34 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand } private File handleFileMessage(final File sourceFile, File tempFile, final File resultFile) throws IOException { + if (!FileExistsMode.APPEND.equals(this.fileExistsMode) && this.deleteSourceFiles) { + if (sourceFile.renameTo(resultFile)) { + return resultFile; + } + if (logger.isInfoEnabled()) { + logger.info(String.format("Failed to move file '%s'. Using copy and delete fallback.", + sourceFile.getAbsolutePath())); + } + } + final BufferedInputStream bis = new BufferedInputStream(new FileInputStream(sourceFile)); + return handleInputStreamMessage(bis, sourceFile, tempFile, resultFile); + } + + private File handleInputStreamMessage(final InputStream sourceFileInputStream, File originalFile, File tempFile, + final File resultFile) throws IOException { if (FileExistsMode.APPEND.equals(this.fileExistsMode)) { File fileToWriteTo = this.determineFileToWrite(resultFile, tempFile); final BufferedOutputStream bos = new BufferedOutputStream(new FileOutputStream(fileToWriteTo, true)); - final BufferedInputStream bis = new BufferedInputStream(new FileInputStream(sourceFile)); - WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry, fileToWriteTo.getAbsolutePath()){ + + WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry, + fileToWriteTo.getAbsolutePath()) { + @Override protected void whileLocked() throws IOException { try { byte[] buffer = new byte[StreamUtils.BUFFER_SIZE]; int bytesRead = -1; - while ((bytesRead = bis.read(buffer)) != -1) { + while ((bytesRead = sourceFileInputStream.read(buffer)) != -1) { bos.write(buffer, 0, bytesRead); } if (FileWritingMessageHandler.this.appendNewLine) { @@ -377,7 +400,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand } finally { try { - bis.close(); + sourceFileInputStream.close(); } catch (IOException ex) { } @@ -388,29 +411,20 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand } } } + }; whileLockedProcessor.doWhileLocked(); - this.cleanUpAfterCopy(fileToWriteTo, resultFile, sourceFile); + cleanUpAfterCopy(fileToWriteTo, resultFile, originalFile); return resultFile; } else { - if (this.deleteSourceFiles) { - if (sourceFile.renameTo(resultFile)) { - return resultFile; - } - if (logger.isInfoEnabled()) { - logger.info(String.format("Failed to move file '%s'. Using copy and delete fallback.", - sourceFile.getAbsolutePath())); - } - } BufferedOutputStream bos = new BufferedOutputStream(new FileOutputStream(tempFile)); - BufferedInputStream bis = new BufferedInputStream(new FileInputStream(sourceFile)); try { byte[] buffer = new byte[StreamUtils.BUFFER_SIZE]; int bytesRead = -1; - while ((bytesRead = bis.read(buffer)) != -1) { + while ((bytesRead = sourceFileInputStream.read(buffer)) != -1) { bos.write(buffer, 0, bytesRead); } if (this.appendNewLine) { @@ -420,7 +434,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand } finally { try { - bis.close(); + sourceFileInputStream.close(); } catch (IOException ex) { } @@ -430,18 +444,21 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand catch (IOException ex) { } } - this.cleanUpAfterCopy(tempFile, resultFile, sourceFile); + cleanUpAfterCopy(tempFile, resultFile, originalFile); return resultFile; } } - private File handleByteArrayMessage(final byte[] bytes, File originalFile, File tempFile, final File resultFile) throws IOException { + private File handleByteArrayMessage(final byte[] bytes, File originalFile, File tempFile, final File resultFile) + throws IOException { File fileToWriteTo = this.determineFileToWrite(resultFile, tempFile); final boolean append = FileExistsMode.APPEND.equals(this.fileExistsMode); final BufferedOutputStream bos = new BufferedOutputStream(new FileOutputStream(fileToWriteTo, append)); - WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry, fileToWriteTo.getAbsolutePath()){ + WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry, + fileToWriteTo.getAbsolutePath()) { + @Override protected void whileLocked() throws IOException { try { @@ -465,13 +482,17 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand return resultFile; } - private File handleStringMessage(final String content, File originalFile, File tempFile, final File resultFile) throws IOException { + private File handleStringMessage(final String content, File originalFile, File tempFile, final File resultFile) + throws IOException { File fileToWriteTo = this.determineFileToWrite(resultFile, tempFile); final boolean append = FileExistsMode.APPEND.equals(this.fileExistsMode); - final BufferedWriter writer = new BufferedWriter(new OutputStreamWriter(new FileOutputStream(fileToWriteTo, append), this.charset)); - WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry, fileToWriteTo.getAbsolutePath()){ + final BufferedWriter writer = + new BufferedWriter(new OutputStreamWriter(new FileOutputStream(fileToWriteTo, append), this.charset)); + WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry, + fileToWriteTo.getAbsolutePath()) { + @Override protected void whileLocked() throws IOException { try { @@ -497,7 +518,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand return resultFile; } - private File determineFileToWrite(File resultFile, File tempFile){ + private File determineFileToWrite(File resultFile, File tempFile) { final File fileToWriteTo; @@ -512,12 +533,12 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand break; default: throw new IllegalStateException("Unsupported FileExistsMode " - + this.fileExistsMode); + + this.fileExistsMode); } return fileToWriteTo; } - private void cleanUpAfterCopy(File fileToWriteTo, File resultFile, File originalFile) throws IOException{ + private void cleanUpAfterCopy(File fileToWriteTo, File resultFile, File originalFile) throws IOException { if (!FileExistsMode.APPEND.equals(this.fileExistsMode) && StringUtils.hasText(this.temporaryFileSuffix)) { this.renameTo(fileToWriteTo, resultFile); } @@ -527,24 +548,27 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand } } - private void renameTo(File tempFile, File resultFile) throws IOException{ + private void renameTo(File tempFile, File resultFile) throws IOException { Assert.notNull(resultFile, "'resultFile' must not be null"); Assert.notNull(tempFile, "'tempFile' must not be null"); if (resultFile.exists()) { - if (resultFile.setWritable(true, false) && resultFile.delete()){ + if (resultFile.setWritable(true, false) && resultFile.delete()) { if (!tempFile.renameTo(resultFile)) { - throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() + "' to '" + resultFile.getAbsolutePath() + "'"); + throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() + + "' to '" + resultFile.getAbsolutePath() + "'"); } } else { - throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() + "' to '" + resultFile.getAbsolutePath() + + throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() + + "' to '" + resultFile.getAbsolutePath() + "' since '" + resultFile.getName() + "' is not writable or can not be deleted"); } } else { if (!tempFile.renameTo(resultFile)) { - throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() + "' to '" + resultFile.getAbsolutePath() + "'"); + throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() + + "' to '" + resultFile.getAbsolutePath() + "'"); } } } @@ -558,7 +582,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand if (destinationDirectoryToUse == null) { throw new IllegalStateException(String.format("The provided " + - "destinationDirectoryExpression (%s) must not resolve to null.", + "destinationDirectoryExpression (%s) must not resolve to null.", this.destinationDirectoryExpression.getExpressionString())); } else if (destinationDirectoryToUse instanceof String) { @@ -572,7 +596,8 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand } else if (destinationDirectoryToUse instanceof File) { destinationDirectory = (File) destinationDirectoryToUse; - } else { + } + else { throw new IllegalStateException(String.format("The provided " + "destinationDirectoryExpression (%s) must be of type " + "java.io.File or be a String.", this.destinationDirectoryExpression.getExpressionString())); diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/FileWritingMessageHandlerTests.java b/spring-integration-file/src/test/java/org/springframework/integration/file/FileWritingMessageHandlerTests.java index 2a697f8241..bcbb49fdcb 100644 --- a/spring-integration-file/src/test/java/org/springframework/integration/file/FileWritingMessageHandlerTests.java +++ b/spring-integration-file/src/test/java/org/springframework/integration/file/FileWritingMessageHandlerTests.java @@ -29,8 +29,10 @@ import static org.junit.Assert.fail; import static org.mockito.Mockito.mock; import java.io.File; +import java.io.FileInputStream; import java.io.FileOutputStream; import java.io.IOException; +import java.io.InputStream; import java.io.UnsupportedEncodingException; import org.junit.Before; @@ -38,7 +40,6 @@ import org.junit.Ignore; import org.junit.Rule; import org.junit.Test; import org.junit.rules.TemporaryFolder; - import org.springframework.beans.factory.BeanFactory; import org.springframework.integration.channel.NullChannel; import org.springframework.integration.channel.QueueChannel; @@ -55,6 +56,7 @@ import org.springframework.util.FileCopyUtils; * @author Alex Peters * @author Gary Russell * @author Tony Falabella + * @author Gunnar Hillert */ public class FileWritingMessageHandlerTests { @@ -171,6 +173,29 @@ public class FileWritingMessageHandlerTests { assertFileContentIs(result, SAMPLE_CONTENT + System.getProperty("line.separator")); } + @Test + public void inputStreamPayloadCopiedToNewFile() throws Exception { + InputStream is = new FileInputStream(sourceFile); + Message message = MessageBuilder.withPayload(is).build(); + QueueChannel output = new QueueChannel(); + handler.setOutputChannel(output); + handler.handleMessage(message); + Message result = output.receive(0); + assertFileContentIsMatching(result); + } + + @Test + public void inputStreamPayloadCopiedToNewFileWithNewLines() throws Exception { + InputStream is = new FileInputStream(sourceFile); + Message message = MessageBuilder.withPayload(is).build(); + QueueChannel output = new QueueChannel(); + handler.setOutputChannel(output); + handler.setAppendNewLine(true); + handler.handleMessage(message); + Message result = output.receive(0); + assertFileContentIs(result, SAMPLE_CONTENT + System.getProperty("line.separator")); + } + @Test @Ignore // INT-3289 ignored because it won't fail on all OS public void testCreateDirFail() { File dir = new File("/foo"); @@ -274,6 +299,44 @@ public class FileWritingMessageHandlerTests { assertFalse(sourceFile.exists()); } + @Test + public void deleteSourceFileWithInputstreamPayloadAndFileInstanceHeader() throws Exception { + QueueChannel output = new QueueChannel(); + handler.setCharset(DEFAULT_ENCODING); + handler.setDeleteSourceFiles(true); + handler.setOutputChannel(output); + + InputStream is = new FileInputStream(sourceFile); + + Message message = MessageBuilder.withPayload(is) + .setHeader(FileHeaders.ORIGINAL_FILE, sourceFile) + .build(); + assertTrue(sourceFile.exists()); + handler.handleMessage(message); + Message result = output.receive(0); + assertFileContentIsMatching(result); + assertFalse(sourceFile.exists()); + } + + @Test + public void deleteSourceFileWithInputstreamPayloadAndFilePathHeader() throws Exception { + QueueChannel output = new QueueChannel(); + handler.setCharset(DEFAULT_ENCODING); + handler.setDeleteSourceFiles(true); + handler.setOutputChannel(output); + + InputStream is = new FileInputStream(sourceFile); + + Message message = MessageBuilder.withPayload(is) + .setHeader(FileHeaders.ORIGINAL_FILE, sourceFile.getAbsolutePath()) + .build(); + assertTrue(sourceFile.exists()); + handler.handleMessage(message); + Message result = output.receive(0); + assertFileContentIsMatching(result); + assertFalse(sourceFile.exists()); + } + @Test public void customFileNameGenerator() throws Exception { final String anyFilename = "fooBar.test"; diff --git a/src/reference/asciidoc/file.adoc b/src/reference/asciidoc/file.adoc index f169c12c44..e86558e2f7 100644 --- a/src/reference/asciidoc/file.adoc +++ b/src/reference/asciidoc/file.adoc @@ -261,7 +261,12 @@ IMPORTANT: Specifying the `delay`, `end` or `reopen` attributes, forces the use === Writing files To write messages to the file system you can use a http://static.springsource.org/spring-integration/api/org/springframework/integration/file/FileWritingMessageHandler.html[FileWritingMessageHandler]. -This class can deal with _File_, _String_, or _byte array_ payloads. +This class can deal with the following payload types: + +* _File_, +* _String_ +* _byte array_ +* _InputStream_ (since _version 4.2_) You can configure the encoding and the charset that will be used in case of a String payload. diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index 0b170c5766..ad6e1d3301 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -48,6 +48,8 @@ The `ignore-hidden` attribute has been introduced for the `> for more information. [[x4.2-class-package-change]]