diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java index 87a2cb82da..8c3fcc489e 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/TcpNioConnectionReadTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2017 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. @@ -34,6 +34,7 @@ import java.util.concurrent.atomic.AtomicReference; import javax.net.SocketFactory; +import org.junit.Rule; import org.junit.Test; import org.springframework.context.ApplicationEventPublisher; @@ -43,16 +44,21 @@ import org.springframework.integration.ip.tcp.serializer.ByteArrayLengthHeaderSe import org.springframework.integration.ip.tcp.serializer.ByteArrayStxEtxSerializer; import org.springframework.integration.ip.util.SocketTestUtils; import org.springframework.integration.ip.util.TestingUtilities; +import org.springframework.integration.test.support.LongRunningIntegrationTest; import org.springframework.messaging.Message; import org.springframework.messaging.support.ErrorMessage; /** * @author Gary Russell * @author Artem Bilan + * * @since 2.0 */ public class TcpNioConnectionReadTests { + @Rule + public LongRunningIntegrationTest longRunningIntegrationTest = new LongRunningIntegrationTest(); + private final CountDownLatch latch = new CountDownLatch(1); private AbstractServerConnectionFactory getConnectionFactory( @@ -81,15 +87,10 @@ public class TcpNioConnectionReadTests { ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer(); final List> responses = new ArrayList>(); final Semaphore semaphore = new Semaphore(0); - AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() { - - @Override - public boolean onMessage(Message message) { - responses.add(message); - semaphore.release(); - return false; - } - + AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> { + responses.add(message); + semaphore.release(); + return false; }); // Fire up the sender. @@ -102,34 +103,28 @@ public class TcpNioConnectionReadTests { assertEquals("Data", SocketTestUtils.TEST_STRING + SocketTestUtils.TEST_STRING, new String((byte[]) responses.get(0).getPayload())); assertEquals("Data", SocketTestUtils.TEST_STRING + SocketTestUtils.TEST_STRING, - new String((byte[]) responses.get(1).getPayload())); + new String((byte[]) responses.get(1).getPayload())); scf.stop(); done.countDown(); } - @SuppressWarnings("unchecked") @Test public void testFragmented() throws Exception { ByteArrayLengthHeaderSerializer serializer = new ByteArrayLengthHeaderSerializer(); final List> responses = new ArrayList>(); final Semaphore semaphore = new Semaphore(0); - AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() { - - @Override - public boolean onMessage(Message message) { - responses.add(message); - try { - Thread.sleep(1000); - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } - semaphore.release(); - return false; + AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> { + responses.add(message); + try { + Thread.sleep(10); } - + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + semaphore.release(); + return false; }); int howMany = 2; @@ -140,7 +135,7 @@ public class TcpNioConnectionReadTests { assertEquals("Expected", howMany, responses.size()); for (int i = 0; i < howMany; i++) { assertEquals("Data", "xx", - new String(((Message) responses.get(0)).getPayload())); + new String(((Message) responses.get(0)).getPayload())); } scf.stop(); done.countDown(); @@ -152,15 +147,10 @@ public class TcpNioConnectionReadTests { ByteArrayStxEtxSerializer serializer = new ByteArrayStxEtxSerializer(); final List> responses = new ArrayList>(); final Semaphore semaphore = new Semaphore(0); - AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() { - - @Override - public boolean onMessage(Message message) { - responses.add(message); - semaphore.release(); - return false; - } - + AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> { + responses.add(message); + semaphore.release(); + return false; }); // Fire up the sender. @@ -171,9 +161,9 @@ public class TcpNioConnectionReadTests { assertTrue(semaphore.tryAcquire(1, 10000, TimeUnit.MILLISECONDS)); assertEquals("Did not receive data", 2, responses.size()); assertEquals("Data", SocketTestUtils.TEST_STRING + SocketTestUtils.TEST_STRING, - new String(((Message) responses.get(0)).getPayload())); + new String(((Message) responses.get(0)).getPayload())); assertEquals("Data", SocketTestUtils.TEST_STRING + SocketTestUtils.TEST_STRING, - new String(((Message) responses.get(1)).getPayload())); + new String(((Message) responses.get(1)).getPayload())); scf.stop(); done.countDown(); } @@ -184,15 +174,10 @@ public class TcpNioConnectionReadTests { ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer(); final List> responses = new ArrayList>(); final Semaphore semaphore = new Semaphore(0); - AbstractServerConnectionFactory scf = getConnectionFactory(serializer, new TcpListener() { - - @Override - public boolean onMessage(Message message) { - responses.add(message); - semaphore.release(); - return false; - } - + AbstractServerConnectionFactory scf = getConnectionFactory(serializer, message -> { + responses.add(message); + semaphore.release(); + return false; }); // Fire up the sender. @@ -203,9 +188,9 @@ public class TcpNioConnectionReadTests { assertTrue(semaphore.tryAcquire(1, 10000, TimeUnit.MILLISECONDS)); assertEquals("Did not receive data", 2, responses.size()); assertEquals("Data", SocketTestUtils.TEST_STRING + SocketTestUtils.TEST_STRING, - new String(((Message) responses.get(0)).getPayload())); + new String(((Message) responses.get(0)).getPayload())); assertEquals("Data", SocketTestUtils.TEST_STRING + SocketTestUtils.TEST_STRING, - new String(((Message) responses.get(1)).getPayload())); + new String(((Message) responses.get(1)).getPayload())); scf.stop(); done.countDown(); } @@ -251,7 +236,8 @@ public class TcpNioConnectionReadTests { assertTrue(errorMessageLetch.await(10, TimeUnit.SECONDS)); assertThat(errorMessageRef.get().getMessage(), - containsString("Message length 2147483647 exceeds max message length: 2048")); + anyOf(containsString("Message length 2147483647 exceeds max message length: 2048"), + containsString("Connection is closed"))); assertTrue(semaphore.tryAcquire(10000, TimeUnit.MILLISECONDS)); assertTrue(removed.size() > 0); @@ -510,11 +496,13 @@ public class TcpNioConnectionReadTests { } return false; }, new TcpSender() { + @Override public void addNewConnection(TcpConnection connection) { added.add(connection); semaphore.release(); } + @Override public void removeDeadConnection(TcpConnection connection) { removed.add(connection);