From 851a2cfdee9cae9cb506cb668a074d63b5798fbc Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 25 May 2017 13:07:57 -0400 Subject: [PATCH] Fix CachingClientConnectionFactoryTests gatewayIntegrationTest() was stealing integrationTest()'s message. Shutdown its executor and wait for the shutdown. Also add a meaningful toString() to TcpConnectionSupport. --- .../ip/tcp/connection/TcpConnectionSupport.java | 5 +++++ .../CachingClientConnectionFactoryTests.java | 10 +++++++--- 2 files changed, 12 insertions(+), 3 deletions(-) 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 9cad55113f..fee3f2c6f0 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 @@ -395,4 +395,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 8bdc3a1f42..a3f5205bfb 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(20_000); 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(20_000); 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(() -> { + ExecutorService exec = Executors.newSingleThreadExecutor(); + exec.execute(() -> { while (okToRun.get()) { Message m = inbound.receive(1000); if (m != null) { @@ -511,6 +513,8 @@ public class CachingClientConnectionFactoryTests { assertEquals(connectionIds.get(0), connectionIds.get(1)); okToRun.set(false); + exec.shutdownNow(); + assertTrue(exec.awaitTermination(20, TimeUnit.SECONDS)); } @Test