INT-1614 factored all FTP and SFTP MessageHandler code into a single FileTransferringMessageHandler in the spring-integration-file module

This commit is contained in:
Mark Fisher
2010-11-21 18:54:12 -05:00
parent fd083e9772
commit 2e3fc66bba
10 changed files with 85 additions and 281 deletions

View File

@@ -45,7 +45,7 @@ public class SftpOutboundChannelAdapterParser extends AbstractOutboundChannelAda
String sessionPoolName = BeanDefinitionReaderUtils.registerWithGeneratedName(
sessionPoolBuilder.getBeanDefinition(), parserContext.getRegistry());
BeanDefinitionBuilder handlerBuilder = BeanDefinitionBuilder.genericBeanDefinition(
"org.springframework.integration.sftp.outbound.SftpSendingMessageHandler");
"org.springframework.integration.file.remote.handler.FileTransferringMessageHandler");
handlerBuilder.addConstructorArgReference(sessionPoolName);
IntegrationNamespaceUtils.setValueIfAttributeDefined(handlerBuilder, element, "charset");
String remoteDirectory = element.getAttribute("remote-directory");

View File

@@ -1,190 +0,0 @@
/*
* Copyright 2002-2010 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.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.sftp.outbound;
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;
import org.springframework.core.io.FileSystemResource;
import org.springframework.core.io.Resource;
import org.springframework.expression.Expression;
import org.springframework.integration.Message;
import org.springframework.integration.MessageDeliveryException;
import org.springframework.integration.MessagingException;
import org.springframework.integration.file.DefaultFileNameGenerator;
import org.springframework.integration.file.FileNameGenerator;
import org.springframework.integration.file.remote.session.Session;
import org.springframework.integration.file.remote.session.SessionFactory;
import org.springframework.integration.handler.AbstractMessageHandler;
import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor;
import org.springframework.util.Assert;
import org.springframework.util.FileCopyUtils;
import org.springframework.util.StringUtils;
/**
* Sends message payloads to a remote SFTP endpoint.
* Assumes that the payload of the inbound message is of type {@link java.io.File}.
*
* @author Josh Long
* @author Oleg Zhurakousky
* @since 2.0
*/
public class SftpSendingMessageHandler extends AbstractMessageHandler {
private static final String TEMPORARY_FILE_SUFFIX = ".writing";
private final SessionFactory sessionFactory;
private volatile ExpressionEvaluatingMessageProcessor<String> directoryExpressionProcesor;
private volatile Expression remoteDirectoryExpression;
private volatile FileNameGenerator fileNameGenerator = new DefaultFileNameGenerator();
private volatile File temporaryBufferFolderFile;
private volatile Resource temporaryBufferFolder =
new FileSystemResource(System.getProperty("java.io.tmpdir"));
private volatile String charset = Charset.defaultCharset().name();
public SftpSendingMessageHandler(SessionFactory sessionFactory) {
Assert.notNull(sessionFactory, "sessionFactory must not be null");
this.sessionFactory = sessionFactory;
}
public void setTemporaryBufferFolder(Resource temporaryBufferFolder) {
this.temporaryBufferFolder = temporaryBufferFolder;
}
public void setFileNameGenerator(FileNameGenerator fileNameGenerator) {
this.fileNameGenerator = fileNameGenerator;
}
public void setRemoteDirectoryExpression(Expression remoteDirectoryExpression) {
this.remoteDirectoryExpression = remoteDirectoryExpression;
}
public void setCharset(String charset) {
this.charset = charset;
}
@Override
protected void onInit() throws Exception {
this.temporaryBufferFolderFile = this.temporaryBufferFolder.getFile();
if (this.remoteDirectoryExpression != null) {
this.directoryExpressionProcesor =
new ExpressionEvaluatingMessageProcessor<String>(this.remoteDirectoryExpression, String.class);
}
}
@Override
protected void handleMessageInternal(Message<?> message) throws Exception {
File inboundFilePayload = this.redeemForStorableFile(message);
try {
if ((inboundFilePayload != null) && inboundFilePayload.exists()) {
this.sendFileToRemoteEndpoint(message, inboundFilePayload);
}
}
catch (Exception e) {
throw new MessageDeliveryException(message, "Failed to transfer '" + message.getPayload() + "' to " +
this.remoteDirectoryExpression.getExpressionString(), e);
}
finally {
if (inboundFilePayload != null && inboundFilePayload.exists()) {
inboundFilePayload.delete();
}
}
}
private File handleFileMessage(File sourceFile, File tempFile, File resultFile) throws IOException {
FileCopyUtils.copy(sourceFile, tempFile);
tempFile.renameTo(resultFile);
return resultFile;
}
private File handleByteArrayMessage(byte[] bytes, File tempFile, File resultFile) throws IOException {
FileCopyUtils.copy(bytes, tempFile);
tempFile.renameTo(resultFile);
return resultFile;
}
private File handleStringMessage(String content, File tempFile, File resultFile, String charset) throws IOException {
OutputStreamWriter writer = new OutputStreamWriter(new FileOutputStream(tempFile), charset);
FileCopyUtils.copy(content, writer);
tempFile.renameTo(resultFile);
return resultFile;
}
private File redeemForStorableFile(Message<?> message) throws MessageDeliveryException {
try {
Object payload = message.getPayload();
String generateFileName = this.fileNameGenerator.generateFileName(message);
File tempFile = new File(this.temporaryBufferFolderFile, generateFileName + TEMPORARY_FILE_SUFFIX);
File resultFile = new File(this.temporaryBufferFolderFile, generateFileName);
File sendableFile = null;
if (payload instanceof String) {
sendableFile = this.handleStringMessage((String) payload, tempFile, resultFile, this.charset);
}
else if (payload instanceof File) {
sendableFile = this.handleFileMessage((File) payload, tempFile, resultFile);
}
else if (payload instanceof byte[]) {
sendableFile = this.handleByteArrayMessage((byte[]) payload, tempFile, resultFile);
}
return sendableFile;
}
catch (Exception e) {
throw new MessageDeliveryException(message, "Failed to create sendable file.", e);
}
}
private boolean sendFileToRemoteEndpoint(Message<?> message, File file) throws Exception {
Session session = this.sessionFactory.getSession();
if (session == null) {
throw new MessagingException("The session returned from the pool is null, cannot proceed.");
}
InputStream fileInputStream = null;
try {
fileInputStream = new FileInputStream(file);
String baseOfRemotePath = "";
if (this.directoryExpressionProcesor != null) {
String result = this.directoryExpressionProcesor.processMessage(message);
if (StringUtils.hasText(result)) {
baseOfRemotePath = result;
}
}
if (!StringUtils.endsWithIgnoreCase(baseOfRemotePath, "/")) {
baseOfRemotePath += "/";
}
session.put(fileInputStream, baseOfRemotePath + file.getName());
return true;
}
finally {
fileInputStream.close();
session.close();
}
}
}

View File

@@ -30,8 +30,8 @@ import org.springframework.expression.common.LiteralExpression;
import org.springframework.expression.spel.standard.SpelExpression;
import org.springframework.integration.endpoint.EventDrivenConsumer;
import org.springframework.integration.file.FileNameGenerator;
import org.springframework.integration.file.remote.handler.FileTransferringMessageHandler;
import org.springframework.integration.file.remote.session.CachingSessionFactory;
import org.springframework.integration.sftp.outbound.SftpSendingMessageHandler;
import org.springframework.integration.sftp.session.DefaultSftpSessionFactory;
import org.springframework.integration.test.util.TestUtils;
@@ -49,8 +49,8 @@ public class OutboundChannelAdapaterParserTests {
assertTrue(consumer instanceof EventDrivenConsumer);
assertEquals(context.getBean("inputChannel"), TestUtils.getPropertyValue(consumer, "inputChannel"));
assertEquals("sftpOutboundAdapter", ((EventDrivenConsumer)consumer).getComponentName());
SftpSendingMessageHandler handler = (SftpSendingMessageHandler) TestUtils.getPropertyValue(consumer, "handler");
Expression remoteDirectoryExpression = (Expression) TestUtils.getPropertyValue(handler, "remoteDirectoryExpression");
FileTransferringMessageHandler handler = (FileTransferringMessageHandler) TestUtils.getPropertyValue(consumer, "handler");
Expression remoteDirectoryExpression = (Expression) TestUtils.getPropertyValue(handler, "directoryExpressionProcessor.expression");
assertNotNull(remoteDirectoryExpression);
assertTrue(remoteDirectoryExpression instanceof LiteralExpression);
assertEquals(context.getBean("fileNameGenerator"), TestUtils.getPropertyValue(handler, "fileNameGenerator"));
@@ -71,8 +71,8 @@ public class OutboundChannelAdapaterParserTests {
assertTrue(consumer instanceof EventDrivenConsumer);
assertEquals(context.getBean("inputChannel"), TestUtils.getPropertyValue(consumer, "inputChannel"));
assertEquals("sftpOutboundAdapterWithExpression", ((EventDrivenConsumer)consumer).getComponentName());
SftpSendingMessageHandler handler = (SftpSendingMessageHandler) TestUtils.getPropertyValue(consumer, "handler");
SpelExpression remoteDirectoryExpression = (SpelExpression) TestUtils.getPropertyValue(handler, "remoteDirectoryExpression");
FileTransferringMessageHandler handler = (FileTransferringMessageHandler) TestUtils.getPropertyValue(consumer, "handler");
SpelExpression remoteDirectoryExpression = (SpelExpression) TestUtils.getPropertyValue(handler, "directoryExpressionProcessor.expression");
assertNotNull(remoteDirectoryExpression);
assertEquals("'foo' + '/' + 'bar'", remoteDirectoryExpression.getExpressionString());
FileNameGenerator generator = (FileNameGenerator) TestUtils.getPropertyValue(handler, "fileNameGenerator");

View File

@@ -27,6 +27,7 @@ import java.io.File;
import org.junit.Test;
import org.springframework.expression.spel.standard.SpelExpressionParser;
import org.springframework.integration.file.remote.handler.FileTransferringMessageHandler;
import org.springframework.integration.file.remote.session.Session;
import org.springframework.integration.file.remote.session.SessionFactory;
import org.springframework.integration.message.GenericMessage;
@@ -43,9 +44,8 @@ public class SftpSendingMessageHandlerTests {
SessionFactory sessionFactory = mock(SessionFactory.class);
Session session = mock(Session.class);
when(sessionFactory.getSession()).thenReturn(session);
SftpSendingMessageHandler handler = new SftpSendingMessageHandler(sessionFactory);
FileTransferringMessageHandler handler = new FileTransferringMessageHandler(sessionFactory);
handler.setRemoteDirectoryExpression(new SpelExpressionParser().parseExpression("'foo.txt'"));
handler.handleMessage(new GenericMessage("hello"));
verify(sessionFactory, times(1)).getSession();
}
@@ -56,7 +56,7 @@ public class SftpSendingMessageHandlerTests {
SessionFactory sessionFactory = mock(SessionFactory.class);
Session session = mock(Session.class);
when(sessionFactory.getSession()).thenReturn(session);
SftpSendingMessageHandler handler = new SftpSendingMessageHandler(sessionFactory);
FileTransferringMessageHandler handler = new FileTransferringMessageHandler(sessionFactory);
handler.setRemoteDirectoryExpression(new SpelExpressionParser().parseExpression("'foo.txt'"));
handler.handleMessage(new GenericMessage("hello".getBytes()));
@@ -69,7 +69,7 @@ public class SftpSendingMessageHandlerTests {
SessionFactory sessionFactory = mock(SessionFactory.class);
Session session = mock(Session.class);
when(sessionFactory.getSession()).thenReturn(session);
SftpSendingMessageHandler handler = new SftpSendingMessageHandler(sessionFactory);
FileTransferringMessageHandler handler = new FileTransferringMessageHandler(sessionFactory);
handler.setRemoteDirectoryExpression(new SpelExpressionParser().parseExpression("'foo.txt'"));
handler.handleMessage(new GenericMessage("hello".getBytes()));