Merge pull request #487 from garyrussell/INT-2590
* INT-2590: INT-2590 Add TestingUtility
This commit is contained in:
@@ -0,0 +1,91 @@
|
||||
/*
|
||||
* Copyright 2002-2012 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.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
package org.springframework.integration.ip.util;
|
||||
|
||||
import org.springframework.integration.ip.tcp.connection.AbstractServerConnectionFactory;
|
||||
|
||||
/**
|
||||
* Convenience class providing methods for testing IP components.
|
||||
* Provided in the main branch so that it is available for
|
||||
* use in user test code, samples etc.
|
||||
*
|
||||
* @author Gary Russell
|
||||
* @since 2.2
|
||||
*
|
||||
*/
|
||||
public class TestingUtilities {
|
||||
|
||||
/**
|
||||
* Wait for a server connection factory to actually start listening before
|
||||
* starting a test. Waits for up to 10 seconds by default.
|
||||
* @param serverConnectionFactory The server connection factory.
|
||||
* @param delay How long to wait in milliseconds; default 10000 (10 seconds) if null.
|
||||
* @throws IllegalStateException
|
||||
*/
|
||||
public static void waitListening(AbstractServerConnectionFactory serverConnectionFactory, Long delay)
|
||||
throws IllegalStateException {
|
||||
if (delay == null) {
|
||||
delay = 100L;
|
||||
}
|
||||
else {
|
||||
delay = delay / 100;
|
||||
}
|
||||
int n = 0;
|
||||
while (!serverConnectionFactory.isListening()) {
|
||||
try {
|
||||
Thread.sleep(100);
|
||||
} catch (InterruptedException e1) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new IllegalStateException(e1);
|
||||
}
|
||||
|
||||
if (n++ > delay) {
|
||||
throw new IllegalStateException ("Server didn't start listening.");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Wait for a server connection factory to stop listening.
|
||||
* Waits for up to 10 seconds by default.
|
||||
* @param serverConnectionFactory The server connection factory.
|
||||
* @param delay How long to wait in milliseconds; default 10000 (10 seconds) if null.
|
||||
* @throws IllegalStateException
|
||||
*/
|
||||
public static void waitStopListening(AbstractServerConnectionFactory serverConnectionFactory, Long delay)
|
||||
throws IllegalStateException {
|
||||
if (delay == null) {
|
||||
delay = 100L;
|
||||
}
|
||||
else {
|
||||
delay = delay / 100;
|
||||
}
|
||||
int n = 0;
|
||||
while (serverConnectionFactory.isListening()) {
|
||||
try {
|
||||
Thread.sleep(100);
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new IllegalStateException(e);
|
||||
}
|
||||
if (n++ > 200) {
|
||||
throw new IllegalStateException ("Server didn't stop listening.");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
@@ -17,13 +17,13 @@ package org.springframework.integration.ip.tcp;
|
||||
|
||||
import static org.junit.Assert.assertFalse;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.junit.Assert.fail;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.springframework.beans.DirectFieldAccessor;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.integration.ip.tcp.connection.AbstractServerConnectionFactory;
|
||||
import org.springframework.integration.ip.util.TestingUtilities;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
|
||||
@@ -58,19 +58,8 @@ public class AutoStartTests {
|
||||
*/
|
||||
private void startAndStop() throws InterruptedException {
|
||||
tcpNetIn.start();
|
||||
int n = 0;
|
||||
while (!cfS1.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to start listening");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(cfS1, null);
|
||||
tcpNetIn.stop();
|
||||
while (cfS1.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to stop listening");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitStopListening(cfS1, null);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -36,6 +36,7 @@ import org.springframework.integration.ip.tcp.connection.AbstractClientConnectio
|
||||
import org.springframework.integration.ip.tcp.connection.AbstractServerConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.TcpConnection;
|
||||
import org.springframework.integration.ip.tcp.serializer.ByteArrayRawSerializer;
|
||||
import org.springframework.integration.ip.util.TestingUtilities;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
@@ -50,16 +51,16 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
@ContextConfiguration
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
public class ConnectionToConnectionTests {
|
||||
|
||||
|
||||
@Autowired
|
||||
AbstractApplicationContext ctx;
|
||||
|
||||
|
||||
@Autowired
|
||||
private AbstractClientConnectionFactory client;
|
||||
|
||||
|
||||
@Autowired
|
||||
private AbstractServerConnectionFactory server;
|
||||
|
||||
|
||||
@Autowired
|
||||
private QueueChannel serverSideChannel;
|
||||
|
||||
@@ -74,18 +75,12 @@ public class ConnectionToConnectionTests {
|
||||
ConnectionToConnectionTests.class.getPackage().getName()
|
||||
.replaceAll("\\.", "/")
|
||||
+ "/ConnectionToConnectionTests-context.xml");
|
||||
ctx.close();
|
||||
ctx.close();
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void testConnect() throws Exception {
|
||||
int n = 0;
|
||||
while (!server.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
throw new Exception("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(server, null);
|
||||
client.start();
|
||||
for (int i = 0; i < 100; i++) {
|
||||
TcpConnection connection = client.getConnection();
|
||||
@@ -118,7 +113,7 @@ public class ConnectionToConnectionTests {
|
||||
assertNotNull(message);
|
||||
assertEquals("Test", new String((byte[]) message.getPayload()));
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void testLookup() throws Exception {
|
||||
client.start();
|
||||
@@ -134,5 +129,5 @@ public class ConnectionToConnectionTests {
|
||||
assertFalse(connection.getConnectionId().contains("localhost"));
|
||||
connection.close();
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
@@ -19,7 +19,6 @@ package org.springframework.integration.ip.tcp;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.fail;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
@@ -50,6 +49,7 @@ import org.springframework.integration.ip.tcp.connection.TcpNetClientConnectionF
|
||||
import org.springframework.integration.ip.tcp.connection.TcpNetServerConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.TcpNioServerConnectionFactory;
|
||||
import org.springframework.integration.ip.util.SocketTestUtils;
|
||||
import org.springframework.integration.ip.util.TestingUtilities;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.support.channel.ChannelResolver;
|
||||
|
||||
@@ -63,13 +63,7 @@ public class TcpInboundGatewayTests {
|
||||
TcpInboundGateway gateway = new TcpInboundGateway();
|
||||
gateway.setConnectionFactory(scf);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 200) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, 20000L);
|
||||
final QueueChannel channel = new QueueChannel();
|
||||
gateway.setRequestChannel(channel);
|
||||
ServiceActivatingHandler handler = new ServiceActivatingHandler(new Service());
|
||||
@@ -99,13 +93,7 @@ public class TcpInboundGatewayTests {
|
||||
TcpInboundGateway gateway = new TcpInboundGateway();
|
||||
gateway.setConnectionFactory(scf);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 200) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, 20000L);
|
||||
final QueueChannel channel = new QueueChannel();
|
||||
gateway.setRequestChannel(channel);
|
||||
ServiceActivatingHandler handler = new ServiceActivatingHandler(new Service());
|
||||
@@ -184,13 +172,7 @@ public class TcpInboundGatewayTests {
|
||||
TcpInboundGateway gateway = new TcpInboundGateway();
|
||||
gateway.setConnectionFactory(scf);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 200) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, 20000L);
|
||||
final QueueChannel channel = new QueueChannel();
|
||||
gateway.setRequestChannel(channel);
|
||||
ServiceActivatingHandler handler = new ServiceActivatingHandler(new Service());
|
||||
@@ -220,13 +202,7 @@ public class TcpInboundGatewayTests {
|
||||
TcpInboundGateway gateway = new TcpInboundGateway();
|
||||
gateway.setConnectionFactory(scf);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 200) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, 20000L);
|
||||
final QueueChannel channel = new QueueChannel();
|
||||
gateway.setRequestChannel(channel);
|
||||
ServiceActivatingHandler handler = new ServiceActivatingHandler(new Service());
|
||||
@@ -262,13 +238,7 @@ public class TcpInboundGatewayTests {
|
||||
});
|
||||
gateway.setErrorChannel(errorChannel);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 200) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, 20000L);
|
||||
final SubscribableChannel channel = new DirectChannel();
|
||||
gateway.setRequestChannel(channel);
|
||||
ServiceActivatingHandler handler = new ServiceActivatingHandler(new FailingService());
|
||||
|
||||
@@ -19,7 +19,6 @@ package org.springframework.integration.ip.tcp;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.fail;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
@@ -58,6 +57,7 @@ import org.springframework.integration.ip.tcp.connection.TcpNetServerConnectionF
|
||||
import org.springframework.integration.ip.tcp.connection.TcpNioServerConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.serializer.ByteArrayCrLfSerializer;
|
||||
import org.springframework.integration.ip.util.SocketTestUtils;
|
||||
import org.springframework.integration.ip.util.TestingUtilities;
|
||||
|
||||
/**
|
||||
* @author Gary Russell
|
||||
@@ -74,13 +74,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
TcpReceivingChannelAdapter adapter = new TcpReceivingChannelAdapter();
|
||||
adapter.setConnectionFactory(scf);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to start listening");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
adapter.afterPropertiesSet();
|
||||
@@ -156,13 +150,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
TcpReceivingChannelAdapter adapter = new TcpReceivingChannelAdapter();
|
||||
adapter.setConnectionFactory(scf);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to start listening");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
@@ -194,13 +182,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
scf.start();
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
socket.setSoTimeout(2000);
|
||||
socket.getOutputStream().write("Test\r\n".getBytes());
|
||||
@@ -232,13 +214,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
scf.start();
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
socket.setSoTimeout(2000);
|
||||
socket.getOutputStream().write("Test\r\n".getBytes());
|
||||
@@ -267,13 +243,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
TcpReceivingChannelAdapter adapter = new TcpReceivingChannelAdapter();
|
||||
adapter.setConnectionFactory(scf);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to start listening");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
@@ -303,13 +273,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
TcpReceivingChannelAdapter adapter = new TcpReceivingChannelAdapter();
|
||||
adapter.setConnectionFactory(scf);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to start listening");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
@@ -353,13 +317,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
scf.start();
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
Socket socket1 = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
socket1.setSoTimeout(2000);
|
||||
socket1.getOutputStream().write("Test1\r\n".getBytes());
|
||||
@@ -394,13 +352,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
scf.start();
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
Socket socket1 = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
socket1.setSoTimeout(2000);
|
||||
socket1.getOutputStream().write("Test1\r\n".getBytes());
|
||||
@@ -438,13 +390,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
scf.start();
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
List<Socket> sockets = new LinkedList<Socket>();
|
||||
for (int i = 100; i < 200; i++) {
|
||||
Socket socket1 = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
@@ -519,13 +465,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
scf.setInterceptorFactoryChain(fc);
|
||||
scf.setSoTimeout(10000);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to start listening");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
@@ -560,13 +500,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
TcpReceivingChannelAdapter adapter = new TcpReceivingChannelAdapter();
|
||||
adapter.setConnectionFactory(scf);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to start listening");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
@@ -612,13 +546,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
scf.start();
|
||||
QueueChannel channel = new QueueChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
Socket socket1 = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
socket1.setSoTimeout(60000);
|
||||
new ObjectOutputStream(socket1.getOutputStream()).writeObject("Hello");
|
||||
@@ -656,13 +584,7 @@ public class TcpReceivingChannelAdapterTests {
|
||||
TcpReceivingChannelAdapter adapter = new TcpReceivingChannelAdapter();
|
||||
adapter.setConnectionFactory(scf);
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 100) {
|
||||
fail("Failed to start listening");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
SubscribableChannel channel = new DirectChannel();
|
||||
adapter.setOutputChannel(channel);
|
||||
ServiceActivatingHandler handler = new ServiceActivatingHandler(new FailingService());
|
||||
|
||||
@@ -52,6 +52,7 @@ import org.springframework.integration.MessageChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.core.PollableChannel;
|
||||
import org.springframework.integration.ip.tcp.connection.AbstractConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.AbstractServerConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.HelloWorldInterceptorFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.TcpConnectionInterceptorFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.TcpConnectionInterceptorFactoryChain;
|
||||
@@ -61,6 +62,7 @@ import org.springframework.integration.ip.tcp.serializer.ByteArrayCrLfSerializer
|
||||
import org.springframework.integration.ip.tcp.serializer.ByteArrayLengthHeaderSerializer;
|
||||
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.message.GenericMessage;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
|
||||
@@ -1068,6 +1070,8 @@ public class TcpSendingMessageHandlerTests {
|
||||
AbstractConnectionFactory ccf = ctx.getBean("ccf", AbstractConnectionFactory.class);
|
||||
// TODO Lifecycle#start() isn't invoked within chain...
|
||||
ccf.start();
|
||||
AbstractServerConnectionFactory scf = ctx.getBean(AbstractServerConnectionFactory.class);
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
MessageChannel channelAdapterWithinChain = ctx.getBean("tcpOutboundChannelAdapterWithinChain", MessageChannel.class);
|
||||
PollableChannel inbound = ctx.getBean("inbound", PollableChannel.class);
|
||||
String testPayload = "Hello, world!";
|
||||
|
||||
@@ -19,7 +19,6 @@ import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertNotSame;
|
||||
import static org.junit.Assert.assertSame;
|
||||
import static org.junit.Assert.fail;
|
||||
import static org.mockito.Mockito.doAnswer;
|
||||
import static org.mockito.Mockito.doThrow;
|
||||
import static org.mockito.Mockito.mock;
|
||||
@@ -38,6 +37,7 @@ import org.springframework.integration.MessagingException;
|
||||
import org.springframework.integration.core.PollableChannel;
|
||||
import org.springframework.integration.core.SubscribableChannel;
|
||||
import org.springframework.integration.ip.IpHeaders;
|
||||
import org.springframework.integration.ip.util.TestingUtilities;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
@@ -270,14 +270,7 @@ public class CachingClientConnectionFactoryTests {
|
||||
|
||||
@Test
|
||||
public void integrationTest() throws Exception {
|
||||
int n = 0;
|
||||
while (!serverCf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
n++;
|
||||
if (n > 10000) {
|
||||
fail("Server didn't begin listening");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(serverCf, null);
|
||||
outbound.send(new GenericMessage<String>("Hello, world!"));
|
||||
Message<?> m = inbound.receive(1000);
|
||||
assertNotNull(m);
|
||||
|
||||
@@ -22,6 +22,7 @@ import java.net.Socket;
|
||||
import org.junit.Test;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.ip.util.SocketTestUtils;
|
||||
import org.springframework.integration.ip.util.TestingUtilities;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
|
||||
/**
|
||||
@@ -53,6 +54,7 @@ public class DefaultTimeoutTests {
|
||||
}
|
||||
});
|
||||
server.start();
|
||||
TestingUtilities.waitListening(server, null);
|
||||
client.start();
|
||||
TcpConnection connection = client.getConnection();
|
||||
Socket socket = TestUtils.getPropertyValue(connection, "socket", Socket.class);
|
||||
|
||||
@@ -50,6 +50,7 @@ import org.springframework.integration.ip.IpHeaders;
|
||||
import org.springframework.integration.ip.tcp.TcpInboundGateway;
|
||||
import org.springframework.integration.ip.tcp.TcpOutboundGateway;
|
||||
import org.springframework.integration.ip.util.SocketTestUtils;
|
||||
import org.springframework.integration.ip.util.TestingUtilities;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
|
||||
@@ -296,8 +297,8 @@ public class FailoverClientConnectionFactoryTests {
|
||||
gateway2.setConnectionFactory(server2);
|
||||
gateway2.setRequestChannel(channel);
|
||||
gateway2.start();
|
||||
waitListening(server1);
|
||||
waitListening(server2);
|
||||
TestingUtilities.waitListening(server1, null);
|
||||
TestingUtilities.waitListening(server2, null);
|
||||
List<AbstractClientConnectionFactory> factories = new ArrayList<AbstractClientConnectionFactory>();
|
||||
factories.add(client1);
|
||||
factories.add(client2);
|
||||
@@ -319,7 +320,7 @@ public class FailoverClientConnectionFactoryTests {
|
||||
Message<?> replyMessage = replyChannel.receive(10000);
|
||||
assertNotNull(replyMessage);
|
||||
server1.stop();
|
||||
waitStopListening(server1);
|
||||
TestingUtilities.waitStopListening(server1, null);
|
||||
outGateway.handleMessage(message);
|
||||
socket = getSocket(client2);
|
||||
port2 = socket.getLocalPort();
|
||||
@@ -339,24 +340,5 @@ public class FailoverClientConnectionFactoryTests {
|
||||
|
||||
}
|
||||
|
||||
private void waitListening(AbstractServerConnectionFactory scf) throws Exception {
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 200) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void waitStopListening(AbstractServerConnectionFactory scf) throws Exception {
|
||||
int n = 0;
|
||||
while (scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 200) {
|
||||
fail("Failed to stop listening");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -29,6 +29,7 @@ import javax.net.SocketFactory;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.integration.ip.util.TestingUtilities;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
|
||||
@@ -43,22 +44,22 @@ public class SOLingerTests {
|
||||
|
||||
@Autowired
|
||||
private AbstractServerConnectionFactory inCFNet;
|
||||
|
||||
|
||||
@Autowired
|
||||
private AbstractServerConnectionFactory inCFNio;
|
||||
|
||||
|
||||
@Autowired
|
||||
private AbstractServerConnectionFactory inCFNetRst;
|
||||
|
||||
|
||||
@Autowired
|
||||
private AbstractServerConnectionFactory inCFNioRst;
|
||||
|
||||
|
||||
@Autowired
|
||||
private AbstractServerConnectionFactory inCFNetLinger;
|
||||
|
||||
|
||||
@Autowired
|
||||
private AbstractServerConnectionFactory inCFNioLinger;
|
||||
|
||||
|
||||
@Test
|
||||
public void configOk() {}
|
||||
|
||||
@@ -71,7 +72,7 @@ public class SOLingerTests {
|
||||
public void finReceivedNio() {
|
||||
finReceived(inCFNio);
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void rstReceivedNet() {
|
||||
rstReceived(inCFNetRst);
|
||||
@@ -81,7 +82,7 @@ public class SOLingerTests {
|
||||
public void rstReceivedNio() {
|
||||
rstReceived(inCFNioRst);
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void finReceivedNetLinger() {
|
||||
finReceived(inCFNetLinger);
|
||||
@@ -91,21 +92,10 @@ public class SOLingerTests {
|
||||
public void finReceivedNioLinger() {
|
||||
finReceived(inCFNioLinger);
|
||||
}
|
||||
|
||||
|
||||
private void finReceived(AbstractServerConnectionFactory inCF) {
|
||||
int port = inCF.getPort();
|
||||
int n = 0;
|
||||
while (!inCF.isListening()) {
|
||||
try {
|
||||
Thread.sleep(100);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
fail("Interrupted");
|
||||
}
|
||||
if (n++ > 100) {
|
||||
fail("Failed to start");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(inCF, null);
|
||||
try {
|
||||
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
String test = "Test\r\n";
|
||||
@@ -113,30 +103,19 @@ public class SOLingerTests {
|
||||
byte[] buff = new byte[test.length() + 5];
|
||||
readFully(socket.getInputStream(), buff);
|
||||
assertEquals("echo:" + test, new String(buff));
|
||||
n = socket.getInputStream().read();
|
||||
int n = socket.getInputStream().read();
|
||||
// we expect an orderly close
|
||||
assertEquals(-1, n);
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
fail("Unexpected Exception " + e.getMessage());
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
|
||||
private void rstReceived(AbstractServerConnectionFactory inCF) {
|
||||
int port = inCF.getPort();
|
||||
int n = 0;
|
||||
while (!inCF.isListening()) {
|
||||
try {
|
||||
Thread.sleep(100);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
fail("Interrupted");
|
||||
}
|
||||
if (n++ > 100) {
|
||||
fail("Failed to start");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(inCF, null);
|
||||
try {
|
||||
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
socket.setSoTimeout(10000);
|
||||
@@ -146,7 +125,7 @@ public class SOLingerTests {
|
||||
readFully(socket.getInputStream(), buff);
|
||||
assertEquals("echo:" + test, new String(buff));
|
||||
try {
|
||||
n = socket.getInputStream().read();
|
||||
socket.getInputStream().read();
|
||||
fail("Expected IOException");
|
||||
} catch (IOException ioe) {
|
||||
assertTrue(ioe instanceof SocketException);
|
||||
@@ -155,12 +134,12 @@ public class SOLingerTests {
|
||||
e.printStackTrace();
|
||||
fail("Unexpected Exception " + e.getMessage());
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
private void readFully(InputStream is, byte[] buff) throws IOException {
|
||||
for (int i = 0; i < buff.length; i++) {
|
||||
buff[i] = (byte) is.read();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
@@ -17,7 +17,6 @@ package org.springframework.integration.ip.tcp.connection;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.fail;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.when;
|
||||
@@ -49,6 +48,7 @@ import org.springframework.integration.ip.tcp.connection.support.TcpSocketFactor
|
||||
import org.springframework.integration.ip.tcp.connection.support.TcpSocketSupport;
|
||||
import org.springframework.integration.ip.tcp.serializer.ByteArrayCrLfSerializer;
|
||||
import org.springframework.integration.ip.util.SocketTestUtils;
|
||||
import org.springframework.integration.ip.util.TestingUtilities;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
|
||||
@@ -156,7 +156,7 @@ public class SocketSupportTests {
|
||||
};
|
||||
serverConnectionFactory.setTcpSocketSupport(serverSocketSupport);
|
||||
serverConnectionFactory.start();
|
||||
waitListening(serverConnectionFactory);
|
||||
TestingUtilities.waitListening(serverConnectionFactory, null);
|
||||
clientConnectionFactory.getConnection().send(new GenericMessage<String>("Hello, world!"));
|
||||
assertTrue(latch.await(10, TimeUnit.SECONDS));
|
||||
assertEquals(0, ppServerSocketCountClient.get());
|
||||
@@ -281,7 +281,7 @@ Certificate fingerprints:
|
||||
}
|
||||
});
|
||||
server.start();
|
||||
waitListening(server);
|
||||
TestingUtilities.waitListening(server, null);
|
||||
|
||||
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", port);
|
||||
client.setTcpSocketFactorySupport(tcpSocketFactorySupport);
|
||||
@@ -313,7 +313,7 @@ Certificate fingerprints:
|
||||
}
|
||||
});
|
||||
server.start();
|
||||
waitListening(server);
|
||||
TestingUtilities.waitListening(server, null);
|
||||
|
||||
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", port);
|
||||
TcpSSLContextSupport clientSslContextSupport = new DefaultTcpSSLContextSupport("client.ks", "client.truststore.ks",
|
||||
@@ -351,7 +351,7 @@ Certificate fingerprints:
|
||||
}
|
||||
});
|
||||
server.start();
|
||||
waitListening(server);
|
||||
TestingUtilities.waitListening(server, null);
|
||||
|
||||
TcpNioClientConnectionFactory client = new TcpNioClientConnectionFactory("localhost", port);
|
||||
client.setTcpNioConnectionSupport(tcpNioConnectionSupport);
|
||||
@@ -400,7 +400,7 @@ Certificate fingerprints:
|
||||
deserializer.setMaxMessageSize(120000);
|
||||
server.setDeserializer(deserializer);
|
||||
server.start();
|
||||
waitListening(server);
|
||||
TestingUtilities.waitListening(server, null);
|
||||
|
||||
TcpNioClientConnectionFactory client = new TcpNioClientConnectionFactory("localhost", port);
|
||||
TcpSSLContextSupport clientSslContextSupport = new DefaultTcpSSLContextSupport("client.ks",
|
||||
@@ -431,16 +431,6 @@ Certificate fingerprints:
|
||||
assertEquals("Hello, world!", new String(payload).substring(0, 13));
|
||||
}
|
||||
|
||||
private void waitListening(AbstractServerConnectionFactory scf) throws Exception {
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (++n > 100) {
|
||||
fail("Server failed to start listening");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private class Replier implements TcpSender {
|
||||
|
||||
private TcpConnection connection;
|
||||
|
||||
@@ -38,6 +38,7 @@ import org.springframework.integration.ip.tcp.serializer.ByteArrayCrLfSerializer
|
||||
import org.springframework.integration.ip.tcp.serializer.ByteArrayLengthHeaderSerializer;
|
||||
import org.springframework.integration.ip.tcp.serializer.ByteArrayStxEtxSerializer;
|
||||
import org.springframework.integration.ip.util.SocketTestUtils;
|
||||
import org.springframework.integration.ip.util.TestingUtilities;
|
||||
|
||||
/**
|
||||
* @author Gary Russell
|
||||
@@ -62,13 +63,7 @@ public class TcpNioConnectionReadTests {
|
||||
scf.registerSender(sender);
|
||||
}
|
||||
scf.start();
|
||||
int n = 0;
|
||||
while (!scf.isListening()) {
|
||||
Thread.sleep(100);
|
||||
if (n++ > 200) {
|
||||
fail("Failed to listen");
|
||||
}
|
||||
}
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
return scf;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user