From 7f92089584bbc0a03c787f851846fdd99487d1c7 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Fri, 14 Oct 2011 14:31:57 -0400 Subject: [PATCH] INT-2146 Expose session limitation on remote file adapters Deprecate 'cache-sessions' attribute Add 'sessionCacheSize' and 'sessionWaitTimeout' attributes on CachingSessionFactory Update documentation --- docs/src/reference/docbook/ftp.xml | 38 +++++-- docs/src/reference/docbook/sftp.xml | 37 +++++-- ...RemoteFileInboundChannelAdapterParser.java | 34 +++++-- ...stractRemoteFileOutboundGatewayParser.java | 35 +++++-- ...emoteFileOutboundChannelAdapterParser.java | 34 +++++-- .../remote/session/CachingSessionFactory.java | 93 ++++++++++-------- .../ftp/config/spring-integration-ftp-2.1.xsd | 4 +- ...boundChannelAdapterParserTests-context.xml | 10 +- .../ftp/session/SessionFactoryTests.java | 98 ++++++++++++++++++- .../config/spring-integration-sftp-2.1.xsd | 6 +- 10 files changed, 292 insertions(+), 97 deletions(-) diff --git a/docs/src/reference/docbook/ftp.xml b/docs/src/reference/docbook/ftp.xml index fbe16e72dd..faad788f9a 100644 --- a/docs/src/reference/docbook/ftp.xml +++ b/docs/src/reference/docbook/ftp.xml @@ -346,17 +346,35 @@ xsi:schemaLocation="http://www.springframework.org/schema/integration/ftp
FTP Session Caching - One of the optimizations implemented by the FTP adapters is session caching. Similar to JDBC pooling of Connections, the FTP Adapters maintain a - pool of Sessions by default. However there are times when this behavior is not desired (e.g., security etc.). - To disable session caching you can set the cache-sessions attribute to false (the default value is true). -]]> + + Since version 2.1 we've exposed more flexibility with regard to session management for remote file adapters (e.g., FTP, SFTP etc). + In previous versions the sessions were cached automatically. And although we did expose cache-session attribute which would + allow you to turn auto caching off it was still not sufficient when it came to other session caching attributes. + For example; One of the requirement we received was to control session limit if remote server imposes client connection limit. + Now such behavior is controlled by the CachingSessionFactory and its sessionCacheSize and + sessionWaitTimeout properties. As its name suggest sessionCacheSize property controls how many active + sessions this adapter will maintain in its cache (DEFAULT unbounded). If sessionCacheSize threshold has been reached any + attempt to get more session will block until session becomes available or until wait time for a session to become available + expires (DEFAULT Integer.MAX_VALUE). The wait time for a session to become available can also be controlled via sessionWaitTimeout + property. + + + Since version 2.1 if you need you session to be cahced you can simply configure your default Session Factory as + described above and than wrap it in the CachingSessionFactory while providing additional properties. + + + + + + + + + + ]]> + + In the above example you see CachingSessionFactory created with + sessionCacheSize set to 10 with sessionWait timeout set to 1 second. -The same attribute can also be used with Outbound Channel Adapters.
diff --git a/docs/src/reference/docbook/sftp.xml b/docs/src/reference/docbook/sftp.xml index 2c4505fca0..738ce65c55 100644 --- a/docs/src/reference/docbook/sftp.xml +++ b/docs/src/reference/docbook/sftp.xml @@ -290,17 +290,34 @@ xsi:schemaLocation="http://www.springframework.org/schema/integration/sftp
SFTP Session Caching - One of the optimizations implemented by the SFTP adapters is session caching. Similar to JDBC pooling of Connections, the SFTP Adapters maintain a - pool of Sessions by default. However there are times when this behavior is not desired (e.g., security etc.). - To disable session caching you can set the cache-sessions attribute to false (the default value is true). -]]> + Since version 2.1 we've exposed more flexibility with regard to session management for remote file adapters (e.g., FTP, SFTP etc). + In previous versions the sessions were cached automatically. And although we did expose cache-session attribute which would + allow you to turn auto caching off it was still not sufficient when it came to other session caching attributes. + For example; One of the requirement we received was to control session limit if remote server imposes client connection limit. + Now such behavior is controlled by the CachingSessionFactory and its sessionCacheSize and + sessionWaitTimeout properties. As its name suggest sessionCacheSize property controls how many active + sessions this adapter will maintain in its cache (DEFAULT unbounded). If sessionCacheSize threshold has been reached any + attempt to get more session will block until session becomes available or until wait time for a session to become available + expires (DEFAULT Integer.MAX_VALUE). The wait time for a session to become available can also be controlled via sessionWaitTimeout + property. + + + Since version 2.1 if you need you session to be cahced you can simply configure your default Session Factory as + described above and than wrap it in the CachingSessionFactory while providing additional properties. + + + + + + + + + + ]]> + + In the above example you see CachingSessionFactory created with + sessionCacheSize set to 10 with sessionWait timeout set to 1 second. -The same attribute can also be used with Outbound Channel Adapters.
diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/config/AbstractRemoteFileInboundChannelAdapterParser.java b/spring-integration-file/src/main/java/org/springframework/integration/file/config/AbstractRemoteFileInboundChannelAdapterParser.java index 5a7f25108e..ecb522ce3b 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/config/AbstractRemoteFileInboundChannelAdapterParser.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/config/AbstractRemoteFileInboundChannelAdapterParser.java @@ -18,12 +18,16 @@ package org.springframework.integration.file.config; import org.w3c.dom.Element; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.springframework.beans.BeanMetadataElement; +import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.ExpressionFactoryBean; import org.springframework.integration.config.xml.AbstractPollingInboundChannelAdapterParser; import org.springframework.integration.config.xml.IntegrationNamespaceUtils; +import org.springframework.integration.file.remote.session.CachingSessionFactory; import org.springframework.util.StringUtils; /** @@ -34,23 +38,35 @@ import org.springframework.util.StringUtils; * @since 2.0 */ public abstract class AbstractRemoteFileInboundChannelAdapterParser extends AbstractPollingInboundChannelAdapterParser { - + private final Log logger = LogFactory.getLog(this.getClass()); + @Override protected final BeanMetadataElement parseSource(Element element, ParserContext parserContext) { BeanDefinitionBuilder synchronizerBuilder = BeanDefinitionBuilder.genericBeanDefinition( this.getInboundFileSynchronizerClassname()); - // build the SessionFactory and provide as a constructor argument - String cacheSessions = element.getAttribute("cache-sessions"); - if ("false".equalsIgnoreCase(cacheSessions)) { - synchronizerBuilder.addConstructorArgReference(element.getAttribute("session-factory")); + // This whole block must be refactored once cache-session attribute is removed + String sessionFactoryName = element.getAttribute("session-factory"); + BeanDefinition sessionFactoryDefinition = parserContext.getReaderContext().getRegistry().getBeanDefinition(sessionFactoryName); + String sessionFactoryClassName = sessionFactoryDefinition.getBeanClassName(); + if (StringUtils.hasText(sessionFactoryClassName) && sessionFactoryClassName.endsWith(CachingSessionFactory.class.getName())){ + synchronizerBuilder.addConstructorArgValue(sessionFactoryDefinition); } else { - BeanDefinitionBuilder sessionFactoryBuilder = BeanDefinitionBuilder.genericBeanDefinition( - "org.springframework.integration.file.remote.session.CachingSessionFactory"); - sessionFactoryBuilder.addConstructorArgReference(element.getAttribute("session-factory")); - synchronizerBuilder.addConstructorArgValue(sessionFactoryBuilder.getBeanDefinition()); + String cacheSessions = element.getAttribute("cache-sessions"); + if (StringUtils.hasText(cacheSessions)){ + logger.warn("The 'cache-sessions' attribute is deprecated since v2.1. Consider configuring CachingSessionFactory explicitly"); + } + if ("false".equalsIgnoreCase(cacheSessions)) { + synchronizerBuilder.addConstructorArgReference(element.getAttribute("session-factory")); + } + else { + BeanDefinitionBuilder sessionFactoryBuilder = BeanDefinitionBuilder.genericBeanDefinition(CachingSessionFactory.class); + sessionFactoryBuilder.addConstructorArgReference(sessionFactoryName); + synchronizerBuilder.addConstructorArgValue(sessionFactoryBuilder.getBeanDefinition()); + } } + // end of what needs to be refactored once cache-session is removed // configure the InboundFileSynchronizer properties IntegrationNamespaceUtils.setValueIfAttributeDefined(synchronizerBuilder, element, "remote-directory"); diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/config/AbstractRemoteFileOutboundGatewayParser.java b/spring-integration-file/src/main/java/org/springframework/integration/file/config/AbstractRemoteFileOutboundGatewayParser.java index d380f4d52a..e7b83fde26 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/config/AbstractRemoteFileOutboundGatewayParser.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/config/AbstractRemoteFileOutboundGatewayParser.java @@ -15,20 +15,27 @@ */ package org.springframework.integration.file.config; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.xml.AbstractConsumerEndpointParser; import org.springframework.integration.config.xml.IntegrationNamespaceUtils; +import org.springframework.integration.file.remote.session.CachingSessionFactory; import org.springframework.util.StringUtils; import org.w3c.dom.Element; /** * @author Gary Russell + * @author Oleg Zhurakousky * @since 2.1 * */ public abstract class AbstractRemoteFileOutboundGatewayParser extends AbstractConsumerEndpointParser { + + private final Log logger = LogFactory.getLog(this.getClass()); @Override protected String getInputChannelAttributeName() { @@ -39,16 +46,30 @@ public abstract class AbstractRemoteFileOutboundGatewayParser extends protected BeanDefinitionBuilder parseHandler(Element element, ParserContext parserContext) { BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(getGatewayClassName()); // build the SessionFactory and provide as a constructor argument - String cacheSessions = element.getAttribute("cache-sessions"); - if ("false".equalsIgnoreCase(cacheSessions)) { - builder.addConstructorArgReference(element.getAttribute("session-factory")); + + // This whole block must be refactored once cache-session attribute is removed + String sessionFactoryName = element.getAttribute("session-factory"); + BeanDefinition sessionFactoryDefinition = parserContext.getReaderContext().getRegistry().getBeanDefinition(sessionFactoryName); + String sessionFactoryClassName = sessionFactoryDefinition.getBeanClassName(); + if (StringUtils.hasText(sessionFactoryClassName) && sessionFactoryClassName.endsWith(CachingSessionFactory.class.getName())){ + builder.addConstructorArgValue(sessionFactoryDefinition); } else { - BeanDefinitionBuilder sessionFactoryBuilder = BeanDefinitionBuilder.genericBeanDefinition( - "org.springframework.integration.file.remote.session.CachingSessionFactory"); - sessionFactoryBuilder.addConstructorArgReference(element.getAttribute("session-factory")); - builder.addConstructorArgValue(sessionFactoryBuilder.getBeanDefinition()); + String cacheSessions = element.getAttribute("cache-sessions"); + if (StringUtils.hasText(cacheSessions)){ + logger.warn("The 'cache-sessions' attribute is deprecated since v2.1. Consider configuring CachingSessionFactory explicitly"); + } + if ("false".equalsIgnoreCase(cacheSessions)) { + builder.addConstructorArgReference(element.getAttribute("session-factory")); + } + else { + BeanDefinitionBuilder sessionFactoryBuilder = BeanDefinitionBuilder.genericBeanDefinition(CachingSessionFactory.class); + sessionFactoryBuilder.addConstructorArgReference(sessionFactoryName); + builder.addConstructorArgValue(sessionFactoryBuilder.getBeanDefinition()); + } } + // end of what needs to be refactored once cache-session is removed + builder.addConstructorArgValue(element.getAttribute("command")); builder.addConstructorArgValue(element.getAttribute(EXPRESSION_ATTRIBUTE)); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "command-options", "options"); diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/config/RemoteFileOutboundChannelAdapterParser.java b/spring-integration-file/src/main/java/org/springframework/integration/file/config/RemoteFileOutboundChannelAdapterParser.java index 4a94afc433..7a3260fa9d 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/config/RemoteFileOutboundChannelAdapterParser.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/config/RemoteFileOutboundChannelAdapterParser.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2010 the original author or authors. + * Copyright 2002-2011 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. @@ -18,6 +18,8 @@ package org.springframework.integration.file.config; import org.w3c.dom.Element; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.springframework.beans.factory.BeanDefinitionStoreException; import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.support.AbstractBeanDefinition; @@ -26,6 +28,7 @@ import org.springframework.beans.factory.support.RootBeanDefinition; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.xml.AbstractOutboundChannelAdapterParser; import org.springframework.integration.config.xml.IntegrationNamespaceUtils; +import org.springframework.integration.file.remote.session.CachingSessionFactory; import org.springframework.util.StringUtils; /** @@ -34,23 +37,34 @@ import org.springframework.util.StringUtils; * @since 2.0 */ public class RemoteFileOutboundChannelAdapterParser extends AbstractOutboundChannelAdapterParser { - + private final Log logger = LogFactory.getLog(this.getClass()); @Override protected AbstractBeanDefinition parseConsumer(Element element, ParserContext parserContext) { BeanDefinitionBuilder handlerBuilder = BeanDefinitionBuilder.genericBeanDefinition( "org.springframework.integration.file.remote.handler.FileTransferringMessageHandler"); - // build the SessionFactory and provide as a constructor argument - String cacheSessions = element.getAttribute("cache-sessions"); - if ("false".equalsIgnoreCase(cacheSessions)) { - handlerBuilder.addConstructorArgReference(element.getAttribute("session-factory")); + // This whole block must be refactored once cache-session attribute is removed + String sessionFactoryName = element.getAttribute("session-factory"); + BeanDefinition sessionFactoryDefinition = parserContext.getReaderContext().getRegistry().getBeanDefinition(sessionFactoryName); + String sessionFactoryClassName = sessionFactoryDefinition.getBeanClassName(); + if (StringUtils.hasText(sessionFactoryClassName) && sessionFactoryClassName.endsWith(CachingSessionFactory.class.getName())){ + handlerBuilder.addConstructorArgValue(sessionFactoryDefinition); } else { - BeanDefinitionBuilder sessionFactoryBuilder = BeanDefinitionBuilder.genericBeanDefinition( - "org.springframework.integration.file.remote.session.CachingSessionFactory"); - sessionFactoryBuilder.addConstructorArgReference(element.getAttribute("session-factory")); - handlerBuilder.addConstructorArgValue(sessionFactoryBuilder.getBeanDefinition()); + String cacheSessions = element.getAttribute("cache-sessions"); + if (StringUtils.hasText(cacheSessions)){ + logger.warn("The 'cache-sessions' attribute is deprecated since v2.1. Consider configuring CachingSessionFactory explicitly"); + } + if ("false".equalsIgnoreCase(cacheSessions)) { + handlerBuilder.addConstructorArgReference(element.getAttribute("session-factory")); + } + else { + BeanDefinitionBuilder sessionFactoryBuilder = BeanDefinitionBuilder.genericBeanDefinition(CachingSessionFactory.class); + sessionFactoryBuilder.addConstructorArgReference(sessionFactoryName); + handlerBuilder.addConstructorArgValue(sessionFactoryBuilder.getBeanDefinition()); + } } + // end of what needs to be refactored once cache-session is removed // configure MessageHandler properties IntegrationNamespaceUtils.setValueIfAttributeDefined(handlerBuilder, element, "temporary-file-suffix"); diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/remote/session/CachingSessionFactory.java b/spring-integration-file/src/main/java/org/springframework/integration/file/remote/session/CachingSessionFactory.java index bb67b4efaa..aa2ef1c04e 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/remote/session/CachingSessionFactory.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/remote/session/CachingSessionFactory.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2010 the original author or authors. + * Copyright 2002-2011 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. @@ -19,13 +19,12 @@ package org.springframework.integration.file.remote.session; import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; -import java.util.Queue; -import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.LinkedBlockingQueue; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; - import org.springframework.beans.factory.DisposableBean; +import org.springframework.integration.util.UpperBound; /** * A {@link SessionFactory} implementation that caches Sessions for reuse without @@ -37,46 +36,47 @@ import org.springframework.beans.factory.DisposableBean; * @author Mark Fisher * @since 2.0 */ -public class CachingSessionFactory implements SessionFactory, DisposableBean { +public class CachingSessionFactory implements SessionFactory, DisposableBean{ private static final Log logger = LogFactory.getLog(CachingSessionFactory.class); - public static final int DEFAULT_POOL_SIZE = 10; + private volatile long sessionWaitTimeout = Integer.MAX_VALUE; - - private final Queue queue; + private volatile LinkedBlockingQueue queue = new LinkedBlockingQueue(); private final SessionFactory sessionFactory; - private final int maxPoolSize; - - + private final UpperBound sessionSizeManager; + public CachingSessionFactory(SessionFactory sessionFactory) { - this(sessionFactory, DEFAULT_POOL_SIZE); + this(sessionFactory, 0); } - - public CachingSessionFactory(SessionFactory sessionFactory, int maxPoolSize) { + + public CachingSessionFactory(SessionFactory sessionFactory, int sessionCacheSize) { this.sessionFactory = sessionFactory; - this.maxPoolSize = maxPoolSize; - this.queue = new ArrayBlockingQueue(this.maxPoolSize, true); + this.sessionSizeManager = new UpperBound(sessionCacheSize); + } + + /** + * Sets the limit of how long it will wait for a session to become available after which + * it will throw {@link IllegalStateException}. + */ + public void setSessionWaitTimeout(long sessionWaitTimeout) { + this.sessionWaitTimeout = sessionWaitTimeout; } - - public Session getSession() { - Session session = this.queue.poll(); - if (session == null || !session.isOpen()) { - if (session != null && logger.isTraceEnabled()) { - logger.trace("Located session in the pool but it is stale, will create new one."); - } - session = this.sessionFactory.getSession(); - if (logger.isTraceEnabled()) { - logger.trace("Created new session"); - } + public Session getSession() { + try { + boolean permitted = this.sessionSizeManager.tryAcquire(this.sessionWaitTimeout); + if (!permitted){ + throw new IllegalStateException("Timed out while waiting to aquire Session"); + } + Session session = this.doGetSession(); + return new CachedSession(session); + } catch (Exception e) { + Thread.currentThread().interrupt(); + throw new IllegalStateException("Exception was received during attempt to obtain Session", e); } - else if (logger.isTraceEnabled()) { - logger.trace("Using session from the pool"); - } - return new CachedSession(session); } public void destroy() { @@ -86,6 +86,21 @@ public class CachingSessionFactory implements SessionFactory, DisposableBean { } } } + + private Session doGetSession() throws InterruptedException { + Session session = this.queue.poll(); + + if (session != null && !session.isOpen()){ + if (logger.isDebugEnabled()) { + logger.debug("Received stale Session, will attempt to get a new one"); + } + return this.doGetSession(); + } + else if (session == null){ + session = this.sessionFactory.getSession(); + } + return session; + } private void closeSession(Session session) { try { @@ -111,18 +126,11 @@ public class CachingSessionFactory implements SessionFactory, DisposableBean { } public void close() { - if (queue.size() < maxPoolSize) { - if (logger.isTraceEnabled()) { - logger.trace("Releasing target session back to the pool"); - } - queue.add(targetSession); - } - else { - if (logger.isTraceEnabled()) { - logger.trace("Disconnecting target session"); - } - targetSession.close(); + if (logger.isDebugEnabled()){ + logger.debug("Releasing Session back into the pool"); } + queue.add(targetSession); + sessionSizeManager.release(); } public boolean remove(String path) throws IOException{ @@ -153,5 +161,4 @@ public class CachingSessionFactory implements SessionFactory, DisposableBean { this.targetSession.mkdir(directory); } } - } diff --git a/spring-integration-ftp/src/main/resources/org/springframework/integration/ftp/config/spring-integration-ftp-2.1.xsd b/spring-integration-ftp/src/main/resources/org/springframework/integration/ftp/config/spring-integration-ftp-2.1.xsd index 73f0b230e4..e2c57f14a0 100644 --- a/spring-integration-ftp/src/main/resources/org/springframework/integration/ftp/config/spring-integration-ftp-2.1.xsd +++ b/spring-integration-ftp/src/main/resources/org/springframework/integration/ftp/config/spring-integration-ftp-2.1.xsd @@ -233,7 +233,7 @@ endpoint itself is a Polling Consumer for a channel with a queue. @@ -384,7 +384,7 @@ endpoint itself is a Polling Consumer for a channel with a queue. diff --git a/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/config/FtpOutboundChannelAdapterParserTests-context.xml b/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/config/FtpOutboundChannelAdapterParserTests-context.xml index 54e292e0c9..19250d1901 100644 --- a/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/config/FtpOutboundChannelAdapterParserTests-context.xml +++ b/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/config/FtpOutboundChannelAdapterParserTests-context.xml @@ -15,6 +15,12 @@ + + + + + + diff --git a/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/session/SessionFactoryTests.java b/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/session/SessionFactoryTests.java index 2de9f0ede5..e8911cfddb 100644 --- a/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/session/SessionFactoryTests.java +++ b/spring-integration-ftp/src/test/java/org/springframework/integration/ftp/session/SessionFactoryTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2010 the original author or authors. + * Copyright 2002-2011 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. @@ -16,12 +16,25 @@ package org.springframework.integration.ftp.session; import java.lang.reflect.Field; +import java.util.Random; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import org.apache.commons.net.ftp.FTPClient; +import org.junit.Ignore; import org.junit.Test; +import org.mockito.Mockito; +import org.springframework.integration.file.remote.session.CachingSessionFactory; +import org.springframework.integration.file.remote.session.Session; +import org.springframework.integration.file.remote.session.SessionFactory; +import org.springframework.integration.test.util.TestUtils; import static junit.framework.Assert.fail; +import static org.junit.Assert.assertEquals; + /** * @author Oleg Zhurakousky * @@ -49,4 +62,87 @@ public class SessionFactoryTests { } } } + + @Test + public void testStaleConnection() throws Exception{ + SessionFactory sessionFactory = Mockito.mock(SessionFactory.class); + Session sessionA = Mockito.mock(Session.class); + Session sessionB = Mockito.mock(Session.class); + Mockito.when(sessionA.isOpen()).thenReturn(true); + Mockito.when(sessionB.isOpen()).thenReturn(false); + + Mockito.when(sessionFactory.getSession()).thenReturn(sessionA); + Mockito.when(sessionFactory.getSession()).thenReturn(sessionB); + + CachingSessionFactory cachingFactory = new CachingSessionFactory(sessionFactory, 2); + + Session firstSession = cachingFactory.getSession(); + Session secondSession = cachingFactory.getSession(); + secondSession.close(); + Session nonStaleSession = cachingFactory.getSession(); + assertEquals(TestUtils.getPropertyValue(firstSession, "targetSession"), TestUtils.getPropertyValue(nonStaleSession, "targetSession")); + } + + @Test + public void testSameSessionFromThePool() throws Exception{ + SessionFactory sessionFactory = Mockito.mock(SessionFactory.class); + Session session = Mockito.mock(Session.class); + Mockito.when(sessionFactory.getSession()).thenReturn(session); + + CachingSessionFactory cachingFactory = new CachingSessionFactory(sessionFactory, 2); + + Session s1 = cachingFactory.getSession(); + s1.close(); + Session s2 = cachingFactory.getSession(); + s2.close(); + assertEquals(TestUtils.getPropertyValue(s1, "targetSession"), TestUtils.getPropertyValue(s2, "targetSession")); + Mockito.verify(sessionFactory, Mockito.times(2)).getSession(); + } + + @Test (expected=IllegalStateException.class) // timeout expire + public void testSessionWaitExpire() throws Exception{ + SessionFactory sessionFactory = Mockito.mock(SessionFactory.class); + Session session = Mockito.mock(Session.class); + Mockito.when(sessionFactory.getSession()).thenReturn(session); + + CachingSessionFactory cachingFactory = new CachingSessionFactory(sessionFactory, 2); + + cachingFactory.setSessionWaitTimeout(3000); + + cachingFactory.getSession(); + cachingFactory.getSession(); + cachingFactory.getSession(); + } + + @Test + @Ignore + public void testConnectionLimit() throws Exception{ + ExecutorService executor = Executors.newCachedThreadPool(); + DefaultFtpSessionFactory sessionFactory = new DefaultFtpSessionFactory(); + sessionFactory.setHost("192.168.28.143"); + sessionFactory.setPassword("password"); + sessionFactory.setUsername("user"); + final CachingSessionFactory factory = new CachingSessionFactory(sessionFactory, 2); + + final Random random = new Random(); + final AtomicInteger failures = new AtomicInteger(); + for (int i = 0; i < 30; i++) { + executor.execute(new Runnable() { + public void run() { + try { + Session session = factory.getSession(); + Thread.sleep(random.nextInt(5000)); + session.close(); + } catch (Exception e) { + e.printStackTrace(); + failures.incrementAndGet(); + } + } + }); + } + executor.shutdown(); + executor.awaitTermination(10000, TimeUnit.SECONDS); + + assertEquals(0, failures.get()); + } } diff --git a/spring-integration-sftp/src/main/resources/org/springframework/integration/sftp/config/spring-integration-sftp-2.1.xsd b/spring-integration-sftp/src/main/resources/org/springframework/integration/sftp/config/spring-integration-sftp-2.1.xsd index 4b69190db9..a4776b30bd 100644 --- a/spring-integration-sftp/src/main/resources/org/springframework/integration/sftp/config/spring-integration-sftp-2.1.xsd +++ b/spring-integration-sftp/src/main/resources/org/springframework/integration/sftp/config/spring-integration-sftp-2.1.xsd @@ -37,7 +37,7 @@ @@ -168,7 +168,7 @@ endpoint itself is a Polling Consumer for a channel with a queue. @@ -302,7 +302,7 @@ endpoint itself is a Polling Consumer for a channel with a queue.