TCP/UDP: Use OS to Select Test Ports

https://build.spring.io/browse/INT-AT42SIO-228/

Conflicts:
	spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactoryTests.java
	spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/MultiClientTests.java
	spring-integration-ip/src/test/java/org/springframework/integration/ip/udp/UdpChannelAdapterTests.java

Fully revert `UdpChannelAdapterTests` changes because we "Support OS Selected UDP Port" only since `4.3`: https://jira.spring.io/browse/INT-3894
This commit is contained in:
Gary Russell
2016-07-29 13:29:11 -04:00
committed by Artem Bilan
parent 94eabb6b2e
commit c1a8180134
3 changed files with 152 additions and 136 deletions

View File

@@ -35,7 +35,6 @@ import org.springframework.context.ApplicationEventPublisher;
import org.springframework.integration.ip.util.TestingUtilities;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.test.support.LongRunningIntegrationTest;
import org.springframework.integration.test.util.SocketUtils;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.ErrorMessage;
@@ -51,12 +50,12 @@ public class ConnectionTimeoutTests {
@Test
public void testDefaultTimeout() throws Exception {
int port = SocketUtils.findAvailableServerSocket();
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(port);
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", port);
this.setupCallbacks(server, client, 0);
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(0);
this.setupServerCallbacks(server, 0);
server.start();
TestingUtilities.waitListening(server, null);
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", server.getPort());
setupClientCallback(client);
client.start();
TcpConnection connection = client.getConnection();
Socket socket = TestUtils.getPropertyValue(connection, "socket", Socket.class);
@@ -69,10 +68,11 @@ public class ConnectionTimeoutTests {
@Test
public void testNetSimpleTimeout() throws Exception {
int port = SocketUtils.findAvailableServerSocket();
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(port);
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", port);
this.setupCallbacks(server, client, 0);
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(0);
this.setupServerCallbacks(server, 0);
server.start();
TestingUtilities.waitListening(server, null);
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", server.getPort());
client.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
@@ -80,9 +80,8 @@ public class ConnectionTimeoutTests {
}
});
client.setSoTimeout(1000);
server.start();
TestingUtilities.waitListening(server, null);
CountDownLatch clientCloseLatch = getCloseLatch(client);
setupClientCallback(client);
client.start();
TcpConnection connection = client.getConnection();
Socket socket = TestUtils.getPropertyValue(connection, "socket", Socket.class);
@@ -100,9 +99,9 @@ public class ConnectionTimeoutTests {
*/
@Test
public void testNetReplyNotTimeout() throws Exception {
int port = SocketUtils.findAvailableServerSocket();
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(port);
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", port);
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(0);
setupAndStartServer(server);
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", server.getPort());
this.notTimeoutGuts(server, client);
}
@@ -113,15 +112,14 @@ public class ConnectionTimeoutTests {
*/
@Test
public void testNioReplyNotTimeout() throws Exception {
int port = SocketUtils.findAvailableServerSocket();
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(port);
TcpNioClientConnectionFactory client = new TcpNioClientConnectionFactory("localhost", port);
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(0);
setupAndStartServer(server);
TcpNioClientConnectionFactory client = new TcpNioClientConnectionFactory("localhost", server.getPort());
this.notTimeoutGuts(server, client);
}
private void notTimeoutGuts(AbstractServerConnectionFactory server, AbstractClientConnectionFactory client)
throws Exception, InterruptedException {
this.setupCallbacks(server, client, 1200);
final AtomicReference<Message<?>> reply = new AtomicReference<Message<?>>();
final CountDownLatch replyLatch = new CountDownLatch(1);
client.registerListener(new TcpListener() {
@@ -135,9 +133,8 @@ public class ConnectionTimeoutTests {
}
});
client.setSoTimeout(2000);
server.start();
TestingUtilities.waitListening(server, null);
CountDownLatch clientClosedLatch = getCloseLatch(client);
setupClientCallback(client);
client.start();
TcpConnection connection = client.getConnection();
Thread.sleep(1000);
@@ -150,6 +147,12 @@ public class ConnectionTimeoutTests {
client.stop();
}
private void setupAndStartServer(AbstractServerConnectionFactory server) {
this.setupServerCallbacks(server, 1200);
server.start();
TestingUtilities.waitListening(server, null);
}
/**
* Ensure we do timeout on the read side (client) if we sent a message within the
* first timeout but the reply takes > 2 timeouts.
@@ -157,11 +160,12 @@ public class ConnectionTimeoutTests {
*/
@Test
public void testNetReplyTimeout() throws Exception {
int port = SocketUtils.findAvailableServerSocket();
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(port);
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", port);
this.setupCallbacks(server, client, 4500);
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(0);
this.setupServerCallbacks(server, 4500);
final AtomicReference<Message<?>> reply = new AtomicReference<Message<?>>();
server.start();
TestingUtilities.waitListening(server, null);
TcpNetClientConnectionFactory client = new TcpNetClientConnectionFactory("localhost", server.getPort());
client.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
@@ -172,9 +176,8 @@ public class ConnectionTimeoutTests {
}
});
client.setSoTimeout(2000);
server.start();
TestingUtilities.waitListening(server, null);
CountDownLatch clientCloseLatch = getCloseLatch(client);
setupClientCallback(client);
client.start();
TcpConnection connection = client.getConnection();
Socket socket = TestUtils.getPropertyValue(connection, "socket", Socket.class);
@@ -197,11 +200,12 @@ public class ConnectionTimeoutTests {
*/
@Test
public void testNioReplyTimeout() throws Exception {
int port = SocketUtils.findAvailableServerSocket();
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(port);
TcpNioClientConnectionFactory client = new TcpNioClientConnectionFactory("localhost", port);
this.setupCallbacks(server, client, 2100);
TcpNetServerConnectionFactory server = new TcpNetServerConnectionFactory(0);
this.setupServerCallbacks(server, 2100);
final AtomicReference<Message<?>> reply = new AtomicReference<Message<?>>();
server.start();
TestingUtilities.waitListening(server, null);
TcpNioClientConnectionFactory client = new TcpNioClientConnectionFactory("localhost", server.getPort());
client.registerListener(new TcpListener() {
@Override
public boolean onMessage(Message<?> message) {
@@ -212,9 +216,8 @@ public class ConnectionTimeoutTests {
}
});
client.setSoTimeout(1000);
server.start();
TestingUtilities.waitListening(server, null);
CountDownLatch clientCloseLatch = getCloseLatch(client);
setupClientCallback(client);
client.start();
TcpConnection connection = client.getConnection();
Thread.sleep(500);
@@ -228,9 +231,7 @@ public class ConnectionTimeoutTests {
client.stop();
}
private void setupCallbacks(AbstractServerConnectionFactory server, AbstractClientConnectionFactory client,
final int serverDelay) {
client.setComponentName("clientFactory");
private void setupServerCallbacks(AbstractServerConnectionFactory server, final int serverDelay) {
server.setComponentName("serverFactory");
final AtomicReference<TcpConnection> serverConnection = new AtomicReference<TcpConnection>();
server.registerListener(new TcpListener() {
@@ -255,6 +256,10 @@ public class ConnectionTimeoutTests {
public void removeDeadConnection(TcpConnection connection) {
}
});
}
public void setupClientCallback(AbstractClientConnectionFactory client) {
client.setComponentName("clientFactory");
client.registerSender(new TcpSender() {
@Override
public void addNewConnection(TcpConnection connection) {

View File

@@ -59,7 +59,6 @@ import org.springframework.integration.ip.tcp.TcpInboundGateway;
import org.springframework.integration.ip.tcp.TcpOutboundGateway;
import org.springframework.integration.ip.util.TestingUtilities;
import org.springframework.integration.test.rule.Log4jLevelAdjuster;
import org.springframework.integration.test.util.SocketUtils;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.integration.util.SimplePool;
import org.springframework.messaging.Message;
@@ -76,6 +75,19 @@ import org.springframework.messaging.support.GenericMessage;
*/
public class FailoverClientConnectionFactoryTests {
private static final ApplicationEventPublisher NULL_PUBLISHER = new ApplicationEventPublisher() {
@Override
public void publishEvent(ApplicationEvent event) {
}
@Override
public void publishEvent(Object event) {
}
};
@Rule
public Log4jLevelAdjuster adjuster = new Log4jLevelAdjuster(Level.TRACE,
"org.springframework.integration.ip.tcp", "org.springframework.integration.util.SimplePool");
@@ -249,65 +261,57 @@ public class FailoverClientConnectionFactoryTests {
@Test
public void testRealNet() throws Exception {
AbstractServerConnectionFactory server1 = new TcpNetServerConnectionFactory(0);
AbstractServerConnectionFactory server2 = new TcpNetServerConnectionFactory(0);
final List<Integer> openPorts = SocketUtils.findAvailableServerSockets(SocketUtils.getRandomSeedPort(), 2);
Holder holder = setupAndStartServers(server1, server2);
int port1 = openPorts.get(0);
int port2 = openPorts.get(1);
AbstractClientConnectionFactory client1 = new TcpNetClientConnectionFactory("localhost", port1);
AbstractClientConnectionFactory client2 = new TcpNetClientConnectionFactory("localhost", port2);
AbstractServerConnectionFactory server1 = new TcpNetServerConnectionFactory(port1);
AbstractServerConnectionFactory server2 = new TcpNetServerConnectionFactory(port2);
testRealGuts(client1, client2, server1, server2);
AbstractClientConnectionFactory client1 = new TcpNetClientConnectionFactory("localhost", server1.getPort());
AbstractClientConnectionFactory client2 = new TcpNetClientConnectionFactory("localhost", server2.getPort());
testRealGuts(client1, client2, holder);
}
@Test
public void testRealNio() throws Exception {
final List<Integer> openPorts = SocketUtils.findAvailableServerSockets(SocketUtils.getRandomSeedPort(), 2);
AbstractServerConnectionFactory server1 = new TcpNioServerConnectionFactory(0);
AbstractServerConnectionFactory server2 = new TcpNioServerConnectionFactory(0);
int port1 = openPorts.get(0);
int port2 = openPorts.get(1);
Holder holder = setupAndStartServers(server1, server2);
AbstractClientConnectionFactory client1 = new TcpNioClientConnectionFactory("localhost", port1);
AbstractClientConnectionFactory client2 = new TcpNioClientConnectionFactory("localhost", port2);
AbstractServerConnectionFactory server1 = new TcpNioServerConnectionFactory(port1);
AbstractServerConnectionFactory server2 = new TcpNioServerConnectionFactory(port2);
testRealGuts(client1, client2, server1, server2);
AbstractClientConnectionFactory client1 = new TcpNioClientConnectionFactory("localhost", server1.getPort());
AbstractClientConnectionFactory client2 = new TcpNioClientConnectionFactory("localhost", server2.getPort());
testRealGuts(client1, client2, holder);
}
@Test
public void testRealNetSingleUse() throws Exception {
final List<Integer> openPorts = SocketUtils.findAvailableServerSockets(SocketUtils.getRandomSeedPort(), 2);
AbstractServerConnectionFactory server1 = new TcpNetServerConnectionFactory(0);
AbstractServerConnectionFactory server2 = new TcpNetServerConnectionFactory(0);
int port1 = openPorts.get(0);
int port2 = openPorts.get(1);
Holder holder = setupAndStartServers(server1, server2);
AbstractClientConnectionFactory client1 = new TcpNetClientConnectionFactory("localhost", port1);
AbstractClientConnectionFactory client2 = new TcpNetClientConnectionFactory("localhost", port2);
AbstractServerConnectionFactory server1 = new TcpNetServerConnectionFactory(port1);
AbstractServerConnectionFactory server2 = new TcpNetServerConnectionFactory(port2);
AbstractClientConnectionFactory client1 = new TcpNetClientConnectionFactory("localhost", server1.getPort());
AbstractClientConnectionFactory client2 = new TcpNetClientConnectionFactory("localhost", server2.getPort());
client1.setSingleUse(true);
client2.setSingleUse(true);
testRealGuts(client1, client2, server1, server2);
testRealGuts(client1, client2, holder);
}
@Test
public void testRealNioSingleUse() throws Exception {
final List<Integer> openPorts = SocketUtils.findAvailableServerSockets(SocketUtils.getRandomSeedPort(), 2);
AbstractServerConnectionFactory server1 = new TcpNioServerConnectionFactory(0);
AbstractServerConnectionFactory server2 = new TcpNioServerConnectionFactory(0);
int port1 = openPorts.get(0);
int port2 = openPorts.get(1);
Holder holder = setupAndStartServers(server1, server2);
AbstractClientConnectionFactory client1 = new TcpNioClientConnectionFactory("localhost", port1);
AbstractClientConnectionFactory client2 = new TcpNioClientConnectionFactory("localhost", port2);
AbstractServerConnectionFactory server1 = new TcpNioServerConnectionFactory(port1);
AbstractServerConnectionFactory server2 = new TcpNioServerConnectionFactory(port2);
AbstractClientConnectionFactory client1 = new TcpNioClientConnectionFactory("localhost", server1.getPort());
AbstractClientConnectionFactory client2 = new TcpNioClientConnectionFactory("localhost", server2.getPort());
client1.setSingleUse(true);
client2.setSingleUse(true);
testRealGuts(client1, client2, server1, server2);
testRealGuts(client1, client2, holder);
}
@Test
@@ -550,34 +554,64 @@ public class FailoverClientConnectionFactoryTests {
}
private void testRealGuts(AbstractClientConnectionFactory client1, AbstractClientConnectionFactory client2,
AbstractServerConnectionFactory server1, AbstractServerConnectionFactory server2) throws Exception {
Holder holder) throws Exception {
int port1 = 0;
int port2 = 0;
Executor exec = Executors.newCachedThreadPool();
client1.setTaskExecutor(exec);
client2.setTaskExecutor(exec);
server1.setTaskExecutor(exec);
server2.setTaskExecutor(exec);
client1.setTaskExecutor(holder.exec);
client2.setTaskExecutor(holder.exec);
client1.setBeanName("client1");
client2.setBeanName("client2");
client1.setApplicationEventPublisher(NULL_PUBLISHER);
client2.setApplicationEventPublisher(NULL_PUBLISHER);
List<AbstractClientConnectionFactory> factories = new ArrayList<AbstractClientConnectionFactory>();
factories.add(client1);
factories.add(client2);
FailoverClientConnectionFactory failFactory = new FailoverClientConnectionFactory(factories);
boolean singleUse = client1.isSingleUse();
failFactory.setSingleUse(singleUse);
failFactory.setBeanFactory(mock(BeanFactory.class));
failFactory.afterPropertiesSet();
TcpOutboundGateway outGateway = new TcpOutboundGateway();
outGateway.setConnectionFactory(failFactory);
outGateway.start();
QueueChannel replyChannel = new QueueChannel();
outGateway.setReplyChannel(replyChannel);
Message<String> message = new GenericMessage<String>("foo");
outGateway.setRemoteTimeout(120000);
outGateway.handleMessage(message);
Socket socket = null;
if (!singleUse) {
socket = getSocket(client1);
port1 = socket.getLocalPort();
}
assertTrue(singleUse | holder.connectionId.get().contains(Integer.toString(port1)));
Message<?> replyMessage = replyChannel.receive(10000);
assertNotNull(replyMessage);
holder.server1.stop();
TestingUtilities.waitStopListening(holder.server1, 10000L);
TestingUtilities.waitUntilFactoryHasThisNumberOfConnections(client1, 0);
outGateway.handleMessage(message);
if (!singleUse) {
socket = getSocket(client2);
port2 = socket.getLocalPort();
}
assertTrue(singleUse | holder.connectionId.get().contains(Integer.toString(port2)));
replyMessage = replyChannel.receive(10000);
assertNotNull(replyMessage);
holder.gateway2.stop();
outGateway.stop();
}
private Holder setupAndStartServers(AbstractServerConnectionFactory server1,
AbstractServerConnectionFactory server2) throws Exception {
Executor exec = Executors.newCachedThreadPool();
server1.setTaskExecutor(exec);
server2.setTaskExecutor(exec);
server1.setBeanName("server1");
server2.setBeanName("server2");
ApplicationEventPublisher pub = new ApplicationEventPublisher() {
@Override
public void publishEvent(ApplicationEvent event) {
}
@Override
public void publishEvent(Object event) {
}
};
client1.setApplicationEventPublisher(pub);
client2.setApplicationEventPublisher(pub);
server1.setApplicationEventPublisher(pub);
server2.setApplicationEventPublisher(pub);
server1.setApplicationEventPublisher(NULL_PUBLISHER);
server2.setApplicationEventPublisher(NULL_PUBLISHER);
TcpInboundGateway gateway1 = new TcpInboundGateway();
gateway1.setConnectionFactory(server1);
SubscribableChannel channel = new DirectChannel();
@@ -601,43 +635,12 @@ public class FailoverClientConnectionFactoryTests {
gateway2.start();
TestingUtilities.waitListening(server1, null);
TestingUtilities.waitListening(server2, null);
List<AbstractClientConnectionFactory> factories = new ArrayList<AbstractClientConnectionFactory>();
factories.add(client1);
factories.add(client2);
FailoverClientConnectionFactory failFactory = new FailoverClientConnectionFactory(factories);
boolean singleUse = client1.isSingleUse();
failFactory.setSingleUse(singleUse);
failFactory.setBeanFactory(mock(BeanFactory.class));
failFactory.afterPropertiesSet();
TcpOutboundGateway outGateway = new TcpOutboundGateway();
outGateway.setConnectionFactory(failFactory);
outGateway.start();
QueueChannel replyChannel = new QueueChannel();
outGateway.setReplyChannel(replyChannel);
Message<String> message = new GenericMessage<String>("foo");
outGateway.setRemoteTimeout(120000);
outGateway.handleMessage(message);
Socket socket = null;
if (!singleUse) {
socket = getSocket(client1);
port1 = socket.getLocalPort();
}
assertTrue(singleUse | connectionId.get().contains(Integer.toString(port1)));
Message<?> replyMessage = replyChannel.receive(10000);
assertNotNull(replyMessage);
server1.stop();
TestingUtilities.waitStopListening(server1, 10000L);
TestingUtilities.waitUntilFactoryHasThisNumberOfConnections(client1, 0);
outGateway.handleMessage(message);
if (!singleUse) {
socket = getSocket(client2);
port2 = socket.getLocalPort();
}
assertTrue(singleUse | connectionId.get().contains(Integer.toString(port2)));
replyMessage = replyChannel.receive(10000);
assertNotNull(replyMessage);
gateway2.stop();
outGateway.stop();
Holder holder = new Holder();
holder.exec = exec;
holder.connectionId = connectionId;
holder.server1 = server1;
holder.gateway2 = gateway2;
return holder;
}
private Socket getSocket(AbstractClientConnectionFactory client) throws Exception {
@@ -650,5 +653,17 @@ public class FailoverClientConnectionFactoryTests {
}
private static class Holder {
private AtomicReference<String> connectionId;
private Executor exec;
private TcpInboundGateway gateway2;
private AbstractServerConnectionFactory server1;
}
}

View File

@@ -26,7 +26,6 @@ import org.junit.Test;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.ip.util.SocketTestUtils;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.test.util.SocketUtils;
import org.springframework.messaging.Message;
@@ -53,7 +52,7 @@ public class MultiClientTests {
public void testNoAck() throws Exception {
final String payload = largePayload(1000);
final UnicastReceivingChannelAdapter adapter =
new UnicastReceivingChannelAdapter(SocketUtils.findAvailableUdpSocket());
new UnicastReceivingChannelAdapter(0);
int drivers = 10;
adapter.setPoolSize(drivers);
QueueChannel queue = new QueueChannel(drivers * 3);
@@ -101,7 +100,7 @@ public class MultiClientTests {
public void testAck() throws Exception {
final String payload = largePayload(1000);
final UnicastReceivingChannelAdapter adapter =
new UnicastReceivingChannelAdapter(SocketUtils.findAvailableUdpSocket(), false);
new UnicastReceivingChannelAdapter(0, false);
int drivers = 5;
adapter.setPoolSize(drivers);
QueueChannel queue = new QueueChannel(drivers * 3);
@@ -113,14 +112,12 @@ public class MultiClientTests {
final AtomicBoolean done = new AtomicBoolean();
for (int i = 0; i < drivers; i++) {
final int j = i;
Thread t = new Thread( new Runnable() {
@Override
public void run() {
UnicastSendingMessageHandler sender = new UnicastSendingMessageHandler(
"localhost", adapter.getPort(),
false, true, "localhost",
SocketUtils.findAvailableUdpSocket(adapter.getPort() + j + 1000),
false, true, "localhost", 0,
10000);
sender.start();
while (true) {
@@ -153,7 +150,7 @@ public class MultiClientTests {
public void testAckWithLength() throws Exception {
final String payload = largePayload(1000);
final UnicastReceivingChannelAdapter adapter =
new UnicastReceivingChannelAdapter(SocketUtils.findAvailableUdpSocket(), true);
new UnicastReceivingChannelAdapter(0, true);
int drivers = 10;
adapter.setPoolSize(drivers);
QueueChannel queue = new QueueChannel(drivers * 3);
@@ -171,8 +168,7 @@ public class MultiClientTests {
public void run() {
UnicastSendingMessageHandler sender = new UnicastSendingMessageHandler(
"localhost", adapter.getPort(),
true, true, "localhost",
SocketUtils.findAvailableUdpSocket(adapter.getPort() + j + 1100),
true, true, "localhost", 0,
10000);
sender.start();
while (true) {