TCP/UDP: Use OS to Select Test Ports
https://build.spring.io/browse/INT-AT42SIO-228/
This commit is contained in:
committed by
Artem Bilan
parent
d00e0cdcc1
commit
2a66f39df1
@@ -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) {
|
||||
|
||||
@@ -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(0, 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(0, 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(0, 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(0, 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 {
|
||||
@@ -649,5 +652,17 @@ public class FailoverClientConnectionFactoryTests {
|
||||
}
|
||||
}
|
||||
|
||||
private static class Holder {
|
||||
|
||||
private AtomicReference<String> connectionId;
|
||||
|
||||
private Executor exec;
|
||||
|
||||
private TcpInboundGateway gateway2;
|
||||
|
||||
private AbstractServerConnectionFactory server1;
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -27,7 +27,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;
|
||||
|
||||
|
||||
@@ -54,7 +53,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);
|
||||
@@ -103,7 +102,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);
|
||||
@@ -115,14 +114,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) {
|
||||
@@ -156,7 +153,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);
|
||||
@@ -174,8 +171,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) {
|
||||
|
||||
@@ -48,7 +48,6 @@ import org.springframework.integration.handler.ServiceActivatingHandler;
|
||||
import org.springframework.integration.ip.IpHeaders;
|
||||
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;
|
||||
import org.springframework.messaging.SubscribableChannel;
|
||||
|
||||
@@ -151,7 +150,7 @@ public class UdpChannelAdapterTests {
|
||||
DatagramPacketMessageMapper mapper = new DatagramPacketMessageMapper();
|
||||
DatagramPacket packet = mapper.fromMessage(message);
|
||||
packet.setSocketAddress(new InetSocketAddress("localhost", port));
|
||||
DatagramSocket datagramSocket = new DatagramSocket(SocketUtils.findAvailableUdpSocket());
|
||||
DatagramSocket datagramSocket = new DatagramSocket(0);
|
||||
datagramSocket.send(packet);
|
||||
datagramSocket.close();
|
||||
@SuppressWarnings("unchecked")
|
||||
@@ -180,7 +179,7 @@ public class UdpChannelAdapterTests {
|
||||
DatagramPacketMessageMapper mapper = new DatagramPacketMessageMapper();
|
||||
DatagramPacket packet = mapper.fromMessage(message);
|
||||
packet.setSocketAddress(new InetSocketAddress("localhost", port));
|
||||
final DatagramSocket socket = new DatagramSocket(SocketUtils.findAvailableUdpSocket());
|
||||
final DatagramSocket socket = new DatagramSocket(0);
|
||||
socket.send(packet);
|
||||
final AtomicReference<DatagramPacket> theAnswer = new AtomicReference<DatagramPacket>();
|
||||
final CountDownLatch receiverReadyLatch = new CountDownLatch(1);
|
||||
@@ -238,7 +237,7 @@ public class UdpChannelAdapterTests {
|
||||
"localhost", port, false, true,
|
||||
"localhost",
|
||||
// whichNic,
|
||||
SocketUtils.findAvailableUdpSocket(), 5000);
|
||||
0, 5000);
|
||||
// handler.setLocalAddress(whichNic);
|
||||
handler.setBeanFactory(mock(BeanFactory.class));
|
||||
handler.afterPropertiesSet();
|
||||
@@ -282,9 +281,8 @@ public class UdpChannelAdapterTests {
|
||||
@Test
|
||||
public void testMulticastSender() throws Exception {
|
||||
QueueChannel channel = new QueueChannel(2);
|
||||
int port = SocketUtils.findAvailableUdpSocket();
|
||||
UnicastReceivingChannelAdapter adapter =
|
||||
new MulticastReceivingChannelAdapter(this.multicastRule.getGroup(), port);
|
||||
new MulticastReceivingChannelAdapter(this.multicastRule.getGroup(), 0);
|
||||
adapter.setOutputChannel(channel);
|
||||
String nic = this.multicastRule.getNic();
|
||||
adapter.setLocalAddress(nic);
|
||||
@@ -292,7 +290,7 @@ public class UdpChannelAdapterTests {
|
||||
SocketTestUtils.waitListening(adapter);
|
||||
|
||||
MulticastSendingMessageHandler handler =
|
||||
new MulticastSendingMessageHandler(this.multicastRule.getGroup(), port);
|
||||
new MulticastSendingMessageHandler(this.multicastRule.getGroup(), adapter.getPort());
|
||||
handler.setLocalAddress(nic);
|
||||
Message<byte[]> message = MessageBuilder.withPayload("ABCD".getBytes()).build();
|
||||
handler.handleMessage(message);
|
||||
@@ -322,7 +320,7 @@ public class UdpChannelAdapterTests {
|
||||
DatagramPacketMessageMapper mapper = new DatagramPacketMessageMapper();
|
||||
DatagramPacket packet = mapper.fromMessage(message);
|
||||
packet.setSocketAddress(new InetSocketAddress("localhost", port));
|
||||
DatagramSocket datagramSocket = new DatagramSocket(SocketUtils.findAvailableUdpSocket());
|
||||
DatagramSocket datagramSocket = new DatagramSocket(0);
|
||||
datagramSocket.send(packet);
|
||||
datagramSocket.close();
|
||||
Message<?> receivedMessage = errorChannel.receive(2000);
|
||||
|
||||
Reference in New Issue
Block a user