diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java index b17a8ffbb8..d29f8f13ec 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java @@ -403,4 +403,9 @@ public abstract class TcpConnectionSupport implements TcpConnection { } } + @Override + public String toString() { + return getClass().getSimpleName() + ":" + this.connectionId; + } + } 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 e4cd50bd2e..5659cbb04f 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 @@ -46,6 +46,7 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.BlockingQueue; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; @@ -458,13 +459,13 @@ public class CachingClientConnectionFactoryTests { new DirectFieldAccessor(this.clientAdapterCf).setPropertyValue("port", this.serverCf.getPort()); this.outbound.send(new GenericMessage<>("Hello, world!")); - Message m = inbound.receive(10000); + Message m = inbound.receive(20000); assertNotNull(m); String connectionId = m.getHeaders().get(IpHeaders.CONNECTION_ID, String.class); // assert we use the same connection from the pool outbound.send(new GenericMessage("Hello, world!")); - m = inbound.receive(10000); + m = inbound.receive(20000); assertNotNull(m); assertEquals(connectionId, m.getHeaders().get(IpHeaders.CONNECTION_ID, String.class)); } @@ -474,7 +475,8 @@ public class CachingClientConnectionFactoryTests { public void gatewayIntegrationTest() throws Exception { final List connectionIds = new ArrayList(); final AtomicBoolean okToRun = new AtomicBoolean(true); - Executors.newSingleThreadExecutor().execute(new Runnable() { + ExecutorService exec = Executors.newSingleThreadExecutor(); + exec.execute(new Runnable() { @Override public void run() { @@ -516,6 +518,8 @@ public class CachingClientConnectionFactoryTests { assertEquals(connectionIds.get(0), connectionIds.get(1)); okToRun.set(false); + exec.shutdownNow(); + assertTrue(exec.awaitTermination(20, TimeUnit.SECONDS)); } @Test