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 8e451b19f1..7f7a123b10 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 @@ -153,7 +153,15 @@ public class CachingClientConnectionFactory extends AbstractClientConnectionFact messageBuilder.setHeader(IpHeaders.ACTUAL_CONNECTION_ID, message.getHeaders().get(IpHeaders.CONNECTION_ID)); } - this.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 } @@ -187,8 +195,18 @@ public class CachingClientConnectionFactory extends AbstractClientConnectionFact } @Override - public boolean equals(Object obj) { - return targetConnectionFactory.equals(obj); + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (o == null || getClass() != o.getClass()) { + return false; + } + + CachingClientConnectionFactory that = (CachingClientConnectionFactory) o; + + return this.targetConnectionFactory.equals(that.targetConnectionFactory); + } @Override 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 820ed1a780..29283789d9 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 @@ -15,16 +15,19 @@ */ 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.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; @@ -39,8 +42,10 @@ import java.util.concurrent.Executors; import java.util.concurrent.Semaphore; 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; @@ -59,6 +64,8 @@ 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; import org.springframework.test.context.ContextConfiguration; @@ -102,6 +109,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();