diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactory.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactory.java index 780c3dc615..69bab12d9f 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactory.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactory.java @@ -193,7 +193,15 @@ public class CachingClientConnectionFactory extends AbstractClientConnectionFact messageBuilder.setHeader(IpHeaders.ACTUAL_CONNECTION_ID, message.getHeaders().get(IpHeaders.CONNECTION_ID)); } - getListener().onMessage(messageBuilder.build()); + TcpListener listener = getListener(); + if (listener != null) { + listener.onMessage(messageBuilder.build()); + } + else { + if (logger.isDebugEnabled()) { + logger.debug("Message discarded; no listener: " + message); + } + } close(); // return to pool after response is received return true; // true so the single-use connection doesn't close itself } @@ -228,8 +236,12 @@ public class CachingClientConnectionFactory extends AbstractClientConnectionFact @Override public boolean equals(Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; + if (this == o) { + return true; + } + if (o == null || getClass() != o.getClass()) { + return false; + } CachingClientConnectionFactory that = (CachingClientConnectionFactory) o; diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactoryTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactoryTests.java index 9e1cd97ba4..277329eb36 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactoryTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/CachingClientConnectionFactoryTests.java @@ -16,11 +16,13 @@ package org.springframework.integration.ip.tcp.connection; +import static org.hamcrest.Matchers.startsWith; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNotSame; import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; import static org.mockito.Matchers.any; @@ -28,6 +30,7 @@ import static org.mockito.Matchers.anyInt; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -44,8 +47,10 @@ import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import org.apache.commons.logging.Log; import org.junit.Test; import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; import org.mockito.Mockito; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; @@ -63,6 +68,7 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessagingException; import org.springframework.messaging.PollableChannel; import org.springframework.messaging.SubscribableChannel; +import org.springframework.messaging.support.ErrorMessage; import org.springframework.messaging.support.GenericMessage; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.annotation.DirtiesContext.ClassMode; @@ -108,6 +114,16 @@ public class CachingClientConnectionFactoryTests { CachingClientConnectionFactory cachingFactory = new CachingClientConnectionFactory(factory, 2); cachingFactory.start(); TcpConnection conn1 = cachingFactory.getConnection(); + // INT-3652 + TcpConnectionInterceptorSupport cachedConn1 = (TcpConnectionInterceptorSupport) conn1; + Log logger = spy(TestUtils.getPropertyValue(cachedConn1, "logger", Log.class)); + when(logger.isDebugEnabled()).thenReturn(true); + new DirectFieldAccessor(cachedConn1).setPropertyValue("logger", logger); + cachedConn1.onMessage(new ErrorMessage(new RuntimeException())); + ArgumentCaptor captor = ArgumentCaptor.forClass(String.class); + verify(logger).debug(captor.capture()); + assertThat(captor.getValue(), startsWith("Message discarded; no listener:")); + // end INT-3652 assertEquals("Cached:" + mockConn1.toString(), conn1.toString()); conn1.close(); conn1 = cachingFactory.getConnection();